ContentCacheManager.js 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441
  1. /**
  2. * ContentCacheManager - maintains a cache of stream hashes from a PDF file
  3. */
  4. const { callbackify } = require('util')
  5. const fs = require('fs')
  6. const crypto = require('crypto')
  7. const Path = require('path')
  8. const Settings = require('@overleaf/settings')
  9. const OError = require('@overleaf/o-error')
  10. const pLimit = require('p-limit')
  11. const { parseXrefTable } = require('./XrefParser')
  12. const {
  13. QueueLimitReachedError,
  14. TimedOutError,
  15. NoXrefTableError,
  16. } = require('./Errors')
  17. const workerpool = require('workerpool')
  18. const Metrics = require('@overleaf/metrics')
  19. /**
  20. * @type {import('workerpool').WorkerPool}
  21. */
  22. let WORKER_POOL
  23. // NOTE: Check for main thread to avoid recursive start of pool.
  24. if (Settings.pdfCachingEnableWorkerPool && workerpool.isMainThread) {
  25. WORKER_POOL = workerpool.pool(Path.join(__dirname, 'ContentCacheWorker.js'), {
  26. // Cap number of worker threads.
  27. maxWorkers: Settings.pdfCachingWorkerPoolSize,
  28. // Warmup workers.
  29. minWorkers: Settings.pdfCachingWorkerPoolSize,
  30. // Limit queue back-log
  31. maxQueueSize: Settings.pdfCachingWorkerPoolBackLogLimit,
  32. })
  33. setInterval(() => {
  34. const {
  35. totalWorkers,
  36. busyWorkers,
  37. idleWorkers,
  38. pendingTasks,
  39. activeTasks,
  40. } = WORKER_POOL.stats()
  41. Metrics.gauge('pdf_caching_total_workers', totalWorkers)
  42. Metrics.gauge('pdf_caching_busy_workers', busyWorkers)
  43. Metrics.gauge('pdf_caching_idle_workers', idleWorkers)
  44. Metrics.gauge('pdf_caching_pending_tasks', pendingTasks)
  45. Metrics.gauge('pdf_caching_active_tasks', activeTasks)
  46. }, 15 * 1000)
  47. }
  48. /**
  49. *
  50. * @param {String} contentDir path to directory where content hash files are cached
  51. * @param {String} filePath the pdf file to scan for streams
  52. * @param {number} pdfSize the pdf size
  53. * @param {number} pdfCachingMinChunkSize per request threshold
  54. * @param {number} compileTime
  55. */
  56. async function update({
  57. contentDir,
  58. filePath,
  59. pdfSize,
  60. pdfCachingMinChunkSize,
  61. compileTime,
  62. }) {
  63. if (pdfSize < pdfCachingMinChunkSize) {
  64. return {
  65. contentRanges: [],
  66. newContentRanges: [],
  67. reclaimedSpace: 0,
  68. startXRefTable: undefined,
  69. }
  70. }
  71. if (Settings.pdfCachingEnableWorkerPool) {
  72. return await updateOtherEventLoop({
  73. contentDir,
  74. filePath,
  75. pdfSize,
  76. pdfCachingMinChunkSize,
  77. compileTime,
  78. })
  79. } else {
  80. return await updateSameEventLoop({
  81. contentDir,
  82. filePath,
  83. pdfSize,
  84. pdfCachingMinChunkSize,
  85. compileTime,
  86. })
  87. }
  88. }
  89. /**
  90. *
  91. * @param {String} contentDir path to directory where content hash files are cached
  92. * @param {String} filePath the pdf file to scan for streams
  93. * @param {number} pdfSize the pdf size
  94. * @param {number} pdfCachingMinChunkSize per request threshold
  95. * @param {number} compileTime
  96. */
  97. async function updateOtherEventLoop({
  98. contentDir,
  99. filePath,
  100. pdfSize,
  101. pdfCachingMinChunkSize,
  102. compileTime,
  103. }) {
  104. const workerLatencyInMs = 100
  105. // Prefer getting the timeout error from the worker vs timing out the worker.
  106. const timeout = getMaxOverhead(compileTime) + workerLatencyInMs
  107. try {
  108. return await WORKER_POOL.exec('updateSameEventLoop', [
  109. {
  110. contentDir,
  111. filePath,
  112. pdfSize,
  113. pdfCachingMinChunkSize,
  114. compileTime,
  115. },
  116. ]).timeout(timeout)
  117. } catch (e) {
  118. if (e instanceof workerpool.Promise.TimeoutError) {
  119. throw new TimedOutError('context-lost-in-worker', { timeout })
  120. }
  121. if (e.message?.includes?.('Max queue size of ')) {
  122. throw new QueueLimitReachedError()
  123. }
  124. if (e.message?.includes?.('xref')) {
  125. throw new NoXrefTableError(e.message)
  126. }
  127. throw e
  128. }
  129. }
  130. /**
  131. *
  132. * @param {String} contentDir path to directory where content hash files are cached
  133. * @param {String} filePath the pdf file to scan for streams
  134. * @param {number} pdfSize the pdf size
  135. * @param {number} pdfCachingMinChunkSize per request threshold
  136. * @param {number} compileTime
  137. */
  138. async function updateSameEventLoop({
  139. contentDir,
  140. filePath,
  141. pdfSize,
  142. pdfCachingMinChunkSize,
  143. compileTime,
  144. }) {
  145. const checkDeadline = getDeadlineChecker(compileTime)
  146. // keep track of hashes expire old ones when they reach a generation > N.
  147. const tracker = await HashFileTracker.from(contentDir)
  148. tracker.updateAge()
  149. checkDeadline('after init HashFileTracker')
  150. const [reclaimedSpace, overheadDeleteStaleHashes] =
  151. await tracker.deleteStaleHashes(5)
  152. checkDeadline('after delete stale hashes')
  153. const { xRefEntries, startXRefTable } = await parseXrefTable(
  154. filePath,
  155. pdfSize
  156. )
  157. xRefEntries.sort((a, b) => {
  158. return a.offset - b.offset
  159. })
  160. xRefEntries.forEach((obj, idx) => {
  161. obj.idx = idx
  162. })
  163. checkDeadline('after parsing')
  164. const uncompressedObjects = []
  165. for (const object of xRefEntries) {
  166. if (!object.uncompressed) {
  167. continue
  168. }
  169. const nextObject = xRefEntries[object.idx + 1]
  170. if (!nextObject) {
  171. // Ignore this possible edge case.
  172. // The last object should be part of the xRef table.
  173. continue
  174. } else {
  175. object.endOffset = nextObject.offset
  176. }
  177. const size = object.endOffset - object.offset
  178. object.size = size
  179. if (size < pdfCachingMinChunkSize) {
  180. continue
  181. }
  182. uncompressedObjects.push({ object, idx: uncompressedObjects.length })
  183. }
  184. checkDeadline('after finding uncompressed')
  185. let timedOutErr = null
  186. const contentRanges = []
  187. const newContentRanges = []
  188. const handle = await fs.promises.open(filePath)
  189. try {
  190. for (const { object, idx } of uncompressedObjects) {
  191. let buffer = Buffer.alloc(object.size, 0)
  192. const { bytesRead } = await handle.read(
  193. buffer,
  194. 0,
  195. object.size,
  196. object.offset
  197. )
  198. checkDeadline('after read ' + idx)
  199. if (bytesRead !== object.size) {
  200. throw new OError('could not read full chunk', {
  201. object,
  202. bytesRead,
  203. })
  204. }
  205. const idxObj = buffer.indexOf('obj')
  206. if (idxObj > 100) {
  207. throw new OError('objectId is too large', {
  208. object,
  209. idxObj,
  210. })
  211. }
  212. const objectIdRaw = buffer.subarray(0, idxObj)
  213. buffer = buffer.subarray(objectIdRaw.byteLength)
  214. const hash = pdfStreamHash(buffer)
  215. checkDeadline('after hash ' + idx)
  216. const range = {
  217. objectId: objectIdRaw.toString(),
  218. start: object.offset + objectIdRaw.byteLength,
  219. end: object.endOffset,
  220. hash,
  221. }
  222. if (tracker.has(range.hash)) {
  223. // Optimization: Skip writing of already seen hashes.
  224. tracker.track(range)
  225. contentRanges.push(range)
  226. continue
  227. }
  228. await writePdfStream(contentDir, hash, buffer)
  229. tracker.track(range)
  230. contentRanges.push(range)
  231. newContentRanges.push(range)
  232. checkDeadline('after write ' + idx)
  233. }
  234. } catch (err) {
  235. if (err instanceof TimedOutError) {
  236. // Let the frontend use ranges that were processed so far.
  237. timedOutErr = err
  238. } else {
  239. throw err
  240. }
  241. } finally {
  242. await handle.close()
  243. // Flush from both success and failure code path. This allows the next
  244. // cycle to complete faster as it can use the already written ranges.
  245. await tracker.flush()
  246. }
  247. return {
  248. contentRanges,
  249. newContentRanges,
  250. reclaimedSpace,
  251. startXRefTable,
  252. overheadDeleteStaleHashes,
  253. timedOutErr,
  254. }
  255. }
  256. function getStatePath(contentDir) {
  257. return Path.join(contentDir, '.state.v0.json')
  258. }
  259. class HashFileTracker {
  260. constructor(contentDir, { hashAge = [], hashSize = [] }) {
  261. this.contentDir = contentDir
  262. this.hashAge = new Map(hashAge)
  263. this.hashSize = new Map(hashSize)
  264. }
  265. static async from(contentDir) {
  266. const statePath = getStatePath(contentDir)
  267. let state = {}
  268. try {
  269. const blob = await fs.promises.readFile(statePath)
  270. state = JSON.parse(blob)
  271. } catch (e) {}
  272. return new HashFileTracker(contentDir, state)
  273. }
  274. has(hash) {
  275. return this.hashAge.has(hash)
  276. }
  277. track(range) {
  278. if (!this.hashSize.has(range.hash)) {
  279. this.hashSize.set(range.hash, range.end - range.start)
  280. }
  281. this.hashAge.set(range.hash, 0)
  282. }
  283. updateAge() {
  284. for (const [hash, age] of this.hashAge) {
  285. this.hashAge.set(hash, age + 1)
  286. }
  287. return this
  288. }
  289. findStale(maxAge) {
  290. const stale = []
  291. for (const [hash, age] of this.hashAge) {
  292. if (age > maxAge) {
  293. stale.push(hash)
  294. }
  295. }
  296. return stale
  297. }
  298. async flush() {
  299. const statePath = getStatePath(this.contentDir)
  300. const blob = JSON.stringify({
  301. hashAge: Array.from(this.hashAge.entries()),
  302. hashSize: Array.from(this.hashSize.entries()),
  303. })
  304. const atomicWrite = statePath + '~'
  305. try {
  306. await fs.promises.writeFile(atomicWrite, blob)
  307. } catch (err) {
  308. try {
  309. await fs.promises.unlink(atomicWrite)
  310. } catch (e) {}
  311. throw err
  312. }
  313. try {
  314. await fs.promises.rename(atomicWrite, statePath)
  315. } catch (err) {
  316. try {
  317. await fs.promises.unlink(atomicWrite)
  318. } catch (e) {}
  319. throw err
  320. }
  321. }
  322. async deleteStaleHashes(n) {
  323. const t0 = Date.now()
  324. // delete any hash file older than N generations
  325. const hashes = this.findStale(n)
  326. let reclaimedSpace = 0
  327. if (hashes.length === 0) {
  328. return [reclaimedSpace, Date.now() - t0]
  329. }
  330. await promiseMapWithLimit(10, hashes, async hash => {
  331. try {
  332. await fs.promises.unlink(Path.join(this.contentDir, hash))
  333. } catch (err) {
  334. if (err?.code === 'ENOENT') {
  335. // Ignore already deleted entries. The previous cleanup cycle may have
  336. // been killed halfway through the deletion process, or before we
  337. // flushed the state to disk.
  338. } else {
  339. throw err
  340. }
  341. }
  342. this.hashAge.delete(hash)
  343. reclaimedSpace += this.hashSize.get(hash)
  344. this.hashSize.delete(hash)
  345. })
  346. return [reclaimedSpace, Date.now() - t0]
  347. }
  348. }
  349. function pdfStreamHash(buffer) {
  350. const hash = crypto.createHash('sha256')
  351. hash.update(buffer)
  352. return hash.digest('hex')
  353. }
  354. async function writePdfStream(dir, hash, buffer) {
  355. const filename = Path.join(dir, hash)
  356. const atomicWriteFilename = filename + '~'
  357. try {
  358. await fs.promises.writeFile(atomicWriteFilename, buffer)
  359. await fs.promises.rename(atomicWriteFilename, filename)
  360. } catch (err) {
  361. try {
  362. await fs.promises.unlink(atomicWriteFilename)
  363. } catch (_) {
  364. throw err
  365. }
  366. }
  367. }
  368. function getMaxOverhead(compileTime) {
  369. return Math.min(
  370. // Adding 10s to a 40s compile time is OK.
  371. // Adding 1s to a 3s compile time is OK.
  372. Math.max(compileTime / 4, 1000),
  373. // Adding 30s to a 120s compile time is not OK, limit to 10s.
  374. Settings.pdfCachingMaxProcessingTime
  375. )
  376. }
  377. function getDeadlineChecker(compileTime) {
  378. const timeout = getMaxOverhead(compileTime)
  379. const deadline = Date.now() + timeout
  380. let lastStage = { stage: 'start', now: Date.now() }
  381. let completedStages = 0
  382. return function (stage) {
  383. const now = Date.now()
  384. if (now > deadline) {
  385. throw new TimedOutError(stage, {
  386. timeout,
  387. completedStages,
  388. lastStage: lastStage.stage,
  389. diffToLastStage: now - lastStage.now,
  390. })
  391. }
  392. completedStages++
  393. lastStage = { stage, now }
  394. }
  395. }
  396. function promiseMapWithLimit(concurrency, array, fn) {
  397. const limit = pLimit(concurrency)
  398. return Promise.all(array.map(x => limit(() => fn(x))))
  399. }
  400. module.exports = {
  401. HASH_REGEX: /^[0-9a-f]{64}$/,
  402. update: callbackify(update),
  403. promises: {
  404. update,
  405. updateSameEventLoop,
  406. },
  407. }