mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-11 14:10:50 +02:00
refactor(server): simplify the wake-queue module (#13139)
## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work > - The server uses a deferred wake queue to release work at the correct time > - The wake queue had repeated decisions, helpers, queries, and recovery data > - This repetition made the module harder to read and left a drain invariant implicit > - This pull request moves pure decisions into the policy layer and simplifies the queue flow > - The benefit is a smaller, clearer module with the same behavior and direct test coverage ## Linked Issues or Issue Description Refs #13136 **Problem:** The deferred wake queue carried repeated logic across the application and database layers. **Expected behavior:** The queue keeps the same wake, release, recovery, and escalation behavior after the refactor. **Solution:** Move pure decisions into the policy layer, share repeated data and predicates, and state the drain invariant in the application layer. ## What Changed - Remove the unused database port method and carry the blocked release notice kind as a typed field. - Move four pre-drain decisions into pure policy functions with table-driven tests. - Resolve the responsible user once in the application layer. - Use the canonical helper for agent invokability checks. - Share string helpers and run predicates across the module. - Split the deferred-wake decision flow and share recovery facts with a discriminator. - Bound the drain loop and throw when it processes a wake identifier twice. - Rename queue ports to describe their behavior. - Share row-loading code between stranded-issue escalation adapters. - Restore the interaction-continuation integration test case. ## Verification - Run the three wake-queue module test files. They pass 61 cases locally. - Run `pnpm check:module-boundaries`. It passes locally. - Run the `server/` type-check and confirm that no error names the changed wake-queue module or `heartbeat.ts`. - Run continuous integration and confirm that `promotes an interaction continuation after removing a coalesced self-authored comment` passes. ## Risks The change refactors queue control flow and database adapter boundaries. The main risk is a behavior change in deferred wake release or recovery. The new policy tests and the restored integration case cover these paths. The local environment cannot run the integration test because the same dependency failure occurs on the base branch. ## Model Used OpenAI Codex, GPT-5, context window not exposed by the runtime, with reasoning, tool use, and code execution. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` / `Refs #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
6dd48cad43
commit
0bff1d5cb6
13 files changed
+1225
-667
No files matched your search
@@ -8,7 +8,7 @@ import type {
|
||||
RunSummary,
|
||||
} from "./types.js";
|
||||
|
||||
export type { InvokableAgentSnapshot, IssueSnapshot, RunSnapshot, RunSummary };
|
||||
export type { InvokableAgentSnapshot, IssueSnapshot, ReleaseRecoveryBlockedNoticeKind, RunSnapshot, RunSummary };
|
||||
|
||||
/** The primary issue a locked release resolves to, plus the finishing run the lock step already loaded. */
|
||||
export type LockedIssueExecution = {
|
||||
@@ -21,9 +21,12 @@ export type ReleaseTransactionResult = {
|
||||
postCommitEffects: PostCommitEffect[];
|
||||
};
|
||||
|
||||
/** Read-only lookups the release use case needs, each scoped to a company. */
|
||||
export interface WakeQueueReader {
|
||||
findInvokableAgent(input: { companyId: string; agentId: string }): Promise<InvokableAgentSnapshot | null>;
|
||||
/**
|
||||
* The three host callbacks the release use case needs. These members do
|
||||
* not run on the module's own transaction, which is why two of them take
|
||||
* the transaction-scoped issue snapshot instead of an issue id.
|
||||
*/
|
||||
export interface WakeQueueHost {
|
||||
/**
|
||||
* Takes the transaction-scoped issue snapshot, not an issue id, so this
|
||||
* port never re-reads the issue on a separate connection while the
|
||||
@@ -64,7 +67,7 @@ export type DeferredWakeCandidate = {
|
||||
reason: string | null;
|
||||
source: string | null;
|
||||
triggerDetail: string | null;
|
||||
requestedByActorType: string | null;
|
||||
requestedByActorType: "user" | "agent" | "system" | null;
|
||||
requestedByActorId: string | null;
|
||||
payload: Record<string, unknown>;
|
||||
/** The queued comment ids the wake's queued-comment context carries, already extracted from the payload. */
|
||||
@@ -94,9 +97,14 @@ export type PromoteDeferredWakeInput = {
|
||||
now: Date;
|
||||
};
|
||||
|
||||
/** The transaction-scoped write operations that drain and resolve the deferred-wake queue. */
|
||||
export interface WakeQueueWriter {
|
||||
claimNextDeferredWake(input: { companyId: string; issueId: string }): Promise<DeferredWakeCandidate | null>;
|
||||
/**
|
||||
* Every member is bound to the one transaction that `withIssueExecutionLock`
|
||||
* owns. The interface holds both reads and writes that drain and resolve
|
||||
* the deferred-wake queue.
|
||||
*/
|
||||
export interface WakeQueueTransaction {
|
||||
findInvokableAgent(input: { companyId: string; agentId: string }): Promise<InvokableAgentSnapshot | null>;
|
||||
findNextDeferredWake(input: { companyId: string; issueId: string }): Promise<DeferredWakeCandidate | null>;
|
||||
getQueuedCommentLiveness(input: {
|
||||
companyId: string;
|
||||
issueId: string;
|
||||
@@ -115,7 +123,7 @@ export interface WakeQueueWriter {
|
||||
normalizeDeferredWakeCommentIds(input: {
|
||||
companyId: string;
|
||||
wakeId: string;
|
||||
/** The wake's current payload, as already read by `claimNextDeferredWake`, used as the rewrite base. */
|
||||
/** The wake's current payload, as already read by `findNextDeferredWake`, used as the rewrite base. */
|
||||
payload: Record<string, unknown>;
|
||||
liveCommentIds: string[];
|
||||
now: Date;
|
||||
@@ -169,12 +177,6 @@ export interface WakeQueueWriter {
|
||||
/** An open, non-hidden issue that still lists this issue as a `blocks` predecessor. */
|
||||
hasExplicitBlockerPath(input: { companyId: string; issueId: string }): Promise<boolean>;
|
||||
isAutomaticRecoverySuppressedByPauseHold(input: { companyId: string; issueId: string }): Promise<boolean>;
|
||||
/** Builds the stranded-recovery notice content for a `blocked` outcome; pure formatting, kept behind the writer so `services/recovery/stranded-notice` stays out of the application layer. */
|
||||
buildBlockedRecoveryNotice(input: {
|
||||
noticeKind: ReleaseRecoveryBlockedNoticeKind;
|
||||
issueStatus: "todo" | "in_progress";
|
||||
finishingRun: RunSnapshot;
|
||||
}): Promise<{ notice: Record<string, unknown>; recoveryCause: string | null }>;
|
||||
queueReviewParticipantRecoveryRun(input: {
|
||||
companyId: string;
|
||||
issue: IssueSnapshot;
|
||||
@@ -184,16 +186,18 @@ export interface WakeQueueWriter {
|
||||
now: Date;
|
||||
}): Promise<RunSummary>;
|
||||
/**
|
||||
* Builds the recovery context snapshot, resolves the responsible user
|
||||
* from it, and queues the run. Throws `WakeQueueApplicationError` with
|
||||
* code `responsible_user_unresolved` when no responsible user resolves,
|
||||
* without queuing anything.
|
||||
* Queues the run with the context snapshot and the responsible user the
|
||||
* caller already resolved.
|
||||
*/
|
||||
queueImmediateRecoveryRun(input: {
|
||||
companyId: string;
|
||||
issue: IssueSnapshot;
|
||||
finishingRun: RunSnapshot;
|
||||
recoveryAgent: InvokableAgentSnapshot;
|
||||
/** The wakeup request's reason and the run's context-snapshot wakeReason; the caller derives it from the issue status. */
|
||||
reason: string;
|
||||
contextSnapshot: Record<string, unknown>;
|
||||
responsibleUserId: string;
|
||||
sessionBefore: string | null;
|
||||
now: Date;
|
||||
}): Promise<RunSummary>;
|
||||
@@ -207,16 +211,15 @@ export interface WakeQueueWriter {
|
||||
* (workspace-validation block, legacy reconciliation, a native-runtime
|
||||
* terminal failure), the adapter returns that outcome directly without
|
||||
* calling `fn`. Otherwise it calls `fn` with the locked issue and run, and
|
||||
* with `reader`/`writer` ports bound to the same transaction, so every
|
||||
* call `fn` makes through them participates in the one transaction this
|
||||
* method owns.
|
||||
* with `host`/`transaction` ports, so every call `fn` makes through the
|
||||
* transaction port participates in the one transaction this method owns.
|
||||
*/
|
||||
export interface IssueLockWriter {
|
||||
withIssueExecutionLock(
|
||||
input: { companyId: string; runId: string; now: Date },
|
||||
fn: (
|
||||
locked: LockedIssueExecution,
|
||||
ports: { reader: WakeQueueReader; writer: WakeQueueWriter },
|
||||
ports: { host: WakeQueueHost; transaction: WakeQueueTransaction },
|
||||
) => Promise<ReleaseTransactionResult>,
|
||||
): Promise<ReleaseTransactionResult & { run: RunSnapshot }>;
|
||||
}
|
||||
@@ -225,8 +228,7 @@ export type StrandedAssignedIssueEscalationInput = {
|
||||
issue: IssueSnapshot;
|
||||
previousStatus: "todo" | "in_progress" | "in_review";
|
||||
latestRun: RunSnapshot;
|
||||
notice: Record<string, unknown>;
|
||||
recoveryCause: string | null;
|
||||
noticeKind: ReleaseRecoveryBlockedNoticeKind;
|
||||
};
|
||||
|
||||
export type StrandedRecoveryInPlaceEscalationInput = {
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
import type { ReleaseRecoveryBlockedNoticeKind } from "../domain/policy.js";
|
||||
|
||||
export type RunSummary = {
|
||||
id: string;
|
||||
companyId: string;
|
||||
@@ -76,8 +78,7 @@ export type ReleaseOutcome =
|
||||
kind: "blocked";
|
||||
issue: IssueSnapshot;
|
||||
previousStatus: "todo" | "in_progress" | "in_review";
|
||||
notice: Record<string, unknown>;
|
||||
recoveryCause: string | null;
|
||||
noticeKind: ReleaseRecoveryBlockedNoticeKind;
|
||||
}
|
||||
| {
|
||||
kind: "blocked_recovery_in_place";
|
||||
@@ -85,7 +86,7 @@ export type ReleaseOutcome =
|
||||
previousStatus: "todo" | "in_progress" | "in_review";
|
||||
};
|
||||
|
||||
export type WakeQueueApplicationErrorCode = "responsible_user_unresolved";
|
||||
export type WakeQueueApplicationErrorCode = "responsible_user_unresolved" | "deferred_wake_not_advanced";
|
||||
|
||||
export class WakeQueueApplicationError extends Error {
|
||||
constructor(
|
||||
|
||||
@@ -10,8 +10,8 @@ import type {
|
||||
RecoveryEscalationPort,
|
||||
RunSnapshot,
|
||||
RunSummary,
|
||||
WakeQueueReader,
|
||||
WakeQueueWriter,
|
||||
WakeQueueHost,
|
||||
WakeQueueTransaction,
|
||||
} from "./ports.js";
|
||||
|
||||
const RUN: RunSnapshot = {
|
||||
@@ -29,7 +29,7 @@ const RUN: RunSnapshot = {
|
||||
const ISSUE: IssueSnapshot = {
|
||||
id: "issue-1",
|
||||
companyId: "company-1",
|
||||
identifier: "PAP-1",
|
||||
identifier: "ISSUE-1",
|
||||
status: "in_progress",
|
||||
assigneeAgentId: "finishing-agent",
|
||||
assigneeUserId: null,
|
||||
@@ -81,9 +81,8 @@ function runSummary(id: string): RunSummary {
|
||||
};
|
||||
}
|
||||
|
||||
function createFakeReader(overrides: Partial<WakeQueueReader> = {}): WakeQueueReader {
|
||||
function createFakeHost(overrides: Partial<WakeQueueHost> = {}): WakeQueueHost {
|
||||
return {
|
||||
findInvokableAgent: vi.fn(async () => AGENT),
|
||||
resolveResponsibleUserId: vi.fn(async () => "user-1"),
|
||||
getRoutineEnv: vi.fn(async () => ({ routineId: null, env: null, responsibleUserId: null })),
|
||||
resolveSessionBeforeForWakeup: vi.fn(async () => null),
|
||||
@@ -91,9 +90,10 @@ function createFakeReader(overrides: Partial<WakeQueueReader> = {}): WakeQueueRe
|
||||
};
|
||||
}
|
||||
|
||||
function createFakeWriter(overrides: Partial<WakeQueueWriter> = {}): WakeQueueWriter {
|
||||
function createFakeTransaction(overrides: Partial<WakeQueueTransaction> = {}): WakeQueueTransaction {
|
||||
return {
|
||||
claimNextDeferredWake: vi.fn(async () => null),
|
||||
findInvokableAgent: vi.fn(async () => AGENT),
|
||||
findNextDeferredWake: vi.fn(async () => null),
|
||||
getQueuedCommentLiveness: vi.fn(async () => ({ liveNonSelfCommentIds: [], containedSelfAuthoredComment: false })),
|
||||
cancelDeferredWake: vi.fn(async () => true),
|
||||
normalizeDeferredWakeCommentIds: vi.fn(async (input) => wakeCandidate({ id: input.wakeId, queuedCommentIds: input.liveCommentIds })),
|
||||
@@ -114,17 +114,16 @@ function createFakeWriter(overrides: Partial<WakeQueueWriter> = {}): WakeQueueWr
|
||||
hasExistingExecutionPath: vi.fn(async () => false),
|
||||
hasExplicitBlockerPath: vi.fn(async () => false),
|
||||
isAutomaticRecoverySuppressedByPauseHold: vi.fn(async () => false),
|
||||
buildBlockedRecoveryNotice: vi.fn(async () => ({ notice: {}, recoveryCause: null })),
|
||||
queueReviewParticipantRecoveryRun: vi.fn(async () => runSummary("review-recovery")),
|
||||
queueImmediateRecoveryRun: vi.fn(async () => runSummary("immediate-recovery")),
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
function createFakeIssueLock(reader: WakeQueueReader, writer: WakeQueueWriter): IssueLockWriter {
|
||||
function createFakeIssueLock(host: WakeQueueHost, transaction: WakeQueueTransaction): IssueLockWriter {
|
||||
return {
|
||||
withIssueExecutionLock: vi.fn(async (_input, fn) => {
|
||||
const result = await fn({ primaryIssue: ISSUE, run: RUN }, { reader, writer });
|
||||
const result = await fn({ primaryIssue: ISSUE, run: RUN }, { host, transaction });
|
||||
return { ...result, run: RUN };
|
||||
}),
|
||||
};
|
||||
@@ -141,35 +140,36 @@ describe("releaseIssueExecution", () => {
|
||||
it("processes the deferred wakes in requestedAt order", async () => {
|
||||
const claimOrder: string[] = [];
|
||||
const queue = [wakeCandidate({ id: "wake-earliest" }), wakeCandidate({ id: "wake-latest" })];
|
||||
const writer = createFakeWriter({
|
||||
claimNextDeferredWake: vi.fn(async () => {
|
||||
const transaction = createFakeTransaction({
|
||||
findNextDeferredWake: vi.fn(async () => {
|
||||
const next = queue.shift() ?? null;
|
||||
if (next) claimOrder.push(next.id);
|
||||
return next;
|
||||
}),
|
||||
// Every wake fails invokability so the loop keeps draining without promoting.
|
||||
findInvokableAgent: vi.fn(async () => null),
|
||||
});
|
||||
// Every wake fails invokability so the loop keeps draining without promoting.
|
||||
const reader = createFakeReader({ findInvokableAgent: vi.fn(async () => null) });
|
||||
const issueLock = createFakeIssueLock(reader, writer);
|
||||
const host = createFakeHost();
|
||||
const issueLock = createFakeIssueLock(host, transaction);
|
||||
const releaseIssueExecution = createReleaseIssueExecution({ issueLock, recovery: createFakeRecovery() });
|
||||
|
||||
await releaseIssueExecution({ companyId: "company-1", runId: "run-1", now: new Date() });
|
||||
|
||||
expect(claimOrder).toEqual(["wake-earliest", "wake-latest"]);
|
||||
expect(writer.failDeferredWake).toHaveBeenCalledTimes(2);
|
||||
expect(transaction.failDeferredWake).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
|
||||
it("stops the loop after the first promotion", async () => {
|
||||
const claimNextDeferredWake = vi.fn(async () => wakeCandidate({ id: "wake-promotes" }));
|
||||
const writer = createFakeWriter({ claimNextDeferredWake });
|
||||
const reader = createFakeReader();
|
||||
const issueLock = createFakeIssueLock(reader, writer);
|
||||
const findNextDeferredWake = vi.fn(async () => wakeCandidate({ id: "wake-promotes" }));
|
||||
const transaction = createFakeTransaction({ findNextDeferredWake });
|
||||
const host = createFakeHost();
|
||||
const issueLock = createFakeIssueLock(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("promoted");
|
||||
expect(claimNextDeferredWake).toHaveBeenCalledTimes(1);
|
||||
expect(findNextDeferredWake).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("continues the loop after a cancel outcome, a fail outcome, and a normalize outcome, then promotes", async () => {
|
||||
@@ -181,7 +181,7 @@ describe("releaseIssueExecution", () => {
|
||||
// normalize: queued comments differ from the live set, then promotes.
|
||||
wakeCandidate({ id: "wake-normalize", queuedCommentIds: ["c1", "c2"] }),
|
||||
];
|
||||
const claimNextDeferredWake = vi.fn(async () => queue.shift() ?? null);
|
||||
const findNextDeferredWake = vi.fn(async () => queue.shift() ?? null);
|
||||
const findInvokableAgent = vi.fn(async (input: { agentId: string }) =>
|
||||
input.agentId === "deferred-agent" ? AGENT : null,
|
||||
);
|
||||
@@ -190,24 +190,45 @@ describe("releaseIssueExecution", () => {
|
||||
? { liveNonSelfCommentIds: [], containedSelfAuthoredComment: false }
|
||||
: { liveNonSelfCommentIds: ["c2"], containedSelfAuthoredComment: false },
|
||||
);
|
||||
const writer = createFakeWriter({ claimNextDeferredWake, getQueuedCommentLiveness });
|
||||
const reader = createFakeReader({ findInvokableAgent });
|
||||
const issueLock = createFakeIssueLock(reader, writer);
|
||||
const transaction = createFakeTransaction({ findNextDeferredWake, findInvokableAgent, getQueuedCommentLiveness });
|
||||
const host = createFakeHost();
|
||||
const issueLock = createFakeIssueLock(host, transaction);
|
||||
const releaseIssueExecution = createReleaseIssueExecution({ issueLock, recovery: createFakeRecovery() });
|
||||
|
||||
const result = await releaseIssueExecution({ companyId: "company-1", runId: "run-1", now: new Date() });
|
||||
|
||||
expect(writer.cancelDeferredWake).toHaveBeenCalledTimes(1);
|
||||
expect(writer.failDeferredWake).toHaveBeenCalledTimes(1);
|
||||
expect(writer.normalizeDeferredWakeCommentIds).toHaveBeenCalledTimes(1);
|
||||
expect(claimNextDeferredWake).toHaveBeenCalledTimes(3);
|
||||
expect(transaction.cancelDeferredWake).toHaveBeenCalledTimes(1);
|
||||
expect(transaction.failDeferredWake).toHaveBeenCalledTimes(1);
|
||||
expect(transaction.normalizeDeferredWakeCommentIds).toHaveBeenCalledTimes(1);
|
||||
expect(findNextDeferredWake).toHaveBeenCalledTimes(3);
|
||||
expect(result.outcome.kind).toBe("promoted");
|
||||
});
|
||||
|
||||
it("rejects with deferred_wake_not_advanced when the queue read returns the same wake id twice, instead of looping forever", async () => {
|
||||
// queuedCommentIds with no live comments and no independent continuation
|
||||
// routes to "cancel_empty", so the drain calls cancelDeferredWake and
|
||||
// discards its result, then reads the queue again for the same row.
|
||||
const repeatedCandidate = wakeCandidate({ id: "wake-repeat", queuedCommentIds: ["c1"] });
|
||||
const findNextDeferredWake = vi.fn(async () => repeatedCandidate);
|
||||
const cancelDeferredWake = vi.fn(async () => false);
|
||||
const transaction = createFakeTransaction({ findNextDeferredWake, cancelDeferredWake });
|
||||
const host = createFakeHost();
|
||||
const issueLock = createFakeIssueLock(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: "deferred_wake_not_advanced",
|
||||
});
|
||||
expect(findNextDeferredWake).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
|
||||
it("returns the post-commit effects as data without running them", async () => {
|
||||
const writer = createFakeWriter({ claimNextDeferredWake: vi.fn(async () => wakeCandidate()) });
|
||||
const reader = createFakeReader();
|
||||
const issueLock = createFakeIssueLock(reader, writer);
|
||||
const transaction = createFakeTransaction({ findNextDeferredWake: vi.fn(async () => wakeCandidate()) });
|
||||
const host = createFakeHost();
|
||||
const issueLock = createFakeIssueLock(host, transaction);
|
||||
const recovery = createFakeRecovery();
|
||||
const releaseIssueExecution = createReleaseIssueExecution({ issueLock, recovery });
|
||||
|
||||
@@ -220,8 +241,8 @@ describe("releaseIssueExecution", () => {
|
||||
|
||||
it("carries the deferred wake's raw issue, interaction, execution-stage, and accepted-plan context onto the promoted run, and clears only the rendered text projections", async () => {
|
||||
const finalizePromotedWake = vi.fn(async (input: PromoteDeferredWakeInput) => runSummary(input.wakeId));
|
||||
const writer = createFakeWriter({
|
||||
claimNextDeferredWake: vi.fn(async () =>
|
||||
const transaction = createFakeTransaction({
|
||||
findNextDeferredWake: vi.fn(async () =>
|
||||
wakeCandidate({
|
||||
deferredContextSeed: {
|
||||
issueId: ISSUE.id,
|
||||
@@ -239,8 +260,8 @@ describe("releaseIssueExecution", () => {
|
||||
),
|
||||
finalizePromotedWake,
|
||||
});
|
||||
const reader = createFakeReader();
|
||||
const issueLock = createFakeIssueLock(reader, writer);
|
||||
const host = createFakeHost();
|
||||
const issueLock = createFakeIssueLock(host, transaction);
|
||||
const releaseIssueExecution = createReleaseIssueExecution({ issueLock, recovery: createFakeRecovery() });
|
||||
|
||||
const result = await releaseIssueExecution({ companyId: "company-1", runId: "run-1", now: new Date() });
|
||||
@@ -275,14 +296,14 @@ describe("releaseIssueExecution", () => {
|
||||
}),
|
||||
wakeCandidate({ id: "wake-promotes" }),
|
||||
];
|
||||
const claimNextDeferredWake = vi.fn(async () => queue.shift() ?? null);
|
||||
const findNextDeferredWake = vi.fn(async () => queue.shift() ?? null);
|
||||
const claimDeferredWakeForPromotion = vi.fn(async ({ wakeId }: { wakeId: string }) => wakeId !== "wake-lost-race");
|
||||
const reopenIssue = vi.fn(async () => null);
|
||||
const writer = createFakeWriter({ claimNextDeferredWake, claimDeferredWakeForPromotion, reopenIssue });
|
||||
const reader = createFakeReader();
|
||||
const transaction = createFakeTransaction({ findNextDeferredWake, claimDeferredWakeForPromotion, reopenIssue });
|
||||
const host = createFakeHost();
|
||||
const issueLock: IssueLockWriter = {
|
||||
withIssueExecutionLock: vi.fn(async (_input, fn) => {
|
||||
const result = await fn({ primaryIssue: doneIssue, run: RUN }, { reader, writer });
|
||||
const result = await fn({ primaryIssue: doneIssue, run: RUN }, { host, transaction });
|
||||
return { ...result, run: RUN };
|
||||
}),
|
||||
};
|
||||
@@ -297,9 +318,9 @@ describe("releaseIssueExecution", () => {
|
||||
});
|
||||
|
||||
it("throws WakeQueueApplicationError with code responsible_user_unresolved when the responsible user cannot resolve", async () => {
|
||||
const writer = createFakeWriter({ claimNextDeferredWake: vi.fn(async () => wakeCandidate()) });
|
||||
const reader = createFakeReader({ resolveResponsibleUserId: vi.fn(async () => null) });
|
||||
const issueLock = createFakeIssueLock(reader, writer);
|
||||
const transaction = createFakeTransaction({ findNextDeferredWake: vi.fn(async () => wakeCandidate()) });
|
||||
const host = createFakeHost({ resolveResponsibleUserId: vi.fn(async () => null) });
|
||||
const issueLock = createFakeIssueLock(host, transaction);
|
||||
const releaseIssueExecution = createReleaseIssueExecution({ issueLock, recovery: createFakeRecovery() });
|
||||
|
||||
await expect(
|
||||
@@ -310,22 +331,77 @@ describe("releaseIssueExecution", () => {
|
||||
});
|
||||
});
|
||||
|
||||
it("resolves the responsible user for an immediate recovery run before queuing it", async () => {
|
||||
const resolveResponsibleUserId = vi.fn(
|
||||
async (_input: Parameters<WakeQueueHost["resolveResponsibleUserId"]>[0]) => "resolved-user",
|
||||
);
|
||||
const queueImmediateRecoveryRun = vi.fn(
|
||||
async (input: Parameters<WakeQueueTransaction["queueImmediateRecoveryRun"]>[0]) => runSummary("immediate-recovery"),
|
||||
);
|
||||
const transaction = createFakeTransaction({ queueImmediateRecoveryRun });
|
||||
const host = createFakeHost({ resolveResponsibleUserId });
|
||||
const issueLock = createFakeIssueLock(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_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(RUN.responsibleUserId);
|
||||
// ISSUE.status is "in_progress", so the stalled-continuation labels apply.
|
||||
expect(resolveCall.contextSnapshot).toEqual({
|
||||
issueId: ISSUE.id,
|
||||
taskId: ISSUE.id,
|
||||
wakeReason: "issue_continuation_needed",
|
||||
retryReason: "issue_continuation_needed",
|
||||
source: "issue.continuation_recovery",
|
||||
retryOfRunId: RUN.id,
|
||||
});
|
||||
|
||||
expect(queueImmediateRecoveryRun).toHaveBeenCalledTimes(1);
|
||||
const queueCall = queueImmediateRecoveryRun.mock.calls[0]![0];
|
||||
expect(queueCall.reason).toBe("issue_continuation_needed");
|
||||
expect(queueCall.responsibleUserId).toBe("resolved-user");
|
||||
expect(queueCall.contextSnapshot).toBe(resolveCall.contextSnapshot);
|
||||
});
|
||||
|
||||
it("throws WakeQueueApplicationError with code responsible_user_unresolved for a recovery run, without queuing it", async () => {
|
||||
const transaction = createFakeTransaction();
|
||||
const host = createFakeHost({ resolveResponsibleUserId: vi.fn(async () => null) });
|
||||
const issueLock = createFakeIssueLock(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.queueImmediateRecoveryRun).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("escalates through the recovery port for a blocked outcome, after the transaction resolves", async () => {
|
||||
const writer = createFakeWriter({
|
||||
claimNextDeferredWake: vi.fn(async () => null),
|
||||
const transaction = createFakeTransaction({
|
||||
findNextDeferredWake: vi.fn(async () => null),
|
||||
hasExistingExecutionPath: vi.fn(async () => false),
|
||||
isAutomaticRecoverySuppressedByPauseHold: vi.fn(async () => false),
|
||||
buildBlockedRecoveryNotice: vi.fn(async () => ({ notice: { kind: "immediate_execution_path" }, recoveryCause: "immediate_execution_path" })),
|
||||
// The recovery agent (the finishing run's own agent) is not invokable, which forces "blocked".
|
||||
findInvokableAgent: vi.fn(async () => null),
|
||||
});
|
||||
// The recovery agent (the finishing run's own agent) is not invokable, which forces "blocked".
|
||||
const reader = createFakeReader({ findInvokableAgent: vi.fn(async () => null) });
|
||||
const issueLock = createFakeIssueLock(reader, writer);
|
||||
const host = createFakeHost();
|
||||
const issueLock = createFakeIssueLock(host, transaction);
|
||||
const recovery = createFakeRecovery();
|
||||
const releaseIssueExecution = createReleaseIssueExecution({ issueLock, recovery });
|
||||
|
||||
const result = await releaseIssueExecution({ companyId: "company-1", runId: "run-1", now: new Date() });
|
||||
|
||||
expect(result.outcome.kind).toBe("blocked");
|
||||
expect(result.outcome.kind === "blocked" && result.outcome.noticeKind).toBe("immediate_execution_path");
|
||||
expect(recovery.escalateStrandedAssignedIssue).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
@@ -1,20 +1,33 @@
|
||||
import { enrichPromotedWakeContext } from "../domain/context.js";
|
||||
import { decideDeferredWake, decideReleaseRecovery } from "../domain/policy.js";
|
||||
import {
|
||||
decideQueuedCommentAction,
|
||||
decideReleaseRecovery,
|
||||
decideWakeOutcome,
|
||||
deriveImmediateRecoveryContextLabels,
|
||||
} from "../domain/policy.js";
|
||||
import {
|
||||
EXECUTION_REVIEW_PARTICIPANT_RECOVERY_RETRY_REASON,
|
||||
isConfigurationIncompleteFailedRun,
|
||||
isWorkspaceValidationFailedRun,
|
||||
readNonEmptyString,
|
||||
} from "../domain/values.js";
|
||||
import { withRecoveryContext } from "../../../services/recovery/status-only-context.js";
|
||||
import type {
|
||||
DeferredWakeCandidate,
|
||||
InvokableAgentSnapshot,
|
||||
IssueLockWriter,
|
||||
IssueSnapshot,
|
||||
LockedIssueExecution,
|
||||
RecoveryEscalationPort,
|
||||
ReleaseTransactionResult,
|
||||
RunSnapshot,
|
||||
WakeQueueReader,
|
||||
WakeQueueWriter,
|
||||
WakeQueueHost,
|
||||
WakeQueueTransaction,
|
||||
} from "./ports.js";
|
||||
import type { PostCommitEffect, ReleaseOutcome } from "./types.js";
|
||||
import { WakeQueueApplicationError } from "./types.js";
|
||||
|
||||
const ISSUE_DISPOSITION_REPAIR_RETRY_REASON = "issue_disposition_repair";
|
||||
const EXECUTION_REVIEW_PARTICIPANT_RECOVERY_RETRY_REASON = "execution_review_participant_recovery";
|
||||
const EXECUTION_REVIEW_PARTICIPANT_RECOVERY_WAKE_REASONS = new Set([
|
||||
"execution_review_requested",
|
||||
"execution_approval_requested",
|
||||
@@ -31,20 +44,6 @@ const UNSUCCESSFUL_HEARTBEAT_RUN_TERMINAL_STATUSES = new Set([
|
||||
"cancelled",
|
||||
]);
|
||||
const STRANDED_ISSUE_RECOVERY_ORIGIN_KIND = "stranded_issue_recovery";
|
||||
const WORKSPACE_VALIDATION_FAILURE_CODE = "workspace_validation_failed";
|
||||
const CONFIGURATION_INCOMPLETE_FAILURE_CODE = "configuration_incomplete";
|
||||
|
||||
function readNonEmptyString(value: unknown): string | null {
|
||||
return typeof value === "string" && value.trim().length > 0 ? value : null;
|
||||
}
|
||||
|
||||
function isWorkspaceValidationFailedRun(run: Pick<RunSnapshot, "errorCode">): boolean {
|
||||
return run.errorCode === WORKSPACE_VALIDATION_FAILURE_CODE;
|
||||
}
|
||||
|
||||
function isConfigurationIncompleteFailedRun(run: Pick<RunSnapshot, "errorCode">): boolean {
|
||||
return run.errorCode === CONFIGURATION_INCOMPLETE_FAILURE_CODE || run.errorCode === "model_not_found";
|
||||
}
|
||||
|
||||
function isExecutionReviewParticipantRecoveryRun(run: Pick<RunSnapshot, "contextSnapshot">): boolean {
|
||||
return readNonEmptyString(run.contextSnapshot.retryReason) === EXECUTION_REVIEW_PARTICIPANT_RECOVERY_RETRY_REASON;
|
||||
@@ -73,6 +72,39 @@ function currentAgentParticipant(issue: IssueSnapshot): { agentId: string } | nu
|
||||
return agentId ? { agentId } : null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolves the responsible user for a heartbeat run the module is about to
|
||||
* queue. The promote path and the immediate-recovery path both call this
|
||||
* one function; each still checks the result and throws its own error with
|
||||
* its own metadata when no responsible user resolves.
|
||||
*/
|
||||
async function resolveResponsibleUserForQueuedRun(
|
||||
host: WakeQueueHost,
|
||||
input: {
|
||||
companyId: string;
|
||||
contextSnapshot: Record<string, unknown>;
|
||||
issue: IssueSnapshot;
|
||||
requestedByActorType: "user" | "agent" | "system" | null;
|
||||
requestedByActorId: string | null;
|
||||
source: string;
|
||||
triggerDetail: string | null;
|
||||
existingRunResponsibleUserId: string | null;
|
||||
},
|
||||
): Promise<string | null> {
|
||||
const routineEnvContext = await host.getRoutineEnv({ companyId: input.companyId, issue: input.issue });
|
||||
return host.resolveResponsibleUserId({
|
||||
companyId: input.companyId,
|
||||
contextSnapshot: input.contextSnapshot,
|
||||
issue: input.issue,
|
||||
routineEnvContext,
|
||||
requestedByActorType: input.requestedByActorType,
|
||||
requestedByActorId: input.requestedByActorId,
|
||||
source: input.source,
|
||||
triggerDetail: input.triggerDetail,
|
||||
existingRunResponsibleUserId: input.existingRunResponsibleUserId,
|
||||
});
|
||||
}
|
||||
|
||||
export type ReleaseIssueExecutionInput = {
|
||||
companyId: string;
|
||||
runId: string;
|
||||
@@ -80,6 +112,8 @@ export type ReleaseIssueExecutionInput = {
|
||||
suppressImmediateRecovery?: boolean;
|
||||
};
|
||||
|
||||
type PauseHoldFacts = Awaited<ReturnType<WakeQueueTransaction["getPauseHoldFacts"]>>;
|
||||
|
||||
/**
|
||||
* Drains the deferred-wake queue for the issue a run just released, in
|
||||
* `requestedAt` order, promoting at most one wake. When the queue empties
|
||||
@@ -89,20 +123,35 @@ export type ReleaseIssueExecutionInput = {
|
||||
*/
|
||||
async function runReleaseDrain(
|
||||
locked: LockedIssueExecution,
|
||||
ports: { reader: WakeQueueReader; writer: WakeQueueWriter },
|
||||
ports: { host: WakeQueueHost; transaction: WakeQueueTransaction },
|
||||
input: ReleaseIssueExecutionInput,
|
||||
): Promise<ReleaseTransactionResult> {
|
||||
const { run } = locked;
|
||||
let issue = locked.primaryIssue;
|
||||
const issue = locked.primaryIssue;
|
||||
const postCommitEffects: PostCommitEffect[] = [];
|
||||
|
||||
// Each `continue` path below leaves the wake row off the
|
||||
// `deferred_issue_execution` status, so the next queue read cannot
|
||||
// return that same row again. That invariant is what ends this loop.
|
||||
// The `processedWakeIds` guard below makes a break of the invariant
|
||||
// fail loudly, instead of holding this transaction open forever.
|
||||
const processedWakeIds = new Set<string>();
|
||||
|
||||
while (true) {
|
||||
const candidate = await ports.writer.claimNextDeferredWake({ companyId: run.companyId, issueId: issue.id });
|
||||
const candidate = await ports.transaction.findNextDeferredWake({ companyId: run.companyId, issueId: issue.id });
|
||||
if (!candidate) break;
|
||||
if (processedWakeIds.has(candidate.id)) {
|
||||
throw new WakeQueueApplicationError(
|
||||
"deferred_wake_not_advanced",
|
||||
"Deferred wake queue read the same wake id twice; the row did not leave the deferred status",
|
||||
{ companyId: run.companyId, issueId: issue.id, wakeId: candidate.id },
|
||||
);
|
||||
}
|
||||
processedWakeIds.add(candidate.id);
|
||||
|
||||
let liveness = { liveNonSelfCommentIds: candidate.queuedCommentIds, containedSelfAuthoredComment: false };
|
||||
if (candidate.queuedCommentIds.length > 0) {
|
||||
liveness = await ports.writer.getQueuedCommentLiveness({
|
||||
liveness = await ports.transaction.getQueuedCommentLiveness({
|
||||
companyId: run.companyId,
|
||||
issueId: issue.id,
|
||||
wakeAgentId: candidate.agentId,
|
||||
@@ -111,12 +160,11 @@ async function runReleaseDrain(
|
||||
queuedCommentIds: candidate.queuedCommentIds,
|
||||
});
|
||||
}
|
||||
const liveCommentIdsChanged =
|
||||
liveness.liveNonSelfCommentIds.length !== candidate.queuedCommentIds.length ||
|
||||
liveness.liveNonSelfCommentIds.some((id, index) => id !== candidate.queuedCommentIds[index]);
|
||||
// A length mismatch is the only way the lists can differ: the adapter derives `liveNonSelfCommentIds` with `.filter`, so it is always a subsequence of `queuedCommentIds`.
|
||||
const liveCommentIdsDiffer = liveness.liveNonSelfCommentIds.length !== candidate.queuedCommentIds.length;
|
||||
|
||||
const deferredAgent = await ports.reader.findInvokableAgent({ companyId: run.companyId, agentId: candidate.agentId });
|
||||
const pauseHold = await ports.writer.getPauseHoldFacts({
|
||||
const deferredAgent = await ports.transaction.findInvokableAgent({ companyId: run.companyId, agentId: candidate.agentId });
|
||||
const pauseHold = await ports.transaction.getPauseHoldFacts({
|
||||
companyId: run.companyId,
|
||||
issueId: issue.id,
|
||||
wakeAgentId: candidate.agentId,
|
||||
@@ -125,24 +173,21 @@ async function runReleaseDrain(
|
||||
requestedByActorId: candidate.requestedByActorId,
|
||||
});
|
||||
|
||||
let decision = decideDeferredWake({
|
||||
queuedComment: {
|
||||
hasQueuedCommentIds: candidate.queuedCommentIds.length > 0,
|
||||
liveNonSelfCommentIdsLength: liveness.liveNonSelfCommentIds.length,
|
||||
queuedCommentIdsLength: candidate.queuedCommentIds.length,
|
||||
liveCommentIdsChanged,
|
||||
containedSelfAuthoredComment: liveness.containedSelfAuthoredComment,
|
||||
preservesIndependentContinuation: candidate.preservesIndependentContinuation,
|
||||
},
|
||||
agent: { agentFound: deferredAgent !== null, invokable: deferredAgent?.invokable ?? false },
|
||||
pauseHold: { activePauseHold: pauseHold.activePauseHold, treeHoldInteractionWake: pauseHold.treeHoldInteractionWake },
|
||||
const commentAction = decideQueuedCommentAction({
|
||||
hasQueuedCommentIds: candidate.queuedCommentIds.length > 0,
|
||||
liveNonSelfCommentIdsLength: liveness.liveNonSelfCommentIds.length,
|
||||
liveCommentIdsDiffer,
|
||||
containedSelfAuthoredComment: liveness.containedSelfAuthoredComment,
|
||||
preservesIndependentContinuation: candidate.preservesIndependentContinuation,
|
||||
});
|
||||
|
||||
if (decision.kind === "cancel_empty") {
|
||||
await ports.writer.cancelDeferredWake({
|
||||
if (commentAction.kind === "cancel_empty") {
|
||||
// A `false` result means another writer already moved this row off
|
||||
// the deferred status, so the next queue read cannot return it again.
|
||||
await ports.transaction.cancelDeferredWake({
|
||||
companyId: run.companyId,
|
||||
wakeId: candidate.id,
|
||||
reason: decision.selfAuthored
|
||||
reason: commentAction.selfAuthored
|
||||
? "Deferred wake contained only comments authored by the finishing run"
|
||||
: "Queued messages were discarded before promotion",
|
||||
now: input.now,
|
||||
@@ -151,8 +196,8 @@ async function runReleaseDrain(
|
||||
}
|
||||
|
||||
let workingCandidate = candidate;
|
||||
if (decision.kind === "normalize") {
|
||||
const normalized = await ports.writer.normalizeDeferredWakeCommentIds({
|
||||
if (commentAction.kind === "normalize") {
|
||||
const normalized = await ports.transaction.normalizeDeferredWakeCommentIds({
|
||||
companyId: run.companyId,
|
||||
wakeId: candidate.id,
|
||||
payload: candidate.payload,
|
||||
@@ -161,29 +206,21 @@ async function runReleaseDrain(
|
||||
});
|
||||
if (!normalized) continue;
|
||||
workingCandidate = normalized;
|
||||
// Re-decide with the same agent/pause-hold facts already fetched above; the
|
||||
// comment-id set now matches, so only fail/cancel-pause-hold/promote can result.
|
||||
decision = decideDeferredWake({
|
||||
queuedComment: {
|
||||
hasQueuedCommentIds: workingCandidate.queuedCommentIds.length > 0,
|
||||
liveNonSelfCommentIdsLength: liveness.liveNonSelfCommentIds.length,
|
||||
queuedCommentIdsLength: liveness.liveNonSelfCommentIds.length,
|
||||
liveCommentIdsChanged: false,
|
||||
containedSelfAuthoredComment: liveness.containedSelfAuthoredComment,
|
||||
preservesIndependentContinuation: workingCandidate.preservesIndependentContinuation,
|
||||
},
|
||||
agent: { agentFound: deferredAgent !== null, invokable: deferredAgent?.invokable ?? false },
|
||||
pauseHold: { activePauseHold: pauseHold.activePauseHold, treeHoldInteractionWake: pauseHold.treeHoldInteractionWake },
|
||||
});
|
||||
}
|
||||
|
||||
if (decision.kind === "fail_not_invokable") {
|
||||
await ports.writer.failDeferredWake({ companyId: run.companyId, wakeId: workingCandidate.id, now: input.now });
|
||||
// A comment-id rewrite cannot change the agent or pause-hold facts already fetched above, so one decision covers both the rewritten and un-rewritten cases.
|
||||
const wakeOutcome = decideWakeOutcome({
|
||||
agent: { agentFound: deferredAgent !== null, invokable: deferredAgent?.invokable ?? false },
|
||||
pauseHold: { activePauseHold: pauseHold.activePauseHold, treeHoldInteractionWake: pauseHold.treeHoldInteractionWake },
|
||||
});
|
||||
|
||||
if (wakeOutcome.kind === "fail_not_invokable") {
|
||||
await ports.transaction.failDeferredWake({ companyId: run.companyId, wakeId: workingCandidate.id, now: input.now });
|
||||
continue;
|
||||
}
|
||||
|
||||
if (decision.kind === "cancel_pause_hold") {
|
||||
await ports.writer.cancelDeferredWake({
|
||||
if (wakeOutcome.kind === "cancel_pause_hold") {
|
||||
await ports.transaction.cancelDeferredWake({
|
||||
companyId: run.companyId,
|
||||
wakeId: workingCandidate.id,
|
||||
reason: "Deferred wake suppressed by active subtree pause hold",
|
||||
@@ -192,150 +229,164 @@ async function runReleaseDrain(
|
||||
continue;
|
||||
}
|
||||
|
||||
// decision.kind === "promote"
|
||||
const invokableAgent = deferredAgent!;
|
||||
// Unreachable: decideWakeOutcome only returns "promote" when agentFound and invokable are both true.
|
||||
if (!deferredAgent) throw new Error("wake-queue: promoted a deferred wake with no invokable agent");
|
||||
|
||||
// Claim the wake for promotion before any other write in this branch
|
||||
// (design choice: claim first, then reopen). A reopen write, or its
|
||||
// `issue_reopened` post-commit effect, must never survive a lost race on
|
||||
// this compare-and-set. When the claim fails, a concurrent writer already
|
||||
// changed the wake's status, so this candidate is gone; move on to the
|
||||
// next one instead of ending the drain.
|
||||
const claimedForPromotion = await ports.writer.claimDeferredWakeForPromotion({
|
||||
companyId: run.companyId,
|
||||
wakeId: workingCandidate.id,
|
||||
now: input.now,
|
||||
});
|
||||
if (!claimedForPromotion) continue;
|
||||
|
||||
let currentIssue = issue;
|
||||
|
||||
if (workingCandidate.deferredCommentIds.length > 0 && (currentIssue.status === "done" || currentIssue.status === "cancelled")) {
|
||||
const selfAuthorship = await ports.writer.getCommentSelfAuthorship({
|
||||
companyId: run.companyId,
|
||||
issueId: currentIssue.id,
|
||||
finishingRunId: run.id,
|
||||
commentIds: workingCandidate.deferredCommentIds,
|
||||
});
|
||||
const shouldReopen =
|
||||
!selfAuthorship.allSelfAuthored &&
|
||||
(workingCandidate.requestedByActorType === "user" || workingCandidate.wakeReason === "issue_reopened_via_comment");
|
||||
if (shouldReopen) {
|
||||
const reopened = await ports.writer.reopenIssue({ companyId: run.companyId, issueId: currentIssue.id, runId: run.id });
|
||||
if (reopened) {
|
||||
postCommitEffects.push({
|
||||
kind: "issue_reopened",
|
||||
companyId: reopened.companyId,
|
||||
agentId: invokableAgent.id,
|
||||
runId: run.id,
|
||||
issueId: reopened.id,
|
||||
identifier: reopened.identifier,
|
||||
reopenedFrom: currentIssue.status,
|
||||
});
|
||||
currentIssue = reopened;
|
||||
issue = reopened;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const promotedReason = workingCandidate.reason ?? "issue_execution_promoted";
|
||||
const promotedSource = workingCandidate.source ?? "automation";
|
||||
const promotedTriggerDetail = workingCandidate.triggerDetail ?? null;
|
||||
const promotedPayload = { ...workingCandidate.payload };
|
||||
delete promotedPayload["_paperclipWakeContext"];
|
||||
|
||||
const promotedContextSeed: Record<string, unknown> = { ...workingCandidate.deferredContextSeed };
|
||||
if (pauseHold.activePauseHold) {
|
||||
promotedContextSeed.treeHoldInteraction = true;
|
||||
promotedContextSeed.activeTreeHold = {
|
||||
holdId: pauseHold.holdId,
|
||||
rootIssueId: pauseHold.rootIssueId,
|
||||
mode: pauseHold.mode,
|
||||
reason: pauseHold.reason,
|
||||
releasePolicy: pauseHold.releasePolicy,
|
||||
interaction: true,
|
||||
};
|
||||
}
|
||||
|
||||
const { contextSnapshot: promotedContextSnapshot, taskKey: promotedTaskKey } = enrichPromotedWakeContext({
|
||||
contextSnapshot: promotedContextSeed,
|
||||
reason: promotedReason,
|
||||
source: promotedSource,
|
||||
triggerDetail: promotedTriggerDetail,
|
||||
payload: promotedPayload,
|
||||
});
|
||||
|
||||
const sessionBefore =
|
||||
readNonEmptyString(promotedContextSnapshot.resumeSessionDisplayId) ??
|
||||
(await ports.reader.resolveSessionBeforeForWakeup({
|
||||
companyId: run.companyId,
|
||||
agentId: invokableAgent.id,
|
||||
taskKey: promotedTaskKey,
|
||||
}));
|
||||
|
||||
const promotedRoutineEnvContext = await ports.reader.getRoutineEnv({
|
||||
companyId: invokableAgent.companyId,
|
||||
issue: currentIssue,
|
||||
});
|
||||
const responsibleUserId = await ports.reader.resolveResponsibleUserId({
|
||||
companyId: invokableAgent.companyId,
|
||||
contextSnapshot: promotedContextSnapshot,
|
||||
issue: currentIssue,
|
||||
routineEnvContext: promotedRoutineEnvContext,
|
||||
requestedByActorType: workingCandidate.requestedByActorType as "user" | "agent" | "system" | null,
|
||||
requestedByActorId: workingCandidate.requestedByActorId,
|
||||
source: promotedSource,
|
||||
triggerDetail: promotedTriggerDetail,
|
||||
existingRunResponsibleUserId: run.responsibleUserId,
|
||||
});
|
||||
if (!responsibleUserId) {
|
||||
throw new WakeQueueApplicationError(
|
||||
"responsible_user_unresolved",
|
||||
"Unable to resolve responsible user for promoted heartbeat run",
|
||||
{
|
||||
runId: run.id,
|
||||
agentId: invokableAgent.id,
|
||||
companyId: invokableAgent.companyId,
|
||||
issueId: currentIssue.id,
|
||||
wakeReason: readNonEmptyString(promotedContextSnapshot.wakeReason),
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
const promotedRun = await ports.writer.finalizePromotedWake({
|
||||
companyId: run.companyId,
|
||||
wakeId: workingCandidate.id,
|
||||
deferredAgent: invokableAgent,
|
||||
issue: currentIssue,
|
||||
finishingRun: run,
|
||||
contextSnapshot: promotedContextSnapshot,
|
||||
reason: promotedReason,
|
||||
source: promotedSource,
|
||||
triggerDetail: promotedTriggerDetail,
|
||||
payload: promotedPayload,
|
||||
responsibleUserId,
|
||||
sessionBefore,
|
||||
now: input.now,
|
||||
});
|
||||
|
||||
postCommitEffects.push({ kind: "run_queued", run: promotedRun });
|
||||
return { outcome: { kind: "promoted", run: promotedRun }, postCommitEffects };
|
||||
const promoted = await promoteDeferredWake(ports, run, issue, workingCandidate, deferredAgent, pauseHold, postCommitEffects, input);
|
||||
if (!promoted) continue;
|
||||
return promoted;
|
||||
}
|
||||
|
||||
return runReleaseRecoveryTail(issue, run, ports.reader, ports.writer, input, postCommitEffects);
|
||||
return runReleaseRecoveryTail(issue, run, ports.host, ports.transaction, input, postCommitEffects);
|
||||
}
|
||||
|
||||
/**
|
||||
* Finalizes one deferred wake that `decideWakeOutcome` chose to promote.
|
||||
* Returns `null` when the promotion claim loses a race, so the caller moves
|
||||
* on to the next queued wake instead of ending the drain.
|
||||
*/
|
||||
async function promoteDeferredWake(
|
||||
ports: { host: WakeQueueHost; transaction: WakeQueueTransaction },
|
||||
run: RunSnapshot,
|
||||
issue: IssueSnapshot,
|
||||
workingCandidate: DeferredWakeCandidate,
|
||||
invokableAgent: InvokableAgentSnapshot,
|
||||
pauseHold: PauseHoldFacts,
|
||||
postCommitEffects: PostCommitEffect[],
|
||||
input: ReleaseIssueExecutionInput,
|
||||
): Promise<ReleaseTransactionResult | null> {
|
||||
// Claim the wake for promotion before any other write in this branch
|
||||
// (design choice: claim first, then reopen). A reopen write, or its
|
||||
// `issue_reopened` post-commit effect, must never survive a lost race on
|
||||
// this compare-and-set. When the claim fails, a concurrent writer already
|
||||
// changed the wake's status, so this candidate is gone; the caller moves
|
||||
// on to the next one instead of ending the drain.
|
||||
const claimedForPromotion = await ports.transaction.claimDeferredWakeForPromotion({
|
||||
companyId: run.companyId,
|
||||
wakeId: workingCandidate.id,
|
||||
now: input.now,
|
||||
});
|
||||
if (!claimedForPromotion) return null;
|
||||
|
||||
let currentIssue = issue;
|
||||
|
||||
if (workingCandidate.deferredCommentIds.length > 0 && (currentIssue.status === "done" || currentIssue.status === "cancelled")) {
|
||||
const selfAuthorship = await ports.transaction.getCommentSelfAuthorship({
|
||||
companyId: run.companyId,
|
||||
issueId: currentIssue.id,
|
||||
finishingRunId: run.id,
|
||||
commentIds: workingCandidate.deferredCommentIds,
|
||||
});
|
||||
const shouldReopen =
|
||||
!selfAuthorship.allSelfAuthored &&
|
||||
(workingCandidate.requestedByActorType === "user" || workingCandidate.wakeReason === "issue_reopened_via_comment");
|
||||
if (shouldReopen) {
|
||||
const reopened = await ports.transaction.reopenIssue({ companyId: run.companyId, issueId: currentIssue.id, runId: run.id });
|
||||
if (reopened) {
|
||||
postCommitEffects.push({
|
||||
kind: "issue_reopened",
|
||||
companyId: reopened.companyId,
|
||||
agentId: invokableAgent.id,
|
||||
runId: run.id,
|
||||
issueId: reopened.id,
|
||||
identifier: reopened.identifier,
|
||||
reopenedFrom: currentIssue.status,
|
||||
});
|
||||
currentIssue = reopened;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const promotedReason = workingCandidate.reason ?? "issue_execution_promoted";
|
||||
const promotedSource = workingCandidate.source ?? "automation";
|
||||
const promotedTriggerDetail = workingCandidate.triggerDetail ?? null;
|
||||
const promotedPayload = { ...workingCandidate.payload };
|
||||
delete promotedPayload["_paperclipWakeContext"];
|
||||
|
||||
const promotedContextSeed: Record<string, unknown> = { ...workingCandidate.deferredContextSeed };
|
||||
if (pauseHold.activePauseHold) {
|
||||
promotedContextSeed.treeHoldInteraction = true;
|
||||
promotedContextSeed.activeTreeHold = {
|
||||
holdId: pauseHold.holdId,
|
||||
rootIssueId: pauseHold.rootIssueId,
|
||||
mode: pauseHold.mode,
|
||||
reason: pauseHold.reason,
|
||||
releasePolicy: pauseHold.releasePolicy,
|
||||
interaction: true,
|
||||
};
|
||||
}
|
||||
|
||||
const { contextSnapshot: promotedContextSnapshot, taskKey: promotedTaskKey } = enrichPromotedWakeContext({
|
||||
contextSnapshot: promotedContextSeed,
|
||||
reason: promotedReason,
|
||||
source: promotedSource,
|
||||
triggerDetail: promotedTriggerDetail,
|
||||
payload: promotedPayload,
|
||||
});
|
||||
|
||||
const sessionBefore =
|
||||
readNonEmptyString(promotedContextSnapshot.resumeSessionDisplayId) ??
|
||||
(await ports.host.resolveSessionBeforeForWakeup({
|
||||
companyId: run.companyId,
|
||||
agentId: invokableAgent.id,
|
||||
taskKey: promotedTaskKey,
|
||||
}));
|
||||
|
||||
const responsibleUserId = await resolveResponsibleUserForQueuedRun(ports.host, {
|
||||
companyId: invokableAgent.companyId,
|
||||
contextSnapshot: promotedContextSnapshot,
|
||||
issue: currentIssue,
|
||||
requestedByActorType: workingCandidate.requestedByActorType,
|
||||
requestedByActorId: workingCandidate.requestedByActorId,
|
||||
source: promotedSource,
|
||||
triggerDetail: promotedTriggerDetail,
|
||||
existingRunResponsibleUserId: run.responsibleUserId,
|
||||
});
|
||||
if (!responsibleUserId) {
|
||||
throw new WakeQueueApplicationError(
|
||||
"responsible_user_unresolved",
|
||||
"Unable to resolve responsible user for promoted heartbeat run",
|
||||
{
|
||||
runId: run.id,
|
||||
agentId: invokableAgent.id,
|
||||
companyId: invokableAgent.companyId,
|
||||
issueId: currentIssue.id,
|
||||
wakeReason: readNonEmptyString(promotedContextSnapshot.wakeReason),
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
const promotedRun = await ports.transaction.finalizePromotedWake({
|
||||
companyId: run.companyId,
|
||||
wakeId: workingCandidate.id,
|
||||
deferredAgent: invokableAgent,
|
||||
issue: currentIssue,
|
||||
finishingRun: run,
|
||||
contextSnapshot: promotedContextSnapshot,
|
||||
reason: promotedReason,
|
||||
source: promotedSource,
|
||||
triggerDetail: promotedTriggerDetail,
|
||||
payload: promotedPayload,
|
||||
responsibleUserId,
|
||||
sessionBefore,
|
||||
now: input.now,
|
||||
});
|
||||
|
||||
postCommitEffects.push({ kind: "run_queued", run: promotedRun });
|
||||
return { outcome: { kind: "promoted", run: promotedRun }, postCommitEffects };
|
||||
}
|
||||
|
||||
async function runReleaseRecoveryTail(
|
||||
issue: IssueSnapshot,
|
||||
run: RunSnapshot,
|
||||
reader: WakeQueueReader,
|
||||
writer: WakeQueueWriter,
|
||||
host: WakeQueueHost,
|
||||
transaction: WakeQueueTransaction,
|
||||
input: ReleaseIssueExecutionInput,
|
||||
postCommitEffects: PostCommitEffect[],
|
||||
): Promise<ReleaseTransactionResult> {
|
||||
const suppressImmediateRecovery = input.suppressImmediateRecovery ?? false;
|
||||
const isStrandedRecoveryOrigin = issue.originKind === STRANDED_ISSUE_RECOVERY_ORIGIN_KIND;
|
||||
const recoveryAgent = await reader.findInvokableAgent({ companyId: issue.companyId, agentId: run.agentId });
|
||||
const recoveryAgent = await transaction.findInvokableAgent({ companyId: issue.companyId, agentId: run.agentId });
|
||||
|
||||
const currentParticipant = currentAgentParticipant(issue);
|
||||
const reviewParticipantApplies =
|
||||
@@ -354,22 +405,22 @@ async function runReleaseRecoveryTail(
|
||||
(run.status === "failed" || run.status === "timed_out" || run.status === "cancelled");
|
||||
|
||||
const suppressedByPauseHold = (reviewParticipantApplies || immediateApplies)
|
||||
? await writer.isAutomaticRecoverySuppressedByPauseHold({ companyId: issue.companyId, issueId: issue.id })
|
||||
? await transaction.isAutomaticRecoverySuppressedByPauseHold({ companyId: issue.companyId, issueId: issue.id })
|
||||
: false;
|
||||
|
||||
const hasExistingExecutionPath = reviewParticipantApplies
|
||||
? await writer.hasExistingExecutionPath({
|
||||
? await transaction.hasExistingExecutionPath({
|
||||
companyId: issue.companyId,
|
||||
issueId: issue.id,
|
||||
excludeRunId: run.id,
|
||||
agentId: currentParticipant?.agentId ?? null,
|
||||
})
|
||||
: immediateApplies
|
||||
? await writer.hasExistingExecutionPath({ companyId: issue.companyId, issueId: issue.id, excludeRunId: run.id, agentId: null })
|
||||
? await transaction.hasExistingExecutionPath({ companyId: issue.companyId, issueId: issue.id, excludeRunId: run.id, agentId: null })
|
||||
: false;
|
||||
|
||||
const hasExplicitBlockerPath = immediateApplies && !reviewParticipantApplies
|
||||
? await writer.hasExplicitBlockerPath({ companyId: issue.companyId, issueId: issue.id })
|
||||
? await transaction.hasExplicitBlockerPath({ companyId: issue.companyId, issueId: issue.id })
|
||||
: false;
|
||||
|
||||
const expectedRetryReason: "assignment_recovery" | "issue_continuation_needed" =
|
||||
@@ -379,27 +430,23 @@ async function runReleaseRecoveryTail(
|
||||
suppressImmediateRecovery,
|
||||
reviewParticipant: {
|
||||
applies: reviewParticipantApplies,
|
||||
hasExistingExecutionPath,
|
||||
hasPersistedMonitor: Boolean(issue.monitorNextCheckAt),
|
||||
suppressedByPauseHold,
|
||||
isStrandedRecoveryOrigin,
|
||||
recoveryAgentPresent: recoveryAgent !== null,
|
||||
recoveryAgentInvokable: recoveryAgent?.invokable ?? false,
|
||||
isExecutionReviewParticipantRecoveryRun: isExecutionReviewParticipantRecoveryRun(run),
|
||||
},
|
||||
immediate: {
|
||||
applies: immediateApplies,
|
||||
isDispositionRepairRetry: readNonEmptyString(run.contextSnapshot.retryReason) === ISSUE_DISPOSITION_REPAIR_RETRY_REASON,
|
||||
hasExplicitBlockerPath,
|
||||
isWorkspaceValidationFailedRun: isWorkspaceValidationFailedRun(run),
|
||||
isConfigurationIncompleteFailedRun: isConfigurationIncompleteFailedRun(run),
|
||||
automaticRecoveryAlreadyFailed: didAutomaticRecoveryFail(run, expectedRetryReason),
|
||||
},
|
||||
shared: {
|
||||
hasExistingExecutionPath,
|
||||
hasPersistedMonitor: Boolean(issue.monitorNextCheckAt),
|
||||
hasExplicitBlockerPath,
|
||||
suppressedByPauseHold,
|
||||
isStrandedRecoveryOrigin,
|
||||
recoveryAgentPresent: recoveryAgent !== null,
|
||||
recoveryAgentInvokable: recoveryAgent?.invokable ?? false,
|
||||
isWorkspaceValidationFailedRun: isWorkspaceValidationFailedRun(run),
|
||||
isConfigurationIncompleteFailedRun: isConfigurationIncompleteFailedRun(run),
|
||||
automaticRecoveryAlreadyFailed: didAutomaticRecoveryFail(run, expectedRetryReason),
|
||||
},
|
||||
});
|
||||
|
||||
@@ -415,35 +462,32 @@ async function runReleaseRecoveryTail(
|
||||
}
|
||||
|
||||
if (decision.kind === "blocked") {
|
||||
const { notice, recoveryCause } = await writer.buildBlockedRecoveryNotice({
|
||||
noticeKind: decision.notice,
|
||||
issueStatus: issue.status === "todo" ? "todo" : "in_progress",
|
||||
finishingRun: run,
|
||||
});
|
||||
return {
|
||||
outcome: {
|
||||
kind: "blocked",
|
||||
issue,
|
||||
previousStatus: statusForBlock(issue),
|
||||
notice,
|
||||
recoveryCause,
|
||||
noticeKind: decision.notice,
|
||||
},
|
||||
postCommitEffects,
|
||||
};
|
||||
}
|
||||
|
||||
const sessionBefore = await reader.resolveSessionBeforeForWakeup({
|
||||
// Unreachable: decideReleaseRecovery only reaches "queue_review_participant_recovery" or "queue_recovery" when the shared recovery-agent facts are both true.
|
||||
if (!recoveryAgent) throw new Error("wake-queue: queued a recovery run with no invokable recovery agent");
|
||||
|
||||
const sessionBefore = await host.resolveSessionBeforeForWakeup({
|
||||
companyId: issue.companyId,
|
||||
agentId: recoveryAgent!.id,
|
||||
agentId: recoveryAgent.id,
|
||||
taskKey: readNonEmptyString(run.contextSnapshot.taskKey) ?? readNonEmptyString(run.contextSnapshot.issueId),
|
||||
});
|
||||
|
||||
if (decision.kind === "queue_review_participant_recovery") {
|
||||
const queuedRun = await writer.queueReviewParticipantRecoveryRun({
|
||||
const queuedRun = await transaction.queueReviewParticipantRecoveryRun({
|
||||
companyId: issue.companyId,
|
||||
issue,
|
||||
finishingRun: run,
|
||||
recoveryAgent: recoveryAgent!,
|
||||
recoveryAgent,
|
||||
sessionBefore,
|
||||
now: input.now,
|
||||
});
|
||||
@@ -451,14 +495,53 @@ async function runReleaseRecoveryTail(
|
||||
return { outcome: { kind: "queued_review_participant_recovery", run: queuedRun }, postCommitEffects };
|
||||
}
|
||||
|
||||
// decision.kind === "queue_recovery"; the adapter builds the recovery
|
||||
// context snapshot and resolves the responsible user from it, throwing
|
||||
// WakeQueueApplicationError when no responsible user resolves.
|
||||
const queuedRun = await writer.queueImmediateRecoveryRun({
|
||||
// decision.kind === "queue_recovery"; resolve the responsible user here,
|
||||
// in the application layer, before the transaction port queues the run.
|
||||
const { retryReason, recoveryReason, recoverySource } = deriveImmediateRecoveryContextLabels(issue.status);
|
||||
const recoveryContextSnapshot = withRecoveryContext(
|
||||
{
|
||||
issueId: issue.id,
|
||||
taskId: issue.id,
|
||||
wakeReason: recoveryReason,
|
||||
retryReason,
|
||||
source: recoverySource,
|
||||
retryOfRunId: run.id,
|
||||
},
|
||||
"normal_model",
|
||||
);
|
||||
|
||||
const recoveryResponsibleUserId = await resolveResponsibleUserForQueuedRun(host, {
|
||||
companyId: issue.companyId,
|
||||
contextSnapshot: recoveryContextSnapshot,
|
||||
issue,
|
||||
requestedByActorType: "system",
|
||||
requestedByActorId: null,
|
||||
source: "automation",
|
||||
triggerDetail: "system",
|
||||
existingRunResponsibleUserId: run.responsibleUserId,
|
||||
});
|
||||
if (!recoveryResponsibleUserId) {
|
||||
throw new WakeQueueApplicationError(
|
||||
"responsible_user_unresolved",
|
||||
"Unable to resolve responsible user for recovery heartbeat run",
|
||||
{
|
||||
runId: run.id,
|
||||
agentId: recoveryAgent.id,
|
||||
companyId: issue.companyId,
|
||||
issueId: issue.id,
|
||||
wakeReason: recoveryReason,
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
const queuedRun = await transaction.queueImmediateRecoveryRun({
|
||||
companyId: issue.companyId,
|
||||
issue,
|
||||
finishingRun: run,
|
||||
recoveryAgent: recoveryAgent!,
|
||||
recoveryAgent,
|
||||
reason: recoveryReason,
|
||||
contextSnapshot: recoveryContextSnapshot,
|
||||
responsibleUserId: recoveryResponsibleUserId,
|
||||
sessionBefore,
|
||||
now: input.now,
|
||||
});
|
||||
@@ -487,8 +570,7 @@ export function createReleaseIssueExecution(deps: {
|
||||
issue: result.outcome.issue,
|
||||
previousStatus: result.outcome.previousStatus,
|
||||
latestRun: result.run,
|
||||
notice: result.outcome.notice,
|
||||
recoveryCause: result.outcome.recoveryCause,
|
||||
noticeKind: result.outcome.noticeKind,
|
||||
});
|
||||
} else if (result.outcome.kind === "blocked_recovery_in_place") {
|
||||
await deps.recovery.escalateStrandedRecoveryIssueInPlace({
|
||||
|
||||
Reference in new issue
Block a user