diff --git a/server/backup/src/backup.ts b/server/backup/src/backup.ts index e80071c773..86763e3c93 100644 --- a/server/backup/src/backup.ts +++ b/server/backup/src/backup.ts @@ -204,10 +204,11 @@ async function verifyDigest ( storage: BackupStorage, snapshots: BackupSnapshot[], domain: Domain -): Promise { +): Promise<{ modified: boolean, modifiedFiles: string[] }> { ctx = ctx.newChild('verify digest', { domain, count: snapshots.length }) ctx.info('verify-digest', { domain, count: snapshots.length }) let modified = false + const modifiedFiles: string[] = [] for (const s of snapshots) { const d = s.domains[domain] if (d === undefined) { @@ -246,35 +247,13 @@ async function verifyDigest ( blobs.set(bname, { doc, buffer: undefined }) } else { blobs.delete(bname) - validDocs.add(bname as Ref) } - } else { - validDocs.add(bname as Ref) } + validDocs.add(bname as Ref) next() }) } else { - const chunks: Buffer[] = [] - stream.on('data', (chunk) => { - chunks.push(chunk) - }) - stream.on('end', () => { - const bf = Buffer.concat(chunks as any) - const d = blobs.get(name) - if (d === undefined) { - blobs.set(name, { doc: undefined, buffer: bf }) - } else { - blobs.delete(name) - const doc = d?.doc as Blob - let sz = doc.size - if (Number.isNaN(sz) || sz !== bf.length) { - sz = bf.length - } - - validDocs.add(name as Ref) - } - next() - }) + next() } stream.resume() // just auto drain the stream }) @@ -310,12 +289,17 @@ async function verifyDigest ( if (storageToRemove.size > 0) { modified = true d.storage = (d.storage ?? []).filter((it) => !storageToRemove.has(it)) + modifiedFiles.push(...Array.from(storageToRemove)) + for (const sf of storageToRemove) { + await storage.delete(sf) + } } - - modified = await updateDigest(d, ctx, storage, validDocs, modified, domain) + let mfiles: string[] = [] + ;({ modified, modifiedFiles: mfiles } = await updateDigest(d, ctx, storage, validDocs, modified, domain)) + modifiedFiles.push(...mfiles) } ctx.end() - return modified + return { modified, modifiedFiles } } async function updateDigest ( @@ -325,8 +309,9 @@ async function updateDigest ( validDocs: Set>, modified: boolean, domain: Domain -): Promise { +): Promise<{ modified: boolean, modifiedFiles: string[] }> { const digestToRemove = new Set() + const modifiedFiles: string[] = [] for (const snapshot of d?.snapshots ?? []) { try { ctx.info('checking', { snapshot }) @@ -365,6 +350,11 @@ async function updateDigest ( const removedCount = parseInt(dataBlob.shift() ?? '0') const removed = dataBlob.splice(0, removedCount) changes.removed = removed as Ref[] + if (addedCount === 0 && removedCount === 0 && updatedCount === 0) { + // Empty digest, need to clean + digestToRemove.add(snapshot) + lmodified = true + } } catch (err: any) { ctx.warn('failed during processing of snapshot file, it will be skipped', { snapshot }) digestToRemove.add(snapshot) @@ -373,17 +363,22 @@ async function updateDigest ( if (lmodified) { modified = true - // Store changes without missing files - await writeChanges(storage, snapshot, changes) + if (digestToRemove.has(snapshot)) { + await storage.delete(snapshot) // No need for digest, lets' remove it + } else { + // Store changes without missing files + await writeChanges(storage, snapshot, changes) + } } } catch (err: any) { digestToRemove.add(snapshot) + modifiedFiles.push(snapshot) ctx.error('digest is broken, will do full backup for', { domain, err: err.message, snapshot }) modified = true } } d.snapshots = (d.snapshots ?? []).filter((it) => !digestToRemove.has(it)) - return modified + return { modified, modifiedFiles } } async function write (chunk: any, stream: Writable): Promise { @@ -846,6 +841,8 @@ export async function backup ( backupInfo.lastTxId = '' // Clear until full backup will be complete + const recheckSizes: string[] = [] + const snapshot: BackupSnapshot = { date: Date.now(), domains: {} @@ -1062,51 +1059,65 @@ export async function backup ( if (d == null) { continue } - const { docs, modified } = await verifyDocsFromSnapshot(ctx, domain, d, s, storage, same) - if (modified) { - changed++ - } - const batchSize = 200 let needRetrieve: Ref[] = [] - for (let i = 0; i < docs.length; i += batchSize) { - const part = docs.slice(i, i + batchSize) - const serverDocs = await connection.loadDocs( - domain, - part.map((it) => it._id) - ) - const smap = toIdMap(serverDocs) - for (const localDoc of part) { - if (TxProcessor.isExtendsCUD(localDoc._class)) { - const tx = localDoc as TxCUD - if (tx.objectSpace == null) { - tx.objectSpace = core.space.Workspace + const batchSize = 200 + const { modified, modifiedFiles } = await verifyDocsFromSnapshot( + ctx, + domain, + d, + s, + storage, + same, + async (docs) => { + const serverDocs = await connection.loadDocs( + domain, + docs.map((it) => it._id) + ) + const smap = toIdMap(serverDocs) + for (const localDoc of docs) { + if (TxProcessor.isExtendsCUD(localDoc._class)) { + const tx = localDoc as TxCUD + if (tx.objectSpace == null) { + tx.objectSpace = core.space.Workspace + } } - } - const serverDoc = smap.get(localDoc._id) - if (serverDoc === undefined) { - // We do not have a doc on server already, ignore it. - } else { - const { '%hash%': _h1, ...dData } = localDoc as any - const { '%hash%': _h2, ...sData } = serverDoc as any + const serverDoc = smap.get(localDoc._id) + if (serverDoc === undefined) { + // We do not have a doc on server already, ignore it. + } else { + const { '%hash%': _h1, ...dData } = localDoc as any + const { '%hash%': _h2, ...sData } = serverDoc as any - const dsame = deepEqual(dData, sData) - if (!dsame) { - needRetrieve.push(localDoc._id) - changes.updated.set(localDoc._id, same.get(localDoc._id) ?? '') - // Docs are not same - if (needRetrieve.length > 200) { - needRetrieveChunks.push(needRetrieve) - needRetrieve = [] + const dsame = deepEqual(dData, sData) + if (!dsame) { + needRetrieve.push(localDoc._id) + changes.updated.set(localDoc._id, same.get(localDoc._id) ?? '') + // Docs are not same + if (needRetrieve.length > 200) { + needRetrieveChunks.push(needRetrieve) + needRetrieve = [] + } } } } - } + }, + batchSize + ) + if (modified) { + changed++ + recheckSizes.push(...modifiedFiles) } if (needRetrieve.length > 0) { needRetrieveChunks.push(needRetrieve) needRetrieve = [] } } + // We need to retrieve all documents from same not matched + const sameArray: Ref[] = Array.from(same.keys()) + while (sameArray.length > 0) { + const docs = sameArray.splice(0, 200) + needRetrieveChunks.push(docs) + } } else { same.clear() } @@ -1144,8 +1155,10 @@ export async function backup ( docs = await ctx.with('<<<< load-docs', {}, async () => await connection.loadDocs(domain, needRetrieve)) lastSize = docs.reduce((p, it) => p + estimateDocSize(it), 0) if (docs.length !== needRetrieve.length) { - const nr = new Set(docs.map((it) => it._id)) - ctx.error('failed to retrieve all documents', { missing: needRetrieve.filter((it) => !nr.has(it)) }) + ctx.error('failed to retrieve all documents', { + docsLen: docs.length, + needRetrieve: needRetrieve.length + }) } ops++ } catch (err: any) { @@ -1416,50 +1429,12 @@ export async function backup ( } result.result = true - const sizeFile = 'backup.size.gz' - - let sizeInfo: Record = {} - - if (await storage.exists(sizeFile)) { - sizeInfo = JSON.parse(gunzipSync(new Uint8Array(await storage.loadFile(sizeFile))).toString()) - } - let processed = 0 - if (!canceled()) { backupInfo.lastTxId = lastTx?._id ?? '0' // We could store last tx, since full backup is complete await storage.writeFile(infoFile, gzipSync(JSON.stringify(backupInfo, undefined, 2), { level: defaultLevel })) } - const addFileSize = async (file: string | undefined | null): Promise => { - if (file != null) { - const sz = sizeInfo[file] - const fileSize = sz ?? (await storage.stat(file)) - if (sz === undefined) { - sizeInfo[file] = fileSize - processed++ - if (processed % 10 === 0) { - ctx.info('Calculate size processed', { processed, size: Math.round(result.backupSize / (1024 * 1024)) }) - } - } - result.backupSize += fileSize - } - } - - // Let's calculate data size for backup - for (const sn of backupInfo.snapshots) { - for (const [, d] of Object.entries(sn.domains)) { - await addFileSize(d.snapshot) - for (const snp of d.snapshots ?? []) { - await addFileSize(snp) - } - for (const snp of d.storage ?? []) { - await addFileSize(snp) - } - } - } - await addFileSize(infoFile) - - await storage.writeFile(sizeFile, gzipSync(JSON.stringify(sizeInfo, undefined, 2), { level: defaultLevel })) + await rebuildSizeInfo(storage, recheckSizes, ctx, result, backupInfo, infoFile) return result } catch (err: any) { @@ -1480,6 +1455,60 @@ export async function backup ( } } +async function rebuildSizeInfo ( + storage: BackupStorage, + recheckSizes: string[], + ctx: MeasureContext, + result: BackupResult, + backupInfo: BackupInfo, + infoFile: string +): Promise { + const sizeFile = 'backup.size.gz' + + let sizeInfo: Record = {} + + if (await storage.exists(sizeFile)) { + sizeInfo = JSON.parse(gunzipSync(new Uint8Array(await storage.loadFile(sizeFile))).toString()) + } + let processed = 0 + + for (const file of recheckSizes) { + // eslint-disable-next-line @typescript-eslint/no-dynamic-delete + delete sizeInfo[file] + } + + const addFileSize = async (file: string | undefined | null): Promise => { + if (file != null) { + const sz = sizeInfo[file] + const fileSize = sz ?? (await storage.stat(file)) + if (sz === undefined) { + sizeInfo[file] = fileSize + processed++ + if (processed % 10 === 0) { + ctx.info('Calculate size processed', { processed, size: Math.round(result.backupSize / (1024 * 1024)) }) + } + } + result.backupSize += fileSize + } + } + + // Let's calculate data size for backup + for (const sn of backupInfo.snapshots) { + for (const [, d] of Object.entries(sn.domains)) { + await addFileSize(d.snapshot) + for (const snp of d.snapshots ?? []) { + await addFileSize(snp) + } + for (const snp of d.storage ?? []) { + await addFileSize(snp) + } + } + } + await addFileSize(infoFile) + + await storage.writeFile(sizeFile, gzipSync(JSON.stringify(sizeInfo, undefined, 2), { level: defaultLevel })) +} + /** * @public */ @@ -2246,11 +2275,14 @@ async function verifyDocsFromSnapshot ( d: DomainData, s: BackupSnapshot, storage: BackupStorage, - digest: Map, string> -): Promise<{ docs: Doc[], modified: boolean }> { - const result: Doc[] = [] + digest: Map, string>, + verify: (docs: Doc[]) => Promise, + chunkSize: number +): Promise<{ modified: boolean, modifiedFiles: string[] }> { + let result: Doc[] = [] const storageToRemove = new Set() const validDocs = new Set>() + const modifiedFiles: string[] = [] if (digest.size > 0) { const sDigest = await loadDigest(ctx, storage, [s], domain) const requiredDocs = new Map(Array.from(sDigest.entries()).filter(([it]) => digest.has(it))) @@ -2269,7 +2301,8 @@ async function verifyDocsFromSnapshot ( ex.on('entry', (headers, stream, next) => { const name = headers.name ?? '' // We found blob data - if (name.endsWith('.json') && requiredDocs.has(name.substring(0, name.length - 5) as Ref)) { + const rdoc = name.substring(0, name.length - 5) as Ref + if (name.endsWith('.json') && requiredDocs.has(rdoc)) { const chunks: Buffer[] = [] const bname = name.substring(0, name.length - 5) stream.on('data', (chunk) => { @@ -2289,11 +2322,19 @@ async function verifyDocsFromSnapshot ( // Skip blob validDocs.add(bname as Ref) } else { - ;(doc as any)['%hash%'] = digest.get(doc._id) - digest.delete(bname as Ref) + ;(doc as any)['%hash%'] = digest.get(rdoc) + digest.delete(rdoc) result.push(doc) validDocs.add(bname as Ref) - next() + + if (result.length > chunkSize) { + void verify(result).then(() => { + result = [] + next() + }) + } else { + next() + } } }) } else { @@ -2321,6 +2362,9 @@ async function verifyDocsFromSnapshot ( }) await endPromise + if (result.length > 0) { + await verify(result) + } } catch (err: any) { storageToRemove.add(sf) ctx.error('failed to processing', { storageFile: sf }) @@ -2330,11 +2374,17 @@ async function verifyDocsFromSnapshot ( } let modified = false if (storageToRemove.size > 0) { + modifiedFiles.push(...Array.from(storageToRemove)) d.storage = (d.storage ?? []).filter((it) => !storageToRemove.has(it)) + for (const sf of storageToRemove) { + await storage.delete(sf) + } modified = true } - modified = await updateDigest(d, ctx, storage, validDocs, modified, domain) - return { docs: result, modified } + let smodifiedFiles: string[] = [] + ;({ modified, modifiedFiles: smodifiedFiles } = await updateDigest(d, ctx, storage, validDocs, modified, domain)) + modifiedFiles.push(...smodifiedFiles) + return { modified, modifiedFiles } } /** @@ -2686,7 +2736,7 @@ function migradeBlobData (blob: Blob, etag: string): string { * @public */ export async function checkBackupIntegrity (ctx: MeasureContext, storage: BackupStorage): Promise { - console.log('starting backup compaction') + console.log('check backup integrity') try { let backupInfo: BackupInfo @@ -2705,6 +2755,8 @@ export async function checkBackupIntegrity (ctx: MeasureContext, storage: Backup return } + const recheckSizes: string[] = [] + const domains: Domain[] = [] for (const sn of backupInfo.snapshots) { for (const d of Object.keys(sn.domains)) { @@ -2717,13 +2769,23 @@ export async function checkBackupIntegrity (ctx: MeasureContext, storage: Backup for (const domain of domains) { console.log('checking domain...', domain) - if (await verifyDigest(ctx, storage, backupInfo.snapshots, domain)) { + const { modified: mm, modifiedFiles } = await verifyDigest(ctx, storage, backupInfo.snapshots, domain) + if (mm) { + recheckSizes.push(...modifiedFiles) modified = true } } if (modified) { await storage.writeFile(infoFile, gzipSync(JSON.stringify(backupInfo, undefined, 2), { level: defaultLevel })) } + + const bresult: BackupResult = { + backupSize: 0, + blobsSize: 0, + dataSize: 0, + result: true + } + await rebuildSizeInfo(storage, recheckSizes, ctx, bresult, backupInfo, infoFile) } catch (err: any) { console.error(err) } finally { diff --git a/server/workspace-service/src/service.ts b/server/workspace-service/src/service.ts index 4974a0233e..14d9f6efcc 100644 --- a/server/workspace-service/src/service.ts +++ b/server/workspace-service/src/service.ts @@ -619,7 +619,7 @@ export class WorkspaceWorker { 50000, ['blob'], sharedPipelineContextVars, - true, + true, // Do a full check (_p: number) => { if (progress !== Math.round(_p)) { progress = Math.round(_p)