mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-02 02:07:25 +08:00
fix: stop remote Grok runs before continuing queued messages (#14100)
## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work. > - The execution service owns each run and saves messages sent while it runs. > - Interrupt must stop the current executor before it delivers those messages. > - Remote Grok commands did not register the host cancellation control. > - A cancelled task run could still write Done and prevent queue recovery. > - This pull request connects remote cancellation and revokes cancelled run writes. > - Saved input can use the existing queue admission rules after verified cleanup. ## Linked Issues or Issue Description **What happened?** Interrupting a queued message marked a remote Grok run cancelled before its sandbox stopped. The old run could still post a reply and mark the task Done. Its saved follow-up remained deferred behind execution recovery. **Expected behavior** Stop revokes run write authority and waits for verified termination. Saved messages remain durable and enter one successor through normal admission after cleanup. **Steps to reproduce** 1. Run a task with `grok_local` in a remote sandbox. 2. Send a follow-up and use Interrupt while the command runs. 3. Let the old command attempt a task status update after cancellation. 4. Observe the task disposition and the saved message queue. **Paperclip version or commit** The gap is present in master at `d3e0f0a238`. **Deployment mode** Authenticated server with a Daytona sandbox. Related work: #14028 and #14046 handle bounded continuation. #13291 covers infrastructure interruption and verified remote cleanup. #13332 addresses atomic recovery holds. This change handles direct Grok operator cancellation and stale task writes. ## What Changed - Register remote Grok cancellation before preparation. Keep command ownership until the host confirms sandbox termination. - Reuse the sandbox cancellation boundary for the direct CLI invocation. Reject fresh attempts after cancellation and preserve workspace restore failure evidence. - Reject writes from cancelled task JWTs and runs with a pending stop. Preserve diagnostic reads and existing conversation error codes. - Recheck run authority under a database lock before task updates and interaction responses commit. - Preserve authorized handoffs that stop their own run. Only the server-issued stop receipt for that request permits the final task update. - Add tests for hung commands, unverified stops, early cancellation, copy-back failures, late Done, late interaction responses, authorized handoffs, exact lease receipts, and one queue successor across concurrent restart sweeps. - Document the cancellation and write-authority contract. ## Verification - Targeted adapter, cancellation-boundary, authentication, queued-message, interaction-service, and activity-route tests passed. The expanded run passed 214 tests; one new test had an incomplete fixture. After correcting the fixture, all 8 selected follow-up cases passed. - `pnpm -r typecheck`: passed on `179c86caf1bf0d89914a503d46e24af7e4b8c557`. - `pnpm build`: passed on the same commit. - `pnpm test:run`: the general-server group completed with 13,521 passed, 99 skipped, and 18 failed tests. It then stopped, so the remaining local groups did not run. Five Slack, email, and wake-batching failures passed on focused reruns after correcting the local environment. The remaining 13 failures reproduce as `EACCES` on rename in unchanged skill-cache code on macOS. Two custom-image suite setup hooks also failed to start embedded PostgreSQL after the machine exhausted shared-memory slots; all 31 tests in that file passed on rerun after the local resource issue was resolved. CI covers all test groups. - CI: 53 checks passed and 2 were skipped on the latest commit, including the aggregate verification gate. The last server shard passed on its single rerun after a preview-server startup timeout. The affected file also passed locally with 28 passed and 3 skipped. - Greptile: 5/5 on the latest commit. Both review threads are resolved. - No live deployment or staging task mutation has been performed. ## Risks - Stopping the sandbox can prevent file copy-back. The result preserves workspace restore failure evidence; termination does not imply restored files. - If provider termination fails, the adapter keeps ownership of its outstanding command and does not acknowledge Stop. - The write restriction now applies to ordinary cancelled tasks. Reads remain allowed. Task and interaction checks add a shared run-row lock to agent mutations. An exact server-issued receipt permits the task request that stopped its own run to complete its handoff. - Existing terminal tasks are not reopened automatically. An operator must correct a historical late Done before its saved queue can continue. - No schema migration or UI change. ## Model Used OpenAI GPT-6 through Codex, with reasoning, repository inspection, code execution, and test tools. The precise backend revision and context-window size are not exposed in this session. ## 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 #` 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:
@@ -981,6 +981,8 @@ Real gates still apply: company and task ownership, active provider ownership, b
|
||||
|
||||
An operator Stop waits for provider termination. Remote sandbox providers may return a stopped/deleted receipt after their control-plane operation completes. Paperclip binds that receipt to the company, run, and exact lease; successful file cleanup, a terminal run row, or an in-sandbox shutdown event is not sufficient. Legacy conversational runs receive their cancellation acknowledgement after all remote leases have confirmed termination. Stop alone never creates a continuation. A user message queued during remote cleanup is reconsidered when the provider confirms termination; it still passes normal admission and adopts pending comment IDs in order. Once stopped, the next explicit wake uses the same queue. A compatible saved ACP session can resume, and an unavailable or incompatible session can start fresh with the full task context. Run credentials and scratch paths remain scoped to the new run. A subtree pause requires Resume; a message does not bypass it.
|
||||
|
||||
Cancelled runs and runs with a recorded pending stop lose run-scoped API write authority, while reads remain available for diagnostics. This applies to ordinary tasks as well as conversations. Task and interaction mutations recheck that authority in the write transaction, so a cancelled run cannot overwrite a recovery disposition with a late Done or resolve an interaction after revocation. A task handoff that intentionally stops its own run can commit only with a server receipt tied to that exact request. Direct sandbox CLI adapters such as Grok register cancellation before preparation and keep ownership until host-owned termination is verified, including a sandbox acquired before adapter registration. A failed stop does not acknowledge cancellation or abandon an outstanding remote command. Interrupted workspace restore failures remain recorded for recovery; a stop receipt proves termination, not successful file restoration.
|
||||
|
||||
For native conversations, an authenticated user message sent after the previous run finishes can retire its execution recovery holds and start a fresh turn. Hold retirement and the new run are atomic. The previous transcript, tool outcomes, and recovery history remain intact. This starts a new conversation; it does not replay tool calls with unknown outcomes.
|
||||
|
||||
Local recovery records a server-authored stop receipt before it clears a verified absent process identity. A new execution request invalidates that receipt before any process can spawn; recording a new process identity also invalidates it. Missing process IDs without a receipt still block admission. Remote execution continues to require termination receipts for every lease.
|
||||
|
||||
@@ -8,10 +8,11 @@ export function cancellableSandboxStartup(ctx: AdapterExecutionContext) {
|
||||
const signal = ctx.signal;
|
||||
const stop = ctx.stopRemoteStartup;
|
||||
if (!signal || !stop || target?.kind !== "remote" || target.transport !== "sandbox" || !target.runner) {
|
||||
return { context: ctx, finish: async () => {} };
|
||||
return { context: ctx, stopAcknowledged: () => false, finish: async () => {} };
|
||||
}
|
||||
let armed = true;
|
||||
let stopping: Promise<void> | undefined;
|
||||
let stopAcknowledged = false;
|
||||
const inFlight = new Set<Promise<unknown>>();
|
||||
let rejectStopped!: (error: unknown) => void;
|
||||
const stopped = new Promise<never>((_, reject) => { rejectStopped = reject; });
|
||||
@@ -21,7 +22,10 @@ export function cancellableSandboxStartup(ctx: AdapterExecutionContext) {
|
||||
if (stopping) return;
|
||||
stopping = Promise.resolve().then(stop);
|
||||
void stopping.then(
|
||||
() => rejectStopped(signal.reason ?? new Error("Stopped during sandbox startup")),
|
||||
() => {
|
||||
stopAcknowledged = true;
|
||||
rejectStopped(signal.reason ?? new Error("Sandbox execution stopped"));
|
||||
},
|
||||
// Without proof, the original operation still owns its resources. Do
|
||||
// not abandon it or release credentials while it could be running.
|
||||
() => {},
|
||||
@@ -78,6 +82,7 @@ export function cancellableSandboxStartup(ctx: AdapterExecutionContext) {
|
||||
};
|
||||
return {
|
||||
context: { ...ctx, executionTarget: { ...target, runner } },
|
||||
stopAcknowledged: () => stopAcknowledged,
|
||||
async finish() {
|
||||
signal.removeEventListener("abort", onAbort);
|
||||
armed = false;
|
||||
|
||||
@@ -199,7 +199,7 @@ export interface AdapterExecutionContext {
|
||||
signal?: AbortSignal;
|
||||
/** Opt in to signal-based cancellation before starting provider work. */
|
||||
onCancellationReady?: () => Promise<void>;
|
||||
/** Host-owned stop of this run's sandbox during setup. Resolves only after
|
||||
/** Host-owned stop of this run's sandbox during setup or direct CLI execution. Resolves only after
|
||||
* provider termination is verified; never accepts an agent-selected lease. */
|
||||
stopRemoteStartup?: () => Promise<void>;
|
||||
/** Server-owned, actor-attributed snapshot also rendered by legacy wake prompts. */
|
||||
|
||||
@@ -62,8 +62,8 @@ vi.mock("@paperclipai/adapter-utils/execution-target", () => ({
|
||||
(mocks.ensureRuntimeInstalledMock as (...args: unknown[]) => unknown)(...args),
|
||||
prepareAdapterExecutionTargetRuntime: (...args: unknown[]) =>
|
||||
(mocks.prepareRuntimeMock as (...args: unknown[]) => unknown)(...args),
|
||||
readAdapterExecutionTarget: () =>
|
||||
mocks.state.isRemote ? { kind: "remote", transport: "ssh" } : { kind: "local" },
|
||||
readAdapterExecutionTarget: (input: { executionTarget?: unknown }) => input.executionTarget ??
|
||||
(mocks.state.isRemote ? { kind: "remote", transport: "ssh" } : { kind: "local" }),
|
||||
resolveAdapterExecutionTargetCommandForLogs: (...args: unknown[]) =>
|
||||
(mocks.resolveCommandForLogsMock as (...args: unknown[]) => unknown)(...args),
|
||||
resolveAdapterExecutionTargetTimeoutSec: (_target: unknown, timeoutSec: number) => timeoutSec,
|
||||
@@ -128,6 +128,8 @@ function makeRestoreWorkspace(
|
||||
|
||||
function makeSuccessfulRunResult(overrides: Partial<{ sessionId: string }> = {}) {
|
||||
return {
|
||||
pid: null,
|
||||
startedAt: new Date().toISOString(),
|
||||
exitCode: 0,
|
||||
signal: null,
|
||||
timedOut: false,
|
||||
@@ -187,6 +189,94 @@ describe("grok_local execute", () => {
|
||||
await Promise.all(tempRoots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true })));
|
||||
});
|
||||
|
||||
async function cancellableContext(stop: () => Promise<void>) {
|
||||
const ctx = await makeCtx("remote-cancellation", await makeTempRoot());
|
||||
const controller = new AbortController();
|
||||
const remoteExecute = vi.fn(async () => makeSuccessfulRunResult());
|
||||
ctx.signal = controller.signal;
|
||||
ctx.stopRemoteStartup = vi.fn(stop);
|
||||
ctx.onCancellationReady = vi.fn(async () => {});
|
||||
ctx.executionTarget = { kind: "remote", transport: "sandbox", providerKey: "daytona", remoteCwd: "/remote/workspace",
|
||||
runner: { execute: remoteExecute } };
|
||||
remoteState.isRemote = true;
|
||||
runProcessMock.mockImplementation(async (_runId, target) => {
|
||||
expect(ctx.onCancellationReady).toHaveBeenCalledOnce();
|
||||
return target.runner.execute({ command: "grok" });
|
||||
});
|
||||
return { ctx, controller, remoteExecute };
|
||||
}
|
||||
|
||||
it("settles remote cancellation only after the sandbox stop receipt, even if the command RPC hangs", async () => {
|
||||
let confirmStop!: () => void;
|
||||
const receipt = new Promise<void>(resolve => { confirmStop = resolve; });
|
||||
const f = await cancellableContext(() => receipt);
|
||||
f.remoteExecute.mockImplementation(() => new Promise(() => {}));
|
||||
let settled = false;
|
||||
const execution = execute(f.ctx).finally(() => { settled = true; });
|
||||
await vi.waitFor(() => expect(f.remoteExecute).toHaveBeenCalledOnce());
|
||||
f.controller.abort(new Error("Interrupted to send queued messages"));
|
||||
await vi.waitFor(() => expect(f.ctx.stopRemoteStartup).toHaveBeenCalledOnce());
|
||||
expect(settled).toBe(false);
|
||||
confirmStop();
|
||||
expect(await execution).toMatchObject({ errorCode: "cancelled",
|
||||
resultJson: { executionCancellation: { state: "acknowledged" } } });
|
||||
expect(runProcessMock).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("keeps ownership of the command when remote termination cannot be verified", async () => {
|
||||
const f = await cancellableContext(async () => { throw new Error("stop unverified"); });
|
||||
let finishCommand!: (result: ReturnType<typeof makeSuccessfulRunResult>) => void;
|
||||
f.remoteExecute.mockImplementation(() => new Promise(resolve => { finishCommand = resolve; }));
|
||||
let settled = false;
|
||||
const execution = execute(f.ctx).catch(error => error).finally(() => { settled = true; });
|
||||
await vi.waitFor(() => expect(f.remoteExecute).toHaveBeenCalledOnce());
|
||||
f.controller.abort();
|
||||
await vi.waitFor(() => expect(f.ctx.stopRemoteStartup).toHaveBeenCalledOnce());
|
||||
expect(settled).toBe(false);
|
||||
finishCommand(makeSuccessfulRunResult());
|
||||
expect(await execution).toEqual(new Error("stop unverified"));
|
||||
expect(runProcessMock).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("does not start Grok when cancellation was requested before registration", async () => {
|
||||
const f = await cancellableContext(async () => {});
|
||||
f.ctx.onCancellationReady = vi.fn(async () => { f.controller.abort(); });
|
||||
expect(await execute(f.ctx)).toMatchObject({ errorCode: "cancelled",
|
||||
executionRecovery: { kind: "bootstrap", providerWorkStarted: false },
|
||||
resultJson: { executionCancellation: { state: "acknowledged" } } });
|
||||
expect(runProcessMock).not.toHaveBeenCalled();
|
||||
expect(prepareRuntimeMock).not.toHaveBeenCalled();
|
||||
expect(f.ctx.stopRemoteStartup).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("does not acknowledge an early cancellation when its acquired sandbox cannot stop", async () => {
|
||||
const f = await cancellableContext(async () => { throw new Error("stop unverified"); });
|
||||
f.ctx.onCancellationReady = vi.fn(async () => { f.controller.abort(); });
|
||||
await expect(execute(f.ctx)).rejects.toThrow("stop unverified");
|
||||
expect(runProcessMock).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("retains workspace recovery evidence when a confirmed stop prevents copy-back", async () => {
|
||||
const f = await cancellableContext(async () => {});
|
||||
prepareRuntimeMock.mockImplementationOnce(async () => ({ workspaceRemoteDir: "/remote/workspace", assetDirs: {},
|
||||
restoreWorkspace: async () => { throw new Error("sandbox stopped during restore"); },
|
||||
}));
|
||||
f.remoteExecute.mockImplementation(() => new Promise(() => {}));
|
||||
const execution = execute(f.ctx);
|
||||
await vi.waitFor(() => expect(f.remoteExecute).toHaveBeenCalledOnce());
|
||||
f.controller.abort();
|
||||
const result = await execution;
|
||||
expect(result.resultJson?.executionCancellation).toMatchObject({ state: "acknowledged" });
|
||||
expect(result.resultJson?.workspaceRestoreFailure).toBeTruthy();
|
||||
});
|
||||
|
||||
it("does not stop the sandbox after a normal completed turn", async () => {
|
||||
const f = await cancellableContext(async () => {});
|
||||
expect(await execute(f.ctx)).toMatchObject({ exitCode: 0 });
|
||||
f.controller.abort();
|
||||
expect(f.ctx.stopRemoteStartup).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("stages Grok-native instructions and skills into the workspace for the run and cleans them up afterward", async () => {
|
||||
const root = await makeTempRoot();
|
||||
const instructionsPath = path.join(root, "managed", "AGENTS.md");
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import { withWorkspaceRestore } from "@paperclipai/adapter-utils/workspace-restore-result";
|
||||
import { cancellableSandboxStartup } from "@paperclipai/adapter-utils/acpx-engine/startup-cancellation";
|
||||
import fs from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import { fileURLToPath } from "node:url";
|
||||
@@ -196,6 +197,55 @@ function resolveBillingType(env: Record<string, string>): "api" | "subscription"
|
||||
}
|
||||
|
||||
export async function execute(ctx: AdapterExecutionContext): Promise<AdapterExecutionResult> {
|
||||
const target = ctx.executionTarget;
|
||||
if (!ctx.signal || !ctx.stopRemoteStartup || target?.kind !== "remote" || target.transport !== "sandbox" || !target.runner) {
|
||||
return executeTurn(ctx);
|
||||
}
|
||||
|
||||
// Direct remote commands have no host child process to kill. Register before
|
||||
// setup and retain ownership until the host verifies this sandbox has stopped.
|
||||
await ctx.onCancellationReady?.();
|
||||
const cancelled = (result?: AdapterExecutionResult): AdapterExecutionResult => ({
|
||||
exitCode: null,
|
||||
signal: null,
|
||||
timedOut: false,
|
||||
...result,
|
||||
errorCode: "cancelled",
|
||||
errorMessage: "Grok execution was cancelled",
|
||||
resultJson: {
|
||||
...result?.resultJson,
|
||||
executionCancellation: { state: "acknowledged", acknowledgedAt: new Date().toISOString() },
|
||||
},
|
||||
});
|
||||
if (ctx.signal.aborted) {
|
||||
// The host may already have acquired a lease before adapter registration.
|
||||
await ctx.stopRemoteStartup();
|
||||
return { ...cancelled(), executionRecovery: { kind: "bootstrap", providerWorkStarted: false } };
|
||||
}
|
||||
// Keep the existing setup boundary armed for the whole direct CLI invocation:
|
||||
// unlike ACP adapters, Grok has no turn-level cancellation protocol.
|
||||
const cancellation = cancellableSandboxStartup(ctx);
|
||||
let result: AdapterExecutionResult | undefined;
|
||||
let failure: unknown;
|
||||
let failed = false;
|
||||
try {
|
||||
result = await executeTurn(cancellation.context);
|
||||
} catch (error) {
|
||||
failure = error;
|
||||
failed = true;
|
||||
}
|
||||
try {
|
||||
await cancellation.finish();
|
||||
} catch (error) {
|
||||
failure = error;
|
||||
failed = true;
|
||||
}
|
||||
if (cancellation.stopAcknowledged()) return cancelled(result);
|
||||
if (failed) throw failure;
|
||||
return result!;
|
||||
}
|
||||
|
||||
async function executeTurn(ctx: AdapterExecutionContext): Promise<AdapterExecutionResult> {
|
||||
const { runId, agent, runtime, config, context, onLog, onMeta, onSpawn, authToken } = ctx;
|
||||
const executionTarget = readAdapterExecutionTarget({
|
||||
executionTarget: ctx.executionTarget,
|
||||
@@ -543,6 +593,7 @@ export async function execute(ctx: AdapterExecutionContext): Promise<AdapterExec
|
||||
};
|
||||
|
||||
const runAttempt = async (resumeSessionId: string | null) => {
|
||||
ctx.signal?.throwIfAborted();
|
||||
const prompt = joinPromptSections([
|
||||
selectInitialCommunicationGuidance(context, { resumedSession: Boolean(resumeSessionId) }),
|
||||
basePrompt,
|
||||
@@ -655,6 +706,7 @@ export async function execute(ctx: AdapterExecutionContext): Promise<AdapterExec
|
||||
};
|
||||
|
||||
const initial = await runAttempt(sessionId);
|
||||
ctx.signal?.throwIfAborted();
|
||||
if (
|
||||
sessionId &&
|
||||
!initial.proc.timedOut &&
|
||||
|
||||
@@ -33,7 +33,8 @@ function createSelectChain(rowsForTable: (table: unknown) => unknown[]) {
|
||||
function createDbState(input: {
|
||||
agent: { id: string; companyId: string; status?: string };
|
||||
agentKey?: { id: string; agentId: string; companyId: string; keyHash: string; responsibleUserId?: string | null };
|
||||
run?: { id: string; companyId: string; agentId: string; responsibleUserId?: string | null };
|
||||
run?: { id: string; companyId: string; agentId: string; responsibleUserId?: string | null;
|
||||
status?: string; contextSnapshot?: Record<string, unknown>; resultJson?: Record<string, unknown> };
|
||||
}) {
|
||||
const activity: Array<Record<string, unknown>> = [];
|
||||
const agentRow = {
|
||||
@@ -58,6 +59,9 @@ function createDbState(input: {
|
||||
companyId: input.run.companyId,
|
||||
agentId: input.run.agentId,
|
||||
responsibleUserId: input.run.responsibleUserId ?? null,
|
||||
status: input.run.status ?? "running",
|
||||
contextSnapshot: input.run.contextSnapshot ?? {},
|
||||
resultJson: input.run.resultJson ?? {},
|
||||
}
|
||||
: null;
|
||||
|
||||
@@ -178,6 +182,30 @@ describe("agent auth middleware", () => {
|
||||
else process.env.PAPERCLIP_INSTANCE_ID = originalInstanceId;
|
||||
});
|
||||
|
||||
it.each([
|
||||
{ status: "cancelled", conversationMode: false, requested: false },
|
||||
{ status: "running", conversationMode: false, requested: true },
|
||||
{ status: "cancelled", conversationMode: true, requested: false },
|
||||
{ status: "running", conversationMode: true, requested: true },
|
||||
])("revokes writes but preserves reads for a stopped run: %j", async ({ status, conversationMode, requested }) => {
|
||||
const agentId = randomUUID();
|
||||
const companyId = randomUUID();
|
||||
const runId = randomUUID();
|
||||
const { db } = createDbState({ agent: { id: agentId, companyId }, run: {
|
||||
id: runId, companyId, agentId, status, contextSnapshot: { conversationMode },
|
||||
resultJson: requested ? { executionCancellation: { state: "requested" } } : {},
|
||||
} });
|
||||
const token = createLocalAgentJwt(agentId, companyId, "grok_local", runId, null);
|
||||
const client = createApp(db);
|
||||
const endpoint = `/companies/${companyId}/issues/${randomUUID()}`;
|
||||
const write = await request(client).patch(endpoint).set("Authorization", `Bearer ${token}`).send({ status: "done" });
|
||||
expect(write.status).toBe(403);
|
||||
expect(write.body.code).toBe(conversationMode ? "conversation_turn_cancelled" : "agent_run_cancelled");
|
||||
const read = await request(client).get(endpoint).set("Authorization", `Bearer ${token}`);
|
||||
expect(read.status).toBe(200);
|
||||
expect(read.body.readable).toBe(true);
|
||||
});
|
||||
|
||||
it("keeps header-less local requests as the implicit board actor with their run id", async () => {
|
||||
const runId = randomUUID();
|
||||
const { db } = createDbState({ agent: { id: randomUUID(), companyId: randomUUID() } });
|
||||
|
||||
@@ -179,6 +179,8 @@ function makeIssue() {
|
||||
function issueUpdateWithReceipt(issue: ReturnType<typeof makeIssue>, patch: Record<string, unknown>) {
|
||||
const {
|
||||
actorAgentId: _actorAgentId,
|
||||
actorRunId: _actorRunId,
|
||||
actorRunStopId: _actorRunStopId,
|
||||
actorUserId: _actorUserId,
|
||||
blockedByIssueIds: _blockedByIssueIds,
|
||||
...issuePatch
|
||||
|
||||
@@ -14,6 +14,7 @@ import {
|
||||
companyMemberships,
|
||||
companySkills,
|
||||
createDb,
|
||||
environmentLeases,
|
||||
heartbeatRuns,
|
||||
heartbeatRunEvents,
|
||||
issueComments,
|
||||
@@ -25,6 +26,9 @@ import {
|
||||
import { errorHandler } from "../middleware/index.js";
|
||||
import { issueRoutes } from "../routes/issues.js";
|
||||
import { heartbeatService } from "../services/heartbeat.js";
|
||||
import { issueService } from "../services/issues.js";
|
||||
import { issueThreadInteractionService } from "../services/issue-thread-interactions.js";
|
||||
import { remoteTerminationReceipt } from "../services/remote-execution-termination.js";
|
||||
import { initializeRunIdentity, reconcileSteeredIdentity } from "../services/run-identity.js";
|
||||
import {
|
||||
getEmbeddedPostgresTestSupport,
|
||||
@@ -187,6 +191,115 @@ describeEmbeddedPostgres("issue queued-comment routes", () => {
|
||||
return { companyId, agentId, issueId, runId, wakeId, commentIds };
|
||||
}
|
||||
|
||||
it.each(["cancelled", "running"])("rejects a late Done from an interrupted %s task run at the write boundary", async status => {
|
||||
const seeded = await seedQueue();
|
||||
await db.update(heartbeatRuns).set({ status, resultJson: { executionCancellation: { state: "requested" } } })
|
||||
.where(eq(heartbeatRuns.id, seeded.runId));
|
||||
// Model a request that passed middleware before cancellation and reached
|
||||
// the service after the run lost its authority. Recovery may have released
|
||||
// checkout and changed the task to Blocked in the meantime.
|
||||
await db.update(issues).set({ status: "blocked", executionRunId: null }).where(eq(issues.id, seeded.issueId));
|
||||
await expect(issueService(db).update(seeded.issueId, { status: "done",
|
||||
actorAgentId: seeded.agentId, actorRunId: seeded.runId,
|
||||
})).rejects.toMatchObject({ status: 403, details: { code: "agent_run_cancelled" } });
|
||||
expect((await db.select().from(issues).where(eq(issues.id, seeded.issueId)))[0].status).toBe("blocked");
|
||||
// A board disposition remains authoritative.
|
||||
expect(await issueService(db).update(seeded.issueId, { status: "done", actorUserId: "queue-owner" }))
|
||||
.toMatchObject({ status: "done" });
|
||||
});
|
||||
|
||||
it("delivers a saved Grok interrupt exactly once after remote cleanup, including concurrent restart sweeps", async () => {
|
||||
const seeded = await seedQueue();
|
||||
await db.update(agents).set({ adapterType: "grok_local", runtimeConfig: { heartbeat: { maxConcurrentRuns: 1 } } })
|
||||
.where(eq(agents.id, seeded.agentId));
|
||||
await db.update(heartbeatRuns).set({ runtimeMode: "legacy", status: "cancelled", errorCode: "operator_interrupted",
|
||||
finishedAt: new Date("2026-08-22T15:03:00.000Z"),
|
||||
contextSnapshot: { issueId: seeded.issueId, paperclipWorkspace: { remoteExecution: { transport: "sandbox" } } },
|
||||
}).where(eq(heartbeatRuns.id, seeded.runId));
|
||||
await db.update(issues).set({ status: "blocked", executionRunId: null }).where(eq(issues.id, seeded.issueId));
|
||||
await db.insert(issueRecoveryActions).values({ companyId: seeded.companyId, sourceIssueId: seeded.issueId,
|
||||
kind: "active_run_watchdog", cause: "legacy_execution_requires_reconciliation", fingerprint: seeded.runId,
|
||||
status: "resolved", outcome: "blocked", nextAction: "Automatic recovery stopped.",
|
||||
evidence: { runId: seeded.runId, automaticRecovery: { replay: "blocked", actionOutcome: "unknown" } },
|
||||
});
|
||||
const [lease] = await db.insert(environmentLeases).values({ companyId: seeded.companyId,
|
||||
issueId: seeded.issueId, heartbeatRunId: seeded.runId, provider: "daytona", providerLeaseId: randomUUID(),
|
||||
status: "active", leasePolicy: "ephemeral",
|
||||
}).returning();
|
||||
// Occupy the agent's slot on another task so this test never invokes a model.
|
||||
await db.insert(heartbeatRuns).values({ companyId: seeded.companyId, agentId: seeded.agentId,
|
||||
status: "running", contextSnapshot: { issueId: randomUUID() },
|
||||
});
|
||||
const client = app(seeded.companyId);
|
||||
const queue = await request(client).get(`/api/issues/${seeded.issueId}/queued-comments`).expect(200);
|
||||
await request(client).post(`/api/issues/${seeded.issueId}/queued-comments/interrupt`).send({
|
||||
queueId: seeded.wakeId, revision: queue.body.revision, targetRunId: null,
|
||||
}).expect(200);
|
||||
expect((await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, seeded.wakeId)))[0].status)
|
||||
.toBe("deferred_issue_execution");
|
||||
// Cleanup timestamps alone are not proof. A mismatched receipt is not proof either.
|
||||
await db.update(environmentLeases).set({ status: "expired", releasedAt: new Date(), cleanupStatus: "success",
|
||||
metadata: { remoteExecutionTermination: { ...remoteTerminationReceipt(lease, {
|
||||
providerLeaseId: lease.providerLeaseId, state: "destroyed",
|
||||
}), runId: randomUUID() } },
|
||||
}).where(eq(environmentLeases.id, lease.id));
|
||||
await heartbeatService(db).resumeQueuedCommentInterrupt(seeded.companyId, seeded.wakeId);
|
||||
expect((await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, seeded.wakeId)))[0].status)
|
||||
.toBe("deferred_issue_execution");
|
||||
await db.update(environmentLeases).set({ metadata: { remoteExecutionTermination: remoteTerminationReceipt(lease, {
|
||||
providerLeaseId: lease.providerLeaseId, state: "destroyed",
|
||||
}) } }).where(eq(environmentLeases.id, lease.id));
|
||||
await db.update(agentWakeupRequests).set({ updatedAt: new Date(0) }).where(eq(agentWakeupRequests.id, seeded.wakeId));
|
||||
await Promise.all([heartbeatService(db).resumeQueuedRuns(), heartbeatService(db).resumeQueuedRuns()]);
|
||||
const [delivered] = await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, seeded.wakeId));
|
||||
expect(delivered.status).toBe("coalesced");
|
||||
const [successor] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, delivered.runId!));
|
||||
expect(successor.contextSnapshot).toMatchObject({ wakeCommentIds: seeded.commentIds, previousRunId: seeded.runId,
|
||||
forceFreshSession: true });
|
||||
await heartbeatService(db).resumeQueuedCommentInterrupt(seeded.companyId, seeded.wakeId);
|
||||
expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.companyId, seeded.companyId))).toHaveLength(3);
|
||||
});
|
||||
|
||||
it("allows only the exact task mutation that stopped its own run to commit", async () => {
|
||||
const seeded = await seedQueue();
|
||||
const stopId = randomUUID();
|
||||
await db.update(heartbeatRuns).set({ status: "cancelled", errorCode: "issue_reassigned",
|
||||
resultJson: { reassignmentStopConfirmed: true, issueMutationStopId: stopId },
|
||||
}).where(eq(heartbeatRuns.id, seeded.runId));
|
||||
const data = { status: "todo", actorAgentId: seeded.agentId, actorRunId: seeded.runId };
|
||||
await expect(issueService(db).update(seeded.issueId, { ...data, actorRunStopId: randomUUID() }))
|
||||
.rejects.toMatchObject({ status: 403, details: { code: "agent_run_cancelled" } });
|
||||
expect(await issueService(db).update(seeded.issueId, { ...data, actorRunStopId: stopId }))
|
||||
.toMatchObject({ status: "todo" });
|
||||
await expect(issueService(db).update(seeded.issueId, { ...data, status: "done" }))
|
||||
.rejects.toMatchObject({ status: 403, details: { code: "agent_run_cancelled" } });
|
||||
});
|
||||
|
||||
it.each(["questions", "verdicts"] as const)("rejects admitted %s responses after Stop revokes the run", async kind => {
|
||||
const seeded = await seedQueue();
|
||||
const service = issueThreadInteractionService(db);
|
||||
const issue = { id: seeded.issueId, companyId: seeded.companyId };
|
||||
const interaction = await service.create(issue, kind === "questions" ? {
|
||||
kind: "ask_user_questions", resolverPolicy: "board_or_agents", payload: { version: 1,
|
||||
questions: [{ id: "scope", prompt: "Choose scope", selectionMode: "single", options: [{ id: "first", label: "First" }] }],
|
||||
},
|
||||
} : {
|
||||
kind: "request_item_verdicts", resolverPolicy: "board_or_agents", payload: { version: 1,
|
||||
prompt: "Review the work",
|
||||
items: [{ id: "work", label: "Work" }],
|
||||
},
|
||||
}, { userId: "queue-owner" });
|
||||
await db.update(heartbeatRuns).set({ resultJson: { executionCancellation: { state: "requested" } } })
|
||||
.where(eq(heartbeatRuns.id, seeded.runId));
|
||||
const actor = { agentId: seeded.agentId, runId: seeded.runId };
|
||||
const responding = kind === "questions"
|
||||
? service.answerQuestions(issue, interaction.id, { answers: [{ questionId: "scope", optionIds: ["first"] }] }, actor)
|
||||
: service.submitItemVerdicts(issue, interaction.id, { verdicts: [{ id: "work", verdict: "approve" }] }, actor);
|
||||
await expect(responding).rejects.toMatchObject({ status: 403, details: { code: "agent_run_cancelled" } });
|
||||
expect((await db.select().from(issueThreadInteractions).where(eq(issueThreadInteractions.id, interaction.id)))[0].status)
|
||||
.toBe("pending");
|
||||
});
|
||||
|
||||
it.each((["request_confirmation", "request_checkbox_confirmation", "ask_user_questions"] as const)
|
||||
.flatMap(kind => (["legacy", "native"] as const).map(runtime => ({ kind, runtime }))))(
|
||||
"queues $kind resolution while its $runtime source run is still running",
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
import { and, eq } from "drizzle-orm";
|
||||
import { heartbeatRuns, type Db } from "@paperclipai/db";
|
||||
import { forbidden } from "./errors.js";
|
||||
|
||||
/** Stop revokes write authority before waiting for the executor to settle. */
|
||||
export function agentRunWritesRevoked(run: {
|
||||
status: string;
|
||||
resultJson?: Record<string, unknown> | null;
|
||||
} | null | undefined): boolean {
|
||||
const cancellation = run?.resultJson?.executionCancellation;
|
||||
return run?.status === "cancelled" || Boolean(cancellation && typeof cancellation === "object"
|
||||
&& "state" in cancellation && cancellation.state === "requested");
|
||||
}
|
||||
|
||||
/** The caller must supply its write transaction, so Stop and commit serialize. */
|
||||
export async function assertAgentRunWriteAllowed(tx: Db, companyId: string, actor: {
|
||||
agentId?: string | null;
|
||||
runId?: string | null;
|
||||
stopId?: string | null;
|
||||
}) {
|
||||
if (!actor.agentId || !actor.runId) return;
|
||||
const [run] = await tx.select({ status: heartbeatRuns.status, resultJson: heartbeatRuns.resultJson })
|
||||
.from(heartbeatRuns).where(and(eq(heartbeatRuns.id, actor.runId),
|
||||
eq(heartbeatRuns.companyId, companyId), eq(heartbeatRuns.agentId, actor.agentId)))
|
||||
.for("share");
|
||||
const stoppedForThisMutation = run?.status === "cancelled" && actor.stopId &&
|
||||
run.resultJson?.issueMutationStopId === actor.stopId;
|
||||
if (agentRunWritesRevoked(run) && !stoppedForThisMutation) {
|
||||
throw forbidden("This run was cancelled", { code: "agent_run_cancelled" });
|
||||
}
|
||||
}
|
||||
@@ -21,6 +21,7 @@ import {
|
||||
rekeyCompanyIssueIdentifiers,
|
||||
} from "../services/issue-prefix.js";
|
||||
import { verifyLocalAgentJwt } from "../agent-auth-jwt.js";
|
||||
import { agentRunWritesRevoked } from "../agent-run-cancellation.js";
|
||||
import { isUuidLike, normalizeAgentApiKeyScope, type DeploymentMode } from "@paperclipai/shared";
|
||||
import type { BetterAuthSessionResult } from "../auth/better-auth.js";
|
||||
import { logger } from "./logger.js";
|
||||
@@ -392,13 +393,15 @@ export function actorMiddleware(db: Db, opts: ActorMiddlewareOptions): RequestHa
|
||||
}
|
||||
|
||||
const [identityRun] = await db.select({ activeIdentityContextId: heartbeatRuns.activeIdentityContextId,
|
||||
responsibleUserId: heartbeatRuns.responsibleUserId, status: heartbeatRuns.status,
|
||||
responsibleUserId: heartbeatRuns.responsibleUserId, status: heartbeatRuns.status, resultJson: heartbeatRuns.resultJson,
|
||||
contextSnapshot: heartbeatRuns.contextSnapshot }).from(heartbeatRuns).where(and(
|
||||
eq(heartbeatRuns.id, claims.run_id), eq(heartbeatRuns.companyId, claims.company_id), eq(heartbeatRuns.agentId, claims.sub),
|
||||
));
|
||||
if (identityRun?.status === "cancelled" && identityRun.contextSnapshot?.conversationMode === true
|
||||
if (agentRunWritesRevoked(identityRun)
|
||||
&& !["GET", "HEAD", "OPTIONS"].includes(req.method)) {
|
||||
_res.status(403).json({ error: "This conversation turn was cancelled", code: "conversation_turn_cancelled" });
|
||||
const conversation = identityRun?.contextSnapshot?.conversationMode === true;
|
||||
_res.status(403).json({ error: conversation ? "This conversation turn was cancelled" : "This run was cancelled",
|
||||
code: conversation ? "conversation_turn_cancelled" : "agent_run_cancelled" });
|
||||
return;
|
||||
}
|
||||
if (identityRun?.activeIdentityContextId && identityRun.status === "running") {
|
||||
|
||||
@@ -13458,6 +13458,9 @@ export function issueRoutes(
|
||||
}
|
||||
}
|
||||
|
||||
// Only this request may finish a mutation that intentionally stops its
|
||||
// own run (for example handing work to a signoff reviewer).
|
||||
const issueMutationStopId = randomUUID();
|
||||
if (assigneeWillChange && existing.assigneeAgentId) {
|
||||
await stopRunnerGoalForOwnershipChange({
|
||||
companyId: existing.companyId,
|
||||
@@ -13471,7 +13474,7 @@ export function issueRoutes(
|
||||
"Cancelled before issue reassignment",
|
||||
{
|
||||
errorCode: "issue_reassigned",
|
||||
resultJson: { reassignmentStopConfirmed: true },
|
||||
resultJson: { reassignmentStopConfirmed: true, issueMutationStopId },
|
||||
eventMessage: "run cancelled before issue reassignment",
|
||||
eventPayload: { issueId: existing.id },
|
||||
},
|
||||
@@ -13512,7 +13515,7 @@ export function issueRoutes(
|
||||
"Cancelled before issue terminalization",
|
||||
{
|
||||
errorCode: "issue_terminalized",
|
||||
resultJson: { terminalizationStopConfirmed: true },
|
||||
resultJson: { terminalizationStopConfirmed: true, issueMutationStopId },
|
||||
eventMessage: "run cancelled before issue terminalization",
|
||||
eventPayload: {
|
||||
issueId: existing.id,
|
||||
@@ -13564,6 +13567,8 @@ export function issueRoutes(
|
||||
const issueUpdateData = {
|
||||
...updateFields,
|
||||
actorAgentId: actor.agentId ?? null,
|
||||
actorRunId: actor.agentId ? actor.runId : null,
|
||||
actorRunStopId: actor.agentId && interruptedRunId === actor.runId ? issueMutationStopId : null,
|
||||
actorUserId: actor.actorType === "user" ? actor.actorId : null,
|
||||
};
|
||||
const shouldCollectCompletionPublication =
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import { currentContinuationOrigins } from "./execution-continuation.js";
|
||||
import { assertAgentRunWriteAllowed } from "../agent-run-cancellation.js";
|
||||
import { connectionIntentDeliveries } from "@paperclipai/db";
|
||||
import { isDeepStrictEqual } from "node:util";
|
||||
import {
|
||||
@@ -141,6 +142,14 @@ type InteractionActor = {
|
||||
resolutionDetails?: Record<string, unknown>;
|
||||
};
|
||||
|
||||
async function assertInteractionRunWriteAllowed(tx: Db, issue: { id: string; companyId: string }, actor: InteractionActor) {
|
||||
if (!actor.agentId || !actor.runId) return;
|
||||
// Keep the same issue -> run lock order as task mutation and checkout.
|
||||
await tx.select({ id: issues.id }).from(issues)
|
||||
.where(and(eq(issues.id, issue.id), eq(issues.companyId, issue.companyId))).for("update");
|
||||
await assertAgentRunWriteAllowed(tx, issue.companyId, actor);
|
||||
}
|
||||
|
||||
type CreateInteractionOptions = {
|
||||
/** Keep independently owned pending cards actionable. Internal runtime bridges use this. */
|
||||
supersedePendingSiblingInteractions?: boolean;
|
||||
@@ -2103,6 +2112,7 @@ export function issueThreadInteractionService(
|
||||
const now = new Date();
|
||||
const postCommitActivityPublications: ActivityPublication[] = [];
|
||||
const result = await db.transaction(async (tx) => {
|
||||
await assertInteractionRunWriteAllowed(tx as unknown as Db, args.issue, args.actor);
|
||||
await args.mutationOptions?.beforeResolveInTransaction?.(tx);
|
||||
// Policy mutations and review transitions use the same issue-row lock,
|
||||
// so the authoritative review policy and requester are stable through
|
||||
@@ -2381,6 +2391,7 @@ export function issueThreadInteractionService(
|
||||
|
||||
const now = new Date();
|
||||
const updated = await db.transaction(async (tx) => {
|
||||
await assertInteractionRunWriteAllowed(tx as unknown as Db, args.issue, args.actor);
|
||||
await args.mutationOptions?.beforeResolveInTransaction?.(tx);
|
||||
const issueContext = await tx
|
||||
.select({
|
||||
@@ -3468,6 +3479,7 @@ export function issueThreadInteractionService(
|
||||
// Idempotent reuse above stays allowed so retries of a pre-close
|
||||
// create keep returning the (by now expired) original.
|
||||
const result = await db.transaction(async (tx) => {
|
||||
await assertInteractionRunWriteAllowed(tx as unknown as Db, issue, actor);
|
||||
const [issueRow] = await tx
|
||||
.select({ status: issues.status })
|
||||
.from(issues)
|
||||
@@ -3808,6 +3820,7 @@ export function issueThreadInteractionService(
|
||||
const createdWakeTargets: IssueWakeTarget[] = [];
|
||||
|
||||
await db.transaction(async (tx) => {
|
||||
await assertInteractionRunWriteAllowed(tx as unknown as Db, issue, actor);
|
||||
const resolvedAt = new Date();
|
||||
const [claimed] = await tx
|
||||
.update(issueThreadInteractions)
|
||||
@@ -3965,6 +3978,7 @@ export function issueThreadInteractionService(
|
||||
assertIssueOpenForInteractionResolution(issue);
|
||||
const data = submitIssueThreadInteractionVerdictsSchema.parse(input);
|
||||
const submission = await db.transaction(async (tx) => {
|
||||
await assertInteractionRunWriteAllowed(tx as unknown as Db, issue, actor);
|
||||
const current = await tx
|
||||
.select()
|
||||
.from(issueThreadInteractions)
|
||||
@@ -4102,8 +4116,9 @@ export function issueThreadInteractionService(
|
||||
throw interactionTerminalError(current);
|
||||
}
|
||||
|
||||
const [updated] = await db
|
||||
.update(issueThreadInteractions)
|
||||
const [updated] = await db.transaction(async (tx) => {
|
||||
await assertInteractionRunWriteAllowed(tx as unknown as Db, issue, actor);
|
||||
return tx.update(issueThreadInteractions)
|
||||
.set({
|
||||
status: "rejected",
|
||||
result: {
|
||||
@@ -4123,6 +4138,7 @@ export function issueThreadInteractionService(
|
||||
),
|
||||
)
|
||||
.returning();
|
||||
});
|
||||
|
||||
if (!updated) {
|
||||
throw interactionAlreadyResolvedError();
|
||||
@@ -4664,6 +4680,7 @@ export function issueThreadInteractionService(
|
||||
// review queue while the card is still pending, and an executable
|
||||
// request must not outlive a withdrawn card.
|
||||
const updated = await db.transaction(async (tx) => {
|
||||
await assertInteractionRunWriteAllowed(tx as unknown as Db, issue, actor);
|
||||
await resolveLinkedToolActionRequests(tx, current, {
|
||||
status: "cancelled",
|
||||
fromStatuses: ["pending", "approved"],
|
||||
@@ -4771,6 +4788,7 @@ export function issueThreadInteractionService(
|
||||
});
|
||||
|
||||
const updated = await db.transaction(async (tx) => {
|
||||
await assertInteractionRunWriteAllowed(tx as unknown as Db, issue, actor);
|
||||
await mutationOptions.beforeResolveInTransaction?.(tx);
|
||||
const resolvedAt = new Date();
|
||||
const [row] = await tx
|
||||
@@ -4846,6 +4864,7 @@ export function issueThreadInteractionService(
|
||||
const reason = data.reason?.trim() || null;
|
||||
const now = new Date();
|
||||
const updated = await db.transaction(async (tx) => {
|
||||
await assertInteractionRunWriteAllowed(tx as unknown as Db, issue, actor);
|
||||
await resolveLinkedToolActionRequests(tx, current, {
|
||||
status: "cancelled",
|
||||
fromStatuses: ["pending", "approved"],
|
||||
@@ -4942,6 +4961,7 @@ export function issueThreadInteractionService(
|
||||
|
||||
const reason = data.reason?.trim() || null;
|
||||
const updated = await db.transaction(async (tx) => {
|
||||
await assertInteractionRunWriteAllowed(tx as unknown as Db, issue, actor);
|
||||
const resolvedAt = new Date();
|
||||
const [row] = await tx
|
||||
.update(issueThreadInteractions)
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import { mirrorSlackBoardComment, slackBoardReplyBindings } from "./slack-board-messages.js";
|
||||
import { assertAgentRunWriteAllowed } from "../agent-run-cancellation.js";
|
||||
import { externalConversationStateSql, nonIdleSlackIssueCondition, resumeSlackConversation } from "./slack-conversation-state.js";
|
||||
import { documentService } from "./documents.js";
|
||||
import { parseTaskSearch, taskSearchCtes, taskSearchScore } from "./task-search.js";
|
||||
@@ -10536,6 +10537,8 @@ export function issueService(db: Db) {
|
||||
labelIds?: string[];
|
||||
blockedByIssueIds?: string[];
|
||||
actorAgentId?: string | null;
|
||||
actorRunId?: string | null;
|
||||
actorRunStopId?: string | null;
|
||||
actorUserId?: string | null;
|
||||
companyGuard?: string;
|
||||
},
|
||||
@@ -10581,6 +10584,8 @@ export function issueService(db: Db) {
|
||||
labelIds: nextLabelIds,
|
||||
blockedByIssueIds,
|
||||
actorAgentId,
|
||||
actorRunId,
|
||||
actorRunStopId,
|
||||
actorUserId,
|
||||
companyGuard,
|
||||
...issueData
|
||||
@@ -10854,6 +10859,13 @@ export function issueService(db: Db) {
|
||||
.for("update")
|
||||
.then((rows: Array<typeof issues.$inferSelect>) => rows[0] ?? null);
|
||||
if (!receiptExisting) return null;
|
||||
if (actorAgentId && actorRunId) {
|
||||
// Recheck under a run lock: a request admitted before Stop must not
|
||||
// commit a late Done after cancellation revoked its credentials.
|
||||
await assertAgentRunWriteAllowed(tx, receiptExisting.companyId, {
|
||||
agentId: actorAgentId, runId: actorRunId, stopId: actorRunStopId,
|
||||
});
|
||||
}
|
||||
if (actorAgentId && patch.status === "done") {
|
||||
const [review] = await tx.select({ id: toolActionRequests.id }).from(toolActionRequests).where(and(eq(toolActionRequests.companyId, existing.companyId), eq(toolActionRequests.issueId, id), inArray(toolActionRequests.status, ["pending", "approved", "executing"]))).limit(1);
|
||||
if (review) throw conflict("This task is waiting for a connection review. Finish unrelated work, then yield in_review without retrying the governed call.", { code: "tool_review_pending", actionRequestId: review.id });
|
||||
|
||||
Reference in New Issue
Block a user