返回源码地图

packages/durable/src/harness/tool.ts

v1.0.0 · a13d35a742c6 · 11:工具意图与 replay 恢复路径

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1import { type Context, copyJson, type JsonValue } from "@earendil-works/chord";
2import { awaitWithContext } from "@earendil-works/chord/context";
3import { overlap } from "@earendil-works/chord/delta";
4import type { ImageContent, TextContent, ToolCall, ToolResultMessage } from "@earendil-works/pi-ai";
5import { validateToolArguments } from "@earendil-works/pi-ai/utils/validation";
6import { AssistantEntry, ToolResultEntry } from "../entries.ts";
7import { defineTask } from "../tasks.ts";
8import { DEFAULT_MAX_BYTES, DEFAULT_MAX_LINES, utf8ByteLength } from "../truncate.ts";
9import type {
10 ConversationId,
11 EntryId,
12 JsonObject,
13 Task,
14 TaskId,
15 TaskOptions,
16 TaskRuntime,
17 Tx,
18 TypedEntry,
19} from "../types.ts";
20import { assignJson } from "./json.ts";
21import { clearProgress, finishSlot, LiveDoc, type ToolSlot, toolSlot } from "./live.ts";
22import { boundOutput, OutputBuffer, type OutputLimits, Progress } from "./output.ts";
23import type {
24 ToolControl,
25 ToolDiagnostic,
26 ToolExecutionApi,
27 ToolExecutionResult,
28 ToolHooks,
29 ToolRegistration,
30} from "./types.ts";
31import { recordUsage } from "./usage.ts";
32
33export type ToolTaskInput = { assistant: EntryId; callId: string };
34
35export type ToolTaskCheckpoint =
36 | { phase: "call" }
37 /** Durable intent: the final arguments and the replay policy recorded before execution. */
38 | { phase: "execute"; arguments: JsonObject; replay: "safe" | "unsafe" };
39
40export type ToolTaskResult = { entryId: EntryId; control?: ToolControl };
41
42type Runtime = TaskRuntime<ToolTaskInput, ToolTaskCheckpoint, ToolTaskResult, ToolHooks>;
43type Content = (TextContent | ImageContent)[];
44
45/**
46 * Built-in tool task: resolves the called tool among its phase agent's tools, validates, runs `beforeTool`, records intent,
47 * executes, runs `afterTool`, and appends the result, all in one `call` handler so nothing separates resolution from
48 * settlement. `execute` is reached only by recovery and applies the replay rule.
49 */
50export const ToolTask = defineTask<ToolTaskInput, ToolTaskCheckpoint, ToolTaskResult, ToolHooks>({
51 name: "pi.tool",
52 version: 1,
53 initial: () => ({ phase: "call" }),
54 phases: {
55 call: async (task, runtime, context) => {
56 const call = await readCall(runtime, task.input, context);
57 const tool = (await runtime.agent(context)).tools.find((each) => each.name === call.name);
58 if (tool === undefined) {
59 const error = harnessError("tool_unavailable", `Tool ${call.name} is not available`);
60 return settle(runtime, call, COMPLETED, () => error, context);
61 }
62 const prepared = prepare(tool, call.arguments as JsonObject);
63 const checked = "error" in prepared ? prepared : validate(tool, call, prepared.args);
64 if ("error" in checked) return settle(runtime, call, COMPLETED, () => invalid(checked.error), context);
65 let args = checked.args;
66 let block: string | undefined;
67 await runtime.hooks.each("beforeTool", async (hook) => {
68 if (block !== undefined) return;
69 try {
70 const decision = await hook({ ...call, arguments: args }, runtime, context);
71 if (decision?.block !== undefined) block = decision.block;
72 else if (decision?.arguments !== undefined) args = decision.arguments;
73 } catch (error) {
74 if (runtime.signal.aborted) throw error;
75 block = errorText(error);
76 }
77 });
78 if (block !== undefined) {
79 const blocked = harnessError("blocked", `Tool call blocked: ${block}`);
80 return settle(runtime, call, COMPLETED, () => blocked, context);
81 }
82 const validated = validate(tool, call, args);
83 if ("error" in validated) return settle(runtime, call, COMPLETED, () => invalid(validated.error), context);
84 const final = validated.args;
85 await runtime.commit(async (tx) => {
86 const slot = toolSlot(await tx.doc(LiveDoc, runtime.conversationId), runtime.taskId);
87 if (slot !== undefined) slot.status = "running";
88 const intent = { phase: "execute", arguments: final, replay: tool.replay ?? "unsafe" } as const;
89 return { status: "running", checkpoint: intent };
90 }, context);
91 await run(runtime, call, tool, final, context);
92 },
93 /** Recovery after intent: rerun only when the stored and the current policy both say `safe`. */
94 execute: async (task, runtime, context) => {
95 const { arguments: args, replay } = task.state.checkpoint;
96 const call = await readCall(runtime, task.input, context);
97 const tool = (await runtime.agent(context)).tools.find((each) => each.name === call.name);
98 if (replay === "safe" && tool?.replay === "safe") {
99 // The rerun reports from scratch; clear what the interrupted attempt published.
100 await runtime.commit(async (tx) => {
101 const slot = toolSlot(await tx.doc(LiveDoc, runtime.conversationId), runtime.taskId);
102 if (slot !== undefined) clearProgress(slot);
103 return undefined;
104 }, context);
105 return run(runtime, call, tool, args, context);
106 }
107 const message = `Tool ${call.name} was interrupted and may have partially run`;
108 // `failed` records cancellation intent, so the call's owned conversations, left unsupervised, are aborted.
109 const ending = { status: "failed", message } as const;
110 await settle(runtime, call, ending, (slot) => fromSlot(slot, "interrupted", message), context);
111 },
112 },
113 abort: async (task, runtime, context) => {
114 const call = await readCall(runtime, task.input, context);
115 const message = `Tool ${call.name} was aborted`;
116 await settle(runtime, call, { status: "aborted" }, (slot) => fromSlot(slot, "aborted", message), context);
117 },
118});
119
120/** The tool call `callId` of the assistant entry. */
121async function readCall(runtime: Runtime, input: ToolTaskInput, context: Context): Promise<ToolCall> {
122 const entry = await runtime.entry(AssistantEntry, input.assistant, context);
123 const message = entry?.model?.[0];
124 const call =
125 message?.role === "assistant"
126 ? message.content.find(
127 (content): content is ToolCall => content.type === "toolCall" && content.id === input.callId,
128 )
129 : undefined;
130 if (call === undefined) throw new Error(`Entry ${input.assistant} has no tool call ${input.callId}`);
131 return call;
132}
133
134/** Arguments, or why they are invalid. */
135type Checked = { readonly args: JsonObject } | { readonly error: string };
136
137/** The call's arguments as repaired by the tool; a throwing repair makes them invalid. */
138function prepare(tool: ToolRegistration, args: JsonObject): Checked {
139 if (tool.prepareArguments === undefined) return { args };
140 try {
141 return { args: tool.prepareArguments(args) as JsonObject };
142 } catch (error) {
143 return { error: errorText(error) };
144 }
145}
146
147/** Arguments validated and coerced against the implementation's schema. */
148function validate(tool: ToolRegistration, call: ToolCall, args: JsonObject): Checked {
149 try {
150 return { args: validateToolArguments(tool, { ...call, arguments: args }) as JsonObject };
151 } catch (error) {
152 return { error: errorText(error) };
153 }
154}
155
156function invalid(message: string): ToolExecutionResult {
157 return harnessError("invalid_arguments", message);
158}
159
160/** What a running tool reported through its api: output, the last details, and diagnostics. */
161type Reported = {
162 readonly output: OutputBuffer;
163 readonly limits: OutputLimits;
164 readonly diagnostics: ToolDiagnostic[];
165 details: JsonValue | undefined;
166};
167
168/** Execute with the resolved implementation, then settle its result. */
169async function run(
170 runtime: Runtime,
171 call: ToolCall,
172 tool: ToolRegistration,
173 args: JsonObject,
174 context: Context,
175): Promise<void> {
176 const limits: OutputLimits = {
177 maxBytes: tool.outputLimits?.maxBytes ?? DEFAULT_MAX_BYTES,
178 maxLines: tool.outputLimits?.maxLines ?? DEFAULT_MAX_LINES,
179 retain: tool.outputLimits?.retain ?? "head",
180 };
181 const reported: Reported = { output: new OutputBuffer(limits), limits, diagnostics: [], details: undefined };
182 const progress = publishProgress(runtime, reported, context);
183 let ended = false;
184 const assertLive = (): void => {
185 if (ended) throw new Error(`Tool call ${call.id} has settled`);
186 };
187 const api: Omit<ToolExecutionApi, "env"> = {
188 taskId: runtime.taskId,
189 conversationId: runtime.conversationId,
190 callId: call.id,
191 registry: runtime.registry,
192 agent: runtime.agent,
193 output: (chunk) => {
194 assertLive();
195 if (reported.output.push(chunk)) progress.mark();
196 },
197 diagnostic: (diagnostic) => {
198 assertLive();
199 reported.diagnostics.push(copyJson(diagnostic, { omitUndefinedProperties: true }) as ToolDiagnostic);
200 progress.mark();
201 },
202 details: async (value, detailsContext) => {
203 assertLive();
204 detailsContext.abortSignal?.throwIfAborted();
205 reported.details = copyJson(value, { omitUndefinedProperties: true });
206 const committed = progress.markAndWait();
207 // Cancelling the wait leaves the update in place; the commit's own outcome stays observed.
208 committed.catch(() => {});
209 return awaitWithContext(committed, detailsContext);
210 },
211 commit: async (change, commitContext) => {
212 let result: Awaited<ReturnType<typeof change>> | undefined;
213 await runtime.commit(async (tx) => {
214 result = await change(tx);
215 return undefined;
216 }, commitContext);
217 return result as Awaited<ReturnType<typeof change>>;
218 },
219 memo: runtime.memo,
220 createTask: async <I, S extends { phase: string }, R, H extends object>(
221 task: Task<I, S, R, H>,
222 input: I,
223 options: Omit<TaskOptions, "conversationId">,
224 taskContext: Context,
225 ): Promise<TaskId<R>> => {
226 let id: TaskId<R> | undefined;
227 await runtime.commit(async (tx) => {
228 id = await tx.createTask(task, input, options);
229 return undefined;
230 }, taskContext);
231 return id!;
232 },
233 getTask: runtime.getTask,
234 waitForTask: runtime.waitForTask,
235 conversation: runtime.conversation,
236 snapshot: runtime.snapshot,
237 snapshotAsOf: runtime.snapshotAsOf,
238 watchDoc: runtime.watchDoc,
239 };
240
241 let result: ToolExecutionResult;
242 let ending = COMPLETED;
243 try {
244 // Built for this call, so a rerun after recovery gets the conversation's environment at that time.
245 const env = await runtime.env(context);
246 result = await tool.execute(args, { ...api, env }, context);
247 } catch (error) {
248 if (runtime.signal.aborted) {
249 ended = true;
250 for (const waiter of await progress.stop()) waiter.reject(error);
251 throw error;
252 }
253 result = { isError: true, diagnostics: [toolDiagnostic("tool_error", errorText(error))] };
254 // A throw, from `execute()` or from building the environment, ends the task `failed`, which cancels what the call owned; it no longer supervises it. The error text
255 // is already in the result entry.
256 ending = { status: "failed", message: `Tool ${call.name} threw` };
257 }
258 ended = true;
259 reported.output.end();
260 // Details still waiting for a progress commit settle with the terminal commit, the final flush.
261 const pending = await progress.stop();
262 try {
263 const settled = await finalResult(runtime, call, result, reported, context);
264 await settle(runtime, call, ending, () => settled, context);
265 } catch (error) {
266 for (const waiter of pending) waiter.reject(error);
267 throw error;
268 }
269 for (const waiter of pending) waiter.resolve();
270}
271
272/**
273 * Throttled commits of what the tool reported into its `pi.live.tools` slot, each writing only what changed since the
274 * last one.
275 */
276function publishProgress(runtime: Runtime, reported: Reported, context: Context): Progress {
277 let written = { text: "", details: undefined as JsonValue | undefined, diagnostics: 0 };
278 return new Progress(
279 async () => {
280 // Capture everything synchronously: the tool keeps reporting while the commit is in flight.
281 const snapshot = reported.output.snapshot();
282 const current = { text: snapshot.text, details: reported.details, diagnostics: reported.diagnostics.length };
283 const added = reported.diagnostics.slice(written.diagnostics, current.diagnostics);
284 const detailsChanged = current.details !== written.details;
285 // What the commit writes, as Chord diffs the string: an append, a trim plus an append of what follows the shared
286 // part, or the whole window when its bounded overlap search finds nothing.
287 let bytes = 0;
288 if (snapshot.text !== written.text) {
289 const shared = snapshot.text.startsWith(written.text)
290 ? written.text.length
291 : overlap(written.text, snapshot.text, 65_536);
292 bytes += utf8ByteLength(snapshot.text.slice(shared));
293 }
294 if (detailsChanged) bytes += utf8ByteLength(JSON.stringify(current.details ?? null));
295 if (added.length > 0) bytes += utf8ByteLength(JSON.stringify(added));
296 await runtime.commit(async (tx) => {
297 const slot = toolSlot(await tx.doc(LiveDoc, runtime.conversationId), runtime.taskId);
298 if (slot === undefined) return undefined;
299 // REMINDER: assign `output` as one string field. Chord then diffs it into an append, or a trim plus an
300 // append for a sliding tail; replacing the slot object would record the whole window on every commit.
301 if ((slot.output ?? "") !== snapshot.text) slot.output = snapshot.text;
302 if (snapshot.droppedBytes > 0) slot.droppedBytes = snapshot.droppedBytes;
303 if (snapshot.droppedLines > 0) slot.droppedLines = snapshot.droppedLines;
304 // Diff details leaf by leaf and append new diagnostics, so each commit writes only what changed.
305 if (detailsChanged && current.details !== undefined) {
306 assignJson(slot as unknown as Record<string, JsonValue>, "details", current.details);
307 }
308 if (added.length > 0) {
309 if (slot.diagnostics === undefined) slot.diagnostics = [];
310 for (const diagnostic of added) slot.diagnostics.push(diagnostic);
311 }
312 return undefined;
313 }, context);
314 written = current;
315 return bytes;
316 },
317 (error) => {
318 // Rejections after an abort mark or close are expected; the committed state stays consistent.
319 if (!runtime.signal.aborted) runtime.report(error);
320 },
321 );
322}
323
324/**
325 * The settled result: the tool's result with the retained output and last details as fallbacks, its diagnostics after
326 * those reported through the api, `afterTool` applied, and explicit text bounded, with the Harness's truncation
327 * diagnostic last.
328 */
329async function finalResult(
330 runtime: Runtime,
331 call: ToolCall,
332 result: ToolExecutionResult,
333 reported: Reported,
334 context: Context,
335): Promise<ToolExecutionResult> {
336 const harness: ToolDiagnostic[] = [];
337 const retained = result.content === undefined ? reported.output.snapshot() : undefined;
338 const content: Content =
339 retained === undefined ? result.content! : retained.text === "" ? [] : [{ type: "text", text: retained.text }];
340 let final: ToolExecutionResult = {
341 ...result,
342 content,
343 details: result.details === undefined ? reported.details : result.details,
344 diagnostics: [...reported.diagnostics, ...(result.diagnostics ?? [])],
345 };
346 await runtime.hooks.each("afterTool", async (hook) => {
347 final = (await hook(call, final, runtime, context)) ?? final;
348 });
349 // The retained output's truncation applies only while afterTool kept that content.
350 if (retained !== undefined && retained.droppedBytes > 0 && final.content === content) {
351 harness.push(truncated(retained, reported.limits.retain));
352 }
353 const bounded = boundContent(final.content ?? [], reported.limits);
354 if (bounded.droppedBytes > 0) harness.push(truncated(bounded, reported.limits.retain));
355 return { ...final, content: bounded.content, diagnostics: [...(final.diagnostics ?? []), ...harness] };
356}
357
358/**
359 * Commit the tool's terminal state: append its result entry, mark its slot done, and complete or end aborted with the
360 * entry ID. `build` receives the slot so interruption and abort can report the durable partial output.
361 */
362async function settle(
363 runtime: Runtime,
364 call: ToolCall,
365 ending: Ending,
366 build: (slot: Readonly<ToolSlot> | undefined) => ToolExecutionResult,
367 context: Context,
368): Promise<void> {
369 await runtime.commit(async (tx) => {
370 const slot = toolSlot(await tx.doc(LiveDoc, runtime.conversationId), runtime.taskId);
371 const result = build(slot);
372 const entry = await appendToolResult(tx, runtime.conversationId, call, result, runtime.now());
373 if (slot !== undefined) finishSlot(slot, entry.id);
374 const entryId = entry.id;
375 if (ending.status === "aborted")
376 return { status: "terminal", outcome: { status: "aborted", result: { entryId } } };
377 if (ending.status === "failed") {
378 const error = { message: ending.message };
379 return { status: "terminal", outcome: { status: "failed", error, result: { entryId } } };
380 }
381 // Tools build control objects freely; drop keys set to undefined so the task result is strict JSON.
382 const control =
383 result.control === undefined
384 ? {}
385 : { control: copyJson(result.control as JsonValue, { omitUndefinedProperties: true }) as ToolControl };
386 return { status: "terminal", outcome: { status: "completed", result: { entryId, ...control } } };
387 }, context);
388}
389
390/**
391 * How a tool task ends; the result entry is appended either way. `failed` (execution threw or was interrupted)
392 * records cancellation intent for the conversations the call owns; a result with `isError` still completes.
393 */
394type Ending = { readonly status: "completed" | "aborted" } | { readonly status: "failed"; readonly message: string };
395
396const COMPLETED: Ending = { status: "completed" };
397
398/** An error result from the slot's durable partial output, details, and diagnostics. */
399function fromSlot(slot: Readonly<ToolSlot> | undefined, code: string, message: string): ToolExecutionResult {
400 const diagnostics = [...(slot?.diagnostics ?? [])];
401 const droppedBytes = slot?.droppedBytes ?? 0;
402 if (droppedBytes > 0) diagnostics.push(truncated({ droppedBytes, droppedLines: slot?.droppedLines ?? 0 }));
403 diagnostics.push(toolDiagnostic(code, message));
404 return {
405 content: slot?.output === undefined || slot.output === "" ? [] : [{ type: "text", text: slot.output }],
406 isError: true,
407 ...(slot?.details === undefined ? {} : { details: slot.details }),
408 diagnostics,
409 };
410}
411
412/** An error result the Harness writes itself: no content and one `error` diagnostic with `code`. */
413export function harnessError(code: string, message: string): ToolExecutionResult {
414 return { content: [], isError: true, diagnostics: [toolDiagnostic(code, message)] };
415}
416
417function toolDiagnostic(code: string, message: string): ToolDiagnostic {
418 return { severity: "error", code, message };
419}
420
421/** The Harness's truncation diagnostic; `retain` is unknown when rebuilt from a slot after recovery. */
422function truncated(
423 dropped: { readonly droppedLines: number; readonly droppedBytes: number },
424 retain?: "head" | "tail",
425): ToolDiagnostic {
426 const kept = retain === undefined ? "" : ` to its ${retain === "head" ? "beginning" : "end"}`;
427 return {
428 severity: "warn",
429 code: "truncated",
430 message: `Output truncated${kept}: ${dropped.droppedLines} lines, ${dropped.droppedBytes} bytes dropped`,
431 };
432}
433
434/**
435 * Append a `pi.tool-result` entry. The content ends with the rendered diagnostics, so the stored message is exactly
436 * what the model sees; `data` keeps the structured list. A result's usage is added to `pi.usage` in the same commit.
437 */
438export async function appendToolResult(
439 tx: Tx,
440 conversationId: ConversationId,
441 call: ToolCall,
442 result: ToolExecutionResult,
443 timestamp: number,
444): Promise<TypedEntry<{ diagnostics: ToolDiagnostic[] }>> {
445 const diagnostics = [...(result.diagnostics ?? [])];
446 const content: Content = [...(result.content ?? [])];
447 if (diagnostics.length > 0) content.push({ type: "text", text: renderDiagnostics(diagnostics) });
448 const message = {
449 role: "toolResult",
450 toolCallId: call.id,
451 toolName: call.name,
452 content,
453 ...(result.details === undefined ? {} : { details: result.details }),
454 ...(result.usage === undefined ? {} : { usage: result.usage }),
455 isError: result.isError ?? false,
456 timestamp,
457 } as ToolResultMessage;
458 if (result.usage !== undefined) await recordUsage(tx, conversationId, "tools", call.name, result.usage);
459 return tx.appendEntry(ToolResultEntry, conversationId, { model: [message], data: { diagnostics } });
460}
461
462function renderDiagnostics(diagnostics: readonly ToolDiagnostic[]): string {
463 return `<harness>\n${diagnostics.map((diagnostic) => `[${diagnostic.severity}] ${diagnostic.message}`).join("\n")}\n</harness>`;
464}
465
466/**
467 * Bound the text of result content. When the joined text exceeds the limits, the text items are replaced by one bounded
468 * item at the position of the first (head) or last (tail) text item; other content is kept.
469 */
470function boundContent(
471 content: Content,
472 limits: OutputLimits,
473): { content: Content; droppedBytes: number; droppedLines: number } {
474 const texts = content.filter((item): item is TextContent => item.type === "text");
475 const bounded = boundOutput(texts.map((item) => item.text).join(""), limits);
476 if (bounded.droppedBytes === 0) return { content, droppedBytes: 0, droppedLines: 0 };
477 const keep = limits.retain === "head" ? texts[0] : texts.at(-1);
478 const result: Content = [];
479 for (const item of content) {
480 if (item.type !== "text") result.push(item);
481 else if (item === keep) result.push({ ...item, text: bounded.text });
482 }
483 return { content: result, droppedBytes: bounded.droppedBytes, droppedLines: bounded.droppedLines };
484}
485
486function errorText(error: unknown): string {
487 return error instanceof Error ? error.message : String(error);
488}