diff --git a/server/src/modules/wake-queue/adapters/postgres.test.ts b/server/src/modules/wake-queue/adapters/postgres.test.ts index 974c592365..c7bb2ca43a 100644 --- a/server/src/modules/wake-queue/adapters/postgres.test.ts +++ b/server/src/modules/wake-queue/adapters/postgres.test.ts @@ -393,6 +393,91 @@ describeEmbeddedPostgres("wake-queue postgres adapter", () => { expect(runs.map((run) => run.id).sort()).toEqual([runId, runA, runB, runC].sort()); }); + // The application layer now resolves the responsible user and builds the + // context snapshot before this write runs (proven in the application-layer + // test). This proves the adapter persists the caller-resolved responsible + // user, and merges the two stage fields it alone can derive onto that same + // snapshot instead of building a new one, so a marker the resolver already + // stamped on it survives into the persisted row. + it("queueReviewParticipantRecoveryRun persists the caller-resolved responsible user and context snapshot, merged with the derived stage fields", async () => { + const companyId = await seedCompany(); + const finishingAgentId = await seedAgent({ companyId, name: "Finishing Agent" }); + const recoveryAgentId = await seedAgent({ companyId, name: "Recovery Agent" }); + const stageId = randomUUID(); + const issueId = await seedIssue({ companyId, assigneeAgentId: null, status: "in_review" }); + await db + .update(issues) + .set({ + executionState: { + status: "pending", + currentStageId: stageId, + currentStageIndex: 0, + currentStageType: "review", + currentParticipant: { type: "agent", agentId: recoveryAgentId }, + returnAssignee: null, + completedStageIds: [], + lastDecisionId: null, + lastDecisionOutcome: null, + }, + }) + .where(eq(issues.id, issueId)); + // A finishing run status other than the legacy-reconciliation set + // (failed, timed_out, interrupted, cancelled) reaches the module's own + // drain logic, so this call runs. + const finishingRunId = await seedRun({ + companyId, + agentId: finishingAgentId, + contextSnapshot: { issueId }, + status: "succeeded", + }); + + const adapter = createPostgresWakeQueueAdapter(db, stubDeps); + const result = await adapter.withIssueExecutionLock( + { companyId, runId: finishingRunId, now: new Date() }, + async (locked, ports) => { + // The caller resolves the responsible user against this exact + // object before this call, and the resolver can stamp a marker on + // it; simulate that stamp here, the same way the real resolver does. + const contextSnapshot: Record = { + issueId, + taskId: issueId, + wakeReason: "execution_review_participant_recovery", + retryReason: "execution_review_participant_recovery", + source: "issue.execution_review_recovery", + retryOfRunId: finishingRunId, + reviewRecoveryInstruction: "Submit the review decision now.", + executionIdentityCause: "company_default", + }; + const run = await ports.transaction.queueReviewParticipantRecoveryRun({ + companyId, + issue: locked.primaryIssue, + finishingRun: locked.run, + recoveryAgent: { id: recoveryAgentId, companyId, name: "Recovery Agent", invokable: true }, + contextSnapshot, + responsibleUserId: "responsible-user", + sessionBefore: null, + now: new Date(), + }); + return { outcome: { kind: "queued_review_participant_recovery" as const, run }, postCommitEffects: [] }; + }, + ); + expect(result.outcome.kind).toBe("queued_review_participant_recovery"); + + const runRow = (await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.companyId, companyId))).find( + (row) => row.agentId === recoveryAgentId, + ); + expect(runRow?.responsibleUserId).toBe("responsible-user"); + expect(runRow?.contextSnapshot).toMatchObject({ + issueId, + retryOfRunId: finishingRunId, + // The marker the resolver stamped survives the merge. + executionIdentityCause: "company_default", + // The two fields only this adapter can derive. + currentStageId: stageId, + currentStageType: "review", + }); + }); + it("locks the context issue and every sibling issue in id order, and two concurrent releases do not deadlock", async () => { const companyId = await seedCompany(); const agentId = await seedAgent({ companyId }); diff --git a/server/src/modules/wake-queue/adapters/postgres.ts b/server/src/modules/wake-queue/adapters/postgres.ts index f9b650bf9e..a27e544adb 100644 --- a/server/src/modules/wake-queue/adapters/postgres.ts +++ b/server/src/modules/wake-queue/adapters/postgres.ts @@ -514,7 +514,16 @@ function buildTransaction(tx: Db, deps: WakeQueuePostgresAdapterDeps, db: Db, ru ); }, - async queueReviewParticipantRecoveryRun({ companyId, issue, finishingRun, recoveryAgent, sessionBefore, now }) { + async queueReviewParticipantRecoveryRun({ + companyId, + issue, + finishingRun, + recoveryAgent, + contextSnapshot, + responsibleUserId, + sessionBefore, + now, + }) { const executionState = parseIssueExecutionState(issue.executionState); const wakeupRequest = await tx .insert(agentWakeupRequests) @@ -542,37 +551,15 @@ function buildTransaction(tx: Db, deps: WakeQueuePostgresAdapterDeps, db: Db, ru .returning() .then((rows) => rows[0]); - // This insert does not set responsibleUserId. When claimQueuedRun claims - // this row, it resolves a responsible user. It writes that value in the - // same update that moves the run from "queued" to "running". But - // claimQueuedRun can also cancel a queued recovery run before it - // resolves that value, and a cancelled row then keeps a null - // responsible user permanently. A claimed row still needs - // initializeRunIdentity to overwrite the value again, from the run - // identity chain; if execution ends before that overwrite runs, the - // claimed value stays. Readers can observe a null value here. - // runsForIssue projects the column with no status filter, so the issue - // run ledger can show a recovery run with no responsible user. That - // projection displays attribution and makes no authorization decision. - // The issue-thread interaction attribution check is an authorization - // read that can also observe a null value. It looks up a - // caller-supplied run id with no status filter, and it compares this - // column against the responsible user of the caller. It makes that - // comparison only when the caller carries a responsible user, and a - // null value never equals one, so the check denies. The audit feed can - // observe a cancelled run: claimQueuedRun cancels a queued run when an - // active subtree pause hold holds the issue, and it writes an activity - // log event for that cancelled run. agentActionAuditService prefers the - // responsible user that the activity log row carries. - // resolveResponsibleUserIdForActivity sets that value, and it finds no - // responsible user on the cancelled run. It falls back to the issue, - // then to the agent API key, then to the company default. - // The immediate-recovery writer takes an already-resolved - // responsibleUserId as an input parameter, because its caller must - // resolve one before it can call that writer. This writer's input - // carries no such parameter. It leaves resolution to claimQueuedRun, - // which resolves and writes a responsible user for every queued run, - // the moment it claims one. + // The caller already resolved the responsible user against this same + // `contextSnapshot` object, so this insert persists that object as + // given. The two stage fields below only exist in this adapter, so + // this is the one place that can add them; the merge writes them onto + // the same object instead of building a new one, so it does not + // disturb any key the resolver already set. + contextSnapshot.currentStageId = executionState?.currentStageId ?? null; + contextSnapshot.currentStageType = executionState?.currentStageType ?? null; + const queuedRun = await tx .insert(heartbeatRuns) .values({ @@ -582,21 +569,8 @@ function buildTransaction(tx: Db, deps: WakeQueuePostgresAdapterDeps, db: Db, ru triggerDetail: "system", status: "queued", wakeupRequestId: wakeupRequest.id, - contextSnapshot: withRecoveryContext( - { - issueId: issue.id, - taskId: issue.id, - wakeReason: EXECUTION_REVIEW_PARTICIPANT_RECOVERY_RETRY_REASON, - retryReason: EXECUTION_REVIEW_PARTICIPANT_RECOVERY_RETRY_REASON, - source: "issue.execution_review_recovery", - retryOfRunId: finishingRun.id, - currentStageId: executionState?.currentStageId ?? null, - currentStageType: executionState?.currentStageType ?? null, - reviewRecoveryInstruction: - "The previous reviewer run ended while this execution-review stage was still pending. Submit the review decision now, or mark the issue blocked with the exact unblock action.", - }, - "normal_model", - ), + contextSnapshot, + responsibleUserId, sessionIdBefore: sessionBefore, retryOfRunId: finishingRun.id, updatedAt: now, diff --git a/server/src/modules/wake-queue/application/ports.ts b/server/src/modules/wake-queue/application/ports.ts index 85a085a94b..82be28315f 100644 --- a/server/src/modules/wake-queue/application/ports.ts +++ b/server/src/modules/wake-queue/application/ports.ts @@ -183,11 +183,17 @@ export interface WakeQueueTransaction { isAutomaticRecoverySuppressedByPauseHold(input: { companyId: string; issueId: string }): Promise; /** Deny-only facts from the exact finishing run and its durable chat wake owner. */ isImmediateRecoverySourceBlocked(input: { companyId: string; runId: string }): Promise; + /** + * Queues the run with the context snapshot and the responsible user the + * caller already resolved. + */ queueReviewParticipantRecoveryRun(input: { companyId: string; issue: IssueSnapshot; finishingRun: RunSnapshot; recoveryAgent: InvokableAgentSnapshot; + contextSnapshot: Record; + responsibleUserId: string; sessionBefore: string | null; now: Date; }): Promise; diff --git a/server/src/modules/wake-queue/application/use-cases.test.ts b/server/src/modules/wake-queue/application/use-cases.test.ts index 3c4e355b1f..6729c0f0d4 100644 --- a/server/src/modules/wake-queue/application/use-cases.test.ts +++ b/server/src/modules/wake-queue/application/use-cases.test.ts @@ -56,6 +56,24 @@ const AGENT: InvokableAgentSnapshot = { invokable: true, }; +// A run and its issue in the shape the release-recovery tail needs to reach +// the review-participant recovery path: the issue is in review with no +// assigned user, the run's agent is the issue's pending stage participant, +// and the run carries a wake reason the eligibility check recognizes. +const REVIEW_PARTICIPANT_RUN: RunSnapshot = { + ...RUN, + contextSnapshot: { wakeReason: "execution_review_requested" }, +}; + +const REVIEW_PARTICIPANT_ISSUE: IssueSnapshot = { + ...ISSUE, + status: "in_review", + executionState: { + status: "pending", + currentParticipant: { type: "agent", agentId: RUN.agentId }, + }, +}; + function wakeCandidate(overrides: Partial = {}): DeferredWakeCandidate { return { id: overrides.id ?? "wake-1", @@ -136,6 +154,15 @@ function createFakeIssueLock(host: WakeQueueHost, transaction: WakeQueueTransact }; } +function createReviewParticipantIssueLock(host: WakeQueueHost, transaction: WakeQueueTransaction): IssueLockWriter { + return { + withIssueExecutionLock: vi.fn(async (_input, fn) => { + const result = await fn({ primaryIssue: REVIEW_PARTICIPANT_ISSUE, run: REVIEW_PARTICIPANT_RUN }, { host, transaction }); + return { ...result, run: REVIEW_PARTICIPANT_RUN }; + }), + }; +} + function createFakeRecovery(): RecoveryEscalationPort { return { escalateStrandedAssignedIssue: vi.fn(async () => {}), @@ -485,6 +512,68 @@ describe("releaseIssueExecution", () => { expect(transaction.queueImmediateRecoveryRun).not.toHaveBeenCalled(); }); + it("resolves the responsible user for a review-participant recovery run before queuing it, and keeps the marker the resolver stamped", async () => { + const resolveResponsibleUserId = vi.fn(async (input: Parameters[0]) => { + // Mirrors the real resolver: it mutates the same context snapshot + // object it received, stamping this marker on the company-default + // fallback path. + input.contextSnapshot.executionIdentityCause = "company_default"; + return "resolved-user"; + }); + const queueReviewParticipantRecoveryRun = vi.fn( + async (_input: Parameters[0]) => runSummary("review-recovery"), + ); + const transaction = createFakeTransaction({ queueReviewParticipantRecoveryRun }); + const host = createFakeHost({ resolveResponsibleUserId }); + const issueLock = createReviewParticipantIssueLock(host, transaction); + const releaseIssueExecution = createReleaseIssueExecution({ issueLock, recovery: createFakeRecovery() }); + + const result = await releaseIssueExecution({ companyId: "company-1", runId: "run-1", now: new Date() }); + + expect(result.outcome.kind).toBe("queued_review_participant_recovery"); + expect(resolveResponsibleUserId).toHaveBeenCalledTimes(1); + const resolveCall = resolveResponsibleUserId.mock.calls[0]![0]; + expect(resolveCall.requestedByActorType).toBe("system"); + expect(resolveCall.requestedByActorId).toBeNull(); + expect(resolveCall.source).toBe("automation"); + expect(resolveCall.triggerDetail).toBe("system"); + expect(resolveCall.existingRunResponsibleUserId).toBe(REVIEW_PARTICIPANT_RUN.responsibleUserId); + expect(resolveCall.contextSnapshot).toEqual({ + issueId: REVIEW_PARTICIPANT_ISSUE.id, + taskId: REVIEW_PARTICIPANT_ISSUE.id, + wakeReason: "execution_review_participant_recovery", + retryReason: "execution_review_participant_recovery", + source: "issue.execution_review_recovery", + retryOfRunId: REVIEW_PARTICIPANT_RUN.id, + reviewRecoveryInstruction: + "The previous reviewer run ended while this execution-review stage was still pending. Submit the review decision now, or mark the issue blocked with the exact unblock action.", + executionIdentityCause: "company_default", + }); + + expect(queueReviewParticipantRecoveryRun).toHaveBeenCalledTimes(1); + const queueCall = queueReviewParticipantRecoveryRun.mock.calls[0]![0]; + expect(queueCall.responsibleUserId).toBe("resolved-user"); + // The object the caller hands the port is the same object it handed the resolver. + expect(queueCall.contextSnapshot).toBe(resolveCall.contextSnapshot); + // The marker the resolver stamped survives onto the persisted call. + expect((queueCall.contextSnapshot as Record).executionIdentityCause).toBe("company_default"); + }); + + it("throws WakeQueueApplicationError with code responsible_user_unresolved for a review-participant recovery run, without queuing it", async () => { + const transaction = createFakeTransaction(); + const host = createFakeHost({ resolveResponsibleUserId: vi.fn(async () => null) }); + const issueLock = createReviewParticipantIssueLock(host, transaction); + const releaseIssueExecution = createReleaseIssueExecution({ issueLock, recovery: createFakeRecovery() }); + + await expect( + releaseIssueExecution({ companyId: "company-1", runId: "run-1", now: new Date() }), + ).rejects.toMatchObject({ + constructor: WakeQueueApplicationError, + code: "responsible_user_unresolved", + }); + expect(transaction.queueReviewParticipantRecoveryRun).not.toHaveBeenCalled(); + }); + it("escalates through the recovery port for a blocked outcome, after the transaction resolves", async () => { const transaction = createFakeTransaction({ findNextDeferredWake: vi.fn(async () => null), diff --git a/server/src/modules/wake-queue/application/use-cases.ts b/server/src/modules/wake-queue/application/use-cases.ts index f4a2a9a666..427b9f95f0 100644 --- a/server/src/modules/wake-queue/application/use-cases.ts +++ b/server/src/modules/wake-queue/application/use-cases.ts @@ -563,11 +563,51 @@ async function runReleaseRecoveryTail( }); if (decision.kind === "queue_review_participant_recovery") { + // Resolve the responsible user here, in the application layer, before + // the transaction port queues the run — the same order the immediate + // recovery path below uses. + const reviewParticipantContextSnapshot: Record = { + issueId: issue.id, + taskId: issue.id, + wakeReason: EXECUTION_REVIEW_PARTICIPANT_RECOVERY_RETRY_REASON, + retryReason: EXECUTION_REVIEW_PARTICIPANT_RECOVERY_RETRY_REASON, + source: "issue.execution_review_recovery", + retryOfRunId: run.id, + reviewRecoveryInstruction: + "The previous reviewer run ended while this execution-review stage was still pending. Submit the review decision now, or mark the issue blocked with the exact unblock action.", + }; + + const reviewParticipantResponsibleUserId = await resolveResponsibleUserForQueuedRun(host, { + companyId: issue.companyId, + contextSnapshot: reviewParticipantContextSnapshot, + issue, + requestedByActorType: "system", + requestedByActorId: null, + source: "automation", + triggerDetail: "system", + existingRunResponsibleUserId: run.responsibleUserId, + }); + if (!reviewParticipantResponsibleUserId) { + throw new WakeQueueApplicationError( + "responsible_user_unresolved", + "Unable to resolve responsible user for review-participant recovery heartbeat run", + { + runId: run.id, + agentId: recoveryAgent.id, + companyId: issue.companyId, + issueId: issue.id, + wakeReason: EXECUTION_REVIEW_PARTICIPANT_RECOVERY_RETRY_REASON, + }, + ); + } + const queuedRun = await transaction.queueReviewParticipantRecoveryRun({ companyId: issue.companyId, issue, finishingRun: run, recoveryAgent, + contextSnapshot: reviewParticipantContextSnapshot, + responsibleUserId: reviewParticipantResponsibleUserId, sessionBefore, now: input.now, });