diff --git a/apps/server/src/orchestration-v2/legacy/LegacyV1Cutover.integration.test.ts b/apps/server/src/orchestration-v2/legacy/LegacyV1Cutover.integration.test.ts index 7ce2e1fe2542..9b98b844c3f1 100644 --- a/apps/server/src/orchestration-v2/legacy/LegacyV1Cutover.integration.test.ts +++ b/apps/server/src/orchestration-v2/legacy/LegacyV1Cutover.integration.test.ts @@ -1,4 +1,5 @@ import { assert, describe, it } from "@effect/vitest"; +import * as NodeCrypto from "node:crypto"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { CommandId, @@ -412,6 +413,37 @@ const seedV1Database = (fixturePath: string, workspace: string) => }); yield* sql`PRAGMA wal_checkpoint(TRUNCATE);`; + + // Deterministic pre-migration snapshot of the legacy recovery tables; + // post-boot assertions compare preserved column contents, not counts. + // Migrations may append columns, so the comparison is restricted to the + // column set recorded here. + const tableColumns = Effect.fn("seedV1Database.tableColumns")(function* (table: string) { + const info = yield* sql<{ + readonly name: string; + }>`SELECT name FROM pragma_table_info(${table})`; + return info.map((row) => row.name); + }); + return { + projects: { + columns: yield* tableColumns("projection_projects"), + rows: yield* sql>` + SELECT * FROM projection_projects ORDER BY project_id + `, + }, + threads: { + columns: yield* tableColumns("projection_threads"), + rows: yield* sql>` + SELECT * FROM projection_threads ORDER BY thread_id + `, + }, + messages: { + columns: yield* tableColumns("projection_thread_messages"), + rows: yield* sql>` + SELECT * FROM projection_thread_messages ORDER BY message_id + `, + }, + }; }).pipe(Effect.provide(NodeSqliteClient.layer({ filename: fixturePath }))), ); @@ -563,21 +595,57 @@ const makeCodexAdapter = (capturedTurns: Ref.Ref>) = }), }) satisfies ProviderAdapterV2Shape; -const waitForIdle = Effect.fn("LegacyV1Cutover.waitForIdle")(function* (threadId: ThreadId) { - const orchestrator = yield* Orchestrator.OrchestratorV2; - for (let attempt = 0; attempt < 1_000; attempt += 1) { - const projection = yield* orchestrator.getThreadProjection(threadId); - if ( - projection.runs.every( - (run) => !["queued", "starting", "running", "waiting"].includes(run.status), - ) - ) { - return projection; - } - yield* Effect.sleep("5 millis"); - } - return yield* Effect.die(new Error("Cutover test timed out waiting for idle")); -}); +// Await the dispatched run's terminal event on the persisted domain-event stream +// instead of polling the projection, then drain remaining outbox effects. +const dispatchMessageAwaitingRun = Effect.fn("LegacyV1Cutover.dispatchMessageAwaitingRun")( + function* (input: { + readonly threadId: ThreadId; + readonly messageId: string; + readonly commandId: string; + readonly text: string; + }) { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const worker = yield* EffectWorker.OrchestrationEffectWorkerV2; + const eventSink = yield* EventSink.EventSinkV2; + const messageId = MessageId.make(input.messageId); + // Capture the cursor before dispatch so even a late subscriber replays + // the terminal event instead of starting its live tail after completion. + const afterSequence = yield* eventSink.latestSequence(); + const terminal = eventSink.stream({ afterSequence }).pipe( + Stream.map((stored) => stored.event), + Stream.filter((event) => event.type === "run.updated"), + Stream.filter( + (event) => + ["completed", "failed", "interrupted", "cancelled", "rolled_back"].includes( + event.payload.status, + ) && + event.payload.threadId === input.threadId && + event.payload.userMessageId === messageId, + ), + Stream.take(1), + Stream.runForEach((event) => + Effect.sync(() => + assert.equal(event.payload.status, "completed", `Run for ${messageId} did not complete`), + ), + ), + ); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(input.commandId), + threadId: input.threadId, + messageId, + text: input.text, + attachments: [], + dispatchMode: { type: "start_immediately" }, + }); + yield* worker.drain(); + yield* terminal; + yield* worker.drain(); + return yield* orchestrator.getThreadProjection(input.threadId); + }, +); const makeBootLayer = (input: { readonly name: string; @@ -634,6 +702,25 @@ const makeCapturingLogger = (logs: CapturedLog[]) => }); }); +const sha256File = Effect.fn("LegacyV1Cutover.sha256File")(function* (filePath: string) { + const fs = yield* FileSystem.FileSystem; + return NodeCrypto.createHash("sha256") + .update(yield* fs.readFile(filePath)) + .digest("hex"); +}); + +interface LegacyTableSnapshot { + readonly columns: ReadonlyArray; + readonly rows: ReadonlyArray>; +} + +// Migrations legitimately append columns; preservation means the columns the +// v1 schema already had keep their exact values. +const restrictToSnapshotColumns = ( + rows: ReadonlyArray>, + snapshot: LegacyTableSnapshot, +) => rows.map((row) => Object.fromEntries(snapshot.columns.map((column) => [column, row[column]]))); + const messageOrdinals = (projection: OrchestrationV2ThreadProjection) => projection.turnItems .filter( @@ -658,7 +745,8 @@ describe("orchestration v2 legacy v1 cutover", () => { const copyPath = path.join(stateDir, "userdata", "state.sqlite"); yield* fs.makeDirectory(path.join(stateDir, "userdata"), { recursive: true }); - yield* seedV1Database(fixturePath, workspace); + const v1Snapshot = yield* seedV1Database(fixturePath, workspace); + const fixtureHashBefore = yield* sha256File(fixturePath); yield* fs.copyFile(fixturePath, copyPath); const capturedTurns = yield* Ref.make>([]); @@ -672,7 +760,6 @@ describe("orchestration v2 legacy v1 cutover", () => { const importer = yield* LegacyV1ThreadImporter.LegacyV1ThreadImporter; const maintenance = yield* ProjectionMaintenance.ProjectionMaintenanceV2; const projections = yield* ProjectionStore.ProjectionStoreV2; - const orchestrator = yield* Orchestrator.OrchestratorV2; assert.equal(yield* importer.pendingThreadCount, ALL_THREADS.length); const shellImport = yield* importer.reconcileShells; @@ -811,6 +898,22 @@ describe("orchestration v2 legacy v1 cutover", () => { `; assert.isNull(snoozedImportRow[0]?.transcript_imported_at); + // A failed hydration leaves the legacy rows and pending marker + // available for retry in a fresh runtime after restart. + yield* sql` + CREATE TRIGGER fail_snoozed_hydration + BEFORE INSERT ON orchestration_events + WHEN NEW.event_id = 'migration:v1:message:message:cutover:snoozed:2' + BEGIN + SELECT RAISE(ABORT, 'injected cutover hydration failure'); + END + `; + const hydrationError = yield* importer + .ensureTranscript(ThreadId.make(SNOOZED_THREAD)) + .pipe(Effect.flip); + assert.equal(hydrationError._tag, "LegacyV1ThreadImportError"); + assert.equal(yield* importer.pendingThreadCount, 1); + const interrupted = yield* projections.getThreadProjection( ThreadId.make(INTERRUPTED_THREAD), ); @@ -821,18 +924,12 @@ describe("orchestration v2 legacy v1 cutover", () => { // First continuation of a migrated thread: a fresh provider session // receives the newest transcript suffix inside the 32k handoff. - yield* orchestrator.dispatch({ - type: "message.dispatch", - createdBy: "user", - creationSource: "web", - commandId: CommandId.make("command:cutover:continue"), + const continued = yield* dispatchMessageAwaitingRun({ threadId: longThreadId, - messageId: MessageId.make("message:cutover:long:continuation"), + messageId: "message:cutover:long:continuation", + commandId: "command:cutover:continue", text: CONTINUATION_PROMPT, - attachments: [], - dispatchMode: { type: "start_immediately" }, }); - const continued = yield* waitForIdle(longThreadId); const runs = continued.runs.filter((run) => run.ordinal >= 1); assert.equal(runs.at(-1)?.status, "completed"); @@ -852,18 +949,12 @@ describe("orchestration v2 legacy v1 cutover", () => { // The next continuation reuses the provider thread without a new // handoff: the imported context is only reissued until a v2 run // completes. - yield* orchestrator.dispatch({ - type: "message.dispatch", - createdBy: "user", - creationSource: "web", - commandId: CommandId.make("command:cutover:continue:again"), + const continuedAgain = yield* dispatchMessageAwaitingRun({ threadId: longThreadId, - messageId: MessageId.make("message:cutover:long:continuation:2"), + messageId: "message:cutover:long:continuation:2", + commandId: "command:cutover:continue:again", text: "One more.", - attachments: [], - dispatchMode: { type: "start_immediately" }, }); - const continuedAgain = yield* waitForIdle(longThreadId); assert.equal(continuedAgain.contextHandoffs.length, 1); assert.equal((yield* Ref.get(capturedTurns)).length, 2); @@ -886,11 +977,14 @@ describe("orchestration v2 legacy v1 cutover", () => { FROM orchestration_v2_legacy_imports ORDER BY thread_id `; - const legacyMessageCount = yield* sql<{ readonly count: number }>` - SELECT COUNT(*) AS count FROM projection_thread_messages + const legacyProjects = yield* sql>` + SELECT * FROM projection_projects ORDER BY project_id `; - const legacyThreadCount = yield* sql<{ readonly count: number }>` - SELECT COUNT(*) AS count FROM projection_threads + const legacyThreads = yield* sql>` + SELECT * FROM projection_threads ORDER BY thread_id + `; + const legacyMessages = yield* sql>` + SELECT * FROM projection_thread_messages ORDER BY message_id `; const recordedMigration41 = yield* sql<{ readonly name: string }>` SELECT name FROM effect_sql_migrations WHERE migration_id = 41 @@ -901,8 +995,9 @@ describe("orchestration v2 legacy v1 cutover", () => { return { migrationEventCount: migrationEventCount[0]?.count ?? 0, importRows, - legacyMessageCount: legacyMessageCount[0]?.count ?? 0, - legacyThreadCount: legacyThreadCount[0]?.count ?? 0, + legacyProjects, + legacyThreads, + legacyMessages, longProjection: continuedAgain, migration41Name: recordedMigration41[0]?.name ?? null, authSessionColumnNames: authSessionColumns.map((column) => column.name), @@ -939,6 +1034,21 @@ describe("orchestration v2 legacy v1 cutover", () => { assert.notInclude(firstBoot.authSessionColumnNames, "client_surface"); assert.notInclude(firstBoot.authSessionColumnNames, "client_app_version"); + // Full row contents of the legacy recovery tables survive migration + // and continuation unchanged, not just their counts. + assert.deepStrictEqual( + restrictToSnapshotColumns(firstBoot.legacyProjects, v1Snapshot.projects), + v1Snapshot.projects.rows, + ); + assert.deepStrictEqual( + restrictToSnapshotColumns(firstBoot.legacyThreads, v1Snapshot.threads), + v1Snapshot.threads.rows, + ); + assert.deepStrictEqual( + restrictToSnapshotColumns(firstBoot.legacyMessages, v1Snapshot.messages), + v1Snapshot.messages.rows, + ); + assert.equal(firstBoot.importRows.length, ALL_THREADS.length); const unhydratedRows = firstBoot.importRows.filter( (row) => row.transcript_imported_at === null, @@ -962,6 +1072,7 @@ describe("orchestration v2 legacy v1 cutover", () => { importedThreadCount: 0, importedMessageCount: 0, }); + yield* sql`DROP TRIGGER fail_snoozed_hydration`; assert.deepStrictEqual( yield* importer.ensureTranscript(ThreadId.make(SNOOZED_THREAD)), { importedThreadCount: 1, importedMessageCount: 1 }, @@ -1002,15 +1113,35 @@ describe("orchestration v2 legacy v1 cutover", () => { ); // The v1 projection tables remain the untouched, read-only - // recovery source after migration and restart. - const legacyMessageCount = yield* sql<{ readonly count: number }>` - SELECT COUNT(*) AS count FROM projection_thread_messages - `; - const legacyThreadCount = yield* sql<{ readonly count: number }>` - SELECT COUNT(*) AS count FROM projection_threads - `; - assert.equal(legacyMessageCount[0]?.count, firstBoot.legacyMessageCount); - assert.equal(legacyThreadCount[0]?.count, firstBoot.legacyThreadCount); + // recovery source after migration and restart: preserved column + // contents still equal the pre-migration snapshot. + assert.deepStrictEqual( + restrictToSnapshotColumns( + yield* sql>` + SELECT * FROM projection_projects ORDER BY project_id + `, + v1Snapshot.projects, + ), + v1Snapshot.projects.rows, + ); + assert.deepStrictEqual( + restrictToSnapshotColumns( + yield* sql>` + SELECT * FROM projection_threads ORDER BY thread_id + `, + v1Snapshot.threads, + ), + v1Snapshot.threads.rows, + ); + assert.deepStrictEqual( + restrictToSnapshotColumns( + yield* sql>` + SELECT * FROM projection_thread_messages ORDER BY message_id + `, + v1Snapshot.messages, + ), + v1Snapshot.messages.rows, + ); const quickCheck = yield* sql<{ readonly quick_check: string }>` PRAGMA quick_check `; @@ -1039,6 +1170,10 @@ describe("orchestration v2 legacy v1 cutover", () => { assert.isTrue( boot2Logs.some((log) => String(log.message).includes("migration history diverges")), ); + + // The seeded source file is the untouched recovery copy: migration + // only ever ran against the file copied out of it. + assert.equal(yield* sha256File(fixturePath), fixtureHashBefore); }).pipe(Effect.provide(NodeServices.layer)), ), );