Connect live native questions and park only after tool results settle

Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
DottaandPaperclip committed 2026-09-18 22:05:48 -05:00
1 parent 9f4133d4ed
commit 814fef7f32
10 files changed
+174 -9

No files matched your search

@@ -29,6 +29,15 @@ describe("published hiring and human-input examples", () => {
for (const { body } of waits) expect(updateIssueSchema.safeParse(substituteIds(body))).toMatchObject({ success: true });
});
it("includes a complete valid text-field recipe in the skill itself", () => {
const skill = readFileSync(new URL("../../../skills/paperclip/SKILL.md", import.meta.url), "utf8");
const section = skill.split("**Asking a free-text question.**")[1]!;
const body = JSON.parse(section.match(/```json\n([\s\S]*?)\n```/)![1]);
expect(createIssueThreadInteractionSchema.safeParse(substituteIds(body))).toMatchObject({ success: true });
expect(body.payload.questionSet.questions[0]).toMatchObject({ answerMode: "text" });
expect(body.payload.questions[0].id).toBe(body.payload.questionSet.questions[0].id);
});
it("keeps these examples in the generated runner reference without displacing confirmations", () => {
for (const example of [...questions, ...hires, ...waits]) {
const key = `${example.method} ${example.path.replace(/\{[^}]+\}/g, "{}")}`;
@@ -0,0 +1,54 @@
import { and, eq } from "drizzle-orm";
import { heartbeatRuns, type Db } from "@paperclipai/db";
import type { HarnessRuntimeRequestResolution, PrpEvent } from "../../vendor/paperclip-runner/index.js";
import { flushNativeQuestionResponses, projectNativeRuntimeRequest, registerNativeQuestionCommandTarget } from "./native-question-bridge.js";
import { readPendingNativeRuntimeRequest } from "./runtime-request-resolution-authority.js";
/** The in-process executor must perform the same card projection and response
* delivery as the durable PRP coordinator. The answer remains durable in DB. */
export function createLocalNativeQuestionBridge(input: {
db: Db;
binding: Parameters<typeof projectNativeRuntimeRequest>[0]["binding"];
resolve: (input: {
runId: string; requestId: string; turnId: string;
resolution: HarnessRuntimeRequestResolution;
authorizeBeforeDispatch: () => Promise<void>;
}) => Promise<{ commandId: string }>;
}) {
let release: (() => void) | undefined;
const close = () => { release?.(); release = undefined; };
return {
close,
async attach() {
close();
release = registerNativeQuestionCommandTarget({
binding: input.binding,
queueCommand: async (type, payload) => {
if (type !== "request.resolve" || typeof payload?.requestId !== "string") throw new Error("native_question_command_invalid");
const requestId = payload.requestId;
const pending = await readPendingNativeRuntimeRequest(input.db, { ...input.binding, requestId });
if (!pending || pending.requestKind !== "runtime") throw new Error("native_question_not_pending");
const result = await input.resolve({
runId: input.binding.runId, requestId, turnId: pending.turnId,
resolution: { action: "submit", response: payload.response as never },
authorizeBeforeDispatch: async () => {
const current = await readPendingNativeRuntimeRequest(input.db, { ...input.binding, requestId });
const [run] = await input.db.select({ status: heartbeatRuns.status }).from(heartbeatRuns).where(and(
eq(heartbeatRuns.id, input.binding.runId), eq(heartbeatRuns.companyId, input.binding.companyId),
eq(heartbeatRuns.nativeIssueId, input.binding.issueId), eq(heartbeatRuns.agentId, input.binding.agentId),
)).limit(1);
if (run?.status !== "running" || current?.turnId !== pending.turnId || current.requestKind !== "runtime") throw new Error("native_question_not_pending");
},
});
return { commandId: result.commandId, controllerSeq: 0 };
},
});
await flushNativeQuestionResponses(input.db, input.binding.runId);
},
async observe(event: PrpEvent) {
const request = event.payload.request as Record<string, unknown> | undefined;
if (event.eventType !== "runtime_request.created" || request?.type !== "input") return;
await projectNativeRuntimeRequest({ db: input.db, binding: input.binding, event });
},
};
}
@@ -1,3 +1,4 @@
import { createLocalNativeQuestionBridge } from "./local-native-question-bridge.js";
import { randomUUID } from "node:crypto";
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "vitest";
import { eq, sql } from "drizzle-orm";
@@ -8,6 +9,7 @@ import {
companies,
createDb,
heartbeatRuns,
heartbeatRunEvents,
issueQuestionResponseDeliveries,
issueThreadInteractions,
issues,
@@ -199,6 +201,34 @@ describeEmbeddedPostgres("native question bridge", () => {
};
}
it("projects an executor question immediately and routes its durable answer into the same live turn", async () => {
await seed();
const event = runtimeRequestEvent();
await db.insert(heartbeatRunEvents).values({ companyId, agentId, runId, seq: 1,
eventType: event.eventType, stream: "system", level: "info", payload: { prpEvent: event } });
const resolve = vi.fn(async (input: any) => { await input.authorizeBeforeDispatch(); return { commandId: "live-response" }; });
const bridge = createLocalNativeQuestionBridge({ db, binding: binding(), resolve });
try {
await bridge.attach();
await bridge.observe(event);
await bridge.observe(event); // replay must not create a second card
const cards = await issueThreadInteractionService(db).listForIssue(issueId);
expect(cards).toHaveLength(1);
expect(cards[0]).toMatchObject({ status: "pending", sourceRunId: runId, continuationPolicy: "none" });
const answered = await issueThreadInteractionService(db).answerQuestions(
{ id: issueId, companyId, status: "in_progress" }, cards[0]!.id,
{ answers: [{ questionId: "color", optionIds: ["green"] }] }, { userId: "operator-1" },
);
if (answered.kind !== "ask_user_questions") throw new Error("wrong question kind");
expect(await deliverNativeQuestionResponse(db, answered)).toBe("queued");
expect(resolve).toHaveBeenCalledWith(expect.objectContaining({ runId, requestId: "request-1", turnId: "turn-1",
resolution: { action: "submit", response: { schema: "paperclip.question_response.v1", answers: { color: { selectedOptionIds: ["green"] } } } },
}));
await db.update(heartbeatRuns).set({ status: "cancelled" }).where(eq(heartbeatRuns.id, runId));
await expect(resolve.mock.calls[0]![0].authorizeBeforeDispatch()).rejects.toThrow("native_question_not_pending");
} finally { bridge.close(); }
});
it.each(["codex", "claude"])("materializes, validates, and durably resumes a %s question response", async (provider) => {
await seed();
const interaction = await projectNativeRuntimeRequest({
@@ -34,7 +34,7 @@ type QueueCommand = (
type: string,
payload?: Record<string, unknown>,
commandId?: string,
) => { readonly commandId: string; readonly controllerSeq: number };
) => { readonly commandId: string; readonly controllerSeq: number } | Promise<{ readonly commandId: string; readonly controllerSeq: number }>;
interface NativeQuestionCommandTarget {
binding: Pick<NativeRunStoreBinding, "companyId" | "issueId" | "runId" | "agentId">;
@@ -203,7 +203,7 @@ async function authorizedNativeRun(
/** Materialize a canonical runtime input request as the existing task-thread card. */
export async function projectNativeRuntimeRequest(input: {
db: Db;
binding: NativeRunStoreBinding;
binding: Pick<NativeRunStoreBinding, "companyId" | "issueId" | "runId" | "agentId" | "normalizedSessionId" | "runnerSourceInstanceId">;
event: PrpEvent;
}): Promise<AskUserQuestionsInteraction | null> {
if (input.event.eventType !== "runtime_request.created") return null;
@@ -319,7 +319,7 @@ export async function deliverNativeQuestionResponse(
return "pending";
}
try {
target.queueCommand(
await target.queueCommand(
"request.resolve",
{ requestId: run.requestId, response: response as unknown as Record<string, unknown> },
`question_${interaction.id}`,
@@ -4125,9 +4125,20 @@ describe("native governed waits", () => {
schemaVersion: 1,
priority: 0 as const,
emittedAt: "2026-08-31T00:00:00.000Z",
payload: {},
payload: { kind: "dynamicToolCall" },
};
// A usage event must not park while the card-creation response is held.
await observation.observe({ ...replayedEvent, eventType: "item.started", itemId: "approval-tool" }, false);
const usage = { ...replayedEvent, payload: { kind: "usage" } };
await observation.observe(usage, true);
expect(observation.consume(usage)).toBeNull();
const other = { ...replayedEvent, itemId: "other-tool" };
await observation.observe(other, true);
expect(observation.consume(other)).toBeNull();
await observation.observe({ ...replayedEvent, itemId: "approval-tool" }, true);
expect(observation.consume({ ...replayedEvent, itemId: "approval-tool" })).toEqual(waitResult);
await observation.observe(replayedEvent, true);
expect(observation.consume(replayedEvent)).toEqual(waitResult);
expect(observation.consume(replayedEvent)).toBeNull();
@@ -4237,6 +4248,7 @@ function leaseDb(
const query = {
then: Promise.resolve(rows).then.bind(Promise.resolve(rows)),
where: () => query,
orderBy: () => query,
for: () => query,
limit: () => Promise.resolve(rows),
};
@@ -4841,7 +4853,7 @@ describe("native runtime request resolution", () => {
snapshot.mockReset().mockResolvedValue({ activeTurnId: "provider-turn-1" });
resolveRuntimeRequest.mockReset().mockResolvedValue(undefined);
state.execute.mockReset().mockImplementation(async (options) => {
options.onSession?.({
await options.onSession?.({
capabilities,
snapshot,
resolveRuntimeRequest,
@@ -4850,7 +4862,7 @@ describe("native runtime request resolution", () => {
await new Promise<void>((resolve) => {
state.release = resolve;
});
options.onSession?.(null);
await options.onSession?.(null);
return {
result: { summary: "completed" },
terminal: { runTerminalState: "succeeded" },
@@ -1,3 +1,4 @@
import { createLocalNativeQuestionBridge } from "./local-native-question-bridge.js";
import { readVerifiedRemoteWorkspaceFile } from "./remote-deliverable-file.js";
import { copyBackCodexAuth } from "@paperclipai/adapter-codex-local/server";
import { nativeCompletionFeedback } from "./native-completion-feedback.js";
@@ -982,6 +983,7 @@ export function nativeConversationReplyResult(input: {
export function createGovernedWaitEventObservation(
resolvePending: () => Promise<PrpStructuredRunResult | null>,
) {
const pendingTools = new Set<string>();
let generation = 0;
let observation: {
sourceInstanceId: string;
@@ -994,6 +996,18 @@ export function createGovernedWaitEventObservation(
async observe(event: PrpEvent, eligible: boolean): Promise<void> {
const currentGeneration = ++generation;
observation = null;
const kind = record(event.payload).kind;
const tool = ["dynamicToolCall", "mcpToolCall", "commandExecution"].includes(String(kind));
if (tool && event.itemId) {
if (event.eventType === "item.started") pendingTools.add(event.itemId);
if (event.eventType === "item.completed") pendingTools.delete(event.itemId);
}
// Usage/model messages can arrive while the tool creating the card is
// still awaiting its response. Parking then interrupts that in-flight
// response and cannot produce a durable suspension checkpoint.
if (event.eventType === "item.completed" && (
pendingTools.size > 0 || (!tool && kind !== "agentMessage")
)) return;
if (!eligible) return;
const result = await resolvePending();
if (generation !== currentGeneration || result === null) return;
@@ -7384,6 +7398,11 @@ async function executePaperclipNativeSessionWithinScope(
payload: event.payload,
},
);
const liveQuestions = createLocalNativeQuestionBridge({
db: input.db,
binding: { ...input.execution.binding, normalizedSessionId: nativeSessionKey(input.execution), runnerSourceInstanceId: effectiveRunnerInstanceId },
resolve: resolveNativeRuntimeRequest,
});
let completedConversationReply: PrpEvent | null = null;
const controlPlane = new PaperclipControlPlanePort(
input.db,
@@ -7405,6 +7424,7 @@ async function executePaperclipNativeSessionWithinScope(
record(event.payload).channel === "final") {
completedConversationReply = event;
}
await liveQuestions.observe(event);
await projectSessionGoalEvent(event);
providerUsageLimitObserved ||= nativeProviderUsageLimitFromEvent(event);
const eventAtMs = Date.parse(event.emittedAt);
@@ -7621,6 +7641,7 @@ async function executePaperclipNativeSessionWithinScope(
// A crash can happen after the event commit but before its callback
// finishes. Recover only idempotent durable projections here; activity,
// publication, logging, trace, and metric effects remain committed-only.
await liveQuestions.observe(event);
await projectSessionGoalEvent(event);
providerUsageLimitObserved ||= nativeProviderUsageLimitFromEvent(event);
const questionFallback = await materializeRuntimeQuestionFallback({
@@ -7763,6 +7784,9 @@ async function executePaperclipNativeSessionWithinScope(
: []),
),
eq(issueThreadInteractions.status, "pending"),
// Live provider questions resume their current turn; only durable
// wake-based cards park it. A timeout creates a separate fallback.
sql`not (${issueThreadInteractions.kind} = 'ask_user_questions' and ${issueThreadInteractions.continuationPolicy} = 'none' and coalesce(${issueThreadInteractions.idempotencyKey}, '') like 'paperclip-runner-question:%')`,
),
)
.orderBy(
@@ -7991,6 +8015,7 @@ async function executePaperclipNativeSessionWithinScope(
},
onSession: async (session) => {
releaseRegisteredGoalController();
liveQuestions.close();
if (session?.goal) {
releaseGoalController = registerLiveRunnerGoalController(
{
@@ -8055,10 +8080,12 @@ async function executePaperclipNativeSessionWithinScope(
session,
cancelRequested: false,
});
if (session.resolveRuntimeRequest) await liveQuestions.attach();
if (nativeRunsDetachingForRestart.has(input.execution.binding.runId)) {
await session.detachControllerForRestart?.();
}
} else {
liveQuestions.close();
activeNativeSessions.delete(input.execution.binding.runId);
clearSteeringDeliveries(input.execution.binding.runId);
clearNativeRuntimeRequestResolutions(
@@ -8094,12 +8121,14 @@ async function executePaperclipNativeSessionWithinScope(
startedAtMs: turnCompletedAtMs ?? nativeSessionExecuteStartedAtMs,
endedAtMs: Date.now(),
});
liveQuestions.close();
activeNativeSessions.delete(input.execution.binding.runId);
clearSteeringDeliveries(input.execution.binding.runId);
clearNativeRuntimeRequestResolutions(input.execution.binding.runId);
} catch (error) {
if (nativeRunsDetachingForRestart.has(input.execution.binding.runId)) {
await leaseRenewal.stop().catch(() => undefined);
liveQuestions.close();
activeNativeSessions.delete(input.execution.binding.runId);
// Disconnecting deliberately ends the old event consumer. It is not a
// provider failure and must not overwrite the shutdown adoption record
@@ -8145,6 +8174,7 @@ async function executePaperclipNativeSessionWithinScope(
});
}
trace.activate(taskSettleScope);
liveQuestions.close();
activeNativeSessions.delete(input.execution.binding.runId);
clearSteeringDeliveries(input.execution.binding.runId);
clearNativeRuntimeRequestResolutions(input.execution.binding.runId);
+20 -1
View File
@@ -690,4 +690,23 @@ Again, rule #1 is: never ask a human to do what an agent could do. Try harder. T
**Asking a free-text question.**
For an open answer, use a text field, not invented choices. Copy the complete **Text answer** interaction example in [references/api-reference.md](references/api-reference.md#questions-and-waiting-for-human-input): it includes both the canonical `payload.questionSet` (`answerMode: "text"`) and required storage fields. The storage fallback alone renders the wrong control.
For an open answer, use a text field, not invented choices. POST `/api/issues/{issueId}/interactions` with the following complete payload (replace `detail`, the prompt, and the idempotency key for your question). `questionSet` controls presentation; the matching `questions` entry is required storage compatibility and must not be sent alone.
```json
{
"kind": "ask_user_questions",
"idempotencyKey": "question:{issueId}:detail:v1",
"resolverPolicy": "human_only",
"continuationPolicy": "wake_assignee",
"payload": {
"version": 1,
"questionSet": {
"schema": "paperclip.question_set.v1",
"questions": [{ "id": "detail", "prompt": "What should I know?", "answerMode": "text", "required": true }]
},
"questions": [{ "id": "detail", "prompt": "What should I know?", "selectionMode": "single", "required": true, "options": [{ "id": "text", "label": "Your answer", "freeText": true }] }]
}
}
```
See [the API reference](references/api-reference.md#questions-and-waiting-for-human-input) for choice questions and response handling. Include the normal Authorization and X-Paperclip-Run-Id headers.
+1
View File
@@ -83,6 +83,7 @@ function commonAgent(
...(adapterType === "paperclip_runner"
? []
: [
"Read the assigned paperclip skill before using the Paperclip API. Use its documented request payloads rather than guessing fields from memory.",
'For a planning task, do not inspect the OpenAPI schema. PUT /api/issues/$PAPERCLIP_TASK_ID/documents/plan with {title:"Plan",format:"markdown",body,changeSummary}; read latestRevisionId and latestRevisionNumber from that response. Then POST /api/issues/$PAPERCLIP_TASK_ID/interactions with {kind:"request_confirmation",continuationPolicy:"wake_assignee",payload:{version:1,prompt,acceptLabel:"Approve",rejectLabel:"Reject",rejectRequiresReason:true,target:{type:"issue_document",key:"plan",revisionId,revisionNumber}}}, and PATCH the issue to {status:"in_review"}. Include Authorization and X-Paperclip-Run-Id on every write.',
]),
"Never print, persist, or expose credential values, and never create unrelated work.",
+9
View File
@@ -237,6 +237,15 @@ describe("chat acceptance contracts", () => {
"Run log not found",
);
});
it("retains events for an unstarted dependency-blocked wake without asking for a nonexistent log", async () => {
const get = vi.fn().mockResolvedValue([]);
const suppressed = { ...run, status: "cancelled", errorCode: "issue_dependencies_blocked", startedAt: null };
expect((await collectChatRunEvidence({ get }, suppressed)).log).toBeNull();
expect(get).toHaveBeenCalledTimes(1);
get.mockRejectedValue(new Error("Run log not found"));
await expect(collectChatRunEvidence({ get }, { ...suppressed, startedAt: "2026-09-18T00:00:00Z" })).rejects.toThrow("Run log not found");
await expect(collectChatRunEvidence({ get }, { ...suppressed, errorCode: "provider_transport_failed" })).rejects.toThrow("Run log not found");
});
it("waits for a newly running provider's log file without swallowing server failures", async () => {
const get = vi.fn().mockResolvedValue({ status: () => 404 });
const api = { request: { get } } as unknown as Pick<RunnerApi, "request">;
+3 -2
View File
@@ -6,6 +6,7 @@ import type {
} from "../../packages/shared/src/types/issue.js";
import type { LiveFixtureValues } from "./live-fixtures.js";
import type { MatrixExecution } from "./types.js";
import { isBlockedUnstartedWake } from "./non-execution-wake.js";
import { chatMarker } from "./chat-cases.js";
// Public API observations only: this driver never fabricates provider results or writes DB state.
@@ -36,7 +37,7 @@ export interface ChatRun {
resultJson?: Record<string, unknown>;
sessionIdBefore?: string | null;
sessionIdAfter?: string | null;
startedAt?: string;
startedAt?: string | null;
}
type Comment = {
id: string;
@@ -212,7 +213,7 @@ export async function collectChatRunEvidence(
) {
return {
runId: run.id,
log: isResetRun(run)
log: isResetRun(run) || isBlockedUnstartedWake({ ...run })
? null
: await api.get(`/api/heartbeat-runs/${run.id}/log?limitBytes=1048576`),
events: await api.get(`/api/heartbeat-runs/${run.id}/events?limit=1000`),