From df0fd6c54367bdfbfbf71dadab6f1693d85ca1f9 Mon Sep 17 00:00:00 2001 From: amanthanvi Date: Thu, 3 Sep 2026 04:52:50 -0400 Subject: [PATCH 1/8] fix(server): settle background agents the provider stopped reporting Native multi-agent children (Codex collab subagents) only reach persisted state as task rows mapped from provider notifications, and nothing on the server ever settled those rows itself. When Codex lost track of a child (context compaction, a lost thread tree, a provider restart), when T3 Code restarted, or when Stop ran for a child missing from the in-memory live-turn map, no terminal row was written. The client fold and the background liveness registry kept counting the ghost, so the composer showed "1 agent working" with a Stop button stuck on "Stopping...". Persisted task rows are now the authority and the server settles them at the three moments it knows background work cannot continue: the provider session exits, the user presses Stop, and startup reconciliation finds a thread whose session this process does not own. Settlement folds the thread's full task history with the same rules as the client, writes one deterministic task.updated interrupted row per still-live agent task, and feeds the same transition to the liveness registry so both authorities agree. Stop drains queued provider events first so a running row from before the interrupt cannot land after the settlement row, and the registry remembers host-settled tasks so a late status-free start row cannot re-arm them. Implemented by Claude Opus 5 and Claude Fable 5.1 via Claude Code, with independent root-cause analysis and adversarial review by GPT-5.6 Sol via Codex. --- .../OrchestrationEngineHarness.integration.ts | 4 + .../Layers/ProviderCommandReactor.test.ts | 216 ++++++++++- .../Layers/ProviderCommandReactor.ts | 33 ++ .../Layers/ProviderRuntimeIngestion.test.ts | 211 +++++++++- .../Layers/ProviderRuntimeIngestion.ts | 13 + .../Services/ProviderRuntimeIngestion.ts | 7 +- .../ThreadBackgroundLiveness.test.ts | 73 ++++ .../orchestration/ThreadBackgroundLiveness.ts | 39 ++ .../src/orchestration/ThreadTaskSettlement.ts | 360 ++++++++++++++++++ .../Layers/ProjectionThreadActivities.ts | 39 ++ .../Services/ProjectionThreadActivities.ts | 11 + apps/server/src/server.ts | 5 +- .../serverRuntimeStartup.reconcile.test.ts | 140 +++++++ apps/server/src/serverRuntimeStartup.ts | 24 ++ docs/user/composer.md | 6 + .../src/state/subagentRuntime.test.ts | 49 +++ 16 files changed, 1226 insertions(+), 4 deletions(-) create mode 100644 apps/server/src/orchestration/ThreadTaskSettlement.ts diff --git a/apps/server/integration/OrchestrationEngineHarness.integration.ts b/apps/server/integration/OrchestrationEngineHarness.integration.ts index be3e21dcbfd4..486d60a125e2 100644 --- a/apps/server/integration/OrchestrationEngineHarness.integration.ts +++ b/apps/server/integration/OrchestrationEngineHarness.integration.ts @@ -30,6 +30,7 @@ import { OrchestrationCommandReceiptRepositoryLive } from "../src/persistence/La import { OrchestrationEventStoreLive } from "../src/persistence/Layers/OrchestrationEventStore.ts"; import { ProjectionCheckpointRepositoryLive } from "../src/persistence/Layers/ProjectionCheckpoints.ts"; import { ProjectionPendingApprovalRepositoryLive } from "../src/persistence/Layers/ProjectionPendingApprovals.ts"; +import { ProjectionThreadActivityRepositoryLive } from "../src/persistence/Layers/ProjectionThreadActivities.ts"; import * as ProviderSessionRuntime from "../src/persistence/ProviderSessionRuntime.ts"; import { makeSqlitePersistenceLive } from "../src/persistence/Layers/Sqlite.ts"; import { ProjectionCheckpointRepository } from "../src/persistence/Services/ProjectionCheckpoints.ts"; @@ -313,6 +314,7 @@ export const makeOrchestrationIntegrationHarness = ( orchestrationLayer.pipe(Layer.provide(projectionSnapshotQueryLayer)), ProjectionCheckpointRepositoryLive, ProjectionPendingApprovalRepositoryLive, + ProjectionThreadActivityRepositoryLive, checkpointStoreLayer, providerLayer, RuntimeReceiptBusTest, @@ -337,6 +339,8 @@ export const makeOrchestrationIntegrationHarness = ( generateThreadTitle: () => Effect.succeed({ title: "New thread" }), } as unknown as TextGeneration["Service"]); const providerCommandReactorLayer = ProviderCommandReactorLive.pipe( + // Stop drains runtime ingestion before settling background tasks. + Layer.provideMerge(runtimeIngestionLayer), Layer.provideMerge(runtimeServicesLayer), Layer.provideMerge(gitWorkflowLayer), Layer.provideMerge(textGenerationLayer), diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index cad80f1d3bca..fed3fb3c4eb5 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -39,6 +39,7 @@ import { TextGenerationError } from "@t3tools/contracts"; import { ProviderAdapterRequestError } from "../../provider/Errors.ts"; import { OrchestrationEventStoreLive } from "../../persistence/Layers/OrchestrationEventStore.ts"; import { OrchestrationCommandReceiptRepositoryLive } from "../../persistence/Layers/OrchestrationCommandReceipts.ts"; +import { ProjectionThreadActivityRepositoryLive } from "../../persistence/Layers/ProjectionThreadActivities.ts"; import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import { ProviderService, @@ -59,6 +60,8 @@ import { } from "./ProviderCommandReactor.ts"; import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts"; import { ProviderCommandReactor } from "../Services/ProviderCommandReactor.ts"; +import { ProviderRuntimeIngestionService } from "../Services/ProviderRuntimeIngestion.ts"; +import { selectLiveAgentTasks } from "../ThreadTaskSettlement.ts"; import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import * as Clock from "effect/Clock"; @@ -109,7 +112,10 @@ async function waitFor( describe("ProviderCommandReactor", () => { let runtime: ManagedRuntime.ManagedRuntime< - OrchestrationEngineService | ProviderCommandReactor | ProjectionSnapshotQuery, + | OrchestrationEngineService + | ProviderCommandReactor + | ProjectionSnapshotQuery + | ThreadBackgroundLiveness.ThreadBackgroundLivenessService, unknown > | null = null; let scope: Scope.Closeable | null = null; @@ -168,6 +174,8 @@ describe("ProviderCommandReactor", () => { readonly titleRegenerationBeforeStart?: "one" | "two"; readonly serverActivation?: Effect.Effect; readonly interruptTurnEffect?: () => Effect.Effect; + /** Stands in for the ingestion queue the reactor drains before settling. */ + readonly ingestionDrain?: Effect.Effect; readonly stopSessionEffect?: () => Effect.Effect; readonly startSessionEffect?: ( session: ProviderSession, @@ -420,6 +428,16 @@ describe("ProviderCommandReactor", () => { const layer = ProviderCommandReactorLive.pipe( Layer.provideMerge(reactorOrchestrationLayer), Layer.provideMerge(projectionSnapshotLayer), + // Same single instance the engine and the snapshot query share. + Layer.provideMerge(ThreadBackgroundLiveness.layer), + Layer.provideMerge(ProjectionThreadActivityRepositoryLive), + Layer.provideMerge(SqlitePersistenceMemory), + Layer.provideMerge( + Layer.succeed(ProviderRuntimeIngestionService, { + start: () => Effect.void, + drain: input?.ingestionDrain ?? Effect.void, + }), + ), Layer.provideMerge(Layer.succeed(ProviderService, service)), Layer.provideMerge(makeProviderRegistryLayer(providerSnapshots as never)), Layer.provideMerge( @@ -453,6 +471,9 @@ describe("ProviderCommandReactor", () => { const engine = await runtime.runPromise(Effect.service(OrchestrationEngineService)); const snapshotQuery = await runtime.runPromise(Effect.service(ProjectionSnapshotQuery)); const reactor = await runtime.runPromise(Effect.service(ProviderCommandReactor)); + const backgroundLiveness = await runtime.runPromise( + Effect.service(ThreadBackgroundLiveness.ThreadBackgroundLivenessService), + ); const runEffect = (effect: Effect.Effect) => runtime!.runPromise(effect); await Effect.runPromise( @@ -531,6 +552,7 @@ describe("ProviderCommandReactor", () => { return { engine, readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()), + backgroundLiveness, startSession, sendTurn, interruptTurn, @@ -2668,6 +2690,198 @@ describe("ProviderCommandReactor", () => { }); }); + it("settles persisted background agent tasks the provider no longer reports on stop", async () => { + // The reactor drains ingestion before settling, so the drain hook is also + // the deterministic signal that it reached the settlement step. + const settlementReached = Effect.runSync(Deferred.make()); + const harness = await createHarness({ + ingestionDrain: Deferred.succeed(settlementReached, undefined).pipe(Effect.asVoid), + }); + const now = "2026-01-01T00:00:00.000Z"; + const threadId = ThreadId.make("thread-1"); + const taskId = "collab-child-1"; + const appendChildRow = (activityId: string, kind: string, payload: Record) => + Effect.runPromise( + harness.engine.dispatch({ + type: "thread.activity.append", + commandId: CommandId.make(`cmd-${activityId}`), + threadId, + activity: { + id: EventId.make(activityId), + createdAt: now, + tone: "info", + kind, + summary: kind, + payload: { + taskId, + agentKind: "agent", + title: "math_one", + role: "general-purpose", + agentPath: "agents/math_one", + timelineBypass: true, + ...payload, + }, + turnId: null, + }, + createdAt: now, + }), + ); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-settle-tasks"), + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: asTurnId("turn-1"), + lastError: null, + updatedAt: now, + }, + createdAt: now, + }), + ); + + // A Codex collab child the provider has since lost: its rows say running + // and the interrupt below produces no child events at all. + await appendChildRow("evt-child-started", "task.started", {}); + await appendChildRow("evt-child-running", "task.updated", { status: "running" }); + harness.backgroundLiveness.recordTaskLiveness({ + threadId, + taskId, + taskType: undefined, + status: "running", + kind: "updated", + }); + expect(harness.backgroundLiveness.getThreadBackgroundLiveness(threadId)).toBe("working"); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: CommandId.make("cmd-turn-interrupt-settle-tasks"), + threadId, + turnId: asTurnId("turn-1"), + createdAt: now, + }), + ); + + await Effect.runPromise(Deferred.await(settlementReached)); + await harness.drain(); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect( + thread?.activities.find((activity) => activity.id === `task-settled:${threadId}:${taskId}`), + ).toMatchObject({ + kind: "task.updated", + payload: { + taskId, + status: "interrupted", + endedAt: now, + agentKind: "agent", + title: "math_one", + role: "general-purpose", + agentPath: "agents/math_one", + timelineBypass: true, + }, + }); + expect(harness.backgroundLiveness.getThreadBackgroundLiveness(threadId)).toBeNull(); + }); + + it("settles only after in-flight provider events have been ingested", async () => { + // A task.updated(running) queued in ingestion before Stop must land BEFORE + // the settlement row; otherwise it is written after it and re-arms both + // the registry and the client fold. The reactor parks in the ingestion + // drain, which this harness controls, so the ordering is deterministic + // rather than a race the test hopes to lose. + const drainEntered = Effect.runSync(Deferred.make()); + const releaseDrain = Effect.runSync(Deferred.make()); + const harness = await createHarness({ + ingestionDrain: Deferred.succeed(drainEntered, undefined).pipe( + Effect.andThen(Deferred.await(releaseDrain)), + ), + }); + const now = "2026-01-01T00:00:00.000Z"; + const threadId = ThreadId.make("thread-1"); + const taskId = "collab-child-2"; + const linkage = { agentKind: "agent", title: "math_two", timelineBypass: true } as const; + const appendChildRow = (activityId: string, payload: Record) => + Effect.runPromise( + harness.engine.dispatch({ + type: "thread.activity.append", + commandId: CommandId.make(`cmd-${activityId}`), + threadId, + activity: { + id: EventId.make(activityId), + createdAt: now, + tone: "info", + kind: "task.updated", + summary: "task.updated", + payload: { taskId, ...linkage, ...payload }, + turnId: null, + }, + createdAt: now, + }), + ); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-drain-order"), + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: asTurnId("turn-1"), + lastError: null, + updatedAt: now, + }, + createdAt: now, + }), + ); + await appendChildRow("evt-drain-child-started", { status: "running" }); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: CommandId.make("cmd-turn-interrupt-drain-order"), + threadId, + turnId: asTurnId("turn-1"), + createdAt: now, + }), + ); + + // The reactor is now parked in the drain, so nothing has been settled yet. + await Effect.runPromise(Deferred.await(drainEntered)); + expect( + (await harness.readModel()).threads + .find((entry) => entry.id === threadId) + ?.activities.some((activity) => activity.id.startsWith("task-settled:")), + ).toBe(false); + + // Ingestion flushes the event it was still holding: a fresh running row + // and the matching registry arm. + await appendChildRow("evt-drain-child-late-running", { status: "running" }); + harness.backgroundLiveness.recordTaskLiveness({ + threadId, + taskId, + taskType: undefined, + status: "running", + kind: "updated", + }); + await Effect.runPromise(Deferred.succeed(releaseDrain, undefined)); + await harness.drain(); + + const activities = + (await harness.readModel()).threads.find((entry) => entry.id === threadId)?.activities ?? []; + expect(selectLiveAgentTasks(activities)).toEqual([]); + expect(harness.backgroundLiveness.getThreadBackgroundLiveness(threadId)).toBeNull(); + }); + effectIt.effect( "stops a running session and records the failure when provider interrupt fails", () => diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 57edb60ff715..b81a3cb1515a 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -41,6 +41,8 @@ import { ProviderCommandReactor, type ProviderCommandReactorShape, } from "../Services/ProviderCommandReactor.ts"; +import { ProviderRuntimeIngestionService } from "../Services/ProviderRuntimeIngestion.ts"; +import { settleThreadTasks } from "../ThreadTaskSettlement.ts"; import { forkParked, ServerActivation } from "../../serverActivation.ts"; import { canReplaceThreadTitle, DEFAULT_THREAD_TITLE } from "../threadTitles.ts"; import { @@ -49,6 +51,9 @@ import { } from "../../serverSettings.ts"; import { VcsStatusBroadcaster } from "../../vcs/VcsStatusBroadcaster.ts"; import { GitWorkflowService } from "../../git/GitWorkflowService.ts"; +/** Bound on waiting for queued provider events before settling on Stop. */ +const INTERRUPT_INGESTION_DRAIN_TIMEOUT = Duration.seconds(5); + const isProviderAdapterRequestError = Schema.is(ProviderAdapterRequestError); const isProviderDriverKind = Schema.is(ProviderDriverKind); @@ -309,6 +314,7 @@ const make = Effect.gen(function* () { const orchestrationEngine = yield* OrchestrationEngineService; const projectionSnapshotQuery = yield* ProjectionSnapshotQuery; const providerService = yield* ProviderService; + const providerRuntimeIngestion = yield* ProviderRuntimeIngestionService; const providerRegistry = yield* ProviderRegistry; const gitWorkflow = yield* GitWorkflowService; const fileSystem = yield* FileSystem.FileSystem; @@ -1336,6 +1342,33 @@ const make = Effect.gen(function* () { yield* providerService .interruptTurn({ threadId: event.payload.threadId }) .pipe(Effect.catchCause(recoverInterruptFailure)); + + // Settlement reads persisted rows, so every provider event that was + // already queued when Stop arrived has to land first — otherwise a + // task.updated(running) from before the interrupt is written after the + // settlement row and re-arms both the registry and the client fold. + // Bounded: a hot event stream must not hold Stop hostage. + const drained = yield* providerRuntimeIngestion.drain.pipe( + Effect.timeoutOption(INTERRUPT_INGESTION_DRAIN_TIMEOUT), + ); + if (Option.isNone(drained)) { + yield* Effect.logWarning( + "provider runtime ingestion did not drain before background task settlement", + { threadId: event.payload.threadId }, + ); + } + + // Stop is a host promise, not a provider request: children the provider + // has already forgotten (compaction, a lost thread tree) never emit a + // terminal event of their own, so settle the persisted rows here. Covers + // both a successful interrupt and the stopSession fallback above; tasks + // the provider does still own emit their own terminal rows afterwards, + // which is harmless. + yield* settleThreadTasks({ + threadId: event.payload.threadId, + status: "interrupted", + createdAt: event.payload.createdAt, + }); }); const processApprovalResponseRequested = Effect.fn("processApprovalResponseRequested")(function* ( diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 26332f9f8c9c..09c9400d879a 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -36,6 +36,7 @@ import { afterEach, describe, expect, it } from "vite-plus/test"; import { OrchestrationEventStoreLive } from "../../persistence/Layers/OrchestrationEventStore.ts"; import { OrchestrationCommandReceiptRepositoryLive } from "../../persistence/Layers/OrchestrationCommandReceipts.ts"; +import { ProjectionThreadActivityRepositoryLive } from "../../persistence/Layers/ProjectionThreadActivities.ts"; import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import { ProviderService, @@ -48,6 +49,7 @@ import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQu import * as ThreadBackgroundLiveness from "../ThreadBackgroundLiveness.ts"; import * as ThreadPlanProgress from "../ThreadPlanProgress.ts"; import { ProviderRuntimeIngestionLive } from "./ProviderRuntimeIngestion.ts"; +import { selectLiveAgentTasks } from "../ThreadTaskSettlement.ts"; import { DEFAULT_THREAD_TITLE } from "../threadTitles.ts"; import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts"; import { ProviderRuntimeIngestionService } from "../Services/ProviderRuntimeIngestion.ts"; @@ -197,7 +199,10 @@ async function waitForThread( describe("ProviderRuntimeIngestion", () => { let runtime: ManagedRuntime.ManagedRuntime< - OrchestrationEngineService | ProviderRuntimeIngestionService | ProjectionSnapshotQuery, + | OrchestrationEngineService + | ProviderRuntimeIngestionService + | ProjectionSnapshotQuery + | ThreadBackgroundLiveness.ThreadBackgroundLivenessService, unknown > | null = null; let scope: Scope.Closeable | null = null; @@ -249,6 +254,9 @@ describe("ProviderRuntimeIngestion", () => { // engine, and the snapshot query (reader). Layer.provideMerge(ThreadBackgroundLiveness.layer), Layer.provideMerge(ThreadPlanProgress.layer), + // Settlement reads the thread's full task history straight from the + // activity projection. + Layer.provideMerge(ProjectionThreadActivityRepositoryLive), Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(Layer.succeed(ProviderService, provider.service)), Layer.provideMerge(makeTestServerSettingsLayer(options?.serverSettings)), @@ -259,6 +267,9 @@ describe("ProviderRuntimeIngestion", () => { const engine = await runtime.runPromise(Effect.service(OrchestrationEngineService)); const snapshotQuery = await runtime.runPromise(Effect.service(ProjectionSnapshotQuery)); const ingestion = await runtime.runPromise(Effect.service(ProviderRuntimeIngestionService)); + const backgroundLiveness = await runtime.runPromise( + Effect.service(ThreadBackgroundLiveness.ThreadBackgroundLivenessService), + ); scope = await Effect.runPromise(Scope.make("sequential")); await Effect.runPromise(ingestion.start().pipe(Scope.provide(scope))); const drain = () => Effect.runPromise(ingestion.drain); @@ -321,6 +332,7 @@ describe("ProviderRuntimeIngestion", () => { engine, dispatch, readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()), + backgroundLiveness, emit: provider.emit, setProviderSession: provider.setSession, drain, @@ -3647,4 +3659,201 @@ describe("ProviderRuntimeIngestion", () => { expect(thread.session?.status).toBe("error"); expect(thread.session?.lastError).toBe("runtime still processed"); }); + it("settles Codex children the provider stopped reporting when the session exits", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const threadId = asThreadId("thread-1"); + const taskId = "collab-child-1"; + // Codex collab children carry no taskType, so ingestion stamps them + // agentKind "agent"; the linkage below is what CodexAdapter emits. + const linkage = { + title: "math_one", + role: "general-purpose", + agentPath: "agents/math_one", + timelineBypass: true, + model: "gpt-5-codex", + effort: "high", + } as const; + + harness.emit({ + type: "task.started", + eventId: asEventId("evt-collab-child-started"), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId, + payload: { taskId, description: "math_one", ...linkage }, + }); + harness.emit({ + type: "task.updated", + eventId: asEventId("evt-collab-child-running"), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId, + payload: { taskId, status: "running", ...linkage }, + }); + + await harness.drain(); + expect(harness.backgroundLiveness.getThreadBackgroundLiveness(threadId)).toBe("working"); + + // No task.updated(idle) ever arrives: the child is gone with the session. + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-collab-session-exited"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:05.000Z", + threadId, + payload: {}, + }); + // Settling twice must not pile up rows: the id is per task, not per event. + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-collab-session-exited-duplicate"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:06.000Z", + threadId, + payload: {}, + }); + + await harness.drain(); + + const thread = (await harness.readModel()).threads.find( + (entry: ProviderRuntimeTestThread) => entry.id === threadId, + ); + const settledRows = + thread?.activities.filter((activity: ProviderRuntimeTestActivity) => + activity.id.startsWith("task-settled:"), + ) ?? []; + expect(settledRows.length).toBe(1); + expect(settledRows[0]?.id).toBe(`task-settled:${threadId}:${taskId}`); + expect(settledRows[0]).toMatchObject({ + kind: "task.updated", + payload: { + taskId, + status: "interrupted", + endedAt: "2026-01-01T00:00:05.000Z", + agentKind: "agent", + ...linkage, + }, + }); + expect(harness.backgroundLiveness.getThreadBackgroundLiveness(threadId)).toBeNull(); + }); +}); + +describe("selectLiveAgentTasks", () => { + const row = ( + id: string, + kind: string, + payload: Record, + ): ProviderRuntimeTestActivity => ({ + id: asEventId(id), + createdAt: "2026-01-01T00:00:00.000Z", + tone: "info", + kind, + summary: kind, + payload, + turnId: null, + }); + const agent = { agentKind: "agent" } as const; + + it("keeps only agent tasks whose latest state is still live", () => { + const live = selectLiveAgentTasks([ + // Live: started, then an explicit running status. + row("a1", "task.started", { taskId: "live-explicit", ...agent, title: "one" }), + row("a2", "task.updated", { taskId: "live-explicit", status: "running", ...agent }), + // Live: started and never given a status of its own. + row("b1", "task.started", { taskId: "live-implicit", ...agent }), + // Not live: a usage-only progress tick is not work. + row("c1", "task.started", { taskId: "usage-only", ...agent }), + row("c2", "task.updated", { taskId: "usage-only", status: "idle", ...agent }), + row("c3", "task.progress", { + taskId: "usage-only", + usageSnapshot: true, + typedUsage: { totalTokens: 10 }, + ...agent, + }), + // Not live: duplicate terminal rows stay terminal. + row("d1", "task.started", { taskId: "terminal", ...agent }), + row("d2", "task.updated", { taskId: "terminal", status: "interrupted", ...agent }), + row("d3", "task.updated", { taskId: "terminal", status: "interrupted", ...agent }), + row("d4", "task.completed", { taskId: "terminal", status: "completed", ...agent }), + // Not live: a resting Codex child is resumable, not working. + row("e1", "task.started", { taskId: "idle", ...agent }), + row("e2", "task.updated", { taskId: "idle", status: "idle", ...agent }), + // Not an agent: shells and monitors never join the roster. + row("f1", "task.started", { taskId: "shell", agentKind: "background" }), + row("f2", "task.updated", { taskId: "shell", status: "running", agentKind: "background" }), + ]); + + expect(live.map((task) => task.taskId).toSorted()).toEqual(["live-explicit", "live-implicit"]); + expect(live.find((task) => task.taskId === "live-explicit")?.linkage).toEqual({ + agentKind: "agent", + }); + }); + + it("a late start row does not reopen a settled task", () => { + const live = selectLiveAgentTasks([ + row("a1", "task.started", { taskId: "settled", ...agent }), + row("a2", "task.updated", { taskId: "settled", status: "failed", ...agent }), + row("a3", "task.started", { taskId: "settled", ...agent }), + ]); + + expect(live).toEqual([]); + }); + + it("a settled workflow coordinator ends its members' run", () => { + const workflow = { taskId: "wf-1", taskType: "local_workflow", ...agent } as const; + const live = selectLiveAgentTasks([ + row("w1", "task.started", { ...workflow }), + row("m1", "task.started", { taskId: "wf-member-1", parentAgentId: "wf-1", ...agent }), + row("m2", "task.updated", { + taskId: "wf-member-1", + status: "running", + parentAgentId: "wf-1", + ...agent, + }), + // A member of another (still running) coordinator keeps working. + row("o1", "task.started", { taskId: "other-1", taskType: "local_workflow", ...agent }), + row("m3", "task.started", { taskId: "other-member-1", parentAgentId: "other-1", ...agent }), + // The coordinator finishes without its member ever getting a terminal + // row: the client already shows that member as completed, so settling + // it would flip a finished run to "interrupted". + row("w2", "task.completed", { taskId: "wf-1", status: "completed", ...agent }), + ]); + + expect(live.map((task) => task.taskId).toSorted()).toEqual(["other-1", "other-member-1"]); + }); + + it("carries the newest linkage onto the live task", () => { + const live = selectLiveAgentTasks([ + row("a1", "task.started", { taskId: "child", ...agent, title: "old-name" }), + row("a2", "task.updated", { + taskId: "child", + status: "running", + ...agent, + title: "math_one", + role: "general-purpose", + agentPath: "agents/math_one", + timelineBypass: true, + model: "gpt-5-codex", + effort: "high", + // Not linkage: dropped from the synthesized row. + detail: "still working", + }), + ]); + + expect(live).toEqual([ + { + taskId: "child", + linkage: { + agentKind: "agent", + title: "math_one", + role: "general-purpose", + agentPath: "agents/math_one", + timelineBypass: true, + model: "gpt-5-codex", + effort: "high", + }, + }, + ]); + }); }); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index a90010f0b6e2..534101166708 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -42,6 +42,7 @@ import { type ProviderRuntimeIngestionShape, } from "../Services/ProviderRuntimeIngestion.ts"; import { projectActivityPayload } from "../ActivityPayloadProjection.ts"; +import { settleThreadTasks } from "../ThreadTaskSettlement.ts"; import { forkParked } from "../../serverActivation.ts"; import { ServerSettingsService } from "../../serverSettings.ts"; import { canReplaceThreadTitle } from "../threadTitles.ts"; @@ -2028,6 +2029,18 @@ const make = Effect.gen(function* () { break; } case "session.exited": + // Rows first, then the registry. Background work dies with its + // provider session, so any task still listed as running gets a + // persisted terminal row before the in-memory mirror is wiped — + // otherwise a restart rehydrates "running" rows with an empty + // registry and the agent reads as working forever. + // No drain needed here: this runs inside the ingestion worker, so + // every earlier provider event for this thread is already persisted. + yield* settleThreadTasks({ + threadId: thread.id, + status: "interrupted", + createdAt: now, + }); threadBackgroundLiveness.clearThreadLiveness(thread.id); break; default: diff --git a/apps/server/src/orchestration/Services/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Services/ProviderRuntimeIngestion.ts index b6fa2711b949..9461845f2ba3 100644 --- a/apps/server/src/orchestration/Services/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Services/ProviderRuntimeIngestion.ts @@ -27,7 +27,12 @@ export interface ProviderRuntimeIngestionShape { /** * Resolves when the internal processing queue is empty and idle. - * Intended for test use to replace timing-sensitive sleeps. + * + * Used to order work against in-flight provider events: the command reactor + * drains before settling a thread's background tasks so an event queued + * before Stop cannot land after the settlement row. Tests use it in place of + * timing-sensitive sleeps. It waits for the queue to reach zero outstanding + * items, so callers on a hot stream should bound it. */ readonly drain: Effect.Effect; } diff --git a/apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts b/apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts index b4c528480fc7..cf6b83174865 100644 --- a/apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts +++ b/apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts @@ -226,4 +226,77 @@ describe("ThreadBackgroundLiveness", () => { a.clearThreadLiveness("t"); expect(a.getThreadBackgroundLiveness("t")).toBeNull(); }); + it("only a host settlement blocks a late status-free start row", () => { + const liveness = ThreadBackgroundLiveness.make(); + const threadId = "t-settled"; + const arm = () => + liveness.recordTaskLiveness({ + threadId, + taskId: "child", + taskType: undefined, + status: undefined, + kind: "started", + }); + + arm(); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("working"); + + // The host settles the task itself (Stop, session death, restart). + liveness.recordTaskLiveness({ + threadId, + taskId: "child", + taskType: undefined, + status: "interrupted", + kind: "updated", + settledByHost: true, + }); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull(); + + // A start row the provider had already queued must not reopen it. + arm(); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull(); + + // An explicit non-terminal status is a real reactivation, and clears the + // tombstone so ordinary lifecycle resumes. + liveness.recordTaskLiveness({ + threadId, + taskId: "child", + taskType: undefined, + status: "running", + kind: "updated", + }); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("working"); + liveness.recordTaskLiveness({ + threadId, + taskId: "child", + taskType: undefined, + status: "idle", + kind: "updated", + }); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull(); + arm(); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("working"); + }); + + it("a provider's own idle does not block a later start row", () => { + const liveness = ThreadBackgroundLiveness.make(); + const threadId = "t-provider-idle"; + liveness.recordTaskLiveness({ + threadId, + taskId: "child", + taskType: undefined, + status: "idle", + kind: "updated", + }); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull(); + // A resumable Codex child that starts a new turn is real work again. + liveness.recordTaskLiveness({ + threadId, + taskId: "child", + taskType: undefined, + status: undefined, + kind: "started", + }); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("working"); + }); }); diff --git a/apps/server/src/orchestration/ThreadBackgroundLiveness.ts b/apps/server/src/orchestration/ThreadBackgroundLiveness.ts index 2781e4981f7c..6ba552adc6b4 100644 --- a/apps/server/src/orchestration/ThreadBackgroundLiveness.ts +++ b/apps/server/src/orchestration/ThreadBackgroundLiveness.ts @@ -33,6 +33,13 @@ interface ThreadLivenessState { // INERT_TASK_TYPES: plan-mode bookkeeping) so this registry, ingestion's // agentKind stamp, and the client fold can never drift apart. +/** + * Recent host settlements, enough to outlast a start row already in flight. + * One bounded FIFO for the whole registry keeps this O(1) regardless of how + * many threads or tasks a long-lived server sees. + */ +const HOST_SETTLED_TASK_MEMORY_LIMIT = 2048; + const TERMINAL_STATUSES: ReadonlySet = new Set([ "completed", "failed", @@ -59,6 +66,12 @@ export class ThreadBackgroundLivenessService extends Context.Service< readonly status: string | undefined; readonly kind: "started" | "progress" | "updated" | "completed"; readonly agentId?: string | undefined; + /** + * Set by host settlement (Stop, session death, startup reconciliation). + * Tombstones the task so a status-free start row already in flight from + * a provider that still thinks it owns the task cannot re-arm it. + */ + readonly settledByHost?: boolean | undefined; }) => void; /** Session death orphans all of a thread's background work. */ @@ -72,8 +85,11 @@ export class ThreadBackgroundLivenessService extends Context.Service< } >()("t3/orchestration/ThreadBackgroundLiveness/ThreadBackgroundLivenessService") {} +const taskKey = (threadId: string, taskId: string) => `${threadId}:${taskId}`; + export function make(): ThreadBackgroundLivenessService["Service"] { const stateByThreadId = new Map(); + const hostSettledTaskKeys = new Set(); const stateFor = (threadId: string): ThreadLivenessState => { const existing = stateByThreadId.get(threadId); @@ -127,6 +143,18 @@ export function make(): ThreadBackgroundLivenessService["Service"] { (input.status !== undefined && TERMINAL_STATUSES.has(input.status)); if (terminal) { drop(input.threadId, input.taskId); + // Only a HOST settlement leaves a tombstone. A provider's own idle or + // terminal event is ordinary lifecycle: a later start row for it is a + // real resumption and must still arm the thread. + if (input.settledByHost === true) { + hostSettledTaskKeys.add(taskKey(input.threadId, input.taskId)); + if (hostSettledTaskKeys.size > HOST_SETTLED_TASK_MEMORY_LIMIT) { + const oldest = hostSettledTaskKeys.values().next().value; + if (oldest !== undefined) { + hostSettledTaskKeys.delete(oldest); + } + } + } return; } @@ -142,7 +170,18 @@ export function make(): ThreadBackgroundLivenessService["Service"] { } } + // A status-free row for a host-settled task is a late delivery, not a + // new run. Only an explicit non-terminal status below proves the task + // really came back, and it clears the tombstone. + if ( + input.status === undefined && + hostSettledTaskKeys.has(taskKey(input.threadId, input.taskId)) + ) { + return; + } + drop(input.threadId, input.taskId); + hostSettledTaskKeys.delete(taskKey(input.threadId, input.taskId)); const state = stateFor(input.threadId); const bucket = taskType !== undefined && MONITOR_TASK_TYPES.has(taskType) ? state.monitors : state.agents; diff --git a/apps/server/src/orchestration/ThreadTaskSettlement.ts b/apps/server/src/orchestration/ThreadTaskSettlement.ts new file mode 100644 index 000000000000..74d622b38ca2 --- /dev/null +++ b/apps/server/src/orchestration/ThreadTaskSettlement.ts @@ -0,0 +1,360 @@ +/** + * Settles a thread's still-running background agent tasks from the persisted + * activity rows, without the provider's cooperation. + * + * Native multi-agent children (Codex collab, workflow members) only leave the + * live set when the provider keeps reporting them. Compaction, a provider + * restart, a host Stop for a child the provider has already forgotten, or a + * T3 restart all lose that reporting, and the last persisted row stays + * "running" forever — the Agents panel and the composer's "N agents working" + * banner never clear. + * + * Persisted rows are the authority here: at the three moments where the server + * knows background work cannot continue (session death, host Stop, startup + * reconciliation) it folds the thread's task rows, synthesizes one terminal + * `task.updated` per still-live agent task, and feeds the same transition to + * the in-memory liveness registry so rows and registry settle together. + * + * @module ThreadTaskSettlement + */ +import { + CommandId, + EventId, + type OrchestrationThreadActivity, + type ThreadId, +} from "@t3tools/contracts"; +import * as Cause from "effect/Cause"; +import * as Crypto from "effect/Crypto"; +import * as Effect from "effect/Effect"; + +import { ProjectionThreadActivityRepository } from "../persistence/Services/ProjectionThreadActivities.ts"; +import { OrchestrationEngineService } from "./Services/OrchestrationEngine.ts"; +import { ThreadBackgroundLivenessService } from "./ThreadBackgroundLiveness.ts"; + +/** + * Fleet size worth a log line. There is no cap on settlement: leaving even one + * live task behind means Stop is visibly incomplete, because the client ranks + * live agents first inside its own roster cap and would keep showing it. + */ +const LARGE_SETTLEMENT_LOG_THRESHOLD = 100; + +/** The fold reads nothing but the kind and the payload of a row. */ +export interface TaskActivityRow { + readonly kind: string; + readonly payload: unknown; +} + +/** + * Linkage the synthesized terminal row carries forward so it stays a + * self-describing agent row: `agentKind` keeps it on the Agents surface and + * the rest keep the card's identity when the start row ages out. + */ +const LINKAGE_FIELDS = [ + "agentKind", + "title", + "role", + "agentPath", + "timelineBypass", + "model", + "effort", +] as const; + +// Collapsed mirror of the client fold (subagentRuntime.foldSubagentActivities): +// only "is this task still live" matters here, so every terminal status folds +// into one state. +type FoldStatus = "pending" | "running" | "waiting" | "idle" | "terminal"; + +const LIVE_STATUSES: ReadonlySet = new Set([ + "pending", + "running", + "waiting", +]); + +// RuntimeTaskStatus, collapsed. `task.completed` carries its own outcome and +// is terminal whatever it says, so its status values are not listed here. +const KNOWN_STATUSES: ReadonlyMap = new Map([ + ["pending", "pending"], + ["running", "running"], + ["waiting", "waiting"], + ["idle", "idle"], + ["completed", "terminal"], + ["failed", "terminal"], + ["cancelled", "terminal"], + ["interrupted", "terminal"], +]); + +interface FoldEntry { + readonly taskId: string; + status: FoldStatus; + /** Latest row's payload, used to copy linkage onto the synthesized row. */ + payload: Record; + /** A workflow coordinator, whose settling ends its members' run. */ + isWorkflow: boolean; + /** The coordinator this task belongs to, if any. */ + parentAgentId: string | undefined; +} + +export interface LiveAgentTask { + readonly taskId: string; + readonly linkage: Record; +} + +function asPayload(activity: TaskActivityRow): Record | undefined { + return typeof activity.payload === "object" && activity.payload !== null + ? (activity.payload as Record) + : undefined; +} + +/** The client fold's `asString`: trimmed, or absent when empty. */ +function asTrimmed(value: unknown): string | undefined { + return typeof value === "string" && value.trim().length > 0 ? value.trim() : undefined; +} + +/** + * Identity the workflow cascade needs, filled from every row and never + * downgraded — the client fold's getOrCreate/fillMetadata rules, narrowed to + * the two fields that decide whether a coordinator owns a task. + */ +function fillIdentity(entry: FoldEntry, payload: Record): void { + if (asTrimmed(payload.taskType) === "local_workflow") { + entry.isWorkflow = true; + } + const parentAgentId = asTrimmed(payload.parentAgentId); + if (parentAgentId !== undefined) { + entry.parentAgentId = parentAgentId; + } +} + +function newEntry(taskId: string, payload: Record, status: FoldStatus): FoldEntry { + return { taskId, status, payload, isWorkflow: false, parentAgentId: undefined }; +} + +/** Rows without the server-stamped agent classification are background work. */ +function isAgentRow(payload: Record): boolean { + return payload.agentKind === "agent"; +} + +function asFoldStatus(value: unknown): FoldStatus | undefined { + return typeof value === "string" ? KNOWN_STATUSES.get(value) : undefined; +} + +/** Terminal is sticky: a duplicate or late terminal row never slides state. */ +function applyStatus(entry: FoldEntry, next: FoldStatus): void { + if (entry.status === "terminal" && next === "terminal") { + return; + } + entry.status = next; +} + +/** + * Folds a thread's persisted task activities into the set of agent tasks that + * are still live, newest linkage attached. Pure and tolerant: malformed rows + * are skipped individually. + * + * Membership is decided once, by the first row for a task id: terminal rows + * often carry only the id and a status, so re-judging them would drop agents + * mid-fold. A late `task.started` after a terminal row is an out-of-order + * delivery and does not reopen the run. A settled workflow coordinator ends + * its members' run, exactly as the client fold has it. + */ +export function selectLiveAgentTasks( + activities: ReadonlyArray, +): ReadonlyArray { + const entries = new Map(); + + for (const activity of activities) { + const payload = asPayload(activity); + if (!payload) { + continue; + } + const taskId = asTrimmed(payload.taskId); + if (!taskId) { + continue; + } + const existing = entries.get(taskId); + if (!existing && !isAgentRow(payload)) { + continue; + } + + switch (activity.kind) { + case "task.started": { + // A start row is judged on its own stamp every time: it is the row + // that puts an agent on the roster. + if (!isAgentRow(payload)) { + break; + } + const entry = existing ?? newEntry(taskId, payload, "running"); + if (existing !== undefined && existing.status === "idle") { + entry.status = "running"; + } + entry.payload = payload; + fillIdentity(entry, payload); + entries.set(taskId, entry); + break; + } + case "task.progress": { + const entry = existing ?? newEntry(taskId, payload, "running"); + const status = asFoldStatus(payload.status); + if (status !== undefined) { + applyStatus(entry, status); + } else if ( + // A usage-only tick must not restart a task that is already known. + // As the FIRST row for a task it does start it: the client fold + // does the same, so a child whose start row aged out of retention + // still reads as live on both sides instead of only one. + (payload.usageSnapshot !== true || existing === undefined) && + entry.status !== "terminal" && + entry.status !== "idle" + ) { + entry.status = "running"; + } + entry.payload = payload; + fillIdentity(entry, payload); + entries.set(taskId, entry); + break; + } + case "task.updated": { + const entry = existing ?? newEntry(taskId, payload, "running"); + const status = asFoldStatus(payload.status); + if (status !== undefined) { + applyStatus(entry, status); + } + entry.payload = payload; + fillIdentity(entry, payload); + entries.set(taskId, entry); + break; + } + case "task.completed": { + const entry = existing ?? newEntry(taskId, payload, "terminal"); + applyStatus(entry, "terminal"); + entry.payload = payload; + fillIdentity(entry, payload); + entries.set(taskId, entry); + break; + } + default: + break; + } + } + + // Consistency pass, mirroring the client fold: once a workflow coordinator + // has settled, members without a terminal row of their own cannot still be + // in flight — the run is over. The client shows them with the coordinator's + // outcome, so settling them here would overwrite a completed run with + // "interrupted". + for (const coordinator of entries.values()) { + if (!coordinator.isWorkflow || coordinator.status !== "terminal") { + continue; + } + for (const member of entries.values()) { + if (member.parentAgentId !== coordinator.taskId) { + continue; + } + if (member.status === "terminal" || member.status === "idle") { + continue; + } + member.status = "terminal"; + } + } + + const live: LiveAgentTask[] = []; + for (const [taskId, entry] of entries) { + if (!LIVE_STATUSES.has(entry.status)) { + continue; + } + const linkage: Record = {}; + for (const field of LINKAGE_FIELDS) { + if (entry.payload[field] !== undefined) { + linkage[field] = entry.payload[field]; + } + } + live.push({ taskId, linkage }); + } + return live; +} + +/** + * Stable per-task id so repeated settlement (Stop, then the session exiting) + * replaces the same row instead of piling up duplicates. + */ +export function settledTaskActivityId(threadId: ThreadId, taskId: string): EventId { + return EventId.make(`task-settled:${threadId}:${taskId}`); +} + +/** + * Marks every still-live agent task on the thread as settled: one persisted + * terminal `task.updated` per task, plus the matching liveness transition. + * + * Never fails the caller — Stop, session teardown, and startup all continue + * when settlement cannot read or write. + */ +export const settleThreadTasks = Effect.fn("settleThreadTasks")(function* (input: { + readonly threadId: ThreadId; + readonly status: "interrupted"; + readonly createdAt: string; +}) { + const activityRepository = yield* ProjectionThreadActivityRepository; + const orchestrationEngine = yield* OrchestrationEngineService; + const threadBackgroundLiveness = yield* ThreadBackgroundLivenessService; + const crypto = yield* Crypto.Crypto; + + const settle = Effect.gen(function* () { + // Complete task history, not the thread detail read's newest-N window: + // that window is shared with every other activity kind, so a busy thread + // could hide a live task's rows and then resurrect it as a settled card. + const rows = yield* activityRepository.listTaskLifecycleByThreadId({ + threadId: input.threadId, + }); + const liveTasks = selectLiveAgentTasks(rows); + if (liveTasks.length > LARGE_SETTLEMENT_LOG_THRESHOLD) { + yield* Effect.logInfo("settling a large background task fleet", { + threadId: input.threadId, + liveTaskCount: liveTasks.length, + }); + } + for (const task of liveTasks) { + const activity: OrchestrationThreadActivity = { + id: settledTaskActivityId(input.threadId, task.taskId), + createdAt: input.createdAt, + tone: "info", + kind: "task.updated", + summary: "Task interrupted", + payload: { + taskId: task.taskId, + status: input.status, + endedAt: input.createdAt, + ...task.linkage, + }, + turnId: null, + }; + yield* orchestrationEngine.dispatch({ + type: "thread.activity.append", + commandId: CommandId.make(`task-settle:${yield* crypto.randomUUIDv4}`), + threadId: input.threadId, + activity, + createdAt: input.createdAt, + }); + // Rows and registry settle together, so the sidebar pill and the + // composer banner can never disagree about the same task. + threadBackgroundLiveness.recordTaskLiveness({ + threadId: input.threadId, + taskId: task.taskId, + taskType: undefined, + status: input.status, + kind: "updated", + settledByHost: true, + }); + } + }); + + yield* settle.pipe( + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.failCause(cause) + : Effect.logWarning("failed to settle thread background tasks", { + threadId: input.threadId, + cause: Cause.pretty(cause), + }), + ), + ); +}); diff --git a/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts b/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts index fa3c948e4f3d..eef60ddf5713 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts @@ -112,6 +112,32 @@ const makeProjectionThreadActivityRepository = Effect.gen(function* () { `, }); + const listTaskLifecycleActivityRows = SqlSchema.findAll({ + Request: ListProjectionThreadActivitiesInput, + Result: ProjectionThreadActivityDbRowSchema, + execute: ({ threadId }) => + sql` + SELECT + activity_id AS "activityId", + thread_id AS "threadId", + turn_id AS "turnId", + tone, + kind, + summary, + payload_json AS "payload", + sequence, + created_at AS "createdAt" + FROM projection_thread_activities + WHERE thread_id = ${threadId} + AND kind IN ('task.started', 'task.progress', 'task.updated', 'task.completed') + ORDER BY + CASE WHEN sequence IS NULL THEN 0 ELSE 1 END ASC, + sequence ASC, + created_at ASC, + activity_id ASC + `, + }); + const listUserInputLifecycleActivityRows = SqlSchema.findAll({ Request: ListProjectionThreadActivitiesInput, Result: ProjectionThreadActivityDbRowSchema, @@ -172,6 +198,18 @@ const makeProjectionThreadActivityRepository = Effect.gen(function* () { Effect.map(mapActivityRows), ); + const listTaskLifecycleByThreadId: ProjectionThreadActivityRepositoryShape["listTaskLifecycleByThreadId"] = + (input) => + listTaskLifecycleActivityRows(input).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "ProjectionThreadActivityRepository.listTaskLifecycleByThreadId:query", + "ProjectionThreadActivityRepository.listTaskLifecycleByThreadId:decodeRows", + ), + ), + Effect.map(mapActivityRows), + ); + const listUserInputLifecycleByThreadId: ProjectionThreadActivityRepositoryShape["listUserInputLifecycleByThreadId"] = (input) => listUserInputLifecycleActivityRows(input).pipe( @@ -194,6 +232,7 @@ const makeProjectionThreadActivityRepository = Effect.gen(function* () { return { upsert, listByThreadId, + listTaskLifecycleByThreadId, listUserInputLifecycleByThreadId, deleteByThreadId, } satisfies ProjectionThreadActivityRepositoryShape; diff --git a/apps/server/src/persistence/Services/ProjectionThreadActivities.ts b/apps/server/src/persistence/Services/ProjectionThreadActivities.ts index e8c1e47a328b..f74b59cab79c 100644 --- a/apps/server/src/persistence/Services/ProjectionThreadActivities.ts +++ b/apps/server/src/persistence/Services/ProjectionThreadActivities.ts @@ -67,6 +67,17 @@ export interface ProjectionThreadActivityRepositoryShape { input: ListProjectionThreadActivitiesInput, ) => Effect.Effect, ProjectionRepositoryError>; + /** + * List a thread's complete task lifecycle history. + * + * Unwindowed on purpose: background-task settlement must see every task + * row, not the newest slice the thread detail read returns. Filters in + * SQLite so unrelated payloads do not enter server memory. + */ + readonly listTaskLifecycleByThreadId: ( + input: ListProjectionThreadActivitiesInput, + ) => Effect.Effect, ProjectionRepositoryError>; + /** * List activity rows used to derive pending user-input state. * diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index e74a4fa3c317..0539d9b2f7a1 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -261,8 +261,11 @@ const PlatformServicesLive = Layer.unwrap( const ReactorLayerLive = Layer.empty.pipe( Layer.provideMerge(OrchestrationReactorLive), - Layer.provideMerge(ProviderRuntimeIngestionLive), + // The command reactor drains runtime ingestion before settling background + // tasks on Stop, so ingestion must be provided to it (nothing in ingestion + // depends on the command reactor, so the order is safe to invert). Layer.provideMerge(ProviderCommandReactorLive), + Layer.provideMerge(ProviderRuntimeIngestionLive), Layer.provideMerge(CheckpointReactorLive), Layer.provideMerge(ThreadDeletionReactorLive), Layer.provideMerge(ThreadSettlementReactor.layer), diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index aa1b1a7f9788..86756c5164c8 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -19,6 +19,9 @@ import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSna import { ProviderSessionDirectoryPersistenceError } from "./provider/Errors.ts"; import * as ProviderService from "./provider/Services/ProviderService.ts"; import * as ProviderSessionDirectory from "./provider/Services/ProviderSessionDirectory.ts"; +import * as ProjectionThreadActivities from "./persistence/Services/ProjectionThreadActivities.ts"; +import * as ThreadBackgroundLiveness from "./orchestration/ThreadBackgroundLiveness.ts"; +import type { TaskActivityRow } from "./orchestration/ThreadTaskSettlement.ts"; import * as ServerRuntimeStartup from "./serverRuntimeStartup.ts"; const providerInstanceId = ProviderInstanceId.make("codex"); @@ -68,18 +71,36 @@ const queryWithThreads = (threads: ReadonlyArray>) getCommandReadModel: () => Effect.succeed({ threads } as never), }) as unknown as ProjectionSnapshotQuery.ProjectionSnapshotQuery["Service"]; +const activityRepositoryWith = ( + activitiesByThreadId: Readonly>>, +) => + ({ + listTaskLifecycleByThreadId: ({ threadId }: { readonly threadId: ThreadId }) => + Effect.succeed(activitiesByThreadId[threadId] ?? []), + }) as unknown as ProjectionThreadActivities.ProjectionThreadActivityRepository["Service"]; + const runReconciliation = (input: { readonly threads: ReadonlyArray>; readonly liveThreadIds?: ReadonlyArray; readonly providerService?: ProviderService.ProviderService["Service"]; readonly directory: ProviderSessionDirectory.ProviderSessionDirectory["Service"]; readonly dispatch: OrchestrationEngine.OrchestrationEngineService["Service"]["dispatch"]; + readonly activitiesByThreadId?: Readonly>>; + readonly backgroundLiveness?: ThreadBackgroundLiveness.ThreadBackgroundLivenessService["Service"]; }) => ServerRuntimeStartup.reconcileProviderSessions.pipe( Effect.provideService( ProjectionSnapshotQuery.ProjectionSnapshotQuery, queryWithThreads(input.threads), ), + Effect.provideService( + ProjectionThreadActivities.ProjectionThreadActivityRepository, + activityRepositoryWith(input.activitiesByThreadId ?? {}), + ), + Effect.provideService( + ThreadBackgroundLiveness.ThreadBackgroundLivenessService, + input.backgroundLiveness ?? ThreadBackgroundLiveness.make(), + ), Effect.provideService( ProviderService.ProviderService, input.providerService ?? makeProviderService(input.liveThreadIds), @@ -654,7 +675,126 @@ it.effect("does not fail startup when the live provider session inventory cannot subscribeDomainEvents: Effect.succeed(Stream.empty), latestSequence: Effect.succeed(0), }), + Effect.provideService( + ThreadBackgroundLiveness.ThreadBackgroundLivenessService, + ThreadBackgroundLiveness.make(), + ), + Effect.provideService( + ProjectionThreadActivities.ProjectionThreadActivityRepository, + activityRepositoryWith({}), + ), Effect.provide(NodeServices.layer), Effect.tap(() => Effect.sync(() => assert.equal(queried, false))), ); }); + +it.effect("settles background agent tasks left running by an orphaned session", () => { + const orphaned = makeThread("thread-settle-tasks", "running", TurnId.make("turn-settle-tasks")); + const taskId = "collab-child-1"; + const linkage = { agentKind: "agent", title: "math_one", timelineBypass: true } as const; + const dispatched: OrchestrationCommand[] = []; + const liveness = ThreadBackgroundLiveness.make(); + // The registry is empty after a real restart; arming it here proves the + // settlement clears rows and registry together. + liveness.recordTaskLiveness({ + threadId: orphaned.id, + taskId, + taskType: undefined, + status: "running", + kind: "updated", + }); + + return runReconciliation({ + threads: [orphaned], + activitiesByThreadId: { + [orphaned.id]: [ + { kind: "task.started", payload: { taskId, ...linkage } }, + { kind: "task.updated", payload: { taskId, status: "running", ...linkage } }, + ], + }, + backgroundLiveness: liveness, + directory: { + getBinding: () => Effect.succeed(Option.none()), + upsert: () => Effect.void, + getProvider: () => Effect.die("unused"), + listThreadIds: () => Effect.die("unused"), + listBindings: () => Effect.die("unused"), + }, + dispatch: (command) => { + dispatched.push(command); + return Effect.succeed({ sequence: dispatched.length }); + }, + }).pipe( + Effect.tap(() => + Effect.sync(() => { + const settlement = dispatched.find((command) => command.type === "thread.activity.append"); + assert.isDefined(settlement); + assert.equal(settlement.activity.id, `task-settled:${orphaned.id}:${taskId}`); + assert.equal(settlement.activity.kind, "task.updated"); + assert.deepStrictEqual(settlement.activity.payload, { + taskId, + status: "interrupted", + endedAt: settlement.createdAt, + ...linkage, + }); + assert.equal(liveness.getThreadBackgroundLiveness(orphaned.id), null); + // The orphaned session itself is still reconciled to error. + assert.isDefined(dispatched.find((command) => command.type === "thread.session.set")); + }), + ), + ); +}); + +it("settles background agents on a ready thread whose turn already ended", () => { + // The headline case: the turn finished, so the session is `ready` with no + // active turn, but its children kept working and did not survive the + // restart. Archived and deleted threads must not be written to. + const taskId = "collab-child-1"; + const linkage = { agentKind: "agent", title: "math_one", timelineBypass: true } as const; + const rows = [ + { kind: "task.started", payload: { taskId, ...linkage } }, + { kind: "task.updated", payload: { taskId, status: "running", ...linkage } }, + ]; + const ready = makeThread("thread-settle-ready", "ready"); + const archived = makeThread("thread-settle-archived", "ready", null, updatedAt); + const deleted = makeThread("thread-settle-deleted", "ready", null, null, updatedAt); + const stopped = makeThread("thread-settle-stopped", "stopped"); + const dispatched: OrchestrationCommand[] = []; + + return runReconciliation({ + threads: [ready, archived, deleted, stopped], + activitiesByThreadId: { + [ready.id]: rows, + [archived.id]: rows, + [deleted.id]: rows, + [stopped.id]: rows, + }, + directory: { + getBinding: () => Effect.succeed(Option.none()), + upsert: () => Effect.void, + getProvider: () => Effect.die("unused"), + listThreadIds: () => Effect.die("unused"), + listBindings: () => Effect.die("unused"), + }, + dispatch: (command) => { + dispatched.push(command); + return Effect.succeed({ sequence: dispatched.length }); + }, + }).pipe( + Effect.tap(() => + Effect.sync(() => { + assert.deepStrictEqual( + dispatched + .filter((command) => command.type === "thread.activity.append") + .map((command) => command.threadId), + [ready.id], + ); + // A ready session is not an orphaned turn, so nothing else changes. + assert.deepStrictEqual( + dispatched.filter((command) => command.type === "thread.session.set"), + [], + ); + }), + ), + ); +}); diff --git a/apps/server/src/serverRuntimeStartup.ts b/apps/server/src/serverRuntimeStartup.ts index 064796b2810b..965158de9d41 100644 --- a/apps/server/src/serverRuntimeStartup.ts +++ b/apps/server/src/serverRuntimeStartup.ts @@ -31,6 +31,7 @@ import * as ExternalLauncher from "./process/externalLauncher.ts"; import * as OrchestrationEngine from "./orchestration/Services/OrchestrationEngine.ts"; import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSnapshotQuery.ts"; import * as OrchestrationReactor from "./orchestration/Services/OrchestrationReactor.ts"; +import { settleThreadTasks } from "./orchestration/ThreadTaskSettlement.ts"; import * as ServerLifecycleEvents from "./serverLifecycleEvents.ts"; import * as ServerSettings from "./serverSettings.ts"; import * as AnalyticsService from "./telemetry/AnalyticsService.ts"; @@ -457,6 +458,29 @@ export const reconcileProviderSessions = Effect.gen(function* () { !liveThreadIds.has(thread.id), ); + // Background work outlives the turn, so the headline case is wider than an + // orphaned turn: a thread whose turn ended is `ready` with no active turn + // while its children keep running. Any thread with a session this process + // does not own has lost that work, and a restart leaves the in-memory + // liveness registry empty while the persisted rows still read "running". + // Archived and deleted threads are skipped — settling writes rows. + for (const thread of threads) { + if ( + thread.session === null || + thread.session.status === "stopped" || + thread.archivedAt !== null || + thread.deletedAt !== null || + liveThreadIds.has(thread.id) + ) { + continue; + } + yield* settleThreadTasks({ + threadId: thread.id, + status: "interrupted", + createdAt: DateTime.formatIso(yield* DateTime.now), + }); + } + for (const thread of orphanedThreads) { const session = thread.session; if (session === null) { diff --git a/docs/user/composer.md b/docs/user/composer.md index 13b4527ad261..8b9cc9932f4f 100644 --- a/docs/user/composer.md +++ b/docs/user/composer.md @@ -160,6 +160,12 @@ On web and desktop, loading and syncing statuses fill the available banner width stash tab. Task progress appears above the composer, while the timeline's working timer shows only elapsed time. +When background agents keep working after a turn ends, a banner above the composer counts them +and offers **Stop**. Pressing Stop, or the provider session ending, marks any agents still listed +as working as interrupted, so the banner clears and the Agents panel stops showing them as +working, even when the provider never reports those agents again. Restarting T3 Code does the +same for threads whose provider session did not survive. + On web and desktop, additional notices peek out above the attached banner. Hover over the peek to reveal them, or focus **Show other notices** with `Tab` and press `Enter` or `Space`. Press `Escape` to close the stack and return focus to that control. On a touchscreen, tap the peek to diff --git a/packages/client-runtime/src/state/subagentRuntime.test.ts b/packages/client-runtime/src/state/subagentRuntime.test.ts index f87531f15808..f6148531ef75 100644 --- a/packages/client-runtime/src/state/subagentRuntime.test.ts +++ b/packages/client-runtime/src/state/subagentRuntime.test.ts @@ -914,3 +914,52 @@ describe("nested agents vs subagent shells", () => { expect(agents.map((agent) => agent.id)).toEqual(["nested-1"]); }); }); + +describe("host-settled background agents", () => { + it("reconciles the live count as soon as a terminal row lands", () => { + // Replay of a real Codex collab child: turnStarted/turnCompleted cycles, + // the host's settlement row, and the provider's own late idle row after + // it. The last word is terminal, so nothing is left working. + const linkage = { + taskId: "collab-child-1", + title: "math_one", + role: "general-purpose", + agentPath: "agents/math_one", + timelineBypass: true, + }; + const agents = fold([ + activity("task.started", { ...linkage }), + activity("task.updated", { ...linkage, status: "running" }), + activity("task.updated", { ...linkage, status: "idle" }), + activity("task.updated", { ...linkage, status: "running" }), + activity("task.updated", { ...linkage, status: "interrupted" }), + activity("task.updated", { ...linkage, status: "interrupted" }), + activity("task.updated", { ...linkage, status: "idle" }), + activity("task.updated", { ...linkage, status: "interrupted" }), + ]); + + const model = deriveAgentPanelModel({ agents }); + expect(model.liveCount).toBe(0); + expect(model.settledCount).toBe(1); + expect(agents[0]?.status).toBe("interrupted"); + }); + + it("a settlement row alone settles a child whose start row aged out", () => { + const agents = fold([ + activity("task.progress", { + taskId: "collab-child-2", + agentKind: "agent", + summary: "still working", + }), + activity("task.updated", { + taskId: "collab-child-2", + agentKind: "agent", + status: "interrupted", + endedAt: "2026-08-01T10:05:00.000Z", + }), + ]); + + expect(deriveAgentPanelModel({ agents }).liveCount).toBe(0); + expect(agents[0]?.completedAt).toBe("2026-08-01T10:05:00.000Z"); + }); +}); From 206340965d8ad6e56af2adde2a9ddad5a4116760 Mon Sep 17 00:00:00 2001 From: amanthanvi Date: Thu, 3 Sep 2026 05:21:02 -0400 Subject: [PATCH 2/8] fix(server): keep agent linkage and settle stopped sessions at startup MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two gaps in background-agent settlement. The fold copied only the newest row's payload onto the synthesized terminal row, so when later lifecycle rows carried just a task id and a status the settled row lost agentKind and its workflow grouping — and once the start row aged out of the client's activity window, that row was the only one left and the agent disappeared from the panel. Startup also skipped threads whose session was already stopped, so a process that died between marking a session stopped and settling its children left those rows running on every later boot. Linkage is now merged across every row for a task, newest non-empty value winning, and covers the full set ingestion stamps minus per-row state. Startup settles any thread it does not own regardless of session status, reading every candidate thread's task history in one batched, chunked query instead of a query per thread. Implemented by Claude Opus 5 via Claude Code. --- .../Layers/ProviderRuntimeIngestion.test.ts | 54 ++++ .../src/orchestration/ThreadTaskSettlement.ts | 255 ++++++++++++------ .../Layers/ProjectionThreadActivities.ts | 60 +++++ .../Services/ProjectionThreadActivities.ts | 18 ++ .../serverRuntimeStartup.reconcile.test.ts | 32 ++- apps/server/src/serverRuntimeStartup.ts | 43 +-- 6 files changed, 353 insertions(+), 109 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 09c9400d879a..a56fa519d18a 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -3785,8 +3785,10 @@ describe("selectLiveAgentTasks", () => { ]); expect(live.map((task) => task.taskId).toSorted()).toEqual(["live-explicit", "live-implicit"]); + // The status row carried no title; the start row's survives the merge. expect(live.find((task) => task.taskId === "live-explicit")?.linkage).toEqual({ agentKind: "agent", + title: "one", }); }); @@ -3823,6 +3825,58 @@ describe("selectLiveAgentTasks", () => { expect(live.map((task) => task.taskId).toSorted()).toEqual(["other-1", "other-member-1"]); }); + it("keeps linkage from earlier rows when a later row carries only a status", () => { + // Terminal and status-patch rows commonly carry nothing but taskId and + // status. If the settled row copied only that newest payload it would + // land with no agentKind — and once the start row falls out of the + // client's activity window, the client would read the settled row as + // background work and drop the agent from the panel entirely. + const live = selectLiveAgentTasks([ + row("a1", "task.started", { + taskId: "wf-member-1", + ...agent, + taskType: "local_agent", + title: "math_one", + role: "general-purpose", + agentPath: "agents/math_one", + timelineBypass: true, + model: "gpt-5-codex", + effort: "high", + parentAgentId: "wf-1", + workflowName: "release", + agentIndex: 2, + phaseIndex: 1, + phaseTitle: "Implement", + attempt: 1, + outputFile: "/tmp/out.md", + }), + row("a2", "task.updated", { taskId: "wf-member-1", status: "running" }), + ]); + + expect(live).toEqual([ + { + taskId: "wf-member-1", + linkage: { + agentKind: "agent", + taskType: "local_agent", + title: "math_one", + role: "general-purpose", + agentPath: "agents/math_one", + timelineBypass: true, + model: "gpt-5-codex", + effort: "high", + parentAgentId: "wf-1", + workflowName: "release", + agentIndex: 2, + phaseIndex: 1, + phaseTitle: "Implement", + attempt: 1, + outputFile: "/tmp/out.md", + }, + }, + ]); + }); + it("carries the newest linkage onto the live task", () => { const live = selectLiveAgentTasks([ row("a1", "task.started", { taskId: "child", ...agent, title: "old-name" }), diff --git a/apps/server/src/orchestration/ThreadTaskSettlement.ts b/apps/server/src/orchestration/ThreadTaskSettlement.ts index 74d622b38ca2..ef2125be3bff 100644 --- a/apps/server/src/orchestration/ThreadTaskSettlement.ts +++ b/apps/server/src/orchestration/ThreadTaskSettlement.ts @@ -46,17 +46,32 @@ export interface TaskActivityRow { /** * Linkage the synthesized terminal row carries forward so it stays a - * self-describing agent row: `agentKind` keeps it on the Agents surface and - * the rest keep the card's identity when the start row ages out. + * self-describing agent row. This is ingestion's `taskLinkageActivityFields` + * minus the per-row state (`status`, `error`, `typedUsage`, `toolUseId`): + * `agentKind` keeps the row on the Agents surface, the workflow fields keep it + * grouped under its coordinator, and the rest keep the card's identity. The + * settled row is often the ONLY row left inside the client's activity window, + * so anything missing here is lost from the panel. */ const LINKAGE_FIELDS = [ "agentKind", + "taskType", + "agentId", "title", "role", - "agentPath", - "timelineBypass", "model", "effort", + "agentPath", + "timelineBypass", + "parentAgentId", + "workflowName", + "agentIndex", + "phaseIndex", + "phaseTitle", + "phases", + "attempt", + "runHandles", + "outputFile", ] as const; // Collapsed mirror of the client fold (subagentRuntime.foldSubagentActivities): @@ -86,8 +101,13 @@ const KNOWN_STATUSES: ReadonlyMap = new Map; + /** + * Linkage accumulated across every row for the task, newest value winning. + * Merged rather than replaced: lifecycle rows commonly carry only a task id + * and a status, and copying just the newest payload would strip the identity + * the start row established. + */ + readonly linkage: Record; /** A workflow coordinator, whose settling ends its members' run. */ isWorkflow: boolean; /** The coordinator this task belongs to, if any. */ @@ -111,22 +131,29 @@ function asTrimmed(value: unknown): string | undefined { } /** - * Identity the workflow cascade needs, filled from every row and never - * downgraded — the client fold's getOrCreate/fillMetadata rules, narrowed to - * the two fields that decide whether a coordinator owns a task. + * Accumulates a row's linkage and the identity the workflow cascade needs. + * Never downgrades: an absent or blank value leaves the known one alone, the + * same rule the client fold's fillMetadata applies. */ -function fillIdentity(entry: FoldEntry, payload: Record): void { - if (asTrimmed(payload.taskType) === "local_workflow") { +function mergeLinkage(entry: FoldEntry, payload: Record): void { + for (const field of LINKAGE_FIELDS) { + const value = payload[field]; + if (value === undefined || (typeof value === "string" && value.trim().length === 0)) { + continue; + } + entry.linkage[field] = value; + } + if (asTrimmed(entry.linkage.taskType) === "local_workflow") { entry.isWorkflow = true; } - const parentAgentId = asTrimmed(payload.parentAgentId); + const parentAgentId = asTrimmed(entry.linkage.parentAgentId); if (parentAgentId !== undefined) { entry.parentAgentId = parentAgentId; } } -function newEntry(taskId: string, payload: Record, status: FoldStatus): FoldEntry { - return { taskId, status, payload, isWorkflow: false, parentAgentId: undefined }; +function newEntry(taskId: string, status: FoldStatus): FoldEntry { + return { taskId, status, linkage: {}, isWorkflow: false, parentAgentId: undefined }; } /** Rows without the server-stamped agent classification are background work. */ @@ -183,17 +210,16 @@ export function selectLiveAgentTasks( if (!isAgentRow(payload)) { break; } - const entry = existing ?? newEntry(taskId, payload, "running"); + const entry = existing ?? newEntry(taskId, "running"); if (existing !== undefined && existing.status === "idle") { entry.status = "running"; } - entry.payload = payload; - fillIdentity(entry, payload); + mergeLinkage(entry, payload); entries.set(taskId, entry); break; } case "task.progress": { - const entry = existing ?? newEntry(taskId, payload, "running"); + const entry = existing ?? newEntry(taskId, "running"); const status = asFoldStatus(payload.status); if (status !== undefined) { applyStatus(entry, status); @@ -208,27 +234,24 @@ export function selectLiveAgentTasks( ) { entry.status = "running"; } - entry.payload = payload; - fillIdentity(entry, payload); + mergeLinkage(entry, payload); entries.set(taskId, entry); break; } case "task.updated": { - const entry = existing ?? newEntry(taskId, payload, "running"); + const entry = existing ?? newEntry(taskId, "running"); const status = asFoldStatus(payload.status); if (status !== undefined) { applyStatus(entry, status); } - entry.payload = payload; - fillIdentity(entry, payload); + mergeLinkage(entry, payload); entries.set(taskId, entry); break; } case "task.completed": { - const entry = existing ?? newEntry(taskId, payload, "terminal"); + const entry = existing ?? newEntry(taskId, "terminal"); applyStatus(entry, "terminal"); - entry.payload = payload; - fillIdentity(entry, payload); + mergeLinkage(entry, payload); entries.set(taskId, entry); break; } @@ -259,16 +282,9 @@ export function selectLiveAgentTasks( const live: LiveAgentTask[] = []; for (const [taskId, entry] of entries) { - if (!LIVE_STATUSES.has(entry.status)) { - continue; - } - const linkage: Record = {}; - for (const field of LINKAGE_FIELDS) { - if (entry.payload[field] !== undefined) { - linkage[field] = entry.payload[field]; - } + if (LIVE_STATUSES.has(entry.status)) { + live.push({ taskId, linkage: entry.linkage }); } - live.push({ taskId, linkage }); } return live; } @@ -281,6 +297,77 @@ export function settledTaskActivityId(threadId: ThreadId, taskId: string): Event return EventId.make(`task-settled:${threadId}:${taskId}`); } +/** + * Writes the terminal rows for one thread's live tasks and mirrors each into + * the liveness registry. Shared by the per-thread and batched entry points so + * the row shape can never diverge between them. + */ +const writeSettledTasks = Effect.fn("writeSettledTasks")(function* (input: { + readonly threadId: ThreadId; + readonly liveTasks: ReadonlyArray; + readonly status: "interrupted"; + readonly createdAt: string; +}) { + const orchestrationEngine = yield* OrchestrationEngineService; + const threadBackgroundLiveness = yield* ThreadBackgroundLivenessService; + const crypto = yield* Crypto.Crypto; + + if (input.liveTasks.length > LARGE_SETTLEMENT_LOG_THRESHOLD) { + yield* Effect.logInfo("settling a large background task fleet", { + threadId: input.threadId, + liveTaskCount: input.liveTasks.length, + }); + } + for (const task of input.liveTasks) { + const activity: OrchestrationThreadActivity = { + id: settledTaskActivityId(input.threadId, task.taskId), + createdAt: input.createdAt, + tone: "info", + kind: "task.updated", + summary: "Task interrupted", + payload: { + taskId: task.taskId, + status: input.status, + endedAt: input.createdAt, + ...task.linkage, + }, + turnId: null, + }; + yield* orchestrationEngine.dispatch({ + type: "thread.activity.append", + commandId: CommandId.make(`task-settle:${yield* crypto.randomUUIDv4}`), + threadId: input.threadId, + activity, + createdAt: input.createdAt, + }); + // Rows and registry settle together, so the sidebar pill and the composer + // banner can never disagree about the same task. + threadBackgroundLiveness.recordTaskLiveness({ + threadId: input.threadId, + taskId: task.taskId, + taskType: undefined, + status: input.status, + kind: "updated", + settledByHost: true, + }); + } +}); + +/** Settlement never fails its caller: Stop and startup continue regardless. */ +const withSettlementRecovery = + (threadIds: ReadonlyArray) => + (effect: Effect.Effect) => + effect.pipe( + Effect.catchCause((cause: Cause.Cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.failCause(cause) + : Effect.logWarning("failed to settle thread background tasks", { + threadIds, + cause: Cause.pretty(cause), + }), + ), + ); + /** * Marks every still-live agent task on the thread as settled: one persisted * terminal `task.updated` per task, plus the matching liveness transition. @@ -294,9 +381,6 @@ export const settleThreadTasks = Effect.fn("settleThreadTasks")(function* (input readonly createdAt: string; }) { const activityRepository = yield* ProjectionThreadActivityRepository; - const orchestrationEngine = yield* OrchestrationEngineService; - const threadBackgroundLiveness = yield* ThreadBackgroundLivenessService; - const crypto = yield* Crypto.Crypto; const settle = Effect.gen(function* () { // Complete task history, not the thread detail read's newest-N window: @@ -305,56 +389,59 @@ export const settleThreadTasks = Effect.fn("settleThreadTasks")(function* (input const rows = yield* activityRepository.listTaskLifecycleByThreadId({ threadId: input.threadId, }); - const liveTasks = selectLiveAgentTasks(rows); - if (liveTasks.length > LARGE_SETTLEMENT_LOG_THRESHOLD) { - yield* Effect.logInfo("settling a large background task fleet", { - threadId: input.threadId, - liveTaskCount: liveTasks.length, - }); + yield* writeSettledTasks({ + threadId: input.threadId, + liveTasks: selectLiveAgentTasks(rows), + status: input.status, + createdAt: input.createdAt, + }); + }); + + yield* settle.pipe(withSettlementRecovery([input.threadId])); +}); + +/** + * Batched settlement for startup reconciliation: one activity read for every + * thread this process does not own, then rows only for the threads that + * actually have live background work. Startup runs across the whole thread + * list, so a query per thread would be the wrong shape. + */ +export const settleThreadsTasks = Effect.fn("settleThreadsTasks")(function* (input: { + readonly threadIds: ReadonlyArray; + readonly status: "interrupted"; + readonly createdAt: string; +}) { + const activityRepository = yield* ProjectionThreadActivityRepository; + + const settle = Effect.gen(function* () { + if (input.threadIds.length === 0) { + return; } - for (const task of liveTasks) { - const activity: OrchestrationThreadActivity = { - id: settledTaskActivityId(input.threadId, task.taskId), - createdAt: input.createdAt, - tone: "info", - kind: "task.updated", - summary: "Task interrupted", - payload: { - taskId: task.taskId, - status: input.status, - endedAt: input.createdAt, - ...task.linkage, - }, - turnId: null, - }; - yield* orchestrationEngine.dispatch({ - type: "thread.activity.append", - commandId: CommandId.make(`task-settle:${yield* crypto.randomUUIDv4}`), - threadId: input.threadId, - activity, - createdAt: input.createdAt, - }); - // Rows and registry settle together, so the sidebar pill and the - // composer banner can never disagree about the same task. - threadBackgroundLiveness.recordTaskLiveness({ - threadId: input.threadId, - taskId: task.taskId, - taskType: undefined, + const rows = yield* activityRepository.listTaskLifecycleByThreadIds({ + threadIds: input.threadIds, + }); + const rowsByThreadId = new Map>(); + for (const row of rows) { + const existing = rowsByThreadId.get(row.threadId); + if (existing) { + existing.push(row); + } else { + rowsByThreadId.set(row.threadId, [row]); + } + } + for (const [threadId, threadRows] of rowsByThreadId) { + const liveTasks = selectLiveAgentTasks(threadRows); + if (liveTasks.length === 0) { + continue; + } + yield* writeSettledTasks({ + threadId, + liveTasks, status: input.status, - kind: "updated", - settledByHost: true, + createdAt: input.createdAt, }); } }); - yield* settle.pipe( - Effect.catchCause((cause) => - Cause.hasInterruptsOnly(cause) - ? Effect.failCause(cause) - : Effect.logWarning("failed to settle thread background tasks", { - threadId: input.threadId, - cause: Cause.pretty(cause), - }), - ), - ); + yield* settle.pipe(withSettlementRecovery(input.threadIds)); }); diff --git a/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts b/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts index eef60ddf5713..f828dc2adc4a 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts @@ -10,6 +10,7 @@ import { toPersistenceDecodeError, toPersistenceSqlError } from "../Errors.ts"; import { DeleteProjectionThreadActivitiesInput, + ListProjectionThreadActivitiesByThreadIdsInput, ListProjectionThreadActivitiesInput, ProjectionThreadActivity, ProjectionThreadActivityRepository, @@ -38,6 +39,17 @@ const mapActivityRows = ( createdAt: row.createdAt, })); +/** SQLite's host-parameter ceiling is 999; stay well inside it. */ +const THREAD_ID_BATCH_SIZE = 500; + +const chunk = (items: ReadonlyArray, size: number): ReadonlyArray> => { + const chunks: A[][] = []; + for (let index = 0; index < items.length; index += size) { + chunks.push(items.slice(index, index + size)); + } + return chunks; +}; + function toPersistenceSqlOrDecodeError(sqlOperation: string, decodeOperation: string) { return (cause: unknown) => Schema.isSchemaError(cause) @@ -138,6 +150,33 @@ const makeProjectionThreadActivityRepository = Effect.gen(function* () { `, }); + const listTaskLifecycleActivityRowsByThreadIds = SqlSchema.findAll({ + Request: ListProjectionThreadActivitiesByThreadIdsInput, + Result: ProjectionThreadActivityDbRowSchema, + execute: ({ threadIds }) => + sql` + SELECT + activity_id AS "activityId", + thread_id AS "threadId", + turn_id AS "turnId", + tone, + kind, + summary, + payload_json AS "payload", + sequence, + created_at AS "createdAt" + FROM projection_thread_activities + WHERE ${sql.in("thread_id", threadIds)} + AND kind IN ('task.started', 'task.progress', 'task.updated', 'task.completed') + ORDER BY + thread_id ASC, + CASE WHEN sequence IS NULL THEN 0 ELSE 1 END ASC, + sequence ASC, + created_at ASC, + activity_id ASC + `, + }); + const listUserInputLifecycleActivityRows = SqlSchema.findAll({ Request: ListProjectionThreadActivitiesInput, Result: ProjectionThreadActivityDbRowSchema, @@ -210,6 +249,26 @@ const makeProjectionThreadActivityRepository = Effect.gen(function* () { Effect.map(mapActivityRows), ); + const listTaskLifecycleByThreadIds: ProjectionThreadActivityRepositoryShape["listTaskLifecycleByThreadIds"] = + ({ threadIds }) => + // SQLite caps host parameters per statement, so long id lists are read + // in chunks. Threads never straddle a chunk, so per-thread ordering is + // preserved by concatenation. + Effect.forEach( + chunk(threadIds, THREAD_ID_BATCH_SIZE), + (threadIdChunk) => + listTaskLifecycleActivityRowsByThreadIds({ threadIds: threadIdChunk }).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "ProjectionThreadActivityRepository.listTaskLifecycleByThreadIds:query", + "ProjectionThreadActivityRepository.listTaskLifecycleByThreadIds:decodeRows", + ), + ), + Effect.map(mapActivityRows), + ), + { concurrency: 1 }, + ).pipe(Effect.map((chunks) => chunks.flat())); + const listUserInputLifecycleByThreadId: ProjectionThreadActivityRepositoryShape["listUserInputLifecycleByThreadId"] = (input) => listUserInputLifecycleActivityRows(input).pipe( @@ -233,6 +292,7 @@ const makeProjectionThreadActivityRepository = Effect.gen(function* () { upsert, listByThreadId, listTaskLifecycleByThreadId, + listTaskLifecycleByThreadIds, listUserInputLifecycleByThreadId, deleteByThreadId, } satisfies ProjectionThreadActivityRepositoryShape; diff --git a/apps/server/src/persistence/Services/ProjectionThreadActivities.ts b/apps/server/src/persistence/Services/ProjectionThreadActivities.ts index f74b59cab79c..e4de19fc8a85 100644 --- a/apps/server/src/persistence/Services/ProjectionThreadActivities.ts +++ b/apps/server/src/persistence/Services/ProjectionThreadActivities.ts @@ -38,6 +38,12 @@ export const ListProjectionThreadActivitiesInput = Schema.Struct({ }); export type ListProjectionThreadActivitiesInput = typeof ListProjectionThreadActivitiesInput.Type; +export const ListProjectionThreadActivitiesByThreadIdsInput = Schema.Struct({ + threadIds: Schema.Array(ThreadId), +}); +export type ListProjectionThreadActivitiesByThreadIdsInput = + typeof ListProjectionThreadActivitiesByThreadIdsInput.Type; + export const DeleteProjectionThreadActivitiesInput = Schema.Struct({ threadId: ThreadId, }); @@ -78,6 +84,18 @@ export interface ProjectionThreadActivityRepositoryShape { input: ListProjectionThreadActivitiesInput, ) => Effect.Effect, ProjectionRepositoryError>; + /** + * List task lifecycle history for many threads in one pass. + * + * Startup reconciliation settles background work across every thread this + * process does not own, so it reads them together instead of issuing one + * query per thread. Rows come back grouped-ready: ordered within each + * thread, threads interleaved. + */ + readonly listTaskLifecycleByThreadIds: ( + input: ListProjectionThreadActivitiesByThreadIdsInput, + ) => Effect.Effect, ProjectionRepositoryError>; + /** * List activity rows used to derive pending user-input state. * diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index 86756c5164c8..b77a53f8d57d 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -73,10 +73,22 @@ const queryWithThreads = (threads: ReadonlyArray>) const activityRepositoryWith = ( activitiesByThreadId: Readonly>>, + batchReads: Array> = [], ) => ({ listTaskLifecycleByThreadId: ({ threadId }: { readonly threadId: ThreadId }) => Effect.succeed(activitiesByThreadId[threadId] ?? []), + listTaskLifecycleByThreadIds: ({ + threadIds, + }: { + readonly threadIds: ReadonlyArray; + }) => + Effect.sync(() => { + batchReads.push(threadIds); + return threadIds.flatMap((threadId) => + (activitiesByThreadId[threadId] ?? []).map((row) => ({ ...row, threadId })), + ); + }), }) as unknown as ProjectionThreadActivities.ProjectionThreadActivityRepository["Service"]; const runReconciliation = (input: { @@ -87,6 +99,8 @@ const runReconciliation = (input: { readonly dispatch: OrchestrationEngine.OrchestrationEngineService["Service"]["dispatch"]; readonly activitiesByThreadId?: Readonly>>; readonly backgroundLiveness?: ThreadBackgroundLiveness.ThreadBackgroundLivenessService["Service"]; + /** Collects the thread-id lists the batched activity read was called with. */ + readonly batchReads?: Array>; }) => ServerRuntimeStartup.reconcileProviderSessions.pipe( Effect.provideService( @@ -95,7 +109,7 @@ const runReconciliation = (input: { ), Effect.provideService( ProjectionThreadActivities.ProjectionThreadActivityRepository, - activityRepositoryWith(input.activitiesByThreadId ?? {}), + activityRepositoryWith(input.activitiesByThreadId ?? {}, input.batchReads), ), Effect.provideService( ThreadBackgroundLiveness.ThreadBackgroundLivenessService, @@ -745,7 +759,7 @@ it.effect("settles background agent tasks left running by an orphaned session", ); }); -it("settles background agents on a ready thread whose turn already ended", () => { +it("settles background agents on ready and stopped threads, and writes nothing when idle", () => { // The headline case: the turn finished, so the session is `ready` with no // active turn, but its children kept working and did not survive the // restart. Archived and deleted threads must not be written to. @@ -758,11 +772,18 @@ it("settles background agents on a ready thread whose turn already ended", () => const ready = makeThread("thread-settle-ready", "ready"); const archived = makeThread("thread-settle-archived", "ready", null, updatedAt); const deleted = makeThread("thread-settle-deleted", "ready", null, null, updatedAt); + // A stopped session still counts: the process can die between marking the + // session stopped and settling its children, and no later boot would ever + // pick them up if stopped threads were skipped. const stopped = makeThread("thread-settle-stopped", "stopped"); + // Nothing to settle here, so it must produce no writes at all. + const quiet = makeThread("thread-settle-quiet", "ready"); const dispatched: OrchestrationCommand[] = []; + const batchReads: Array> = []; return runReconciliation({ - threads: [ready, archived, deleted, stopped], + threads: [ready, archived, deleted, stopped, quiet], + batchReads, activitiesByThreadId: { [ready.id]: rows, [archived.id]: rows, @@ -787,8 +808,11 @@ it("settles background agents on a ready thread whose turn already ended", () => dispatched .filter((command) => command.type === "thread.activity.append") .map((command) => command.threadId), - [ready.id], + [ready.id, stopped.id], ); + // One batched read for every candidate thread, not a query per thread, + // and archived/deleted threads are never even read. + assert.deepStrictEqual(batchReads, [[ready.id, stopped.id, quiet.id]]); // A ready session is not an orphaned turn, so nothing else changes. assert.deepStrictEqual( dispatched.filter((command) => command.type === "thread.session.set"), diff --git a/apps/server/src/serverRuntimeStartup.ts b/apps/server/src/serverRuntimeStartup.ts index 965158de9d41..082814c4f0d3 100644 --- a/apps/server/src/serverRuntimeStartup.ts +++ b/apps/server/src/serverRuntimeStartup.ts @@ -31,7 +31,7 @@ import * as ExternalLauncher from "./process/externalLauncher.ts"; import * as OrchestrationEngine from "./orchestration/Services/OrchestrationEngine.ts"; import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSnapshotQuery.ts"; import * as OrchestrationReactor from "./orchestration/Services/OrchestrationReactor.ts"; -import { settleThreadTasks } from "./orchestration/ThreadTaskSettlement.ts"; +import { settleThreadsTasks } from "./orchestration/ThreadTaskSettlement.ts"; import * as ServerLifecycleEvents from "./serverLifecycleEvents.ts"; import * as ServerSettings from "./serverSettings.ts"; import * as AnalyticsService from "./telemetry/AnalyticsService.ts"; @@ -460,26 +460,27 @@ export const reconcileProviderSessions = Effect.gen(function* () { // Background work outlives the turn, so the headline case is wider than an // orphaned turn: a thread whose turn ended is `ready` with no active turn - // while its children keep running. Any thread with a session this process - // does not own has lost that work, and a restart leaves the in-memory - // liveness registry empty while the persisted rows still read "running". - // Archived and deleted threads are skipped — settling writes rows. - for (const thread of threads) { - if ( - thread.session === null || - thread.session.status === "stopped" || - thread.archivedAt !== null || - thread.deletedAt !== null || - liveThreadIds.has(thread.id) - ) { - continue; - } - yield* settleThreadTasks({ - threadId: thread.id, - status: "interrupted", - createdAt: DateTime.formatIso(yield* DateTime.now), - }); - } + // while its children keep running. A stopped session qualifies too — the + // process can die between marking the session stopped and settling, and + // that thread would otherwise never be settled on any later boot. Any + // thread this process does not own has lost its background work, and a + // restart leaves the in-memory liveness registry empty while the persisted + // rows still read "running". Archived and deleted threads are skipped — + // settling writes rows. One batched read, not a query per thread. + const settleableThreadIds = threads + .filter( + (thread) => + thread.session !== null && + thread.archivedAt === null && + thread.deletedAt === null && + !liveThreadIds.has(thread.id), + ) + .map((thread) => thread.id); + yield* settleThreadsTasks({ + threadIds: settleableThreadIds, + status: "interrupted", + createdAt: DateTime.formatIso(yield* DateTime.now), + }); for (const thread of orphanedThreads) { const session = thread.session; From 8292e39bbf551bcb828cdca66bc72cf836a6e1a6 Mon Sep 17 00:00:00 2001 From: amanthanvi Date: Thu, 3 Sep 2026 05:41:19 -0400 Subject: [PATCH 3/8] fix(server): isolate settlement failures per task and per thread Background-task settlement gave up too easily. A failed activity append for one task aborted the whole loop, so every later live task on that thread was left neither persisted nor recorded in the liveness registry, and the startup batch ran all threads under one recovery, so the first failing thread skipped every thread after it. Each append now recovers on its own: log a warning with the thread and task, skip that task's tombstone so the registry never claims a settlement no row backs, and carry on with the rest of the fleet. The batched startup path wraps each thread's writes in its own recovery while keeping the single batched read. Interrupt causes still propagate. Implemented by Claude Opus 5 via Claude Code. --- .../ThreadTaskSettlement.test.ts | 170 ++++++++++++++++++ .../src/orchestration/ThreadTaskSettlement.ts | 39 +++- 2 files changed, 201 insertions(+), 8 deletions(-) create mode 100644 apps/server/src/orchestration/ThreadTaskSettlement.test.ts diff --git a/apps/server/src/orchestration/ThreadTaskSettlement.test.ts b/apps/server/src/orchestration/ThreadTaskSettlement.test.ts new file mode 100644 index 000000000000..726bd1dba572 --- /dev/null +++ b/apps/server/src/orchestration/ThreadTaskSettlement.test.ts @@ -0,0 +1,170 @@ +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { type OrchestrationCommand, ThreadId } from "@t3tools/contracts"; +import { assert, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Stream from "effect/Stream"; + +import { OrchestrationCommandInvariantError } from "./Errors.ts"; +import * as ProjectionThreadActivities from "../persistence/Services/ProjectionThreadActivities.ts"; +import * as OrchestrationEngine from "./Services/OrchestrationEngine.ts"; +import * as ThreadBackgroundLiveness from "./ThreadBackgroundLiveness.ts"; +import { + settleThreadsTasks, + settleThreadTasks, + type TaskActivityRow, +} from "./ThreadTaskSettlement.ts"; + +const createdAt = "2026-01-01T00:00:00.000Z"; + +const liveRowsFor = (taskId: string): ReadonlyArray => [ + { kind: "task.started", payload: { taskId, agentKind: "agent", title: taskId } }, + { kind: "task.updated", payload: { taskId, status: "running", agentKind: "agent" } }, +]; + +const settledTaskIds = (dispatched: ReadonlyArray): ReadonlyArray => + dispatched.flatMap((command) => + command.type === "thread.activity.append" && + typeof command.activity.payload === "object" && + command.activity.payload !== null + ? [String((command.activity.payload as { taskId?: unknown }).taskId)] + : [], + ); + +/** + * Provides the settlement dependencies. `failFor` decides which appends blow + * up, so a test can prove the loop keeps going past a failure. + */ +const withSettlementServices = (input: { + readonly activitiesByThreadId: Readonly>>; + readonly dispatched: Array; + readonly liveness: ThreadBackgroundLiveness.ThreadBackgroundLivenessService["Service"]; + readonly failFor?: (command: OrchestrationCommand) => boolean; +}) => { + const repository = { + listTaskLifecycleByThreadId: ({ threadId }: { readonly threadId: ThreadId }) => + Effect.succeed(input.activitiesByThreadId[threadId] ?? []), + listTaskLifecycleByThreadIds: ({ + threadIds, + }: { + readonly threadIds: ReadonlyArray; + }) => + Effect.succeed( + threadIds.flatMap((threadId) => + (input.activitiesByThreadId[threadId] ?? []).map((row) => ({ ...row, threadId })), + ), + ), + } as unknown as ProjectionThreadActivities.ProjectionThreadActivityRepository["Service"]; + + return (effect: Effect.Effect) => + effect.pipe( + Effect.provideService( + ProjectionThreadActivities.ProjectionThreadActivityRepository, + repository, + ), + Effect.provideService(OrchestrationEngine.OrchestrationEngineService, { + readEvents: () => Stream.empty, + dispatch: (command) => { + if (input.failFor?.(command) === true) { + return Effect.fail( + new OrchestrationCommandInvariantError({ + commandType: command.type, + detail: "simulated settlement append failure", + }), + ); + } + input.dispatched.push(command); + return Effect.succeed({ sequence: input.dispatched.length }); + }, + streamDomainEvents: Stream.empty, + subscribeDomainEvents: Effect.succeed(Stream.empty), + latestSequence: Effect.succeed(0), + }), + Effect.provideService( + ThreadBackgroundLiveness.ThreadBackgroundLivenessService, + input.liveness, + ), + Effect.provide(NodeServices.layer), + ); +}; + +it.effect("settles the rest of the fleet when one task's append fails", () => { + const threadId = ThreadId.make("thread-partial-failure"); + const dispatched: Array = []; + const liveness = ThreadBackgroundLiveness.make(); + for (const taskId of ["child-one", "child-two"]) { + liveness.recordTaskLiveness({ + threadId, + taskId, + taskType: undefined, + status: "running", + kind: "updated", + }); + } + + return settleThreadTasks({ threadId, status: "interrupted", createdAt }).pipe( + withSettlementServices({ + activitiesByThreadId: { + [threadId]: [...liveRowsFor("child-one"), ...liveRowsFor("child-two")], + }, + dispatched, + liveness, + failFor: (command) => + command.type === "thread.activity.append" && + command.activity.id === `task-settled:${threadId}:child-one`, + }), + Effect.tap(() => + Effect.sync(() => { + // The failed task keeps no row, so it must keep no tombstone either; + // the second task settles regardless. + assert.deepStrictEqual(settledTaskIds(dispatched), ["child-two"]); + assert.equal(liveness.getThreadBackgroundLiveness(threadId), "working"); + liveness.recordTaskLiveness({ + threadId, + taskId: "child-one", + taskType: undefined, + status: "interrupted", + kind: "updated", + settledByHost: true, + }); + assert.equal(liveness.getThreadBackgroundLiveness(threadId), null); + }), + ), + ); +}); + +it.effect("a failing thread does not block the threads after it in a batch", () => { + const failing = ThreadId.make("thread-batch-failing"); + const healthy = ThreadId.make("thread-batch-healthy"); + const dispatched: Array = []; + const liveness = ThreadBackgroundLiveness.make(); + liveness.recordTaskLiveness({ + threadId: healthy, + taskId: "child-two", + taskType: undefined, + status: "running", + kind: "updated", + }); + + return settleThreadsTasks({ + threadIds: [failing, healthy], + status: "interrupted", + createdAt, + }).pipe( + withSettlementServices({ + activitiesByThreadId: { + [failing]: liveRowsFor("child-one"), + [healthy]: liveRowsFor("child-two"), + }, + dispatched, + liveness, + failFor: (command) => + command.type === "thread.activity.append" && command.threadId === failing, + }), + Effect.tap(() => + Effect.sync(() => { + assert.deepStrictEqual(settledTaskIds(dispatched), ["child-two"]); + assert.equal(liveness.getThreadBackgroundLiveness(healthy), null); + }), + ), + ); +}); diff --git a/apps/server/src/orchestration/ThreadTaskSettlement.ts b/apps/server/src/orchestration/ThreadTaskSettlement.ts index ef2125be3bff..b22a75a91c96 100644 --- a/apps/server/src/orchestration/ThreadTaskSettlement.ts +++ b/apps/server/src/orchestration/ThreadTaskSettlement.ts @@ -333,13 +333,34 @@ const writeSettledTasks = Effect.fn("writeSettledTasks")(function* (input: { }, turnId: null, }; - yield* orchestrationEngine.dispatch({ - type: "thread.activity.append", - commandId: CommandId.make(`task-settle:${yield* crypto.randomUUIDv4}`), - threadId: input.threadId, - activity, - createdAt: input.createdAt, - }); + const commandId = CommandId.make(`task-settle:${yield* crypto.randomUUIDv4}`); + // One task's append must not abandon the rest of the fleet: a partially + // settled thread is exactly the state this whole module exists to avoid. + const appended = yield* orchestrationEngine + .dispatch({ + type: "thread.activity.append", + commandId, + threadId: input.threadId, + activity, + createdAt: input.createdAt, + }) + .pipe( + Effect.as(true), + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.failCause(cause) + : Effect.logWarning("failed to settle background task", { + threadId: input.threadId, + taskId: task.taskId, + cause: Cause.pretty(cause), + }).pipe(Effect.as(false)), + ), + ); + if (!appended) { + // No row, no tombstone: the registry must not claim a task is settled + // when nothing was persisted to prove it. + continue; + } // Rows and registry settle together, so the sidebar pill and the composer // banner can never disagree about the same task. threadBackgroundLiveness.recordTaskLiveness({ @@ -434,12 +455,14 @@ export const settleThreadsTasks = Effect.fn("settleThreadsTasks")(function* (inp if (liveTasks.length === 0) { continue; } + // Per thread, so one bad thread cannot cost every thread after it its + // settlement for the rest of this process's life. yield* writeSettledTasks({ threadId, liveTasks, status: input.status, createdAt: input.createdAt, - }); + }).pipe(withSettlementRecovery([threadId])); } }); From 22006073ba7f5e1824a867a6d684eb72c76869d2 Mon Sep 17 00:00:00 2001 From: amanthanvi Date: Sat, 5 Sep 2026 21:17:06 -0400 Subject: [PATCH 4/8] test(server): satisfy the engine and directory mock shapes after merging main --- .../src/orchestration/ThreadTaskSettlement.test.ts | 2 ++ apps/server/src/project/AgentSessionImporter.test.ts | 9 +++++++++ apps/server/src/serverRuntimeStartup.reconcile.test.ts | 6 ++++-- 3 files changed, 15 insertions(+), 2 deletions(-) diff --git a/apps/server/src/orchestration/ThreadTaskSettlement.test.ts b/apps/server/src/orchestration/ThreadTaskSettlement.test.ts index 726bd1dba572..102db9e9b14d 100644 --- a/apps/server/src/orchestration/ThreadTaskSettlement.test.ts +++ b/apps/server/src/orchestration/ThreadTaskSettlement.test.ts @@ -63,6 +63,8 @@ const withSettlementServices = (input: { ), Effect.provideService(OrchestrationEngine.OrchestrationEngineService, { readEvents: () => Stream.empty, + readThreadEvents: () => Stream.empty, + getThreadReplayStats: () => Effect.die("unused thread replay stats"), dispatch: (command) => { if (input.failFor?.(command) === true) { return Effect.fail( diff --git a/apps/server/src/project/AgentSessionImporter.test.ts b/apps/server/src/project/AgentSessionImporter.test.ts index 38d6ad1d4331..4f62e16260f8 100644 --- a/apps/server/src/project/AgentSessionImporter.test.ts +++ b/apps/server/src/project/AgentSessionImporter.test.ts @@ -37,6 +37,8 @@ import { OrchestrationProjectionSnapshotQueryLive } from "../orchestration/Layer import { ProviderCommandReactorLive } from "../orchestration/Layers/ProviderCommandReactor.ts"; import { OrchestrationCommandInvariantError } from "../orchestration/Errors.ts"; import * as ThreadBackgroundLiveness from "../orchestration/ThreadBackgroundLiveness.ts"; +import * as ProjectionThreadActivities from "../persistence/Services/ProjectionThreadActivities.ts"; +import { ProviderRuntimeIngestionService } from "../orchestration/Services/ProviderRuntimeIngestion.ts"; import * as ThreadPlanProgress from "../orchestration/ThreadPlanProgress.ts"; import * as OrchestrationEngine from "../orchestration/Services/OrchestrationEngine.ts"; import * as ProjectionSnapshotQuery from "../orchestration/Services/ProjectionSnapshotQuery.ts"; @@ -931,6 +933,13 @@ it.layer(integrationLayer)("AgentSessionImporter integration", (it) => { Layer.provide(Layer.mock(VcsStatusBroadcaster)({})), Layer.provide(Layer.mock(TextGeneration)({})), Layer.provide(ServerSettingsService.layerTest()), + // Stop settles background tasks after draining runtime ingestion; + // neither path runs in this test, so inert stand-ins suffice. + Layer.provide(Layer.mock(ProviderRuntimeIngestionService)({})), + Layer.provide( + Layer.mock(ProjectionThreadActivities.ProjectionThreadActivityRepository)({}), + ), + Layer.provide(ThreadBackgroundLiveness.layer), ); yield* engine.dispatch({ diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index da6d87f5093b..b87abae6f42d 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -781,7 +781,8 @@ it.effect("settles background agent tasks left running by an orphaned session", upsert: () => Effect.void, getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), - listBindings: () => Effect.die("unused"), + listBindings: () => Effect.succeed([]), + recordImportedTranscript: () => Effect.die("unused"), }, dispatch: (command) => { dispatched.push(command); @@ -844,7 +845,8 @@ it("settles background agents on ready and stopped threads, and writes nothing w upsert: () => Effect.void, getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), - listBindings: () => Effect.die("unused"), + listBindings: () => Effect.succeed([]), + recordImportedTranscript: () => Effect.die("unused"), }, dispatch: (command) => { dispatched.push(command); From f4b6bdd2f02fb9acf906bcaca37ec038d29c8dfb Mon Sep 17 00:00:00 2001 From: amanthanvi Date: Mon, 7 Sep 2026 18:50:13 -0400 Subject: [PATCH 5/8] refactor(server): type the activity repository doubles instead of asserting The settlement tests reached the repository shape through `as unknown as`, which hid the fact that the doubles returned narrow rows. They now build complete activity rows and check the shape with `satisfies`, so a contract change fails the typecheck instead of the test run. Also tightens the composer doc sentence about background agents. --- .../ThreadTaskSettlement.test.ts | 40 ++++++++++++++----- .../serverRuntimeStartup.reconcile.test.ts | 39 +++++++++++++----- docs/user/composer.md | 8 ++-- 3 files changed, 64 insertions(+), 23 deletions(-) diff --git a/apps/server/src/orchestration/ThreadTaskSettlement.test.ts b/apps/server/src/orchestration/ThreadTaskSettlement.test.ts index 102db9e9b14d..335b38931670 100644 --- a/apps/server/src/orchestration/ThreadTaskSettlement.test.ts +++ b/apps/server/src/orchestration/ThreadTaskSettlement.test.ts @@ -1,5 +1,5 @@ import * as NodeServices from "@effect/platform-node/NodeServices"; -import { type OrchestrationCommand, ThreadId } from "@t3tools/contracts"; +import { EventId, type OrchestrationCommand, ThreadId } from "@t3tools/contracts"; import { assert, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; import * as Stream from "effect/Stream"; @@ -21,6 +21,25 @@ const liveRowsFor = (taskId: string): ReadonlyArray => [ { kind: "task.updated", payload: { taskId, status: "running", agentKind: "agent" } }, ]; +/** + * Fills the columns settlement never reads, so the repository double returns + * the same rows the real one does. + */ +const toActivityRows = ( + threadId: ThreadId, + rows: ReadonlyArray, +): ReadonlyArray => + rows.map((row, index) => ({ + activityId: EventId.make(`activity-${threadId}-${index}`), + threadId, + turnId: null, + tone: "info", + kind: row.kind, + summary: row.kind, + payload: row.payload, + createdAt, + })); + const settledTaskIds = (dispatched: ReadonlyArray): ReadonlyArray => dispatched.flatMap((command) => command.type === "thread.activity.append" && @@ -41,19 +60,20 @@ const withSettlementServices = (input: { readonly failFor?: (command: OrchestrationCommand) => boolean; }) => { const repository = { - listTaskLifecycleByThreadId: ({ threadId }: { readonly threadId: ThreadId }) => - Effect.succeed(input.activitiesByThreadId[threadId] ?? []), - listTaskLifecycleByThreadIds: ({ - threadIds, - }: { - readonly threadIds: ReadonlyArray; - }) => + upsert: () => Effect.die("unused"), + listByThreadId: () => Effect.die("unused"), + listUserInputLifecycleByThreadId: () => Effect.die("unused"), + getLatestTaskActivity: () => Effect.die("unused"), + deleteByThreadId: () => Effect.die("unused"), + listTaskLifecycleByThreadId: ({ threadId }) => + Effect.succeed(toActivityRows(threadId, input.activitiesByThreadId[threadId] ?? [])), + listTaskLifecycleByThreadIds: ({ threadIds }) => Effect.succeed( threadIds.flatMap((threadId) => - (input.activitiesByThreadId[threadId] ?? []).map((row) => ({ ...row, threadId })), + toActivityRows(threadId, input.activitiesByThreadId[threadId] ?? []), ), ), - } as unknown as ProjectionThreadActivities.ProjectionThreadActivityRepository["Service"]; + } satisfies ProjectionThreadActivities.ProjectionThreadActivityRepository["Service"]; return (effect: Effect.Effect) => effect.pipe( diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index 2443b97f855d..1268e84f652f 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -1,5 +1,6 @@ import * as NodeServices from "@effect/platform-node/NodeServices"; import { + EventId, type OrchestrationCommand, type OrchestrationSessionStatus, ProviderDriverKind, @@ -81,25 +82,45 @@ const queryWithThreads = (threads: ReadonlyArray>) getCommandReadModel: () => Effect.succeed({ threads } as never), }) as unknown as ProjectionSnapshotQuery.ProjectionSnapshotQuery["Service"]; +/** + * Fills the columns settlement never reads, so the repository double returns + * the same rows the real one does. + */ +const toActivityRows = ( + threadId: ThreadId, + rows: ReadonlyArray, +): ReadonlyArray => + rows.map((row, index) => ({ + activityId: EventId.make(`activity-${threadId}-${index}`), + threadId, + turnId: null, + tone: "info", + kind: row.kind, + summary: row.kind, + payload: row.payload, + createdAt: updatedAt, + })); + const activityRepositoryWith = ( activitiesByThreadId: Readonly>>, batchReads: Array> = [], ) => ({ - listTaskLifecycleByThreadId: ({ threadId }: { readonly threadId: ThreadId }) => - Effect.succeed(activitiesByThreadId[threadId] ?? []), - listTaskLifecycleByThreadIds: ({ - threadIds, - }: { - readonly threadIds: ReadonlyArray; - }) => + upsert: () => Effect.die("unused"), + listByThreadId: () => Effect.die("unused"), + listUserInputLifecycleByThreadId: () => Effect.die("unused"), + getLatestTaskActivity: () => Effect.die("unused"), + deleteByThreadId: () => Effect.die("unused"), + listTaskLifecycleByThreadId: ({ threadId }) => + Effect.succeed(toActivityRows(threadId, activitiesByThreadId[threadId] ?? [])), + listTaskLifecycleByThreadIds: ({ threadIds }) => Effect.sync(() => { batchReads.push(threadIds); return threadIds.flatMap((threadId) => - (activitiesByThreadId[threadId] ?? []).map((row) => ({ ...row, threadId })), + toActivityRows(threadId, activitiesByThreadId[threadId] ?? []), ); }), - }) as unknown as ProjectionThreadActivities.ProjectionThreadActivityRepository["Service"]; + }) satisfies ProjectionThreadActivities.ProjectionThreadActivityRepository["Service"]; const runReconciliation = (input: { readonly threads: ReadonlyArray>; diff --git a/docs/user/composer.md b/docs/user/composer.md index 2050d10defcc..4a4ef6738628 100644 --- a/docs/user/composer.md +++ b/docs/user/composer.md @@ -112,10 +112,10 @@ provider supports it. Web and desktop also offer compaction from the context met ## Background agents after a turn When background agents keep working after a turn ends, a banner above the composer counts them -and offers **Stop**. Pressing Stop, or the provider session ending, marks any agents still listed -as working as interrupted, so the banner clears and the Agents panel stops showing them as -working, even when the provider never reports those agents again. Restarting T3 Code does the -same for threads whose provider session did not survive. +and offers **Stop**. Pressing Stop, or the provider session ending, interrupts every agent still +shown as working, so the banner clears and the Agents panel stops counting them, even when the +provider never reports those agents again. Restarting T3 Code does the same for threads whose +provider session did not survive. ## Images and videos in messages From dc6fa2ce21da71e7e812a1ca573b5398ac0fd837 Mon Sep 17 00:00:00 2001 From: amanthanvi Date: Mon, 7 Sep 2026 19:03:40 -0400 Subject: [PATCH 6/8] test(server): run the last reconciliation case and drop a dead settlement guard The final case in the startup reconciliation suite used bare `it`, so the Effect it returned was never executed and none of its assertions ran. It now uses `it.effect`; the assertions were already correct. `applyStatus` carried a guard that only skipped writing "terminal" over "terminal", a provable no-op, under a comment claiming terminal is sticky. It is not: a later running row reopens the entry, matching the client fold. `listTaskLifecycleByThreadIds` orders by `thread_id` first, so its doc no longer claims threads come back interleaved. Adds repository coverage for the three behaviors that read path relies on: the lifecycle-kind filter, unsequenced rows sorting ahead of sequenced ones, and thread ids spanning more than one query chunk. --- .../src/orchestration/ThreadTaskSettlement.ts | 5 +- .../Layers/ProjectionThreadActivities.test.ts | 171 +++++++++++++++++- .../Layers/ProjectionThreadActivities.ts | 2 +- .../Services/ProjectionThreadActivities.ts | 4 +- .../serverRuntimeStartup.reconcile.test.ts | 129 ++++++------- 5 files changed, 240 insertions(+), 71 deletions(-) diff --git a/apps/server/src/orchestration/ThreadTaskSettlement.ts b/apps/server/src/orchestration/ThreadTaskSettlement.ts index b22a75a91c96..aeb90238ab10 100644 --- a/apps/server/src/orchestration/ThreadTaskSettlement.ts +++ b/apps/server/src/orchestration/ThreadTaskSettlement.ts @@ -165,11 +165,8 @@ function asFoldStatus(value: unknown): FoldStatus | undefined { return typeof value === "string" ? KNOWN_STATUSES.get(value) : undefined; } -/** Terminal is sticky: a duplicate or late terminal row never slides state. */ +/** Last status row wins, mirroring the client fold's reactivation rule. */ function applyStatus(entry: FoldEntry, next: FoldStatus): void { - if (entry.status === "terminal" && next === "terminal") { - return; - } entry.status = next; } diff --git a/apps/server/src/persistence/Layers/ProjectionThreadActivities.test.ts b/apps/server/src/persistence/Layers/ProjectionThreadActivities.test.ts index d92ed97ea6fa..1780a9dc073d 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadActivities.test.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadActivities.test.ts @@ -5,7 +5,10 @@ import * as Layer from "effect/Layer"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { ProjectionThreadActivityRepository } from "../Services/ProjectionThreadActivities.ts"; -import { ProjectionThreadActivityRepositoryLive } from "./ProjectionThreadActivities.ts"; +import { + ProjectionThreadActivityRepositoryLive, + THREAD_ID_BATCH_SIZE, +} from "./ProjectionThreadActivities.ts"; import { SqlitePersistenceMemory } from "./Sqlite.ts"; const layer = it.layer( @@ -97,4 +100,170 @@ layer("ProjectionThreadActivityRepository", (it) => { ); }), ); + + it.effect("reads only task lifecycle rows, scoped to the requested threads", () => + Effect.gen(function* () { + const repository = yield* ProjectionThreadActivityRepository; + const sql = yield* SqlClient.SqlClient; + const threadId = ThreadId.make("thread-task-lifecycle-kinds"); + const otherThreadId = ThreadId.make("thread-task-lifecycle-kinds-other"); + + // The excluded rows carry unparseable payloads: if the kind filter ever + // stops running in SQLite, decoding them fails the read outright. + yield* sql` + INSERT INTO projection_thread_activities ( + activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at + ) + VALUES + ( + 'lifecycle-started', ${threadId}, NULL, 'info', 'task.started', + 'started', '{"taskId":"task-1","agentKind":"agent"}', 1, + '2026-03-02T00:00:00.000Z' + ), + ( + 'lifecycle-progress', ${threadId}, NULL, 'info', 'task.progress', + 'progress', '{"taskId":"task-1"}', 2, '2026-03-02T00:00:01.000Z' + ), + ( + 'lifecycle-updated', ${threadId}, NULL, 'info', 'task.updated', + 'updated', '{"taskId":"task-1","status":"running"}', 3, + '2026-03-02T00:00:02.000Z' + ), + ( + 'lifecycle-completed', ${threadId}, NULL, 'info', 'task.completed', + 'completed', '{"taskId":"task-1"}', 4, '2026-03-02T00:00:03.000Z' + ), + ( + 'lifecycle-tool', ${threadId}, NULL, 'tool', 'tool.completed', + 'tool output', 'not-json', 5, '2026-03-02T00:00:04.000Z' + ), + ( + 'lifecycle-user-input', ${threadId}, NULL, 'info', 'user-input.requested', + 'input requested', 'not-json', 6, '2026-03-02T00:00:05.000Z' + ), + ( + 'lifecycle-other-thread', ${otherThreadId}, NULL, 'info', 'task.started', + 'started', '{"taskId":"task-2","agentKind":"agent"}', 1, + '2026-03-02T00:00:06.000Z' + ) + `; + + const threadRows = yield* repository.listTaskLifecycleByThreadId({ threadId }); + assert.deepEqual( + threadRows.map((entry) => entry.activityId), + ["lifecycle-started", "lifecycle-progress", "lifecycle-updated", "lifecycle-completed"], + ); + + const batchedRows = yield* repository.listTaskLifecycleByThreadIds({ + threadIds: [threadId, otherThreadId], + }); + assert.deepEqual( + batchedRows.map((entry) => entry.activityId), + [ + "lifecycle-started", + "lifecycle-progress", + "lifecycle-updated", + "lifecycle-completed", + "lifecycle-other-thread", + ], + ); + }), + ); + + it.effect("orders unsequenced task rows ahead of sequenced ones", () => + Effect.gen(function* () { + const repository = yield* ProjectionThreadActivityRepository; + const sql = yield* SqlClient.SqlClient; + const threadId = ThreadId.make("thread-task-lifecycle-order"); + + yield* sql` + INSERT INTO projection_thread_activities ( + activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at + ) + VALUES + ( + 'order-sequenced-second', ${threadId}, NULL, 'info', 'task.progress', + 'progress', '{"taskId":"task-1"}', 2, '2026-03-03T00:00:00.000Z' + ), + ( + 'order-sequenced-first', ${threadId}, NULL, 'info', 'task.started', + 'started', '{"taskId":"task-1"}', 1, '2026-03-03T00:00:01.000Z' + ), + ( + 'order-unsequenced-newer', ${threadId}, NULL, 'info', 'task.updated', + 'updated', '{"taskId":"task-1"}', NULL, '2026-03-03T00:00:03.000Z' + ), + ( + 'order-unsequenced-older', ${threadId}, NULL, 'info', 'task.updated', + 'updated', '{"taskId":"task-1"}', NULL, '2026-03-03T00:00:02.000Z' + ) + `; + + // Unsequenced rows first, ordered by created_at; sequenced rows after, + // ordered by sequence even when created_at disagrees. + const expected = [ + "order-unsequenced-older", + "order-unsequenced-newer", + "order-sequenced-first", + "order-sequenced-second", + ]; + assert.deepEqual( + (yield* repository.listTaskLifecycleByThreadId({ threadId })).map( + (entry) => entry.activityId, + ), + expected, + ); + assert.deepEqual( + (yield* repository.listTaskLifecycleByThreadIds({ threadIds: [threadId] })).map( + (entry) => entry.activityId, + ), + expected, + ); + }), + ); + + it.effect("reads every thread when the id list spans more than one query chunk", () => + Effect.gen(function* () { + const repository = yield* ProjectionThreadActivityRepository; + const sql = yield* SqlClient.SqlClient; + + // One thread past the chunk boundary, so the batched read has to stitch + // two statements together. Ids are zero-padded, keeping thread_id order + // identical to the requested order. + const threadIds = Array.from({ length: THREAD_ID_BATCH_SIZE + 1 }, (_, index) => + ThreadId.make(`thread-task-chunk-${String(index).padStart(4, "0")}`), + ); + // The threads on either side of the boundary carry two rows each, so a + // regression that drops or reorders a straddling thread shows up here. + const rows = threadIds.flatMap((threadId, index) => + Array.from( + { length: index === THREAD_ID_BATCH_SIZE - 1 || index === THREAD_ID_BATCH_SIZE ? 2 : 1 }, + (_, rowIndex) => ({ + activity_id: `activity-${threadId}-${rowIndex}`, + thread_id: threadId, + turn_id: null, + tone: "info", + kind: rowIndex === 0 ? "task.started" : "task.completed", + summary: "chunked", + payload_json: '{"taskId":"task-1","agentKind":"agent"}', + sequence: rowIndex + 1, + created_at: "2026-03-04T00:00:00.000Z", + }), + ), + ); + // Small insert batches keep this fixture under SQLite's parameter limit. + for (let index = 0; index < rows.length; index += 100) { + yield* sql` + INSERT INTO projection_thread_activities ${sql.insert(rows.slice(index, index + 100))} + `; + } + + assert.deepEqual( + (yield* repository.listTaskLifecycleByThreadIds({ threadIds })).map( + (entry) => entry.activityId, + ), + rows.map((row) => row.activity_id), + ); + }), + ); }); diff --git a/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts b/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts index 306774df400b..8710e1de09c1 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts @@ -43,7 +43,7 @@ function toProjectionThreadActivity( } /** SQLite's host-parameter ceiling is 999; stay well inside it. */ -const THREAD_ID_BATCH_SIZE = 500; +export const THREAD_ID_BATCH_SIZE = 500; const chunk = (items: ReadonlyArray, size: number): ReadonlyArray> => { const chunks: A[][] = []; diff --git a/apps/server/src/persistence/Services/ProjectionThreadActivities.ts b/apps/server/src/persistence/Services/ProjectionThreadActivities.ts index cbc6dc24def4..e724cf58cbd1 100644 --- a/apps/server/src/persistence/Services/ProjectionThreadActivities.ts +++ b/apps/server/src/persistence/Services/ProjectionThreadActivities.ts @@ -99,8 +99,8 @@ export interface ProjectionThreadActivityRepositoryShape { * * Startup reconciliation settles background work across every thread this * process does not own, so it reads them together instead of issuing one - * query per thread. Rows come back grouped-ready: ordered within each - * thread, threads interleaved. + * query per thread. Rows come back grouped-ready: each thread's rows are + * contiguous and ordered within the thread. */ readonly listTaskLifecycleByThreadIds: ( input: ListProjectionThreadActivitiesByThreadIdsInput, diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index 1268e84f652f..e64015f6f8c9 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -850,70 +850,73 @@ it.effect("settles background agent tasks left running by an orphaned session", ); }); -it("settles background agents on ready and stopped threads, and writes nothing when idle", () => { - // The headline case: the turn finished, so the session is `ready` with no - // active turn, but its children kept working and did not survive the - // restart. Archived and deleted threads must not be written to. - const taskId = "collab-child-1"; - const linkage = { agentKind: "agent", title: "math_one", timelineBypass: true } as const; - const rows = [ - { kind: "task.started", payload: { taskId, ...linkage } }, - { kind: "task.updated", payload: { taskId, status: "running", ...linkage } }, - ]; - const ready = makeThread("thread-settle-ready", "ready"); - const archived = makeThread("thread-settle-archived", "ready", null, updatedAt); - const deleted = makeThread("thread-settle-deleted", "ready", null, null, updatedAt); - // A stopped session still counts: the process can die between marking the - // session stopped and settling its children, and no later boot would ever - // pick them up if stopped threads were skipped. - const stopped = makeThread("thread-settle-stopped", "stopped"); - // Nothing to settle here, so it must produce no writes at all. - const quiet = makeThread("thread-settle-quiet", "ready"); - const dispatched: OrchestrationCommand[] = []; - const batchReads: Array> = []; +it.effect( + "settles background agents on ready and stopped threads, and writes nothing when idle", + () => { + // The headline case: the turn finished, so the session is `ready` with no + // active turn, but its children kept working and did not survive the + // restart. Archived and deleted threads must not be written to. + const taskId = "collab-child-1"; + const linkage = { agentKind: "agent", title: "math_one", timelineBypass: true } as const; + const rows = [ + { kind: "task.started", payload: { taskId, ...linkage } }, + { kind: "task.updated", payload: { taskId, status: "running", ...linkage } }, + ]; + const ready = makeThread("thread-settle-ready", "ready"); + const archived = makeThread("thread-settle-archived", "ready", null, updatedAt); + const deleted = makeThread("thread-settle-deleted", "ready", null, null, updatedAt); + // A stopped session still counts: the process can die between marking the + // session stopped and settling its children, and no later boot would ever + // pick them up if stopped threads were skipped. + const stopped = makeThread("thread-settle-stopped", "stopped"); + // Nothing to settle here, so it must produce no writes at all. + const quiet = makeThread("thread-settle-quiet", "ready"); + const dispatched: OrchestrationCommand[] = []; + const batchReads: Array> = []; - return runReconciliation({ - threads: [ready, archived, deleted, stopped, quiet], - batchReads, - activitiesByThreadId: { - [ready.id]: rows, - [archived.id]: rows, - [deleted.id]: rows, - [stopped.id]: rows, - }, - directory: { - getBinding: () => Effect.succeed(Option.none()), - upsert: () => Effect.void, - getProvider: () => Effect.die("unused"), - listThreadIds: () => Effect.die("unused"), - listBindings: () => Effect.succeed([]), - recordImportedTranscript: () => Effect.die("unused"), - }, - dispatch: (command) => { - dispatched.push(command); - return Effect.succeed({ sequence: dispatched.length }); - }, - }).pipe( - Effect.tap(() => - Effect.sync(() => { - assert.deepStrictEqual( - dispatched - .filter((command) => command.type === "thread.activity.append") - .map((command) => command.threadId), - [ready.id, stopped.id], - ); - // One batched read for every candidate thread, not a query per thread, - // and archived/deleted threads are never even read. - assert.deepStrictEqual(batchReads, [[ready.id, stopped.id, quiet.id]]); - // A ready session is not an orphaned turn, so nothing else changes. - assert.deepStrictEqual( - dispatched.filter((command) => command.type === "thread.session.set"), - [], - ); - }), - ), - ); -}); + return runReconciliation({ + threads: [ready, archived, deleted, stopped, quiet], + batchReads, + activitiesByThreadId: { + [ready.id]: rows, + [archived.id]: rows, + [deleted.id]: rows, + [stopped.id]: rows, + }, + directory: { + getBinding: () => Effect.succeed(Option.none()), + upsert: () => Effect.void, + getProvider: () => Effect.die("unused"), + listThreadIds: () => Effect.die("unused"), + listBindings: () => Effect.succeed([]), + recordImportedTranscript: () => Effect.die("unused"), + }, + dispatch: (command) => { + dispatched.push(command); + return Effect.succeed({ sequence: dispatched.length }); + }, + }).pipe( + Effect.tap(() => + Effect.sync(() => { + assert.deepStrictEqual( + dispatched + .filter((command) => command.type === "thread.activity.append") + .map((command) => command.threadId), + [ready.id, stopped.id], + ); + // One batched read for every candidate thread, not a query per thread, + // and archived/deleted threads are never even read. + assert.deepStrictEqual(batchReads, [[ready.id, stopped.id, quiet.id]]); + // A ready session is not an orphaned turn, so nothing else changes. + assert.deepStrictEqual( + dispatched.filter((command) => command.type === "thread.session.set"), + [], + ); + }), + ), + ); + }, +); for (const scenario of [ "disabled", From bd854bba0d59f1d28ca4137b02e5177bcc1c3985 Mon Sep 17 00:00:00 2001 From: amanthanvi Date: Mon, 7 Sep 2026 19:38:32 -0400 Subject: [PATCH 7/8] fix(server): date late provider rows against the host settlement The liveness registry's host-settled tombstone only blocked status-free rows, so an explicit task.updated(running) stamped before a Stop but delivered after it re-armed the thread. The sidebar pill and composer banner then showed work in flight while the persisted fold counted zero, and a second Stop could not clear it because settlement found nothing live to settle. The registry now stores the settlement timestamp instead of a bare key and drops any row the provider stamped at or before it, which is the same created_at rule the persisted activity fold already applies. --- .../Layers/ProviderRuntimeIngestion.ts | 3 ++ .../ThreadBackgroundLiveness.test.ts | 52 ++++++++++++++++++ .../orchestration/ThreadBackgroundLiveness.ts | 54 ++++++++++++++----- .../src/orchestration/ThreadTaskSettlement.ts | 3 ++ 4 files changed, 98 insertions(+), 14 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 65342ea5bbe8..00d590dc9bcc 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -2037,6 +2037,9 @@ const make = Effect.gen(function* () { taskType: payload.taskType, status: payload.status, agentId: payload.agentId, + // The adapter's stamp, not arrival time: this worker can be far + // behind, and a host settlement dates late rows against it. + occurredAt: event.createdAt, kind: event.type === "task.started" ? "started" diff --git a/apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts b/apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts index cf6b83174865..4f840eeb614a 100644 --- a/apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts +++ b/apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts @@ -299,4 +299,56 @@ describe("ThreadBackgroundLiveness", () => { }); expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("working"); }); + + it("a running row the provider stamped before a host settlement stays settled", () => { + // Stop bounds its wait on the ingestion drain, so the worker can still be + // holding a task.updated(running) the provider emitted BEFORE the Stop and + // deliver it after settlement. Persisted rows order it behind the + // settlement row by createdAt; the registry has to reach the same answer. + const liveness = ThreadBackgroundLiveness.make(); + const threadId = "t-late-running"; + liveness.recordTaskLiveness({ + threadId, + taskId: "child", + taskType: undefined, + status: "running", + kind: "updated", + occurredAt: "2026-01-01T00:00:00.000Z", + }); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("working"); + + liveness.recordTaskLiveness({ + threadId, + taskId: "child", + taskType: undefined, + status: "interrupted", + kind: "updated", + occurredAt: "2026-01-01T00:00:10.000Z", + settledByHost: true, + }); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull(); + + // The drained-too-late row. Stale, so it must not re-arm the thread. + liveness.recordTaskLiveness({ + threadId, + taskId: "child", + taskType: undefined, + status: "running", + kind: "updated", + occurredAt: "2026-01-01T00:00:05.000Z", + }); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull(); + + // A row stamped after the settlement is the provider proving it still owns + // the task, which the persisted fold also honours. + liveness.recordTaskLiveness({ + threadId, + taskId: "child", + taskType: undefined, + status: "running", + kind: "updated", + occurredAt: "2026-01-01T00:00:11.000Z", + }); + expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("working"); + }); }); diff --git a/apps/server/src/orchestration/ThreadBackgroundLiveness.ts b/apps/server/src/orchestration/ThreadBackgroundLiveness.ts index 6ba552adc6b4..3925f1db3600 100644 --- a/apps/server/src/orchestration/ThreadBackgroundLiveness.ts +++ b/apps/server/src/orchestration/ThreadBackgroundLiveness.ts @@ -66,6 +66,14 @@ export class ThreadBackgroundLivenessService extends Context.Service< readonly status: string | undefined; readonly kind: "started" | "progress" | "updated" | "completed"; readonly agentId?: string | undefined; + /** + * When the provider stamped this transition (the settlement time on a + * host settlement). The registry is fed in arrival order, but persisted + * rows are read in `createdAt` order, so without this a row the ingestion + * worker delivered late would re-arm a task the rows already show as + * settled. Absent means "unordered": the timestamp comparison is skipped. + */ + readonly occurredAt?: string | undefined; /** * Set by host settlement (Stop, session death, startup reconciliation). * Tombstones the task so a status-free start row already in flight from @@ -89,7 +97,8 @@ const taskKey = (threadId: string, taskId: string) => `${threadId}:${taskId}`; export function make(): ThreadBackgroundLivenessService["Service"] { const stateByThreadId = new Map(); - const hostSettledTaskKeys = new Set(); + /** Task key to the settlement's timestamp, or null when it carried none. */ + const hostSettledAtByTaskKey = new Map(); const stateFor = (threadId: string): ThreadLivenessState => { const existing = stateByThreadId.get(threadId); @@ -147,11 +156,14 @@ export function make(): ThreadBackgroundLivenessService["Service"] { // terminal event is ordinary lifecycle: a later start row for it is a // real resumption and must still arm the thread. if (input.settledByHost === true) { - hostSettledTaskKeys.add(taskKey(input.threadId, input.taskId)); - if (hostSettledTaskKeys.size > HOST_SETTLED_TASK_MEMORY_LIMIT) { - const oldest = hostSettledTaskKeys.values().next().value; + hostSettledAtByTaskKey.set( + taskKey(input.threadId, input.taskId), + input.occurredAt ?? null, + ); + if (hostSettledAtByTaskKey.size > HOST_SETTLED_TASK_MEMORY_LIMIT) { + const oldest = hostSettledAtByTaskKey.keys().next().value; if (oldest !== undefined) { - hostSettledTaskKeys.delete(oldest); + hostSettledAtByTaskKey.delete(oldest); } } } @@ -170,18 +182,32 @@ export function make(): ThreadBackgroundLivenessService["Service"] { } } - // A status-free row for a host-settled task is a late delivery, not a - // new run. Only an explicit non-terminal status below proves the task - // really came back, and it clears the tombstone. - if ( - input.status === undefined && - hostSettledTaskKeys.has(taskKey(input.threadId, input.taskId)) - ) { - return; + // A row for a host-settled task has to prove the task really came back. + // A status-free one never does. An explicit non-terminal status does, + // but only if the provider stamped it AFTER the settlement: the drain + // that holds queued provider events ahead of the settlement row is + // bounded, so an older row can still arrive here once the ingestion + // worker gets to it. The persisted fold orders that row behind the + // settlement row by createdAt, and the registry has to agree or the + // sidebar pill and the composer banner contradict the rows (#9391). + const settledKey = taskKey(input.threadId, input.taskId); + if (hostSettledAtByTaskKey.has(settledKey)) { + const settledAt = hostSettledAtByTaskKey.get(settledKey); + if (input.status === undefined) { + return; + } + if ( + settledAt !== undefined && + settledAt !== null && + input.occurredAt !== undefined && + input.occurredAt <= settledAt + ) { + return; + } } drop(input.threadId, input.taskId); - hostSettledTaskKeys.delete(taskKey(input.threadId, input.taskId)); + hostSettledAtByTaskKey.delete(settledKey); const state = stateFor(input.threadId); const bucket = taskType !== undefined && MONITOR_TASK_TYPES.has(taskType) ? state.monitors : state.agents; diff --git a/apps/server/src/orchestration/ThreadTaskSettlement.ts b/apps/server/src/orchestration/ThreadTaskSettlement.ts index aeb90238ab10..d51b9eb379ce 100644 --- a/apps/server/src/orchestration/ThreadTaskSettlement.ts +++ b/apps/server/src/orchestration/ThreadTaskSettlement.ts @@ -366,6 +366,9 @@ const writeSettledTasks = Effect.fn("writeSettledTasks")(function* (input: { taskType: undefined, status: input.status, kind: "updated", + // Same stamp the row carries, so the registry can date a late provider + // row against the settlement exactly as the persisted fold does. + occurredAt: input.createdAt, settledByHost: true, }); } From 72d20380a52f3c7e466bd9b4ad747212c5c86c54 Mon Sep 17 00:00:00 2001 From: amanthanvi Date: Wed, 9 Sep 2026 20:06:39 -0400 Subject: [PATCH 8/8] chore(server): plainer comments and composer wording for settlement Rewrites the comments this branch adds so they say what the code does without em dashes, label colons, or figures of speech. Two of them were also wrong: the layer comment claimed the provideMerge order was safe to invert (it is not, the later entry provides to the earlier), and the projection test comment described a straddling thread the chunker never produces. --- .../Layers/ProviderCommandReactor.ts | 7 +++-- .../Layers/ProviderRuntimeIngestion.test.ts | 2 +- .../Layers/ProviderRuntimeIngestion.ts | 8 ++--- .../src/orchestration/ThreadTaskSettlement.ts | 29 ++++++++++--------- .../Layers/ProjectionThreadActivities.test.ts | 3 +- apps/server/src/server.ts | 4 +-- apps/server/src/serverRuntimeStartup.ts | 6 ++-- docs/user/composer.md | 6 ++-- .../src/state/subagentRuntime.test.ts | 2 +- 9 files changed, 35 insertions(+), 32 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 40e996e2fc62..2ee372536988 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -1546,10 +1546,11 @@ const make = Effect.gen(function* () { .pipe(Effect.catchCause(recoverInterruptFailure)); // Settlement reads persisted rows, so every provider event that was - // already queued when Stop arrived has to land first — otherwise a + // already queued when Stop arrived has to land first. Otherwise a // task.updated(running) from before the interrupt is written after the // settlement row and re-arms both the registry and the client fold. - // Bounded: a hot event stream must not hold Stop hostage. + // The wait is bounded so a busy event stream cannot delay Stop past the + // timeout. const drained = yield* providerRuntimeIngestion.drain.pipe( Effect.timeoutOption(INTERRUPT_INGESTION_DRAIN_TIMEOUT), ); @@ -1560,7 +1561,7 @@ const make = Effect.gen(function* () { ); } - // Stop is a host promise, not a provider request: children the provider + // The host guarantees Stop; the provider does not. Children the provider // has already forgotten (compaction, a lost thread tree) never emit a // terminal event of their own, so settle the persisted rows here. Covers // both a successful interrupt and the stopSession fallback above; tasks diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 1d2794e3bdb9..b8dff394fdff 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -4547,7 +4547,7 @@ describe("selectLiveAgentTasks", () => { it("keeps linkage from earlier rows when a later row carries only a status", () => { // Terminal and status-patch rows commonly carry nothing but taskId and // status. If the settled row copied only that newest payload it would - // land with no agentKind — and once the start row falls out of the + // land with no agentKind. Once the start row falls out of the // client's activity window, the client would read the settled row as // background work and drop the agent from the panel entirely. const live = selectLiveAgentTasks([ diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 00d590dc9bcc..5cc6acbe7eb0 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -2054,11 +2054,11 @@ const make = Effect.gen(function* () { case "session.exited": // Rows first, then the registry. Background work dies with its // provider session, so any task still listed as running gets a - // persisted terminal row before the in-memory mirror is wiped — - // otherwise a restart rehydrates "running" rows with an empty + // persisted terminal row before the in-memory mirror is wiped. + // Otherwise a restart rehydrates "running" rows with an empty // registry and the agent reads as working forever. - // No drain needed here: this runs inside the ingestion worker, so - // every earlier provider event for this thread is already persisted. + // This runs inside the ingestion worker, so every earlier provider + // event for this thread is already persisted and no drain is needed. yield* settleThreadTasks({ threadId: thread.id, status: "interrupted", diff --git a/apps/server/src/orchestration/ThreadTaskSettlement.ts b/apps/server/src/orchestration/ThreadTaskSettlement.ts index d51b9eb379ce..4dbcc3f09468 100644 --- a/apps/server/src/orchestration/ThreadTaskSettlement.ts +++ b/apps/server/src/orchestration/ThreadTaskSettlement.ts @@ -2,18 +2,19 @@ * Settles a thread's still-running background agent tasks from the persisted * activity rows, without the provider's cooperation. * - * Native multi-agent children (Codex collab, workflow members) only leave the - * live set when the provider keeps reporting them. Compaction, a provider - * restart, a host Stop for a child the provider has already forgotten, or a - * T3 restart all lose that reporting, and the last persisted row stays - * "running" forever — the Agents panel and the composer's "N agents working" + * Native multi-agent children (Codex collab, workflow members) leave the live + * set only when the provider reports a terminal event for them. Compaction, a + * provider restart, a host Stop for a child the provider has already forgotten, + * or a T3 restart all lose that reporting, and the last persisted row stays + * "running" forever. The Agents panel and the composer's "N agents working" * banner never clear. * - * Persisted rows are the authority here: at the three moments where the server - * knows background work cannot continue (session death, host Stop, startup - * reconciliation) it folds the thread's task rows, synthesizes one terminal - * `task.updated` per still-live agent task, and feeds the same transition to - * the in-memory liveness registry so rows and registry settle together. + * Persisted rows are the authority here. The server settles at the three + * moments when it knows background work cannot continue: session death, host + * Stop, and startup reconciliation. At each one it folds the thread's task + * rows, synthesizes one terminal `task.updated` per still-live agent task, and + * feeds the same transition to the in-memory liveness registry so rows and + * registry settle together. * * @module ThreadTaskSettlement */ @@ -259,7 +260,7 @@ export function selectLiveAgentTasks( // Consistency pass, mirroring the client fold: once a workflow coordinator // has settled, members without a terminal row of their own cannot still be - // in flight — the run is over. The client shows them with the coordinator's + // in flight. The run is over. The client shows them with the coordinator's // outcome, so settling them here would overwrite a completed run with // "interrupted". for (const coordinator of entries.values()) { @@ -331,8 +332,8 @@ const writeSettledTasks = Effect.fn("writeSettledTasks")(function* (input: { turnId: null, }; const commandId = CommandId.make(`task-settle:${yield* crypto.randomUUIDv4}`); - // One task's append must not abandon the rest of the fleet: a partially - // settled thread is exactly the state this whole module exists to avoid. + // A failed append for one task must not skip the tasks after it. Their + // rows would stay "running" and the Agents panel would never clear. const appended = yield* orchestrationEngine .dispatch({ type: "thread.activity.append", @@ -393,7 +394,7 @@ const withSettlementRecovery = * Marks every still-live agent task on the thread as settled: one persisted * terminal `task.updated` per task, plus the matching liveness transition. * - * Never fails the caller — Stop, session teardown, and startup all continue + * Never fails the caller. Stop, session teardown, and startup all continue * when settlement cannot read or write. */ export const settleThreadTasks = Effect.fn("settleThreadTasks")(function* (input: { diff --git a/apps/server/src/persistence/Layers/ProjectionThreadActivities.test.ts b/apps/server/src/persistence/Layers/ProjectionThreadActivities.test.ts index 1780a9dc073d..dd7d774b00d1 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadActivities.test.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadActivities.test.ts @@ -234,7 +234,8 @@ layer("ProjectionThreadActivityRepository", (it) => { ThreadId.make(`thread-task-chunk-${String(index).padStart(4, "0")}`), ); // The threads on either side of the boundary carry two rows each, so a - // regression that drops or reorders a straddling thread shows up here. + // regression that drops rows at the boundary or reorders them inside a + // thread shows up here. const rows = threadIds.flatMap((threadId, index) => Array.from( { length: index === THREAD_ID_BATCH_SIZE - 1 || index === THREAD_ID_BATCH_SIZE ? 2 : 1 }, diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 2bcffd2c5270..8a42062dbdc7 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -276,8 +276,8 @@ const PlatformServicesLive = Layer.unwrap( const ReactorLayerLive = Layer.empty.pipe( Layer.provideMerge(OrchestrationReactorLive), // The command reactor drains runtime ingestion before settling background - // tasks on Stop, so ingestion must be provided to it (nothing in ingestion - // depends on the command reactor, so the order is safe to invert). + // tasks on Stop, so ingestion must be provided to it. Nothing in ingestion + // depends on the command reactor, so there is no cycle. Layer.provideMerge(ProviderCommandReactorLive), Layer.provideMerge(ProviderRuntimeIngestionLive), Layer.provideMerge(CheckpointReactorLive), diff --git a/apps/server/src/serverRuntimeStartup.ts b/apps/server/src/serverRuntimeStartup.ts index 4c9d24b2fd87..ea44a27ef061 100644 --- a/apps/server/src/serverRuntimeStartup.ts +++ b/apps/server/src/serverRuntimeStartup.ts @@ -534,13 +534,13 @@ export const reconcileProviderSessions = Effect.gen(function* () { // Background work outlives the turn, so the headline case is wider than an // orphaned turn: a thread whose turn ended is `ready` with no active turn - // while its children keep running. A stopped session qualifies too — the + // while its children keep running. A stopped session qualifies too. The // process can die between marking the session stopped and settling, and // that thread would otherwise never be settled on any later boot. Any // thread this process does not own has lost its background work, and a // restart leaves the in-memory liveness registry empty while the persisted - // rows still read "running". Archived and deleted threads are skipped — - // settling writes rows. One batched read, not a query per thread. + // rows still read "running". Archived and deleted threads are skipped + // because settling writes rows. One batched read, not a query per thread. const settleableThreadIds = threads .filter( (thread) => diff --git a/docs/user/composer.md b/docs/user/composer.md index 4a4ef6738628..2eb377bbb159 100644 --- a/docs/user/composer.md +++ b/docs/user/composer.md @@ -113,9 +113,9 @@ provider supports it. Web and desktop also offer compaction from the context met When background agents keep working after a turn ends, a banner above the composer counts them and offers **Stop**. Pressing Stop, or the provider session ending, interrupts every agent still -shown as working, so the banner clears and the Agents panel stops counting them, even when the -provider never reports those agents again. Restarting T3 Code does the same for threads whose -provider session did not survive. +shown as working. The banner then clears and the Agents panel stops counting those agents, even +when the provider never reports them again. Restarting T3 Code does the same for threads whose +provider session ended. ## Images and videos in messages diff --git a/packages/client-runtime/src/state/subagentRuntime.test.ts b/packages/client-runtime/src/state/subagentRuntime.test.ts index e9f7b271600d..6f54fb33cdf1 100644 --- a/packages/client-runtime/src/state/subagentRuntime.test.ts +++ b/packages/client-runtime/src/state/subagentRuntime.test.ts @@ -896,7 +896,7 @@ describe("host-settled background agents", () => { it("reconciles the live count as soon as a terminal row lands", () => { // Replay of a real Codex collab child: turnStarted/turnCompleted cycles, // the host's settlement row, and the provider's own late idle row after - // it. The last word is terminal, so nothing is left working. + // it. The last row is terminal, so nothing is left working. const linkage = { taskId: "collab-child-1", title: "math_one",