| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081 |
- // @ts-check
- import Crypto from 'node:crypto'
- import Events from 'node:events'
- import fs from 'node:fs'
- import Path from 'node:path'
- import { performance } from 'node:perf_hooks'
- import Stream from 'node:stream'
- import zLib from 'node:zlib'
- import { setTimeout } from 'node:timers/promises'
- import { Binary, ObjectId } from 'mongodb'
- import logger from '@overleaf/logger'
- import {
- batchedUpdate,
- READ_PREFERENCE_SECONDARY,
- } from '@overleaf/mongo-utils/batchedUpdate.js'
- import OError from '@overleaf/o-error'
- import {
- AlreadyWrittenError,
- NoKEKMatchedError,
- NotFoundError,
- } from '@overleaf/object-persistor/src/Errors.js'
- import { promiseMapWithLimit } from '@overleaf/promise-utils'
- import { backupPersistor, projectBlobsBucket } from '../lib/backupPersistor.mjs'
- import {
- BlobStore,
- GLOBAL_BLOBS,
- loadGlobalBlobs,
- getStringLengthOfFile,
- makeBlobForFile,
- makeProjectKey,
- } from '../lib/blob_store/index.js'
- import { backedUpBlobs, db } from '../lib/mongodb.js'
- import filestorePersistor from '../lib/persistor.js'
- // Silence warning.
- Events.setMaxListeners(20)
- // Enable caching for ObjectId.toString()
- ObjectId.cacheHexString = true
- /**
- * @typedef {import("overleaf-editor-core").Blob} Blob
- * @typedef {import("perf_hooks").EventLoopUtilization} EventLoopUtilization
- * @typedef {import("mongodb").Collection} Collection
- * @typedef {import("@overleaf/object-persistor/src/PerProjectEncryptedS3Persistor").CachedPerProjectEncryptedS3Persistor} CachedPerProjectEncryptedS3Persistor
- */
- /**
- * @typedef {Object} FileRef
- * @property {ObjectId} _id
- * @property {string} hash
- */
- /**
- * @typedef {Object} Folder
- * @property {Array<Folder>} folders
- * @property {Array<FileRef>} fileRefs
- */
- /**
- * @typedef {Object} DeletedFileRef
- * @property {ObjectId} _id
- * @property {ObjectId} projectId
- * @property {string} hash
- */
- /**
- * @typedef {Object} Project
- * @property {ObjectId} _id
- * @property {Array<Folder>} rootFolder
- * @property {Array<string>} deletedFileIds
- * @property {Array<Blob>} blobs
- * @property {{history: {id: string}}} overleaf
- * @property {Array<string>} [backedUpBlobs]
- */
- /**
- * @typedef {Object} QueueEntry
- * @property {ProjectContext} ctx
- * @property {string} cacheKey
- * @property {string} [fileId]
- * @property {string} path
- * @property {string} [hash]
- * @property {Blob} [blob]
- */
- const COLLECT_BLOBS = process.argv.includes('blobs')
- // Time of closing the ticket for adding hashes: https://github.com/overleaf/internal/issues/464#issuecomment-492668129
- const ALL_PROJECTS_HAVE_FILE_HASHES_AFTER = new Date('2019-05-15T14:02:00Z')
- const PUBLIC_LAUNCH_DATE = new Date('2012-01-01T00:00:00Z')
- const BATCH_RANGE_START =
- process.env.BATCH_RANGE_START ||
- ObjectId.createFromTime(PUBLIC_LAUNCH_DATE.getTime() / 1000).toString()
- const BATCH_RANGE_END =
- process.env.BATCH_RANGE_END ||
- ObjectId.createFromTime(
- ALL_PROJECTS_HAVE_FILE_HASHES_AFTER.getTime() / 1000
- ).toString()
- // We need to control the start and end as ids of deleted projects are created at time of deletion.
- delete process.env.BATCH_RANGE_START
- delete process.env.BATCH_RANGE_END
- // Concurrency for downloading from GCS and updating hashes in mongo
- const CONCURRENCY = parseInt(process.env.CONCURRENCY || '100', 10)
- // Retries for processing a given file
- const RETRIES = parseInt(process.env.RETRIES || '10', 10)
- const RETRY_DELAY_MS = parseInt(process.env.RETRY_DELAY_MS || '100', 10)
- const USER_FILES_BUCKET_NAME = process.env.USER_FILES_BUCKET_NAME || ''
- if (!USER_FILES_BUCKET_NAME) {
- throw new Error('env var USER_FILES_BUCKET_NAME is missing')
- }
- const RETRY_FILESTORE_404 = process.env.RETRY_FILESTORE_404 === 'true'
- const BUFFER_DIR = fs.mkdtempSync(
- process.env.BUFFER_DIR_PREFIX || '/tmp/back_fill_file_hash-'
- )
- // https://nodejs.org/api/stream.html#streamgetdefaulthighwatermarkobjectmode
- const STREAM_HIGH_WATER_MARK = parseInt(
- process.env.STREAM_HIGH_WATER_MARK || (64 * 1024).toString(),
- 10
- )
- const LOGGING_INTERVAL = parseInt(process.env.LOGGING_INTERVAL || '60000', 10)
- const projectsCollection = db.collection('projects')
- const deletedProjectsCollection = db.collection('deletedProjects')
- const deletedFilesCollection = db.collection('deletedFiles')
- const STATS = {
- projects: 0,
- blobs: 0,
- backedUpBlobs: 0,
- filesWithoutHash: 0,
- filesDuplicated: 0,
- filesRetries: 0,
- filesFailed: 0,
- fileTreeUpdated: 0,
- globalBlobsCount: 0,
- globalBlobsEgress: 0,
- projectDeleted: 0,
- projectHardDeleted: 0,
- fileHardDeleted: 0,
- mongoUpdates: 0,
- deduplicatedWriteToAWSLocalCount: 0,
- deduplicatedWriteToAWSLocalEgress: 0,
- deduplicatedWriteToAWSRemoteCount: 0,
- deduplicatedWriteToAWSRemoteEgress: 0,
- readFromGCSCount: 0,
- readFromGCSIngress: 0,
- writeToAWSCount: 0,
- writeToAWSEgress: 0,
- writeToGCSCount: 0,
- writeToGCSEgress: 0,
- }
- const processStart = performance.now()
- let lastLogTS = processStart
- let lastLog = Object.assign({}, STATS)
- let lastEventLoopStats = performance.eventLoopUtilization()
- /**
- * @param {number} v
- * @param {number} ms
- */
- function toMiBPerSecond(v, ms) {
- const ONE_MiB = 1024 * 1024
- return v / ONE_MiB / (ms / 1000)
- }
- /**
- * @param {any} stats
- * @param {number} ms
- * @return {{writeToAWSThroughputMiBPerSecond: number, readFromGCSThroughputMiBPerSecond: number}}
- */
- function bandwidthStats(stats, ms) {
- return {
- readFromGCSThroughputMiBPerSecond: toMiBPerSecond(
- stats.readFromGCSIngress,
- ms
- ),
- writeToAWSThroughputMiBPerSecond: toMiBPerSecond(
- stats.writeToAWSEgress,
- ms
- ),
- }
- }
- /**
- * @param {EventLoopUtilization} nextEventLoopStats
- * @param {number} now
- * @return {Object}
- */
- function computeDiff(nextEventLoopStats, now) {
- const ms = now - lastLogTS
- lastLogTS = now
- const diff = {
- eventLoop: performance.eventLoopUtilization(
- nextEventLoopStats,
- lastEventLoopStats
- ),
- }
- for (const [name, v] of Object.entries(STATS)) {
- diff[name] = v - lastLog[name]
- }
- return Object.assign(diff, bandwidthStats(diff, ms))
- }
- function printStats() {
- const now = performance.now()
- const nextEventLoopStats = performance.eventLoopUtilization()
- console.log(
- JSON.stringify({
- time: new Date(),
- ...STATS,
- ...bandwidthStats(STATS, now - processStart),
- eventLoop: nextEventLoopStats,
- diff: computeDiff(nextEventLoopStats, now),
- })
- )
- lastEventLoopStats = nextEventLoopStats
- lastLog = Object.assign({}, STATS)
- }
- setInterval(printStats, LOGGING_INTERVAL)
- /**
- * @param {QueueEntry} entry
- * @return {Promise<string>}
- */
- async function processFile(entry) {
- for (let attempt = 0; attempt < RETRIES; attempt++) {
- try {
- return await processFileOnce(entry)
- } catch (err) {
- if (err instanceof NotFoundError) {
- const { bucketName } = OError.getFullInfo(err)
- if (bucketName === USER_FILES_BUCKET_NAME && !RETRY_FILESTORE_404) {
- throw err // disable retries for not found in filestore bucket case
- }
- }
- if (err instanceof NoKEKMatchedError) {
- throw err // disable retries when upload to S3 will fail again
- }
- STATS.filesRetries++
- const {
- ctx: { projectId },
- fileId,
- path,
- } = entry
- logger.warn(
- { err, projectId, fileId, path, attempt },
- 'failed to process file, trying again'
- )
- await setTimeout(RETRY_DELAY_MS)
- }
- }
- return await processFileOnce(entry)
- }
- /**
- * @param {QueueEntry} entry
- * @return {Promise<string>}
- */
- async function processFileOnce(entry) {
- const { projectId, historyId } = entry.ctx
- const { fileId, cacheKey } = entry
- const filePath = Path.join(BUFFER_DIR, projectId.toString() + cacheKey)
- const blobStore = new BlobStore(historyId)
- if (entry.blob) {
- const { blob } = entry
- const hash = blob.getHash()
- if (entry.ctx.hasBackedUpBlob(hash)) {
- STATS.deduplicatedWriteToAWSLocalCount++
- STATS.deduplicatedWriteToAWSLocalEgress += estimateBlobSize(blob)
- return hash
- }
- entry.ctx.recordPendingBlob(hash)
- STATS.readFromGCSCount++
- const src = await blobStore.getStream(hash)
- const dst = fs.createWriteStream(filePath, {
- highWaterMark: STREAM_HIGH_WATER_MARK,
- })
- try {
- await Stream.promises.pipeline(src, dst)
- } finally {
- STATS.readFromGCSIngress += dst.bytesWritten
- }
- await uploadBlobToAWS(entry, blob, filePath)
- return hash
- }
- STATS.readFromGCSCount++
- const src = await filestorePersistor.getObjectStream(
- USER_FILES_BUCKET_NAME,
- `${projectId}/${fileId}`
- )
- const dst = fs.createWriteStream(filePath, {
- highWaterMark: STREAM_HIGH_WATER_MARK,
- })
- try {
- await Stream.promises.pipeline(src, dst)
- } finally {
- STATS.readFromGCSIngress += dst.bytesWritten
- }
- const blob = await makeBlobForFile(filePath)
- blob.setStringLength(
- await getStringLengthOfFile(blob.getByteLength(), filePath)
- )
- const hash = blob.getHash()
- if (GLOBAL_BLOBS.has(hash)) {
- STATS.globalBlobsCount++
- STATS.globalBlobsEgress += estimateBlobSize(blob)
- return hash
- }
- if (entry.ctx.hasBackedUpBlob(hash)) {
- STATS.deduplicatedWriteToAWSLocalCount++
- STATS.deduplicatedWriteToAWSLocalEgress += estimateBlobSize(blob)
- return hash
- }
- entry.ctx.recordPendingBlob(hash)
- try {
- await uploadBlobToGCS(blobStore, entry, blob, hash, filePath)
- await uploadBlobToAWS(entry, blob, filePath)
- } catch (err) {
- entry.ctx.recordFailedBlob(hash)
- throw err
- }
- return hash
- }
- /**
- * @param {BlobStore} blobStore
- * @param {QueueEntry} entry
- * @param {Blob} blob
- * @param {string} hash
- * @param {string} filePath
- * @return {Promise<void>}
- */
- async function uploadBlobToGCS(blobStore, entry, blob, hash, filePath) {
- if (entry.ctx.hasHistoryBlob(hash)) {
- return // fast-path using hint from pre-fetched blobs
- }
- if (!COLLECT_BLOBS && (await blobStore.getBlob(hash))) {
- entry.ctx.recordHistoryBlob(hash)
- return // round trip to postgres/mongo when not pre-fetched
- }
- // blob missing in history-v1, create in GCS and persist in postgres/mongo
- STATS.writeToGCSCount++
- STATS.writeToGCSEgress += blob.getByteLength()
- await blobStore.putBlob(filePath, blob)
- entry.ctx.recordHistoryBlob(hash)
- }
- /**
- * @param {QueueEntry} entry
- * @param {Blob} blob
- * @param {string} filePath
- * @return {Promise<void>}
- */
- async function uploadBlobToAWS(entry, blob, filePath) {
- const { historyId } = entry.ctx
- let backupSource
- let contentEncoding
- const md5 = Crypto.createHash('md5')
- let size
- if (blob.getStringLength()) {
- const filePathCompressed = filePath + '.gz'
- backupSource = filePathCompressed
- contentEncoding = 'gzip'
- size = 0
- await Stream.promises.pipeline(
- fs.createReadStream(filePath, { highWaterMark: STREAM_HIGH_WATER_MARK }),
- zLib.createGzip(),
- async function* (source) {
- for await (const chunk of source) {
- size += chunk.byteLength
- md5.update(chunk)
- yield chunk
- }
- },
- fs.createWriteStream(filePathCompressed, {
- highWaterMark: STREAM_HIGH_WATER_MARK,
- })
- )
- } else {
- backupSource = filePath
- size = blob.getByteLength()
- await Stream.promises.pipeline(
- fs.createReadStream(filePath, { highWaterMark: STREAM_HIGH_WATER_MARK }),
- md5
- )
- }
- const backendKeyPath = makeProjectKey(historyId, blob.getHash())
- const persistor = await entry.ctx.getCachedPersistor(backendKeyPath)
- try {
- STATS.writeToAWSCount++
- await persistor.sendStream(
- projectBlobsBucket,
- backendKeyPath,
- fs.createReadStream(backupSource, {
- highWaterMark: STREAM_HIGH_WATER_MARK,
- }),
- {
- contentEncoding,
- contentType: 'application/octet-stream',
- contentLength: size,
- sourceMd5: md5.digest('hex'),
- ifNoneMatch: '*', // de-duplicate write (we pay for the request, but avoid egress)
- }
- )
- STATS.writeToAWSEgress += size
- } catch (err) {
- if (err instanceof AlreadyWrittenError) {
- STATS.deduplicatedWriteToAWSRemoteCount++
- STATS.deduplicatedWriteToAWSRemoteEgress += size
- } else {
- STATS.writeToAWSEgress += size
- throw err
- }
- }
- entry.ctx.recordBackedUpBlob(blob.getHash())
- }
- /**
- * @param {Array<QueueEntry>} files
- * @return {Promise<void>}
- */
- async function processFiles(files) {
- if (files.length === 0) return // all processed
- await fs.promises.mkdir(BUFFER_DIR, { recursive: true })
- try {
- await promiseMapWithLimit(
- CONCURRENCY,
- files,
- /**
- * @param {QueueEntry} entry
- * @return {Promise<void>}
- */
- async function (entry) {
- try {
- await entry.ctx.processFile(entry)
- } catch (err) {
- STATS.filesFailed++
- const {
- ctx: { projectId },
- fileId,
- path,
- } = entry
- logger.error(
- { err, projectId, fileId, path },
- 'failed to process file'
- )
- }
- }
- )
- } finally {
- await fs.promises.rm(BUFFER_DIR, { recursive: true, force: true })
- }
- }
- /**
- * @param {Array<Project>} batch
- * @param {string} prefix
- * @return {Promise<void>}
- */
- async function handleLiveTreeBatch(batch, prefix = 'rootFolder.0') {
- let nBackedUpBlobs = 0
- if (process.argv.includes('collectBackedUpBlobs')) {
- nBackedUpBlobs = await collectBackedUpBlobs(batch)
- }
- if (process.argv.includes('deletedFiles')) {
- await collectDeletedFiles(batch)
- }
- let blobs = 0
- if (COLLECT_BLOBS) {
- blobs = await collectBlobs(batch)
- }
- const files = Array.from(findFileInBatch(batch, prefix))
- STATS.projects += batch.length
- STATS.blobs += blobs
- STATS.backedUpBlobs += nBackedUpBlobs
- STATS.filesWithoutHash += files.length - (blobs - nBackedUpBlobs)
- batch.length = 0 // GC
- // The files are currently ordered by project-id.
- // Order them by file-id ASC then blobs ASC to
- // - process files before blobs
- // - avoid head-of-line blocking from many project-files waiting on the generation of the projects DEK (round trip to AWS)
- // - bonus: increase chance of de-duplicating write to AWS
- files.sort(
- /**
- * @param {QueueEntry} a
- * @param {QueueEntry} b
- * @return {number}
- */
- function (a, b) {
- if (a.fileId && b.fileId) return a.fileId > b.fileId ? 1 : -1
- if (a.hash && b.hash) return a.hash > b.hash ? 1 : -1
- if (a.fileId) return -1
- return 1
- }
- )
- await processFiles(files)
- await promiseMapWithLimit(
- CONCURRENCY,
- files,
- /**
- * @param {QueueEntry} entry
- * @return {Promise<void>}
- */
- async function (entry) {
- await entry.ctx.flushMongoQueues()
- }
- )
- }
- /**
- * @param {Array<{project: Project}>} batch
- * @return {Promise<void>}
- */
- async function handleDeletedFileTreeBatch(batch) {
- await handleLiveTreeBatch(
- batch.map(d => d.project),
- 'project.rootFolder.0'
- )
- }
- /**
- * @param {QueueEntry} entry
- * @return {Promise<boolean>}
- */
- async function tryUpdateFileRefInMongo(entry) {
- if (entry.path === '') {
- return await tryUpdateDeletedFileRefInMongo(entry)
- } else if (entry.path.startsWith('project.')) {
- return await tryUpdateFileRefInMongoInDeletedProject(entry)
- }
- STATS.mongoUpdates++
- const result = await projectsCollection.updateOne(
- {
- _id: entry.ctx.projectId,
- [`${entry.path}._id`]: new ObjectId(entry.fileId),
- },
- {
- $set: { [`${entry.path}.hash`]: entry.hash },
- }
- )
- return result.matchedCount === 1
- }
- /**
- * @param {QueueEntry} entry
- * @return {Promise<boolean>}
- */
- async function tryUpdateDeletedFileRefInMongo(entry) {
- STATS.mongoUpdates++
- const result = await deletedFilesCollection.updateOne(
- {
- _id: new ObjectId(entry.fileId),
- projectId: entry.ctx.projectId,
- },
- { $set: { hash: entry.hash } }
- )
- return result.matchedCount === 1
- }
- /**
- * @param {QueueEntry} entry
- * @return {Promise<boolean>}
- */
- async function tryUpdateFileRefInMongoInDeletedProject(entry) {
- STATS.mongoUpdates++
- const result = await deletedProjectsCollection.updateOne(
- {
- 'deleterData.deletedProjectId': entry.ctx.projectId,
- [`${entry.path}._id`]: new ObjectId(entry.fileId),
- },
- {
- $set: { [`${entry.path}.hash`]: entry.hash },
- }
- )
- return result.matchedCount === 1
- }
- const RETRY_UPDATE_HASH = 100
- /**
- * @param {QueueEntry} entry
- * @return {Promise<void>}
- */
- async function updateFileRefInMongo(entry) {
- if (await tryUpdateFileRefInMongo(entry)) return
- const { fileId } = entry
- const { projectId } = entry.ctx
- for (let i = 0; i < RETRY_UPDATE_HASH; i++) {
- let prefix = 'rootFolder.0'
- let p = await projectsCollection.findOne(
- { _id: projectId },
- { projection: { rootFolder: 1 } }
- )
- if (!p) {
- STATS.projectDeleted++
- prefix = 'project.rootFolder.0'
- const deletedProject = await deletedProjectsCollection.findOne(
- {
- 'deleterData.deletedProjectId': projectId,
- project: { $exists: true },
- },
- { projection: { 'project.rootFolder': 1 } }
- )
- p = deletedProject?.project
- if (!p) {
- STATS.projectHardDeleted++
- console.warn(
- 'bug: project hard-deleted while processing',
- projectId,
- fileId
- )
- return
- }
- }
- let found = false
- for (const e of findFiles(entry.ctx, p.rootFolder[0], prefix)) {
- found = e.fileId === fileId
- if (!found) continue
- if (await tryUpdateFileRefInMongo(e)) return
- break
- }
- if (!found) {
- if (await tryUpdateDeletedFileRefInMongo(entry)) return
- STATS.fileHardDeleted++
- console.warn('bug: file hard-deleted while processing', projectId, fileId)
- return
- }
- STATS.fileTreeUpdated++
- }
- throw new OError(
- 'file-tree updated repeatedly while trying to add hash',
- entry
- )
- }
- /**
- * @param {ProjectContext} ctx
- * @param {Folder} folder
- * @param {string} path
- * @return Generator<QueueEntry>
- */
- function* findFiles(ctx, folder, path) {
- let i = 0
- for (const child of folder.folders) {
- yield* findFiles(ctx, child, `${path}.folders.${i}`)
- i++
- }
- i = 0
- for (const fileRef of folder.fileRefs) {
- if (!fileRef.hash) {
- yield {
- ctx,
- cacheKey: fileRef._id.toString(),
- fileId: fileRef._id.toString(),
- path: `${path}.fileRefs.${i}`,
- }
- }
- i++
- }
- }
- /**
- * @param {Array<Project>} projects
- * @param {string} prefix
- * @return Generator<QueueEntry>
- */
- function* findFileInBatch(projects, prefix) {
- for (const project of projects) {
- const ctx = new ProjectContext(project)
- yield* findFiles(ctx, project.rootFolder[0], prefix)
- for (const fileId of project.deletedFileIds || []) {
- yield { ctx, cacheKey: fileId, fileId, path: '' }
- }
- for (const blob of project.blobs || []) {
- if (ctx.hasBackedUpBlob(blob.getHash())) continue
- yield {
- ctx,
- cacheKey: blob.getHash(),
- path: 'blob',
- blob,
- hash: blob.getHash(),
- }
- }
- }
- }
- /**
- * @param {Array<Project>} projects
- * @return {Promise<number>}
- */
- async function collectBlobs(projects) {
- let blobs = 0
- for (const project of projects) {
- const historyId = project.overleaf.history.id.toString()
- const blobStore = new BlobStore(historyId)
- project.blobs = await blobStore.getProjectBlobs()
- blobs += project.blobs.length
- }
- return blobs
- }
- /**
- * @param {Array<Project>} projects
- * @return {Promise<void>}
- */
- async function collectDeletedFiles(projects) {
- const cursor = deletedFilesCollection.find(
- {
- projectId: { $in: projects.map(p => p._id) },
- hash: { $exists: false },
- },
- {
- projection: { _id: 1, projectId: 1 },
- readPreference: READ_PREFERENCE_SECONDARY,
- sort: { projectId: 1 },
- }
- )
- const processed = projects.slice()
- for await (const deletedFileRef of cursor) {
- const idx = processed.findIndex(
- p => p._id.toString() === deletedFileRef.projectId.toString()
- )
- if (idx === -1) {
- throw new Error(
- `bug: order of deletedFiles mongo records does not match batch of projects (${deletedFileRef.projectId} out of order)`
- )
- }
- processed.splice(0, idx)
- const project = processed[0]
- project.deletedFileIds = project.deletedFileIds || []
- project.deletedFileIds.push(deletedFileRef._id.toString())
- }
- }
- /**
- * @param {Array<Project>} projects
- * @return {Promise<number>}
- */
- async function collectBackedUpBlobs(projects) {
- const cursor = backedUpBlobs.find(
- { _id: { $in: projects.map(p => p._id) } },
- {
- readPreference: READ_PREFERENCE_SECONDARY,
- sort: { _id: 1 },
- }
- )
- let nBackedUpBlobs = 0
- const processed = projects.slice()
- for await (const record of cursor) {
- const idx = processed.findIndex(
- p => p._id.toString() === record._id.toString()
- )
- if (idx === -1) {
- throw new Error(
- `bug: order of backedUpBlobs mongo records does not match batch of projects (${record._id} out of order)`
- )
- }
- processed.splice(0, idx)
- const project = processed[0]
- project.backedUpBlobs = record.blobs.map(b => b.toString('hex'))
- nBackedUpBlobs += record.blobs.length
- }
- return nBackedUpBlobs
- }
- const BATCH_HASH_WRITES = 1_000
- const BATCH_FILE_UPDATES = 100
- class ProjectContext {
- /** @type {Promise<CachedPerProjectEncryptedS3Persistor> | null} */
- #cachedPersistorPromise = null
- /** @type {Set<string>} */
- #backedUpBlobs
- /** @type {Set<string>} */
- #historyBlobs
- /**
- * @param {Project} project
- */
- constructor(project) {
- this.projectId = project._id
- this.historyId = project.overleaf.history.id.toString()
- this.#backedUpBlobs = new Set(project.backedUpBlobs || [])
- this.#historyBlobs = new Set((project.blobs || []).map(b => b.getHash()))
- }
- hasHistoryBlob(hash) {
- return this.#historyBlobs.has(hash)
- }
- recordHistoryBlob(hash) {
- this.#historyBlobs.add(hash)
- }
- /**
- * @param {string} key
- * @return {Promise<CachedPerProjectEncryptedS3Persistor>}
- */
- getCachedPersistor(key) {
- if (!this.#cachedPersistorPromise) {
- // Fetch DEK once, but only if needed -- upon the first use
- this.#cachedPersistorPromise = this.#getCachedPersistorWithRetries(key)
- }
- return this.#cachedPersistorPromise
- }
- /**
- * @param {string} key
- * @return {Promise<CachedPerProjectEncryptedS3Persistor>}
- */
- async #getCachedPersistorWithRetries(key) {
- for (let attempt = 0; attempt < RETRIES; attempt++) {
- try {
- return await backupPersistor.forProject(projectBlobsBucket, key)
- } catch (err) {
- if (err instanceof NoKEKMatchedError) {
- throw err
- } else {
- logger.warn(
- { err, projectId: this.projectId, attempt },
- 'failed to get DEK, trying again'
- )
- await setTimeout(RETRY_DELAY_MS)
- }
- }
- }
- return await backupPersistor.forProject(projectBlobsBucket, key)
- }
- async flushMongoQueuesIfNeeded() {
- if (this.#completedBlobs.size > BATCH_HASH_WRITES) {
- await this.#storeBackedUpBlobs()
- }
- if (this.#pendingFileWrites.length > BATCH_FILE_UPDATES) {
- await this.#storeFileHashes()
- }
- }
- async flushMongoQueues() {
- await this.#storeBackedUpBlobs()
- await this.#storeFileHashes()
- }
- /** @type {Set<string>} */
- #pendingBlobs = new Set()
- /** @type {Set<string>} */
- #completedBlobs = new Set()
- async #storeBackedUpBlobs() {
- if (this.#completedBlobs.size === 0) return
- const blobs = Array.from(this.#completedBlobs).map(
- hash => new Binary(Buffer.from(hash, 'hex'))
- )
- this.#completedBlobs.clear()
- STATS.mongoUpdates++
- await backedUpBlobs.updateOne(
- { _id: this.projectId },
- { $addToSet: { blobs: { $each: blobs } } },
- { upsert: true }
- )
- }
- /**
- * @param {string} hash
- */
- recordPendingBlob(hash) {
- this.#pendingBlobs.add(hash)
- }
- /**
- * @param {string} hash
- */
- recordFailedBlob(hash) {
- this.#pendingBlobs.delete(hash)
- }
- /**
- * @param {string} hash
- */
- recordBackedUpBlob(hash) {
- this.#backedUpBlobs.add(hash)
- this.#completedBlobs.add(hash)
- this.#pendingBlobs.delete(hash)
- }
- /**
- * @param {string} hash
- * @return {boolean}
- */
- hasBackedUpBlob(hash) {
- return (
- this.#pendingBlobs.has(hash) ||
- this.#completedBlobs.has(hash) ||
- this.#backedUpBlobs.has(hash)
- )
- }
- /** @type {Array<QueueEntry>} */
- #pendingFileWrites = []
- /**
- * @param {QueueEntry} entry
- */
- queueFileForWritingHash(entry) {
- if (entry.path === 'blob') return
- this.#pendingFileWrites.push(entry)
- }
- /**
- * @param {Collection} collection
- * @param {Array<QueueEntry>} entries
- * @param {Object} query
- * @return {Promise<Array<QueueEntry>>}
- */
- async #tryBatchHashWrites(collection, entries, query) {
- if (entries.length === 0) return []
- const update = {}
- for (const entry of entries) {
- query[`${entry.path}._id`] = new ObjectId(entry.fileId)
- update[`${entry.path}.hash`] = entry.hash
- }
- STATS.mongoUpdates++
- const result = await collection.updateOne(query, { $set: update })
- if (result.matchedCount === 1) {
- return [] // all updated
- }
- return entries
- }
- async #storeFileHashes() {
- if (this.#pendingFileWrites.length === 0) return
- const individualUpdates = []
- const projectEntries = []
- const deletedProjectEntries = []
- for (const entry of this.#pendingFileWrites) {
- if (entry.path === '') {
- individualUpdates.push(entry)
- } else if (entry.path.startsWith('project.')) {
- deletedProjectEntries.push(entry)
- } else {
- projectEntries.push(entry)
- }
- }
- this.#pendingFileWrites.length = 0
- // Try to process them together, otherwise fallback to individual updates and retries.
- individualUpdates.push(
- ...(await this.#tryBatchHashWrites(projectsCollection, projectEntries, {
- _id: this.projectId,
- }))
- )
- individualUpdates.push(
- ...(await this.#tryBatchHashWrites(
- deletedProjectsCollection,
- deletedProjectEntries,
- { 'deleterData.deletedProjectId': this.projectId }
- ))
- )
- for (const entry of individualUpdates) {
- await updateFileRefInMongo(entry)
- }
- }
- /** @type {Map<string, Promise<string>>} */
- #pendingFiles = new Map()
- /**
- * @param {QueueEntry} entry
- */
- async processFile(entry) {
- if (this.#pendingFiles.has(entry.cacheKey)) {
- STATS.filesDuplicated++
- } else {
- this.#pendingFiles.set(entry.cacheKey, processFile(entry))
- }
- entry.hash = await this.#pendingFiles.get(entry.cacheKey)
- this.queueFileForWritingHash(entry)
- await this.flushMongoQueuesIfNeeded()
- }
- }
- /**
- * @param {Blob} blob
- * @return {number}
- */
- function estimateBlobSize(blob) {
- let size = blob.getByteLength()
- if (blob.getStringLength()) {
- // approximation for gzip (25 bytes gzip overhead and 20% compression ratio)
- size = 25 + Math.ceil(size * 0.2)
- }
- return size
- }
- async function updateLiveFileTrees() {
- await batchedUpdate(
- projectsCollection,
- { 'overleaf.history.id': { $exists: true } },
- handleLiveTreeBatch,
- { rootFolder: 1, _id: 1, 'overleaf.history.id': 1 },
- {},
- {
- BATCH_RANGE_START,
- BATCH_RANGE_END,
- }
- )
- console.warn('Done updating live projects')
- }
- async function updateDeletedFileTrees() {
- await batchedUpdate(
- deletedProjectsCollection,
- {
- 'deleterData.deletedProjectId': {
- $gt: new ObjectId(BATCH_RANGE_START),
- $lte: new ObjectId(BATCH_RANGE_END),
- },
- 'project.overleaf.history.id': { $exists: true },
- },
- handleDeletedFileTreeBatch,
- {
- 'project.rootFolder': 1,
- 'project._id': 1,
- 'project.overleaf.history.id': 1,
- }
- )
- console.warn('Done updating deleted projects')
- }
- async function main() {
- await loadGlobalBlobs()
- if (process.argv.includes('live')) {
- await updateLiveFileTrees()
- }
- if (process.argv.includes('deleted')) {
- await updateDeletedFileTrees()
- }
- console.warn('Done.')
- }
- try {
- try {
- await main()
- } finally {
- printStats()
- }
- let code = 0
- if (STATS.filesFailed > 0) {
- console.warn('Some files could not be processed, see logs and try again')
- code++
- }
- if (STATS.fileHardDeleted > 0) {
- console.warn(
- 'Some hashes could not be updated as the files were hard-deleted, this should not happen'
- )
- code++
- }
- if (STATS.projectHardDeleted > 0) {
- console.warn(
- 'Some hashes could not be updated as the project was hard-deleted, this should not happen'
- )
- code++
- }
- process.exit(code)
- } catch (err) {
- console.error(err)
- process.exit(1)
- }
|