/** * KHAELOR * File: tests/agent/fixtures.ts * Description: Agent-kernel test fixtures — scripted ModelClient (no network), full real-service harness on a temp dir. * * Author: Simon-Pierre Boucher * Contact: contact@spboucher.ai */ import { mkdtempSync } from "node:fs"; import { tmpdir } from "node:os"; import * as path from "node:path"; import { expect } from "vitest"; import { AgentKernel, EventLogSession, InterruptionController, SteeringQueue, ToolExecutor, VerificationGate, } from "../../src/agent/index.js"; import { ModelError } from "../../src/anthropic/index.js"; import type { ModelClient, ModelEvent, ModelRequest } from "../../src/anthropic/index.js"; import { ContextBudget, KhaelorContextEngine } from "../../src/context/index.js"; import { PermissionService } from "../../src/permissions/index.js"; import type { PermissionRule } from "../../src/permissions/index.js"; import { SessionEventBus, SessionLog, buildConversation, findDanglingToolUseIds } from "../../src/session/index.js"; import type { DurableEvent } from "../../src/session/index.js"; import { createDefaultToolRegistry } from "../../src/tools/index.js"; import { InMemoryFileTimeRegistry, LocalProcessManager, LocalWorkspace, } from "../../src/workspace/index.js"; // ───────────────────────── scripted model client ───────────────────────── export interface ScriptedTurn { /** Called after the request is captured, before any event is yielded. */ onStart?: () => void; events: ModelEvent[]; /** After yielding events, wait for the abort signal, then throw cancelled. */ hangUntilAbort?: boolean; /** Throw this after yielding events. */ error?: Error; } /** A ModelClient that replays scripted ModelEvent sequences — no network, honest signals. */ export class ScriptedModelClient implements ModelClient { readonly requests: ModelRequest[] = []; readonly #turns: ScriptedTurn[]; #index = 0; constructor(turns: ScriptedTurn[]) { this.#turns = turns; } async *stream(request: ModelRequest, signal?: AbortSignal): AsyncIterable { const turn = this.#turns[this.#index]; this.#index += 1; if (turn === undefined) { throw new ModelError("invalid-request", "ScriptedModelClient: no more scripted turns."); } this.requests.push(request); turn.onStart?.(); for (const event of turn.events) { if (signal?.aborted === true) throw new ModelError("cancelled", "Request cancelled."); await Promise.resolve(); yield event; } if (turn.hangUntilAbort === true) { await new Promise((resolve) => { if (signal?.aborted === true) { resolve(); return; } signal?.addEventListener("abort", () => resolve(), { once: true }); }); throw new ModelError("cancelled", "Request cancelled."); } if (turn.error !== undefined) throw turn.error; } } /** Standard scripted blocks. */ export const USAGE = { inputTokens: 100, outputTokens: 20, cacheReadTokens: 0, cacheWriteTokens: 0 }; export function textTurn(requestId: string, text: string): ScriptedTurn { return { events: [ { type: "started", requestId }, { type: "text-block-completed", blockIndex: 0, text }, { type: "completed", stopReason: "end_turn", usage: { ...USAGE }, durationMs: 5 }, ], }; } export function toolTurn( requestId: string, text: string, calls: { toolUseId: string; toolName: string; input: unknown }[], ): ScriptedTurn { const events: ModelEvent[] = [ { type: "started", requestId }, { type: "text-block-completed", blockIndex: 0, text }, ]; calls.forEach((call, i) => { events.push({ type: "tool-call-completed", blockIndex: i + 1, toolUseId: call.toolUseId, toolName: call.toolName, input: call.input, }); }); events.push({ type: "completed", stopReason: "tool_use", usage: { ...USAGE }, durationMs: 5 }); return { events }; } // ───────────────────────── harness ───────────────────────── export const ALLOW_ALL_RULES: PermissionRule[] = [ { capability: "*", pattern: "*", action: "allow", source: "user" }, ]; export interface AgentHarness { dir: string; sessionsDir: string; workspace: LocalWorkspace; log: SessionLog; session: EventLogSession; client: ScriptedModelClient; kernel: AgentKernel; steering: SteeringQueue; interruption: InterruptionController; verifier: VerificationGate; events(): readonly DurableEvent[]; user(text: string): void; } export interface HarnessOptions { turns: ScriptedTurn[]; /** User-layer permission rules. Defaults to allow-everything. Pass [] for shipped defaults only. */ rules?: PermissionRule[]; /** Reuse an existing workspace dir (resume tests). */ dir?: string; /** Resume an existing session log instead of creating one. */ resume?: { sessionsDir: string; sessionId: string }; /** Kernel model-call budget override. */ maxIterations?: number; } /** Build a kernel wired to REAL services (workspace, tools, permissions, context) in a temp dir. */ export async function makeHarness(options: HarnessOptions): Promise { const dir = options.dir ?? mkdtempSync(path.join(tmpdir(), "khaelor-agent-")); const sessionsDir = options.resume?.sessionsDir ?? path.join(dir, ".khaelor-sessions"); const log = options.resume !== undefined ? await SessionLog.open({ projectHash: "agenttest", sessionsDir, sessionId: options.resume.sessionId, }) : await SessionLog.create({ projectHash: "agenttest", sessionsDir }); const bus = new SessionEventBus({ sessionId: log.sessionId, appender: log }); const session = new EventLogSession({ sessionId: log.sessionId, bus, replayed: log.replayedEvents, }); const workspace = new LocalWorkspace(dir); const fileTimes = new InMemoryFileTimeRegistry(); const processes = new LocalProcessManager({ logDir: path.join(dir, ".khaelor-process-logs"), shell: "/bin/sh", }); const registry = createDefaultToolRegistry(); const permissions = new PermissionService({ rules: { user: options.rules ?? ALLOW_ALL_RULES }, publish: (event) => { session.publishDurable(event); }, }); const executor = new ToolExecutor({ session, registry, permissions, workspace, fileTimes, processes, spillDir: path.join(dir, ".khaelor-spill"), }); const client = new ScriptedModelClient(options.turns); const context = new KhaelorContextEngine({ model: "test-model", auxModel: "test-aux", maxOutputTokens: 4096, systemTiers: [{ name: "identity", text: "You are KHAELOR, a terminal engineering agent." }], tools: registry.list().map((tool) => ({ name: tool.name, description: tool.description, inputSchema: tool.inputSchema, })), modelClient: client, budget: new ContextBudget({ model: "test-model", reservedOutputTokens: 4096 }), }); const verifier = new VerificationGate({ workspace }); const steering = new SteeringQueue(session); const interruption = new InterruptionController(session); const kernel = new AgentKernel({ session, context, model: client, executor, verifier, steering, interruption, ...(options.maxIterations !== undefined ? { maxIterations: options.maxIterations } : {}), }); return { dir, sessionsDir, workspace, log, session, client, kernel, steering, interruption, verifier, events: () => session.events(), user: (text: string) => { session.publishDurable({ type: "user.message-created", payload: { text, mentions: [] } }); }, }; } // ───────────────────────── assertions and helpers ───────────────────────── /** Durable event type strings, in seq order. */ export function durableTypes(events: readonly DurableEvent[]): string[] { return events.map((event) => event.type); } /** * EVENT_MODEL invariants over a recorded log: * - every tool_use has exactly one terminal tool_result (§6.5); * - per toolUseId, ToolApproved/Started precede the terminal event; * - the rebuilt conversation pairs every tool_use with a tool_result, * tool_use in assistant messages, tool_result in user messages. */ export function assertPairingSafe(events: readonly DurableEvent[]): void { expect(findDanglingToolUseIds(events)).toEqual([]); const requestedAt = new Map(); const startedAt = new Map(); const terminalAt = new Map(); for (const event of events) { switch (event.type) { case "tool.requested": requestedAt.set(event.payload.toolUseId, event.seq); break; case "tool.started": startedAt.set(event.payload.toolUseId, event.seq); break; case "tool.completed": case "tool.failed": case "tool.cancelled": expect(terminalAt.has(event.payload.toolUseId)).toBe(false); // exactly one terminal terminalAt.set(event.payload.toolUseId, event.seq); break; default: break; } } for (const [toolUseId, seq] of terminalAt) { const requested = requestedAt.get(toolUseId); expect(requested).toBeDefined(); expect(requested as number).toBeLessThan(seq); const started = startedAt.get(toolUseId); if (started !== undefined) expect(started).toBeLessThan(seq); } const messages = buildConversation(events); if (messages.length > 0) expect(messages[0]?.role).toBe("user"); const useIds = new Set(); const resultIds = new Set(); for (const message of messages) { for (const block of message.content) { if (block.type === "tool_use") { expect(message.role).toBe("assistant"); useIds.add(block.id); } else if (block.type === "tool_result") { expect(message.role).toBe("user"); expect(useIds.has(block.tool_use_id)).toBe(true); // result after its use resultIds.add(block.tool_use_id); } } } expect([...useIds].sort()).toEqual([...resultIds].sort()); } /** Poll until `predicate` holds (macro-task cadence) — for interrupt timing. */ export async function waitFor(predicate: () => boolean, timeoutMs = 2000): Promise { const deadline = Date.now() + timeoutMs; while (!predicate()) { if (Date.now() > deadline) throw new Error("waitFor: condition not met in time"); await new Promise((resolve) => setTimeout(resolve, 5)); } }