From 34c6c03058b527e4e823ad80e5cc9414e5ccd1f8 Mon Sep 17 00:00:00 2001 From: "wizzoapp[bot]" <254688279+wizzoapp[bot]@users.noreply.github.com> Date: Wed, 22 Jul 2026 11:11:21 +0100 Subject: [PATCH 1/4] fix(server): backfill delivered terminal claims (ADA-194) --- .../Layers/ChildThreadCoordinator.test.ts | 194 +++++++++ apps/server/src/persistence/Migrations.ts | 2 + ...ingLocalSubagentTerminalDeliveries.test.ts | 372 ++++++++++++++++++ ...existingLocalSubagentTerminalDeliveries.ts | 165 ++++++++ 4 files changed, 733 insertions(+) create mode 100644 apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts create mode 100644 apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts diff --git a/apps/server/src/orchestration/Layers/ChildThreadCoordinator.test.ts b/apps/server/src/orchestration/Layers/ChildThreadCoordinator.test.ts index e6cac1a3789..35744f8866a 100644 --- a/apps/server/src/orchestration/Layers/ChildThreadCoordinator.test.ts +++ b/apps/server/src/orchestration/Layers/ChildThreadCoordinator.test.ts @@ -35,6 +35,7 @@ import { SqlClient } from "effect/unstable/sql/SqlClient"; import { afterEach, describe, expect, it } from "vite-plus/test"; import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; +import { backfillPreexistingLocalSubagentTerminalDeliveries } from "../../persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts"; import { PendingDispatchRepositoryLive } from "../../persistence/Layers/PendingDispatches.ts"; import { PendingDispatchId, @@ -414,6 +415,8 @@ describe("ChildThreadCoordinator", () => { readonly claimedSequence: number; readonly terminalKind?: "completed" | "failed" | "killed" | "archived"; }>; + /** Run migration 054's backfill after seeding projection rows but before coordinator boot. */ + readonly backfillPreexistingTerminalDeliveriesBeforeStart?: boolean; /** subagent_promoted_children rows inserted BEFORE start() (simulated restart). */ readonly seedPromotedChildren?: ReadonlyArray<{ readonly childThreadId: ThreadId; @@ -942,6 +945,127 @@ describe("ChildThreadCoordinator", () => { ); } } + if (input?.backfillPreexistingTerminalDeliveriesBeforeStart) { + for (const event of input.persistedEvents ?? []) { + const sequence = (event as { readonly sequence?: number }).sequence; + if (sequence === undefined) continue; + await activeRuntime.runPromise( + Effect.flatMap( + Effect.service(SqlClient), + (sql) => sql` + INSERT INTO orchestration_events ( + sequence, + event_id, + aggregate_kind, + stream_id, + stream_version, + event_type, + occurred_at, + actor_kind, + payload_json, + metadata_json + ) + VALUES ( + ${sequence}, + ${event.eventId}, + ${event.aggregateKind}, + ${event.aggregateId}, + ${sequence}, + ${event.type}, + ${event.occurredAt}, + ${"server"}, + ${JSON.stringify(event.payload)}, + ${"{}"} + ) + `, + ), + ); + } + for (const [rowIndex, row] of (input.seedChildRows ?? []).entries()) { + const state = threadStates.get(row.threadId); + if (state === undefined) continue; + const latestTurn = state.shell.latestTurn; + await activeRuntime.runPromise( + Effect.flatMap(Effect.service(SqlClient), (sql) => + Effect.gen(function* () { + yield* sql` + UPDATE projection_threads + SET + latest_turn_id = ${latestTurn?.turnId ?? null}, + archived_at = ${state.shell.archivedAt ?? null}, + deleted_at = ${state.detail.deletedAt ?? null} + WHERE thread_id = ${row.threadId} + `; + if (latestTurn !== null) { + yield* sql` + INSERT INTO projection_turns ( + thread_id, + turn_id, + state, + requested_at, + started_at, + completed_at, + checkpoint_files_json + ) + VALUES ( + ${row.threadId}, + ${latestTurn.turnId}, + ${latestTurn.state}, + ${latestTurn.requestedAt}, + ${latestTurn.startedAt}, + ${latestTurn.completedAt}, + ${"[]"} + ) + `; + } + if (state.shell.session !== null) { + yield* sql` + INSERT INTO projection_thread_sessions ( + thread_id, + status, + active_turn_id, + updated_at + ) + VALUES ( + ${row.threadId}, + ${state.shell.session.status}, + ${state.shell.session.activeTurnId}, + ${state.shell.session.updatedAt} + ) + `; + } + yield* sql` + INSERT INTO orchestration_events ( + sequence, + event_id, + aggregate_kind, + stream_id, + stream_version, + event_type, + occurred_at, + actor_kind, + payload_json, + metadata_json + ) + VALUES ( + ${10_000 + rowIndex}, + ${`terminal-backfill-delivered-${row.threadId}`}, + ${"thread"}, + ${row.parentThreadId}, + ${10_000 + rowIndex}, + ${"thread.message-sent"}, + ${now}, + ${"server"}, + ${`{"role":"system","text":"[sub-agent ${row.threadId} completed]"}`}, + ${"{}"} + ) + `; + }), + ), + ); + } + await activeRuntime.runPromise(backfillPreexistingLocalSubagentTerminalDeliveries()); + } if (input?.seedPromotedChildren) { for (const row of input.seedPromotedChildren) { await activeRuntime.runPromise( @@ -1303,6 +1427,76 @@ describe("ChildThreadCoordinator", () => { ).toEqual([]); }; + it("backfills a pre-existing terminal child and emits zero wakes across restarts", async () => { + const child = ThreadId.make("terminal-delivery-preexisting-backfill-child"); + const parent = ThreadId.make("terminal-delivery-preexisting-backfill-parent"); + const childTurn = TurnId.make("terminal-delivery-preexisting-backfill-turn"); + const parentState = makeThreadState({ + threadId: parent, + latestTurn: makeLatestTurn("completed", TurnId.make("preexisting-backfill-parent-turn")), + session: makeSession(parent, "ready"), + }); + const childState = makeThreadState({ + threadId: child, + parentThreadId: parent, + latestTurn: makeLatestTurn("completed", childTurn), + session: makeSession(child, "ready"), + assistantText: "completed before durable claims deployed", + }); + const persistedEvents = [ + turnStartRequestedEvent(child), + sessionSetEvent(child, "running", childTurn), + turnDiffEvent(child, "ready", childTurn), + ].map((event, index) => ({ ...event, sequence: index + 1 })); + const harnessInput = { + threads: [parentState, childState], + seedChildRows: [{ threadId: child, parentThreadId: parent }], + persistedEvents, + }; + + let harness = await createHarness({ + ...harnessInput, + backfillPreexistingTerminalDeliveriesBeforeStart: true, + }); + const [claim] = await harness.listTerminalDeliveries(); + expect(claim).toMatchObject({ + childThreadId: String(child), + parentThreadId: String(parent), + claimId: `migration:054:${child}`, + claimedSequence: 3, + terminalKind: "completed", + }); + + let wakeCount = harness.dispatched.filter( + (command) => command.type === "thread.turn.start" && command.threadId === parent, + ).length; + expect(await harness.listPendingDispatches()).toEqual([]); + + const simulatedRestarts = 3; + for (let restart = 0; restart < simulatedRestarts; restart += 1) { + await disposeActiveHarness(); + harness = await createHarness({ + ...harnessInput, + seedTerminalDeliveries: [ + { + childThreadId: child, + parentThreadId: parent, + claimId: claim!.claimId, + claimedAt: claim!.claimedAt, + claimedSequence: claim!.claimedSequence, + terminalKind: claim!.terminalKind, + }, + ], + }); + wakeCount += harness.dispatched.filter( + (command) => command.type === "thread.turn.start" && command.threadId === parent, + ).length; + expect(await harness.listPendingDispatches()).toEqual([]); + } + + expect(wakeCount).toBe(0); + }); + for (const terminalKind of ["completed", "failed", "killed", "archived"] as const) { it(`delivers one ${terminalKind} parent wake across repeated restart replay`, async () => { const child = ThreadId.make(`terminal-delivery-${terminalKind}-child`); diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index 3b46d011a2b..67f1054c37a 100644 --- a/apps/server/src/persistence/Migrations.ts +++ b/apps/server/src/persistence/Migrations.ts @@ -65,6 +65,7 @@ import Migration0050 from "./Migrations/050_ProjectionProjectDataAudience.ts"; import Migration0051 from "./Migrations/051_AuthAudienceCeilings.ts"; import Migration0052 from "./Migrations/052_ProjectionTurnsEffectiveModel.ts"; import Migration0053 from "./Migrations/053_LocalSubagentTerminalDeliveries.ts"; +import Migration0054 from "./Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts"; /** * Migration loader with all migrations defined inline. @@ -129,6 +130,7 @@ export const migrationEntries = [ [51, "AuthAudienceCeilings", Migration0051], [52, "ProjectionTurnsEffectiveModel", Migration0052], [53, "LocalSubagentTerminalDeliveries", Migration0053], + [54, "BackfillPreexistingLocalSubagentTerminalDeliveries", Migration0054], ] as const; export const makeMigrationLoader = (throughId?: number) => diff --git a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts new file mode 100644 index 00000000000..b6392978571 --- /dev/null +++ b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts @@ -0,0 +1,372 @@ +import { assert, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +import { runMigrations } from "../Migrations.ts"; +import * as NodeSqliteClient from "../NodeSqliteClient.ts"; +import { backfillPreexistingLocalSubagentTerminalDeliveries } from "./054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts"; + +const layer = it.layer(Layer.mergeAll(NodeSqliteClient.layerMemory())); + +const parentThreadId = "terminal-backfill-parent"; +const timestamp = "2000-01-01T08:00:00.000Z"; + +layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { + it.effect("backfills every delivered pre-existing terminal local child and is idempotent", () => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* runMigrations({ toMigrationInclusive: 53 }); + + const seedChild = Effect.fn("seedTerminalBackfillChild")(function* (input: { + readonly threadId: string; + readonly turnState?: "completed" | "error" | "interrupted" | "pending" | "running"; + readonly turnAt?: string; + readonly archivedAt?: string; + readonly deletedAt?: string; + readonly latestUserMessageAt?: string; + readonly sessionStatus?: "ready" | "running" | "waiting" | "error" | "stopped"; + readonly sessionUpdatedAt?: string; + readonly activeTurn?: boolean; + readonly parentEnvironmentId?: string; + readonly sequence?: number; + }) { + const turnId = input.turnState === undefined ? null : `${input.threadId}-turn`; + yield* sql` + INSERT INTO projection_threads ( + thread_id, + project_id, + title, + latest_turn_id, + model_selection_json, + runtime_mode, + interaction_mode, + created_at, + updated_at, + deleted_at, + archived_at, + latest_user_message_at, + parent_thread_id, + parent_environment_id + ) + VALUES ( + ${input.threadId}, + ${"terminal-backfill-project"}, + ${"terminal child"}, + ${turnId}, + ${'{"instanceId":"codex","model":"gpt-5-codex"}'}, + ${"full-access"}, + ${"default"}, + ${timestamp}, + ${timestamp}, + ${input.deletedAt ?? null}, + ${input.archivedAt ?? null}, + ${input.latestUserMessageAt ?? null}, + ${parentThreadId}, + ${input.parentEnvironmentId ?? null} + ) + `; + if (turnId !== null) { + const turnAt = input.turnAt ?? timestamp; + yield* sql` + INSERT INTO projection_turns ( + thread_id, + turn_id, + state, + requested_at, + started_at, + completed_at, + checkpoint_files_json + ) + VALUES ( + ${input.threadId}, + ${turnId}, + ${input.turnState}, + ${turnAt}, + ${turnAt}, + ${input.turnState === "pending" || input.turnState === "running" ? null : turnAt}, + ${"[]"} + ) + `; + } + if (input.sessionStatus !== undefined) { + yield* sql` + INSERT INTO projection_thread_sessions ( + thread_id, + status, + active_turn_id, + updated_at + ) + VALUES ( + ${input.threadId}, + ${input.sessionStatus}, + ${input.activeTurn === true ? turnId : null}, + ${input.sessionUpdatedAt ?? timestamp} + ) + `; + } + if (input.sequence !== undefined) { + yield* sql` + INSERT INTO orchestration_events ( + sequence, + event_id, + aggregate_kind, + stream_id, + stream_version, + event_type, + occurred_at, + actor_kind, + payload_json, + metadata_json + ) + VALUES ( + ${input.sequence}, + ${`terminal-backfill-event-${input.threadId}`}, + ${"thread"}, + ${input.threadId}, + ${0}, + ${"thread.turn-diff-completed"}, + ${timestamp}, + ${"server"}, + ${"{}"}, + ${"{}"} + ) + `; + } + }); + + yield* Effect.forEach( + [ + { + threadId: "completed-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + sequence: 101, + }, + { + threadId: "failed-local", + turnState: "error" as const, + sessionStatus: "error" as const, + sequence: 102, + }, + { threadId: "stopped-local", sessionStatus: "stopped" as const }, + { threadId: "archived-local", archivedAt: "2000-01-01T09:00:00.000Z" }, + { + threadId: "archived-completed-local", + turnState: "completed" as const, + turnAt: "2000-01-01T08:00:00.000Z", + archivedAt: "2000-01-01T09:00:00.000Z", + }, + { + threadId: "archived-before-terminal-local", + turnState: "completed" as const, + turnAt: "2000-01-01T10:00:00.000Z", + archivedAt: "2000-01-01T09:00:00.000Z", + }, + { + threadId: "deleted-active-local", + turnState: "running" as const, + deletedAt: "2000-01-01T09:00:00.000Z", + sessionStatus: "running" as const, + activeTurn: true, + }, + { + threadId: "already-claimed-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + }, + { + threadId: "live-local", + turnState: "running" as const, + sessionStatus: "running" as const, + activeTurn: true, + }, + { + threadId: "stale-completed-local", + turnState: "completed" as const, + turnAt: "2000-01-01T08:00:00.000Z", + latestUserMessageAt: "2000-01-01T09:00:00.000Z", + sessionStatus: "ready" as const, + }, + { + threadId: "undelivered-pre-053-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + }, + { + threadId: "wait-delivered-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + }, + { + threadId: "post-053-completed-local", + turnState: "completed" as const, + turnAt: "2999-01-01T08:00:00.000Z", + sessionStatus: "ready" as const, + }, + { + threadId: "post-053-archived-local", + archivedAt: "2999-01-01T08:00:00.000Z", + }, + { + threadId: "pre-053-archive-post-053-stop-local", + archivedAt: "2000-01-01T09:00:00.000Z", + sessionStatus: "stopped" as const, + sessionUpdatedAt: "2999-01-01T08:00:00.000Z", + }, + { + threadId: "remote-terminal", + turnState: "completed" as const, + sessionStatus: "ready" as const, + parentEnvironmentId: "remote-environment", + }, + ], + seedChild, + ); + + const messageDeliveredChildren = [ + "archived-before-terminal-local", + "archived-completed-local", + "archived-local", + "completed-local", + "deleted-active-local", + "failed-local", + "post-053-archived-local", + "post-053-completed-local", + "stopped-local", + ] as const; + for (const [index, childThreadId] of messageDeliveredChildren.entries()) { + const deliveredAt = childThreadId.startsWith("post-053-") + ? "3000-01-01T08:00:00.000Z" + : "2001-01-01T08:00:00.000Z"; + yield* sql` + INSERT INTO orchestration_events ( + sequence, + event_id, + aggregate_kind, + stream_id, + stream_version, + event_type, + occurred_at, + actor_kind, + payload_json, + metadata_json + ) + VALUES ( + ${1_001 + index}, + ${`terminal-backfill-delivered-${childThreadId}`}, + ${"thread"}, + ${parentThreadId}, + ${1_001 + index}, + ${"thread.message-sent"}, + ${deliveredAt}, + ${"server"}, + ${`{"role":"system","text":"[sub-agent ${childThreadId} completed]"}`}, + ${"{}"} + ) + `; + } + yield* sql` + INSERT INTO subagent_wait_deliveries ( + child_thread_id, + parent_thread_id, + delivered_at, + parent_turn_id_at_delivery + ) + VALUES ( + ${"wait-delivered-local"}, + ${parentThreadId}, + ${"2001-01-01T08:00:00.000Z"}, + ${null} + ) + `; + + yield* sql` + INSERT INTO subagent_terminal_deliveries ( + child_thread_id, + parent_thread_id, + terminal_delivery_claim_id, + terminal_delivery_claimed_at, + terminal_delivery_claimed_sequence, + terminal_kind + ) + VALUES ( + ${"already-claimed-local"}, + ${parentThreadId}, + ${"existing-claim"}, + ${timestamp}, + ${99}, + ${"completed"} + ) + `; + + yield* runMigrations({ toMigrationInclusive: 54 }); + yield* backfillPreexistingLocalSubagentTerminalDeliveries(); + + const rows = yield* sql<{ + readonly childThreadId: string; + readonly claimId: string; + readonly claimedSequence: number; + readonly terminalKind: string; + }>` + SELECT + child_thread_id AS "childThreadId", + terminal_delivery_claim_id AS "claimId", + terminal_delivery_claimed_sequence AS "claimedSequence", + terminal_kind AS "terminalKind" + FROM subagent_terminal_deliveries + ORDER BY child_thread_id + `; + + assert.deepEqual( + rows.map(({ childThreadId }) => childThreadId), + [ + "already-claimed-local", + "archived-before-terminal-local", + "archived-completed-local", + "archived-local", + "completed-local", + "deleted-active-local", + "failed-local", + "stopped-local", + "wait-delivered-local", + ], + ); + assert.deepEqual( + new Map(rows.map(({ childThreadId, terminalKind }) => [childThreadId, terminalKind])), + new Map([ + ["already-claimed-local", "completed"], + ["archived-before-terminal-local", "archived"], + ["archived-completed-local", "completed"], + ["archived-local", "archived"], + ["completed-local", "completed"], + ["deleted-active-local", "killed"], + ["failed-local", "failed"], + ["stopped-local", "failed"], + ["wait-delivered-local", "completed"], + ]), + ); + assert.equal( + rows.find(({ childThreadId }) => childThreadId === "completed-local")?.claimedSequence, + 101, + ); + assert.equal( + rows.find(({ childThreadId }) => childThreadId === "failed-local")?.claimedSequence, + 102, + ); + assert.equal( + rows.find(({ childThreadId }) => childThreadId === "archived-local")?.claimedSequence, + 0, + ); + assert.equal( + rows.find(({ childThreadId }) => childThreadId === "already-claimed-local")?.claimId, + "existing-claim", + ); + assert.isTrue( + rows + .filter(({ childThreadId }) => childThreadId !== "already-claimed-local") + .every(({ childThreadId, claimId }) => claimId === `migration:054:${childThreadId}`), + ); + }), + ); +}); diff --git a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts new file mode 100644 index 00000000000..ad322a7b891 --- /dev/null +++ b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts @@ -0,0 +1,165 @@ +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +/** + * Mark local child lifecycles that were already terminal when durable delivery + * claims shipped, and whose parent log proves the wake was delivered. Migration + * 053's durable creation time is the cutoff, so a child that terminates in the + * later upgrade window remains unclaimed and is recovered by the normal replay + * path. Per-child delivery evidence keeps skip-level and crash-interrupted + * upgrades from turning undelivered results into tombstones. This is + * deliberately one-time: doing the same reconciliation on every boot could + * suppress a genuinely undelivered wake after a later crash between terminal + * persistence and parent delivery. + * + * Exported so the migration and restart regression exercise the same query. + */ +export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( + "backfillPreexistingLocalSubagentTerminalDeliveries", +)(function* () { + const sql = yield* SqlClient.SqlClient; + + yield* sql` + WITH durable_claims_cutoff AS ( + SELECT created_at AS cutoff_at + FROM effect_sql_migrations + WHERE migration_id = 53 + ), + local_children AS ( + SELECT + threads.thread_id, + threads.parent_thread_id, + threads.deleted_at, + threads.archived_at, + threads.latest_user_message_at, + turns.state AS turn_state, + COALESCE(turns.completed_at, turns.started_at, turns.requested_at) AS terminal_at, + sessions.status AS session_status, + sessions.active_turn_id, + sessions.updated_at AS session_updated_at, + durable_claims_cutoff.cutoff_at, + CASE + WHEN turns.state IN ('completed', 'error', 'interrupted') + AND ( + threads.latest_user_message_at IS NULL + OR COALESCE(turns.completed_at, turns.started_at, turns.requested_at) IS NULL + OR threads.latest_user_message_at <= + COALESCE(turns.completed_at, turns.started_at, turns.requested_at) + ) + THEN 1 + ELSE 0 + END AS has_fresh_terminal_turn + FROM projection_threads AS threads + LEFT JOIN projection_turns AS turns + ON turns.thread_id = threads.thread_id + AND turns.turn_id = threads.latest_turn_id + LEFT JOIN projection_thread_sessions AS sessions + ON sessions.thread_id = threads.thread_id + CROSS JOIN durable_claims_cutoff + WHERE threads.parent_thread_id IS NOT NULL + AND threads.parent_environment_id IS NULL + ), + terminal_children AS ( + SELECT + thread_id, + parent_thread_id, + cutoff_at, + CASE + WHEN deleted_at IS NOT NULL THEN 'killed' + WHEN archived_at IS NOT NULL + AND ( + has_fresh_terminal_turn = 0 + OR (terminal_at IS NOT NULL AND terminal_at > archived_at) + ) + THEN 'archived' + WHEN has_fresh_terminal_turn = 1 AND turn_state = 'completed' THEN 'completed' + WHEN has_fresh_terminal_turn = 1 THEN 'failed' + WHEN archived_at IS NOT NULL THEN 'archived' + ELSE 'failed' + END AS terminal_kind, + CASE + WHEN deleted_at IS NOT NULL THEN deleted_at + WHEN has_fresh_terminal_turn = 1 THEN terminal_at + WHEN archived_at IS NOT NULL AND session_status IN ('error', 'stopped') + THEN CASE + WHEN session_updated_at IS NULL OR archived_at >= session_updated_at THEN archived_at + ELSE session_updated_at + END + WHEN archived_at IS NOT NULL THEN archived_at + ELSE session_updated_at + END AS terminal_evidence_at + FROM local_children + WHERE deleted_at IS NOT NULL + OR ( + NOT ( + session_status IN ('running', 'waiting') + AND active_turn_id IS NOT NULL + ) + AND ( + archived_at IS NOT NULL + OR has_fresh_terminal_turn = 1 + OR session_status IN ('error', 'stopped') + ) + ) + ), + delivered_parent_wakes AS MATERIALIZED ( + SELECT + parent_thread_id, + occurred_at, + substr(wake_text, 12, instr(substr(wake_text, 12), ' ') - 1) AS child_thread_id + FROM ( + SELECT + stream_id AS parent_thread_id, + occurred_at, + json_extract(payload_json, '$.text') AS wake_text + FROM orchestration_events + WHERE event_type = 'thread.message-sent' + AND json_extract(payload_json, '$.role') = 'system' + ) + WHERE wake_text LIKE '[sub-agent % %' + ) + INSERT INTO subagent_terminal_deliveries ( + child_thread_id, + parent_thread_id, + terminal_delivery_claim_id, + terminal_delivery_claimed_at, + terminal_delivery_claimed_sequence, + terminal_kind + ) + SELECT + terminal_children.thread_id, + terminal_children.parent_thread_id, + 'migration:054:' || terminal_children.thread_id, + strftime('%Y-%m-%dT%H:%M:%fZ', 'now'), + COALESCE(( + SELECT MAX(events.sequence) + FROM orchestration_events AS events + WHERE events.stream_id = terminal_children.thread_id + AND julianday(events.occurred_at) <= julianday(terminal_children.cutoff_at) + ), 0), + terminal_children.terminal_kind + FROM terminal_children + WHERE julianday(terminal_children.terminal_evidence_at) <= + julianday(terminal_children.cutoff_at) + AND ( + EXISTS ( + SELECT 1 + FROM delivered_parent_wakes AS delivered_event + WHERE delivered_event.parent_thread_id = terminal_children.parent_thread_id + AND delivered_event.child_thread_id = terminal_children.thread_id + AND julianday(delivered_event.occurred_at) >= + julianday(terminal_children.terminal_evidence_at) + ) + OR EXISTS ( + SELECT 1 + FROM subagent_wait_deliveries AS wait_delivery + WHERE wait_delivery.child_thread_id = terminal_children.thread_id + AND julianday(wait_delivery.delivered_at) >= + julianday(terminal_children.terminal_evidence_at) + ) + ) + ON CONFLICT (child_thread_id) DO NOTHING + `; +}); + +export default backfillPreexistingLocalSubagentTerminalDeliveries(); From 2fb8bd76fbeaa24578649398c1731ee472368c66 Mon Sep 17 00:00:00 2001 From: "wizzoapp[bot]" <254688279+wizzoapp[bot]@users.noreply.github.com> Date: Wed, 22 Jul 2026 11:51:51 +0100 Subject: [PATCH 2/4] codex: address PR review feedback (#235) --- .../Layers/ChildThreadCoordinator.test.ts | 2 +- ...ingLocalSubagentTerminalDeliveries.test.ts | 57 ++++++++++++++++++- ...existingLocalSubagentTerminalDeliveries.ts | 34 ++++++++++- 3 files changed, 87 insertions(+), 6 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ChildThreadCoordinator.test.ts b/apps/server/src/orchestration/Layers/ChildThreadCoordinator.test.ts index 35744f8866a..464573208f4 100644 --- a/apps/server/src/orchestration/Layers/ChildThreadCoordinator.test.ts +++ b/apps/server/src/orchestration/Layers/ChildThreadCoordinator.test.ts @@ -1056,7 +1056,7 @@ describe("ChildThreadCoordinator", () => { ${"thread.message-sent"}, ${now}, ${"server"}, - ${`{"role":"system","text":"[sub-agent ${row.threadId} completed]"}`}, + ${`{"role":"system","text":"[sub-agent ${row.threadId} completed] delivered"}`}, ${"{}"} ) `; diff --git a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts index b6392978571..09889fcd62f 100644 --- a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts +++ b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts @@ -198,6 +198,26 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { turnState: "completed" as const, sessionStatus: "ready" as const, }, + { + threadId: "consolidated-second-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + }, + { + threadId: "wrong-parent-wait-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + }, + { + threadId: "wild_card-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + }, + { + threadId: "wild%card-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + }, { threadId: "post-053-completed-local", turnState: "completed" as const, @@ -238,7 +258,9 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { for (const [index, childThreadId] of messageDeliveredChildren.entries()) { const deliveredAt = childThreadId.startsWith("post-053-") ? "3000-01-01T08:00:00.000Z" - : "2001-01-01T08:00:00.000Z"; + : childThreadId === "archived-before-terminal-local" + ? "2000-01-01T09:30:00.000Z" + : "2001-01-01T08:00:00.000Z"; yield* sql` INSERT INTO orchestration_events ( sequence, @@ -261,11 +283,37 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ${"thread.message-sent"}, ${deliveredAt}, ${"server"}, - ${`{"role":"system","text":"[sub-agent ${childThreadId} completed]"}`}, + ${`{"role":"system","text":"[sub-agent ${childThreadId} completed] delivered"}`}, ${"{}"} ) `; } + yield* sql` + INSERT INTO orchestration_events ( + sequence, + event_id, + aggregate_kind, + stream_id, + stream_version, + event_type, + occurred_at, + actor_kind, + payload_json, + metadata_json + ) + VALUES ( + ${1_100}, + ${"terminal-backfill-delivered-consolidated"}, + ${"thread"}, + ${parentThreadId}, + ${1_100}, + ${"thread.message-sent"}, + ${"2001-01-01T08:00:00.000Z"}, + ${"server"}, + ${'{"role":"system","text":"[sub-agent already-delivered-first completed] delivered\\n[sub-agent consolidated-second-local completed] delivered\\n[sub-agent wildXcard-local completed] delivered\\n[sub-agent wildXYZcard-local completed] delivered"}'}, + ${"{}"} + ) + `; yield* sql` INSERT INTO subagent_wait_deliveries ( child_thread_id, @@ -278,6 +326,11 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ${parentThreadId}, ${"2001-01-01T08:00:00.000Z"}, ${null} + ), ( + ${"wrong-parent-wait-local"}, + ${"previous-parent"}, + ${"2001-01-01T08:00:00.000Z"}, + ${null} ) `; diff --git a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts index ad322a7b891..12179d318f3 100644 --- a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts +++ b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts @@ -79,6 +79,12 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( END AS terminal_kind, CASE WHEN deleted_at IS NOT NULL THEN deleted_at + WHEN archived_at IS NOT NULL + AND ( + has_fresh_terminal_turn = 0 + OR (terminal_at IS NOT NULL AND terminal_at > archived_at) + ) + THEN archived_at WHEN has_fresh_terminal_turn = 1 THEN terminal_at WHEN archived_at IS NOT NULL AND session_status IN ('error', 'stopped') THEN CASE @@ -106,7 +112,7 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( SELECT parent_thread_id, occurred_at, - substr(wake_text, 12, instr(substr(wake_text, 12), ' ') - 1) AS child_thread_id + wake_text FROM ( SELECT stream_id AS parent_thread_id, @@ -116,7 +122,7 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( WHERE event_type = 'thread.message-sent' AND json_extract(payload_json, '$.role') = 'system' ) - WHERE wake_text LIKE '[sub-agent % %' + WHERE wake_text LIKE '%[sub-agent % %' ) INSERT INTO subagent_terminal_deliveries ( child_thread_id, @@ -146,7 +152,28 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( SELECT 1 FROM delivered_parent_wakes AS delivered_event WHERE delivered_event.parent_thread_id = terminal_children.parent_thread_id - AND delivered_event.child_thread_id = terminal_children.thread_id + AND substr( + delivered_event.wake_text, + 1, + length('[sub-agent ' || terminal_children.thread_id || ' ') + ) = '[sub-agent ' || terminal_children.thread_id || ' ' + AND ( + substr( + delivered_event.wake_text, + length('[sub-agent ' || terminal_children.thread_id || ' ') + 1, + length('completed] ') + ) = 'completed] ' + OR substr( + delivered_event.wake_text, + length('[sub-agent ' || terminal_children.thread_id || ' ') + 1, + length('failed] ') + ) = 'failed] ' + OR substr( + delivered_event.wake_text, + length('[sub-agent ' || terminal_children.thread_id || ' ') + 1, + length('killed] ') + ) = 'killed] ' + ) AND julianday(delivered_event.occurred_at) >= julianday(terminal_children.terminal_evidence_at) ) @@ -154,6 +181,7 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( SELECT 1 FROM subagent_wait_deliveries AS wait_delivery WHERE wait_delivery.child_thread_id = terminal_children.thread_id + AND wait_delivery.parent_thread_id = terminal_children.parent_thread_id AND julianday(wait_delivery.delivered_at) >= julianday(terminal_children.terminal_evidence_at) ) From 804d9fa487459c2ae25c9e0d728bb1abaa8c6f19 Mon Sep 17 00:00:00 2001 From: "wizzoapp[bot]" <254688279+wizzoapp[bot]@users.noreply.github.com> Date: Wed, 22 Jul 2026 12:07:21 +0100 Subject: [PATCH 3/4] codex: harden lifecycle backfill evidence (#235) --- ...ingLocalSubagentTerminalDeliveries.test.ts | 59 ++++++++++++++++++- ...existingLocalSubagentTerminalDeliveries.ts | 46 +++++++++------ 2 files changed, 84 insertions(+), 21 deletions(-) diff --git a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts index 09889fcd62f..727807ec801 100644 --- a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts +++ b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts @@ -149,6 +149,13 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { sessionStatus: "error" as const, sequence: 102, }, + { + threadId: "session-failed-before-turn-projection-local", + turnState: "error" as const, + turnAt: "2000-01-01T10:00:00.000Z", + sessionStatus: "error" as const, + sessionUpdatedAt: "2000-01-01T09:00:00.000Z", + }, { threadId: "stopped-local", sessionStatus: "stopped" as const }, { threadId: "archived-local", archivedAt: "2000-01-01T09:00:00.000Z" }, { @@ -175,6 +182,17 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { turnState: "completed" as const, sessionStatus: "ready" as const, }, + { + threadId: "stale-claim-newer-lifecycle-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + sequence: 103, + }, + { + threadId: "status-mismatch-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + }, { threadId: "live-local", turnState: "running" as const, @@ -253,6 +271,9 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { "failed-local", "post-053-archived-local", "post-053-completed-local", + "session-failed-before-turn-projection-local", + "stale-claim-newer-lifecycle-local", + "status-mismatch-local", "stopped-local", ] as const; for (const [index, childThreadId] of messageDeliveredChildren.entries()) { @@ -260,7 +281,21 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ? "3000-01-01T08:00:00.000Z" : childThreadId === "archived-before-terminal-local" ? "2000-01-01T09:30:00.000Z" - : "2001-01-01T08:00:00.000Z"; + : childThreadId === "session-failed-before-turn-projection-local" + ? "2000-01-01T09:30:00.000Z" + : "2001-01-01T08:00:00.000Z"; + const deliveredStatus = + childThreadId === "failed-local" || + childThreadId === "session-failed-before-turn-projection-local" || + childThreadId === "stopped-local" || + childThreadId === "status-mismatch-local" + ? "failed" + : childThreadId === "archived-before-terminal-local" || + childThreadId === "archived-local" || + childThreadId === "deleted-active-local" || + childThreadId === "post-053-archived-local" + ? "killed" + : "completed"; yield* sql` INSERT INTO orchestration_events ( sequence, @@ -283,7 +318,7 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ${"thread.message-sent"}, ${deliveredAt}, ${"server"}, - ${`{"role":"system","text":"[sub-agent ${childThreadId} completed] delivered"}`}, + ${`{"role":"system","text":"[sub-agent ${childThreadId} ${deliveredStatus}] delivered"}`}, ${"{}"} ) `; @@ -350,6 +385,13 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ${timestamp}, ${99}, ${"completed"} + ), ( + ${"stale-claim-newer-lifecycle-local"}, + ${parentThreadId}, + ${"migration:053:stale-claim"}, + ${timestamp}, + ${1}, + ${"failed"} ) `; @@ -381,6 +423,8 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { "completed-local", "deleted-active-local", "failed-local", + "session-failed-before-turn-projection-local", + "stale-claim-newer-lifecycle-local", "stopped-local", "wait-delivered-local", ], @@ -395,6 +439,8 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ["completed-local", "completed"], ["deleted-active-local", "killed"], ["failed-local", "failed"], + ["session-failed-before-turn-projection-local", "failed"], + ["stale-claim-newer-lifecycle-local", "completed"], ["stopped-local", "failed"], ["wait-delivered-local", "completed"], ]), @@ -415,6 +461,15 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { rows.find(({ childThreadId }) => childThreadId === "already-claimed-local")?.claimId, "existing-claim", ); + assert.deepEqual( + rows.find(({ childThreadId }) => childThreadId === "stale-claim-newer-lifecycle-local"), + { + childThreadId: "stale-claim-newer-lifecycle-local", + claimId: "migration:054:stale-claim-newer-lifecycle-local", + claimedSequence: 103, + terminalKind: "completed", + }, + ); assert.isTrue( rows .filter(({ childThreadId }) => childThreadId !== "already-claimed-local") diff --git a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts index 12179d318f3..ee357b89176 100644 --- a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts +++ b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts @@ -85,6 +85,10 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( OR (terminal_at IS NOT NULL AND terminal_at > archived_at) ) THEN archived_at + WHEN has_fresh_terminal_turn = 1 + AND turn_state <> 'completed' + AND session_status IN ('error', 'stopped') + THEN COALESCE(session_updated_at, terminal_at) WHEN has_fresh_terminal_turn = 1 THEN terminal_at WHEN archived_at IS NOT NULL AND session_status IN ('error', 'stopped') THEN CASE @@ -155,24 +159,21 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( AND substr( delivered_event.wake_text, 1, - length('[sub-agent ' || terminal_children.thread_id || ' ') - ) = '[sub-agent ' || terminal_children.thread_id || ' ' - AND ( - substr( - delivered_event.wake_text, - length('[sub-agent ' || terminal_children.thread_id || ' ') + 1, - length('completed] ') - ) = 'completed] ' - OR substr( - delivered_event.wake_text, - length('[sub-agent ' || terminal_children.thread_id || ' ') + 1, - length('failed] ') - ) = 'failed] ' - OR substr( - delivered_event.wake_text, - length('[sub-agent ' || terminal_children.thread_id || ' ') + 1, - length('killed] ') - ) = 'killed] ' + length( + '[sub-agent ' || terminal_children.thread_id || ' ' || + CASE terminal_children.terminal_kind + WHEN 'completed' THEN 'completed' + WHEN 'failed' THEN 'failed' + ELSE 'killed' + END || '] ' + ) + ) = ( + '[sub-agent ' || terminal_children.thread_id || ' ' || + CASE terminal_children.terminal_kind + WHEN 'completed' THEN 'completed' + WHEN 'failed' THEN 'failed' + ELSE 'killed' + END || '] ' ) AND julianday(delivered_event.occurred_at) >= julianday(terminal_children.terminal_evidence_at) @@ -186,7 +187,14 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( julianday(terminal_children.terminal_evidence_at) ) ) - ON CONFLICT (child_thread_id) DO NOTHING + ON CONFLICT (child_thread_id) DO UPDATE SET + parent_thread_id = excluded.parent_thread_id, + terminal_delivery_claim_id = excluded.terminal_delivery_claim_id, + terminal_delivery_claimed_at = excluded.terminal_delivery_claimed_at, + terminal_delivery_claimed_sequence = excluded.terminal_delivery_claimed_sequence, + terminal_kind = excluded.terminal_kind + WHERE excluded.terminal_delivery_claimed_sequence > + subagent_terminal_deliveries.terminal_delivery_claimed_sequence `; }); From f18a792918b0d3d0c29bfff366481228d48e1601 Mon Sep 17 00:00:00 2001 From: "wizzoapp[bot]" <254688279+wizzoapp[bot]@users.noreply.github.com> Date: Wed, 22 Jul 2026 12:40:28 +0100 Subject: [PATCH 4/4] codex: bind backfill to terminal event order (#235) --- ...ingLocalSubagentTerminalDeliveries.test.ts | 143 +++++++++++++++++- ...existingLocalSubagentTerminalDeliveries.ts | 131 +++++++++++++--- 2 files changed, 244 insertions(+), 30 deletions(-) diff --git a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts index 727807ec801..bd4033fbbeb 100644 --- a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts +++ b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts @@ -1,6 +1,7 @@ import { assert, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Schema from "effect/Schema"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { runMigrations } from "../Migrations.ts"; @@ -11,16 +12,19 @@ const layer = it.layer(Layer.mergeAll(NodeSqliteClient.layerMemory())); const parentThreadId = "terminal-backfill-parent"; const timestamp = "2000-01-01T08:00:00.000Z"; +const encodeUnknownJson = Schema.encodeUnknownSync(Schema.UnknownFromJsonString); layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { it.effect("backfills every delivered pre-existing terminal local child and is idempotent", () => Effect.gen(function* () { const sql = yield* SqlClient.SqlClient; yield* runMigrations({ toMigrationInclusive: 53 }); + let nextTerminalSequence = 200; const seedChild = Effect.fn("seedTerminalBackfillChild")(function* (input: { readonly threadId: string; readonly turnState?: "completed" | "error" | "interrupted" | "pending" | "running"; + readonly terminalDiffStatus?: "ready" | "error" | "missing"; readonly turnAt?: string; readonly archivedAt?: string; readonly deletedAt?: string; @@ -105,7 +109,76 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ) `; } - if (input.sequence !== undefined) { + const turnAt = input.turnAt ?? timestamp; + const sessionUpdatedAt = input.sessionUpdatedAt ?? timestamp; + const archiveOwnsTerminal = + input.archivedAt !== undefined && + (input.turnState === undefined || + input.turnState === "pending" || + input.turnState === "running" || + turnAt > input.archivedAt); + const terminalEvent = + input.deletedAt !== undefined + ? { + type: "thread.deleted", + occurredAt: input.deletedAt, + payload: { deletedAt: input.deletedAt }, + } + : archiveOwnsTerminal + ? { + type: "thread.archived", + occurredAt: input.archivedAt, + payload: { archivedAt: input.archivedAt }, + } + : input.turnState === "error" && + (input.sessionStatus === "error" || input.sessionStatus === "stopped") + ? { + type: "thread.session-set", + occurredAt: sessionUpdatedAt, + payload: { + session: { status: input.sessionStatus, updatedAt: sessionUpdatedAt }, + }, + } + : input.turnState === "completed" || + input.turnState === "error" || + (input.turnState === "interrupted" && input.terminalDiffStatus !== undefined) + ? { + type: "thread.turn-diff-completed", + occurredAt: turnAt, + payload: { + turnId, + completedAt: turnAt, + status: + input.terminalDiffStatus ?? + (input.turnState === "error" ? "error" : "ready"), + }, + } + : input.turnState === "interrupted" + ? { + type: "thread.turn-interrupt-requested", + occurredAt: turnAt, + payload: { turnId, createdAt: turnAt }, + } + : input.archivedAt !== undefined + ? { + type: "thread.archived", + occurredAt: input.archivedAt, + payload: { archivedAt: input.archivedAt }, + } + : input.sessionStatus === "error" || input.sessionStatus === "stopped" + ? { + type: "thread.session-set", + occurredAt: sessionUpdatedAt, + payload: { + session: { + status: input.sessionStatus, + updatedAt: sessionUpdatedAt, + }, + }, + } + : null; + if (terminalEvent !== null) { + const sequence = input.sequence ?? nextTerminalSequence++; yield* sql` INSERT INTO orchestration_events ( sequence, @@ -120,15 +193,15 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { metadata_json ) VALUES ( - ${input.sequence}, + ${sequence}, ${`terminal-backfill-event-${input.threadId}`}, ${"thread"}, ${input.threadId}, ${0}, - ${"thread.turn-diff-completed"}, - ${timestamp}, + ${terminalEvent.type}, + ${terminalEvent.occurredAt}, ${"server"}, - ${"{}"}, + ${encodeUnknownJson(terminalEvent.payload)}, ${"{}"} ) `; @@ -193,6 +266,18 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { turnState: "completed" as const, sessionStatus: "ready" as const, }, + { + threadId: "missing-diff-interrupted-local", + turnState: "interrupted" as const, + terminalDiffStatus: "missing" as const, + sessionStatus: "ready" as const, + }, + { + threadId: "same-millisecond-old-wake-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + sequence: 1_200, + }, { threadId: "live-local", turnState: "running" as const, @@ -216,6 +301,11 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { turnState: "completed" as const, sessionStatus: "ready" as const, }, + { + threadId: "same-millisecond-wait-local", + turnState: "completed" as const, + sessionStatus: "ready" as const, + }, { threadId: "consolidated-second-local", turnState: "completed" as const, @@ -262,6 +352,33 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { seedChild, ); + yield* sql` + INSERT INTO orchestration_events ( + sequence, + event_id, + aggregate_kind, + stream_id, + stream_version, + event_type, + occurred_at, + actor_kind, + payload_json, + metadata_json + ) + VALUES ( + ${104}, + ${"terminal-backfill-newer-start-after-stale-claim-terminal"}, + ${"thread"}, + ${"stale-claim-newer-lifecycle-local"}, + ${1}, + ${"thread.turn-start-requested"}, + ${"2000-01-01T08:30:00.000Z"}, + ${"server"}, + ${"{}"}, + ${"{}"} + ) + `; + const messageDeliveredChildren = [ "archived-before-terminal-local", "archived-completed-local", @@ -269,9 +386,11 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { "completed-local", "deleted-active-local", "failed-local", + "missing-diff-interrupted-local", "post-053-archived-local", "post-053-completed-local", "session-failed-before-turn-projection-local", + "same-millisecond-old-wake-local", "stale-claim-newer-lifecycle-local", "status-mismatch-local", "stopped-local", @@ -283,9 +402,12 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ? "2000-01-01T09:30:00.000Z" : childThreadId === "session-failed-before-turn-projection-local" ? "2000-01-01T09:30:00.000Z" - : "2001-01-01T08:00:00.000Z"; + : childThreadId === "same-millisecond-old-wake-local" + ? timestamp + : "2001-01-01T08:00:00.000Z"; const deliveredStatus = childThreadId === "failed-local" || + childThreadId === "missing-diff-interrupted-local" || childThreadId === "session-failed-before-turn-projection-local" || childThreadId === "stopped-local" || childThreadId === "status-mismatch-local" @@ -366,6 +488,11 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ${"previous-parent"}, ${"2001-01-01T08:00:00.000Z"}, ${null} + ), ( + ${"same-millisecond-wait-local"}, + ${parentThreadId}, + ${timestamp}, + ${null} ) `; @@ -423,6 +550,7 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { "completed-local", "deleted-active-local", "failed-local", + "missing-diff-interrupted-local", "session-failed-before-turn-projection-local", "stale-claim-newer-lifecycle-local", "stopped-local", @@ -439,6 +567,7 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ["completed-local", "completed"], ["deleted-active-local", "killed"], ["failed-local", "failed"], + ["missing-diff-interrupted-local", "failed"], ["session-failed-before-turn-projection-local", "failed"], ["stale-claim-newer-lifecycle-local", "completed"], ["stopped-local", "failed"], @@ -455,7 +584,7 @@ layer("054_BackfillPreexistingLocalSubagentTerminalDeliveries", (it) => { ); assert.equal( rows.find(({ childThreadId }) => childThreadId === "archived-local")?.claimedSequence, - 0, + 202, ); assert.equal( rows.find(({ childThreadId }) => childThreadId === "already-claimed-local")?.claimId, diff --git a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts index ee357b89176..9fb89b24ca2 100644 --- a/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts +++ b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts @@ -32,6 +32,7 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( threads.deleted_at, threads.archived_at, threads.latest_user_message_at, + threads.latest_turn_id, turns.state AS turn_state, COALESCE(turns.completed_at, turns.started_at, turns.requested_at) AS terminal_at, sessions.status AS session_status, @@ -64,6 +65,12 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( thread_id, parent_thread_id, cutoff_at, + deleted_at, + archived_at, + latest_turn_id, + turn_state, + session_status, + session_updated_at, CASE WHEN deleted_at IS NOT NULL THEN 'killed' WHEN archived_at IS NOT NULL @@ -112,14 +119,95 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( ) ) ), + terminal_evidence AS ( + SELECT + terminal_children.*, + ( + SELECT MAX(terminal_event.sequence) + FROM orchestration_events AS terminal_event + WHERE terminal_event.stream_id = terminal_children.thread_id + AND julianday(terminal_event.occurred_at) <= julianday(terminal_children.cutoff_at) + AND ( + ( + terminal_children.terminal_kind = 'killed' + AND terminal_event.event_type = 'thread.deleted' + AND julianday(json_extract(terminal_event.payload_json, '$.deletedAt')) = + julianday(terminal_children.terminal_evidence_at) + ) + OR ( + terminal_children.terminal_kind = 'archived' + AND terminal_event.event_type = 'thread.archived' + AND julianday(json_extract(terminal_event.payload_json, '$.archivedAt')) = + julianday(terminal_children.terminal_evidence_at) + ) + OR ( + terminal_event.event_type = 'thread.turn-diff-completed' + AND json_extract(terminal_event.payload_json, '$.turnId') = + terminal_children.latest_turn_id + AND julianday(json_extract(terminal_event.payload_json, '$.completedAt')) = + julianday(terminal_children.terminal_evidence_at) + AND ( + ( + terminal_children.terminal_kind = 'completed' + AND json_extract(terminal_event.payload_json, '$.status') <> 'error' + ) + OR ( + terminal_children.terminal_kind = 'failed' + AND json_extract(terminal_event.payload_json, '$.status') + IN ('error', 'missing') + ) + ) + ) + OR ( + terminal_children.terminal_kind = 'completed' + AND terminal_event.event_type = 'thread.message-sent' + AND json_extract(terminal_event.payload_json, '$.role') = 'assistant' + AND json_extract(terminal_event.payload_json, '$.streaming') = 0 + AND json_extract(terminal_event.payload_json, '$.turnId') = + terminal_children.latest_turn_id + AND julianday(json_extract(terminal_event.payload_json, '$.updatedAt')) = + julianday(terminal_children.terminal_evidence_at) + ) + OR ( + terminal_children.terminal_kind = 'failed' + AND terminal_event.event_type = 'thread.turn-interrupt-requested' + AND json_extract(terminal_event.payload_json, '$.turnId') = + terminal_children.latest_turn_id + AND julianday(json_extract(terminal_event.payload_json, '$.createdAt')) = + julianday(terminal_children.terminal_evidence_at) + ) + OR ( + terminal_event.event_type = 'thread.session-set' + AND julianday( + json_extract(terminal_event.payload_json, '$.session.updatedAt') + ) = julianday(terminal_children.terminal_evidence_at) + AND ( + ( + terminal_children.terminal_kind = 'completed' + AND json_extract(terminal_event.payload_json, '$.session.status') + IN ('idle', 'ready') + ) + OR ( + terminal_children.terminal_kind = 'failed' + AND json_extract(terminal_event.payload_json, '$.session.status') + IN ('error', 'interrupted', 'stopped') + ) + ) + ) + ) + ) AS terminal_sequence + FROM terminal_children + ), delivered_parent_wakes AS MATERIALIZED ( SELECT parent_thread_id, + sequence, occurred_at, wake_text FROM ( SELECT stream_id AS parent_thread_id, + sequence, occurred_at, json_extract(payload_json, '$.text') AS wake_text FROM orchestration_events @@ -137,54 +225,51 @@ export const backfillPreexistingLocalSubagentTerminalDeliveries = Effect.fn( terminal_kind ) SELECT - terminal_children.thread_id, - terminal_children.parent_thread_id, - 'migration:054:' || terminal_children.thread_id, + terminal_evidence.thread_id, + terminal_evidence.parent_thread_id, + 'migration:054:' || terminal_evidence.thread_id, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'), - COALESCE(( - SELECT MAX(events.sequence) - FROM orchestration_events AS events - WHERE events.stream_id = terminal_children.thread_id - AND julianday(events.occurred_at) <= julianday(terminal_children.cutoff_at) - ), 0), - terminal_children.terminal_kind - FROM terminal_children - WHERE julianday(terminal_children.terminal_evidence_at) <= - julianday(terminal_children.cutoff_at) + terminal_evidence.terminal_sequence, + terminal_evidence.terminal_kind + FROM terminal_evidence + WHERE terminal_evidence.terminal_sequence IS NOT NULL + AND julianday(terminal_evidence.terminal_evidence_at) <= + julianday(terminal_evidence.cutoff_at) AND ( EXISTS ( SELECT 1 FROM delivered_parent_wakes AS delivered_event - WHERE delivered_event.parent_thread_id = terminal_children.parent_thread_id + WHERE delivered_event.parent_thread_id = terminal_evidence.parent_thread_id AND substr( delivered_event.wake_text, 1, length( - '[sub-agent ' || terminal_children.thread_id || ' ' || - CASE terminal_children.terminal_kind + '[sub-agent ' || terminal_evidence.thread_id || ' ' || + CASE terminal_evidence.terminal_kind WHEN 'completed' THEN 'completed' WHEN 'failed' THEN 'failed' ELSE 'killed' END || '] ' ) ) = ( - '[sub-agent ' || terminal_children.thread_id || ' ' || - CASE terminal_children.terminal_kind + '[sub-agent ' || terminal_evidence.thread_id || ' ' || + CASE terminal_evidence.terminal_kind WHEN 'completed' THEN 'completed' WHEN 'failed' THEN 'failed' ELSE 'killed' END || '] ' ) + AND delivered_event.sequence > terminal_evidence.terminal_sequence AND julianday(delivered_event.occurred_at) >= - julianday(terminal_children.terminal_evidence_at) + julianday(terminal_evidence.terminal_evidence_at) ) OR EXISTS ( SELECT 1 FROM subagent_wait_deliveries AS wait_delivery - WHERE wait_delivery.child_thread_id = terminal_children.thread_id - AND wait_delivery.parent_thread_id = terminal_children.parent_thread_id - AND julianday(wait_delivery.delivered_at) >= - julianday(terminal_children.terminal_evidence_at) + WHERE wait_delivery.child_thread_id = terminal_evidence.thread_id + AND wait_delivery.parent_thread_id = terminal_evidence.parent_thread_id + AND julianday(wait_delivery.delivered_at) > + julianday(terminal_evidence.terminal_evidence_at) ) ) ON CONFLICT (child_thread_id) DO UPDATE SET