Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- `http_request` is no longer open to DNS rebinding (N13a). It resolved a host name once to check it against the private-address list and then let `fetch` resolve it again to connect, so a name that answered a public address to the check and `127.0.0.1` (or a cloud metadata or intranet address) to the connection reached the private address; the docs and a comment in the Worker runtime described the check as rebinding-safe. Now the one resolution is the check: the in-process path connects through an undici `Agent` whose `connect.lookup` checks every address and hands the socket the address it checked (each redirect hop included), and the sandboxed path (`sandboxExecute`, which agents use) resolves and checks the host in the agent's process and has the sandboxed process connect to that address (node:http/https with a fixed lookup; the host name still goes in `Host` and TLS SNI). The private-address list is now the shared one (below), which adds `0.0.0.0/8`, `100.64.0.0/10`, `192.0.0.0/24`, `198.18.0.0/15`, multicast and reserved `224.0.0.0/3`, `::`, `ff00::/8`, NAT64 `64:ff9b::/96` and 6to4 `2002::/16` to what `http_request` refused. Behavior changes: a host name that cannot be resolved now fails with the resolver's error instead of being passed to `fetch`; `undici` is loaded on the first `http_request` request (not only with `validateSSL: false`). New option `createHttpTool({ allowPrivate })` for hosts that may resolve to private addresses. Cloudflare Workers have no DNS hook, so neither tool exists in the Worker build; docs/deployment.md now says what a Worker can and cannot check.

