mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-10 20:50:08 +02:00
fix(runner): restore task monitors and durable timed waits (#15446)
## Thinking Path > - Paperclip manages AI agents and their task execution. > - Agents need a durable way to return to work after a delayed check. > - The issue monitor scheduler already provides a one-shot wake for an assignee. > - Native runners reject generic execution-policy writes and had no bound monitor tool. > - Scheduling alone is insufficient because native completion also needs to accept a timed wait. > - This pull request adds an authorized monitor tool and connects it to completion and the existing scheduler. > - An agent can now schedule its next check, end the run, and resume on the same task. ## Linked Issues or Issue Description **What happened?** A native runner could not set its own task monitor. `call_api` correctly rejected execution-policy writes, while `schedule_wake` had no production binding. `paperclip_finish` also rejected monitor waits. **Expected behavior** A standard native run can set a one-shot monitor on its current task or another accessible task assigned to the same agent. After a confirmed schedule on the current task, it can yield. The scheduler later delivers `issue_monitor_due`. **Steps to reproduce** 1. Start a standard native task. 2. Ask the agent to check the task again later and end its current run. 3. Inspect available tools and try the generic issue execution-policy update. 4. Observe the missing native tool and the lifecycle-write denial. Related PRs: #14680 concerns monitor notes in the shared wake prompt. #11919 changes attempt-limit scope. This PR adds native scheduling and completion authority and retains the existing cumulative attempt bounds. It does not depend on either PR. ## What Changed - Add provider-neutral `set_task_monitor` with a default current-task target, future timestamp, required notes, existing bounds, and explicit clearing. - Check company, task visibility, ownership, runtime permissions, work mode, and active-run authority. Preserve review-only restrictions. Reject the reserved server-owned quota-recovery name before saving or accepting a native wait. - Commit the monitor, audit event, and retry receipt together. Retry receipts survive a successor run without re-arming cleared or consumed timers. - Permit `paperclip_finish` to yield to a persisted monitor. Recheck ownership and the schedule when committing final disposition. Release execution without an immediate continuation. - Preserve due monitors during native execution. Fence wake admission and consumption against replacement, clearing, reassignment, and completion. Preserve unrelated review policy. - Expose scheduled and consumed monitor instructions in task context. Update provider schemas, Rust validation, generated contracts, and execution documentation. - Add an opt-in live Codex smoke script with isolated data and explicit run/session/runner/process evidence. ## Verification - Repository `pnpm -r typecheck` and `pnpm build` passed after rebase. Server typecheck passed again after review fixes. All CI test shards pass on `3def77b1b`, including runner TypeScript/Rust, server, serialized server, workspace, and browser tests. All CI gates are green, including the canary dry run. Greptile is 5/5 on the same commit with zero unresolved threads. - The local monolithic `pnpm test:run`, started before the rebase, was interrupted after current-head CI test coverage passed. It is not counted as a standalone full-suite pass; the focused local regression suites passed. - Targeted server tests cover scheduling, replacement, clearing, policy preservation, cumulative bounds, cross-run retries, permissions, provider-neutral discovery, review restrictions, completion authority, and scheduler/finalizer races. - Runner contract/catalog/semantic tests and Rust terminal-tool tests cover the new operation and monitor completion. - Live Codex test passed twice (latest live run on `e0bcd63e6`) in a temporary database and workspace, with a 300,000 ms warm window. First run `67bd7709-c089-4d4a-9d2b-0d6b618a34b0` yielded at `2026-10-07T13:15:03.274Z`. Second run `97c323f9-595a-4cc5-a007-db5a2fbb937c` started at `13:15:30.952Z`, received `issue_monitor_due`, and completed the same task. Exactly one monitor wake was recorded. - Both live runs used native session `7b1dd753-1c9b-4e7a-b22f-a125dbc3748c`, runner `e90d9a1b-3502-4ee3-b15e-edc024c555d4`, provider session `01a11680-6d01-70c0-9a55-db7246ed66c3`, and PID `64218` with the same process start time. This proves warm reuse for that local Codex test, not only successful scheduling. - Reproduce the paid live test with `node --import ./server/node_modules/tsx/dist/loader.mjs server/scripts/smoke-native-task-monitor.ts --run`, with the installed Codex binary on `PATH` and a valid local login. ## Risks - The scheduler now defers monitor dispatch while the task has an active native run. A stuck run still depends on the existing recovery lifecycle. - Idempotency uses the existing run ledger; no table or migration is added. - Other providers share the tested tool and completion contracts. Only Codex received a live model test. - Existing `call_api` lifecycle restrictions remain enforced. Monitor waits do not bypass task blockers, reviews, or approvals. ## Model Used OpenAI Codex, GPT-6 family, with tool use, code execution, and TypeScript/Rust editing. The session does not expose the exact deployed model ID or context-window size. ## 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
b315580640
commit
8cfedd7df8
45 files changed
+1157
-207
No files matched your search
@@ -0,0 +1,98 @@
|
||||
/** Opt-in live qualification: creates an isolated database/workspace and spends two Codex turns.
|
||||
* PATH=<pnpm-bin>:<codex-bin>:$PATH node --import tsx server/scripts/smoke-native-task-monitor.ts --run
|
||||
*/
|
||||
import assert from "node:assert/strict";
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { mkdtemp, mkdir, writeFile } from "node:fs/promises";
|
||||
import { createServer } from "node:http";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { setTimeout as delay } from "node:timers/promises";
|
||||
import { and, eq } from "drizzle-orm";
|
||||
|
||||
if (!process.argv.includes("--run")) throw new Error("Pass --run to authorize the live two-turn Codex smoke.");
|
||||
const root = await mkdtemp(join(tmpdir(), "paperclip-live-task-monitor-"));
|
||||
process.env.PAPERCLIP_HOME = join(root, "paperclip");
|
||||
process.env.PAPERCLIP_INSTANCE_ID = "monitor-smoke";
|
||||
process.env.PAPERCLIP_TELEMETRY_ENABLED = "false";
|
||||
const { createDb, companies, agents, authUsers, companyMemberships, issues, heartbeatRuns, agentWakeupRequests, nativeRunResults, statusDecisions } = await import("@paperclipai/db");
|
||||
const { startEmbeddedPostgresTestDatabase } = await import("../src/__tests__/helpers/embedded-postgres.js");
|
||||
const { heartbeatService } = await import("../src/services/heartbeat.js");
|
||||
const { setupRunnerPrpWebSocketServer, runnerPrpWebSocketInternals } = await import("../src/realtime/runner-prp-ws.js");
|
||||
const { closeIdleWarmNativeSessionsForRestart } = await import("../src/services/native-runtime/native-session-executor.js");
|
||||
const temporary = await startEmbeddedPostgresTestDatabase("live-task-monitor-");
|
||||
const db = createDb(temporary.connectionString);
|
||||
const server = createServer();
|
||||
await new Promise<void>(resolve => server.listen(0, "127.0.0.1", resolve));
|
||||
const address = server.address();
|
||||
if (!address || typeof address === "string") throw new Error("TCP listener missing");
|
||||
process.env.PAPERCLIP_API_URL = `http://127.0.0.1:${address.port}`;
|
||||
setupRunnerPrpWebSocketServer(server, { apiUrl: process.env.PAPERCLIP_API_URL });
|
||||
const companyId = randomUUID(), agentId = randomUUID(), issueId = randomUUID();
|
||||
const heartbeat = heartbeatService(db);
|
||||
const readRuns = () => db.select().from(heartbeatRuns).where(eq(heartbeatRuns.companyId, companyId)).orderBy(heartbeatRuns.createdAt);
|
||||
console.log(JSON.stringify({ root, companyId, agentId, issueId }));
|
||||
try {
|
||||
const cwd = join(root, "workspace");
|
||||
await mkdir(cwd);
|
||||
await db.insert(companies).values({ id: companyId, name: "Isolated monitor smoke", issuePrefix: "MON", defaultResponsibleUserId: "monitor-smoke", requireBoardApprovalForNewAgents: false });
|
||||
await db.insert(authUsers).values({ id: "monitor-smoke", name: "Monitor smoke", email: "monitor-smoke@example.test", createdAt: new Date(), updatedAt: new Date() });
|
||||
await db.insert(companyMemberships).values({ companyId, principalType: "user", principalId: "monitor-smoke", status: "active", membershipRole: "owner" });
|
||||
await db.insert(agents).values({ id: agentId, companyId, name: "Monitor verifier", status: "active", adapterType: "paperclip_runner",
|
||||
adapterConfig: { provider: "codex", cwd, lifecycleMode: "warm", idleTimeoutMs: 300_000, timeoutSeconds: 180 },
|
||||
runtimeConfig: { heartbeat: { enabled: false, wakeOnDemand: true } } });
|
||||
await db.insert(issues).values({ id: issueId, companyId, title: "Verify a one-shot native task monitor", identifier: "MON-1", status: "in_progress", assigneeAgentId: agentId,
|
||||
description: "This is an isolated two-run integration test. First run: use set_task_monitor on this task with a fresh nextCheckAt about 45 seconds in the future (compute the current UTC time if needed), notes saying to verify issue_monitor_due, and idempotencyKey monitor-live-first. Confirm the receipt. Then call paperclip_finish with yielded and continuation.kind monitor, accurately retaining the second run as blocking remaining work. End immediately after acceptance. Do not sleep or poll. Second run: if the wake reason is issue_monitor_due, the requested check succeeded. Report that reason and finish done using the current completion contract. Do not schedule another monitor or request human review. No files or other deliverables are required." });
|
||||
await heartbeat.wakeup(agentId, { source: "on_demand", triggerDetail: "manual", reason: "monitor_smoke",
|
||||
requestedByActorType: "system", requestedByActorId: "monitor-smoke", payload: { issueId }, contextSnapshot: { issueId, wakeReason: "monitor_smoke" } });
|
||||
const deadline = Date.now() + 8 * 60_000;
|
||||
let firstReleasedAt: number | null = null;
|
||||
while (Date.now() < deadline) {
|
||||
const runs = await readRuns();
|
||||
if (runs.some(run => ["failed", "timed_out", "cancelled"].includes(run.status))) {
|
||||
throw new Error(`Live run failed: ${JSON.stringify(runs.map(run => ({ id: run.id, status: run.status, error: run.error, errorCode: run.errorCode })))}`);
|
||||
}
|
||||
if (runs.length === 1 && runs[0].status === "succeeded") {
|
||||
firstReleasedAt ??= Date.now();
|
||||
const [issue] = await db.select().from(issues).where(eq(issues.id, issueId));
|
||||
assert.equal(issue.status, "in_progress");
|
||||
assert.ok(issue.monitorNextCheckAt, "first run must leave a persisted monitor");
|
||||
assert.equal(issue.executionRunId, null);
|
||||
}
|
||||
if (runs.length >= 2 && runs[1].status === "succeeded") break;
|
||||
await heartbeat.tickTimers();
|
||||
await delay(1_000);
|
||||
}
|
||||
await heartbeat.drainActiveRunExecutions();
|
||||
const runs = await readRuns();
|
||||
assert.equal(runs.length, 2, "exactly two server runs must execute");
|
||||
assert.ok(runs.every(run => run.status === "succeeded"));
|
||||
const [issue] = await db.select().from(issues).where(eq(issues.id, issueId));
|
||||
const wakes = await db.select().from(agentWakeupRequests).where(and(eq(agentWakeupRequests.companyId, companyId), eq(agentWakeupRequests.reason, "issue_monitor_due")));
|
||||
assert.equal(wakes.length, 1);
|
||||
assert.equal(runs[1].contextSnapshot?.wakeReason, "issue_monitor_due");
|
||||
assert.equal(issue.status, "done");
|
||||
assert.equal(issue.monitorNextCheckAt, null);
|
||||
assert.ok(firstReleasedAt && runs[1].createdAt.getTime() - firstReleasedAt < 300_000, "wake must arrive within the configured warm window");
|
||||
const results = await db.select().from(nativeRunResults).where(eq(nativeRunResults.companyId, companyId));
|
||||
const decisions = await db.select().from(statusDecisions).where(eq(statusDecisions.companyId, companyId));
|
||||
const evidence = { companyId, issueId, agentId, configuredIdleTimeoutMs: 300_000, firstReleasedAt, status: issue.status,
|
||||
runs: runs.map(run => ({ id: run.id, status: run.status, createdAt: run.createdAt, finishedAt: run.finishedAt,
|
||||
wakeReason: run.contextSnapshot?.wakeReason, nativeSessionId: run.nativeSessionId, runnerInstanceId: run.runnerInstanceId,
|
||||
processPid: run.processPid, processStartedAt: run.processStartedAt, providerSessionId: run.sessionIdAfter,
|
||||
semanticToolReceipts: run.resultJson?.semanticToolReceipts })),
|
||||
warmReuse: { sameNativeSession: runs[0].nativeSessionId === runs[1].nativeSessionId,
|
||||
sameRunnerInstance: runs[0].runnerInstanceId === runs[1].runnerInstanceId,
|
||||
sameProcess: runs[0].processPid != null && runs[0].processPid === runs[1].processPid && runs[0].processStartedAt?.getTime() === runs[1].processStartedAt?.getTime() },
|
||||
results: results.map(row => ({ runId: row.runId, result: row.resultJson })),
|
||||
decisions: decisions.map(row => ({ runId: row.runId, reasonCode: row.reasonCode, toStatus: row.toStatus })),
|
||||
monitorWakeIds: wakes.map(row => row.id) };
|
||||
await writeFile(join(root, "evidence.json"), JSON.stringify(evidence, null, 2));
|
||||
console.log(JSON.stringify({ passed: true, evidence: join(root, "evidence.json"), warmReuse: evidence.warmReuse }));
|
||||
} finally {
|
||||
await writeFile(join(root, "run-records.json"), JSON.stringify(await readRuns(), null, 2));
|
||||
await closeIdleWarmNativeSessionsForRestart();
|
||||
runnerPrpWebSocketInternals.resetForTests();
|
||||
await new Promise<void>(resolve => server.close(() => resolve()));
|
||||
await temporary.cleanup();
|
||||
}
|
||||
Reference in new issue
Block a user