mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-10 20:50:08 +02:00
Fix confirmed Dot Runner shutdown and revoked binding cache
Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
5201456bfc
commit
fccf35f90f
5 files changed
+89
-2
No files matched your search
@@ -147,3 +147,11 @@ After those fixes, full workspace typecheck and token gates passed again.
|
||||
Focused verification passed 89 gateway, Dot, OpenAPI, connection-instruction
|
||||
and pairing UI tests, 40 protocol/runtime compatibility tests, and the real
|
||||
Runner protocol-upgrade/replacement test.
|
||||
|
||||
PR CI exposed a clean-shutdown race: Rust could exit after the durable shutdown
|
||||
receipt was acknowledged but before the Dot SDK's next poll. The adapter now
|
||||
recognizes that confirmed clean exit and still requires reconciliation after
|
||||
an unexpected exit. A real Rust regression reproduces the failing order and
|
||||
passes with the fix. Four driver tests and 11 integration/pairing tests passed,
|
||||
including a failed connection refresh after successful revocation. Revocation
|
||||
clears the cached binding before refetch so that failure cannot restore it.
|
||||
@@ -4,6 +4,7 @@ import { tmpdir } from "node:os";
|
||||
import { join, resolve } from "node:path";
|
||||
import { expect, it, vi } from "vitest";
|
||||
import { RunnerdDotDriver, type RunnerdDotDriverOptions } from "./runnerd-dot-driver.js";
|
||||
import * as controlPlane from "../../control-plane/durable-prp-control-plane.js";
|
||||
import { externalOperationDigest, type ExternalProviderOperation } from "../../contracts/external-provider.js";
|
||||
import { buildNativeModelEnvelope, parseNativeExecutionInput, type NativeExecutionInputV6 } from "../../contracts/native-execution.js";
|
||||
import { NATIVE_RUNTIME_ASSET_SCHEMA, PAPERCLIP_EXECUTION_PROMPT, PAPERCLIP_EXECUTION_PROMPT_REVISION,
|
||||
@@ -47,6 +48,54 @@ it("v6 closes Dot identity, workspace, model and credential fields", () => {
|
||||
]) expect(() => parseNativeExecutionInput(changed)).toThrow();
|
||||
});
|
||||
|
||||
it.each(["shutdown", "unexpected exit"] as const)("classifies a real Rust %s before the next command poll", async kind => {
|
||||
const root = await mkdtemp(join(tmpdir(), "dot-driver-exit-"));
|
||||
await writeFile(join(root, "AGENTS.md"), "Use only the synthetic counter.");
|
||||
const input = execution(root);
|
||||
let handle: controlPlane.RunnerProcessHandle | undefined;
|
||||
let exited = false;
|
||||
const spawn = controlPlane.spawnRunner;
|
||||
const spawnSpy = vi.spyOn(controlPlane, "spawnRunner").mockImplementation(options => {
|
||||
handle = spawn(options);
|
||||
void handle.completion.then(() => { exited = true; });
|
||||
return handle;
|
||||
});
|
||||
const getCommand = controlPlane.DurablePrpControlPlane.prototype.getCommand;
|
||||
const commandSpy = vi.spyOn(controlPlane.DurablePrpControlPlane.prototype, "getCommand").mockImplementation(function (id) {
|
||||
const command = getCommand.call(this, id);
|
||||
// Hold the SDK poll until the real process has exited. The authenticated
|
||||
// shutdown receipt and its ACK remain committed in the transport journal.
|
||||
if (id === "dot_shutdown" && command?.status === "completed" && !exited) return { ...command, status: "pending" };
|
||||
return command;
|
||||
});
|
||||
const driver = new RunnerdDotDriver({
|
||||
execution: input, stateDirectory: join(root, "state"),
|
||||
identity: { runnerInstanceId: randomUUID(), environmentLeaseId: randomUUID(), runId: input.binding.runId,
|
||||
normalizedSessionId: input.session.normalizedSessionId!, turnId: "turn-" + input.binding.runId, itemId: "item-" + input.binding.runId },
|
||||
runnerBinary: resolve("runner/target/debug/paperclip-runnerd"),
|
||||
port: { dispatch: async () => {}, settle: async () => {}, attach: async () => async () => {} },
|
||||
});
|
||||
let session: Awaited<ReturnType<typeof driver.openSession>> | undefined;
|
||||
try {
|
||||
session = await driver.openSession({ runId: input.binding.runId, normalizedSessionId: input.session.normalizedSessionId! });
|
||||
if (kind === "shutdown") {
|
||||
await expect(session.close({ reason: "Test complete" })).resolves.toBeUndefined();
|
||||
expect(exited).toBe(true);
|
||||
expect((await handle!.completion).code).toBe(0);
|
||||
} else {
|
||||
handle!.child.kill("SIGTERM");
|
||||
await handle!.completion;
|
||||
await expect(session.read!()).rejects.toThrow("dot_runner_process_exited_recovery_required");
|
||||
}
|
||||
} finally {
|
||||
commandSpy.mockRestore(); spawnSpy.mockRestore();
|
||||
await session?.detachControllerForRestart?.();
|
||||
if (!exited) handle?.child.kill("SIGTERM");
|
||||
await handle?.completion;
|
||||
await rm(root, { recursive: true, force: true });
|
||||
}
|
||||
}, 45000);
|
||||
|
||||
it("reattaches the same Rust bridge without duplicating a settled operation and refuses a lost checkpoint", async () => {
|
||||
const root = await mkdtemp(join(tmpdir(), "dot-driver-recovery-"));
|
||||
await writeFile(join(root, "AGENTS.md"), "Use only the synthetic counter.");
|
||||
|
||||
@@ -132,7 +132,18 @@ class RunnerdDotSession implements HarnessSession {
|
||||
const pid = this.#process.child.pid;
|
||||
if (pid) await o.onSpawn?.({ pid, processGroupId: this.#process.processGroupId ?? null,
|
||||
startedAt: this.#process.startedAt ?? new Date().toISOString() });
|
||||
void this.#process.completion.then(() => { if (!this.#closed) { this.#failure = new Error("dot_runner_process_exited_recovery_required"); this.#wake(); } });
|
||||
void this.#process.completion.then(result => {
|
||||
if (this.#closed) return;
|
||||
// Rust exits after the authenticated shutdown receipt is committed and
|
||||
// ACKed. That exit can precede the SDK's next command poll.
|
||||
if (result.code === 0 && core.getCommand("dot_shutdown")?.status === "completed") return;
|
||||
this.#failure ??= new Error(`dot_runner_process_exited_recovery_required: code=${result.code} signal=${result.signal}`);
|
||||
this.#wake();
|
||||
}, () => {
|
||||
if (this.#closed) return;
|
||||
this.#failure ??= new Error("dot_runner_process_exited_recovery_required");
|
||||
this.#wake();
|
||||
});
|
||||
} else if (!await o.adoptExistingRunner.isAlive()) { throw new Error("dot_runner_adoption_failed"); }
|
||||
await registration?.activate?.();
|
||||
await registration?.ready?.();
|
||||
|
||||
@@ -61,3 +61,17 @@ it("does not restore a stale binding while revocation refreshes the connection",
|
||||
expect(onBinding).toHaveBeenLastCalledWith("");
|
||||
expect(Array.from(container.querySelectorAll("button")).some(button => button.textContent === "Pair Dot")).toBe(true);
|
||||
});
|
||||
|
||||
it("keeps a revoked binding cleared when the follow-up connection read fails", async () => {
|
||||
api.get.mockResolvedValue({ enabled: true, resourceUrl: "https://paperclip.example/mcp/runner", binding });
|
||||
api.delete.mockImplementation(async () => { api.get.mockRejectedValue(new Error("Connection refresh failed")); });
|
||||
await render();
|
||||
expect(container.querySelector("output")?.textContent).toBe(binding.id);
|
||||
await act(async () => Array.from(container.querySelectorAll("button")).find(button => button.textContent === "Revoke connection")!.click());
|
||||
await flush();
|
||||
expect(container.querySelector("output")?.textContent).toBe("");
|
||||
expect(onBinding).toHaveBeenLastCalledWith("");
|
||||
expect(container.querySelector("[role=alert]")?.textContent).toBe("Connection refresh failed");
|
||||
expect(client.getQueryData(["dot-binding", "company", "agent"])).toMatchObject({ binding: null });
|
||||
expect(Array.from(container.querySelectorAll("button")).some(button => button.textContent === "Revoke connection")).toBe(false);
|
||||
});
|
||||
@@ -26,7 +26,12 @@ export function DotRunnerConnection({ companyId, agentId, bindingId, onBinding }
|
||||
const test = useMutation({ mutationFn: () => api.post(path + "/event-test", {}),
|
||||
onSuccess: () => { void client.invalidateQueries({ queryKey: key }); } });
|
||||
const revoke = useMutation({ mutationFn: () => api.delete(path),
|
||||
onSuccess: async () => { await client.invalidateQueries({ queryKey: key }); setPairing(null); onBinding(""); } });
|
||||
onSuccess: async () => {
|
||||
await client.cancelQueries({ queryKey: key });
|
||||
client.setQueryData<Connection>(key, previous => previous ? { ...previous, binding: null } : previous);
|
||||
setPairing(null); onBinding("");
|
||||
await client.invalidateQueries({ queryKey: key });
|
||||
} });
|
||||
const existingBindingId = state.data?.binding?.id;
|
||||
useEffect(() => {
|
||||
if (existingBindingId && existingBindingId !== bindingId && !revoke.isPending) onBinding(existingBindingId);
|
||||
|
||||
Reference in new issue
Block a user