mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-07 16:11:46 +02:00
fix(dot): preserve completion and admission receipts
Normalize accepted completion reports across the control plane and Rust bridge, serialize Dot assignments, replay work admission receipts, and continue subscription polling. Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
ad1ff9b329
commit
78f4c4f5e6
8 files changed
+227
-30
No files matched your search
@@ -64,6 +64,16 @@ structured report. Paperclip's ordinary result and status finalizers decide
|
||||
the task disposition. Dot can also discover its assigned tasks and request
|
||||
normal admission with `paperclip_dot_tasks` and `paperclip_dot_request_work`.
|
||||
|
||||
The completion tool returns its canonical accepted report. The final turn can
|
||||
repeat either that canonical report or the exact original accepted arguments;
|
||||
the Runner persists and emits the canonical report. Changed reports are rejected.
|
||||
Work requests use a durable admission receipt keyed by binding, generation and
|
||||
request ID. Retries read that receipt, even after a task finishes, and cannot
|
||||
enqueue another run or move the same request to another task. A Dot runs one
|
||||
assignment at a time; competing tasks stay queued regardless of the agent's
|
||||
configured concurrency. During setup the connection panel continues checking
|
||||
for a verified event subscription without requiring a manual refresh.
|
||||
|
||||
## Limits and recovery
|
||||
|
||||
- One active assignment per binding. Acceptance expires after 10 minutes;
|
||||
|
||||
@@ -81,6 +81,8 @@ struct State {
|
||||
text: Option<String>,
|
||||
operations: BTreeMap<String, Receipt>,
|
||||
completion: Option<Value>,
|
||||
#[serde(default)]
|
||||
completion_input: Option<Value>,
|
||||
events: VecDeque<PolledEvent>,
|
||||
next_event: u64,
|
||||
}
|
||||
@@ -335,6 +337,7 @@ impl DotCommandExecutor {
|
||||
text: None,
|
||||
operations: BTreeMap::new(),
|
||||
completion: None,
|
||||
completion_input: None,
|
||||
events: VecDeque::new(),
|
||||
next_event: 1,
|
||||
});
|
||||
@@ -421,10 +424,16 @@ impl DotCommandExecutor {
|
||||
if state.lifecycle != "running" || state.tools.pending_calls().next().is_some() {
|
||||
return Err(invalid("Dot cannot finish with unsettled tool calls"));
|
||||
}
|
||||
let result = op.input["result"].clone();
|
||||
if state.completion.as_ref() != Some(&result) {
|
||||
let submitted = &op.input["result"];
|
||||
if state.completion.as_ref() != Some(submitted)
|
||||
&& state.completion_input.as_ref() != Some(submitted)
|
||||
{
|
||||
return Err(invalid("Dot must successfully invoke paperclip_finish or paperclip_block with this result first"));
|
||||
}
|
||||
let result = state
|
||||
.completion
|
||||
.clone()
|
||||
.ok_or_else(|| invalid("Dot completion is missing"))?;
|
||||
state.lifecycle = "completed".to_owned();
|
||||
state
|
||||
.tools
|
||||
@@ -472,7 +481,15 @@ impl DotCommandExecutor {
|
||||
.and_then(Value::as_str)
|
||||
.is_none_or(|o| matches!(o, "completed" | "succeeded" | "success"))
|
||||
{
|
||||
let mut completion = receipt.operation.input["arguments"].clone();
|
||||
// The control plane normalizes provider aliases and optional fields.
|
||||
// Persist that exact accepted report; do not revalidate raw arguments
|
||||
// against a different contract. Legacy canonical receipts still work.
|
||||
let accepted_input = receipt.operation.input["arguments"].clone();
|
||||
let mut completion = result
|
||||
.result
|
||||
.get("completionReport")
|
||||
.cloned()
|
||||
.unwrap_or_else(|| accepted_input.clone());
|
||||
completion["schema"] = json!("paperclip.run_result.v1");
|
||||
let schema: Value = serde_json::from_str(include_str!(
|
||||
"../../../../protocol/schemas/result.schema.json"
|
||||
@@ -504,6 +521,7 @@ impl DotCommandExecutor {
|
||||
));
|
||||
}
|
||||
state.completion = Some(completion);
|
||||
state.completion_input = Some(accepted_input);
|
||||
}
|
||||
let outcome =
|
||||
json!({"status":"completed","isError":result.is_error,"result":result.result});
|
||||
@@ -938,6 +956,59 @@ mod tests {
|
||||
);
|
||||
}
|
||||
#[test]
|
||||
fn normalized_completion_survives_restart_and_accepts_the_original_report() {
|
||||
let dir = TestDirectory::new();
|
||||
let mut e = prepared(dir.path());
|
||||
accept(&mut e);
|
||||
let input = json!({"reportedWorkDisposition":"completed","summary":"Saved report",
|
||||
"completionClaim":{"contractRevision":"contract-1","objectiveSatisfied":true,"criteria":[{"criterionId":"criterion-1","status":"passed","evidenceRefs":[]}],"remainingWork":[]},
|
||||
"evidence":[],"verification":[]});
|
||||
let mut canonical = input.clone();
|
||||
canonical["schema"] = json!("paperclip.run_result.v1");
|
||||
canonical["reportedWorkDisposition"] = json!("done");
|
||||
canonical["completionClaim"]["criteria"][0]["status"] = json!("satisfied");
|
||||
canonical["attentionRequests"] = json!([]);
|
||||
canonical["artifacts"] = json!([]);
|
||||
call(
|
||||
&mut e,
|
||||
"external_provider.operation",
|
||||
operation(
|
||||
"completion-1",
|
||||
"tool",
|
||||
json!({"name":"paperclip_finish","arguments":input}),
|
||||
),
|
||||
);
|
||||
call(
|
||||
&mut e,
|
||||
"semantic_tool.result",
|
||||
json!({"callId":"completion-1","operationId":"paperclip_finish","result":{"accepted":true,"completionReport":canonical},"isError":false}),
|
||||
);
|
||||
drop(e);
|
||||
let mut e = DotCommandExecutor::with_runner_config(dir.path(), &config(dir.path()));
|
||||
let mut changed = input.clone();
|
||||
changed["summary"] = json!("Changed after acceptance");
|
||||
assert!(e
|
||||
.execute(&command(
|
||||
"external_provider.operation",
|
||||
operation("finish-1", "finish", json!({"result":changed}))
|
||||
))
|
||||
.is_err());
|
||||
call(
|
||||
&mut e,
|
||||
"external_provider.operation",
|
||||
operation("finish-2", "finish", json!({"result":input})),
|
||||
);
|
||||
let events = e.poll_events().unwrap();
|
||||
assert_eq!(
|
||||
events
|
||||
.iter()
|
||||
.find(|event| event.event_type == "run.result.proposed")
|
||||
.unwrap()
|
||||
.payload,
|
||||
canonical
|
||||
);
|
||||
}
|
||||
#[test]
|
||||
fn wrong_epoch_and_expired_offers_fail_closed() {
|
||||
let dir = TestDirectory::new();
|
||||
let mut e = prepared(dir.path());
|
||||
|
||||
@@ -102,7 +102,7 @@ class RunnerdDotSession implements HarnessSession {
|
||||
if (!result.ok || (call.operationId === "paperclip_block") !== (result.ok && result.result.reportedWorkDisposition === "blocked")) {
|
||||
return { result: { error: "Invalid completion report. Use the admitted completion contract and the correct completion tool." }, isError: true };
|
||||
}
|
||||
try { return { result: { accepted: true, feedback: await o.completionFeedback?.(result.result) ?? "Completion accepted; finish the external turn." } }; }
|
||||
try { return { result: { accepted: true, completionReport: result.result, feedback: await o.completionFeedback?.(result.result) ?? "Completion accepted; finish the external turn." } }; }
|
||||
catch (error) { return { result: { error: error instanceof Error ? error.message : "Completion rejected" }, isError: true }; }
|
||||
}
|
||||
if (!o.dynamicToolHandler) throw new Error("dot_semantic_authority_unavailable");
|
||||
|
||||
@@ -5,12 +5,12 @@ import { fileURLToPath } from "node:url";
|
||||
import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { eq, sql } from "drizzle-orm";
|
||||
import { and, eq, sql } from "drizzle-orm";
|
||||
import express, { type Request } from "express";
|
||||
import request from "supertest";
|
||||
import { afterAll, beforeAll, describe, expect, it, vi } from "vitest";
|
||||
import { agents, authUsers, companies, companyMemberships, createDb, dotMailboxItems, dotRunnerAssignments, dotRunnerOperations,
|
||||
heartbeatRuns, issues, nativeRunResults, nativeRunFinalizations, completionContracts, mcpEventDeliveries, workspaceOperations } from "@paperclipai/db";
|
||||
heartbeatRuns, agentWakeupRequests, issues, nativeRunResults, nativeRunFinalizations, completionContracts, mcpEventDeliveries, workspaceOperations } from "@paperclipai/db";
|
||||
import { startEmbeddedPostgresTestDatabase } from "./helpers/embedded-postgres.js";
|
||||
import { createPublicMcpOAuth, DEVICE_GRANT } from "../services/public-mcp/oauth.js";
|
||||
import { createPublicMcpExecutor } from "../services/public-mcp/capabilities.js";
|
||||
@@ -26,6 +26,7 @@ import { executePaperclipNativeSession } from "../services/native-runtime/native
|
||||
import { documentService } from "../services/documents.js";
|
||||
import { setupRunnerPrpWebSocketServer } from "../realtime/runner-prp-ws.js";
|
||||
import { finalizeNativeRun } from "../services/native-runtime/native-run-finalizer.js";
|
||||
import { heartbeatService } from "../services/heartbeat.js";
|
||||
|
||||
describe("durable Dot Runner integration", () => {
|
||||
let temporary: Awaited<ReturnType<typeof startEmbeddedPostgresTestDatabase>>;
|
||||
@@ -38,6 +39,7 @@ describe("durable Dot Runner integration", () => {
|
||||
db = createDb(temporary.connectionString);
|
||||
root = await mkdtemp(join(tmpdir(), "paperclip-dot-state-"));
|
||||
vi.stubEnv("PAPERCLIP_ENABLE_OPENAI_DOT", "1");
|
||||
vi.stubEnv("PAPERCLIP_IN_WORKTREE", "false");
|
||||
vi.stubEnv("PAPERCLIP_RUNNER_BINARY", join(runnerRoot, "runner/target/release", process.platform === "win32" ? "paperclip-runnerd.exe" : "paperclip-runnerd"));
|
||||
vi.stubEnv("PAPERCLIP_RUNNER_STATE_DIR", join(root, "runner-state"));
|
||||
vi.stubEnv("PAPERCLIP_SECRETS_MASTER_KEY", randomBytes(32).toString("base64"));
|
||||
@@ -128,6 +130,45 @@ describe("durable Dot Runner integration", () => {
|
||||
}
|
||||
});
|
||||
|
||||
it("keeps competing Dot assignments queued and replays durable work requests", async () => {
|
||||
const f = await fixture();
|
||||
await db.update(companies).set({ defaultResponsibleUserId: f.userId }).where(eq(companies.id, f.company.id));
|
||||
await db.update(agents).set({ runtimeConfig: { heartbeat: { maxConcurrentRuns: 20, wakeOnDemand: true } } }).where(eq(agents.id, f.agent.id));
|
||||
const holding = await offeredWork(f);
|
||||
await db.update(heartbeatRuns).set({ status: "running", startedAt: new Date() }).where(eq(heartbeatRuns.id, holding.run.id));
|
||||
const task = async (title: string) => (await db.insert(issues).values({ companyId: f.company.id, title, status: "todo", assigneeAgentId: f.agent.id, responsibleUserId: f.userId }).returning())[0]!;
|
||||
const first = await task("First queued task");
|
||||
const second = await task("Second queued task");
|
||||
try {
|
||||
const id = randomUUID();
|
||||
const responses = await Promise.all(Array.from({ length: 3 }, () => f.broker.requestWork(f.principal, first.id, id)));
|
||||
expect(responses[0]?.runId).toBeTruthy();
|
||||
expect(responses).toEqual([responses[0], responses[0], responses[0]]);
|
||||
const wakes = await db.select().from(agentWakeupRequests).where(and(eq(agentWakeupRequests.agentId, f.agent.id), eq(agentWakeupRequests.idempotencyKey, `dot-work:${f.snapshot.bindingId}:${f.snapshot.bindingGeneration}:${id}`)));
|
||||
expect(wakes).toHaveLength(1);
|
||||
expect((await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, responses[0]!.runId!)))[0]?.status).toBe("queued");
|
||||
await heartbeatService(db).resumeQueuedRuns();
|
||||
expect((await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, responses[0]!.runId!)))[0]?.status).toBe("queued");
|
||||
await db.update(heartbeatRuns).set({ status: "succeeded", finishedAt: new Date() }).where(eq(heartbeatRuns.id, responses[0]!.runId!));
|
||||
await db.update(issues).set({ status: "done" }).where(eq(issues.id, first.id));
|
||||
expect(await f.broker.requestWork(f.principal, first.id, id)).toEqual(responses[0]);
|
||||
await expect(f.broker.requestWork(f.principal, second.id, id)).rejects.toThrow("another task");
|
||||
const third = await task("Third queued task");
|
||||
const raceId = randomUUID();
|
||||
const race = await Promise.allSettled([f.broker.requestWork(f.principal, second.id, raceId), f.broker.requestWork(f.principal, third.id, raceId)]);
|
||||
expect(race.filter(r => r.status === "fulfilled")).toHaveLength(1);
|
||||
expect(race.filter(r => r.status === "rejected")).toHaveLength(1);
|
||||
const raceWakes = await db.select().from(agentWakeupRequests).where(and(eq(agentWakeupRequests.agentId, f.agent.id), eq(agentWakeupRequests.idempotencyKey, `dot-work:${f.snapshot.bindingId}:${f.snapshot.bindingGeneration}:${raceId}`)));
|
||||
expect(raceWakes).toHaveLength(1);
|
||||
const allRuns = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.agentId, f.agent.id));
|
||||
expect(allRuns).toHaveLength(3); // Holding turn plus one admitted run per unique request.
|
||||
expect(allRuns.filter(r => r.status === "queued")).toHaveLength(1);
|
||||
} finally {
|
||||
await db.update(heartbeatRuns).set({ status: "cancelled", finishedAt: new Date() }).where(eq(heartbeatRuns.agentId, f.agent.id));
|
||||
await f.events.unsubscribe(f.principal, f.subscription); await f.events.stop();
|
||||
}
|
||||
}, 30000);
|
||||
|
||||
it("verifies reconnects and wakes the same outstanding assignment after exhausted delivery", async () => {
|
||||
const f = await fixture();
|
||||
const { assignment } = await offeredWork(f);
|
||||
@@ -308,13 +349,17 @@ describe("durable Dot Runner integration", () => {
|
||||
await expect(f.broker.operation(f.principal, assignment!.id, writeId, "tool", { ...args, arguments: { ...args.arguments, body: "changed" } })).rejects.toThrow("different arguments");
|
||||
const document = await documentService(db).getIssueDocumentByKey(issue!.id, "report");
|
||||
expect(document?.body).toBe("17 + 25 = 42.");
|
||||
const result = { schema: "paperclip.run_result.v1", reportedWorkDisposition: "done", summary: "Saved the arithmetic report.",
|
||||
const result = { reportedWorkDisposition: "completed", summary: "Saved the arithmetic report.",
|
||||
completionClaim: { contractRevision: execution.completionContract.contract.revision, objectiveSatisfied: true,
|
||||
criteria: execution.completionContract.contract.criteria.map(c => ({ criterionId: c.id, status: "satisfied", evidenceRefs: [] })), remainingWork: [] },
|
||||
evidence: [], verification: [{ commandOrCheck: "17 + 25", status: "passed" }], attentionRequests: [], artifacts: [] };
|
||||
criteria: execution.completionContract.contract.criteria.map(c => ({ criterionId: c.id, status: "passed", evidenceRefs: [] })), remainingWork: [] },
|
||||
evidence: [], verification: [{ commandOrCheck: "17 + 25", status: "pass" }] };
|
||||
const completionId = randomUUID();
|
||||
await f.broker.operation(f.principal, assignment!.id, completionId, "tool", { name: "paperclip_finish", arguments: result });
|
||||
await vi.waitFor(async () => expect(await f.broker.operationStatus(f.principal, assignment!.id, completionId)).toMatchObject({ status: "completed", isError: false }), { timeout: 10000 });
|
||||
expect(await f.broker.operationStatus(f.principal, assignment!.id, completionId)).toMatchObject({ result: { completionReport: {
|
||||
schema: "paperclip.run_result.v1", reportedWorkDisposition: "done", artifacts: [], attentionRequests: [],
|
||||
verification: [{ commandOrCheck: "17 + 25", status: "passed" }],
|
||||
} } });
|
||||
await f.broker.operation(f.principal, assignment!.id, randomUUID(), "finish", { result });
|
||||
await resultPromise;
|
||||
// The heartbeat records its controller workspace settlement before the
|
||||
|
||||
@@ -197,17 +197,42 @@ function createBroker(db: Db) {
|
||||
const b = await principalBinding(principal);
|
||||
if (!z.uuid().safeParse(requestId).success) throw fail("Use a stable UUID requestId.");
|
||||
if (!enabled() || !(await this.bindingForAgent(b.companyId, b.agentId))?.subscriptionVerified) throw fail("Dot admission is unavailable.");
|
||||
const key = `dot-work:${b.id}:${b.generation}:${requestId}`;
|
||||
const receipt = async () => (await db.select().from(agentWakeupRequests).where(and(eq(agentWakeupRequests.companyId, b.companyId), eq(agentWakeupRequests.agentId, b.agentId), eq(agentWakeupRequests.idempotencyKey, key))).limit(1))[0];
|
||||
const replay = (row: typeof agentWakeupRequests.$inferSelect) => {
|
||||
if (row.payload?.issueId !== issueId) throw fail("requestId was reused for another task.");
|
||||
return { status: "requested", runId: row.runId, message: "Normal admission determines when this task can run. Read the mailbox after its event." };
|
||||
};
|
||||
// A retry is a receipt read, including after the task yielded or finished.
|
||||
const old = await receipt();
|
||||
if (old) return replay(old);
|
||||
const [issue] = await db.select().from(issues).where(and(eq(issues.companyId, b.companyId), eq(issues.id, issueId), eq(issues.assigneeAgentId, b.agentId)));
|
||||
if (!issue || !["todo", "in_progress"].includes(issue.status)) throw fail("Request work only for an eligible task assigned to this agent.");
|
||||
const key = `dot-work:${b.id}:${b.generation}:${requestId}`;
|
||||
const [old] = await db.select().from(agentWakeupRequests).where(and(eq(agentWakeupRequests.companyId, b.companyId), eq(agentWakeupRequests.agentId, b.agentId), eq(agentWakeupRequests.idempotencyKey, key)));
|
||||
if (old && old.payload?.issueId !== issueId) throw fail("requestId was reused for another task.");
|
||||
// The receipt PK reserves this request in the same transaction as its run.
|
||||
// Different issue locks cannot admit the same request concurrently.
|
||||
const hex = hash(key);
|
||||
const receiptId = `${hex.slice(0, 8)}-${hex.slice(8, 12)}-8${hex.slice(13, 16)}-${((parseInt(hex[16]!, 16) & 3) | 8).toString(16)}${hex.slice(17, 20)}-${hex.slice(20, 32)}`;
|
||||
const { heartbeatService } = await import("./heartbeat.js");
|
||||
const run = await heartbeatService(db).wakeup(b.agentId, { source: "assignment", triggerDetail: "system", reason: "issue_assigned",
|
||||
payload: { issueId, dotRequestId: requestId }, contextSnapshot: { issueId }, idempotencyKey: key, requestedByActorType: "agent", requestedByActorId: b.agentId });
|
||||
const [reserved] = await db.select().from(agentWakeupRequests).where(and(eq(agentWakeupRequests.companyId, b.companyId), eq(agentWakeupRequests.agentId, b.agentId), eq(agentWakeupRequests.idempotencyKey, key)));
|
||||
if (reserved && reserved.payload?.issueId !== issueId) throw fail("requestId was reused for another task.");
|
||||
return { status: "requested", runId: run?.id ?? null, message: "Normal admission determines when this task can run. Read the mailbox after its event." };
|
||||
try {
|
||||
const run = await heartbeatService(db).wakeup(b.agentId, { source: "assignment", triggerDetail: "system", reason: "issue_assigned",
|
||||
payload: { issueId, dotRequestId: requestId }, contextSnapshot: { issueId }, idempotencyKey: key, requestedByActorType: "agent", requestedByActorId: b.agentId,
|
||||
allowRunCoalescing: false, durableDotRequest: { id: receiptId, companyId: b.companyId, agentId: b.agentId, issueId, requestId, idempotencyKey: key, requestedAt: new Date() } });
|
||||
const reserved = await receipt();
|
||||
return reserved ? replay(reserved) : { status: "requested", runId: run?.id ?? null, message: "Normal admission determines when this task can run. Read the mailbox after its event." };
|
||||
} catch (error) {
|
||||
// A concurrent reservation can lose the PK race on a different task.
|
||||
// Replay only that durable receipt; do not mask unrelated failures.
|
||||
let cause: unknown = error;
|
||||
for (let depth = 0; depth < 5 && cause && typeof cause === "object"; depth++) {
|
||||
if ((cause as { code?: string }).code === "23505") {
|
||||
const reserved = await receipt();
|
||||
if (reserved?.id === receiptId) return replay(reserved);
|
||||
break;
|
||||
}
|
||||
cause = (cause as { cause?: unknown }).cause;
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
},
|
||||
async operationStatus(principal: McpPrincipal, assignmentId: string, requestId: string) {
|
||||
await authorizeAssignment(principal, assignmentId, true);
|
||||
|
||||
@@ -3779,6 +3779,16 @@ interface WakeupOptions {
|
||||
/** Exact failed run selected by an authenticated board Retry request. */
|
||||
failedRunId?: string | null;
|
||||
durableChatRequest?: DurableChatWakeupRequest;
|
||||
/** Server-owned Dot admission receipt; never accepted from API payloads. */
|
||||
durableDotRequest?: {
|
||||
id: string;
|
||||
companyId: string;
|
||||
agentId: string;
|
||||
issueId: string;
|
||||
requestId: string;
|
||||
idempotencyKey: string;
|
||||
requestedAt: Date;
|
||||
};
|
||||
source?: "timer" | "assignment" | "on_demand" | "automation";
|
||||
triggerDetail?: "manual" | "ping" | "callback" | "system";
|
||||
reason?: string | null;
|
||||
@@ -16861,9 +16871,11 @@ export function heartbeatService(
|
||||
enabled: asBoolean(heartbeat.enabled, false),
|
||||
intervalSec: Math.max(0, asNumber(heartbeat.intervalSec, 0)),
|
||||
wakeOnDemand: isHeartbeatWakeOnDemandEnabled(agent),
|
||||
maxConcurrentRuns: normalizeMaxConcurrentRuns(
|
||||
heartbeat.maxConcurrentRuns,
|
||||
),
|
||||
// A Dot binding has one external turn. Competing assignments must retain
|
||||
// their queue position instead of claiming a second run that cannot bind.
|
||||
maxConcurrentRuns: agent.adapterType === "paperclip_runner" &&
|
||||
parseObject(agent.adapterConfig).provider === "openai_dot"
|
||||
? 1 : normalizeMaxConcurrentRuns(heartbeat.maxConcurrentRuns),
|
||||
skipTimerWhenNoActionableWork: asBoolean(
|
||||
heartbeat.skipTimerWhenNoActionableWork ??
|
||||
heartbeat.requireActionableTimerWork ??
|
||||
@@ -27176,18 +27188,35 @@ export function heartbeatService(
|
||||
if (durableRequest?.failedRunRetry) {
|
||||
opts = { ...opts, allowRunCoalescing: false };
|
||||
}
|
||||
const durableReceiptFields = durableRequest
|
||||
? { id: durableRequest.id, requestedAt: durableRequest.requestedAt }
|
||||
const dotRequest = opts.durableDotRequest;
|
||||
if (dotRequest && (durableRequest || dotRequest.agentId !== agentId ||
|
||||
dotRequest.companyId !== agent.companyId || dotRequest.issueId !== issueId ||
|
||||
dotRequest.requestId !== payload?.dotRequestId || source !== "assignment" ||
|
||||
opts.requestedByActorType !== "agent" || opts.requestedByActorId !== agentId ||
|
||||
agent.adapterType !== "paperclip_runner" || parseObject(agent.adapterConfig).provider !== "openai_dot" ||
|
||||
opts.idempotencyKey !== dotRequest.idempotencyKey)) {
|
||||
throw conflict("Dot work request does not match its admission authority.");
|
||||
}
|
||||
const receiptRequest = durableRequest ?? dotRequest;
|
||||
const durableReceiptFields = receiptRequest
|
||||
? { id: receiptRequest.id, requestedAt: receiptRequest.requestedAt }
|
||||
: {};
|
||||
const existingDurableReceipt = async (queryDb: Db) => {
|
||||
if (!durableRequest) return null;
|
||||
if (!receiptRequest) return null;
|
||||
const receipt = await queryDb
|
||||
.select()
|
||||
.from(agentWakeupRequests)
|
||||
.where(eq(agentWakeupRequests.id, durableRequest.id))
|
||||
.where(eq(agentWakeupRequests.id, receiptRequest.id))
|
||||
.limit(1)
|
||||
.then((rows) => rows[0] ?? null);
|
||||
if (receipt) assertDurableChatWakeupReceipt(durableRequest, receipt);
|
||||
if (receipt && durableRequest) assertDurableChatWakeupReceipt(durableRequest, receipt);
|
||||
if (receipt && dotRequest && (receipt.companyId !== dotRequest.companyId ||
|
||||
receipt.agentId !== dotRequest.agentId || receipt.source !== "assignment" ||
|
||||
receipt.requestedByActorType !== "agent" || receipt.requestedByActorId !== dotRequest.agentId ||
|
||||
receipt.idempotencyKey !== dotRequest.idempotencyKey || receipt.payload?.issueId !== dotRequest.issueId ||
|
||||
receipt.payload?.dotRequestId !== dotRequest.requestId)) {
|
||||
throw conflict("requestId was reused for another task.");
|
||||
}
|
||||
return receipt;
|
||||
};
|
||||
const priorReceipt = await existingDurableReceipt(db);
|
||||
@@ -27216,7 +27245,7 @@ export function heartbeatService(
|
||||
// retain distinct durable receipts even when the same gate blocks them.
|
||||
const coalesceExecutionWait =
|
||||
opts.requestedByActorType === "system" &&
|
||||
!durableRequest &&
|
||||
!receiptRequest &&
|
||||
!wakeCommentId &&
|
||||
queuedCommentIdsFromRunContext(enrichedContextSnapshot).length === 0 &&
|
||||
!isInteractionResolutionWakePayload(payload ?? {}) &&
|
||||
@@ -28603,10 +28632,10 @@ export function heartbeatService(
|
||||
wakeupRequestId: activeExecutionRun.wakeupRequestId,
|
||||
},
|
||||
allowRunCoalescing: isConversation(issue) ? false : opts.allowRunCoalescing,
|
||||
durableReceipt: durableRequest
|
||||
durableReceipt: receiptRequest
|
||||
? {
|
||||
id: durableRequest.id,
|
||||
requestedAt: durableRequest.requestedAt,
|
||||
id: receiptRequest.id,
|
||||
requestedAt: receiptRequest.requestedAt,
|
||||
}
|
||||
: undefined,
|
||||
reason,
|
||||
|
||||
@@ -75,3 +75,18 @@ it("keeps a revoked binding cleared when the follow-up connection read fails", a
|
||||
expect(client.getQueryData(["dot-binding", "company", "agent"])).toMatchObject({ binding: null });
|
||||
expect(Array.from(container.querySelectorAll("button")).some(button => button.textContent === "Revoke connection")).toBe(false);
|
||||
});
|
||||
|
||||
it("shows event testing when a connected Dot subscribes without a manual refresh", async () => {
|
||||
api.get.mockResolvedValue({ enabled: true, resourceUrl: "https://paperclip.example/mcp/runner", binding: { ...binding, subscriptionVerified: false } });
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
await act(async () => root.render(<QueryClientProvider client={client}><Form /></QueryClientProvider>));
|
||||
await act(async () => { await vi.advanceTimersByTimeAsync(20); });
|
||||
expect(container.textContent).toContain("Event subscription: required");
|
||||
expect(container.textContent).not.toContain("Test event delivery");
|
||||
api.get.mockResolvedValue({ enabled: true, resourceUrl: "https://paperclip.example/mcp/runner", binding });
|
||||
await act(async () => { await vi.advanceTimersByTimeAsync(5020); });
|
||||
expect(container.textContent).toContain("Test event delivery");
|
||||
expect(container.textContent).toContain("Event subscription: verified");
|
||||
} finally { vi.useRealTimers(); }
|
||||
});
|
||||
@@ -20,7 +20,9 @@ export function DotRunnerConnection({ companyId, agentId, bindingId, onBinding }
|
||||
const key = ["dot-binding", companyId, agentId];
|
||||
const state = useQuery({ queryKey: key, enabled: !!companyId && !!agentId,
|
||||
queryFn: () => api.get<Connection>(path),
|
||||
refetchInterval: query => query.state.data?.binding?.status === "pairing" || query.state.data?.binding?.hasPendingChallenge || query.state.data?.binding?.assignment ? 5000 : false });
|
||||
refetchInterval: query => query.state.data?.binding?.status === "pairing" ||
|
||||
(query.state.data?.binding?.connected && !query.state.data.binding.subscriptionVerified) ||
|
||||
query.state.data?.binding?.hasPendingChallenge || query.state.data?.binding?.assignment ? 5000 : false });
|
||||
const pair = useMutation({ mutationFn: () => api.post<{ bindingId: string; pairingCode: string; expiresAt: string }>(path, {}),
|
||||
onSuccess: result => { setPairing(result); onBinding(result.bindingId); void client.invalidateQueries({ queryKey: key }); } });
|
||||
const test = useMutation({ mutationFn: () => api.post(path + "/event-test", {}),
|
||||
|
||||
Reference in new issue
Block a user