diff --git a/packages/junior/src/chat/agent/index.ts b/packages/junior/src/chat/agent/index.ts index ffee90649..8ebcd12d4 100644 --- a/packages/junior/src/chat/agent/index.ts +++ b/packages/junior/src/chat/agent/index.ts @@ -1141,6 +1141,10 @@ async function executeAgentRunInPrivacyContext( if (promptPersisted) { await runResume.commitInput(); } + } else if (durability.onInputCommitted) { + // Prompt already owned by the running record. Still ack mailbox + // ownership when redelivery is carrying a pending inbound. + await runResume.commitInput(); } /** Race one provider operation against the turn deadline and abort its owner. */ diff --git a/packages/junior/src/chat/agent/prompt.ts b/packages/junior/src/chat/agent/prompt.ts index 01a85ec2b..9fd611417 100644 --- a/packages/junior/src/chat/agent/prompt.ts +++ b/packages/junior/src/chat/agent/prompt.ts @@ -435,11 +435,11 @@ export async function assemblePrompt(args: { userContentParts: UserContentPart[]; }): Promise { const source = args.routing.source; - const hasPromptCheckpoint = - args.resumedFromSessionRecord && - args.existingTurnStartMessageIndex !== undefined; - const shouldPromptAgent = - !args.resumedFromSessionRecord || !hasPromptCheckpoint; + // The turn-start cursor is the durable ownership signal for the prompt. + // Resume classification can lag (still-running redelivery), but once that + // cursor exists the prompt is already committed and must not be rebuilt. + const hasPromptCheckpoint = args.existingTurnStartMessageIndex !== undefined; + const shouldPromptAgent = !hasPromptCheckpoint; const requestContentParts: UserContentPart[] = [ ...(args.explicitSkill ? [ diff --git a/packages/junior/src/chat/services/turn-session-record.ts b/packages/junior/src/chat/services/turn-session-record.ts index 7c8f82693..a2425d9ad 100644 --- a/packages/junior/src/chat/services/turn-session-record.ts +++ b/packages/junior/src/chat/services/turn-session-record.ts @@ -85,12 +85,19 @@ export async function loadTurnSessionRecord( ctx.conversationId, ctx.sessionId, ); - const hasAwaitingResumeRecord = Boolean( - existingSessionRecord && existingSessionRecord.state === "awaiting_resume", + // A still-running record with a committed turn-start cursor already owns its + // prompt. Treat that as resume so redelivery continues instead of rebuilding + // and re-checkpointing a non-identical user message. + const hasCommittedPrompt = + existingSessionRecord?.turnStartMessageIndex !== undefined; + const resumedFromSessionRecord = Boolean( + existingSessionRecord && + (existingSessionRecord.state === "awaiting_resume" || + (existingSessionRecord.state === "running" && hasCommittedPrompt)), ); return { - resumedFromSessionRecord: hasAwaitingResumeRecord, - currentSliceId: hasAwaitingResumeRecord + resumedFromSessionRecord, + currentSliceId: resumedFromSessionRecord ? existingSessionRecord!.sliceId : 1, existingSessionRecord, @@ -177,9 +184,22 @@ export async function persistRunningSessionRecord(args: { }); return true; } catch (recordError) { + // Boundary mismatch is permanent for this attempt's message shape. Log it + // distinctly, but still return false: the poison-turn farm is stopped by + // treating already-checkpointed running turns as resume (no rebuild / + // recheckpoint), not by failing closed here. Handoff/compaction follow-ups + // and no-checkpoint resume still rely on false when no mailbox ack is + // pending. + const message = + recordError instanceof Error ? recordError.message : String(recordError); + const boundaryMismatch = message.includes( + "changed before its committed boundary", + ); logSessionRecordError( recordError, - "agent.turn.running_session_record.failed", + boundaryMismatch + ? "agent.turn.running_session_record.boundary_mismatch" + : "agent.turn.running_session_record.failed", args, { "app.ai.resume_slice_id": args.sliceId, diff --git a/packages/junior/src/chat/task-execution/state.ts b/packages/junior/src/chat/task-execution/state.ts index aa8065268..f77b84dc1 100644 --- a/packages/junior/src/chat/task-execution/state.ts +++ b/packages/junior/src/chat/task-execution/state.ts @@ -60,6 +60,8 @@ export const CONVERSATION_WORK_LEASE_TTL_MS = 90_000; export const CONVERSATION_WORK_CHECK_IN_INTERVAL_MS = 15_000; export const CONVERSATION_WORK_STALE_ENQUEUE_MS = 60_000; export const CONVERSATION_WORK_MAX_DELIVERY_ATTEMPTS = 5; +/** Empty continue/lost-lease wakes with no mailbox progress before fail-closed. */ +export const CONVERSATION_WORK_MAX_EMPTY_WAKES = 5; const inboundMessageSourceSchema = z.enum([ "api", @@ -127,6 +129,8 @@ export interface Lease { } export interface ConversationExecution { + /** Consecutive empty wakes (no mailbox attempt) that made no durable progress. */ + emptyWakeCount?: number; inboundMessageIds: string[]; lastCheckpointAtMs?: number; lastEnqueuedAtMs?: number; @@ -432,6 +436,7 @@ function normalizeExecution( pendingCount: pendingMessages.length, pendingMessages, lease, + emptyWakeCount: toOptionalNumber(value.emptyWakeCount), lastCheckpointAtMs: toOptionalNumber(value.lastCheckpointAtMs), lastEnqueuedAtMs: toOptionalNumber(value.lastEnqueuedAtMs), runId: toOptionalString(value.runId), @@ -1138,6 +1143,8 @@ export async function appendInboundMessage(args: { next, { ...current.execution, + // Fresh mailbox work resets the empty-wake poison budget. + emptyWakeCount: undefined, status, inboundMessageIds: [ ...current.execution.inboundMessageIds, @@ -1741,11 +1748,102 @@ export async function completeConversationWork(args: { }); } +/** + * Record one empty wake that made no mailbox progress (continue / lost-lease + * recovery with nothing to ack). Caps poison requeue farms that bypass the + * per-message delivery attempt counter. + */ +export async function recordEmptyWakeFailure(args: { + conversationId: string; + leaseToken: string; + nowMs?: number; + state?: StateAdapter; +}): Promise { + const nowMs = args.nowMs ?? now(); + return await withConversationMutation(args, async (state, lock) => { + const current = await readConversation(state, args.conversationId); + if (!current || current.execution.lease?.token !== args.leaseToken) { + return { + status: "lost_lease", + pendingCount: 0, + deadLetteredMessages: [], + }; + } + const emptyWakeCount = (current.execution.emptyWakeCount ?? 0) + 1; + const terminal = emptyWakeCount >= CONVERSATION_WORK_MAX_EMPTY_WAKES; + await writeConversation( + state, + lock, + withExecutionUpdate( + current, + { + ...current.execution, + emptyWakeCount, + ...(terminal + ? { + lease: undefined, + status: + pendingMessages(current).length > 0 ? "pending" : "failed", + } + : {}), + }, + nowMs, + ), + ); + return { + status: "recorded", + pendingCount: current.execution.pendingMessages.length, + deadLetteredMessages: [], + ...(terminal ? { terminal: true as const } : {}), + }; + }); +} + +/** Clear empty-wake streak after durable progress (ack, complete, new work). */ +export async function clearEmptyWakeCount(args: { + conversationId: string; + leaseToken?: string; + nowMs?: number; + state?: StateAdapter; +}): Promise { + const nowMs = args.nowMs ?? now(); + return await withConversationMutation(args, async (state, lock) => { + const current = await readConversation(state, args.conversationId); + if (!current) { + return false; + } + if ( + args.leaseToken !== undefined && + current.execution.lease?.token !== args.leaseToken + ) { + return false; + } + if ((current.execution.emptyWakeCount ?? 0) === 0) { + return true; + } + await writeConversation( + state, + lock, + withExecutionUpdate( + current, + { + ...current.execution, + emptyWakeCount: undefined, + }, + nowMs, + ), + ); + return true; + }); +} + /** Failure outcome: `lost_lease` (another owner took over), `recorded` (attempt counted), or `skipped` (durable progress was made). */ export interface AttemptFailure { pendingCount: number; deadLetteredMessages: InboundMessage[]; status: "lost_lease" | "recorded" | "skipped"; + /** True when empty-wake failures hit the fail-closed ceiling. */ + terminal?: boolean; } /** diff --git a/packages/junior/src/chat/task-execution/store.ts b/packages/junior/src/chat/task-execution/store.ts index 05c256c61..a7c3c5ec0 100644 --- a/packages/junior/src/chat/task-execution/store.ts +++ b/packages/junior/src/chat/task-execution/store.ts @@ -9,6 +9,7 @@ export { CONVERSATION_WORK_CHECK_IN_INTERVAL_MS, CONVERSATION_WORK_LEASE_TTL_MS, CONVERSATION_WORK_MAX_DELIVERY_ATTEMPTS, + CONVERSATION_WORK_MAX_EMPTY_WAKES, CONVERSATION_WORK_STALE_ENQUEUE_MS, isFinalAttempt, isInvalidConversationRecordError, @@ -453,6 +454,32 @@ export async function completeConversationWork(args: { return result; } +/** Record one empty wake with no mailbox progress before fail-closed. */ +export async function recordEmptyWakeFailure(args: { + conversationId: string; + leaseToken: string; + conversationStore?: ConversationStore; + nowMs?: number; + state?: StateAdapter; +}) { + const result = await workState.recordEmptyWakeFailure(args); + if (result.status === "recorded") { + await recordExecutionMetadata(args); + } + return result; +} + +/** Clear empty-wake streak after durable progress. */ +export async function clearEmptyWakeCount(args: { + conversationId: string; + leaseToken?: string; + conversationStore?: ConversationStore; + nowMs?: number; + state?: StateAdapter; +}) { + return await workState.clearEmptyWakeCount(args); +} + /** Record one failed delivery attempt and dead-letter messages at their limit. */ export async function recordAttemptFailure(args: { conversationId: string; diff --git a/packages/junior/src/chat/task-execution/worker.ts b/packages/junior/src/chat/task-execution/worker.ts index 79248a890..ba8f80d26 100644 --- a/packages/junior/src/chat/task-execution/worker.ts +++ b/packages/junior/src/chat/task-execution/worker.ts @@ -14,6 +14,7 @@ import { beginConversationResume, checkInConversationWork, clearConsumedConversationWake, + clearEmptyWakeCount, completeConversationWork, CONVERSATION_WORK_CHECK_IN_INTERVAL_MS, countPendingConversationMessages, @@ -24,6 +25,7 @@ import { isFinalAttempt, isInvalidConversationRecordError, recordAttemptFailure, + recordEmptyWakeFailure, releaseConversationWork, requestConversationContinuation, startConversationWork, @@ -120,13 +122,37 @@ function nudgeIdempotencyKey( return `${reason}:${conversationId}:${nowMs}`; } +/** + * Requeue after lease loss. Empty wakes (no mailbox attempt) count toward + * CONVERSATION_WORK_MAX_EMPTY_WAKES so poison continue loops fail closed. + */ async function requestLostLeaseRecovery(args: { conversationId: string; destination: Destination; + /** True when this wake had no mailbox messages to attempt. */ + emptyWake: boolean; leaseToken: string; nowMs: number; options: ProcessConversationWorkOptions; -}): Promise { +}): Promise<"requeued" | "terminal" | "skipped"> { + if (args.emptyWake) { + const emptyFailure = await recordEmptyWakeFailure({ + conversationId: args.conversationId, + leaseToken: args.leaseToken, + nowMs: args.nowMs, + state: args.options.state, + }); + if (emptyFailure.status === "lost_lease") { + return "skipped"; + } + if (emptyFailure.terminal) { + logWarn("conversation.work.empty_wake.terminal", { + "app.conversation.empty_wake_count": "max", + }); + return "terminal"; + } + } + const resumeRequested = await requestConversationContinuation({ conversationId: args.conversationId, destination: args.destination, @@ -136,7 +162,7 @@ async function requestLostLeaseRecovery(args: { state: args.options.state, }); if (!resumeRequested) { - return; + return "skipped"; } const released = await releaseConversationWork({ conversationId: args.conversationId, @@ -146,7 +172,7 @@ async function requestLostLeaseRecovery(args: { state: args.options.state, }); if (!released) { - return; + return "skipped"; } await ensureConversationWake({ conversationId: args.conversationId, @@ -161,6 +187,7 @@ async function requestLostLeaseRecovery(args: { replaceExistingWake: true, state: args.options.state, }); + return "requeued"; } /** @@ -434,14 +461,17 @@ async function processConversationWorkInContext( leaseLost ) { markLeaseLost(); - await requestLostLeaseRecovery({ + const recovery = await requestLostLeaseRecovery({ conversationId, destination, + emptyWake: true, leaseToken: lease.leaseToken, nowMs: now(options), options, }); - return { status: "lost_lease" }; + return { + status: recovery === "terminal" ? "failed" : "lost_lease", + }; } const resumePending = leasedWork.execution.status === "awaiting_resume"; @@ -467,14 +497,17 @@ async function processConversationWorkInContext( }); if (!resumeStarted) { markLeaseLost(); - await requestLostLeaseRecovery({ + const recovery = await requestLostLeaseRecovery({ conversationId, destination, + emptyWake: true, leaseToken: lease.leaseToken, nowMs: now(options), options, }); - return { status: "lost_lease" }; + return { + status: recovery === "terminal" ? "failed" : "lost_lease", + }; } } @@ -517,26 +550,40 @@ async function processConversationWorkInContext( const result = await options.run(workerContext); hasRun = true; + const emptyWake = attemptMessageIds.length === 0; if (result.status === "lost_lease") { - await requestLostLeaseRecovery({ + const recovery = await requestLostLeaseRecovery({ conversationId, destination, + emptyWake, leaseToken: lease.leaseToken, nowMs: now(options), options, }); - return { status: "lost_lease" }; + return { + status: recovery === "terminal" ? "failed" : "lost_lease", + }; } if (leaseLost) { - await requestLostLeaseRecovery({ + const recovery = await requestLostLeaseRecovery({ conversationId, destination, + emptyWake, leaseToken: lease.leaseToken, nowMs: now(options), options, }); - return { status: "lost_lease" }; + return { + status: recovery === "terminal" ? "failed" : "lost_lease", + }; } + // Successful slices (including empty continue resumes) made progress. + await clearEmptyWakeCount({ + conversationId, + leaseToken: lease.leaseToken, + nowMs: now(options), + state: options.state, + }); if (result.status === "yielded") { const resumeRequested = await requestConversationContinuation({ conversationId, @@ -653,6 +700,12 @@ async function processConversationWorkInContext( : { status: "completed" }; } + // Lease may already be released by completeConversationWork. + await clearEmptyWakeCount({ + conversationId, + nowMs: now(options), + state: options.state, + }); logInfo("conversation.work.completed", { "app.worker.elapsed_ms": now(options) - startedAtMs, }); @@ -684,28 +737,47 @@ async function processConversationWorkInContext( state: options.state, }); } else { - const resumeRequested = await requestConversationContinuation({ - conversationId, - destination, - leaseToken: lease.leaseToken, - conversationStore: options.conversationStore, - nowMs: errorNowMs, - state: options.state, - }); - if (resumeRequested) { - await ensureConversationWake({ + let skipWake = false; + if (attemptMessageIds.length === 0) { + const emptyFailure = await recordEmptyWakeFailure({ conversationId, + leaseToken: lease.leaseToken, + nowMs: errorNowMs, + state: options.state, + }); + if (emptyFailure.status === "lost_lease") { + skipWake = true; + } else if (emptyFailure.terminal) { + logWarn("conversation.work.empty_wake.terminal", { + "app.conversation.empty_wake_count": "max", + }); + skipWake = true; + } + } + if (!skipWake) { + const resumeRequested = await requestConversationContinuation({ + conversationId, + destination, + leaseToken: lease.leaseToken, conversationStore: options.conversationStore, - idempotencyKey: nudgeIdempotencyKey( - "error", - conversationId, - errorNowMs, - ), nowMs: errorNowMs, - queue: options.queue, - replaceExistingWake: true, state: options.state, }); + if (resumeRequested) { + await ensureConversationWake({ + conversationId, + conversationStore: options.conversationStore, + idempotencyKey: nudgeIdempotencyKey( + "error", + conversationId, + errorNowMs, + ), + nowMs: errorNowMs, + queue: options.queue, + replaceExistingWake: true, + state: options.state, + }); + } } await releaseConversationWork({ conversationId, diff --git a/packages/junior/tests/component/plugins/plugin-prompt-hooks.test.ts b/packages/junior/tests/component/plugins/plugin-prompt-hooks.test.ts index 24fd24df2..263c8f101 100644 --- a/packages/junior/tests/component/plugins/plugin-prompt-hooks.test.ts +++ b/packages/junior/tests/component/plugins/plugin-prompt-hooks.test.ts @@ -117,7 +117,10 @@ import { z } from "zod"; import { executeAgentRun } from "@/chat/agent"; import { setPlugins } from "@/chat/plugins/agent-hooks"; import { disconnectStateAdapter } from "@/chat/state/adapter"; -import { upsertAgentTurnSessionRecord } from "@/chat/state/turn-session"; +import { + getAgentTurnSessionRecord, + upsertAgentTurnSessionRecord, +} from "@/chat/state/turn-session"; import { getConversationEventStore } from "@/chat/db"; import { TurnInputCommitLostError } from "@/chat/runtime/turn"; @@ -313,15 +316,12 @@ describe("plugin prompt hook composition", () => { }), ).rejects.toBeInstanceOf(TurnInputCommitLostError); + // Prompt is already owned by the still-running record. Redelivery continues + // that history instead of rebuilding / re-running user-prompt hooks. await executeAgentRun(request); expect(recallCount).toBe(1); - expect(JSON.stringify(captured.promptMessages[0])).toContain( - "Use pnpm snapshot 1.", - ); - expect(JSON.stringify(captured.promptMessages[0])).not.toContain( - "Use pnpm snapshot 2.", - ); + expect(captured.promptMessages).toEqual([]); const stored = await getConversationEventStore().loadByIdempotencyKey( LOCAL_DESTINATION.conversationId, `turn:${turnId}:context:memory:0`, @@ -332,6 +332,16 @@ describe("plugin prompt hook composition", () => { memories: [{ id: "memory-1", content: "Use pnpm snapshot 1." }], }, }); + const session = await getAgentTurnSessionRecord( + LOCAL_DESTINATION.conversationId, + turnId, + ); + expect(JSON.stringify(session?.piMessages)).toContain( + "Use pnpm snapshot 1.", + ); + expect(JSON.stringify(session?.piMessages)).not.toContain( + "Use pnpm snapshot 2.", + ); }); it("runs user prompt hooks for non-bootstrap follow-up prompts", async () => { diff --git a/packages/junior/tests/component/runtime/agent-run-provider-retry.test.ts b/packages/junior/tests/component/runtime/agent-run-provider-retry.test.ts index 32ad917bb..0295827cb 100644 --- a/packages/junior/tests/component/runtime/agent-run-provider-retry.test.ts +++ b/packages/junior/tests/component/runtime/agent-run-provider-retry.test.ts @@ -1351,6 +1351,8 @@ describe("agent run continuation", () => { await executeAgentRun({ conversationId, turnId: sessionId, + // A running record with turnStartMessageIndex already owns the prompt. + // Redelivery must continue that history, not rebuild/recheckpoint it. input: { messageText: "help me", piMessages: [checkpointedPrompt] }, routing: { destinationVisibility: "private", @@ -1358,6 +1360,9 @@ describe("agent run continuation", () => { source: TEST_SOURCE, actor: { platform: "slack", teamId: "T123", userId: "U123" }, }, + durability: { + onInputCommitted: async () => undefined, + }, }), ); @@ -1370,6 +1375,10 @@ describe("agent run continuation", () => { sessionRecord?.piMessages.filter((message) => message.role === "user") ?? []; expect(userMessages).toHaveLength(1); + expect(userMessages[0]).toMatchObject({ + role: "user", + timestamp: 5, + }); expect( JSON.stringify(sessionRecord?.piMessages).split("help me"), ).toHaveLength(2); diff --git a/packages/junior/tests/component/task-execution/conversation-work.test.ts b/packages/junior/tests/component/task-execution/conversation-work.test.ts index 6f3538421..1c1e290d2 100644 --- a/packages/junior/tests/component/task-execution/conversation-work.test.ts +++ b/packages/junior/tests/component/task-execution/conversation-work.test.ts @@ -13,6 +13,7 @@ import { completeConversationWork, CONVERSATION_WORK_LEASE_TTL_MS, CONVERSATION_WORK_MAX_DELIVERY_ATTEMPTS, + CONVERSATION_WORK_MAX_EMPTY_WAKES, countPendingConversationMessages, drainConversationMailbox, getConversationWorkState, @@ -1291,6 +1292,8 @@ describe("conversation work execution", () => { expect(state?.needsRun).toBe(true); expect(state ? countPendingConversationMessages(state) : 0).toBe(1); expect(state?.lastEnqueuedAtMs).toBe(2_000); + // Pending mailbox work is not an empty wake — do not burn the empty-wake budget. + expect(state?.execution.emptyWakeCount).toBeUndefined(); expect(queue.sentRecords()).toEqual([ { conversationId: CONVERSATION_ID, @@ -1299,6 +1302,77 @@ describe("conversation work execution", () => { ]); }); + it("fails closed after consecutive empty continue wakes without mailbox progress", async () => { + const queue = createConversationWorkQueueTestAdapter(); + let currentNowMs = 1_000; + // Seed an awaiting_resume conversation with no pending mailbox messages — + // the poison continue / lost_lease farm path. + await requestConversationWork({ + conversationId: CONVERSATION_ID, + destination: SLACK_DESTINATION, + nowMs: 1_000, + }); + const started = await startConversationWork({ + conversationId: CONVERSATION_ID, + nowMs: 1_000, + }); + expect(started.status).toBe("acquired"); + if (started.status !== "acquired") return; + await expect( + requestConversationContinuation({ + conversationId: CONVERSATION_ID, + destination: SLACK_DESTINATION, + leaseToken: started.leaseToken, + nowMs: 1_000, + }), + ).resolves.toBe(true); + await releaseConversationWork({ + conversationId: CONVERSATION_ID, + leaseToken: started.leaseToken, + nowMs: 1_000, + }); + + const results: Array<{ status: string }> = []; + for (let attempt = 0; attempt < CONVERSATION_WORK_MAX_EMPTY_WAKES; attempt++) { + currentNowMs = 2_000 + attempt; + results.push( + await processConversationWork(conversationQueueMessage(), { + nowMs: () => currentNowMs, + queue, + run: async () => { + currentNowMs += 1; + return { status: "lost_lease" }; + }, + }), + ); + } + expect(results.slice(0, -1)).toEqual( + Array.from({ length: CONVERSATION_WORK_MAX_EMPTY_WAKES - 1 }, () => ({ + status: "lost_lease", + })), + ); + expect(results.at(-1)).toEqual({ status: "failed" }); + + const state = await getConversationWorkState({ + conversationId: CONVERSATION_ID, + }); + expect(state?.execution.emptyWakeCount).toBe( + CONVERSATION_WORK_MAX_EMPTY_WAKES, + ); + expect(state?.execution.status).toBe("failed"); + expect(state?.lease).toBeUndefined(); + // Final terminal wake must not schedule another recovery nudge. + expect( + queue + .sentRecords() + .filter((record) => + (record.idempotencyKey ?? "").startsWith( + `lost_lease:${CONVERSATION_ID}:`, + ), + ), + ).toHaveLength(CONVERSATION_WORK_MAX_EMPTY_WAKES - 1); + }); + it("drains pending messages and completes the leased conversation", async () => { const queue = createConversationWorkQueueTestAdapter(); await appendInboundMessage({ message: inboundMessage("m1"), nowMs: 1_000 });