Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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<Record<string, unknown>>`
SELECT * FROM projection_projects ORDER BY project_id
`,
},
threads: {
columns: yield* tableColumns("projection_threads"),
rows: yield* sql<Record<string, unknown>>`
SELECT * FROM projection_threads ORDER BY thread_id
`,
},
messages: {
columns: yield* tableColumns("projection_thread_messages"),
rows: yield* sql<Record<string, unknown>>`
SELECT * FROM projection_thread_messages ORDER BY message_id
`,
},
};
}).pipe(Effect.provide(NodeSqliteClient.layer({ filename: fixturePath }))),
);

Expand Down Expand Up @@ -563,21 +595,57 @@ const makeCodexAdapter = (capturedTurns: Ref.Ref<ReadonlyArray<CapturedTurn>>) =
}),
}) 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;
Expand Down Expand Up @@ -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<string>;
readonly rows: ReadonlyArray<Record<string, unknown>>;
}

// Migrations legitimately append columns; preservation means the columns the
// v1 schema already had keep their exact values.
const restrictToSnapshotColumns = (
rows: ReadonlyArray<Record<string, unknown>>,
snapshot: LegacyTableSnapshot,
) => rows.map((row) => Object.fromEntries(snapshot.columns.map((column) => [column, row[column]])));

const messageOrdinals = (projection: OrchestrationV2ThreadProjection) =>
projection.turnItems
.filter(
Expand All @@ -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<ReadonlyArray<CapturedTurn>>([]);
Expand All @@ -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;
Expand Down Expand Up @@ -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),
);
Expand All @@ -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");
Expand All @@ -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);

Expand All @@ -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<Record<string, unknown>>`
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<Record<string, unknown>>`
SELECT * FROM projection_threads ORDER BY thread_id
`;
const legacyMessages = yield* sql<Record<string, unknown>>`
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
Expand All @@ -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),
Expand Down Expand Up @@ -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,
Expand All @@ -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 },
Expand Down Expand Up @@ -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<Record<string, unknown>>`
SELECT * FROM projection_projects ORDER BY project_id
`,
v1Snapshot.projects,
),
v1Snapshot.projects.rows,
);
assert.deepStrictEqual(
restrictToSnapshotColumns(
yield* sql<Record<string, unknown>>`
SELECT * FROM projection_threads ORDER BY thread_id
`,
v1Snapshot.threads,
),
v1Snapshot.threads.rows,
);
assert.deepStrictEqual(
restrictToSnapshotColumns(
yield* sql<Record<string, unknown>>`
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
`;
Expand Down Expand Up @@ -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)),
),
);
Expand Down
Loading