ソースを参照

add script to finalise broken history-v1 chunks (#34005)

* add script to finalise broken history-v1 chunks

* use history-id instead of project-id

* update project-id to history-id in tests

* silence unwanted event emitter warnings

* fix up test for historyId

GitOrigin-RevId: 58d2a768f1eff296e921e2ed985f6faf3929f619
Brian Gough 2 ヶ月 前
コミット
0da93aaab3

+ 129 - 0
services/history-v1/storage/scripts/finalise_chunk.mjs

@@ -0,0 +1,129 @@
+// Finalise the current chunk for a project and start a new empty chunk
+// whose starting snapshot is the end snapshot of the (now-closed) current
+// chunk.
+//
+// This is intended as a recovery tool for projects whose current chunk has
+// become corrupted in such a way that further changes can no longer be
+// persisted, but where the end snapshot of the current chunk can still be
+// computed.
+
+import logger from '@overleaf/logger'
+import commandLineArgs from 'command-line-args'
+import { Change, Chunk, History, NoOperation } from 'overleaf-editor-core'
+import * as redis from '../lib/redis.js'
+import knex from '../lib/knex.js'
+import knexReadOnly from '../lib/knex_read_only.js'
+import { client as mongoClient } from '../lib/mongodb.js'
+import chunkStore from '../lib/chunk_store/index.js'
+import redisBackend from '../lib/chunk_store/redis.js'
+import { loadGlobalBlobs } from '../lib/blob_store/index.js'
+import { fileURLToPath } from 'node:url'
+import { EventEmitter } from 'node:events'
+
+EventEmitter.defaultMaxListeners = 20
+
+logger.initialize('finalise-chunk')
+
+const optionDefinitions = [
+  { name: 'historyId', type: String },
+  { name: 'dry-run', alias: 'd', type: Boolean },
+]
+const options = commandLineArgs(optionDefinitions)
+const HISTORY_ID = options.historyId
+const DRY_RUN = options['dry-run'] || false
+
+if (!HISTORY_ID) {
+  console.error('Usage: finalise_chunk.mjs --historyId <id> [--dry-run]')
+  process.exit(2)
+}
+
+async function finaliseCurrentChunk(historyId) {
+  // Validates the history id and selects the backend (postgres or mongo).
+  chunkStore.getBackend(historyId)
+
+  await loadGlobalBlobs()
+
+  const currentChunk = await chunkStore.loadLatest(historyId, {
+    persistedOnly: true,
+  })
+  const startVersion = currentChunk.getStartVersion()
+  const endVersion = currentChunk.getEndVersion()
+  const numChanges = currentChunk.getChanges().length
+
+  logger.info(
+    { historyId, startVersion, endVersion, numChanges },
+    'loaded current chunk'
+  )
+
+  if (endVersion === startVersion) {
+    throw new Error(
+      `current chunk for history ${historyId} is already empty (no changes); refusing to create another empty chunk`
+    )
+  }
+
+  let nonPersistedChanges
+  try {
+    nonPersistedChanges = await redisBackend.getNonPersistedChanges(
+      historyId,
+      endVersion
+    )
+  } catch (err) {
+    throw new Error(
+      `unable to read non-persisted changes from redis for history ${historyId}: ${err.message}`
+    )
+  }
+  if (nonPersistedChanges.length > 0) {
+    throw new Error(
+      `history ${historyId} has ${nonPersistedChanges.length} non-persisted change(s) in redis; persist or expire them before running this script`
+    )
+  }
+
+  const endSnapshot = currentChunk.getSnapshot().clone()
+  endSnapshot.applyAll(currentChunk.getChanges())
+
+  // The chunks table has a unique constraint on (doc_id, end_version), so the
+  // new chunk cannot share an end_version with the chunk we are closing. Add a
+  // single NoOperation change to bump end_version by 1 without mutating the
+  // snapshot.
+  const recoveryChange = new Change([new NoOperation()], new Date(), [])
+  const newChunk = new Chunk(
+    new History(endSnapshot, [recoveryChange]),
+    endVersion
+  )
+
+  if (DRY_RUN) {
+    logger.info(
+      { historyId, endVersion },
+      'dry run: would close current chunk and create new empty chunk'
+    )
+    return
+  }
+
+  await chunkStore.create(historyId, newChunk)
+
+  logger.info(
+    { historyId, endVersion },
+    'closed current chunk and created new empty chunk'
+  )
+}
+
+async function main() {
+  try {
+    await finaliseCurrentChunk(HISTORY_ID)
+  } catch (err) {
+    logger.fatal({ err, historyId: HISTORY_ID }, 'failed to finalise chunk')
+    process.exitCode = 1
+  } finally {
+    await redis.disconnect()
+    await mongoClient.close()
+    await knex.destroy()
+    await knexReadOnly.destroy()
+  }
+}
+
+const currentScriptPath = fileURLToPath(import.meta.url)
+if (process.argv[1] === currentScriptPath) {
+  main()
+}
+
+export { finaliseCurrentChunk }

+ 205 - 0
services/history-v1/test/acceptance/js/storage/finalise_chunk.test.js

@@ -0,0 +1,205 @@
+'use strict'
+
+const { promisify } = require('node:util')
+const { execFile } = require('node:child_process')
+const { expect } = require('chai')
+const {
+  Change,
+  AddFileOperation,
+  EditFileOperation,
+  TextOperation,
+  File,
+} = require('overleaf-editor-core')
+const cleanup = require('./support/cleanup')
+const fixtures = require('./support/fixtures')
+const { setupProjectState } = require('./support/redis')
+const chunkStore = require('../../../../storage/lib/chunk_store')
+const persistChanges = require('../../../../storage/lib/persist_changes')
+const knex = require('../../../../storage/lib/knex')
+
+const SCRIPT_PATH = 'storage/scripts/finalise_chunk.mjs'
+
+async function runScript(args) {
+  try {
+    const result = await promisify(execFile)('node', [SCRIPT_PATH, ...args], {
+      encoding: 'utf-8',
+      timeout: 10000,
+      env: { ...process.env, LOG_LEVEL: 'debug' },
+    })
+    return { ...result, status: 0 }
+  } catch (err) {
+    if (typeof err.code !== 'number') {
+      throw err
+    }
+    return { stdout: err.stdout, stderr: err.stderr, status: err.code }
+  }
+}
+
+async function getChunkRows(projectId) {
+  return await knex('chunks')
+    .select('id', 'start_version', 'end_version', 'closed')
+    .where('doc_id', parseInt(projectId, 10))
+    .orderBy('end_version')
+}
+
+describe('finalise_chunk script', function () {
+  before(cleanup.everything)
+
+  let limitsToPersistImmediately
+  before(function () {
+    const farFuture = new Date()
+    farFuture.setTime(farFuture.getTime() + 7 * 24 * 3600 * 1000)
+    limitsToPersistImmediately = {
+      minChangeTimestamp: farFuture,
+      maxChangeTimestamp: farFuture,
+      maxChunkChanges: 100,
+    }
+  })
+
+  beforeEach(async function () {
+    await cleanup.everything()
+    await fixtures.create()
+  })
+
+  describe('with a populated current chunk', function () {
+    let projectId
+    let initialContent
+    let initialEndVersion
+    let initialChunkId
+
+    beforeEach(async function () {
+      projectId = await chunkStore.initializeProject()
+      initialContent = 'Hello world.'
+      const change1 = new Change(
+        [new AddFileOperation('main.tex', File.fromString(initialContent))],
+        new Date(Date.now() - 30000),
+        []
+      )
+      const change2 = new Change(
+        [
+          new EditFileOperation(
+            'main.tex',
+            new TextOperation().retain(initialContent.length).insert(' More.')
+          ),
+        ],
+        new Date(Date.now() - 20000),
+        []
+      )
+      await persistChanges(
+        projectId,
+        [change1, change2],
+        limitsToPersistImmediately,
+        0
+      )
+      const metadata = await chunkStore.getLatestChunkMetadata(projectId)
+      initialEndVersion = metadata.endVersion
+      initialChunkId = metadata.id
+      expect(initialEndVersion).to.equal(2)
+    })
+
+    it('closes the current chunk and creates a new empty chunk', async function () {
+      const result = await runScript(['--historyId', projectId])
+      expect(result.status, result.stderr).to.equal(0)
+
+      const rows = await getChunkRows(projectId)
+      expect(rows).to.have.length(2)
+
+      const oldRow = rows.find(r => r.id.toString() === initialChunkId)
+      expect(oldRow, 'old chunk still in DB').to.exist
+      expect(oldRow.closed).to.equal(true)
+      expect(oldRow.end_version).to.equal(initialEndVersion)
+
+      const newRow = rows.find(r => r.id.toString() !== initialChunkId)
+      expect(newRow, 'new chunk inserted').to.exist
+      expect(newRow.start_version).to.equal(initialEndVersion)
+      expect(newRow.end_version).to.equal(initialEndVersion + 1)
+      expect(newRow.closed).to.equal(false)
+
+      const newChunk = await chunkStore.loadLatest(projectId, {
+        persistedOnly: true,
+      })
+      expect(newChunk.getStartVersion()).to.equal(initialEndVersion)
+      expect(newChunk.getEndVersion()).to.equal(initialEndVersion + 1)
+      // The new chunk holds a single NoOperation change as a placeholder so
+      // its end_version is distinct from the closed chunk's end_version.
+      expect(newChunk.getChanges()).to.have.length(1)
+
+      const newSnapshot = newChunk.getSnapshot()
+      const file = newSnapshot.getFile('main.tex')
+      expect(file, 'main.tex carried over to new chunk snapshot').to.exist
+    })
+
+    it('does not modify state on a dry-run', async function () {
+      const result = await runScript(['--historyId', projectId, '--dry-run'])
+      expect(result.status, result.stderr).to.equal(0)
+
+      const rows = await getChunkRows(projectId)
+      expect(rows).to.have.length(1)
+      expect(rows[0].id.toString()).to.equal(initialChunkId)
+      expect(rows[0].closed).to.equal(false)
+    })
+  })
+
+  describe('safety checks', function () {
+    it('refuses to act on an already-empty current chunk', async function () {
+      const projectId = await chunkStore.initializeProject()
+
+      const result = await runScript(['--historyId', projectId])
+      expect(result.status).to.not.equal(0)
+      expect(result.stdout + result.stderr).to.match(/already empty/)
+
+      const rows = await getChunkRows(projectId)
+      expect(rows).to.have.length(1)
+      expect(rows[0].start_version).to.equal(0)
+      expect(rows[0].end_version).to.equal(0)
+      expect(rows[0].closed).to.equal(false)
+    })
+
+    it('refuses when there are non-persisted changes in redis', async function () {
+      const projectId = await chunkStore.initializeProject()
+      const initialContent = 'Initial.'
+      const initialChange = new Change(
+        [new AddFileOperation('main.tex', File.fromString(initialContent))],
+        new Date(Date.now() - 30000),
+        []
+      )
+      await persistChanges(
+        projectId,
+        [initialChange],
+        limitsToPersistImmediately,
+        0
+      )
+      const metadata = await chunkStore.getLatestChunkMetadata(projectId)
+      const initialChunkId = metadata.id
+
+      const bufferedChange = new Change(
+        [
+          new EditFileOperation(
+            'main.tex',
+            new TextOperation()
+              .retain(initialContent.length)
+              .insert(' Buffered.')
+          ),
+        ],
+        new Date(Date.now() - 10000),
+        []
+      )
+      await setupProjectState(projectId, {
+        persistTime: Date.now() - 1000,
+        headVersion: 2,
+        persistedVersion: 1,
+        changes: [bufferedChange],
+        expireTimeFuture: true,
+      })
+
+      const result = await runScript(['--historyId', projectId])
+      expect(result.status).to.not.equal(0)
+      expect(result.stdout + result.stderr).to.match(/non-persisted/)
+
+      const rows = await getChunkRows(projectId)
+      expect(rows).to.have.length(1)
+      expect(rows[0].id.toString()).to.equal(initialChunkId)
+      expect(rows[0].closed).to.equal(false)
+    })
+  })
+})