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";4849function 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}5455function 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}7576function cloneStreamOptions(streamOptions?: AgentHarnessStreamOptions): AgentHarnessStreamOptions {77 return {78 ...streamOptions,79 headers: streamOptions?.headers ? { ...streamOptions.headers } : undefined,80 metadata: streamOptions?.metadata ? { ...streamOptions.metadata } : undefined,81 };82}8384function 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}9394function applyStreamOptionsPatch(95 base: AgentHarnessStreamOptions,96 patch?: AgentHarnessStreamOptionsPatch,97): AgentHarnessStreamOptions {98 const result = cloneStreamOptions(base);99 if (!patch) return result;100101 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;106107 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 }119120 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 }132133 return result;134}135136const SUBSCRIBER_EVENT_TYPE = "*";137138type AgentHarnessHandler = (event: any, signal?: AbortSignal) => Promise<any> | any;139140function 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}148149function normalizeHookError(error: unknown): AgentHarnessError {150 return normalizeHarnessError(error, "hook");151}152153interface 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}170171export 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>>();198199 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.activeToolNames217 ? [...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 }224225 private getHandlers(type: string): Set<AgentHarnessHandler> | undefined {226 return this.handlers.get(type);227 }228229 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 }238239 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 }248249 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 }267268 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 }276277 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 }302303 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 }319320 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 }328329 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 }339340 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 }346347 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 }353354 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.activeToolNames361 .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 }388389 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 }399400 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 }429430 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 }441442 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 patch476 ? {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 }499500 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 }505506 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 }511512 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 }537538 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 }566567 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 }580581 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];606607 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 }657658 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 }672673 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 }689690 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 }706707 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 }712713 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 }718719 async nextTurn(text: string, options?: { images?: ImageContent[] }): Promise<void> {720 this.nextTurnQueue.push(createUserMessage(text, options?.images));721 await this.emitQueueUpdate();722 }723724 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 }735736 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 = provided757 ? { 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 }790791 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 summaryText857 ? {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 }883884 getModel(): Model<any> {885 return this.model;886 }887888 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 }902903 getThinkingLevel(): ThinkingLevel {904 return this.thinkingLevel;905 }906907 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 }921922 getTools(): TTool[] {923 return [...this.tools.values()];924 }925926 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 }956957 getActiveTools(): TTool[] {958 return this.activeToolNames.map((name) => this.tools.get(name)!);959 }960961 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 }984985 getSteeringMode(): QueueMode {986 return this.steeringQueueMode;987 }988989 async setSteeringMode(mode: QueueMode): Promise<void> {990 this.steeringQueueMode = mode;991 }992993 getFollowUpMode(): QueueMode {994 return this.followUpQueueMode;995 }996997 async setFollowUpMode(mode: QueueMode): Promise<void> {998 this.followUpQueueMode = mode;999 }10001001 getResources(): AgentHarnessResources<TSkill, TPromptTemplate> {1002 return {1003 skills: this.resources.skills?.slice(),1004 promptTemplates: this.resources.promptTemplates?.slice(),1005 };1006 }10071008 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 }10161017 getStreamOptions(): AgentHarnessStreamOptions {1018 return cloneStreamOptions(this.streamOptions);1019 }10201021 async setStreamOptions(streamOptions: AgentHarnessStreamOptions): Promise<void> {1022 this.streamOptions = cloneStreamOptions(streamOptions);1023 }10241025 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 }10531054 async waitForIdle(): Promise<void> {1055 await this.runPromise;1056 }10571058 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 }10691070 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