From 2f56c6ddb4e280c74fa5363020913ac9f076f27a Mon Sep 17 00:00:00 2001 From: Alexander Onnikov Date: Wed, 18 Jun 2025 21:00:03 +0700 Subject: [PATCH] fix: restore markup refs corrupted by github (#9284) Signed-off-by: Alexander Onnikov --- dev/tool/package.json | 2 + dev/tool/src/index.ts | 13 +++ dev/tool/src/markup.ts | 244 ++++++++++++++++++++++++++++++++++++++++- 3 files changed, 257 insertions(+), 2 deletions(-) diff --git a/dev/tool/package.json b/dev/tool/package.json index 150920e4e5..eab2a8b930 100644 --- a/dev/tool/package.json +++ b/dev/tool/package.json @@ -156,11 +156,13 @@ "@hcengineering/tags": "^0.6.16", "@hcengineering/task": "^0.6.20", "@hcengineering/text": "^0.6.5", + "@hcengineering/text-core": "^0.6.0", "@hcengineering/text-ydoc": "^0.6.0", "@hcengineering/telegram": "^0.6.21", "@hcengineering/tracker": "^0.6.24", "@hcengineering/collaboration": "^0.6.0", "@hcengineering/datalake": "^0.6.0", + "@hcengineering/retry": "^0.6.0", "@hcengineering/s3": "^0.6.0", "@hcengineering/kvs-client": "^0.6.0", "commander": "^8.1.0", diff --git a/dev/tool/src/index.ts b/dev/tool/src/index.ts index b190d0e11e..fc20e926b1 100644 --- a/dev/tool/src/index.ts +++ b/dev/tool/src/index.ts @@ -120,6 +120,7 @@ import { createRestClient } from '@hcengineering/api-client' import { mkdir, writeFile } from 'fs/promises' import { basename, dirname } from 'path' import { existsSync } from 'fs' +import { restoreMarkupRefs } from './markup' const colorConstants = { colorRed: '\u001b[31m', @@ -2672,6 +2673,18 @@ export function devTool ( client.close() }) + program + .command('restore-markup-refs') + .option('--region ', 'DB region') + .action(async (cmd: { region?: string }) => { + const { dbUrl, txes } = prepareTools() + const region = cmd.region ?? null + + await withStorage(async (adapter) => { + await restoreMarkupRefs(dbUrl, txes, adapter, region) + }) + }) + extendProgram?.(program) process.on('unhandledRejection', (reason, promise) => { diff --git a/dev/tool/src/markup.ts b/dev/tool/src/markup.ts index fec581e44a..ec806797b5 100644 --- a/dev/tool/src/markup.ts +++ b/dev/tool/src/markup.ts @@ -20,12 +20,14 @@ import { yDocCopyXmlField, yDocFromBuffer } from '@hcengineering/collaboration' +import { withRetry } from '@hcengineering/retry' import core, { type Blob, type Doc, type Hierarchy, type MeasureContext, type Ref, + type Tx, type TxCreateDoc, type TxUpdateDoc, DOMAIN_TX, @@ -33,14 +35,32 @@ import core, { type WorkspaceIds, makeCollabId, makeCollabYdocId, - makeDocCollabId + makeDocCollabId, + MeasureMetricsContext, + systemAccountUuid, + isArchivingMode, + isDeletingMode, + type Domain, + type AnyAttribute, + type LowLevelStorage, + type Class, + RateLimiter, + type WorkspaceUuid, + groupByArray } from '@hcengineering/core' import document, { type Document } from '@hcengineering/document' import documents from '@hcengineering/controlled-documents' import { DOMAIN_DOCUMENT } from '@hcengineering/model-document' import { DOMAIN_DOCUMENTS } from '@hcengineering/model-controlled-documents' -import { type StorageAdapter } from '@hcengineering/server-core' +import { getDBClient } from '@hcengineering/postgres' +import { type PipelineFactory, type StorageAdapter, createDummyStorageAdapter } from '@hcengineering/server-core' +import { getAccountClient } from '@hcengineering/server-client' +import { createBackupPipeline, sharedPipelineContextVars } from '@hcengineering/server-pipeline' +import { generateToken } from '@hcengineering/server-token' +import { isEmptyMarkup } from '@hcengineering/text-core' + import { type Db } from 'mongodb' +import { type Sql } from 'postgres' export interface RestoreWikiContentParams { dryRun: boolean @@ -351,3 +371,223 @@ export async function restoreMarkupRefsMongo ( } } } + +export async function restoreMarkupRefs ( + dbUrl: string, + txes: Tx[], + storageAdapter: StorageAdapter, + region: string | null +): Promise { + const token = generateToken(systemAccountUuid, undefined, { service: 'admin', admin: 'true' }) + const ctx = new MeasureMetricsContext('restore-markup-ref', {}) + + const accountClient = getAccountClient(token) + ctx.info('fetching workspaces in region', { region }) + const workspaces = await accountClient.listWorkspaces() + const workspacesById = workspaces.reduce((map, ws) => map.set(ws.uuid, ws), new Map()) + ctx.info('found workspaces in region', { region, count: workspaces.length }) + + const factory: PipelineFactory = createBackupPipeline(ctx, dbUrl, txes, { + externalStorage: createDummyStorageAdapter(), + usePassedCtx: true + }) + + const pg = getDBClient(sharedPipelineContextVars, dbUrl) + const pgClient = await pg.getClient() + + try { + ctx.info('fetching workspaces to fix', { region }) + const targets = await pgClient<{ workspaceId: WorkspaceUuid, _class: Ref> }[]>` + SELECT "workspaceId", _class, COUNT(*) + FROM task + WHERE data->>'description' LIKE '{%' + GROUP BY "workspaceId", _class + ` + + const workspaceTargets = groupByArray(targets, (it) => it.workspaceId) + ctx.info('found workspaces to fix', { region, count: workspaceTargets.size }) + + for (const [wsUuid, targets] of workspaceTargets) { + const workspace = workspacesById.get(wsUuid) + if (workspace === undefined) { + ctx.warn('workspace not found', { wsUuid, region }) + continue + } + + const { uuid, name, url, mode } = workspace + if (isArchivingMode(mode) || isDeletingMode(mode)) { + ctx.warn('skipping workspace', { uuid, name, url, region, mode }) + continue + } + + const classes = targets.map((it) => it._class) + ctx.info('processing workspace', { uuid, name, url, region, classes }) + + try { + const pipeline = await factory(ctx, workspace, (): void => {}, null, null) + + try { + const { hierarchy, lowLevelStorage } = pipeline.context + if (lowLevelStorage === undefined) { + ctx.error('Low level storage not available', { region }) + continue + } + + for (const _class of classes) { + await restoreMarkupRefsForClass( + ctx, + _class, + workspace, + hierarchy, + lowLevelStorage, + storageAdapter, + pgClient + ) + } + } finally { + await pipeline.close() + } + } catch (err: any) { + ctx.error('failed to process workspace', { err, uuid, name, url, region }) + } + } + } finally { + pg.close() + } +} + +async function restoreMarkupRefsForClass ( + ctx: MeasureContext, + _class: Ref>, + wsIds: WorkspaceIds, + hierarchy: Hierarchy, + lowLevelStorage: LowLevelStorage, + storageAdapter: StorageAdapter, + pgClient: Sql +): Promise { + const workspace = wsIds.uuid + const rateLimiter = new RateLimiter(10) + + const domain = hierarchy.findDomain(_class) + if (domain === undefined) return + + const attributes = findCollabAttributes(hierarchy, _class) + if (attributes.length === 0) return + if (hierarchy.isMixin(_class) && attributes.every((p) => p.attributeOf !== _class)) return + + const query = hierarchy.isMixin(_class) ? { [_class]: { $exists: true } } : { _class } + const iterator = await lowLevelStorage.traverse(domain, query) + + let processed = 0 + + try { + while (true) { + const docs = await iterator.next(100) + if (docs === null || docs.length === 0) { + break + } + + for (const doc of docs) { + await rateLimiter.add(async () => { + try { + await withRetry(() => + restoreMarkupRefsForDoc( + ctx, + doc, + domain, + attributes, + wsIds, + hierarchy, + lowLevelStorage, + storageAdapter, + pgClient + ) + ) + processed++ + } catch (err: any) { + ctx.error('failed to restore markup refs', { doc: doc._id, class: doc._class, workspace }) + } + }) + } + + await rateLimiter.waitProcessing() + ctx.info('...', { _class, workspace, processed }) + } + + await rateLimiter.waitProcessing() + ctx.info('done', { _class, workspace, processed }) + } finally { + await iterator.close() + } +} + +async function restoreMarkupRefsForDoc ( + ctx: MeasureContext, + doc: Doc, + domain: Domain, + attributes: AnyAttribute[], + wsIds: WorkspaceIds, + hierarchy: Hierarchy, + lowLevelStorage: LowLevelStorage, + storageAdapter: StorageAdapter, + pgClient: Sql +): Promise { + const workspace = wsIds.uuid + + for (const attribute of attributes) { + const value = hierarchy.isMixin(attribute.attributeOf) + ? ((doc as any)[attribute.attributeOf]?.[attribute.name] as string) + : ((doc as any)[attribute.name] as string) + + if (value == null || value === '') continue + if (!value.startsWith('{')) continue + + const attributeName = hierarchy.isMixin(attribute.attributeOf) + ? `${attribute.attributeOf}.${attribute.name}` + : attribute.name + + if (isEmptyMarkup(value)) { + // If the value is empty, we can just remove the attribute + await lowLevelStorage.rawUpdate(domain, { _id: doc._id }, { [attributeName]: null }) + } else { + // Otherwise try to find the most recent blob for this attribute + const blobId = await findRecentBlobId(pgClient, workspace, doc, attribute) + if (blobId !== undefined) { + // If found, we need to update the document with the blobId + ctx.info('updating blob', { doc: doc._id, class: doc._class, workspace, attribute: attribute.name, blobId }) + await lowLevelStorage.rawUpdate(domain, { _id: doc._id }, { [attributeName]: blobId }) + } else { + // If not found, and the content is not empty, we need to save it and update the document with the blobId + const collabId = makeDocCollabId(doc, attribute.name) + const blobId = await withRetry(() => saveCollabJson(ctx, storageAdapter, wsIds, collabId, value)) + await lowLevelStorage.rawUpdate(domain, { _id: doc._id }, { [attributeName]: blobId }) + ctx.info('uploading blob', { doc: doc._id, class: doc._class, workspace, attribute: attribute.name, blobId }) + } + } + } +} + +function findCollabAttributes (hierarchy: Hierarchy, _class: Ref>): AnyAttribute[] { + const allAttributes = hierarchy.getAllAttributes(_class) + return Array.from(allAttributes.values()).filter((attribute) => { + return hierarchy.isDerived(attribute.type._class, core.class.TypeCollaborativeDoc) + }) +} + +async function findRecentBlobId ( + sql: Sql, + workspace: string, + doc: Doc, + attribute: AnyAttribute +): Promise { + const prefix = `${doc._id}-${attribute.name}-%` + + const [blobId] = await sql` + SELECT name + FROM blob.blob + WHERE workspace = ${workspace} AND name LIKE ${prefix} AND deleted_at IS NULL + ORDER BY created_at DESC + LIMIT 1 + ` + return blobId?.name +}