ππ Agent

agent-harness.ts

37 KB1085 lines
agent-harness.ts
1import {
2 type AssistantMessage,
3 contentText,
4 type ImageContent,
5 type Model,
6 type Models,
7 type RetryCallbacks,
8 type RetryPolicy,
9 type UserMessage,
10} from "@earendil-works/pi-ai";
11import { runAgentLoop } from "../agent-loop.ts";
12import type {
13 AgentContext,
14 AgentEvent,
15 AgentLoopConfig,
16 AgentMessage,
17 AgentTool,
18 QueueMode,
19 StreamFn,
20 ThinkingLevel,
21} from "../types.ts";
22import { collectEntriesForBranchSummary, generateBranchSummary } from "./compaction/branch-summarization.ts";
23import { compact, DEFAULT_COMPACTION_SETTINGS, prepareCompaction } from "./compaction/compaction.ts";
24import { convertToLlm } from "./messages.ts";
25import { formatPromptTemplateInvocation } from "./prompt-templates.ts";
26import { formatSkillInvocation } from "./skills.ts";
27import type {
28 AbortResult,
29 AgentHarnessEvent,
30 AgentHarnessEventResultMap,
31 AgentHarnessOptions,
32 AgentHarnessOwnEvent,
33 AgentHarnessPhase,
34 AgentHarnessResources,
35 AgentHarnessStreamOptions,
36 AgentHarnessStreamOptionsPatch,
37 AgentHarnessSystemPrompt,
38 AgentHarnessTool,
39 AgentHarnessToolContextSource,
40 CompactResult,
41 NavigateTreeResult,
42 PendingSessionWrite,
43 PromptTemplate,
44 Session,
45 Skill,
46} from "./types.ts";
47import { AgentHarnessError, BranchSummaryError, CompactionError, SessionError, toError } from "./types.ts";
48
49function createUserMessage(text: string, images?: ImageContent[]): UserMessage {
50 const content: Array<{ type: "text"; text: string } | ImageContent> = [{ type: "text", text }];
51 if (images) content.push(...images);
52 return { role: "user", content, timestamp: Date.now() };
53}
54
55function createFailureMessage(model: Model<any>, error: unknown, aborted: boolean): AssistantMessage {
56 return {
57 role: "assistant",
58 content: [{ type: "text", text: "" }],
59 api: model.api,
60 provider: model.provider,
61 model: model.id,
62 stopReason: aborted ? "aborted" : "error",
63 errorMessage: error instanceof Error ? error.message : String(error),
64 timestamp: Date.now(),
65 usage: {
66 input: 0,
67 output: 0,
68 cacheRead: 0,
69 cacheWrite: 0,
70 totalTokens: 0,
71 cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
72 },
73 };
74}
75
76function cloneStreamOptions(streamOptions?: AgentHarnessStreamOptions): AgentHarnessStreamOptions {
77 return {
78 ...streamOptions,
79 headers: streamOptions?.headers ? { ...streamOptions.headers } : undefined,
80 metadata: streamOptions?.metadata ? { ...streamOptions.metadata } : undefined,
81 };
82}
83
84function findDuplicateNames(names: string[]): string[] {
85 const seen = new Set<string>();
86 const duplicates = new Set<string>();
87 for (const name of names) {
88 if (seen.has(name)) duplicates.add(name);
89 seen.add(name);
90 }
91 return [...duplicates];
92}
93
94function applyStreamOptionsPatch(
95 base: AgentHarnessStreamOptions,
96 patch?: AgentHarnessStreamOptionsPatch,
97): AgentHarnessStreamOptions {
98 const result = cloneStreamOptions(base);
99 if (!patch) return result;
100
101 if (Object.hasOwn(patch, "transport")) result.transport = patch.transport;
102 if (Object.hasOwn(patch, "timeoutMs")) result.timeoutMs = patch.timeoutMs;
103 if (Object.hasOwn(patch, "maxRetries")) result.maxRetries = patch.maxRetries;
104 if (Object.hasOwn(patch, "maxRetryDelayMs")) result.maxRetryDelayMs = patch.maxRetryDelayMs;
105 if (Object.hasOwn(patch, "cacheRetention")) result.cacheRetention = patch.cacheRetention;
106
107 if (Object.hasOwn(patch, "headers")) {
108 if (patch.headers === undefined) {
109 result.headers = undefined;
110 } else {
111 const headers = { ...(result.headers ?? {}) };
112 for (const [key, value] of Object.entries(patch.headers)) {
113 if (value === undefined) delete headers[key];
114 else headers[key] = value;
115 }
116 result.headers = Object.keys(headers).length > 0 ? headers : undefined;
117 }
118 }
119
120 if (Object.hasOwn(patch, "metadata")) {
121 if (patch.metadata === undefined) {
122 result.metadata = undefined;
123 } else {
124 const metadata = { ...(result.metadata ?? {}) };
125 for (const [key, value] of Object.entries(patch.metadata)) {
126 if (value === undefined) delete metadata[key];
127 else metadata[key] = value;
128 }
129 result.metadata = Object.keys(metadata).length > 0 ? metadata : undefined;
130 }
131 }
132
133 return result;
134}
135
136const SUBSCRIBER_EVENT_TYPE = "*";
137
138type AgentHarnessHandler = (event: any, signal?: AbortSignal) => Promise<any> | any;
139
140function normalizeHarnessError(error: unknown, fallbackCode: AgentHarnessError["code"]): AgentHarnessError {
141 if (error instanceof AgentHarnessError) return error;
142 const cause = toError(error);
143 if (cause instanceof SessionError) return new AgentHarnessError("session", cause.message, cause);
144 if (cause instanceof CompactionError) return new AgentHarnessError("compaction", cause.message, cause);
145 if (cause instanceof BranchSummaryError) return new AgentHarnessError("branch_summary", cause.message, cause);
146 return new AgentHarnessError(fallbackCode, cause.message, cause);
147}
148
149function normalizeHookError(error: unknown): AgentHarnessError {
150 return normalizeHarnessError(error, "hook");
151}
152
153interface AgentHarnessTurnState<
154 TContext extends object | undefined,
155 TSkill extends Skill = Skill,
156 TPromptTemplate extends PromptTemplate = PromptTemplate,
157 TTool extends AgentHarnessTool<TContext> = AgentHarnessTool<TContext>,
158> {
159 messages: AgentMessage[];
160 resources: AgentHarnessResources<TSkill, TPromptTemplate>;
161 toolContext: TContext;
162 streamOptions: AgentHarnessStreamOptions;
163 sessionId: string;
164 systemPrompt: string;
165 model: Model<any>;
166 thinkingLevel: ThinkingLevel;
167 tools: TTool[];
168 activeTools: TTool[];
169}
170
171export class AgentHarness<
172 TContext extends object | undefined = undefined,
173 TSkill extends Skill = Skill,
174 TPromptTemplate extends PromptTemplate = PromptTemplate,
175 TTool extends AgentHarnessTool<TContext> = AgentHarnessTool<TContext>,
176> {
177 private session: Session;
178 readonly models: Models;
179 private phase: AgentHarnessPhase = "idle";
180 private runAbortController?: AbortController;
181 private runPromise?: Promise<void>;
182 private pendingSessionWrites: PendingSessionWrite[] = [];
183 private model: Model<any>;
184 private thinkingLevel: ThinkingLevel;
185 private systemPrompt: AgentHarnessSystemPrompt<TContext, TSkill, TPromptTemplate, TTool> | undefined;
186 private toolContext: AgentHarnessToolContextSource<TContext> | undefined;
187 private streamOptions: AgentHarnessStreamOptions;
188 private retry: RetryPolicy | undefined;
189 private resources: AgentHarnessResources<TSkill, TPromptTemplate>;
190 private tools = new Map<string, TTool>();
191 private activeToolNames: string[];
192 private steerQueue: UserMessage[] = [];
193 private steeringQueueMode: QueueMode;
194 private followUpQueue: UserMessage[] = [];
195 private followUpQueueMode: QueueMode;
196 private nextTurnQueue: AgentMessage[] = [];
197 private handlers = new Map<string, Set<AgentHarnessHandler>>();
198
199 constructor(options: AgentHarnessOptions<TContext, TSkill, TPromptTemplate, TTool>) {
200 this.session = options.session;
201 this.models = options.models;
202 this.resources = options.resources ?? {};
203 this.streamOptions = cloneStreamOptions(options.streamOptions);
204 this.retry = options.retry;
205 this.systemPrompt = options.systemPrompt;
206 this.toolContext = options.toolContext;
207 this.validateUniqueNames(
208 (options.tools ?? []).map((tool) => tool.name),
209 "Duplicate tool name(s)",
210 );
211 for (const tool of options.tools ?? []) {
212 this.tools.set(tool.name, tool);
213 }
214 this.model = options.model;
215 this.thinkingLevel = options.thinkingLevel ?? "off";
216 this.activeToolNames = options.activeToolNames
217 ? [...options.activeToolNames]
218 : (options.tools ?? []).map((tool) => tool.name);
219 this.validateUniqueNames(this.activeToolNames, "Duplicate active tool name(s)");
220 this.validateToolNames(this.activeToolNames);
221 this.steeringQueueMode = options.steeringMode ?? "one-at-a-time";
222 this.followUpQueueMode = options.followUpMode ?? "one-at-a-time";
223 }
224
225 private getHandlers(type: string): Set<AgentHarnessHandler> | undefined {
226 return this.handlers.get(type);
227 }
228
229 private async emitOwn(event: AgentHarnessOwnEvent<TSkill, TPromptTemplate>, signal?: AbortSignal): Promise<void> {
230 for (const listener of this.getHandlers(SUBSCRIBER_EVENT_TYPE) ?? []) {
231 try {
232 await listener(event, signal);
233 } catch (error) {
234 throw normalizeHookError(error);
235 }
236 }
237 }
238
239 private async emitAny(event: AgentHarnessEvent<TSkill, TPromptTemplate>, signal?: AbortSignal): Promise<void> {
240 for (const listener of this.getHandlers(SUBSCRIBER_EVENT_TYPE) ?? []) {
241 try {
242 await listener(event, signal);
243 } catch (error) {
244 throw normalizeHookError(error);
245 }
246 }
247 }
248
249 private async emitHook<TType extends keyof AgentHarnessEventResultMap>(
250 event: Extract<AgentHarnessOwnEvent, { type: TType }>,
251 ): Promise<AgentHarnessEventResultMap[TType] | undefined> {
252 const handlers = this.getHandlers(event.type as TType);
253 if (!handlers || handlers.size === 0) return undefined;
254 let lastResult: AgentHarnessEventResultMap[TType] | undefined;
255 for (const handler of handlers) {
256 try {
257 const result = await handler(event);
258 if (result !== undefined) {
259 lastResult = result;
260 }
261 } catch (error) {
262 throw normalizeHookError(error);
263 }
264 }
265 return lastResult;
266 }
267
268 private retryCallbacks(operation: "compaction" | "branch_summary"): RetryCallbacks {
269 return {
270 onRetryScheduled: (attempt, maxAttempts, delayMs, errorMessage) =>
271 this.emitOwn({ type: "retry_scheduled", operation, attempt, maxAttempts, delayMs, errorMessage }),
272 onRetryAttemptStart: () => this.emitOwn({ type: "retry_attempt_start", operation }),
273 onRetryFinished: () => this.emitOwn({ type: "retry_finished", operation }),
274 };
275 }
276
277 private async emitBeforeProviderRequest(
278 model: Model<any>,
279 sessionId: string,
280 streamOptions: AgentHarnessStreamOptions,
281 ): Promise<AgentHarnessStreamOptions> {
282 const handlers = this.getHandlers("before_provider_request");
283 let current = cloneStreamOptions(streamOptions);
284 if (!handlers || handlers.size === 0) return current;
285 for (const handler of handlers) {
286 try {
287 const result = await handler({
288 type: "before_provider_request",
289 model,
290 sessionId,
291 streamOptions: cloneStreamOptions(current),
292 });
293 if (result?.streamOptions) {
294 current = applyStreamOptionsPatch(current, result.streamOptions);
295 }
296 } catch (error) {
297 throw normalizeHookError(error);
298 }
299 }
300 return current;
301 }
302
303 private async emitBeforeProviderPayload(model: Model<any>, payload: unknown): Promise<unknown> {
304 const handlers = this.getHandlers("before_provider_payload");
305 let current = payload;
306 if (!handlers || handlers.size === 0) return current;
307 for (const handler of handlers) {
308 try {
309 const result = await handler({ type: "before_provider_payload", model, payload: current });
310 if (result !== undefined) {
311 current = result.payload;
312 }
313 } catch (error) {
314 throw normalizeHookError(error);
315 }
316 }
317 return current;
318 }
319
320 private async emitQueueUpdate(): Promise<void> {
321 await this.emitOwn({
322 type: "queue_update",
323 steer: [...this.steerQueue],
324 followUp: [...this.followUpQueue],
325 nextTurn: [...this.nextTurnQueue],
326 });
327 }
328
329 private startRunPromise(): () => void {
330 let finish = () => {};
331 this.runPromise = new Promise<void>((resolve) => {
332 finish = resolve;
333 });
334 return () => {
335 this.runPromise = undefined;
336 finish();
337 };
338 }
339
340 private async resolveToolContext(): Promise<TContext> {
341 if (typeof this.toolContext === "function") {
342 return await (this.toolContext as () => TContext | Promise<TContext>)();
343 }
344 return this.toolContext as TContext;
345 }
346
347 private bindToolContext(tool: TTool, context: TContext): AgentTool {
348 return {
349 ...tool,
350 execute: (toolCallId, params, signal, onUpdate) => tool.execute(toolCallId, params, signal, onUpdate, context),
351 };
352 }
353
354 private async createTurnState(): Promise<AgentHarnessTurnState<TContext, TSkill, TPromptTemplate, TTool>> {
355 const context = await this.session.buildContext();
356 const resources = this.getResources();
357 const sessionMetadata = await this.session.getMetadata();
358 const toolContext = await this.resolveToolContext();
359 const tools = [...this.tools.values()];
360 const activeTools = this.activeToolNames
361 .map((name) => this.tools.get(name))
362 .filter((tool): tool is TTool => tool !== undefined);
363 let systemPrompt = "You are a helpful assistant.";
364 if (typeof this.systemPrompt === "string") {
365 systemPrompt = this.systemPrompt;
366 } else if (this.systemPrompt) {
367 systemPrompt = await this.systemPrompt({
368 session: this.session,
369 model: this.model,
370 thinkingLevel: this.thinkingLevel,
371 activeTools,
372 resources,
373 });
374 }
375 return {
376 messages: context.messages,
377 resources,
378 toolContext,
379 streamOptions: cloneStreamOptions(this.streamOptions),
380 sessionId: sessionMetadata.id,
381 systemPrompt,
382 model: this.model,
383 thinkingLevel: this.thinkingLevel,
384 tools,
385 activeTools,
386 };
387 }
388
389 private createContext(
390 turnState: AgentHarnessTurnState<TContext, TSkill, TPromptTemplate, TTool>,
391 systemPrompt?: string,
392 ): AgentContext {
393 return {
394 systemPrompt: systemPrompt ?? turnState.systemPrompt,
395 messages: turnState.messages.slice(),
396 tools: turnState.activeTools.map((tool) => this.bindToolContext(tool, turnState.toolContext)),
397 };
398 }
399
400 private createStreamFn(
401 getTurnState: () => AgentHarnessTurnState<TContext, TSkill, TPromptTemplate, TTool>,
402 ): StreamFn {
403 return async (model, context, streamOptions) => {
404 const turnState = getTurnState();
405 const snapshotOptions: AgentHarnessStreamOptions = { ...turnState.streamOptions };
406 const requestOptions = await this.emitBeforeProviderRequest(model, turnState.sessionId, snapshotOptions);
407 return this.models.streamSimple(model, context, {
408 cacheRetention: requestOptions.cacheRetention,
409 headers: requestOptions.headers,
410 maxRetries: requestOptions.maxRetries,
411 maxRetryDelayMs: requestOptions.maxRetryDelayMs,
412 metadata: requestOptions.metadata,
413 onPayload: async (payload) => await this.emitBeforeProviderPayload(model, payload),
414 onResponse: async (response) => {
415 const headers = { ...(response.headers as Record<string, string>) };
416 await this.emitOwn(
417 { type: "after_provider_response", status: response.status, headers },
418 streamOptions?.signal,
419 );
420 },
421 reasoning: streamOptions?.reasoning,
422 signal: streamOptions?.signal,
423 sessionId: turnState.sessionId,
424 timeoutMs: requestOptions.timeoutMs,
425 transport: requestOptions.transport,
426 });
427 };
428 }
429
430 private async drainQueuedMessages(queue: AgentMessage[], mode: QueueMode): Promise<AgentMessage[]> {
431 const messages = mode === "all" ? queue.splice(0) : queue.splice(0, 1);
432 if (messages.length === 0) return messages;
433 try {
434 await this.emitQueueUpdate();
435 return messages;
436 } catch (error) {
437 queue.unshift(...messages);
438 throw normalizeHookError(error);
439 }
440 }
441
442 private createLoopConfig(
443 getTurnState: () => AgentHarnessTurnState<TContext, TSkill, TPromptTemplate, TTool>,
444 setTurnState: (turnState: AgentHarnessTurnState<TContext, TSkill, TPromptTemplate, TTool>) => void,
445 ): AgentLoopConfig {
446 const turnState = getTurnState();
447 return {
448 model: turnState.model,
449 reasoning: turnState.thinkingLevel === "off" ? undefined : turnState.thinkingLevel,
450 convertToLlm,
451 transformContext: async (messages) => {
452 const result = await this.emitHook({ type: "context", messages: [...messages] });
453 return result?.messages ?? messages;
454 },
455 beforeToolCall: async ({ toolCall, args }) => {
456 const result = await this.emitHook({
457 type: "tool_call",
458 toolCallId: toolCall.id,
459 toolName: toolCall.name,
460 input: args as Record<string, unknown>,
461 });
462 return result ? { block: result.block, reason: result.reason } : undefined;
463 },
464 afterToolCall: async ({ toolCall, args, result, isError }) => {
465 const patch = await this.emitHook({
466 type: "tool_result",
467 toolCallId: toolCall.id,
468 toolName: toolCall.name,
469 input: args as Record<string, unknown>,
470 content: result.content,
471 details: result.details,
472 isError,
473 usage: result.usage,
474 });
475 return patch
476 ? {
477 content: patch.content,
478 details: patch.details,
479 isError: patch.isError,
480 usage: patch.usage,
481 terminate: patch.terminate,
482 }
483 : undefined;
484 },
485 prepareNextTurn: async () => {
486 await this.flushPendingSessionWrites();
487 const nextTurnState = await this.createTurnState();
488 setTurnState(nextTurnState);
489 return {
490 context: this.createContext(nextTurnState),
491 model: nextTurnState.model,
492 thinkingLevel: nextTurnState.thinkingLevel,
493 };
494 },
495 getSteeringMessages: async () => this.drainQueuedMessages(this.steerQueue, this.steeringQueueMode),
496 getFollowUpMessages: async () => this.drainQueuedMessages(this.followUpQueue, this.followUpQueueMode),
497 };
498 }
499
500 private validateUniqueNames(names: string[], message: string): void {
501 const duplicates = findDuplicateNames(names);
502 if (duplicates.length > 0)
503 throw new AgentHarnessError("invalid_argument", `${message}: ${duplicates.join(", ")}`);
504 }
505
506 private validateToolNames(toolNames: string[], tools: Map<string, TTool> = this.tools): void {
507 this.validateUniqueNames(toolNames, "Duplicate active tool name(s)");
508 const missing = toolNames.filter((name) => !tools.has(name));
509 if (missing.length > 0) throw new AgentHarnessError("invalid_argument", `Unknown tool(s): ${missing.join(", ")}`);
510 }
511
512 private async flushPendingSessionWrites(): Promise<void> {
513 while (this.pendingSessionWrites.length > 0) {
514 const write = this.pendingSessionWrites[0]!;
515 if (write.type === "message") {
516 await this.session.appendMessage(write.message);
517 } else if (write.type === "model_change") {
518 await this.session.appendModelChange(write.provider, write.modelId);
519 } else if (write.type === "thinking_level_change") {
520 await this.session.appendThinkingLevelChange(write.thinkingLevel);
521 } else if (write.type === "active_tools_change") {
522 await this.session.appendActiveToolsChange(write.activeToolNames);
523 } else if (write.type === "custom") {
524 await this.session.appendCustomEntry(write.customType, write.data);
525 } else if (write.type === "custom_message") {
526 await this.session.appendCustomMessageEntry(write.customType, write.content, write.display, write.details);
527 } else if (write.type === "label") {
528 await this.session.appendLabel(write.targetId, write.label);
529 } else if (write.type === "session_info") {
530 await this.session.appendSessionName(write.name ?? "");
531 } else if (write.type === "leaf") {
532 await this.session.getStorage().setLeafId(write.targetId);
533 }
534 this.pendingSessionWrites.shift();
535 }
536 }
537
538 private async handleAgentEvent(event: AgentEvent, signal?: AbortSignal): Promise<void> {
539 if (event.type === "message_end") {
540 await this.session.appendMessage(event.message);
541 await this.emitAny(event, signal);
542 return;
543 }
544 if (event.type === "turn_end") {
545 let eventError: unknown;
546 try {
547 await this.emitAny(event, signal);
548 } catch (error) {
549 eventError = error;
550 }
551 const hadPendingMutations = this.pendingSessionWrites.length > 0;
552 await this.flushPendingSessionWrites();
553 if (eventError) throw eventError;
554 await this.emitOwn({ type: "save_point", hadPendingMutations });
555 return;
556 }
557 if (event.type === "agent_end") {
558 await this.flushPendingSessionWrites();
559 this.phase = "idle";
560 await this.emitAny(event, signal);
561 await this.emitOwn({ type: "settled", nextTurnCount: this.nextTurnQueue.length }, signal);
562 return;
563 }
564 await this.emitAny(event, signal);
565 }
566
567 private async emitRunFailure(
568 model: Model<any>,
569 error: unknown,
570 aborted: boolean,
571 signal: AbortSignal,
572 ): Promise<AgentMessage[]> {
573 const failureMessage = createFailureMessage(model, error, aborted);
574 await this.handleAgentEvent({ type: "message_start", message: failureMessage }, signal);
575 await this.handleAgentEvent({ type: "message_end", message: failureMessage }, signal);
576 await this.handleAgentEvent({ type: "turn_end", message: failureMessage, toolResults: [] }, signal);
577 await this.handleAgentEvent({ type: "agent_end", messages: [failureMessage] }, signal);
578 return [failureMessage];
579 }
580
581 private async executeTurn(
582 turnState: AgentHarnessTurnState<TContext, TSkill, TPromptTemplate, TTool>,
583 text: string,
584 options?: { images?: ImageContent[] },
585 ): Promise<AssistantMessage> {
586 let activeTurnState = turnState;
587 let messages: AgentMessage[] = [createUserMessage(text, options?.images)];
588 if (this.nextTurnQueue.length > 0) {
589 const queuedMessages = this.nextTurnQueue.splice(0);
590 try {
591 await this.emitQueueUpdate();
592 } catch (error) {
593 this.nextTurnQueue.unshift(...queuedMessages);
594 throw normalizeHookError(error);
595 }
596 messages = [...queuedMessages, messages[0]!];
597 }
598 const beforeResult = await this.emitHook({
599 type: "before_agent_start",
600 prompt: text,
601 images: options?.images,
602 systemPrompt: turnState.systemPrompt,
603 resources: turnState.resources,
604 });
605 if (beforeResult?.messages) messages = [...messages, ...beforeResult.messages];
606
607 const abortController = new AbortController();
608 const getTurnState = () => activeTurnState;
609 const setTurnState = (nextTurnState: AgentHarnessTurnState<TContext, TSkill, TPromptTemplate, TTool>) => {
610 activeTurnState = nextTurnState;
611 };
612 this.runAbortController = abortController;
613 const runResultPromise = (async () => {
614 try {
615 return await runAgentLoop(
616 messages,
617 this.createContext(turnState, beforeResult?.systemPrompt),
618 this.createLoopConfig(getTurnState, setTurnState),
619 (event) => this.handleAgentEvent(event, abortController.signal),
620 abortController.signal,
621 this.createStreamFn(getTurnState),
622 );
623 } catch (error) {
624 try {
625 return await this.emitRunFailure(
626 activeTurnState.model,
627 error,
628 abortController.signal.aborted,
629 abortController.signal,
630 );
631 } catch (failureError) {
632 const cause = new AggregateError(
633 [toError(error), toError(failureError)],
634 "Agent run failed and failure reporting failed",
635 );
636 throw new AgentHarnessError("unknown", cause.message, cause);
637 }
638 }
639 })();
640 try {
641 const newMessages = await runResultPromise;
642 for (let i = newMessages.length - 1; i >= 0; i--) {
643 const message = newMessages[i]!;
644 if (message.role === "assistant") {
645 return message;
646 }
647 }
648 throw new AgentHarnessError("invalid_state", "AgentHarness prompt completed without an assistant message");
649 } finally {
650 try {
651 await this.flushPendingSessionWrites();
652 } finally {
653 this.runAbortController = undefined;
654 }
655 }
656 }
657
658 async prompt(text: string, options?: { images?: ImageContent[] }): Promise<AssistantMessage> {
659 if (this.phase !== "idle") throw new AgentHarnessError("busy", "AgentHarness is busy");
660 this.phase = "turn";
661 const finishRunPromise = this.startRunPromise();
662 try {
663 const turnState = await this.createTurnState();
664 return await this.executeTurn(turnState, text, options);
665 } catch (error) {
666 this.phase = "idle";
667 throw normalizeHarnessError(error, "unknown");
668 } finally {
669 finishRunPromise();
670 }
671 }
672
673 async skill(name: string, additionalInstructions?: string): Promise<AssistantMessage> {
674 if (this.phase !== "idle") throw new AgentHarnessError("busy", "AgentHarness is busy");
675 this.phase = "turn";
676 const finishRunPromise = this.startRunPromise();
677 try {
678 const turnState = await this.createTurnState();
679 const skill = (turnState.resources.skills ?? []).find((candidate) => candidate.name === name);
680 if (!skill) throw new AgentHarnessError("invalid_argument", `Unknown skill: ${name}`);
681 return await this.executeTurn(turnState, formatSkillInvocation(skill, additionalInstructions));
682 } catch (error) {
683 this.phase = "idle";
684 throw normalizeHarnessError(error, "unknown");
685 } finally {
686 finishRunPromise();
687 }
688 }
689
690 async promptFromTemplate(name: string, args: string[] = []): Promise<AssistantMessage> {
691 if (this.phase !== "idle") throw new AgentHarnessError("busy", "AgentHarness is busy");
692 this.phase = "turn";
693 const finishRunPromise = this.startRunPromise();
694 try {
695 const turnState = await this.createTurnState();
696 const template = (turnState.resources.promptTemplates ?? []).find((candidate) => candidate.name === name);
697 if (!template) throw new AgentHarnessError("invalid_argument", `Unknown prompt template: ${name}`);
698 return await this.executeTurn(turnState, formatPromptTemplateInvocation(template, args));
699 } catch (error) {
700 this.phase = "idle";
701 throw normalizeHarnessError(error, "unknown");
702 } finally {
703 finishRunPromise();
704 }
705 }
706
707 async steer(text: string, options?: { images?: ImageContent[] }): Promise<void> {
708 if (this.phase === "idle") throw new AgentHarnessError("invalid_state", "Cannot steer while idle");
709 this.steerQueue.push(createUserMessage(text, options?.images));
710 await this.emitQueueUpdate();
711 }
712
713 async followUp(text: string, options?: { images?: ImageContent[] }): Promise<void> {
714 if (this.phase === "idle") throw new AgentHarnessError("invalid_state", "Cannot follow up while idle");
715 this.followUpQueue.push(createUserMessage(text, options?.images));
716 await this.emitQueueUpdate();
717 }
718
719 async nextTurn(text: string, options?: { images?: ImageContent[] }): Promise<void> {
720 this.nextTurnQueue.push(createUserMessage(text, options?.images));
721 await this.emitQueueUpdate();
722 }
723
724 async appendMessage(message: AgentMessage): Promise<void> {
725 try {
726 if (this.phase === "idle") {
727 await this.session.appendMessage(message);
728 } else {
729 this.pendingSessionWrites.push({ type: "message", message });
730 }
731 } catch (error) {
732 throw normalizeHarnessError(error, "session");
733 }
734 }
735
736 async compact(customInstructions?: string): Promise<CompactResult> {
737 if (this.phase !== "idle") throw new AgentHarnessError("busy", "compact() requires idle harness");
738 this.phase = "compaction";
739 try {
740 const model = this.model;
741 if (!model) throw new AgentHarnessError("invalid_state", "No model set for compaction");
742 const branchEntries = await this.session.getBranch();
743 const preparationResult = prepareCompaction(branchEntries, DEFAULT_COMPACTION_SETTINGS);
744 if (!preparationResult.ok) throw preparationResult.error;
745 const preparation = preparationResult.value;
746 if (!preparation) throw new AgentHarnessError("compaction", "Nothing to compact");
747 const hookResult = await this.emitHook({
748 type: "session_before_compact",
749 preparation,
750 branchEntries,
751 customInstructions,
752 signal: new AbortController().signal,
753 });
754 if (hookResult?.cancel) throw new AgentHarnessError("compaction", "Compaction cancelled");
755 const provided = hookResult?.compaction;
756 const compactResult = provided
757 ? { ok: true as const, value: provided }
758 : await compact(
759 preparation,
760 this.models,
761 model,
762 customInstructions,
763 undefined,
764 this.thinkingLevel,
765 this.retry,
766 this.retryCallbacks("compaction"),
767 );
768 if (!compactResult.ok) throw compactResult.error;
769 const result = compactResult.value;
770 const entryId = await this.session.appendCompaction(
771 result.summary,
772 result.firstKeptEntryId,
773 result.tokensBefore,
774 result.details,
775 provided !== undefined,
776 result.usage,
777 result.retainedTail,
778 );
779 const entry = await this.session.getEntry(entryId);
780 if (entry?.type === "compaction") {
781 await this.emitOwn({ type: "session_compact", compactionEntry: entry, fromHook: provided !== undefined });
782 }
783 return result;
784 } catch (error) {
785 throw normalizeHarnessError(error, "compaction");
786 } finally {
787 this.phase = "idle";
788 }
789 }
790
791 async navigateTree(
792 targetId: string,
793 options?: { summarize?: boolean; customInstructions?: string; replaceInstructions?: boolean; label?: string },
794 ): Promise<NavigateTreeResult> {
795 if (this.phase !== "idle") throw new AgentHarnessError("busy", "navigateTree() requires idle harness");
796 this.phase = "branch_summary";
797 try {
798 const oldLeafId = await this.session.getLeafId();
799 if (oldLeafId === targetId) return { cancelled: false };
800 const targetEntry = await this.session.getEntry(targetId);
801 if (!targetEntry) throw new AgentHarnessError("invalid_argument", `Entry ${targetId} not found`);
802 const { entries, commonAncestorId } = await collectEntriesForBranchSummary(this.session, oldLeafId, targetId);
803 const preparation = {
804 targetId,
805 oldLeafId,
806 commonAncestorId,
807 entriesToSummarize: entries,
808 userWantsSummary: options?.summarize ?? false,
809 customInstructions: options?.customInstructions,
810 replaceInstructions: options?.replaceInstructions,
811 label: options?.label,
812 };
813 const signal = new AbortController().signal;
814 const hookResult = await this.emitHook({ type: "session_before_tree", preparation, signal });
815 if (hookResult?.cancel) return { cancelled: true };
816 let summaryEntry: NavigateTreeResult["summaryEntry"];
817 let summaryText: string | undefined = hookResult?.summary?.summary;
818 let summaryDetails: unknown = hookResult?.summary?.details;
819 let summaryUsage = hookResult?.summary?.usage;
820 if (!summaryText && options?.summarize && entries.length > 0) {
821 const model = this.model;
822 if (!model) throw new AgentHarnessError("invalid_state", "No model set for branch summary");
823 const branchSummary = await generateBranchSummary(entries, {
824 models: this.models,
825 model,
826 signal: new AbortController().signal,
827 customInstructions: hookResult?.customInstructions ?? options?.customInstructions,
828 replaceInstructions: hookResult?.replaceInstructions ?? options?.replaceInstructions,
829 retry: this.retry,
830 callbacks: this.retryCallbacks("branch_summary"),
831 });
832 if (!branchSummary.ok) {
833 if (branchSummary.error.code === "aborted") return { cancelled: true };
834 throw new AgentHarnessError("branch_summary", branchSummary.error.message, branchSummary.error);
835 }
836 summaryText = branchSummary.value.summary;
837 summaryUsage = branchSummary.value.usage;
838 summaryDetails = {
839 readFiles: branchSummary.value.readFiles,
840 modifiedFiles: branchSummary.value.modifiedFiles,
841 };
842 }
843 let editorText: string | undefined;
844 let newLeafId: string | null;
845 if (targetEntry.type === "message" && targetEntry.message.role === "user") {
846 newLeafId = targetEntry.parentId;
847 editorText = contentText(targetEntry.message.content, "");
848 } else if (targetEntry.type === "custom_message") {
849 newLeafId = targetEntry.parentId;
850 editorText = contentText(targetEntry.content, "");
851 } else {
852 newLeafId = targetId;
853 }
854 const summaryId = await this.session.moveTo(
855 newLeafId,
856 summaryText
857 ? {
858 summary: summaryText,
859 details: summaryDetails,
860 usage: summaryUsage,
861 fromHook: hookResult?.summary !== undefined,
862 }
863 : undefined,
864 );
865 if (summaryId) {
866 const entry = await this.session.getEntry(summaryId);
867 if (entry?.type === "branch_summary") summaryEntry = entry;
868 }
869 await this.emitOwn({
870 type: "session_tree",
871 newLeafId: await this.session.getLeafId(),
872 oldLeafId,
873 summaryEntry,
874 fromHook: hookResult?.summary !== undefined,
875 });
876 return { cancelled: false, editorText, summaryEntry };
877 } catch (error) {
878 throw normalizeHarnessError(error, "branch_summary");
879 } finally {
880 this.phase = "idle";
881 }
882 }
883
884 getModel(): Model<any> {
885 return this.model;
886 }
887
888 async setModel(model: Model<any>): Promise<void> {
889 try {
890 const previousModel = this.model;
891 if (this.phase === "idle") {
892 await this.session.appendModelChange(model.provider, model.id);
893 } else {
894 this.pendingSessionWrites.push({ type: "model_change", provider: model.provider, modelId: model.id });
895 }
896 this.model = model;
897 await this.emitOwn({ type: "model_update", model, previousModel, source: "set" });
898 } catch (error) {
899 throw normalizeHarnessError(error, "session");
900 }
901 }
902
903 getThinkingLevel(): ThinkingLevel {
904 return this.thinkingLevel;
905 }
906
907 async setThinkingLevel(level: ThinkingLevel): Promise<void> {
908 try {
909 const previousLevel = this.thinkingLevel;
910 if (this.phase === "idle") {
911 await this.session.appendThinkingLevelChange(level);
912 } else {
913 this.pendingSessionWrites.push({ type: "thinking_level_change", thinkingLevel: level });
914 }
915 this.thinkingLevel = level;
916 await this.emitOwn({ type: "thinking_level_update", level, previousLevel });
917 } catch (error) {
918 throw normalizeHarnessError(error, "session");
919 }
920 }
921
922 getTools(): TTool[] {
923 return [...this.tools.values()];
924 }
925
926 async setTools(tools: TTool[], activeToolNames?: string[]): Promise<void> {
927 try {
928 this.validateUniqueNames(
929 tools.map((tool) => tool.name),
930 "Duplicate tool name(s)",
931 );
932 const nextTools = new Map(tools.map((tool) => [tool.name, tool]));
933 const nextActiveToolNames = activeToolNames ? [...activeToolNames] : this.activeToolNames;
934 this.validateToolNames(nextActiveToolNames, nextTools);
935 const previousToolNames = [...this.tools.keys()];
936 const previousActiveToolNames = [...this.activeToolNames];
937 if (this.phase === "idle") {
938 await this.session.appendActiveToolsChange(nextActiveToolNames);
939 } else {
940 this.pendingSessionWrites.push({ type: "active_tools_change", activeToolNames: [...nextActiveToolNames] });
941 }
942 this.tools = nextTools;
943 this.activeToolNames = [...nextActiveToolNames];
944 await this.emitOwn({
945 type: "tools_update",
946 toolNames: [...this.tools.keys()],
947 previousToolNames,
948 activeToolNames: [...this.activeToolNames],
949 previousActiveToolNames,
950 source: "set",
951 });
952 } catch (error) {
953 throw normalizeHarnessError(error, "invalid_argument");
954 }
955 }
956
957 getActiveTools(): TTool[] {
958 return this.activeToolNames.map((name) => this.tools.get(name)!);
959 }
960
961 async setActiveTools(toolNames: string[]): Promise<void> {
962 try {
963 this.validateToolNames(toolNames);
964 const previousToolNames = [...this.tools.keys()];
965 const previousActiveToolNames = [...this.activeToolNames];
966 if (this.phase === "idle") {
967 await this.session.appendActiveToolsChange(toolNames);
968 } else {
969 this.pendingSessionWrites.push({ type: "active_tools_change", activeToolNames: [...toolNames] });
970 }
971 this.activeToolNames = [...toolNames];
972 await this.emitOwn({
973 type: "tools_update",
974 toolNames: [...this.tools.keys()],
975 previousToolNames,
976 activeToolNames: [...this.activeToolNames],
977 previousActiveToolNames,
978 source: "set",
979 });
980 } catch (error) {
981 throw normalizeHarnessError(error, "invalid_argument");
982 }
983 }
984
985 getSteeringMode(): QueueMode {
986 return this.steeringQueueMode;
987 }
988
989 async setSteeringMode(mode: QueueMode): Promise<void> {
990 this.steeringQueueMode = mode;
991 }
992
993 getFollowUpMode(): QueueMode {
994 return this.followUpQueueMode;
995 }
996
997 async setFollowUpMode(mode: QueueMode): Promise<void> {
998 this.followUpQueueMode = mode;
999 }
1000
1001 getResources(): AgentHarnessResources<TSkill, TPromptTemplate> {
1002 return {
1003 skills: this.resources.skills?.slice(),
1004 promptTemplates: this.resources.promptTemplates?.slice(),
1005 };
1006 }
1007
1008 async setResources(resources: AgentHarnessResources<TSkill, TPromptTemplate>): Promise<void> {
1009 const previousResources = this.getResources();
1010 this.resources = {
1011 skills: resources.skills?.slice(),
1012 promptTemplates: resources.promptTemplates?.slice(),
1013 };
1014 await this.emitOwn({ type: "resources_update", resources: this.getResources(), previousResources });
1015 }
1016
1017 getStreamOptions(): AgentHarnessStreamOptions {
1018 return cloneStreamOptions(this.streamOptions);
1019 }
1020
1021 async setStreamOptions(streamOptions: AgentHarnessStreamOptions): Promise<void> {
1022 this.streamOptions = cloneStreamOptions(streamOptions);
1023 }
1024
1025 async abort(): Promise<AbortResult> {
1026 const clearedSteer = [...this.steerQueue];
1027 const clearedFollowUp = [...this.followUpQueue];
1028 this.steerQueue = [];
1029 this.followUpQueue = [];
1030 this.runAbortController?.abort();
1031 const errors: Error[] = [];
1032 try {
1033 await this.emitQueueUpdate();
1034 } catch (error) {
1035 errors.push(toError(error));
1036 }
1037 try {
1038 await this.waitForIdle();
1039 } catch (error) {
1040 errors.push(toError(error));
1041 }
1042 try {
1043 await this.emitOwn({ type: "abort", clearedSteer, clearedFollowUp });
1044 } catch (error) {
1045 errors.push(toError(error));
1046 }
1047 if (errors.length > 0) {
1048 const cause = errors.length === 1 ? errors[0]! : new AggregateError(errors, "Abort completed with errors");
1049 throw normalizeHarnessError(cause, "hook");
1050 }
1051 return { clearedSteer, clearedFollowUp };
1052 }
1053
1054 async waitForIdle(): Promise<void> {
1055 await this.runPromise;
1056 }
1057
1058 subscribe(
1059 listener: (event: AgentHarnessEvent<TSkill, TPromptTemplate>, signal?: AbortSignal) => Promise<void> | void,
1060 ): () => void {
1061 let handlers = this.handlers.get(SUBSCRIBER_EVENT_TYPE);
1062 if (!handlers) {
1063 handlers = new Set();
1064 this.handlers.set(SUBSCRIBER_EVENT_TYPE, handlers);
1065 }
1066 handlers.add(listener as AgentHarnessHandler);
1067 return () => handlers!.delete(listener as AgentHarnessHandler);
1068 }
1069
1070 on<TType extends keyof AgentHarnessEventResultMap>(
1071 type: TType,
1072 handler: (
1073 event: Extract<AgentHarnessOwnEvent, { type: TType }>,
1074 ) => Promise<AgentHarnessEventResultMap[TType]> | AgentHarnessEventResultMap[TType],
1075 ): () => void {
1076 let handlers = this.handlers.get(type);
1077 if (!handlers) {
1078 handlers = new Set();
1079 this.handlers.set(type, handlers);
1080 }
1081 handlers.add(handler as AgentHarnessHandler);
1082 return () => handlers!.delete(handler as AgentHarnessHandler);
1083 }
1084}
1085