diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index e1019634fc99..8bf392ef1e81 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -545,6 +545,20 @@ describe("ProviderCommandReactor", () => { readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()), readTurns: (threadId: ThreadId) => Effect.runPromise(turnRepository.listByThreadId({ threadId })), + replacePendingTurnStart: (row: { + readonly threadId: ThreadId; + readonly messageId: MessageId; + readonly requestedAt: string; + }) => + harnessRuntime.runPromise( + turnRepository.replacePendingTurnStart({ + threadId: row.threadId, + messageId: row.messageId, + sourceProposedPlanThreadId: null, + sourceProposedPlanId: null, + requestedAt: row.requestedAt, + }), + ), settleSession, startSession, sendTurn, @@ -648,6 +662,35 @@ describe("ProviderCommandReactor", () => { expect(harness.sendTurn).toHaveBeenCalledTimes(1); }); + it("abandons a ghost pending turn start when its event cannot be recovered", async () => { + // Regression: a pending-start row with no matching turn-start-requested event + // used to stick forever, so every follow-up was message-queued and Discord + // stayed on Working… with no drain path (Protect Agents / 77ff9acd). + const harness = await createHarness({ deferReactorStart: true }); + const threadId = ThreadId.make("thread-1"); + const ghostMessageId = asMessageId("ghost-missing-user-message"); + const now = "2026-01-01T00:00:00.000Z"; + + await harness.replacePendingTurnStart({ + threadId, + messageId: ghostMessageId, + requestedAt: now, + }); + const pendingBefore = (await harness.readTurns(threadId)).filter( + (turn) => turn.turnId === null && turn.state === "pending", + ); + expect(pendingBefore).toHaveLength(1); + expect(pendingBefore[0]?.pendingMessageId).toBe(String(ghostMessageId)); + + await harness.startReactor(); + await harness.drain(); + + const pendingAfter = (await harness.readTurns(threadId)).filter( + (turn) => turn.turnId === null && turn.state === "pending", + ); + expect(pendingAfter).toEqual([]); + }); + it("continues an interrupted running turn without replaying its user message", async () => { const modelSelection: ModelSelection = { instanceId: ProviderInstanceId.make("codex"), diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 88d38d886644..f1ab1e9fbb31 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -400,6 +400,76 @@ const make = Effect.gen(function* () { }); }); + /** + * Abandon a pending turn start that can never adopt (missing user message, + * missing persisted event, etc.). Without this, the SQL pending-start row + + * in-memory `pendingTurnStart` flag stay forever, so every follow-up is + * `message-queued` and the queue never drains (Discord stuck on Working…). + * + * `thread.session.set` with a settled status clears the pending-start flag + * in the event-sourced read model and deletes the pending SQL placeholder + * (see ProjectionPipeline session-set handling). We then attempt a queue + * drain so any messages that piled up while stuck can run. + */ + const abandonUnadoptablePendingTurnStart = Effect.fnUntraced(function* (input: { + readonly threadId: ThreadId; + readonly createdAt: string; + readonly reason: string; + }) { + const thread = yield* resolveThread(input.threadId); + if (!thread) { + return; + } + const session = thread.session; + // Prefer ready over error: this is an orchestration glitch, not a live + // provider failure. Keep stopped sessions stopped. + const nextStatus = session?.status === "stopped" ? "stopped" : "ready"; + yield* setThreadSession({ + threadId: input.threadId, + session: { + ...(session ?? { + threadId: input.threadId, + providerName: null, + providerInstanceId: thread.modelSelection.instanceId, + runtimeMode: thread.runtimeMode, + }), + status: nextStatus, + activeTurnId: null, + // Do not sticky-page this internal failure via lastError forever. + lastError: nextStatus === "stopped" ? (session?.lastError ?? null) : null, + updatedAt: input.createdAt, + }, + createdAt: input.createdAt, + }); + // Belt-and-suspenders if session was already ready (session-set may no-op + // some paths) — always clear the pending SQL placeholder. + yield* projectionTurnRepository + .deletePendingTurnStartByThreadId({ + threadId: input.threadId, + }) + .pipe(Effect.catch(() => Effect.void)); + + yield* Effect.logWarning("Abandoned unadoptable pending turn start", { + threadId: input.threadId, + reason: input.reason, + }); + + // Drain any follow-ups that queued while the ghost pending start blocked. + // May no-op (empty queue / already busy); next natural completion retries. + yield* serverCommandId("queue-drain-after-abandon").pipe( + Effect.flatMap((commandId) => + orchestrationEngine + .dispatch({ + type: "thread.queue.drain", + commandId, + threadId: input.threadId, + createdAt: input.createdAt, + }) + .pipe(Effect.catchTag("OrchestrationCommandInvariantError", () => Effect.void)), + ), + ); + }); + const resolveProject = Effect.fnUntraced(function* (projectId: ProjectId) { return yield* projectionSnapshotQuery .getProjectShellById(projectId) @@ -1059,14 +1129,22 @@ const make = Effect.gen(function* () { const message = thread.messages.find((entry) => entry.id === event.payload.messageId); if (!message || message.role !== "user") { + const detail = `User message '${event.payload.messageId}' was not found for turn start request.`; yield* appendProviderFailureActivity({ threadId: event.payload.threadId, kind: "provider.turn.start.failed", summary: "Provider turn start failed", - detail: `User message '${event.payload.messageId}' was not found for turn start request.`, + detail, turnId: null, createdAt: event.payload.createdAt, }); + // Without abandon, the pending-start placeholder remains forever and + // every Discord follow-up is queued with no drain path (Working… forever). + yield* abandonUnadoptablePendingTurnStart({ + threadId: event.payload.threadId, + createdAt: event.payload.createdAt, + reason: detail, + }); return; } @@ -1833,13 +1911,20 @@ const make = Effect.gen(function* () { `${pending.threadId}\u0000${pending.messageId}\u0000${pending.requestedAt}`, ); if (event === undefined) { + const createdAt = DateTime.formatIso(yield* DateTime.now); + const detail = `Persisted turn start event for user message '${pending.messageId}' could not be found.`; yield* appendProviderFailureActivity({ threadId: pending.threadId, kind: "provider.turn.start.failed", summary: "Provider turn start recovery failed", - detail: `Persisted turn start event for user message '${pending.messageId}' could not be found.`, + detail, turnId: null, - createdAt: DateTime.formatIso(yield* DateTime.now), + createdAt, + }); + yield* abandonUnadoptablePendingTurnStart({ + threadId: pending.threadId, + createdAt, + reason: detail, }); yield* increment(providerTurnRecoveriesTotal, { outcome: "failed",