diff --git a/apps/server/integration/OrchestrationEngineHarness.integration.ts b/apps/server/integration/OrchestrationEngineHarness.integration.ts index 4d8a384ee997..baf57a62cc68 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"; @@ -314,6 +315,7 @@ export const makeOrchestrationIntegrationHarness = ( orchestrationLayer.pipe(Layer.provide(projectionSnapshotQueryLayer)), ProjectionCheckpointRepositoryLive, ProjectionPendingApprovalRepositoryLive, + ProjectionThreadActivityRepositoryLive, checkpointStoreLayer, providerLayer, RuntimeReceiptBusTest, @@ -343,6 +345,8 @@ export const makeOrchestrationIntegrationHarness = ( tryHandlePromptCommand: () => Effect.succeed(false), }), ), + // 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 c927e72fe737..f348d1e3cae1 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -46,6 +46,7 @@ import { } 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, @@ -66,6 +67,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"; @@ -119,6 +122,7 @@ describe("ProviderCommandReactor", () => { | OrchestrationEngineService | ProviderCommandReactor | ProjectionSnapshotQuery + | ThreadBackgroundLiveness.ThreadBackgroundLivenessService | SqlClient.SqlClient, unknown > | null = null; @@ -179,6 +183,8 @@ describe("ProviderCommandReactor", () => { readonly afterTurnStartDispatch?: () => Effect.Effect; readonly compactThreadEffect?: () => 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, @@ -460,6 +466,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.provide(Layer.mock(ProviderAuthService, { tryHandlePromptCommand })), Layer.provideMerge(makeProviderRegistryLayer(providerSnapshots as never)), @@ -497,6 +513,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( @@ -593,6 +612,7 @@ describe("ProviderCommandReactor", () => { engine, snapshotQuery, readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()), + backgroundLiveness, readPendingTurnStarts: () => runtime!.runPromise( Effect.gen(function* () { @@ -3460,6 +3480,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 c5d120106a19..26c9338191fe 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -48,6 +48,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 { @@ -57,6 +59,9 @@ import { import { resolveProjectSettings } from "@t3tools/shared/projectSettings"; 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 isProviderAdapterValidationError = Schema.is(ProviderAdapterValidationError); const isProviderWorkspaceMissingError = Schema.is(ProviderWorkspaceMissingError); @@ -324,6 +329,7 @@ const make = Effect.gen(function* () { const projectionSnapshotQuery = yield* ProjectionSnapshotQuery; const providerAuthService = yield* ProviderAuthService; const providerService = yield* ProviderService; + const providerRuntimeIngestion = yield* ProviderRuntimeIngestionService; const providerRegistry = yield* ProviderRegistry; const gitWorkflow = yield* GitWorkflowService; const fileSystem = yield* FileSystem.FileSystem; @@ -1669,6 +1675,34 @@ 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. + // 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), + ); + if (Option.isNone(drained)) { + yield* Effect.logWarning( + "provider runtime ingestion did not drain before background task settlement", + { threadId: event.payload.threadId }, + ); + } + + // 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 + // 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 1094ab48b7ac..b8dff394fdff 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -41,6 +41,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, @@ -56,6 +57,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"; @@ -234,7 +236,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; @@ -304,6 +309,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)), @@ -318,6 +326,9 @@ describe("ProviderRuntimeIngestion", () => { const engine = await testRuntime.runPromise(Effect.service(OrchestrationEngineService)); const snapshotQuery = await testRuntime.runPromise(Effect.service(ProjectionSnapshotQuery)); const ingestion = await testRuntime.runPromise(Effect.service(ProviderRuntimeIngestionService)); + const backgroundLiveness = await testRuntime.runPromise( + Effect.service(ThreadBackgroundLiveness.ThreadBackgroundLivenessService), + ); scope = await Effect.runPromise(Scope.make("sequential")); await testRuntime.runPromise(ingestion.start().pipe(Scope.provide(scope))); const drain = () => testRuntime.runPromise(ingestion.drain); @@ -385,6 +396,7 @@ describe("ProviderRuntimeIngestion", () => { engine, dispatch, readModel: () => testRuntime.runPromise(snapshotQuery.getSnapshot()), + backgroundLiveness, readThreadShell: () => testRuntime.runPromise( snapshotQuery @@ -4366,4 +4378,255 @@ 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"]); + // 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", + }); + }); + + 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("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. 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" }), + 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 964f60d3a306..4d40c5f25149 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -49,6 +49,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 { resolveProjectSettings } from "@t3tools/shared/projectSettings"; @@ -2043,6 +2044,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" @@ -2055,6 +2059,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. + // 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", + 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..4f840eeb614a 100644 --- a/apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts +++ b/apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts @@ -226,4 +226,129 @@ 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"); + }); + + 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 2781e4981f7c..3925f1db3600 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,20 @@ 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 + * 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 +93,12 @@ 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(); + /** 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); @@ -127,6 +152,21 @@ 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) { + 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) { + hostSettledAtByTaskKey.delete(oldest); + } + } + } return; } @@ -142,7 +182,32 @@ export function make(): ThreadBackgroundLivenessService["Service"] { } } + // 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); + 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.test.ts b/apps/server/src/orchestration/ThreadTaskSettlement.test.ts new file mode 100644 index 000000000000..335b38931670 --- /dev/null +++ b/apps/server/src/orchestration/ThreadTaskSettlement.test.ts @@ -0,0 +1,192 @@ +import * as NodeServices from "@effect/platform-node/NodeServices"; +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"; + +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" } }, +]; + +/** + * 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" && + 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 = { + 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) => + toActivityRows(threadId, input.activitiesByThreadId[threadId] ?? []), + ), + ), + } satisfies ProjectionThreadActivities.ProjectionThreadActivityRepository["Service"]; + + return (effect: Effect.Effect) => + effect.pipe( + Effect.provideService( + ProjectionThreadActivities.ProjectionThreadActivityRepository, + repository, + ), + 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( + 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 new file mode 100644 index 000000000000..4dbcc3f09468 --- /dev/null +++ b/apps/server/src/orchestration/ThreadTaskSettlement.ts @@ -0,0 +1,471 @@ +/** + * 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) 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. 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 + */ +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. 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", + "model", + "effort", + "agentPath", + "timelineBypass", + "parentAgentId", + "workflowName", + "agentIndex", + "phaseIndex", + "phaseTitle", + "phases", + "attempt", + "runHandles", + "outputFile", +] 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; + /** + * 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. */ + 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; +} + +/** + * 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 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(entry.linkage.parentAgentId); + if (parentAgentId !== undefined) { + entry.parentAgentId = parentAgentId; + } +} + +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. */ +function isAgentRow(payload: Record): boolean { + return payload.agentKind === "agent"; +} + +function asFoldStatus(value: unknown): FoldStatus | undefined { + return typeof value === "string" ? KNOWN_STATUSES.get(value) : undefined; +} + +/** Last status row wins, mirroring the client fold's reactivation rule. */ +function applyStatus(entry: FoldEntry, next: FoldStatus): void { + 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, "running"); + if (existing !== undefined && existing.status === "idle") { + entry.status = "running"; + } + mergeLinkage(entry, payload); + entries.set(taskId, entry); + break; + } + case "task.progress": { + const entry = existing ?? newEntry(taskId, "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"; + } + mergeLinkage(entry, payload); + entries.set(taskId, entry); + break; + } + case "task.updated": { + const entry = existing ?? newEntry(taskId, "running"); + const status = asFoldStatus(payload.status); + if (status !== undefined) { + applyStatus(entry, status); + } + mergeLinkage(entry, payload); + entries.set(taskId, entry); + break; + } + case "task.completed": { + const entry = existing ?? newEntry(taskId, "terminal"); + applyStatus(entry, "terminal"); + mergeLinkage(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)) { + live.push({ taskId, linkage: entry.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}`); +} + +/** + * 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, + }; + const commandId = CommandId.make(`task-settle:${yield* crypto.randomUUIDv4}`); + // 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", + 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({ + threadId: input.threadId, + taskId: task.taskId, + 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, + }); + } +}); + +/** 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. + * + * 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 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, + }); + 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; + } + 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; + } + // 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])); + } + }); + + yield* settle.pipe(withSettlementRecovery(input.threadIds)); +}); diff --git a/apps/server/src/persistence/Layers/ProjectionThreadActivities.test.ts b/apps/server/src/persistence/Layers/ProjectionThreadActivities.test.ts index d92ed97ea6fa..dd7d774b00d1 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,171 @@ 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 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 }, + (_, 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 9b6d44d170ee..8710e1de09c1 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadActivities.ts @@ -11,6 +11,7 @@ import { toPersistenceDecodeError, toPersistenceSqlError } from "../Errors.ts"; import { DeleteProjectionThreadActivitiesInput, + ListProjectionThreadActivitiesByThreadIdsInput, ListProjectionThreadActivitiesInput, GetLatestProjectionThreadTaskActivityInput, ProjectionThreadActivity, @@ -41,6 +42,17 @@ function toProjectionThreadActivity( }; } +/** SQLite's host-parameter ceiling is 999; stay well inside it. */ +export 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) @@ -125,6 +137,59 @@ 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 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, @@ -219,6 +284,38 @@ const makeProjectionThreadActivityRepository = Effect.gen(function* () { Effect.map((rows) => rows.map(toProjectionThreadActivity)), ); + const listTaskLifecycleByThreadId: ProjectionThreadActivityRepositoryShape["listTaskLifecycleByThreadId"] = + (input) => + listTaskLifecycleActivityRows(input).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "ProjectionThreadActivityRepository.listTaskLifecycleByThreadId:query", + "ProjectionThreadActivityRepository.listTaskLifecycleByThreadId:decodeRows", + ), + ), + Effect.map((rows) => rows.map(toProjectionThreadActivity)), + ); + + 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((rows) => rows.map(toProjectionThreadActivity)), + ), + { concurrency: 1 }, + ).pipe(Effect.map((chunks) => chunks.flat())); + const listUserInputLifecycleByThreadId: ProjectionThreadActivityRepositoryShape["listUserInputLifecycleByThreadId"] = (input) => listUserInputLifecycleActivityRows(input).pipe( @@ -254,6 +351,8 @@ const makeProjectionThreadActivityRepository = Effect.gen(function* () { return { upsert, listByThreadId, + listTaskLifecycleByThreadId, + listTaskLifecycleByThreadIds, listUserInputLifecycleByThreadId, getLatestTaskActivity, deleteByThreadId, diff --git a/apps/server/src/persistence/Services/ProjectionThreadActivities.ts b/apps/server/src/persistence/Services/ProjectionThreadActivities.ts index 85e9d368df4c..e724cf58cbd1 100644 --- a/apps/server/src/persistence/Services/ProjectionThreadActivities.ts +++ b/apps/server/src/persistence/Services/ProjectionThreadActivities.ts @@ -41,6 +41,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 GetLatestProjectionThreadTaskActivityInput = Schema.Struct({ threadId: ThreadId, taskId: Schema.String, @@ -77,6 +83,29 @@ 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 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: each thread's rows are + * contiguous and ordered within the thread. + */ + readonly listTaskLifecycleByThreadIds: ( + input: ListProjectionThreadActivitiesByThreadIdsInput, + ) => Effect.Effect, ProjectionRepositoryError>; + /** * List activity rows used to derive pending user-input state. * diff --git a/apps/server/src/project/AgentSessionImporter.test.ts b/apps/server/src/project/AgentSessionImporter.test.ts index 4eb03a5cc036..c8871f816cda 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"; @@ -932,6 +934,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/server.ts b/apps/server/src/server.ts index e3b42a2637a1..30f268614549 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -279,8 +279,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 there is no cycle. 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 37fd210ee6da..e64015f6f8c9 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, @@ -24,6 +25,9 @@ import { } 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 { ServerActivation } from "./serverActivation.ts"; import * as ServerSettings from "./serverSettings.ts"; import * as ServerRuntimeStartup from "./serverRuntimeStartup.ts"; @@ -78,6 +82,46 @@ 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> = [], +) => + ({ + 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) => + toActivityRows(threadId, activitiesByThreadId[threadId] ?? []), + ); + }), + }) satisfies ProjectionThreadActivities.ProjectionThreadActivityRepository["Service"]; + const runReconciliation = (input: { readonly threads: ReadonlyArray>; readonly continueAfterRestart?: boolean; @@ -85,12 +129,24 @@ const runReconciliation = (input: { readonly providerService?: ProviderService.ProviderService["Service"]; readonly directory: ProviderSessionDirectory.ProviderSessionDirectory["Service"]; 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( ProjectionSnapshotQuery.ProjectionSnapshotQuery, queryWithThreads(input.threads), ), + Effect.provideService( + ProjectionThreadActivities.ProjectionThreadActivityRepository, + activityRepositoryWith(input.activitiesByThreadId ?? {}, input.batchReads), + ), + Effect.provideService( + ThreadBackgroundLiveness.ThreadBackgroundLivenessService, + input.backgroundLiveness ?? ThreadBackgroundLiveness.make(), + ), Effect.provideService( ProviderService.ProviderService, input.providerService ?? makeProviderService(input.liveThreadIds), @@ -723,11 +779,145 @@ 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(Layer.mergeAll(NodeServices.layer, ServerSettings.layerTest())), 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.succeed([]), + recordImportedTranscript: () => 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.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"), + [], + ); + }), + ), + ); + }, +); + for (const scenario of [ "disabled", "stopped projection", diff --git a/apps/server/src/serverRuntimeStartup.ts b/apps/server/src/serverRuntimeStartup.ts index 90d6c3c576a0..7c8f277294cd 100644 --- a/apps/server/src/serverRuntimeStartup.ts +++ b/apps/server/src/serverRuntimeStartup.ts @@ -34,6 +34,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 { settleThreadsTasks } from "./orchestration/ThreadTaskSettlement.ts"; import * as ServerLifecycleEvents from "./serverLifecycleEvents.ts"; import * as ServerSettings from "./serverSettings.ts"; import * as AnalyticsService from "./telemetry/AnalyticsService.ts"; @@ -538,6 +539,30 @@ 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. 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 + // because 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; if (session === null) { diff --git a/docs/user/composer.md b/docs/user/composer.md index 4a8df5333664..2eb377bbb159 100644 --- a/docs/user/composer.md +++ b/docs/user/composer.md @@ -109,6 +109,14 @@ Provider commands must start the message to run. T3 Code commands such as Send `/compact` in an existing conversation to reduce context usage when the provider supports it. Web and desktop also offer compaction from the context meter. +## 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, interrupts every agent still +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 Select an image or video attachment or link to preview it. Playback support depends diff --git a/packages/client-runtime/src/state/subagentRuntime.test.ts b/packages/client-runtime/src/state/subagentRuntime.test.ts index d366d7f0d4ee..6f54fb33cdf1 100644 --- a/packages/client-runtime/src/state/subagentRuntime.test.ts +++ b/packages/client-runtime/src/state/subagentRuntime.test.ts @@ -891,3 +891,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 row 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"); + }); +});