diff --git a/packages/plugins/paperclip-plugin-escalation/src/worker.ts b/packages/plugins/paperclip-plugin-escalation/src/worker.ts index b09f224b70b5..9af8adcbf8a8 100644 --- a/packages/plugins/paperclip-plugin-escalation/src/worker.ts +++ b/packages/plugins/paperclip-plugin-escalation/src/worker.ts @@ -57,6 +57,23 @@ async function guard(ctx: PluginContext, what: string, run: () => Promise) const plugin = definePlugin({ async setup(ctx) { + const pendingIssueUpdates = new Map>(); + + /** + * Delivery can overlap for one task. Keep the read-modify-write escalation + * transition in order so concurrent reviewer returns cannot escalate twice. + */ + const serializeIssueUpdate = async (key: string, run: () => Promise): Promise => { + const previous = pendingIssueUpdates.get(key) ?? Promise.resolve(); + const next = previous.catch(() => undefined).then(run); + pendingIssueUpdates.set(key, next); + try { + await next; + } finally { + if (pendingIssueUpdates.get(key) === next) pendingIssueUpdates.delete(key); + } + }; + /** Parks the task for a human and tells other plugins about it. */ const escalate = async ( event: PluginEvent, @@ -94,41 +111,43 @@ const plugin = definePlugin({ guard(ctx, "issue.updated", async () => { const issueId = event.entityId; if (!issueId) return; - const config = await readConfig(ctx, event.companyId); - if (!config) return; - - const issue = await ctx.issues.get(issueId, event.companyId); - if (!issue) return; - - const stored = await ctx.state.get({ ...scope(issueId), stateKey: STATE_KEYS.lastStatus }); - const previous = typeof stored === "string" ? (stored as IssueStatus) : null; - if (previous === issue.status) return; - - await ctx.state.set({ ...scope(issueId), stateKey: STATE_KEYS.lastStatus }, issue.status); - - if (isTerminal(issue.status)) { - await clearCounters(ctx, issueId); - return; - } - - if (!isReviewReturn(previous, issue.status)) return; - - const alreadyEscalated = await ctx.state.get({ - ...scope(issueId), - stateKey: STATE_KEYS.escalated, + await serializeIssueUpdate(`${event.companyId}\u0000${issueId}`, async () => { + const config = await readConfig(ctx, event.companyId); + if (!config) return; + + const issue = await ctx.issues.get(issueId, event.companyId); + if (!issue) return; + + const stored = await ctx.state.get({ ...scope(issueId), stateKey: STATE_KEYS.lastStatus }); + const previous = typeof stored === "string" ? (stored as IssueStatus) : null; + if (previous === issue.status) return; + + await ctx.state.set({ ...scope(issueId), stateKey: STATE_KEYS.lastStatus }, issue.status); + + if (isTerminal(issue.status)) { + await clearCounters(ctx, issueId); + return; + } + + if (!isReviewReturn(previous, issue.status)) return; + + const alreadyEscalated = await ctx.state.get({ + ...scope(issueId), + stateKey: STATE_KEYS.escalated, + }); + if (alreadyEscalated === true) return; + + const reviewReturns = (await readNumber(ctx, issueId, STATE_KEYS.reviewReturns)) + 1; + await ctx.state.set({ ...scope(issueId), stateKey: STATE_KEYS.reviewReturns }, reviewReturns); + + const counters: Counters = { + gateFailures: await readNumber(ctx, issueId, STATE_KEYS.gateFailures), + reviewReturns, + }; + if (shouldEscalate(counters, config.reviewReturnThreshold)) { + await escalate(event, counters, config.reviewReturnThreshold); + } }); - if (alreadyEscalated === true) return; - - const reviewReturns = (await readNumber(ctx, issueId, STATE_KEYS.reviewReturns)) + 1; - await ctx.state.set({ ...scope(issueId), stateKey: STATE_KEYS.reviewReturns }, reviewReturns); - - const counters: Counters = { - gateFailures: await readNumber(ctx, issueId, STATE_KEYS.gateFailures), - reviewReturns, - }; - if (shouldEscalate(counters, config.reviewReturnThreshold)) { - await escalate(event, counters, config.reviewReturnThreshold); - } }), ); diff --git a/packages/plugins/paperclip-plugin-escalation/tests/plugin.spec.ts b/packages/plugins/paperclip-plugin-escalation/tests/plugin.spec.ts index 91a3a3f9a4fe..af5b37c9316c 100644 --- a/packages/plugins/paperclip-plugin-escalation/tests/plugin.spec.ts +++ b/packages/plugins/paperclip-plugin-escalation/tests/plugin.spec.ts @@ -147,6 +147,42 @@ describe("escalation behaviour", () => { expect(emitted).toHaveLength(1); }); + it("serializes concurrent reviewer returns at the escalation threshold", async () => { + const { harness, comments, emitted } = await setup(); + await harness.ctx.state.set( + { scopeKind: "issue", scopeId: ISSUE_ID, stateKey: STATE_KEYS.lastStatus }, + "in_review", + ); + await harness.ctx.state.set( + { scopeKind: "issue", scopeId: ISSUE_ID, stateKey: STATE_KEYS.reviewReturns }, + 2, + ); + + const originalUpdate = harness.ctx.issues.update; + let releaseUpdate: (() => void) | undefined; + let markUpdateStarted: (() => void) | undefined; + const updateStarted = new Promise((resolve) => { + markUpdateStarted = resolve; + }); + const updateMayFinish = new Promise((resolve) => { + releaseUpdate = resolve; + }); + harness.ctx.issues.update = (async (...args) => { + markUpdateStarted?.(); + await updateMayFinish; + return originalUpdate(...args); + }) as typeof harness.ctx.issues.update; + + const first = harness.emit("issue.updated", {}, { entityId: ISSUE_ID, companyId: COMPANY_ID }); + await updateStarted; + const second = harness.emit("issue.updated", {}, { entityId: ISSUE_ID, companyId: COMPANY_ID }); + releaseUpdate?.(); + await Promise.all([first, second]); + + expect(comments).toHaveLength(1); + expect(emitted).toHaveLength(1); + }); + it("respects a configured threshold", async () => { const { emitted, reviewerReturn } = await setup({ reviewReturnThreshold: 1 }); diff --git a/server/src/__tests__/issue-execution-policy.test.ts b/server/src/__tests__/issue-execution-policy.test.ts index 8b9d560191c0..a75e97b311af 100644 --- a/server/src/__tests__/issue-execution-policy.test.ts +++ b/server/src/__tests__/issue-execution-policy.test.ts @@ -779,6 +779,25 @@ describe("issue execution policy transitions", () => { const policy = twoStagePolicy(); const reviewStageId = policy.stages[0].id; + it("requires executor evidence before routing an agent into review", () => { + expect(() => + applyIssueExecutionPolicyTransition({ + issue: { + status: "in_progress", + assigneeAgentId: coderAgentId, + assigneeUserId: null, + executionPolicy: policy, + executionState: null, + }, + policy, + requestedStatus: "done", + requestedAssigneePatch: {}, + actor: { agentId: coderAgentId }, + commentBody: " ", + }), + ).toThrow("requires a comment"); + }); + it("approval without comment throws", () => { expect(() => applyIssueExecutionPolicyTransition({ @@ -1721,6 +1740,7 @@ describe("issue execution policy transitions", () => { requestedStatus: "done", requestedAssigneePatch: {}, actor: { agentId: coderAgentId }, + commentBody: "Deployment monitor completed successfully.", }); expect(result.patch.executionPolicy).toBeNull(); diff --git a/server/src/services/issue-execution-policy.ts b/server/src/services/issue-execution-policy.ts index 859dc5fe1e31..859687bc2be0 100644 --- a/server/src/services/issue-execution-policy.ts +++ b/server/src/services/issue-execution-policy.ts @@ -971,6 +971,15 @@ function applyIssueExecutionStageTransition(input: TransitionInput): TransitionR return { patch }; } + if ( + requestedStatus === "done" && + input.policy.commentRequired && + input.actor.agentId && + !input.commentBody?.trim() + ) { + throw unprocessable("Completing an issue with an execution policy requires a comment"); + } + // A workflow whose execution already completed is terminal for approve/done: // closing the issue must not restart the chain at the first stage (#7893). if (requestedStatus === "done" && existingState?.status === COMPLETED_STATUS) {