diff --git a/apps/server/src/orchestration/Layers/ChildThreadCoordinator.test.ts b/apps/server/src/orchestration/Layers/ChildThreadCoordinator.test.ts index e6cac1a3789..464573208f4 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] delivered"}`}, + ${"{}"} + ) + `; + }), + ), + ); + } + 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..bd4033fbbeb --- /dev/null +++ b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.test.ts @@ -0,0 +1,609 @@ +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"; +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"; +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; + 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} + ) + `; + } + 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, + event_id, + aggregate_kind, + stream_id, + stream_version, + event_type, + occurred_at, + actor_kind, + payload_json, + metadata_json + ) + VALUES ( + ${sequence}, + ${`terminal-backfill-event-${input.threadId}`}, + ${"thread"}, + ${input.threadId}, + ${0}, + ${terminalEvent.type}, + ${terminalEvent.occurredAt}, + ${"server"}, + ${encodeUnknownJson(terminalEvent.payload)}, + ${"{}"} + ) + `; + } + }); + + 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: "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" }, + { + 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: "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: "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, + 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: "same-millisecond-wait-local", + 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, + 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, + ); + + 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", + "archived-local", + "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", + ] as const; + for (const [index, childThreadId] of messageDeliveredChildren.entries()) { + const deliveredAt = childThreadId.startsWith("post-053-") + ? "3000-01-01T08:00:00.000Z" + : childThreadId === "archived-before-terminal-local" + ? "2000-01-01T09:30:00.000Z" + : childThreadId === "session-failed-before-turn-projection-local" + ? "2000-01-01T09:30: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" + ? "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, + 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} ${deliveredStatus}] 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, + parent_thread_id, + delivered_at, + parent_turn_id_at_delivery + ) + VALUES ( + ${"wait-delivered-local"}, + ${parentThreadId}, + ${"2001-01-01T08:00:00.000Z"}, + ${null} + ), ( + ${"wrong-parent-wait-local"}, + ${"previous-parent"}, + ${"2001-01-01T08:00:00.000Z"}, + ${null} + ), ( + ${"same-millisecond-wait-local"}, + ${parentThreadId}, + ${timestamp}, + ${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"} + ), ( + ${"stale-claim-newer-lifecycle-local"}, + ${parentThreadId}, + ${"migration:053:stale-claim"}, + ${timestamp}, + ${1}, + ${"failed"} + ) + `; + + 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", + "missing-diff-interrupted-local", + "session-failed-before-turn-projection-local", + "stale-claim-newer-lifecycle-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"], + ["missing-diff-interrupted-local", "failed"], + ["session-failed-before-turn-projection-local", "failed"], + ["stale-claim-newer-lifecycle-local", "completed"], + ["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, + 202, + ); + assert.equal( + 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") + .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..9fb89b24ca2 --- /dev/null +++ b/apps/server/src/persistence/Migrations/054_BackfillPreexistingLocalSubagentTerminalDeliveries.ts @@ -0,0 +1,286 @@ +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, + 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, + 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, + 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 + 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 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 + 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 + 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') + ) + ) + ), + 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 + 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_evidence.thread_id, + terminal_evidence.parent_thread_id, + 'migration:054:' || terminal_evidence.thread_id, + strftime('%Y-%m-%dT%H:%M:%fZ', 'now'), + 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_evidence.parent_thread_id + AND substr( + delivered_event.wake_text, + 1, + length( + '[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_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_evidence.terminal_evidence_at) + ) + OR EXISTS ( + SELECT 1 + FROM subagent_wait_deliveries AS wait_delivery + 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 + 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 + `; +}); + +export default backfillPreexistingLocalSubagentTerminalDeliveries();