diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index a9d7317999a..926182a3ef0 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -2464,6 +2464,74 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { ); }); +it.layer(makeProjectionPipelinePrefixedTestLayer("t3-pending-turn-terminal-test-"))( + "OrchestrationProjectionPipeline pending turn cleanup", + (it) => { + it.effect("clears pending turn starts when startup reaches a terminal session state", () => + Effect.gen(function* () { + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const sql = yield* SqlClient.SqlClient; + + for (const [index, status] of (["error", "interrupted", "stopped"] as const).entries()) { + const threadId = ThreadId.make(`thread-terminal-${status}`); + const requestedAt = `2026-02-26T14:00:0${index}.000Z`; + yield* eventStore.append({ + type: "thread.turn-start-requested", + eventId: EventId.make(`evt-terminal-pending-${status}`), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: requestedAt, + commandId: CommandId.make(`cmd-terminal-pending-${status}`), + causationEventId: null, + correlationId: CorrelationId.make(`cmd-terminal-pending-${status}`), + metadata: {}, + payload: { + threadId, + messageId: MessageId.make(`message-terminal-${status}`), + runtimeMode: "approval-required", + createdAt: requestedAt, + }, + }); + yield* eventStore.append({ + type: "thread.session-set", + eventId: EventId.make(`evt-terminal-session-${status}`), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: requestedAt, + commandId: CommandId.make(`cmd-terminal-session-${status}`), + causationEventId: null, + correlationId: CorrelationId.make(`cmd-terminal-session-${status}`), + metadata: {}, + payload: { + threadId, + session: { + threadId, + status, + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: status === "error" ? "startup failed" : null, + updatedAt: requestedAt, + }, + }, + }); + } + + yield* projectionPipeline.bootstrap; + + const pendingRows = yield* sql<{ readonly threadId: string }>` + SELECT thread_id AS "threadId" + FROM projection_turns + WHERE turn_id IS NULL + AND state = 'pending' + `; + assert.deepEqual(pendingRows, []); + }), + ); + }, +); + it.effect("restores pending turn-start metadata across projection pipeline restart", () => Effect.gen(function* () { const { dbPath } = yield* ServerConfig; diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index 3ceae0ea43b..cfb88a06cd2 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -1059,6 +1059,15 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti case "thread.session-set": { const turnId = event.payload.session.activeTurnId; if (turnId === null || event.payload.session.status !== "running") { + if ( + event.payload.session.status === "error" || + event.payload.session.status === "stopped" || + event.payload.session.status === "interrupted" + ) { + yield* projectionTurnRepository.deletePendingTurnStartByThreadId({ + threadId: event.payload.threadId, + }); + } // Leaving the "running" session status is the turn-end signal: // settle still-running turns so their duration reflects the whole // turn rather than the last assistant message. diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index ce464565dc5..c49646b7a4b 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -22,6 +22,7 @@ import { TurnId, } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; +import * as Deferred from "effect/Deferred"; import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; import * as ManagedRuntime from "effect/ManagedRuntime"; @@ -145,6 +146,9 @@ describe("ProviderCommandReactor", () => { readonly threadModelSelection?: ModelSelection; readonly sessionModelSwitch?: "unsupported" | "in-session"; readonly requiresNewThreadForModelChange?: boolean; + readonly startSessionEffect?: ( + session: ProviderSession, + ) => Effect.Effect; }) { const now = "2026-01-01T00:00:00.000Z"; const baseDir = @@ -159,6 +163,7 @@ describe("ProviderCommandReactor", () => { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5-codex", }; + const startSessionEffect = input?.startSessionEffect; const startSession = vi.fn((_: unknown, input: unknown) => { const sessionIndex = nextSessionIndex++; const resumeCursor = @@ -212,8 +217,13 @@ describe("ProviderCommandReactor", () => { createdAt: now, updatedAt: now, }; - runtimeSessions.push(session); - return Effect.succeed(session); + return (startSessionEffect?.(session) ?? Effect.succeed(session)).pipe( + Effect.tap((startedSession) => + Effect.sync(() => { + runtimeSessions.push(startedSession); + }), + ), + ); }); const sendTurn = vi.fn((_: unknown) => Effect.succeed({ @@ -463,9 +473,120 @@ describe("ProviderCommandReactor", () => { const readModel = await harness.readModel(); const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); expect(thread?.session?.threadId).toBe("thread-1"); + expect(thread?.session?.status).toBe("starting"); expect(thread?.session?.runtimeMode).toBe("approval-required"); }); + effectIt.effect("projects starting before a slow provider session finishes", () => + Effect.gen(function* () { + const releaseStart = yield* Deferred.make(); + const harness = yield* Effect.promise(() => + createHarness({ + startSessionEffect: (session) => Deferred.await(releaseStart).pipe(Effect.as(session)), + }), + ); + const now = "2026-01-01T00:00:00.000Z"; + + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-slow-provider"), + threadId: ThreadId.make("thread-1"), + message: { + messageId: asMessageId("user-message-slow-provider"), + role: "user", + text: "start slowly", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: now, + }); + + yield* Effect.promise(() => waitFor(() => harness.startSession.mock.calls.length === 1)); + const duringStartup = yield* Effect.promise(() => harness.readModel()); + expect( + duringStartup.threads.find((entry) => entry.id === ThreadId.make("thread-1"))?.session + ?.status, + ).toBe("starting"); + expect(harness.sendTurn).not.toHaveBeenCalled(); + + yield* Deferred.succeed(releaseStart, undefined); + yield* Effect.promise(() => waitFor(() => harness.sendTurn.mock.calls.length === 1)); + }), + ); + + effectIt.effect("settles a failed provider startup and allows a clean retry", () => + Effect.gen(function* () { + let failStartup = true; + const harness = yield* Effect.promise(() => + createHarness({ + startSessionEffect: (session) => + failStartup + ? Effect.fail( + new ProviderAdapterRequestError({ + provider: "codex", + method: "thread.start", + detail: "deterministic startup failure", + }), + ) + : Effect.succeed(session), + }), + ); + const now = "2026-01-01T00:00:00.000Z"; + + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-provider-failure"), + threadId: ThreadId.make("thread-1"), + message: { + messageId: asMessageId("user-message-provider-failure"), + role: "user", + text: "fail once", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: now, + }); + + yield* Effect.promise(() => + waitFor(async () => { + const readModel = await harness.readModel(); + return ( + readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1"))?.session + ?.status === "error" + ); + }), + ); + let readModel = yield* Effect.promise(() => harness.readModel()); + let thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); + expect(thread?.session?.lastError).toContain("deterministic startup failure"); + expect(harness.sendTurn).not.toHaveBeenCalled(); + + failStartup = false; + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-provider-retry"), + threadId: ThreadId.make("thread-1"), + message: { + messageId: asMessageId("user-message-provider-retry"), + role: "user", + text: "retry", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: "2026-01-01T00:00:01.000Z", + }); + + yield* Effect.promise(() => waitFor(() => harness.sendTurn.mock.calls.length === 1)); + readModel = yield* Effect.promise(() => harness.readModel()); + thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); + expect(thread?.session?.status).toBe("starting"); + expect(thread?.session?.lastError).toBeNull(); + }), + ); + it("generates a thread title on the first turn", async () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 9c7a7c94bb1..b6bff8c766a 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -288,15 +288,20 @@ const make = Effect.gen(function* () { readonly createdAt: string; }) { const thread = yield* resolveThread(input.threadId); - const session = thread?.session; - if (!session) { + if (!thread) { return; } + const session = thread.session; yield* setThreadSession({ threadId: input.threadId, session: { - ...session, - status: session.status === "stopped" ? "stopped" : "ready", + ...(session ?? { + threadId: input.threadId, + providerName: null, + providerInstanceId: thread.modelSelection.instanceId, + runtimeMode: thread.runtimeMode, + }), + status: session?.status === "stopped" ? "stopped" : "error", activeTurnId: null, lastError: input.detail, updatedAt: input.createdAt, @@ -354,6 +359,7 @@ const make = Effect.gen(function* () { createdAt: string, options?: { readonly modelSelection?: ModelSelection; + readonly pendingTurnStart?: boolean; }, ) { const thread = yield* resolveThread(threadId); @@ -428,6 +434,22 @@ const make = Effect.gen(function* () { }); } const preferredProvider: ProviderDriverKind = desiredDriverKind; + if (options?.pendingTurnStart === true && thread.session?.status !== "running") { + yield* setThreadSession({ + threadId, + session: { + threadId, + status: "starting", + providerName: activeSession?.provider ?? preferredProvider, + providerInstanceId: activeSession?.providerInstanceId ?? desiredInstanceId, + runtimeMode: desiredRuntimeMode, + activeTurnId: null, + lastError: null, + updatedAt: createdAt, + }, + createdAt, + }); + } if (thread.session !== null) { yield* rejectStartedThreadModelChangeIfRequired({ threadId, @@ -498,7 +520,10 @@ const make = Effect.gen(function* () { threadId, session: { threadId, - status: mapProviderSessionStatusToOrchestrationStatus(session.status), + status: + options?.pendingTurnStart === true && session.status === "ready" + ? "starting" + : mapProviderSessionStatusToOrchestrationStatus(session.status), providerName: session.provider, providerInstanceId: session.providerInstanceId, runtimeMode: desiredRuntimeMode, @@ -597,11 +622,10 @@ const make = Effect.gen(function* () { new Error(`Thread '${input.threadId}' was not found in read model.`), ); } - yield* ensureSessionForThread( - input.threadId, - input.createdAt, - input.modelSelection !== undefined ? { modelSelection: input.modelSelection } : {}, - ); + yield* ensureSessionForThread(input.threadId, input.createdAt, { + ...(input.modelSelection !== undefined ? { modelSelection: input.modelSelection } : {}), + pendingTurnStart: true, + }); if (input.modelSelection !== undefined) { threadModelSelections.set(input.threadId, input.modelSelection); } diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 20611e1ee75..74ece50cd31 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -30,6 +30,7 @@ import * as ManagedRuntime from "effect/ManagedRuntime"; import * as PubSub from "effect/PubSub"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; +import { it as effectIt } from "@effect/vitest"; import { afterEach, describe, expect, it } from "vite-plus/test"; import { OrchestrationEventStoreLive } from "../../persistence/Layers/OrchestrationEventStore.ts"; @@ -493,6 +494,185 @@ describe("ProviderRuntimeIngestion", () => { expect(thread.session?.lastError).toBeNull(); }); + effectIt.effect( + "keeps a reconnecting pending turn starting while ready clears stale active state", + () => + Effect.gen(function* () { + const harness = yield* Effect.promise(() => createHarness()); + const threadId = asThreadId("thread-1"); + const staleTurnId = asTurnId("turn-stale-before-reconnect"); + + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-pending-reconnect"), + threadId, + message: { + messageId: MessageId.make("message-pending-reconnect"), + role: "user", + text: "resume after reconnect", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: "2026-01-01T00:00:01.000Z", + }); + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-starting-pending-reconnect"), + threadId, + session: { + threadId, + status: "starting", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: staleTurnId, + lastError: null, + updatedAt: "2026-01-01T00:00:01.000Z", + }, + createdAt: "2026-01-01T00:00:01.000Z", + }); + + harness.emit({ + type: "session.state.changed", + eventId: asEventId("evt-session-ready-pending-reconnect"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:02.000Z", + payload: { state: "ready" }, + }); + + let thread = yield* Effect.promise(() => + waitForThread( + harness.readModel, + (entry) => entry.session?.status === "starting" && entry.session.activeTurnId === null, + ), + ); + expect(thread.session?.status).toBe("starting"); + expect(thread.session?.activeTurnId).toBeNull(); + + harness.emit({ + type: "session.started", + eventId: asEventId("evt-session-started-pending-reconnect"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:03.000Z", + }); + yield* Effect.promise(() => harness.drain()); + thread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + )!; + expect(thread.session?.status).toBe("starting"); + expect(thread.session?.activeTurnId).toBeNull(); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-pending-reconnect"), + provider: ProviderDriverKind.make("codex"), + threadId, + turnId: asTurnId("turn-after-reconnect"), + createdAt: "2026-01-01T00:00:04.000Z", + }); + thread = yield* Effect.promise(() => + waitForThread( + harness.readModel, + (entry) => + entry.session?.status === "running" && + entry.session.activeTurnId === asTurnId("turn-after-reconnect"), + ), + ); + expect(thread.session?.status).toBe("running"); + + harness.emit({ + type: "session.started", + eventId: asEventId("evt-session-started-duplicate-midturn"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:05.000Z", + }); + yield* Effect.promise(() => harness.drain()); + thread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + )!; + expect(thread.session?.status).toBe("running"); + expect(thread.session?.activeTurnId).toBe(asTurnId("turn-after-reconnect")); + }), + ); + + effectIt.effect("keeps an aborted pending start stopped across duplicate exit events", () => + Effect.gen(function* () { + const harness = yield* Effect.promise(() => createHarness()); + const threadId = asThreadId("thread-1"); + const stoppedAt = "2026-01-01T00:00:02.000Z"; + + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-before-stop"), + threadId, + message: { + messageId: MessageId.make("message-before-stop"), + role: "user", + text: "stop this startup", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: "2026-01-01T00:00:01.000Z", + }); + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-starting-before-stop"), + threadId, + session: { + threadId, + status: "starting", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: "2026-01-01T00:00:01.000Z", + }, + createdAt: "2026-01-01T00:00:01.000Z", + }); + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-stop-pending-start"), + threadId, + session: { + threadId, + status: "stopped", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: stoppedAt, + }, + createdAt: stoppedAt, + }); + + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-session-exited-after-stop"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:03.000Z", + }); + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-duplicate-session-exited-after-stop"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:04.000Z", + }); + + yield* Effect.promise(() => harness.drain()); + const thread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + ); + expect(thread?.session?.status).toBe("stopped"); + expect(thread?.session?.activeTurnId).toBeNull(); + }), + ); + it("does not clear active turn when session/thread started arrives mid-turn", async () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 09566feb2b2..a8a51b30260 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -1308,6 +1308,11 @@ const make = Effect.gen(function* () { const now = event.createdAt; const eventTurnId = toTurnId(event.turnId); const activeTurnId = thread.session?.activeTurnId ?? null; + const pendingTurnStart = yield* projectionTurnRepository.getPendingTurnStartByThreadId({ + threadId: thread.id, + }); + const hasPendingTurnStart = + Option.isSome(pendingTurnStart) && thread.session?.status === "starting"; const conflictsWithActiveTurn = activeTurnId !== null && eventTurnId !== undefined && !sameId(activeTurnId, eventTurnId); @@ -1322,11 +1327,7 @@ const make = Effect.gen(function* () { const conflictingTurnStartIsPendingTurnStart = event.type === "turn.started" && conflictsWithActiveTurn ? sameId(yield* getExpectedProviderTurnIdForThread(thread.id), eventTurnId) && - Option.isSome( - yield* projectionTurnRepository.getPendingTurnStartByThreadId({ - threadId: thread.id, - }), - ) + Option.isSome(pendingTurnStart) : false; const shouldApplyThreadLifecycle = (() => { @@ -1370,8 +1371,10 @@ const make = Effect.gen(function* () { ) { const status = (() => { switch (event.type) { - case "session.state.changed": - return orchestrationSessionStatusFromRuntimeState(event.payload.state); + case "session.state.changed": { + const runtimeStatus = orchestrationSessionStatusFromRuntimeState(event.payload.state); + return hasPendingTurnStart && runtimeStatus === "ready" ? "starting" : runtimeStatus; + } case "turn.started": return "running"; case "session.exited": @@ -1383,8 +1386,8 @@ const make = Effect.gen(function* () { case "session.started": case "thread.started": // Provider thread/session start notifications can arrive during an - // active turn; preserve turn-running state in that case. - return activeTurnId !== null ? "running" : "ready"; + // active or pending turn; preserve that lifecycle state. + return activeTurnId !== null ? "running" : hasPendingTurnStart ? "starting" : "ready"; } })(); const nextActiveTurnId = @@ -1392,7 +1395,10 @@ const make = Effect.gen(function* () { ? (eventTurnId ?? null) : event.type === "turn.completed" || event.type === "session.exited" ? null - : event.type === "session.state.changed" && !sessionStatusAllowsActiveTurn(status) + : event.type === "session.state.changed" && + !sessionStatusAllowsActiveTurn( + orchestrationSessionStatusFromRuntimeState(event.payload.state), + ) ? null : activeTurnId; const lastError = diff --git a/packages/client-runtime/src/state/threadSettled.test.ts b/packages/client-runtime/src/state/threadSettled.test.ts index 50f3c2d3912..59f9f1d6f55 100644 --- a/packages/client-runtime/src/state/threadSettled.test.ts +++ b/packages/client-runtime/src/state/threadSettled.test.ts @@ -178,6 +178,51 @@ describe("effectiveSettled", () => { ).toBe(false); }); + it("keeps a new turn active from queued through starting and running", () => { + const requestedAt = "2026-04-09T12:00:00.000Z"; + const transitionNow = "2026-04-09T12:00:30.000Z"; + const base = makeShell({ + settledOverride: null, + activityAt: STALE, + }); + const queued: OrchestrationThreadShell = { + ...base, + latestUserMessageAt: requestedAt, + latestTurn: null, + session: null, + }; + const starting: OrchestrationThreadShell = { + ...queued, + session: { + threadId: queued.id, + status: "starting", + providerName: "Codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: requestedAt, + }, + }; + const running: OrchestrationThreadShell = { + ...starting, + session: { + ...starting.session!, + status: "running", + activeTurnId: TurnId.make("turn-new"), + }, + }; + + for (const shell of [queued, starting, running]) { + expect( + effectiveSettled(shell, { + now: transitionNow, + autoSettleAfterDays: 3, + changeRequestState: "merged", + }), + ).toBe(false); + } + }); + it("uses a strict inactivity boundary and honors a null threshold", () => { const boundary = makeShell({ activityAt: "2026-04-07T00:00:00.000Z",