Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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"),
Expand Down
91 changes: 88 additions & 3 deletions apps/server/src/orchestration/Layers/ProviderCommandReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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;
}

Expand Down Expand Up @@ -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",
Expand Down
Loading