diff --git a/doc/execution-semantics.md b/doc/execution-semantics.md index ca0eeeb806..ecabd67698 100644 --- a/doc/execution-semantics.md +++ b/doc/execution-semantics.md @@ -169,6 +169,8 @@ Use it for: Blocked issues should stay idle while blockers remain unresolved. Paperclip should not create a queued heartbeat run for that issue until the final blocker is done and the `issue_blockers_resolved` wake can start real work. +`cancelled` is terminal for the blocker issue itself, but it does not satisfy the dependency. A cancelled blocker edge remains unresolved until the edge is removed or replaced, and Paperclip must surface blocker attention on the dependent regardless of whether that dependent is currently displayed as `blocked`, `todo`, `backlog`, or another non-terminal agent-owned status. + If a parent is truly waiting on a child, model that with blockers. Do not rely on the parent/child relationship alone. ## 7. Accepted-Plan Decomposition @@ -266,6 +268,8 @@ An external wait counts as a live or waiting path only when the next move surviv - a first-class blocker or `blocked` disposition that names the external owner and concrete action required to unblock the issue - a delegated child issue with a responsible owner and its own healthy action path, plus a blocker edge when the source issue must wait for that child; `parentId` alone is not a dependency +A one-shot issue monitor consumes its persisted `nextCheckAt` when it dispatches the assignee wake. If that monitor-consuming run is lost before it records a new disposition or future monitor, Paperclip restores exactly one bounded continuation using the existing process-loss retry limit; if that continuation is also lost, the normal recovery-action escalation owns the next step instead of creating another monitor loop. + An unmanaged local process is not a durable action path. Shell jobs started with `&`, `nohup`, local polling loops, detached PTY sessions, adapter child processes, or similar background watchers do not keep an issue live unless Paperclip persists them as a run or pairs a managed runtime service with a monitor, scheduled wake, blocker, or delegated issue that owns the next check. A PID, session id, log file, comment, or promise to check later is evidence only. The process may be killed when the adapter invocation or heartbeat exits and cannot be assumed observable or recoverable by another worker. Before a heartbeat finalizes, its issue disposition must therefore be evaluated from durable Paperclip state, not from processes still visible only to that heartbeat. An agent-owned issue may remain `in_progress` after the heartbeat only when another valid action-path primitive already exists. If the only claimed continuation is a local/background watcher, finalization treats the issue as having no live path even when the process has not yet been observed exiting. diff --git a/server/src/__tests__/heartbeat-process-recovery.test.ts b/server/src/__tests__/heartbeat-process-recovery.test.ts index 49e87b2fd5..78ca629263 100644 --- a/server/src/__tests__/heartbeat-process-recovery.test.ts +++ b/server/src/__tests__/heartbeat-process-recovery.test.ts @@ -1303,6 +1303,106 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { expect(checkoutReleasedIssue?.checkoutRunId).toBeNull(); }); + it("restores one lost monitor dispatch before escalating a second process loss", async () => { + const { companyId, agentId, runId, issueId } = await seedRunFixture({ + adapterType: "openclaw_gateway", + agentStatus: "idle", + processPid: null, + processGroupId: null, + contextSnapshot: { + wakeReason: "issue_monitor_due", + nextCheckAt: "2026-03-19T00:00:00.000Z", + }, + }); + const heartbeat = heartbeatService(db); + + const firstLoss = await heartbeat.reapOrphanedRuns(); + expect(firstLoss).toEqual({ reaped: 1, runIds: [runId] }); + + const firstRetry = await db + .select() + .from(heartbeatRuns) + .where(and(eq(heartbeatRuns.companyId, companyId), eq(heartbeatRuns.retryOfRunId, runId))) + .then((rows) => rows[0] ?? null); + expect(firstRetry).toMatchObject({ processLossRetryCount: 1 }); + expect(firstRetry?.contextSnapshot).toMatchObject({ + wakeReason: "process_lost_retry", + retryReason: "issue_continuation_needed", + retryOfRunId: runId, + }); + + const secondAttempt = await seedRunFixture({ + adapterType: "openclaw_gateway", + agentStatus: "idle", + processPid: null, + processGroupId: null, + processLossRetryCount: 1, + contextSnapshot: { + wakeReason: "process_lost_retry", + retryReason: "issue_continuation_needed", + retryOfRunId: runId, + }, + }); + + const secondLoss = await heartbeat.reapOrphanedRuns(); + expect(secondLoss).toEqual({ reaped: 1, runIds: [secondAttempt.runId] }); + + const secondAttemptRuns = await db + .select() + .from(heartbeatRuns) + .where(eq(heartbeatRuns.agentId, secondAttempt.agentId)); + expect(secondAttemptRuns.find((run) => run.id === secondAttempt.runId)).toMatchObject({ + id: secondAttempt.runId, + status: "failed", + errorCode: "process_lost", + processLossRetryCount: 1, + }); + expect(secondAttemptRuns.some((run) => run.processLossRetryCount > 1)).toBe(false); + + const issue = await waitForValue(async () => + db.select().from(issues).where(eq(issues.id, secondAttempt.issueId)).then((rows) => { + const row = rows[0] ?? null; + return row?.status === "blocked" ? row : null; + }) + ); + expect(issue?.monitorNextCheckAt).toBeNull(); + + await expectSourceScopedStrandedRecoveryAction({ + companyId: secondAttempt.companyId, + agentId: secondAttempt.agentId, + issueId: secondAttempt.issueId, + runId: secondAttempt.runId, + previousStatus: "in_progress", + retryReason: "issue_continuation_needed", + }); + }); + + it("does not retry a lost monitor dispatch while another monitor wake remains scheduled", async () => { + const { companyId, runId, issueId } = await seedRunFixture({ + adapterType: "openclaw_gateway", + agentStatus: "idle", + processPid: null, + processGroupId: null, + contextSnapshot: { + wakeReason: "issue_monitor_due", + }, + }); + await db + .update(issues) + .set({ monitorNextCheckAt: new Date("2099-03-19T00:00:00.000Z") }) + .where(and(eq(issues.id, issueId), eq(issues.companyId, companyId))); + + const heartbeat = heartbeatService(db); + const result = await heartbeat.reapOrphanedRuns(); + + expect(result).toEqual({ reaped: 1, runIds: [runId] }); + const retries = await db + .select() + .from(heartbeatRuns) + .where(and(eq(heartbeatRuns.companyId, companyId), eq(heartbeatRuns.retryOfRunId, runId))); + expect(retries).toHaveLength(0); + }); + it("interrupts running runs on graceful shutdown and queues restart recovery without recording a failure", async () => { const { agentId, runId, issueId, wakeupRequestId } = await seedRunFixture({ agentStatus: "running", diff --git a/server/src/__tests__/issue-blocker-attention.test.ts b/server/src/__tests__/issue-blocker-attention.test.ts index 13480b6649..0d63f1a6fb 100644 --- a/server/src/__tests__/issue-blocker-attention.test.ts +++ b/server/src/__tests__/issue-blocker-attention.test.ts @@ -525,6 +525,42 @@ describeEmbeddedPostgres("issue blocker attention", () => { }); }); + it("does not treat cancelled liveness escalation issues as covered waiting paths", async () => { + const { companyId, agentId } = await createCompany("PBLX"); + const parentId = await insertIssue({ companyId, identifier: "PBLX-1", title: "Parent", status: "blocked" }); + const cancelledLeafId = await insertIssue({ + companyId, + identifier: "PBLX-2", + title: "Cancelled blocker", + status: "cancelled", + assigneeAgentId: agentId, + }); + await insertIssue({ + companyId, + identifier: "PBLX-3", + title: "Cancelled liveness escalation", + status: "cancelled", + assigneeAgentId: agentId, + originKind: "harness_liveness_escalation", + originId: [ + "harness_liveness", + companyId, + parentId, + "blocked_by_cancelled_issue", + cancelledLeafId, + ].join(":"), + }); + await block({ companyId, blockerIssueId: cancelledLeafId, blockedIssueId: parentId }); + + const parent = (await svc.list(companyId, { attention: "blocked" })).find((issue) => issue.id === parentId); + + expect(parent?.blockedInboxAttention).toMatchObject({ + state: "needs_attention", + reason: "blocked_by_cancelled_issue", + leafIssue: { id: cancelledLeafId, identifier: "PBLX-2" }, + }); + }); + it("does not treat a scheduled retry as actively covered work", async () => { const { companyId, agentId } = await createCompany("PBY"); const parentId = await insertIssue({ companyId, identifier: "PBY-1", title: "Parent", status: "blocked" }); @@ -581,6 +617,37 @@ describeEmbeddedPostgres("issue blocker attention", () => { await expect(svc.count(companyId, { attention: "blocked" })).resolves.toBe(1); }); + it("surfaces cancelled-blocker attention on an assigned todo source", async () => { + const { companyId, agentId } = await createCompany("BICX"); + const sourceId = await insertIssue({ + companyId, + identifier: "BICX-1", + title: "Dispatch-suppressed source", + status: "todo", + assigneeAgentId: agentId, + }); + const blockerId = await insertIssue({ + companyId, + identifier: "BICX-2", + title: "Cancelled dependency", + status: "cancelled", + assigneeAgentId: agentId, + }); + await block({ companyId, blockerIssueId: blockerId, blockedIssueId: sourceId }); + + const rows = await svc.list(companyId, { attention: "blocked" }); + const source = rows.find((issue) => issue.id === sourceId); + + expect(source?.blockedInboxAttention).toMatchObject({ + kind: "blocked", + state: "needs_attention", + reason: "blocked_by_cancelled_issue", + owner: { type: "agent", agentId }, + action: { label: "Replace blocker" }, + leafIssue: { id: blockerId, identifier: "BICX-2" }, + }); + }); + it("redacts external wait details from blocked inbox payloads and search", async () => { const { companyId } = await createCompany("BIX"); const owner = "Private Vendor Security Team"; diff --git a/server/src/__tests__/issue-liveness.test.ts b/server/src/__tests__/issue-liveness.test.ts index 8dba1eb583..40ec85bf3e 100644 --- a/server/src/__tests__/issue-liveness.test.ts +++ b/server/src/__tests__/issue-liveness.test.ts @@ -273,6 +273,50 @@ describe("issue graph liveness classifier", () => { expect(paused[0]?.state).toBe("blocked_by_uninvokable_assignee"); }); + it("detects a cancelled blocker on an assigned todo source", () => { + const findings = classifyIssueGraphLiveness({ + issues: [ + issue({ status: "todo" }), + issue({ + id: blockerId, + identifier: "PAP-1704", + title: "Cancelled unblock work", + status: "cancelled", + assigneeAgentId: "blocker-agent", + }), + ], + relations: blocks, + agents: [agent(), manager, agent({ id: "blocker-agent", name: "Cancelled owner" })], + }); + + expect(findings).toHaveLength(1); + expect(findings[0]).toMatchObject({ + issueId: blockedId, + state: "blocked_by_cancelled_issue", + recoveryIssueId: blockerId, + }); + }); + + it("prefers the blocker finding for an in-review source with a cancelled blocker", () => { + const findings = classifyIssueGraphLiveness({ + issues: [ + issue({ status: "in_review" }), + issue({ + id: blockerId, + identifier: "PAP-1704", + title: "Cancelled unblock work", + status: "cancelled", + assigneeAgentId: "blocker-agent", + }), + ], + relations: blocks, + agents: [agent(), manager, agent({ id: "blocker-agent", name: "Cancelled owner" })], + }); + + expect(findings).toHaveLength(1); + expect(findings[0]?.state).toBe("blocked_by_cancelled_issue"); + }); + it("detects blocker assignees under terminated org ancestors as uninvokable", () => { const findings = classifyIssueGraphLiveness({ issues: [ @@ -492,6 +536,41 @@ describe("issue graph liveness classifier", () => { } }); + it("still flags a stalled in_review issue when its blocker has an active run", () => { + const reviewIssueId = "review-1"; + const activeBlockerId = "active-blocker-1"; + + const findings = classifyIssueGraphLiveness({ + issues: [ + issue({ + id: reviewIssueId, + identifier: "PAP-2279", + title: "Screenshot acceptance review", + status: "in_review", + assigneeAgentId: coderId, + executionState: null, + }), + issue({ + id: activeBlockerId, + identifier: "PAP-2280", + title: "Active blocker", + status: "in_progress", + assigneeAgentId: coderId, + }), + ], + relations: [{ companyId, blockerIssueId: activeBlockerId, blockedIssueId: reviewIssueId }], + agents: [agent(), manager], + activeRuns: [{ companyId, issueId: activeBlockerId, agentId: coderId, status: "running" }], + }); + + expect(findings).toHaveLength(1); + expect(findings[0]).toMatchObject({ + issueId: reviewIssueId, + state: "in_review_without_action_path", + recoveryIssueId: reviewIssueId, + }); + }); + it("ignores cross-company waiting paths for stalled in_review issues", () => { const reviewIssueId = "review-1"; diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 62a3f26904..a284d3a145 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -8509,13 +8509,16 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {}) const contextSnapshot = parseObject(run.contextSnapshot); const issueId = readNonEmptyString(contextSnapshot.issueId); + const retryReason = readNonEmptyString(contextSnapshot.wakeReason) === "issue_monitor_due" + ? "issue_continuation_needed" + : "process_lost"; const taskKey = deriveTaskKeyWithHeartbeatFallback(contextSnapshot, null); const sessionBefore = await resolveSessionBeforeForWakeup(agent, taskKey); const retryContextSnapshot = withRecoveryModelProfileHint({ ...contextSnapshot, retryOfRunId: run.id, wakeReason: "process_lost_retry", - retryReason: "process_lost", + retryReason, }, "normal_model"); const responsibleUserId = await resolveResponsibleUserIdForRunContext(run, retryContextSnapshot); @@ -10920,6 +10923,29 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {}) .innerJoin(agents, eq(heartbeatRuns.agentId, agents.id)) .where(eq(heartbeatRuns.status, "running")); + const monitorIssueIds = [...new Set(activeRuns.flatMap(({ run }) => { + const runContext = parseObject(run.contextSnapshot); + if (readNonEmptyString(runContext.wakeReason) !== "issue_monitor_due") return []; + const issueId = readNonEmptyString(runContext.issueId); + return issueId ? [issueId] : []; + }))]; + const monitorIssues = monitorIssueIds.length > 0 + ? await db + .select({ + id: issues.id, + companyId: issues.companyId, + monitorNextCheckAt: issues.monitorNextCheckAt, + }) + .from(issues) + .where(inArray(issues.id, monitorIssueIds)) + : []; + const monitorNextCheckAtByIssue = new Map( + monitorIssues.map((issue) => [ + `${issue.companyId}:${issue.id}`, + issue.monitorNextCheckAt, + ]), + ); + const reaped: string[] = []; for (const { run, adapterType, adapterConfig } of activeRuns) { @@ -10965,7 +10991,19 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {}) }); } - const shouldRetry = tracksLocalChild && (!!run.processPid || !!run.processGroupId) && (run.processLossRetryCount ?? 0) < 1; + const runContext = parseObject(run.contextSnapshot); + const monitorIssueId = readNonEmptyString(runContext.issueId); + const monitorNextCheckAt = monitorIssueId + ? monitorNextCheckAtByIssue.get(`${run.companyId}:${monitorIssueId}`) + : undefined; + const monitorDispatchLostWithoutFutureWake = + readNonEmptyString(runContext.wakeReason) === "issue_monitor_due" && + monitorNextCheckAt !== undefined && + (!monitorNextCheckAt || monitorNextCheckAt.getTime() <= now.getTime()); + const shouldRetry = (run.processLossRetryCount ?? 0) < 1 && ( + (tracksLocalChild && (!!run.processPid || !!run.processGroupId)) || + monitorDispatchLostWithoutFutureWake + ); const baseMessage = buildProcessLossMessage(run, descendantOnlyCleanup ? { descendantOnly: true } : undefined); const unmanagedBackgroundTaskEvidence = descendantOnlyCleanup ? { diff --git a/server/src/services/issues.ts b/server/src/services/issues.ts index 2f75a672bf..3d93671edf 100644 --- a/server/src/services/issues.ts +++ b/server/src/services/issues.ts @@ -2985,7 +2985,7 @@ async function listIssueBlockedInboxAttentionMap( .where(and( eq(issues.companyId, companyId), visibleIssueCondition(), - notInArray(issues.status, [...BLOCKED_INBOX_TERMINAL_STATUSES]), + ne(issues.status, "done"), )), dbOrTx .select({ @@ -3117,6 +3117,7 @@ async function listIssueBlockedInboxAttentionMap( const openRecoveryIssues = graphIssues .filter((issue) => BLOCKED_INBOX_RECOVERY_ORIGIN_KINDS.includes(issue.originKind as typeof BLOCKED_INBOX_RECOVERY_ORIGIN_KINDS[number])) + .filter((issue) => !BLOCKED_INBOX_TERMINAL_STATUSES.includes(issue.status as typeof BLOCKED_INBOX_TERMINAL_STATUSES[number])) .flatMap((issue) => { const entries = [{ companyId, issueId: issue.id, status: issue.status }]; if (issue.originKind === "harness_liveness_escalation") { diff --git a/server/src/services/recovery/issue-graph-liveness.ts b/server/src/services/recovery/issue-graph-liveness.ts index 91cf2e73dd..c041667aa7 100644 --- a/server/src/services/recovery/issue-graph-liveness.ts +++ b/server/src/services/recovery/issue-graph-liveness.ts @@ -590,13 +590,26 @@ export function classifyIssueGraphLiveness(input: IssueGraphLivenessInput): Issu } for (const issue of input.issues) { - if (issue.status === "blocked") { + const hasUnresolvedBlockerEdge = (blockersByBlockedIssueId.get(issue.id) ?? []).some((relation) => { + if (relation.companyId !== issue.companyId) return false; + const blocker = issuesById.get(relation.blockerIssueId); + return Boolean(blocker && blocker.companyId === issue.companyId && blocker.status !== "done"); + }); + const shouldInspectBlockedChain = issue.status === "blocked" || ( + issue.status !== "done" && + issue.status !== "cancelled" && + Boolean(issue.assigneeAgentId) && + hasUnresolvedBlockerEdge + ); + + let chainFinding: IssueLivenessFinding | null = null; + if (shouldInspectBlockedChain) { if (unresolvedBlockers.has(issue.id)) continue; - const chainFinding = firstBlockedChainFinding(issue, issue, [issue], new Set()); + chainFinding = firstBlockedChainFinding(issue, issue, [issue], new Set()); if (chainFinding) findings.push(chainFinding); } - if (issue.status === "in_review" && !unresolvedBlockers.has(issue.id)) { + if (issue.status === "in_review" && !chainFinding && !unresolvedBlockers.has(issue.id)) { const review = reviewFinding(issue, issue, [issue]); if (review) findings.push(review); }