### Changed
- BREAKING: `send()` and `AgentExecutor.execute()` with a listener (`createAgent({ onEvent })`, `onAgentEvent`, or the deprecated `onEvent` option) now stream model calls (M9, #232): a step's text arrives as several `text.delta` events, as in `stream()`, instead of one `text.delta` per step after the step finished. What a listener receives otherwise does not change: the same event types in the same order, exactly one `text.done` per step with the step's whole text, the same `step.done` usage and `run.done`, the same `ExecutionResult` (text, tool calls, usage, steps). A provider without `stream()`, or whose `supportsStreaming(model)` is `false`, still gives one `text.delta` per step. Output guardrails on `send()` still check only the final reply (a `stream()` you iterate checks every step's text, as before). `send()` without a listener still calls `provider.generate()`. Hooks' `ctx.emit` is set whenever the run has listeners (it was already; the docs said "only when streamed"). Migration: concatenate `text.delta`, or listen to `text.done` for whole steps; `execute({ streamModelCalls: false })` (new option, default `true`) restores whole steps.
- `MockLLMProvider` / `createMockProvider()`: `stream()` now streams the same step `generate()` returns (M9): the same tool calls (it streamed none, so a streamed mock run never called a tool), finish reason and usage, and text chunks that add up to exactly the generated text (the old chunks ended with an extra space: `'This is a mock response. '`). Runs on the mock provider with a listener (Agent Forge's, for one) therefore give the same results as before, and `stream()` runs on it now call tools as `send()` does.
- Agent Forge (M9): runs now stream their model calls, so the run WebSocket (`WS /agents/:id/stream`) carries several `text.delta` events per step; the chat, logs and trace views are unchanged (they use `text.done` and the reconciled messages). Stop is checked between streamed chunks and after the last one, as it was after each `generate()` call.
- CI peer matrix (LOU-M8, #231): the `typecheck-ai7` and `typecheck-zod4` jobs are replaced by one `peers` job with five entries, each installed for real on top of the default install and run through `tsc`, `test:types`, both builds and `npx vitest run`: `ai4-zod4`, `ai6-zod3`, `ai6-zod4`, `ai7-zod3`, `ai7-zod4`. The two zod 4 entries on `ai` 6/7 also install `ollama-ai-provider-v2` (3.x / 4.x), and the new `src/providers/ollamaV2.contract.test.ts` runs `OllamaProvider` against the real package and a local fake Ollama server (`generate()`, `stream()`, a tool-call turn through `createAgent().send()` and `.stream()`). The `ai-v6` dev alias replaces the hand-made `ai` 6 stand-in in `aiMajorPeers.test.ts`.
- `lousho init --provider ollama` now scaffolds `ai@^7.0.0` with `ollama-ai-provider-v2@^4.0.0` and `zod@^4.0.0` (it was `ai@^4.3.19` with `ollama-ai-provider@^1.2.0` and zod 3). Existing projects are not touched; to stay on the old pairing keep `ai@^4.3.19` and `ollama-ai-provider@^1.2.0`.
- A run paused inside a sub-agent compares the sub-agent with its current definition on resume, under the lead's `onAgentDrift` (M10c). Behaviour change: with `onAgentDrift: 'error'`, `agent.approvals.resolve()` (and `resumeAfterApproval()`) now also rejects with `LOUSHO_AGENT_DRIFT` when a paused sub-agent's instructions, model or tools changed; before, only tool and provider default-model changes of a sub-agent were seen, and always as a warning. This holds at any depth (a sub-agent of a sub-agent uses the lead's mode). With `'warn'` (still the default) the `agent.drift` event carries `subagent`. A rejected resume puts the lead's approval record and its `'awaiting-approval'` checkpoint back, so fixing the sub-agent and resolving again finishes the run; before, a sub-agent's missing-tool error became an error result of the `task` call and the run continued. Migration: none for `'warn'` / `'ignore'`; with `'error'`, resolve approvals paused before a sub-agent changed with the old definition, or use `'warn'`. See docs/durable-execution.md#resuming-with-a-changed-agent.
Expand Down
47 changes: 47 additions & 0 deletions apps/agent-forge/server/__tests__/runRegistry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,12 @@ import { describe, it, expect, beforeEach, afterEach } from 'vitest';
import * as fs from 'node:fs';
import * as path from 'node:path';
import * as os from 'node:os';
import * as http from 'node:http';
import type { AddressInfo } from 'node:net';
import { WebSocket } from 'ws';
import type { AgentSpec } from '@lousho/build-ai-agent';
import { RunManager } from '../runRegistry';
import { attachWebSocketServer } from '../wsServer';
import { FileCheckpointStore } from '../checkpointStore';
import { FileApprovalStore } from '../approvalStore';
import { graphToSpec } from '../../src/graph/graphToSpec';
Expand Down Expand Up @@ -120,6 +124,49 @@ describe('RunManager', () => {
expect(afterResume?.messages.map((m) => m.content)).toContain('a follow-up appended on resume');
});

it('streams each model step over the WebSocket: several text.delta, one text.done per step (M9)', async () => {
const server = http.createServer();
const wss = attachWebSocketServer(server, runManager);
await new Promise<void>((resolve) => server.listen(0, '127.0.0.1', resolve));
const { port } = server.address() as AddressInfo;
const ws = new WebSocket(`ws://127.0.0.1:${port}/agents/agent-ws/stream`);
const events: any[] = [];
try {
ws.on('message', (data) => {
const message = JSON.parse(String(data));
if (message.type === 'event') events.push(message.payload);
});
await new Promise((resolve, reject) => {
ws.once('open', resolve);
ws.once('error', reject);
});

// A tool call, then the final reply: two model steps, each streamed by the mock provider.
await runManager.run('agent-ws', 'please use current-date', SPEC);
const final = await waitForStatus(runManager, 'agent-ws', (s) => s.status === 'stopped');
await new Promise((resolve) => setTimeout(resolve, 50));

expect(final.resultText).toBe('This is a mock response.');
const steps: { deltas: string[]; done: string[] }[] = [];
for (const event of events) {
if (event.type === 'step.start') steps.push({ deltas: [], done: [] });
if (event.type === 'text.delta') steps.at(-1)!.deltas.push(event.text);
if (event.type === 'text.done') steps.at(-1)!.done.push(event.text);
}
expect(steps).toHaveLength(2);
for (const step of steps) {
expect(step.deltas.length).toBeGreaterThan(1);
expect(step.done).toEqual(['This is a mock response.']);
expect(step.deltas.join('')).toBe(step.done[0]);
}
expect(events.filter((e) => e.type === 'tool.done').map((e) => e.toolName)).toEqual(['current-date']);
} finally {
ws.close();
wss.close();
await new Promise((resolve) => server.close(resolve));
}
});

it('emits structured log entries derived from the AgentEvent stream (O1)', async () => {
const logs: any[] = [];
runManager.on('log', (agentId: string, entry: any) => {
Expand Down
26 changes: 26 additions & 0 deletions apps/agent-forge/server/abortableProvider.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
import { describe, expect, it } from 'vitest';
import { MockLLMProvider } from '@lousho/build-ai-agent';
import { RunAbortedError, withAbortSignal } from './abortableProvider';

const request = { model: 'mock-1', messages: [{ role: 'user' as const, content: 'hi' }] };

describe('withAbortSignal (M9: streamed calls)', () => {
it('passes a stream through chunk by chunk when not stopped', async () => {
const provider = withAbortSignal(new MockLLMProvider({ name: 'mock' }), new AbortController().signal);
const streamed = await provider.stream(request);
const deltas: string[] = [];
for await (const chunk of streamed.fullStream) if (chunk.type === 'text-delta') deltas.push(chunk.textDelta ?? '');
expect(deltas.length).toBeGreaterThan(1);
expect(deltas.join('')).toBe(await streamed.text);
});

it('throws RunAbortedError from a stream stopped while it is read', async () => {
const controller = new AbortController();
const provider = withAbortSignal(new MockLLMProvider({ name: 'mock' }), controller.signal);
const streamed = await provider.stream(request);
const read = async () => {
for await (const _chunk of streamed.fullStream) controller.abort();
};
await expect(read()).rejects.toBeInstanceOf(RunAbortedError);
});
});
22 changes: 20 additions & 2 deletions apps/agent-forge/server/abortableProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@
* cancellation; this wrapper's shape (signal checked, forwarded if the
* real provider accepts one) is what that wiring would build on.
*/
import type { LLMProvider, GenerateOptions, GenerateResult, StreamResult } from '@lousho/build-ai-agent';
import type { LLMProvider, GenerateOptions, GenerateResult, StreamChunk, StreamResult } from '@lousho/build-ai-agent';

/**
* Named `AbortError` (the fetch/AbortSignal convention) on purpose: the SDK's
Expand Down Expand Up @@ -69,9 +69,27 @@ export function withAbortSignal(provider: LLMProvider, signal: AbortSignal): LLM
checkAborted();
return result;
},
// M9: AgentExecutor.execute({ onAgentEvent }) streams model calls, so a
// stop() during a streamed call is checked between chunks and after the
// last one - the same point generate() checks after its call returns.
async stream(options: GenerateOptions): Promise<StreamResult> {
checkAborted();
return provider.stream(options);
const streamed = await provider.stream(options);
async function* fullStream(): AsyncGenerator<StreamChunk> {
for await (const chunk of streamed.fullStream) {
checkAborted();
yield chunk;
}
checkAborted();
}
return {
fullStream: fullStream(),
textStream: streamed.textStream,
text: streamed.text,
usage: streamed.usage,
finishReason: streamed.finishReason,
toolCalls: streamed.toolCalls,
};
},
supportsTools(model: string): boolean {
return provider.supportsTools(model);
Expand Down
14 changes: 7 additions & 7 deletions apps/agent-forge/src/components/ChatPanel.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,13 @@
* ported from `.design-ref/agent-forge-mockup.html`'s bubble/avatar/
* timestamp layout.
*
* Streaming granularity: MESSAGE-level, not token-level. AgentExecutor's
* execution loop (src/execution/AgentExecutor.ts) only ever calls
* `provider.generate()` - never `provider.stream()` - so there is no
* per-token event to render incrementally even for a provider (like
* MockLLMProvider) that DOES implement `.stream()`. A message appears in
* the thread once its run turn settles (completes, or pauses for
* approval); the "typing" indicator below fills the gap, driven by the
* Streaming granularity: MESSAGE-level, not token-level. Since M9 the
* server's `AgentExecutor.execute({ onAgentEvent })` streams each model
* call through `provider.stream()` when the provider can, so the
* WebSocket carries several `text.delta` events per step - but this panel
* renders the reconciled `Message[]` history, not the deltas. A message
* appears in the thread once its run turn settles (completes, or pauses
* for approval); the "typing" indicator below fills the gap, driven by the
* real `status: 'running'` from LOU-N rather than a fake timeout.
*/
import { useEffect, useRef, useState } from 'react';
Expand Down
5 changes: 3 additions & 2 deletions docs/compaction.md
Original file line number Diff line number Diff line change
Expand Up @@ -58,8 +58,9 @@ strategy could not shrink anything, `tokensAfter` equals `tokensBefore`; when it
failed (or the summarizer failed and the hook fell back to pruning), `error` is
set and the run continues. `summary` is `true` when old turns were replaced by
a summary (the text itself is not sent; use `onCompaction` for it). A
non-streaming `send()` emits no events. Hooks add their own events with
`ctx.emit?.(...)` on the `preGenerate` context, which exists only in streamed runs.
`send()` without a listener emits no events. Hooks add their own events with
`ctx.emit?.(...)` on the `preGenerate` context, which exists when the run has
listeners (`createAgent({ onEvent })`, `onAgentEvent`) or is streamed.

## Compacting a run

Expand Down
1 change: 1 addition & 0 deletions docs/executor-api.md
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ console.log(events); // includes 'run.start' and 'run.done'
| `limits` | Token, cost, time and step budgets of the run; see [Budgets](./configuration.md#budgets). |
| `temperature`, `maxTokens` | Generation parameters. |
| `onAgentEvent` | Listener for the run's `AgentEvent`s (`run.start`, `tool.start`, `run.done`, ...); see [Streaming](./streaming.md#listening-without-iterating). |
| `streamModelCalls` | With a listener, stream each model call so text arrives as several `text.delta` events (default `true`); `false` generates whole steps, one `text.delta` each. Ignored by `stream()`. |
| `approvalStore`, `sessionId` | Human-in-the-loop approvals (see `resumeAfterApproval()`). |
| `checkpointStore` | Persist/resume execution checkpoints. |
| `exporter` | A `TraceExporter` for tracing spans (OpenTelemetry GenAI conventions, see [observability](observability.md)). |
Expand Down
2 changes: 1 addition & 1 deletion docs/guardrails.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ three points, each list in order:
| List | Runs on | When |
| ---- | ------- | ---- |
| `input` | Each new user message (the user messages that end the transcript) | Before the first model call. A block makes no model call. |
| `output` | The final assistant text; in a streamed run, every step's text | Before it is emitted: before its `text.done`, and before `run.done`. |
| `output` | The final assistant text; in a `stream()` run (one you iterate), every step's text. `send()` with a listener checks the final text only. | Before it is emitted: before its `text.done`, and before `run.done`. |
| `tools` | A tool call's arguments (`text` is them as JSON, `args` the object) | After the [permission rules](./approvals.md#permission-policies) (skipped when a rule denies the call) and before `needsApproval`. |

A guardrail is `{ name, check(ctx) }`. `ctx` has `kind` (`'input'`, `'output'`
Expand Down
2 changes: 1 addition & 1 deletion docs/hooks.md
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ Every method is optional: implement only the points you need. A hook that return
| `preGenerate(ctx)` | Before each model call. | `GenerateHookContext`: `request` (the live `GenerateOptions` about to be sent) and `emit?`, plus the common fields. | Nothing; mutate `ctx.messages` or `ctx.request`, or throw. |
| `postGenerate(ctx, result)` | After each model call resolves, with its `GenerateResult`. | `GenerateHookContext`. | Nothing; or throw. |

The common fields on every context are `agentId`, `agentName`, `sessionId`, `messages` (the live conversation history), `metadata` (a free-form bag) and `subagent` (set inside a sub-agent). `ctx.emit` exists only in streamed runs and adds `compaction.start` / `compaction.done` events to the stream.
The common fields on every context are `agentId`, `agentName`, `sessionId`, `messages` (the live conversation history), `metadata` (a free-form bag) and `subagent` (set inside a sub-agent). `ctx.emit` exists when the run has listeners (`createAgent({ onEvent })`, `onAgentEvent`) or is streamed, and adds `compaction.start` / `compaction.done` events to the run's events.

## Hook outcomes

Expand Down
14 changes: 9 additions & 5 deletions docs/streaming.md
Original file line number Diff line number Diff line change
Expand Up @@ -77,10 +77,13 @@ To observe every run without iterating one, give the agent a listener:
/ `stream()` / `resumeAfterApproval()` options. It is called synchronously
with each `AgentEvent` as it happens, on `send()` as on `stream()`, for
session turns and for runs resumed after an approval. It gets the same events,
in the same order, as iterating the stream would, with one difference:
`send()` / `execute()` generate each model step whole, so a step's text comes
as a single `text.delta` (as it does on `stream()` with a provider that cannot
stream). Sub-agents' events arrive tagged with `subagent`.
in the same order, as iterating the stream would: with a listener, `send()` /
`execute()` stream each model call too, so a step's text arrives as several
`text.delta` events (one `text.delta` per step when the provider cannot
stream). Concatenate them, or use `text.done` for each step's whole text;
`AgentExecutor.execute({ streamModelCalls: false })` generates whole steps.
Output guardrails still check only the final reply of a `send()`, as without
a listener. Sub-agents' events arrive tagged with `subagent`.

```ts
import { createAgent } from '@lousho/build-ai-agent';
Expand Down Expand Up @@ -234,7 +237,8 @@ provider has no `stream()`, or its `supportsStreaming(model)` returns
Streaming changes only how one model step is obtained. Hooks, argument
validation, approvals, parallel tool calls, checkpoints, cancellation,
tracing spans and the other `execute()` callbacks behave exactly as in
`send()`. `send()` and `execute()` themselves still use `generate()`.
`send()`. `send()` and `execute()` stream the same way when they have a
listener (`onEvent` / `onAgentEvent`); without one they use `generate()`.

A custom provider's `stream()` should yield `text-delta` chunks and a final
`finish` chunk with `finishReason` and `usage`. Tool calls can be yielded as
Expand Down
Loading
Loading