Skip to content

Commit b03fb45

Browse files
author
Test
committed
fix: restore execution lock reaper scheduling
1 parent 26aa3a1 commit b03fb45

7 files changed

Lines changed: 131 additions & 46 deletions

File tree

server/src/__tests__/execution-lock-acquisition-helper.test.ts

Lines changed: 31 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,15 +1,26 @@
11
import { describe, expect, it } from "vitest";
22
import { readFile } from "node:fs/promises";
33
import path from "node:path";
4-
import { executionLockAcquisitionFields } from "../services/issues.js";
4+
import { ISSUE_EXECUTION_LOCK_TTL_MS, executionLockAcquisitionFields } from "../services/issues.js";
55

66
describe("execution lock acquisition helper usage", () => {
7-
it("does not create a scheduled review monitor when acquiring an execution lock", () => {
7+
it("always schedules the exact execution-lock reaper deadline", () => {
88
const now = new Date("2026-08-08T03:15:00.000Z");
99

1010
expect(executionLockAcquisitionFields("run-1", now)).toEqual({
1111
executionRunId: "run-1",
1212
executionLockedAt: now,
13+
monitorNextCheckAt: new Date(now.getTime() + ISSUE_EXECUTION_LOCK_TTL_MS),
14+
monitorWakeRequestedAt: null,
15+
});
16+
});
17+
18+
it("keeps an earlier validated review-monitor deadline ahead of the lock reaper", () => {
19+
const now = new Date("2026-08-08T03:15:00.000Z");
20+
const reviewDeadline = new Date("2026-08-08T03:17:00.000Z");
21+
22+
expect(executionLockAcquisitionFields("run-1", now, reviewDeadline)).toMatchObject({
23+
monitorNextCheckAt: reviewDeadline,
1324
});
1425
});
1526

@@ -34,4 +45,22 @@ describe("execution lock acquisition helper usage", () => {
3445

3546
expect(directAcquisitions).toEqual([]);
3647
});
48+
49+
it("removes scheduleMonitor option plumbing from every production acquisition path", async () => {
50+
const root = process.cwd();
51+
const productionFiles = [
52+
"server/src/services/issues.ts",
53+
"server/src/services/heartbeat.ts",
54+
];
55+
56+
const optionUses: string[] = [];
57+
for (const file of productionFiles) {
58+
const source = await readFile(path.join(root, file), "utf8");
59+
source.split(/\r?\n/).forEach((line, index) => {
60+
if (/scheduleMonitor/.test(line)) optionUses.push(`${file}:${index + 1}:${line.trim()}`);
61+
});
62+
}
63+
64+
expect(optionUses).toEqual([]);
65+
});
3766
});

server/src/__tests__/issue-liveness.test.ts

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import { describe, expect, it } from "vitest";
22
import { classifyIssueGraphLiveness } from "../services/issue-liveness.ts";
3+
import { hasScheduledIssueMonitorPath } from "../services/recovery/issue-graph-liveness.ts";
34

45
const companyId = "company-1";
56
const managerId = "manager-1";
@@ -46,6 +47,43 @@ const manager = agent({
4647
const blocks = [{ companyId, blockerIssueId: blockerId, blockedIssueId: blockedId }];
4748

4849
describe("issue graph liveness classifier", () => {
50+
it("counts only validated future typed review monitors, never a raw lock-reaper timestamp", () => {
51+
const now = new Date("2026-08-08T03:15:00.000Z");
52+
const future = new Date("2026-08-08T03:20:00.000Z").toISOString();
53+
const expired = new Date("2026-08-08T03:10:00.000Z").toISOString();
54+
const executionState = (monitor: Record<string, unknown>) => ({
55+
status: "idle",
56+
currentStageId: null,
57+
currentStageIndex: null,
58+
currentStageType: null,
59+
currentParticipant: null,
60+
returnAssignee: null,
61+
reviewRequest: null,
62+
completedStageIds: [],
63+
lastDecisionId: null,
64+
lastDecisionOutcome: null,
65+
monitor,
66+
changesRequestedCount: 0,
67+
});
68+
69+
const scheduled = executionState({
70+
status: "scheduled", nextCheckAt: future, lastTriggeredAt: null, attemptCount: 0,
71+
notes: null, scheduledBy: "assignee", kind: null, serviceName: null, externalRef: null,
72+
timeoutAt: null, maxAttempts: null, recoveryPolicy: null, clearedAt: null, clearReason: null,
73+
});
74+
const triggered = { ...scheduled, monitor: { ...scheduled.monitor, status: "triggered", nextCheckAt: null, lastTriggeredAt: expired } };
75+
const cleared = { ...scheduled, monitor: { ...scheduled.monitor, status: "cleared", nextCheckAt: null, clearedAt: expired, clearReason: "manual" } };
76+
const expiredState = { ...scheduled, monitor: { ...scheduled.monitor, timeoutAt: expired } };
77+
const exhaustedState = { ...scheduled, monitor: { ...scheduled.monitor, maxAttempts: 1, attemptCount: 1 } };
78+
const invalid = { ...scheduled, monitor: { ...scheduled.monitor, nextCheckAt: "not-a-date" } };
79+
80+
expect(hasScheduledIssueMonitorPath(issue({ status: "in_review", monitorNextCheckAt: future }), now)).toBe(false);
81+
expect(hasScheduledIssueMonitorPath(issue({ status: "in_review", monitorNextCheckAt: future, executionState: scheduled }), now)).toBe(true);
82+
for (const state of [triggered, cleared, expiredState, exhaustedState, invalid]) {
83+
expect(hasScheduledIssueMonitorPath(issue({ status: "in_review", monitorNextCheckAt: future, executionState: state }), now)).toBe(false);
84+
}
85+
});
86+
4987
it("detects a PAP-1703-style blocked chain with an unassigned blocker and stable incident key", () => {
5088
const findings = classifyIssueGraphLiveness({
5189
issues: [

server/src/services/heartbeat.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13617,7 +13617,7 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
1361713617
})
1361813618
) {
1361913619
try {
13620-
await issuesSvc.checkout(issueId, agent.id, ["todo", "backlog", "blocked"], run.id, { scheduleMonitor: false });
13620+
await issuesSvc.checkout(issueId, agent.id, ["todo", "backlog", "blocked"], run.id);
1362113621
context[PAPERCLIP_HARNESS_CHECKOUT_KEY] = true;
1362213622
} catch (error) {
1362313623
if (!isCheckoutConflictError(error)) throw error;

server/src/services/issue-execution-policy.ts

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -421,6 +421,39 @@ export function parseIssueExecutionState(input: unknown): IssueExecutionState |
421421
return parsed.data;
422422
}
423423

424+
export function activeTypedIssueMonitorDeadline(
425+
input: Pick<IssueLike, "executionPolicy" | "executionState">,
426+
now: Date | string | number = new Date(),
427+
): Date | null {
428+
const nowMs = now instanceof Date ? now.getTime() : typeof now === "number" ? now : new Date(now).getTime();
429+
if (!Number.isFinite(nowMs)) return null;
430+
const state = parseIssueExecutionState(input.executionState);
431+
const policyResult = issueExecutionPolicySchema.safeParse(input.executionPolicy);
432+
const candidates: Array<{ monitor: IssueExecutionMonitorPolicy | IssueExecutionMonitorState; attempts: number }> = [];
433+
434+
if (policyResult.success && policyResult.data.monitor) {
435+
candidates.push({ monitor: policyResult.data.monitor, attempts: state?.monitor?.attemptCount ?? 0 });
436+
}
437+
if (state?.monitor?.status === "scheduled") {
438+
candidates.push({ monitor: state.monitor, attempts: state.monitor.attemptCount });
439+
}
440+
441+
const deadlines = candidates.flatMap(({ monitor, attempts }) => {
442+
if (!monitor.nextCheckAt) return [];
443+
const nextCheckAt = new Date(monitor.nextCheckAt);
444+
const nextCheckAtMs = nextCheckAt.getTime();
445+
if (!Number.isFinite(nextCheckAtMs) || nextCheckAtMs <= nowMs) return [];
446+
if (monitor.timeoutAt) {
447+
const timeoutAtMs = new Date(monitor.timeoutAt).getTime();
448+
if (!Number.isFinite(timeoutAtMs) || timeoutAtMs <= nowMs) return [];
449+
}
450+
if (monitor.maxAttempts != null && attempts >= monitor.maxAttempts) return [];
451+
return [nextCheckAt];
452+
});
453+
454+
return deadlines.sort((left, right) => left.getTime() - right.getTime())[0] ?? null;
455+
}
456+
424457
export function assigneePrincipal(input: AssigneeLike): IssueExecutionStagePrincipal | null {
425458
if (input.assigneeAgentId) {
426459
return { type: "agent", agentId: input.assigneeAgentId, userId: null };

server/src/services/issues.ts

Lines changed: 24 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,7 @@ import {
8686
type ParsedExecutionWorkspaceMode,
8787
} from "./execution-workspace-policy.js";
8888
import { mergeExecutionWorkspaceConfig } from "./execution-workspaces.js";
89-
import { buildInitialIssueMonitorFields, normalizeIssueExecutionPolicy } from "./issue-execution-policy.js";
89+
import { activeTypedIssueMonitorDeadline, buildInitialIssueMonitorFields, normalizeIssueExecutionPolicy } from "./issue-execution-policy.js";
9090
import { instanceSettingsService } from "./instance-settings.js";
9191
import { redactCurrentUserText } from "../log-redaction.js";
9292
import { redactSensitiveText } from "../redaction.js";
@@ -783,25 +783,18 @@ function issueExecutionLockMonitorNextCheckAt(now: Date) {
783783
export function executionLockAcquisitionFields(
784784
runId: string,
785785
now: Date,
786-
options: { scheduleMonitor?: boolean } = {},
786+
activeReviewMonitorDeadline: Date | null = null,
787787
) {
788-
const monitorFields = options.scheduleMonitor
789-
? {
790-
monitorNextCheckAt: sql<Date | null>`case
791-
when ${issues.monitorNextCheckAt} is null
792-
and ${issues.monitorLastTriggeredAt} is null
793-
and ${issues.monitorAttemptCount} = 0
794-
then ${issueExecutionLockMonitorNextCheckAt(now).toISOString()}::timestamptz
795-
else ${issues.monitorNextCheckAt}
796-
end`,
797-
monitorWakeRequestedAt: null,
798-
}
799-
: {};
788+
const lockDeadline = issueExecutionLockMonitorNextCheckAt(now);
789+
const monitorNextCheckAt = activeReviewMonitorDeadline && activeReviewMonitorDeadline.getTime() < lockDeadline.getTime()
790+
? activeReviewMonitorDeadline
791+
: lockDeadline;
800792

801793
return {
802794
executionRunId: runId,
803795
executionLockedAt: now,
804-
...monitorFields,
796+
monitorNextCheckAt,
797+
monitorWakeRequestedAt: null,
805798
};
806799
}
807800

@@ -5111,7 +5104,6 @@ export function issueService(db: Db) {
51115104
actorAgentId: string;
51125105
actorRunId: string;
51135106
expectedCheckoutRunId: string;
5114-
scheduleMonitor: boolean;
51155107
}) {
51165108
return db.transaction(async (tx) => {
51175109
const lockedIssue = await tx
@@ -5121,6 +5113,8 @@ export function issueService(db: Db) {
51215113
assigneeAgentId: issues.assigneeAgentId,
51225114
checkoutRunId: issues.checkoutRunId,
51235115
executionRunId: issues.executionRunId,
5116+
executionPolicy: issues.executionPolicy,
5117+
executionState: issues.executionState,
51245118
})
51255119
.from(issues)
51265120
.where(eq(issues.id, input.issueId))
@@ -5169,7 +5163,7 @@ export function issueService(db: Db) {
51695163
.update(issues)
51705164
.set({
51715165
checkoutRunId: input.actorRunId,
5172-
...executionLockAcquisitionFields(input.actorRunId, now, { scheduleMonitor: input.scheduleMonitor }),
5166+
...executionLockAcquisitionFields(input.actorRunId, now, activeTypedIssueMonitorDeadline(lockedIssue, now)),
51735167
updatedAt: now,
51745168
})
51755169
.where(
@@ -5211,7 +5205,6 @@ export function issueService(db: Db) {
52115205
issueId: string;
52125206
actorAgentId: string;
52135207
actorRunId: string;
5214-
scheduleMonitor: boolean;
52155208
}) {
52165209
return db.transaction(async (tx) => {
52175210
await tx.execute(
@@ -5225,11 +5218,16 @@ export function issueService(db: Db) {
52255218
if (!actorRun || heartbeatRunIsDead(actorRun)) return null;
52265219

52275220
const now = new Date();
5221+
const issue = await tx
5222+
.select({ executionPolicy: issues.executionPolicy, executionState: issues.executionState })
5223+
.from(issues)
5224+
.where(eq(issues.id, input.issueId))
5225+
.then((rows) => rows[0] ?? null);
52285226
const adopted = await tx
52295227
.update(issues)
52305228
.set({
52315229
checkoutRunId: input.actorRunId,
5232-
...executionLockAcquisitionFields(input.actorRunId, now, { scheduleMonitor: input.scheduleMonitor }),
5230+
...executionLockAcquisitionFields(input.actorRunId, now, activeTypedIssueMonitorDeadline(issue ?? {}, now)),
52335231
updatedAt: now,
52345232
})
52355233
.where(
@@ -7992,18 +7990,17 @@ export function issueService(db: Db) {
79927990
agentId: string,
79937991
expectedStatuses: string[],
79947992
checkoutRunId: string | null,
7995-
options: { scheduleMonitor?: boolean } = {},
79967993
) => {
79977994
const issueCompany = await db
7998-
.select({ companyId: issues.companyId })
7995+
.select({ companyId: issues.companyId, executionPolicy: issues.executionPolicy, executionState: issues.executionState })
79997996
.from(issues)
80007997
.where(eq(issues.id, id))
80017998
.then((rows) => rows[0] ?? null);
80027999
if (!issueCompany) throw notFound("Issue not found");
80038000
await assertAssignableAgent(db, issueCompany.companyId, agentId, { kind: "work" });
80048001

80058002
const now = new Date();
8006-
const scheduleMonitor = options.scheduleMonitor ?? true;
8003+
const initialActiveMonitorDeadline = activeTypedIssueMonitorDeadline(issueCompany, now);
80078004
const activePauseHold = await treeControlSvc.getActivePauseHoldGate(issueCompany.companyId, id);
80088005
if (
80098006
activePauseHold &&
@@ -8053,7 +8050,7 @@ export function issueService(db: Db) {
80538050
assigneeUserId: null,
80548051
checkoutRunId,
80558052
...(checkoutRunId
8056-
? executionLockAcquisitionFields(checkoutRunId, now, { scheduleMonitor })
8053+
? executionLockAcquisitionFields(checkoutRunId, now, initialActiveMonitorDeadline)
80578054
: { executionRunId: null }),
80588055
status: "in_progress",
80598056
startedAt: now,
@@ -8082,6 +8079,8 @@ export function issueService(db: Db) {
80828079
assigneeAgentId: issues.assigneeAgentId,
80838080
checkoutRunId: issues.checkoutRunId,
80848081
executionRunId: issues.executionRunId,
8082+
executionPolicy: issues.executionPolicy,
8083+
executionState: issues.executionState,
80858084
})
80868085
.from(issues)
80878086
.where(eq(issues.id, id))
@@ -8101,7 +8100,7 @@ export function issueService(db: Db) {
81018100
.set({
81028101
checkoutRunId,
81038102
...(checkoutRunId
8104-
? executionLockAcquisitionFields(checkoutRunId, now, { scheduleMonitor })
8103+
? executionLockAcquisitionFields(checkoutRunId, now, activeTypedIssueMonitorDeadline(current, now))
81058104
: { executionRunId: null }),
81068105
updatedAt: new Date(),
81078106
})
@@ -8131,7 +8130,6 @@ export function issueService(db: Db) {
81318130
actorAgentId: agentId,
81328131
actorRunId: checkoutRunId,
81338132
expectedCheckoutRunId: current.checkoutRunId,
8134-
scheduleMonitor,
81358133
});
81368134
if (staleAdoption.adopted) {
81378135
const row = await db.select().from(issues).where(eq(issues.id, id)).then((rows) => rows[0] ?? null);
@@ -8156,7 +8154,7 @@ export function issueService(db: Db) {
81568154
const adoptionSet: Record<string, unknown> = {
81578155
assigneeAgentId: agentId,
81588156
checkoutRunId,
8159-
...executionLockAcquisitionFields(checkoutRunId, now, { scheduleMonitor }),
8157+
...executionLockAcquisitionFields(checkoutRunId, now, activeTypedIssueMonitorDeadline(current, now)),
81608158
executionAgentNameKey: null,
81618159
status: "in_progress",
81628160
updatedAt: now,
@@ -8271,7 +8269,6 @@ export function issueService(db: Db) {
82718269
issueId: id,
82728270
actorAgentId,
82738271
actorRunId: actorRunId!,
8274-
scheduleMonitor: true,
82758272
});
82768273

82778274
if (adopted) {
@@ -8298,7 +8295,6 @@ export function issueService(db: Db) {
82988295
actorAgentId,
82998296
actorRunId,
83008297
expectedCheckoutRunId: previousCheckoutRunId,
8301-
scheduleMonitor: true,
83028298
});
83038299

83048300
if (staleAdoption.adopted) {

server/src/services/recovery/issue-graph-liveness.ts

Lines changed: 2 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import { getAgentWorkEligibility, isAgentInvokable } from "@paperclipai/shared";
22
import { buildIssueGraphLivenessIncidentKey } from "./origins.js";
3+
import { activeTypedIssueMonitorDeadline } from "../issue-execution-policy.js";
34

45
export type IssueLivenessSeverity = "warning" | "critical";
56

@@ -192,20 +193,7 @@ function monitorFromIssue(issue: IssueLivenessIssueInput) {
192193
}
193194

194195
export function hasScheduledIssueMonitorPath(issue: IssueLivenessIssueInput, now: Date | string | number) {
195-
const nowMs = typeof now === "number" ? now : readDateMs(now) ?? Date.now();
196-
const nextCheckAtMs = readDateMs(issue.monitorNextCheckAt);
197-
if (nextCheckAtMs === null || nextCheckAtMs <= nowMs) return false;
198-
199-
const { policyMonitor, stateMonitor } = monitorFromIssue(issue);
200-
const timeoutAtMs = readDateMs(policyMonitor?.timeoutAt ?? stateMonitor?.timeoutAt);
201-
if (timeoutAtMs !== null && timeoutAtMs <= nowMs) return false;
202-
203-
const maxAttempts = readPositiveInteger(policyMonitor?.maxAttempts ?? stateMonitor?.maxAttempts);
204-
const stateAttemptCount = readPositiveInteger(stateMonitor?.attemptCount) ?? 0;
205-
const attemptCount = issue.monitorAttemptCount ?? stateAttemptCount;
206-
if (maxAttempts !== null && attemptCount >= maxAttempts) return false;
207-
208-
return true;
196+
return activeTypedIssueMonitorDeadline(issue, now) !== null;
209197
}
210198

211199
export function classifyIssueReviewPaths(

server/src/services/recovery/service.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ import { issueTreeControlService } from "../issue-tree-control.js";
4343
import { ISSUE_EXECUTION_LOCK_TTL_MS, TERMINAL_HEARTBEAT_RUN_STATUSES, issueService } from "../issues.js";
4444
import {
4545
applyIssueMonitorPolicyTransition,
46+
activeTypedIssueMonitorDeadline,
4647
normalizeIssueExecutionPolicy,
4748
parseIssueExecutionState,
4849
} from "../issue-execution-policy.js";
@@ -942,7 +943,7 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup })
942943
}
943944

944945
async function hasPersistedDurableWaitPath(issue: typeof issues.$inferSelect) {
945-
if (issue.monitorNextCheckAt) return true;
946+
if (activeTypedIssueMonitorDeadline(issue)) return true;
946947

947948
return db
948949
.select({ id: issueRelations.issueId })

0 commit comments

Comments
 (0)