1
import { type Context, copyJson, type JsonValue } from "@earendil-works/chord";2
import { awaitWithContext } from "@earendil-works/chord/context";3
import { overlap } from "@earendil-works/chord/delta";4
import type { ImageContent, TextContent, ToolCall, ToolResultMessage } from "@earendil-works/pi-ai";5
import { validateToolArguments } from "@earendil-works/pi-ai/utils/validation";6
import { AssistantEntry, ToolResultEntry } from "../entries.ts";7
import { defineTask } from "../tasks.ts";8
import { DEFAULT_MAX_BYTES, DEFAULT_MAX_LINES, utf8ByteLength } from "../truncate.ts";9
import type {10
ConversationId,11
EntryId,12
JsonObject,13
Task,14
TaskId,15
TaskOptions,16
TaskRuntime,17
Tx,18
TypedEntry,19
} from "../types.ts";20
import { assignJson } from "./json.ts";21
import { clearProgress, finishSlot, LiveDoc, type ToolSlot, toolSlot } from "./live.ts";22
import { boundOutput, OutputBuffer, type OutputLimits, Progress } from "./output.ts";23
import type {24
ToolControl,25
ToolDiagnostic,26
ToolExecutionApi,27
ToolExecutionResult,28
ToolHooks,29
ToolRegistration,30
} from "./types.ts";31
import { recordUsage } from "./usage.ts";33
export type ToolTaskInput = { assistant: EntryId; callId: string };35
export 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" };40
export type ToolTaskResult = { entryId: EntryId; control?: ToolControl };42
type Runtime = TaskRuntime<ToolTaskInput, ToolTaskCheckpoint, ToolTaskResult, ToolHooks>;43
type Content = (TextContent | ImageContent)[];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 from48
* settlement. `execute` is reached only by recovery and applies the replay rule.49
*/50
export 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
});120
/** The tool call `callId` of the assistant entry. */121
async 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
}134
/** Arguments, or why they are invalid. */135
type Checked = { readonly args: JsonObject } | { readonly error: string };137
/** The call's arguments as repaired by the tool; a throwing repair makes them invalid. */138
function 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
}147
/** Arguments validated and coerced against the implementation's schema. */148
function 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
}156
function invalid(message: string): ToolExecutionResult {157
return harnessError("invalid_arguments", message);158
}160
/** What a running tool reported through its api: output, the last details, and diagnostics. */161
type Reported = {162
readonly output: OutputBuffer;163
readonly limits: OutputLimits;164
readonly diagnostics: ToolDiagnostic[];165
details: JsonValue | undefined;166
};168
/** Execute with the resolved implementation, then settle its result. */169
async 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
};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 text255
// 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
}272
/**273
* Throttled commits of what the tool reported into its `pi.live.tools` slot, each writing only what changed since the274
* last one.275
*/276
function 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 shared286
// 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.length291
: 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 an300
// 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
}324
/**325
* The settled result: the tool's result with the retained output and last details as fallbacks, its diagnostics after326
* those reported through the api, `afterTool` applied, and explicit text bounded, with the Harness's truncation327
* diagnostic last.328
*/329
async 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
}358
/**359
* Commit the tool's terminal state: append its result entry, mark its slot done, and complete or end aborted with the360
* entry ID. `build` receives the slot so interruption and abort can report the durable partial output.361
*/362
async 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 === undefined384
? {}385
: { control: copyJson(result.control as JsonValue, { omitUndefinedProperties: true }) as ToolControl };386
return { status: "terminal", outcome: { status: "completed", result: { entryId, ...control } } };387
}, context);388
}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
*/394
type Ending = { readonly status: "completed" | "aborted" } | { readonly status: "failed"; readonly message: string };396
const COMPLETED: Ending = { status: "completed" };398
/** An error result from the slot's durable partial output, details, and diagnostics. */399
function 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
}412
/** An error result the Harness writes itself: no content and one `error` diagnostic with `code`. */413
export function harnessError(code: string, message: string): ToolExecutionResult {414
return { content: [], isError: true, diagnostics: [toolDiagnostic(code, message)] };415
}417
function toolDiagnostic(code: string, message: string): ToolDiagnostic {418
return { severity: "error", code, message };419
}421
/** The Harness's truncation diagnostic; `retain` is unknown when rebuilt from a slot after recovery. */422
function 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
}434
/**435
* Append a `pi.tool-result` entry. The content ends with the rendered diagnostics, so the stored message is exactly436
* what the model sees; `data` keeps the structured list. A result's usage is added to `pi.usage` in the same commit.437
*/438
export 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
}462
function renderDiagnostics(diagnostics: readonly ToolDiagnostic[]): string {463
return `<harness>\n${diagnostics.map((diagnostic) => `[${diagnostic.severity}] ${diagnostic.message}`).join("\n")}\n</harness>`;464
}466
/**467
* Bound the text of result content. When the joined text exceeds the limits, the text items are replaced by one bounded468
* item at the position of the first (head) or last (tail) text item; other content is kept.469
*/470
function 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
}486
function errorText(error: unknown): string {487
return error instanceof Error ? error.message : String(error);488
}