refactor(agent): harden harness session semantics

This commit is contained in:
Mario Zechner
2026-05-16 00:32:16 +02:00
parent a8af0b5e99
commit 4f40f62b7b
23 changed files with 1112 additions and 873 deletions
+325 -231
View File
@@ -38,6 +38,7 @@ import type {
Session,
Skill,
} from "./types.js";
import { AgentHarnessError, BranchSummaryError, CompactionError, SessionError, toError } from "./types.js";
function createUserMessage(text: string, images?: ImageContent[]): UserMessage {
const content: Array<{ type: "text"; text: string } | ImageContent> = [{ type: "text", text }];
@@ -131,6 +132,19 @@ const SUBSCRIBER_EVENT_TYPE = "*";
type AgentHarnessHandler = (event: any, signal?: AbortSignal) => Promise<any> | any;
function normalizeHarnessError(error: unknown, fallbackCode: AgentHarnessError["code"]): AgentHarnessError {
if (error instanceof AgentHarnessError) return error;
const cause = toError(error);
if (cause instanceof SessionError) return new AgentHarnessError("session", cause.message, cause);
if (cause instanceof CompactionError) return new AgentHarnessError("compaction", cause.message, cause);
if (cause instanceof BranchSummaryError) return new AgentHarnessError("branch_summary", cause.message, cause);
return new AgentHarnessError(fallbackCode, cause.message, cause);
}
function normalizeHookError(error: unknown): AgentHarnessError {
return normalizeHarnessError(error, "hook");
}
interface AgentHarnessTurnState<
TSkill extends Skill = Skill,
TPromptTemplate extends PromptTemplate = PromptTemplate,
@@ -196,13 +210,21 @@ export class AgentHarness<
private async emitOwn(event: AgentHarnessOwnEvent<TSkill, TPromptTemplate>, signal?: AbortSignal): Promise<void> {
for (const listener of this.getHandlers(SUBSCRIBER_EVENT_TYPE) ?? []) {
await listener(event, signal);
try {
await listener(event, signal);
} catch (error) {
throw normalizeHookError(error);
}
}
}
private async emitAny(event: AgentHarnessEvent<TSkill, TPromptTemplate>, signal?: AbortSignal): Promise<void> {
for (const listener of this.getHandlers(SUBSCRIBER_EVENT_TYPE) ?? []) {
await listener(event, signal);
try {
await listener(event, signal);
} catch (error) {
throw normalizeHookError(error);
}
}
}
@@ -213,9 +235,13 @@ export class AgentHarness<
if (!handlers || handlers.size === 0) return undefined;
let lastResult: AgentHarnessEventResultMap[TType] | undefined;
for (const handler of handlers) {
const result = await handler(event);
if (result !== undefined) {
lastResult = result;
try {
const result = await handler(event);
if (result !== undefined) {
lastResult = result;
}
} catch (error) {
throw normalizeHookError(error);
}
}
return lastResult;
@@ -230,14 +256,18 @@ export class AgentHarness<
let current = cloneStreamOptions(streamOptions);
if (!handlers || handlers.size === 0) return current;
for (const handler of handlers) {
const result = await handler({
type: "before_provider_request",
model,
sessionId,
streamOptions: cloneStreamOptions(current),
});
if (result?.streamOptions) {
current = applyStreamOptionsPatch(current, result.streamOptions);
try {
const result = await handler({
type: "before_provider_request",
model,
sessionId,
streamOptions: cloneStreamOptions(current),
});
if (result?.streamOptions) {
current = applyStreamOptionsPatch(current, result.streamOptions);
}
} catch (error) {
throw normalizeHookError(error);
}
}
return current;
@@ -248,9 +278,13 @@ export class AgentHarness<
let current = payload;
if (!handlers || handlers.size === 0) return current;
for (const handler of handlers) {
const result = await handler({ type: "before_provider_payload", model, payload: current });
if (result !== undefined) {
current = result.payload;
try {
const result = await handler({ type: "before_provider_payload", model, payload: current });
if (result !== undefined) {
current = result.payload;
}
} catch (error) {
throw normalizeHookError(error);
}
}
return current;
@@ -356,8 +390,14 @@ export class AgentHarness<
private async drainQueuedMessages(queue: AgentMessage[], mode: QueueMode): Promise<AgentMessage[]> {
const messages = mode === "all" ? queue.splice(0) : queue.splice(0, 1);
if (messages.length > 0) await this.emitQueueUpdate();
return messages;
if (messages.length === 0) return messages;
try {
await this.emitQueueUpdate();
return messages;
} catch (error) {
queue.unshift(...messages);
throw normalizeHookError(error);
}
}
private createLoopConfig(
@@ -411,15 +451,14 @@ export class AgentHarness<
};
}
private validateToolNames(toolNames: string[]): void {
const missing = toolNames.filter((name) => !this.tools.has(name));
if (missing.length > 0) throw new Error(`Unknown tool(s): ${missing.join(", ")}`);
private validateToolNames(toolNames: string[], tools: Map<string, TTool> = this.tools): void {
const missing = toolNames.filter((name) => !tools.has(name));
if (missing.length > 0) throw new AgentHarnessError("invalid_argument", `Unknown tool(s): ${missing.join(", ")}`);
}
private async flushPendingSessionWrites(): Promise<void> {
const writes = this.pendingSessionWrites;
this.pendingSessionWrites = [];
for (const write of writes) {
while (this.pendingSessionWrites.length > 0) {
const write = this.pendingSessionWrites[0]!;
if (write.type === "message") {
await this.session.appendMessage(write.message);
} else if (write.type === "model_change") {
@@ -434,28 +473,40 @@ export class AgentHarness<
await this.session.appendLabel(write.targetId, write.label);
} else if (write.type === "session_info") {
await this.session.appendSessionName(write.name ?? "");
} else if (write.type === "leaf") {
await this.session.getStorage().setLeafId(write.targetId);
}
this.pendingSessionWrites.shift();
}
}
private async handleAgentEvent(event: AgentEvent, signal?: AbortSignal): Promise<void> {
await this.emitAny(event, signal);
if (event.type === "message_end") {
await this.session.appendMessage(event.message);
await this.emitAny(event, signal);
return;
}
if (event.type === "turn_end") {
let eventError: unknown;
try {
await this.emitAny(event, signal);
} catch (error) {
eventError = error;
}
const hadPendingMutations = this.pendingSessionWrites.length > 0;
await this.flushPendingSessionWrites();
await this.emitOwn({
type: "save_point",
hadPendingMutations,
});
if (eventError) throw eventError;
await this.emitOwn({ type: "save_point", hadPendingMutations });
return;
}
if (event.type === "agent_end") {
await this.flushPendingSessionWrites();
this.phase = "idle";
await this.emitAny(event, signal);
await this.emitOwn({ type: "settled", nextTurnCount: this.nextTurnQueue.length }, signal);
return;
}
await this.emitAny(event, signal);
}
private async emitRunFailure(
@@ -480,9 +531,14 @@ export class AgentHarness<
let activeTurnState = turnState;
let messages: AgentMessage[] = [createUserMessage(text, options?.images)];
if (this.nextTurnQueue.length > 0) {
messages = [...this.nextTurnQueue, messages[0]!];
this.nextTurnQueue = [];
await this.emitQueueUpdate();
const queuedMessages = this.nextTurnQueue.splice(0);
try {
await this.emitQueueUpdate();
} catch (error) {
this.nextTurnQueue.unshift(...queuedMessages);
throw normalizeHookError(error);
}
messages = [...queuedMessages, messages[0]!];
}
const beforeResult = await this.emitHook({
type: "before_agent_start",
@@ -510,12 +566,20 @@ export class AgentHarness<
this.createStreamFn(getTurnState),
);
} catch (error) {
return await this.emitRunFailure(
activeTurnState.model,
error,
abortController.signal.aborted,
abortController.signal,
);
try {
return await this.emitRunFailure(
activeTurnState.model,
error,
abortController.signal.aborted,
abortController.signal,
);
} catch (failureError) {
const cause = new AggregateError(
[toError(error), toError(failureError)],
"Agent run failed and failure reporting failed",
);
throw new AgentHarnessError("unknown", cause.message, cause);
}
}
})();
try {
@@ -526,7 +590,7 @@ export class AgentHarness<
return message;
}
}
throw new Error("AgentHarness prompt completed without an assistant message");
throw new AgentHarnessError("invalid_state", "AgentHarness prompt completed without an assistant message");
} finally {
try {
await this.flushPendingSessionWrites();
@@ -537,7 +601,7 @@ export class AgentHarness<
}
async prompt(text: string, options?: { images?: ImageContent[] }): Promise<AssistantMessage> {
if (this.phase !== "idle") throw new Error("AgentHarness is busy");
if (this.phase !== "idle") throw new AgentHarnessError("busy", "AgentHarness is busy");
this.phase = "turn";
const finishRunPromise = this.startRunPromise();
try {
@@ -545,230 +609,229 @@ export class AgentHarness<
return await this.executeTurn(turnState, text, options);
} catch (error) {
this.phase = "idle";
throw error;
throw normalizeHarnessError(error, "unknown");
} finally {
finishRunPromise();
}
}
async skill(name: string, additionalInstructions?: string): Promise<AssistantMessage> {
if (this.phase !== "idle") throw new Error("AgentHarness is busy");
if (this.phase !== "idle") throw new AgentHarnessError("busy", "AgentHarness is busy");
this.phase = "turn";
const finishRunPromise = this.startRunPromise();
try {
const turnState = await this.createTurnState();
const skill = (turnState.resources.skills ?? []).find((candidate) => candidate.name === name);
if (!skill) throw new Error(`Unknown skill: ${name}`);
if (!skill) throw new AgentHarnessError("invalid_argument", `Unknown skill: ${name}`);
return await this.executeTurn(turnState, formatSkillInvocation(skill, additionalInstructions));
} catch (error) {
this.phase = "idle";
throw error;
throw normalizeHarnessError(error, "unknown");
} finally {
finishRunPromise();
}
}
async promptFromTemplate(name: string, args: string[] = []): Promise<AssistantMessage> {
if (this.phase !== "idle") throw new Error("AgentHarness is busy");
if (this.phase !== "idle") throw new AgentHarnessError("busy", "AgentHarness is busy");
this.phase = "turn";
const finishRunPromise = this.startRunPromise();
try {
const turnState = await this.createTurnState();
const template = (turnState.resources.promptTemplates ?? []).find((candidate) => candidate.name === name);
if (!template) throw new Error(`Unknown prompt template: ${name}`);
if (!template) throw new AgentHarnessError("invalid_argument", `Unknown prompt template: ${name}`);
return await this.executeTurn(turnState, formatPromptTemplateInvocation(template, args));
} catch (error) {
this.phase = "idle";
throw error;
throw normalizeHarnessError(error, "unknown");
} finally {
finishRunPromise();
}
}
steer(text: string, options?: { images?: ImageContent[] }): void {
if (this.phase === "idle") throw new Error("Cannot steer while idle");
async steer(text: string, options?: { images?: ImageContent[] }): Promise<void> {
if (this.phase === "idle") throw new AgentHarnessError("invalid_state", "Cannot steer while idle");
this.steerQueue.push(createUserMessage(text, options?.images));
void this.emitQueueUpdate();
await this.emitQueueUpdate();
}
followUp(text: string, options?: { images?: ImageContent[] }): void {
if (this.phase === "idle") throw new Error("Cannot follow up while idle");
async followUp(text: string, options?: { images?: ImageContent[] }): Promise<void> {
if (this.phase === "idle") throw new AgentHarnessError("invalid_state", "Cannot follow up while idle");
this.followUpQueue.push(createUserMessage(text, options?.images));
void this.emitQueueUpdate();
await this.emitQueueUpdate();
}
nextTurn(text: string, options?: { images?: ImageContent[] }): void {
async nextTurn(text: string, options?: { images?: ImageContent[] }): Promise<void> {
this.nextTurnQueue.push(createUserMessage(text, options?.images));
void this.emitQueueUpdate();
await this.emitQueueUpdate();
}
async appendMessage(message: AgentMessage): Promise<void> {
if (this.phase === "idle") {
await this.session.appendMessage(message);
} else {
this.pendingSessionWrites.push({ type: "message", message });
try {
if (this.phase === "idle") {
await this.session.appendMessage(message);
} else {
this.pendingSessionWrites.push({ type: "message", message });
}
} catch (error) {
throw normalizeHarnessError(error, "session");
}
}
async compact(
customInstructions?: string,
): Promise<{ summary: string; firstKeptEntryId: string; tokensBefore: number; details?: unknown }> {
if (this.phase !== "idle") throw new Error("compact() requires idle harness");
if (this.phase !== "idle") throw new AgentHarnessError("busy", "compact() requires idle harness");
this.phase = "compaction";
const model = this.model;
if (!model) throw new Error("No model set for compaction");
const auth = await this.getApiKeyAndHeaders?.(model);
if (!auth) throw new Error("No auth available for compaction");
const branchEntries = await this.session.getBranch();
const preparation = prepareCompaction(branchEntries, DEFAULT_COMPACTION_SETTINGS);
if (!preparation) throw new Error("Nothing to compact");
const hookResult = await this.emitHook({
type: "session_before_compact",
preparation,
branchEntries,
customInstructions,
signal: new AbortController().signal,
});
if (hookResult?.cancel) {
try {
const model = this.model;
if (!model) throw new AgentHarnessError("invalid_state", "No model set for compaction");
const auth = await this.getApiKeyAndHeaders?.(model);
if (!auth) throw new AgentHarnessError("auth", "No auth available for compaction");
const branchEntries = await this.session.getBranch();
const preparationResult = prepareCompaction(branchEntries, DEFAULT_COMPACTION_SETTINGS);
if (!preparationResult.ok) throw preparationResult.error;
const preparation = preparationResult.value;
if (!preparation) throw new AgentHarnessError("compaction", "Nothing to compact");
const hookResult = await this.emitHook({
type: "session_before_compact",
preparation,
branchEntries,
customInstructions,
signal: new AbortController().signal,
});
if (hookResult?.cancel) throw new AgentHarnessError("compaction", "Compaction cancelled");
const provided = hookResult?.compaction;
const compactResult = provided
? { ok: true as const, value: provided }
: await compact(
preparation,
model,
auth.apiKey,
auth.headers,
customInstructions,
undefined,
this.thinkingLevel,
);
if (!compactResult.ok) throw compactResult.error;
const result = compactResult.value;
const entryId = await this.session.appendCompaction(
result.summary,
result.firstKeptEntryId,
result.tokensBefore,
result.details,
provided !== undefined,
);
const entry = await this.session.getEntry(entryId);
if (entry?.type === "compaction") {
await this.emitOwn({ type: "session_compact", compactionEntry: entry, fromHook: provided !== undefined });
}
return result;
} catch (error) {
throw normalizeHarnessError(error, "compaction");
} finally {
this.phase = "idle";
throw new Error("Compaction cancelled");
}
const provided = hookResult?.compaction;
const compactResult = provided
? { ok: true as const, value: provided }
: await compact(
preparation,
model,
auth.apiKey,
auth.headers,
customInstructions,
undefined,
this.thinkingLevel,
);
if (!compactResult.ok) throw compactResult.error;
const result = compactResult.value;
const entryId = await this.session.appendCompaction(
result.summary,
result.firstKeptEntryId,
result.tokensBefore,
result.details,
provided !== undefined,
);
const entry = await this.session.getEntry(entryId);
if (entry?.type === "compaction") {
await this.emitOwn({ type: "session_compact", compactionEntry: entry, fromHook: provided !== undefined });
}
this.phase = "idle";
return result;
}
async navigateTree(
targetId: string,
options?: { summarize?: boolean; customInstructions?: string; replaceInstructions?: boolean; label?: string },
): Promise<NavigateTreeResult> {
if (this.phase !== "idle") throw new Error("navigateTree() requires idle harness");
if (this.phase !== "idle") throw new AgentHarnessError("busy", "navigateTree() requires idle harness");
this.phase = "branch_summary";
const oldLeafId = await this.session.getLeafId();
if (oldLeafId === targetId) {
this.phase = "idle";
return { cancelled: false };
}
const targetEntry = await this.session.getEntry(targetId);
if (!targetEntry) throw new Error(`Entry ${targetId} not found`);
const { entries, commonAncestorId } = await collectEntriesForBranchSummary(this.session, oldLeafId, targetId);
const preparation = {
targetId,
oldLeafId,
commonAncestorId,
entriesToSummarize: entries,
userWantsSummary: options?.summarize ?? false,
customInstructions: options?.customInstructions,
replaceInstructions: options?.replaceInstructions,
label: options?.label,
};
const signal = new AbortController().signal;
const hookResult = await this.emitHook({
type: "session_before_tree",
preparation,
signal,
});
if (hookResult?.cancel) {
this.phase = "idle";
return { cancelled: true };
}
let summaryEntry: any | undefined;
let summaryText: string | undefined = hookResult?.summary?.summary;
let summaryDetails: unknown = hookResult?.summary?.details;
if (!summaryText && options?.summarize && entries.length > 0) {
const model = this.model;
if (!model) throw new Error("No model set for branch summary");
const auth = await this.getApiKeyAndHeaders?.(model);
if (!auth) throw new Error("No auth available for branch summary");
const branchSummary = await generateBranchSummary(entries, {
model,
apiKey: auth.apiKey,
headers: auth.headers,
signal: new AbortController().signal,
customInstructions: hookResult?.customInstructions ?? options?.customInstructions,
replaceInstructions: hookResult?.replaceInstructions ?? options?.replaceInstructions,
});
if (branchSummary.aborted) {
this.phase = "idle";
return { cancelled: true };
}
if (branchSummary.error) throw new Error(branchSummary.error);
summaryText = branchSummary.summary;
summaryDetails = {
readFiles: branchSummary.readFiles ?? [],
modifiedFiles: branchSummary.modifiedFiles ?? [],
try {
const oldLeafId = await this.session.getLeafId();
if (oldLeafId === targetId) return { cancelled: false };
const targetEntry = await this.session.getEntry(targetId);
if (!targetEntry) throw new AgentHarnessError("invalid_argument", `Entry ${targetId} not found`);
const { entries, commonAncestorId } = await collectEntriesForBranchSummary(this.session, oldLeafId, targetId);
const preparation = {
targetId,
oldLeafId,
commonAncestorId,
entriesToSummarize: entries,
userWantsSummary: options?.summarize ?? false,
customInstructions: options?.customInstructions,
replaceInstructions: options?.replaceInstructions,
label: options?.label,
};
const signal = new AbortController().signal;
const hookResult = await this.emitHook({ type: "session_before_tree", preparation, signal });
if (hookResult?.cancel) return { cancelled: true };
let summaryEntry: NavigateTreeResult["summaryEntry"];
let summaryText: string | undefined = hookResult?.summary?.summary;
let summaryDetails: unknown = hookResult?.summary?.details;
if (!summaryText && options?.summarize && entries.length > 0) {
const model = this.model;
if (!model) throw new AgentHarnessError("invalid_state", "No model set for branch summary");
const auth = await this.getApiKeyAndHeaders?.(model);
if (!auth) throw new AgentHarnessError("auth", "No auth available for branch summary");
const branchSummary = await generateBranchSummary(entries, {
model,
apiKey: auth.apiKey,
headers: auth.headers,
signal: new AbortController().signal,
customInstructions: hookResult?.customInstructions ?? options?.customInstructions,
replaceInstructions: hookResult?.replaceInstructions ?? options?.replaceInstructions,
});
if (!branchSummary.ok) {
if (branchSummary.error.code === "aborted") return { cancelled: true };
throw new AgentHarnessError("branch_summary", branchSummary.error.message, branchSummary.error);
}
summaryText = branchSummary.value.summary;
summaryDetails = {
readFiles: branchSummary.value.readFiles,
modifiedFiles: branchSummary.value.modifiedFiles,
};
}
let editorText: string | undefined;
let newLeafId: string | null;
if (targetEntry.type === "message" && targetEntry.message.role === "user") {
newLeafId = targetEntry.parentId;
const content = targetEntry.message.content;
editorText =
typeof content === "string"
? content
: content
.filter((c): c is { readonly type: "text"; readonly text: string } => c.type === "text")
.map((c) => c.text)
.join("");
} else if (targetEntry.type === "custom_message") {
newLeafId = targetEntry.parentId;
editorText =
typeof targetEntry.content === "string"
? targetEntry.content
: targetEntry.content
.filter((c): c is { readonly type: "text"; readonly text: string } => c.type === "text")
.map((c) => c.text)
.join("");
} else {
newLeafId = targetId;
}
const summaryId = await this.session.moveTo(
newLeafId,
summaryText
? { summary: summaryText, details: summaryDetails, fromHook: hookResult?.summary !== undefined }
: undefined,
);
if (summaryId) {
const entry = await this.session.getEntry(summaryId);
if (entry?.type === "branch_summary") summaryEntry = entry;
}
await this.emitOwn({
type: "session_tree",
newLeafId: await this.session.getLeafId(),
oldLeafId,
summaryEntry,
fromHook: hookResult?.summary !== undefined,
});
return { cancelled: false, editorText, summaryEntry };
} catch (error) {
throw normalizeHarnessError(error, "branch_summary");
} finally {
this.phase = "idle";
}
let editorText: string | undefined;
let newLeafId: string | null;
if (targetEntry.type === "message" && targetEntry.message.role === "user") {
newLeafId = targetEntry.parentId;
const content = targetEntry.message.content;
editorText =
typeof content === "string"
? content
: content
.filter((c): c is { readonly type: "text"; readonly text: string } => c.type === "text")
.map((c) => c.text)
.join("");
} else if (targetEntry.type === "custom_message") {
newLeafId = targetEntry.parentId;
editorText =
typeof targetEntry.content === "string"
? targetEntry.content
: targetEntry.content
.filter((c): c is { readonly type: "text"; readonly text: string } => c.type === "text")
.map((c) => c.text)
.join("");
} else {
newLeafId = targetId;
}
const summaryId = await this.session.moveTo(
newLeafId,
summaryText
? {
summary: summaryText,
details: summaryDetails,
fromHook: hookResult?.summary !== undefined,
}
: undefined,
);
if (summaryId) {
summaryEntry = await this.session.getEntry(summaryId);
}
await this.emitOwn({
type: "session_tree",
newLeafId: await this.session.getLeafId(),
oldLeafId,
summaryEntry,
fromHook: hookResult?.summary !== undefined,
});
this.phase = "idle";
return { cancelled: false, editorText, summaryEntry };
}
getModel(): Model<any> {
@@ -780,37 +843,49 @@ export class AgentHarness<
}
async setModel(model: Model<any>): Promise<void> {
const previousModel = this.model;
this.model = model;
if (this.phase === "idle") {
await this.session.appendModelChange(model.provider, model.id);
} else {
this.pendingSessionWrites.push({ type: "model_change", provider: model.provider, modelId: model.id });
try {
const previousModel = this.model;
if (this.phase === "idle") {
await this.session.appendModelChange(model.provider, model.id);
} else {
this.pendingSessionWrites.push({ type: "model_change", provider: model.provider, modelId: model.id });
}
this.model = model;
await this.emitOwn({ type: "model_select", model, previousModel, source: "set" });
} catch (error) {
throw normalizeHarnessError(error, "session");
}
await this.emitOwn({ type: "model_select", model, previousModel, source: "set" });
}
async setThinkingLevel(level: ThinkingLevel): Promise<void> {
const previousLevel = this.thinkingLevel;
this.thinkingLevel = level;
if (this.phase === "idle") {
await this.session.appendThinkingLevelChange(level);
} else {
this.pendingSessionWrites.push({ type: "thinking_level_change", thinkingLevel: level });
try {
const previousLevel = this.thinkingLevel;
if (this.phase === "idle") {
await this.session.appendThinkingLevelChange(level);
} else {
this.pendingSessionWrites.push({ type: "thinking_level_change", thinkingLevel: level });
}
this.thinkingLevel = level;
await this.emitOwn({ type: "thinking_level_select", level, previousLevel });
} catch (error) {
throw normalizeHarnessError(error, "session");
}
await this.emitOwn({ type: "thinking_level_select", level, previousLevel });
}
async setActiveTools(toolNames: string[]): Promise<void> {
this.validateToolNames(toolNames);
this.activeToolNames = [...toolNames];
try {
this.validateToolNames(toolNames);
this.activeToolNames = [...toolNames];
} catch (error) {
throw normalizeHarnessError(error, "invalid_argument");
}
}
getSteeringMode(): QueueMode {
return this.steeringQueueMode;
}
setSteeringMode(mode: QueueMode): void {
async setSteeringMode(mode: QueueMode): Promise<void> {
this.steeringQueueMode = mode;
}
@@ -818,7 +893,7 @@ export class AgentHarness<
return this.followUpQueueMode;
}
setFollowUpMode(mode: QueueMode): void {
async setFollowUpMode(mode: QueueMode): Promise<void> {
this.followUpQueueMode = mode;
}
@@ -842,17 +917,19 @@ export class AgentHarness<
return cloneStreamOptions(this.streamOptions);
}
setStreamOptions(streamOptions: AgentHarnessStreamOptions): void {
async setStreamOptions(streamOptions: AgentHarnessStreamOptions): Promise<void> {
this.streamOptions = cloneStreamOptions(streamOptions);
}
async setTools(tools: TTool[], activeToolNames?: string[]): Promise<void> {
this.tools = new Map(tools.map((tool) => [tool.name, tool]));
if (activeToolNames) {
this.validateToolNames(activeToolNames);
this.activeToolNames = [...activeToolNames];
} else {
this.validateToolNames(this.activeToolNames);
try {
const nextTools = new Map(tools.map((tool) => [tool.name, tool]));
const nextActiveToolNames = activeToolNames ? [...activeToolNames] : this.activeToolNames;
this.validateToolNames(nextActiveToolNames, nextTools);
this.tools = nextTools;
this.activeToolNames = [...nextActiveToolNames];
} catch (error) {
throw normalizeHarnessError(error, "invalid_argument");
}
}
@@ -861,10 +938,27 @@ export class AgentHarness<
const clearedFollowUp = [...this.followUpQueue];
this.steerQueue = [];
this.followUpQueue = [];
await this.emitQueueUpdate();
this.runAbortController?.abort();
await this.waitForIdle();
await this.emitOwn({ type: "abort", clearedSteer, clearedFollowUp });
const errors: Error[] = [];
try {
await this.emitQueueUpdate();
} catch (error) {
errors.push(toError(error));
}
try {
await this.waitForIdle();
} catch (error) {
errors.push(toError(error));
}
try {
await this.emitOwn({ type: "abort", clearedSteer, clearedFollowUp });
} catch (error) {
errors.push(toError(error));
}
if (errors.length > 0) {
const cause = errors.length === 1 ? errors[0]! : new AggregateError(errors, "Abort completed with errors");
throw normalizeHarnessError(cause, "hook");
}
return { clearedSteer, clearedFollowUp };
}