projects.js 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590
  1. 'use strict'
  2. const _ = require('lodash')
  3. const Path = require('node:path')
  4. const Stream = require('node:stream')
  5. const HTTPStatus = require('http-status')
  6. const fs = require('node:fs')
  7. const { promisify } = require('node:util')
  8. const config = require('config')
  9. const OError = require('@overleaf/o-error')
  10. const { expressify } = require('@overleaf/promise-utils')
  11. const { parseReq } = require('@overleaf/validation-tools')
  12. const logger = require('@overleaf/logger')
  13. const { Chunk, ChunkResponse, Blob } = require('overleaf-editor-core')
  14. const {
  15. BlobStore,
  16. BatchBlobStore,
  17. blobHash,
  18. chunkStore,
  19. redisBuffer,
  20. HashCheckBlobStore,
  21. ProjectArchive,
  22. zipStore,
  23. persistBuffer,
  24. } = require('../../storage')
  25. const render = require('./render')
  26. const schemas = require('../schema')
  27. const withTmpDir = require('./with_tmp_dir')
  28. const StreamSizeLimit = require('./stream_size_limit')
  29. const { getProjectBlobsBatch } = require('../../storage/lib/blob_store')
  30. const assert = require('../../storage/lib/assert')
  31. const { getChunkMetadataForVersion } = require('../../storage/lib/chunk_store')
  32. const { IncrementalResponse } = require('@overleaf/stream-utils')
  33. const pipeline = promisify(Stream.pipeline)
  34. async function initializeProject(req, res, next) {
  35. const { body } = parseReq(req, schemas.initializeProject)
  36. let projectId = body?.projectId
  37. try {
  38. projectId = await chunkStore.initializeProject(projectId)
  39. res.status(HTTPStatus.OK).json({ projectId })
  40. } catch (err) {
  41. if (err instanceof chunkStore.AlreadyInitialized) {
  42. logger.warn({ err, projectId }, 'failed to initialize')
  43. render.conflict(res)
  44. } else {
  45. throw err
  46. }
  47. }
  48. }
  49. async function cloneProject(req, res) {
  50. const {
  51. body: { targetProjectId },
  52. params: { project_id: sourceProjectId },
  53. } = parseReq(req, schemas.cloneProject)
  54. const incrResp = new IncrementalResponse({
  55. res,
  56. timeout: 10 * 60_000 - 5_000,
  57. logger,
  58. label: 'clone history in history-v1',
  59. info: { targetProjectId, sourceProjectId },
  60. })
  61. const signal = incrResp.signal()
  62. try {
  63. try {
  64. // Use the same limits importChanges, since these are passed to persistChanges
  65. const farFuture = new Date()
  66. farFuture.setTime(farFuture.getTime() + 7 * 24 * 3600 * 1000)
  67. const limits = {
  68. maxChanges: 0,
  69. minChangeTimestamp: farFuture,
  70. maxChangeTimestamp: farFuture,
  71. autoResync: true,
  72. }
  73. incrResp.sendUpdate('flushing redis buffer: pending')
  74. await persistBuffer(sourceProjectId, limits)
  75. incrResp.sendUpdate('flushing redis buffer: done')
  76. } catch (err) {
  77. incrResp.sendUpdate('failed to flush redis buffer')
  78. logger.error(
  79. { err, targetProjectId, sourceProjectId },
  80. 'failed to persist buffer during clone'
  81. )
  82. }
  83. await chunkStore.cloneProject(
  84. sourceProjectId,
  85. targetProjectId,
  86. progress => {
  87. if (signal.aborted) return
  88. incrResp.sendUpdate(progress)
  89. },
  90. signal
  91. )
  92. if (!signal.aborted) {
  93. incrResp.sendUpdate('cloning full project history data: done')
  94. }
  95. } catch (err) {
  96. incrResp.fail(err)
  97. } finally {
  98. incrResp.end()
  99. }
  100. }
  101. async function getLatestContent(req, res, next) {
  102. const { params } = parseReq(req, schemas.getLatestContent)
  103. const projectId = params.project_id
  104. const blobStore = new BlobStore(projectId)
  105. const chunk = await chunkStore.loadLatest(projectId)
  106. const snapshot = chunk.getSnapshot()
  107. snapshot.applyAll(chunk.getChanges())
  108. await snapshot.loadFiles('eager', blobStore)
  109. res.json(snapshot.toRaw())
  110. }
  111. async function getContentAtVersion(req, res, next) {
  112. const { params } = parseReq(req, schemas.getContentAtVersion)
  113. const projectId = params.project_id
  114. const version = params.version
  115. const blobStore = new BlobStore(projectId)
  116. const snapshot = await getSnapshotAtVersion(projectId, version)
  117. await snapshot.loadFiles('eager', blobStore)
  118. res.json(snapshot.toRaw())
  119. }
  120. async function getLatestHashedContent(req, res, next) {
  121. const { params } = parseReq(req, schemas.getLatestHashedContent)
  122. const projectId = params.project_id
  123. const blobStore = new HashCheckBlobStore(new BlobStore(projectId))
  124. const chunk = await chunkStore.loadLatest(projectId)
  125. const snapshot = chunk.getSnapshot()
  126. snapshot.applyAll(chunk.getChanges())
  127. await snapshot.loadFiles('eager', blobStore)
  128. const rawSnapshot = await snapshot.store(blobStore)
  129. res.json(rawSnapshot)
  130. }
  131. async function getLatestHistory(req, res, next) {
  132. const { params } = parseReq(req, schemas.getLatestHistory)
  133. const projectId = params.project_id
  134. try {
  135. const chunk = await chunkStore.loadLatest(projectId)
  136. const chunkResponse = new ChunkResponse(chunk)
  137. res.json(chunkResponse.toRaw())
  138. } catch (err) {
  139. if (err instanceof Chunk.NotFoundError) {
  140. render.notFound(res)
  141. } else {
  142. throw err
  143. }
  144. }
  145. }
  146. async function getLatestHistoryRaw(req, res, next) {
  147. const { params, query } = parseReq(req, schemas.getLatestHistoryRaw)
  148. const projectId = params.project_id
  149. const readOnly = query.readOnly
  150. try {
  151. const { startVersion, endVersion, endTimestamp } =
  152. await chunkStore.getLatestChunkMetadata(projectId, { readOnly })
  153. res.json({
  154. startVersion,
  155. endVersion,
  156. endTimestamp,
  157. })
  158. } catch (err) {
  159. if (err instanceof Chunk.NotFoundError) {
  160. render.notFound(res)
  161. } else {
  162. throw err
  163. }
  164. }
  165. }
  166. async function getHistory(req, res, next) {
  167. const { params } = parseReq(req, schemas.getHistory)
  168. const projectId = params.project_id
  169. const version = params.version
  170. try {
  171. const chunk = await chunkStore.loadAtVersion(projectId, version)
  172. const chunkResponse = new ChunkResponse(chunk)
  173. res.json(chunkResponse.toRaw())
  174. } catch (err) {
  175. if (err instanceof Chunk.NotFoundError) {
  176. render.notFound(res)
  177. } else {
  178. throw err
  179. }
  180. }
  181. }
  182. async function getHistoryBefore(req, res, next) {
  183. const { params } = parseReq(req, schemas.getHistoryBefore)
  184. const projectId = params.project_id
  185. const timestamp = params.timestamp
  186. try {
  187. const chunk = await chunkStore.loadAtTimestamp(projectId, timestamp)
  188. const chunkResponse = new ChunkResponse(chunk)
  189. res.json(chunkResponse.toRaw())
  190. } catch (err) {
  191. if (err instanceof Chunk.NotFoundError) {
  192. render.notFound(res)
  193. } else {
  194. throw err
  195. }
  196. }
  197. }
  198. /**
  199. * Get all changes since the beginning of history or since a given version
  200. */
  201. async function getChanges(req, res, next) {
  202. const { params, query } = parseReq(req, schemas.getChanges)
  203. const projectId = params.project_id
  204. const sinceParam = query.since
  205. const since = sinceParam == null ? 0 : sinceParam
  206. if (since < 0) {
  207. // Negative values would cause an infinite loop
  208. return res.status(400).json({
  209. error: `Version out of bounds: ${since}`,
  210. })
  211. }
  212. try {
  213. const { changes, hasMore } = await chunkStore.getChangesSinceVersion(
  214. projectId,
  215. since
  216. )
  217. res.json({ changes: changes.map(change => change.toRaw()), hasMore })
  218. } catch (err) {
  219. if (err instanceof Chunk.VersionNotFoundError) {
  220. return res.status(400).json({
  221. error: `Version out of bounds: ${since}`,
  222. })
  223. }
  224. throw err
  225. }
  226. }
  227. async function getZip(req, res, next) {
  228. const { params } = parseReq(req, schemas.getZip)
  229. const projectId = params.project_id
  230. const version = params.version
  231. const blobStore = new BlobStore(projectId)
  232. let snapshot
  233. try {
  234. snapshot = await getSnapshotAtVersion(projectId, version)
  235. } catch (err) {
  236. if (err instanceof Chunk.NotFoundError) {
  237. return render.notFound(res)
  238. } else {
  239. throw err
  240. }
  241. }
  242. await withTmpDir('get-zip-', async tmpDir => {
  243. const tmpFilename = Path.join(tmpDir, 'project.zip')
  244. const archive = new ProjectArchive(snapshot)
  245. await archive.writeZip(blobStore, tmpFilename)
  246. res.set('Content-Type', 'application/octet-stream')
  247. res.set('Content-Disposition', 'attachment; filename=project.zip')
  248. const stream = fs.createReadStream(tmpFilename)
  249. await pipeline(stream, res)
  250. })
  251. }
  252. async function createZip(req, res, next) {
  253. const { params } = parseReq(req, schemas.createZip)
  254. const projectId = params.project_id
  255. const version = params.version
  256. try {
  257. const snapshot = await getSnapshotAtVersion(projectId, version)
  258. const zipUrl = await zipStore.getSignedUrl(projectId, version)
  259. // Do not await this; run it in the background.
  260. zipStore.storeZip(projectId, version, snapshot).catch(err => {
  261. logger.error({ err, projectId, version }, 'createZip: storeZip failed')
  262. })
  263. res.status(HTTPStatus.OK).json({ zipUrl })
  264. } catch (error) {
  265. if (error instanceof Chunk.NotFoundError) {
  266. render.notFound(res)
  267. } else {
  268. next(error)
  269. }
  270. }
  271. }
  272. async function deleteProject(req, res, next) {
  273. const { params } = parseReq(req, schemas.deleteProject)
  274. const projectId = params.project_id
  275. const blobStore = new BlobStore(projectId)
  276. await Promise.all([
  277. redisBuffer.hardDeleteProject(projectId),
  278. chunkStore.deleteProjectChunks(projectId),
  279. blobStore.deleteBlobs(),
  280. ])
  281. res.status(HTTPStatus.NO_CONTENT).send()
  282. }
  283. async function createProjectBlob(req, res, next) {
  284. const { params } = parseReq(req, schemas.createProjectBlob)
  285. const projectId = params.project_id
  286. const expectedHash = params.hash
  287. const maxUploadSize = parseInt(config.get('maxFileUploadSize'), 10)
  288. await withTmpDir('blob-', async tmpDir => {
  289. const tmpPath = Path.join(tmpDir, 'content')
  290. const sizeLimit = new StreamSizeLimit(maxUploadSize)
  291. await pipeline(req, sizeLimit, fs.createWriteStream(tmpPath))
  292. if (sizeLimit.sizeLimitExceeded) {
  293. logger.warn(
  294. { projectId, expectedHash, maxUploadSize },
  295. 'blob exceeds size threshold'
  296. )
  297. return render.requestEntityTooLarge(res)
  298. }
  299. const hash = await blobHash.fromFile(tmpPath)
  300. if (hash !== expectedHash) {
  301. logger.warn({ projectId, hash, expectedHash }, 'Hash mismatch')
  302. return render.conflict(res, 'File hash mismatch')
  303. }
  304. const blobStore = new BlobStore(projectId)
  305. const newBlob = await blobStore.putFile(tmpPath)
  306. if (config.has('backupStore')) {
  307. try {
  308. const { backupBlob } = await import('../../storage/lib/backupBlob.mjs')
  309. await backupBlob(projectId, newBlob, tmpPath)
  310. } catch (error) {
  311. logger.warn({ error, projectId, hash }, 'Failed to backup blob')
  312. }
  313. }
  314. res.status(HTTPStatus.CREATED).end()
  315. })
  316. }
  317. async function headProjectBlob(req, res) {
  318. const { params } = parseReq(req, schemas.headProjectBlob)
  319. const projectId = params.project_id
  320. const hash = params.hash
  321. const blobStore = new BlobStore(projectId)
  322. const blob = await blobStore.getBlob(hash)
  323. if (blob) {
  324. res.set('Content-Length', blob.getByteLength())
  325. res.status(200).end()
  326. } else {
  327. res.status(404).end()
  328. }
  329. }
  330. // Support simple, singular ranges starting from zero only, up-to 2MB = 2_000_000, 7 digits
  331. const RANGE_HEADER = /^bytes=(\d{1,7})-(\d{1,7})$/
  332. /**
  333. * @param {string} header
  334. * @return {undefined | {start: number, end: number}}
  335. * @private
  336. */
  337. function _getRangeOpts(header) {
  338. if (!header) return undefined
  339. const match = header.match(RANGE_HEADER)
  340. if (match) {
  341. const start = parseInt(match[1], 10)
  342. const end = parseInt(match[2], 10)
  343. return { start, end }
  344. }
  345. return undefined
  346. }
  347. async function getProjectBlob(req, res, next) {
  348. const { params, headers } = parseReq(req, schemas.getProjectBlob)
  349. const projectId = params.project_id
  350. const hash = params.hash
  351. const rangeHeader = headers.range || ''
  352. const opts = _getRangeOpts(rangeHeader)
  353. const blobStore = new BlobStore(projectId)
  354. logger.debug({ projectId, hash }, 'getProjectBlob started')
  355. try {
  356. if (req.method === 'HEAD') {
  357. return await headProjectBlob(req, res)
  358. }
  359. let stream
  360. try {
  361. if (opts) {
  362. // This is a range request, so we need to set the appropriate headers
  363. // Browser caching only works if the total size is known, so we have
  364. // to fetch the blob metadata first.
  365. const metaData = await blobStore.getBlob(hash)
  366. if (metaData) {
  367. const blobLength = metaData.getByteLength()
  368. if (opts.start > opts.end || opts.start >= blobLength) {
  369. return res
  370. .status(416) // Requested Range Not Satisfiable
  371. .set('Content-Range', `bytes */${blobLength}`)
  372. .set('Content-Length', '0')
  373. .end()
  374. }
  375. // Valid range request
  376. const actualEnd = Math.min(opts.end, blobLength - 1)
  377. const returnedSize = actualEnd - opts.start + 1
  378. res.set('Content-Length', returnedSize)
  379. res.set(
  380. 'Content-Range',
  381. `bytes ${opts.start}-${actualEnd}/${blobLength}`
  382. )
  383. res.status(206)
  384. }
  385. }
  386. stream = await blobStore.getStream(hash, opts)
  387. } catch (err) {
  388. if (err instanceof Blob.NotFoundError) {
  389. logger.warn({ projectId, hash }, 'Blob not found')
  390. return res.status(404).end()
  391. } else {
  392. throw err
  393. }
  394. }
  395. res.set('Content-Type', 'application/octet-stream')
  396. try {
  397. await pipeline(stream, res)
  398. } catch (err) {
  399. if (
  400. err?.code === 'ERR_STREAM_PREMATURE_CLOSE' ||
  401. err?.code === 'ERR_STREAM_UNABLE_TO_PIPE'
  402. ) {
  403. res.end()
  404. } else {
  405. throw OError.tag(err, 'error transferring stream', { projectId, hash })
  406. }
  407. }
  408. } finally {
  409. logger.debug({ projectId, hash }, 'getProjectBlob finished')
  410. }
  411. }
  412. async function copyProjectBlob(req, res, next) {
  413. const { params, query } = parseReq(req, schemas.copyProjectBlob)
  414. const sourceProjectId = query.copyFrom
  415. const targetProjectId = params.project_id
  416. const blobHash = params.hash
  417. // Check that blob exists in source project
  418. const sourceBlobStore = new BlobStore(sourceProjectId)
  419. const targetBlobStore = new BlobStore(targetProjectId)
  420. const [sourceBlob, targetBlob] = await Promise.all([
  421. sourceBlobStore.getBlob(blobHash),
  422. targetBlobStore.getBlob(blobHash),
  423. ])
  424. if (!sourceBlob) {
  425. logger.warn(
  426. { sourceProjectId, targetProjectId, blobHash },
  427. 'missing source blob when copying across projects'
  428. )
  429. return render.notFound(res)
  430. }
  431. // Exit early if the blob exists in the target project.
  432. // This will also catch global blobs, which always exist.
  433. if (targetBlob) {
  434. return res.status(HTTPStatus.NO_CONTENT).end()
  435. }
  436. // Otherwise, copy blob from source project to target project
  437. await sourceBlobStore.copyBlob(sourceBlob, targetProjectId)
  438. res.status(HTTPStatus.CREATED).end()
  439. }
  440. async function getSnapshotAtVersion(projectId, version) {
  441. const chunk = await chunkStore.loadAtVersion(projectId, version)
  442. const snapshot = chunk.getSnapshot()
  443. const changes = _.dropRight(
  444. chunk.getChanges(),
  445. chunk.getEndVersion() - version
  446. )
  447. if (changes.length > 0) {
  448. snapshot.applyAll(changes)
  449. } else {
  450. // There are no changes in this chunk; we need to look at the previous chunk
  451. // to get the snapshot's timestamp
  452. let chunkMetadata
  453. try {
  454. chunkMetadata = await getChunkMetadataForVersion(projectId, version)
  455. } catch (err) {
  456. if (err instanceof Chunk.VersionNotFoundError) {
  457. // The snapshot is the first snapshot of the first chunk, so we can't
  458. // find a timestamp. This shouldn't happen often. Ignore the error and
  459. // leave the timestamp empty.
  460. } else {
  461. throw err
  462. }
  463. }
  464. snapshot.setTimestamp(chunkMetadata.endTimestamp)
  465. }
  466. return snapshot
  467. }
  468. function sumUpByteLength(blobs) {
  469. return blobs.reduce((sum, blob) => sum + blob.getByteLength(), 0)
  470. }
  471. async function getBlobStats(req, res) {
  472. const { params, body } = parseReq(req, schemas.getBlobStats)
  473. const projectId = params.project_id
  474. const blobHashes = body.blobHashes || []
  475. for (const hash of blobHashes) {
  476. assert.blobHash(hash, 'bad hash')
  477. }
  478. const blobStore = new BlobStore(projectId)
  479. const batchBlobStore = new BatchBlobStore(blobStore)
  480. await batchBlobStore.preload(Array.from(blobHashes))
  481. const blobs = Array.from(batchBlobStore.blobs.values()).filter(Boolean)
  482. const textBlobs = blobs.filter(b => b.getStringLength() !== null)
  483. const binaryBlobs = blobs.filter(b => b.getStringLength() === null)
  484. const textBlobBytes = sumUpByteLength(textBlobs)
  485. const binaryBlobBytes = sumUpByteLength(binaryBlobs)
  486. res.json({
  487. projectId,
  488. textBlobBytes,
  489. binaryBlobBytes,
  490. totalBytes: textBlobBytes + binaryBlobBytes,
  491. nTextBlobs: textBlobs.length,
  492. nBinaryBlobs: binaryBlobs.length,
  493. })
  494. }
  495. async function getProjectBlobsStats(req, res) {
  496. const { body } = parseReq(req, schemas.getProjectBlobsStats)
  497. const projectIds = body.projectIds
  498. const { blobs } = await getProjectBlobsBatch(
  499. projectIds.map(id => {
  500. if (assert.POSTGRES_ID_REGEXP.test(id)) {
  501. return parseInt(id, 10)
  502. } else {
  503. return id
  504. }
  505. })
  506. )
  507. const sizes = []
  508. for (const projectId of projectIds) {
  509. const projectBlobs = blobs.get(projectId) || []
  510. const textBlobs = projectBlobs.filter(b => b.getStringLength() !== null)
  511. const binaryBlobs = projectBlobs.filter(b => b.getStringLength() === null)
  512. const textBlobBytes = sumUpByteLength(textBlobs)
  513. const binaryBlobBytes = sumUpByteLength(binaryBlobs)
  514. sizes.push({
  515. projectId,
  516. textBlobBytes,
  517. binaryBlobBytes,
  518. totalBytes: textBlobBytes + binaryBlobBytes,
  519. nTextBlobs: textBlobs.length,
  520. nBinaryBlobs: binaryBlobs.length,
  521. })
  522. }
  523. res.json(sizes)
  524. }
  525. module.exports = {
  526. initializeProject: expressify(initializeProject),
  527. cloneProject: expressify(cloneProject),
  528. getLatestContent: expressify(getLatestContent),
  529. getContentAtVersion: expressify(getContentAtVersion),
  530. getLatestHashedContent: expressify(getLatestHashedContent),
  531. getLatestPersistedHistory: expressify(getLatestHistory),
  532. getLatestHistory: expressify(getLatestHistory),
  533. getLatestHistoryRaw: expressify(getLatestHistoryRaw),
  534. getHistory: expressify(getHistory),
  535. getHistoryBefore: expressify(getHistoryBefore),
  536. getChanges: expressify(getChanges),
  537. getZip: expressify(getZip),
  538. createZip: expressify(createZip),
  539. deleteProject: expressify(deleteProject),
  540. createProjectBlob: expressify(createProjectBlob),
  541. getProjectBlob: expressify(getProjectBlob),
  542. headProjectBlob: expressify(headProjectBlob),
  543. copyProjectBlob: expressify(copyProjectBlob),
  544. getBlobStats: expressify(getBlobStats),
  545. getProjectBlobsStats: expressify(getProjectBlobsStats),
  546. }