From 81c9227198c0f7bbaf87f01590bffe428a440881 Mon Sep 17 00:00:00 2001 From: dimakis Date: Sat, 4 Jul 2026 10:32:56 +0100 Subject: [PATCH 1/4] feat(providers): add ModelSession interface and AnthropicSession adapter MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Provider-agnostic abstraction for the agentic loop. ModelSession.turn() replaces the Agent SDK's query() — Mitzo owns the loop, ModelSession handles a single LLM call with streaming events. AnthropicSession speaks Anthropic Messages API format, pointed at praxis-proxy (127.0.0.1:9090) which translates to any backend provider. Co-Authored-By: Claude Opus 4.6 --- .../__tests__/anthropic-session.test.ts | 459 ++++++++++++++++++ packages/harness/src/index.ts | 13 + .../src/providers/anthropic-session.ts | 209 ++++++++ packages/harness/src/providers/index.ts | 18 + .../harness/src/providers/session-types.ts | 143 ++++++ 5 files changed, 842 insertions(+) create mode 100644 packages/harness/__tests__/anthropic-session.test.ts create mode 100644 packages/harness/src/providers/anthropic-session.ts create mode 100644 packages/harness/src/providers/session-types.ts diff --git a/packages/harness/__tests__/anthropic-session.test.ts b/packages/harness/__tests__/anthropic-session.test.ts new file mode 100644 index 00000000..5fb53813 --- /dev/null +++ b/packages/harness/__tests__/anthropic-session.test.ts @@ -0,0 +1,459 @@ +import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest'; +import { AnthropicSession } from '../src/providers/anthropic-session.js'; +import type { + ModelSessionConfig, + StreamEvent, + ConversationMessage, +} from '../src/providers/session-types.js'; + +/** Build a minimal session config for testing. */ +function testConfig(overrides: Partial = {}): ModelSessionConfig { + return { + model: 'claude-haiku-4-5', + systemPrompt: 'You are a test assistant.', + maxTokens: 100, + ...overrides, + }; +} + +/** Encode text as a ReadableStream chunk. */ +function encodeChunk(text: string): Uint8Array { + return new TextEncoder().encode(text); +} + +/** Build a mock SSE response body from events. */ +function mockSSEStream(events: Array<{ event: string; data: string }>): ReadableStream { + const chunks = events.map((e) => `event: ${e.event}\ndata: ${e.data}\n\n`).join(''); + return new ReadableStream({ + start(controller) { + controller.enqueue(encodeChunk(chunks)); + controller.close(); + }, + }); +} + +/** Collect all events from an async iterable. */ +async function collectEvents(iter: AsyncIterable): Promise { + const events: StreamEvent[] = []; + for await (const event of iter) { + events.push(event); + } + return events; +} + +describe('AnthropicSession', () => { + let fetchSpy: ReturnType; + + beforeEach(() => { + fetchSpy = vi.spyOn(globalThis, 'fetch'); + }); + + afterEach(() => { + fetchSpy.mockRestore(); + }); + + describe('constructor', () => { + it('creates session with default options', () => { + const session = new AnthropicSession(testConfig()); + expect(session.provider).toBe('anthropic'); + }); + + it('uses praxis proxy when MITZO_USE_PRAXIS is set', () => { + vi.stubEnv('MITZO_USE_PRAXIS', '1'); + const session = new AnthropicSession(testConfig()); + expect(session.provider).toBe('anthropic'); + vi.unstubAllEnvs(); + }); + + it('uses explicit baseUrl over env var', () => { + const session = new AnthropicSession(testConfig(), { + baseUrl: 'http://custom:8080', + }); + expect(session.provider).toBe('anthropic'); + }); + }); + + describe('turn()', () => { + it('streams message_start event', async () => { + const sseBody = mockSSEStream([ + { + event: 'message_start', + data: JSON.stringify({ + type: 'message_start', + message: { + id: 'msg_123', + model: 'claude-haiku-4-5', + role: 'assistant', + usage: { input_tokens: 10, output_tokens: 0 }, + }, + }), + }, + { + event: 'message_delta', + data: JSON.stringify({ + type: 'message_delta', + delta: { stop_reason: 'end_turn', stop_sequence: null }, + usage: { output_tokens: 5 }, + }), + }, + ]); + + fetchSpy.mockResolvedValueOnce( + new Response(sseBody, { + status: 200, + headers: { 'Content-Type': 'text/event-stream' }, + }), + ); + + const session = new AnthropicSession(testConfig(), { + baseUrl: 'http://test:9090', + apiKey: 'test-key', + }); + + const messages: ConversationMessage[] = [{ role: 'user', content: 'hello' }]; + const events = await collectEvents(session.turn(messages)); + + expect(events).toHaveLength(2); + expect(events[0].type).toBe('message_start'); + expect(events[1].type).toBe('message_delta'); + }); + + it('streams content blocks with text deltas', async () => { + const sseBody = mockSSEStream([ + { + event: 'message_start', + data: JSON.stringify({ + type: 'message_start', + message: { + id: 'msg_456', + model: 'claude-haiku-4-5', + role: 'assistant', + usage: { input_tokens: 10, output_tokens: 0 }, + }, + }), + }, + { + event: 'content_block_start', + data: JSON.stringify({ + type: 'content_block_start', + index: 0, + content_block: { type: 'text', text: '' }, + }), + }, + { + event: 'content_block_delta', + data: JSON.stringify({ + type: 'content_block_delta', + index: 0, + delta: { type: 'text_delta', text: 'Hello!' }, + }), + }, + { + event: 'content_block_stop', + data: JSON.stringify({ + type: 'content_block_stop', + index: 0, + }), + }, + { + event: 'message_delta', + data: JSON.stringify({ + type: 'message_delta', + delta: { stop_reason: 'end_turn', stop_sequence: null }, + usage: { output_tokens: 3 }, + }), + }, + ]); + + fetchSpy.mockResolvedValueOnce( + new Response(sseBody, { + status: 200, + headers: { 'Content-Type': 'text/event-stream' }, + }), + ); + + const session = new AnthropicSession(testConfig(), { + baseUrl: 'http://test:9090', + apiKey: 'test-key', + }); + + const events = await collectEvents(session.turn([{ role: 'user', content: 'hi' }])); + + expect(events).toHaveLength(5); + expect(events[0].type).toBe('message_start'); + expect(events[1].type).toBe('content_block_start'); + expect(events[2].type).toBe('content_block_delta'); + if (events[2].type === 'content_block_delta' && events[2].delta.type === 'text_delta') { + expect(events[2].delta.text).toBe('Hello!'); + } + expect(events[3].type).toBe('content_block_stop'); + expect(events[4].type).toBe('message_delta'); + }); + + it('streams tool_use blocks', async () => { + const sseBody = mockSSEStream([ + { + event: 'message_start', + data: JSON.stringify({ + type: 'message_start', + message: { + id: 'msg_789', + model: 'gpt-5.5', + role: 'assistant', + usage: { input_tokens: 20, output_tokens: 0 }, + }, + }), + }, + { + event: 'content_block_start', + data: JSON.stringify({ + type: 'content_block_start', + index: 0, + content_block: { + type: 'tool_use', + id: 'toolu_abc', + name: 'Read', + input: {}, + }, + }), + }, + { + event: 'content_block_delta', + data: JSON.stringify({ + type: 'content_block_delta', + index: 0, + delta: { type: 'input_json_delta', partial_json: '{"file_path":"/tmp/test"}' }, + }), + }, + { + event: 'content_block_stop', + data: JSON.stringify({ type: 'content_block_stop', index: 0 }), + }, + { + event: 'message_delta', + data: JSON.stringify({ + type: 'message_delta', + delta: { stop_reason: 'tool_use', stop_sequence: null }, + usage: { output_tokens: 15 }, + }), + }, + ]); + + fetchSpy.mockResolvedValueOnce( + new Response(sseBody, { + status: 200, + headers: { 'Content-Type': 'text/event-stream' }, + }), + ); + + const session = new AnthropicSession( + testConfig({ + model: 'gpt-5.5', + tools: [ + { + name: 'Read', + description: 'Read a file', + input_schema: { type: 'object', properties: { file_path: { type: 'string' } } }, + }, + ], + }), + { baseUrl: 'http://test:9090', apiKey: 'test-key' }, + ); + + const events = await collectEvents(session.turn([{ role: 'user', content: 'read /tmp/test' }])); + + expect(events).toHaveLength(5); + + // Verify tool_use block + const blockStart = events[1]; + expect(blockStart.type).toBe('content_block_start'); + if (blockStart.type === 'content_block_start') { + expect(blockStart.content_block.type).toBe('tool_use'); + if (blockStart.content_block.type === 'tool_use') { + expect(blockStart.content_block.name).toBe('Read'); + expect(blockStart.content_block.id).toBe('toolu_abc'); + } + } + + // Verify stop_reason is tool_use + const msgDelta = events[4]; + if (msgDelta.type === 'message_delta') { + expect(msgDelta.delta.stop_reason).toBe('tool_use'); + } + }); + + it('sends correct request body with system prompt and tools', async () => { + fetchSpy.mockResolvedValueOnce( + new Response( + mockSSEStream([ + { + event: 'message_delta', + data: JSON.stringify({ + type: 'message_delta', + delta: { stop_reason: 'end_turn', stop_sequence: null }, + usage: { output_tokens: 0 }, + }), + }, + ]), + { status: 200, headers: { 'Content-Type': 'text/event-stream' } }, + ), + ); + + const tools = [ + { + name: 'Bash', + description: 'Run a command', + input_schema: { type: 'object', properties: { command: { type: 'string' } } }, + }, + ]; + + const session = new AnthropicSession( + testConfig({ model: 'gpt-5.5', tools }), + { baseUrl: 'http://test:9090', apiKey: 'test-key' }, + ); + + await collectEvents(session.turn([{ role: 'user', content: 'list files' }])); + + expect(fetchSpy).toHaveBeenCalledOnce(); + const [url, init] = fetchSpy.mock.calls[0]; + expect(url).toBe('http://test:9090/v1/messages'); + + const body = JSON.parse(init?.body as string); + expect(body.model).toBe('gpt-5.5'); + expect(body.max_tokens).toBe(100); + expect(body.stream).toBe(true); + expect(body.system).toBe('You are a test assistant.'); + expect(body.tools).toEqual(tools); + expect(body.messages).toEqual([{ role: 'user', content: 'list files' }]); + }); + + it('throws on API error', async () => { + fetchSpy.mockResolvedValueOnce( + new Response('{"error": "rate limited"}', { status: 429 }), + ); + + const session = new AnthropicSession(testConfig(), { + baseUrl: 'http://test:9090', + apiKey: 'test-key', + }); + + await expect( + collectEvents(session.turn([{ role: 'user', content: 'hello' }])), + ).rejects.toThrow('Anthropic API error 429'); + }); + + it('skips ping and message_stop SSE events', async () => { + const sseBody = mockSSEStream([ + { event: 'ping', data: '{}' }, + { + event: 'message_start', + data: JSON.stringify({ + type: 'message_start', + message: { + id: 'msg_ping', + model: 'claude-haiku-4-5', + role: 'assistant', + usage: { input_tokens: 5, output_tokens: 0 }, + }, + }), + }, + { event: 'message_stop', data: '{}' }, + { + event: 'message_delta', + data: JSON.stringify({ + type: 'message_delta', + delta: { stop_reason: 'end_turn', stop_sequence: null }, + usage: { output_tokens: 1 }, + }), + }, + ]); + + fetchSpy.mockResolvedValueOnce( + new Response(sseBody, { + status: 200, + headers: { 'Content-Type': 'text/event-stream' }, + }), + ); + + const session = new AnthropicSession(testConfig(), { + baseUrl: 'http://test:9090', + apiKey: 'test-key', + }); + + const events = await collectEvents(session.turn([{ role: 'user', content: 'hi' }])); + + // Only message_start + message_delta (ping and message_stop filtered) + expect(events).toHaveLength(2); + expect(events[0].type).toBe('message_start'); + expect(events[1].type).toBe('message_delta'); + }); + + it('handles chunked SSE delivery (split across reads)', async () => { + // Simulate SSE data arriving in multiple chunks, with event split across reads + const chunk1 = 'event: message_start\ndata: {"type":"message_start","message":'; + const chunk2 = + '{"id":"msg_chunked","model":"claude-haiku-4-5","role":"assistant","usage":{"input_tokens":5,"output_tokens":0}}}\n\n'; + const chunk3 = + 'event: message_delta\ndata: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":1}}\n\n'; + + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(encodeChunk(chunk1)); + controller.enqueue(encodeChunk(chunk2)); + controller.enqueue(encodeChunk(chunk3)); + controller.close(); + }, + }); + + fetchSpy.mockResolvedValueOnce( + new Response(stream, { + status: 200, + headers: { 'Content-Type': 'text/event-stream' }, + }), + ); + + const session = new AnthropicSession(testConfig(), { + baseUrl: 'http://test:9090', + apiKey: 'test-key', + }); + + const events = await collectEvents(session.turn([{ role: 'user', content: 'hi' }])); + expect(events).toHaveLength(2); + expect(events[0].type).toBe('message_start'); + expect(events[1].type).toBe('message_delta'); + }); + }); + + describe('headers', () => { + it('sends correct auth and version headers', async () => { + fetchSpy.mockResolvedValueOnce( + new Response( + mockSSEStream([ + { + event: 'message_delta', + data: JSON.stringify({ + type: 'message_delta', + delta: { stop_reason: 'end_turn', stop_sequence: null }, + usage: { output_tokens: 0 }, + }), + }, + ]), + { status: 200, headers: { 'Content-Type': 'text/event-stream' } }, + ), + ); + + const session = new AnthropicSession(testConfig(), { + baseUrl: 'http://test:9090', + apiKey: 'my-key-123', + apiVersion: '2024-01-01', + }); + + await collectEvents(session.turn([{ role: 'user', content: 'hi' }])); + + const headers = (fetchSpy.mock.calls[0][1]?.headers ?? {}) as Record; + expect(headers['x-api-key']).toBe('my-key-123'); + expect(headers['anthropic-version']).toBe('2024-01-01'); + expect(headers['Content-Type']).toBe('application/json'); + }); + }); +}); diff --git a/packages/harness/src/index.ts b/packages/harness/src/index.ts index ba7eb243..0dc6bf62 100644 --- a/packages/harness/src/index.ts +++ b/packages/harness/src/index.ts @@ -95,6 +95,19 @@ export type { } from './providers/index.js'; export { MODEL_COSTS, calculateCost } from './providers/index.js'; +// Agentic session — provider-agnostic streaming model interface +export { AnthropicSession, createAnthropicSession } from './providers/index.js'; +export type { + ModelSession, + ModelSessionConfig, + ModelSessionFactory, + StreamEvent, + ConversationMessage, + ContentBlock, + ToolDefinition, + AnthropicSessionOptions, +} from './providers/index.js'; + // Reasoning — deliberation + fusion orchestrators export { DeliberationOrchestrator, diff --git a/packages/harness/src/providers/anthropic-session.ts b/packages/harness/src/providers/anthropic-session.ts new file mode 100644 index 00000000..977b280d --- /dev/null +++ b/packages/harness/src/providers/anthropic-session.ts @@ -0,0 +1,209 @@ +/** + * Anthropic Messages API session — calls Claude (or any provider via praxis-proxy). + * + * Sends requests in Anthropic Messages format and streams back SSE events. + * When pointed at praxis-proxy (:9090), the proxy handles translation to + * whatever backend provider is configured (OpenAI, local models, etc.). + * + * This means Mitzo always speaks Anthropic wire format — praxis-proxy + * is the translation layer. + */ + +import { createLogger } from '../logger.js'; +import type { + ConversationMessage, + ModelSession, + ModelSessionConfig, + StreamEvent, +} from './session-types.js'; + +const log = createLogger('provider:anthropic-session'); + +/** Default Anthropic API base URL. */ +const DEFAULT_BASE_URL = 'https://api.anthropic.com'; + +/** Praxis proxy URL for provider-agnostic routing. */ +const PRAXIS_URL = 'http://127.0.0.1:9090'; + +export interface AnthropicSessionOptions { + /** Base URL for the API. Defaults to Anthropic API or praxis-proxy. */ + baseUrl?: string; + /** API key. Falls back to ANTHROPIC_API_KEY env var. */ + apiKey?: string; + /** Use praxis-proxy instead of direct Anthropic API. */ + useProxy?: boolean; + /** Anthropic API version header. */ + apiVersion?: string; +} + +/** + * Parse an SSE line pair into a StreamEvent. + * SSE format: "event: \ndata: \n\n" + */ +function parseSSE(eventType: string, data: string): StreamEvent | null { + if (eventType === 'ping' || eventType === 'message_stop') return null; + if (data === '[DONE]') return null; + + try { + const parsed = JSON.parse(data); + // The event type in the SSE matches our StreamEvent.type + if ( + parsed.type === 'message_start' || + parsed.type === 'content_block_start' || + parsed.type === 'content_block_delta' || + parsed.type === 'content_block_stop' || + parsed.type === 'message_delta' + ) { + return parsed as StreamEvent; + } + return null; + } catch { + log.warn('failed to parse SSE data', { eventType, data: data.slice(0, 200) }); + return null; + } +} + +export class AnthropicSession implements ModelSession { + readonly provider = 'anthropic'; + private baseUrl: string; + private apiKey: string; + private apiVersion: string; + private config: ModelSessionConfig; + + constructor(config: ModelSessionConfig, options: AnthropicSessionOptions = {}) { + this.config = config; + + const useProxy = options.useProxy ?? !!process.env.MITZO_USE_PRAXIS; + this.baseUrl = options.baseUrl ?? (useProxy ? PRAXIS_URL : DEFAULT_BASE_URL); + this.apiKey = options.apiKey ?? process.env.ANTHROPIC_API_KEY ?? 'dummy'; + this.apiVersion = options.apiVersion ?? '2023-06-01'; + + log.info('session created', { + model: config.model, + baseUrl: this.baseUrl, + useProxy, + toolCount: config.tools?.length ?? 0, + }); + } + + async *turn(messages: ConversationMessage[]): AsyncIterable { + const body: Record = { + model: this.config.model, + max_tokens: this.config.maxTokens, + stream: true, + messages, + }; + + if (this.config.systemPrompt) { + body.system = this.config.systemPrompt; + } + + if (this.config.tools && this.config.tools.length > 0) { + body.tools = this.config.tools; + } + + if (this.config.thinking) { + body.thinking = this.config.thinking; + // When thinking is enabled, Anthropic requires removing max_tokens + // and using budget_tokens in the thinking block instead + } + + const url = `${this.baseUrl}/v1/messages`; + + log.debug('turn request', { + model: this.config.model, + messageCount: messages.length, + url, + }); + + const response = await fetch(url, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + 'x-api-key': this.apiKey, + 'anthropic-version': this.apiVersion, + }, + body: JSON.stringify(body), + signal: this.config.signal, + }); + + if (!response.ok) { + const errorBody = await response.text().catch(() => ''); + const err = new Error( + `Anthropic API error ${response.status}: ${errorBody.slice(0, 500)}`, + ); + log.error('API request failed', { + status: response.status, + body: errorBody.slice(0, 200), + }); + throw err; + } + + if (!response.body) { + throw new Error('Response body is null — streaming not supported'); + } + + // Parse SSE stream + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ''; + let currentEventType = ''; + + try { + while (true) { + const { done, value } = await reader.read(); + if (done) break; + + buffer += decoder.decode(value, { stream: true }); + const lines = buffer.split('\n'); + // Keep the last incomplete line in the buffer + buffer = lines.pop() ?? ''; + + for (const line of lines) { + if (line.startsWith('event: ')) { + currentEventType = line.slice(7).trim(); + } else if (line.startsWith('data: ')) { + const data = line.slice(6); + const event = parseSSE(currentEventType, data); + if (event) { + yield event; + } + currentEventType = ''; + } + // Empty lines and other lines are ignored (SSE separators) + } + } + + // Process any remaining buffer + if (buffer.trim()) { + const remaining = buffer.split('\n'); + for (const line of remaining) { + if (line.startsWith('event: ')) { + currentEventType = line.slice(7).trim(); + } else if (line.startsWith('data: ')) { + const data = line.slice(6); + const event = parseSSE(currentEventType, data); + if (event) { + yield event; + } + } + } + } + } finally { + reader.releaseLock(); + } + } +} + +/** + * Factory: create an AnthropicSession. + * + * Reads MITZO_USE_PRAXIS env var to decide whether to route through + * praxis-proxy (provider-agnostic) or direct to Anthropic. + */ +export function createAnthropicSession( + config: ModelSessionConfig, + options?: AnthropicSessionOptions, +): AnthropicSession { + return new AnthropicSession(config, options); +} diff --git a/packages/harness/src/providers/index.ts b/packages/harness/src/providers/index.ts index 14b2435e..88b8da1d 100644 --- a/packages/harness/src/providers/index.ts +++ b/packages/harness/src/providers/index.ts @@ -72,3 +72,21 @@ export type { export { calculateCost, MODEL_COSTS } from './types.js'; export { AnthropicVertexModelProvider } from './anthropic-vertex.js'; export { GoogleVertexModelProvider } from './google-vertex.js'; + +// Agentic session abstractions (provider-agnostic loop) +export type { + ModelSession, + ModelSessionConfig, + ModelSessionFactory, + StreamEvent, + ConversationMessage, + ContentBlock, + ToolDefinition, + MessageStartEvent, + ContentBlockStartEvent, + ContentBlockDeltaEvent, + ContentBlockStopEvent, + MessageDeltaEvent, +} from './session-types.js'; +export { AnthropicSession, createAnthropicSession } from './anthropic-session.js'; +export type { AnthropicSessionOptions } from './anthropic-session.js'; diff --git a/packages/harness/src/providers/session-types.ts b/packages/harness/src/providers/session-types.ts new file mode 100644 index 00000000..0cbae3b8 --- /dev/null +++ b/packages/harness/src/providers/session-types.ts @@ -0,0 +1,143 @@ +/** + * Provider-agnostic agentic session types. + * + * ModelSession is the abstraction that replaces the Agent SDK's query() function. + * It handles a single LLM call (one turn of the agentic loop), streaming back + * normalized events. Mitzo owns the agentic loop and calls ModelSession.turn() + * repeatedly until the model produces no tool_use blocks. + * + * The event types map closely to the Anthropic Messages API streaming format + * because praxis-proxy normalizes all providers to that wire format. + */ + +// ── Tool definitions (sent to the model) ───────────────────────── + +/** A tool the model can invoke. Provider-agnostic schema. */ +export interface ToolDefinition { + name: string; + description: string; + input_schema: Record; +} + +// ── Conversation messages ──────────────────────────────────────── + +/** A content block in a message. */ +export type ContentBlock = + | { type: 'text'; text: string } + | { type: 'thinking'; thinking: string } + | { type: 'tool_use'; id: string; name: string; input: Record } + | { type: 'tool_result'; tool_use_id: string; content: string; is_error?: boolean }; + +/** A message in the conversation history. */ +export interface ConversationMessage { + role: 'user' | 'assistant'; + content: string | ContentBlock[]; +} + +// ── Streaming events (from model) ──────────────────────────────── + +/** Fired when the model starts a new message. */ +export interface MessageStartEvent { + type: 'message_start'; + message: { + id: string; + model: string; + role: 'assistant'; + usage: { input_tokens: number; output_tokens: number }; + }; +} + +/** Fired when a content block begins. */ +export interface ContentBlockStartEvent { + type: 'content_block_start'; + index: number; + content_block: + | { type: 'text'; text: string } + | { type: 'thinking'; thinking: string } + | { type: 'tool_use'; id: string; name: string; input: Record }; +} + +/** Incremental content within a block. */ +export interface ContentBlockDeltaEvent { + type: 'content_block_delta'; + index: number; + delta: + | { type: 'text_delta'; text: string } + | { type: 'thinking_delta'; thinking: string } + | { type: 'input_json_delta'; partial_json: string }; +} + +/** Fired when a content block ends. */ +export interface ContentBlockStopEvent { + type: 'content_block_stop'; + index: number; +} + +/** Fired when the message is complete. */ +export interface MessageDeltaEvent { + type: 'message_delta'; + delta: { + stop_reason: 'end_turn' | 'tool_use' | 'max_tokens' | 'stop_sequence' | null; + stop_sequence?: string | null; + }; + usage: { output_tokens: number }; +} + +/** Union of all streaming events from a single turn. */ +export type StreamEvent = + | MessageStartEvent + | ContentBlockStartEvent + | ContentBlockDeltaEvent + | ContentBlockStopEvent + | MessageDeltaEvent; + +// ── Session configuration ──────────────────────────────────────── + +/** Configuration for a model session. */ +export interface ModelSessionConfig { + /** Model identifier (e.g. 'claude-opus-4-6', 'gpt-5.5'). */ + model: string; + /** System prompt text. */ + systemPrompt: string; + /** Maximum output tokens per turn. */ + maxTokens: number; + /** Available tools. */ + tools?: ToolDefinition[]; + /** Enable extended thinking. */ + thinking?: { type: 'enabled'; budget_tokens: number }; + /** Abort signal for cancellation. */ + signal?: AbortSignal; +} + +// ── ModelSession interface ─────────────────────────────────────── + +/** + * A provider-agnostic streaming model session. + * + * Each call to turn() sends the conversation history to the model and + * streams back events. The caller (Mitzo's agentic loop) inspects the + * response for tool_use blocks, executes them, appends tool_result + * messages, and calls turn() again. + * + * Implementations: + * - AnthropicSession: calls Anthropic Messages API (direct or via praxis-proxy) + * - (future) ResponsesSession: calls OpenAI Responses API with server-side state + */ +export interface ModelSession { + /** Provider name for logging/tracing. */ + readonly provider: string; + + /** + * Execute one turn of the conversation. + * + * @param messages - Full conversation history (client-owned state). + * @returns An async iterable of streaming events for this turn. + */ + turn(messages: ConversationMessage[]): AsyncIterable; +} + +/** + * Factory function to create a ModelSession from config. + * Different providers have different construction needs. + */ +export type ModelSessionFactory = (config: ModelSessionConfig) => ModelSession; From 7017441f3909d34931a866e84a8c001372ff29cd Mon Sep 17 00:00:00 2001 From: dimakis Date: Sat, 4 Jul 2026 11:07:12 +0100 Subject: [PATCH 2/4] style: format with prettier Co-Authored-By: Claude Opus 4.6 --- .../harness/__tests__/anthropic-session.test.ts | 16 ++++++++-------- .../harness/src/providers/anthropic-session.ts | 4 +--- 2 files changed, 9 insertions(+), 11 deletions(-) diff --git a/packages/harness/__tests__/anthropic-session.test.ts b/packages/harness/__tests__/anthropic-session.test.ts index 5fb53813..928e7235 100644 --- a/packages/harness/__tests__/anthropic-session.test.ts +++ b/packages/harness/__tests__/anthropic-session.test.ts @@ -260,7 +260,9 @@ describe('AnthropicSession', () => { { baseUrl: 'http://test:9090', apiKey: 'test-key' }, ); - const events = await collectEvents(session.turn([{ role: 'user', content: 'read /tmp/test' }])); + const events = await collectEvents( + session.turn([{ role: 'user', content: 'read /tmp/test' }]), + ); expect(events).toHaveLength(5); @@ -307,10 +309,10 @@ describe('AnthropicSession', () => { }, ]; - const session = new AnthropicSession( - testConfig({ model: 'gpt-5.5', tools }), - { baseUrl: 'http://test:9090', apiKey: 'test-key' }, - ); + const session = new AnthropicSession(testConfig({ model: 'gpt-5.5', tools }), { + baseUrl: 'http://test:9090', + apiKey: 'test-key', + }); await collectEvents(session.turn([{ role: 'user', content: 'list files' }])); @@ -328,9 +330,7 @@ describe('AnthropicSession', () => { }); it('throws on API error', async () => { - fetchSpy.mockResolvedValueOnce( - new Response('{"error": "rate limited"}', { status: 429 }), - ); + fetchSpy.mockResolvedValueOnce(new Response('{"error": "rate limited"}', { status: 429 })); const session = new AnthropicSession(testConfig(), { baseUrl: 'http://test:9090', diff --git a/packages/harness/src/providers/anthropic-session.ts b/packages/harness/src/providers/anthropic-session.ts index 977b280d..2477c494 100644 --- a/packages/harness/src/providers/anthropic-session.ts +++ b/packages/harness/src/providers/anthropic-session.ts @@ -129,9 +129,7 @@ export class AnthropicSession implements ModelSession { if (!response.ok) { const errorBody = await response.text().catch(() => ''); - const err = new Error( - `Anthropic API error ${response.status}: ${errorBody.slice(0, 500)}`, - ); + const err = new Error(`Anthropic API error ${response.status}: ${errorBody.slice(0, 500)}`); log.error('API request failed', { status: response.status, body: errorBody.slice(0, 200), From 5d1a20f97f57597b9f34b942eab82c497e968d4f Mon Sep 17 00:00:00 2001 From: dimakis Date: Sat, 4 Jul 2026 12:14:29 +0100 Subject: [PATCH 3/4] =?UTF-8?q?fix:=20address=20Centaur=20review=20?= =?UTF-8?q?=E2=80=94=20shape=20validation,=20API=20key=20warning,=20remove?= =?UTF-8?q?=20misleading=20comment?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.6 --- .../src/providers/anthropic-session.ts | 40 +++++++++++++------ 1 file changed, 27 insertions(+), 13 deletions(-) diff --git a/packages/harness/src/providers/anthropic-session.ts b/packages/harness/src/providers/anthropic-session.ts index 2477c494..ea4aa2f1 100644 --- a/packages/harness/src/providers/anthropic-session.ts +++ b/packages/harness/src/providers/anthropic-session.ts @@ -46,17 +46,29 @@ function parseSSE(eventType: string, data: string): StreamEvent | null { try { const parsed = JSON.parse(data); - // The event type in the SSE matches our StreamEvent.type - if ( - parsed.type === 'message_start' || - parsed.type === 'content_block_start' || - parsed.type === 'content_block_delta' || - parsed.type === 'content_block_stop' || - parsed.type === 'message_delta' - ) { - return parsed as StreamEvent; + if (typeof parsed !== 'object' || parsed === null || typeof parsed.type !== 'string') { + return null; + } + + switch (parsed.type) { + case 'message_start': + if (!parsed.message?.id || !parsed.message?.role) return null; + return parsed as StreamEvent; + case 'content_block_start': + if (typeof parsed.index !== 'number' || !parsed.content_block?.type) return null; + return parsed as StreamEvent; + case 'content_block_delta': + if (typeof parsed.index !== 'number' || !parsed.delta?.type) return null; + return parsed as StreamEvent; + case 'content_block_stop': + if (typeof parsed.index !== 'number') return null; + return parsed as StreamEvent; + case 'message_delta': + if (!parsed.delta) return null; + return parsed as StreamEvent; + default: + return null; } - return null; } catch { log.warn('failed to parse SSE data', { eventType, data: data.slice(0, 200) }); return null; @@ -75,9 +87,13 @@ export class AnthropicSession implements ModelSession { const useProxy = options.useProxy ?? !!process.env.MITZO_USE_PRAXIS; this.baseUrl = options.baseUrl ?? (useProxy ? PRAXIS_URL : DEFAULT_BASE_URL); - this.apiKey = options.apiKey ?? process.env.ANTHROPIC_API_KEY ?? 'dummy'; + this.apiKey = options.apiKey ?? process.env.ANTHROPIC_API_KEY ?? ''; this.apiVersion = options.apiVersion ?? '2023-06-01'; + if (!this.apiKey && !useProxy) { + log.warn('no API key configured and not using praxis-proxy — requests will fail with 401'); + } + log.info('session created', { model: config.model, baseUrl: this.baseUrl, @@ -104,8 +120,6 @@ export class AnthropicSession implements ModelSession { if (this.config.thinking) { body.thinking = this.config.thinking; - // When thinking is enabled, Anthropic requires removing max_tokens - // and using budget_tokens in the thinking block instead } const url = `${this.baseUrl}/v1/messages`; From a4e005d5ed5a0b3b806576349dbd1e7de0bf96bd Mon Sep 17 00:00:00 2001 From: dimakis Date: Sat, 4 Jul 2026 12:59:36 +0100 Subject: [PATCH 4/4] fix: add usage check to message_delta shape validation Co-Authored-By: Claude Opus 4.6 --- packages/harness/src/providers/anthropic-session.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/harness/src/providers/anthropic-session.ts b/packages/harness/src/providers/anthropic-session.ts index ea4aa2f1..9a11cd8c 100644 --- a/packages/harness/src/providers/anthropic-session.ts +++ b/packages/harness/src/providers/anthropic-session.ts @@ -64,7 +64,7 @@ function parseSSE(eventType: string, data: string): StreamEvent | null { if (typeof parsed.index !== 'number') return null; return parsed as StreamEvent; case 'message_delta': - if (!parsed.delta) return null; + if (!parsed.delta || !parsed.usage) return null; return parsed as StreamEvent; default: return null;