Skip to content

[M9] send() and execute() with a listener stream model calls like stream() (breaking) - #310

Merged
LinuxDevil merged 3 commits into
mainfrom
lou-m9-send-streams-with-listener
Oct 2, 2026
Merged

LinuxDevil merged 3 commits into
mainfrom
lou-m9-send-streams-with-listener

Conversation

@LinuxDevil

@LinuxDevil LinuxDevil commented Oct 2, 2026 •

Copy link
Copy Markdown
Owner

Closes #232

Breaking. send() and AgentExecutor.execute() with a listener (createAgent({ onEvent }), onAgentEvent, or the deprecated onEvent option) now stream model calls. A step's text arrives as several text.delta events, as it does in stream().

What changed

  • RunEvents (src/execution/agentRun.ts) now has two flags instead of streamed:

    • streamModelCalls: stream each model call when the provider can.
    • iterated: the caller iterates a stream().

    stream() sets both to true. observeRun (listeners, no iteration) sets streamModelCalls (default true) and iterated: false. RunEventSink.streamed was renamed to iterated. It is internal: no entry point exports RunEventSink.

  • AgentExecutor.guardOutput reads iterated, so output guardrails on send() still check only the final reply.

  • New ExecuteOptions.streamModelCalls?: boolean. It defaults to true with listeners, and false generates whole steps. stream() ignores it, and it is not added to createAgent().

  • The ctx.emit doc in hooks.ts now reads "set when the run has listeners or is streamed". The notes in createAgent.ts, AgentExecutor.ts, agentRun.ts and generateStep.ts are updated too.

  • MockLLMProvider.stream() now streams the same step that generate() returns. This is why round 1 changed Agent Forge's results. The old stream() never emitted tool calls, used different usage numbers, and its chunks were word + ' ', so the text ended in an extra space ('This is a mock response. '). Before this fix, Forge's mock runs stopped calling tools once they streamed. I checked this: with the new run loop and the old mock.ts, 9 Forge server tests fail (runRegistry ×6, chat, localTools, timeTravel). With the fixed mock, all of them pass. Two pinned expectations changed because of it: llm.test.ts ('Hello, world! ' → 'Hello, world!') and cloudflare.test.ts ("This is a mock response. " → "This is a mock response.").

  • Agent Forge:

    • abortableProvider.stream() now checks stop() between streamed chunks and after the last one. That is the same point at which generate() checked after its call returned.
    • The ChatPanel.tsx comment is updated.
    • No client code renders text.delta. The logs use text.done and the chat uses the reconciled messages, so no rendering change was needed.

Tests

  • src/execution/sendStreams.test.ts (new) covers:
    • send() with onEvent emits more than one text.delta per step, and their concatenation equals the text.done text.
    • result.usage, toolCalls, steps and the step.done usage equal those of the same run without a listener (no double counting).
    • execute({ onAgentEvent }) streams, and execute({ streamModelCalls: false }) gives one delta per step.
    • A provider without stream(), and a provider with supportsStreaming: false, each give one delta per step.
    • Output guardrails: send() with a listener and a tool-calling step does not check the intermediate text, while stream() does and trips.
    • createAgent({ retry, fallbackModels }) with a listener retries a failed streamed call and then falls back. The listener gets provider.retry and provider.fallback.
  • legacyEvents.test.ts: the supportsStreaming: () => false override is gone, and the assertion is unchanged.
  • compactionEvents.test.ts: tests the corrected ctx.emit rule (send() with a listener gets the hook's event).
  • cancellation.test.ts: the scripted provider stubs stream: vi.fn() (it returns undefined). Its supportsStreaming is now false ("generate-only"). The assertions are unchanged.
  • llm.test.ts: a new case checks that the mock's stream equals its generate(): text, tool calls, finish reason and usage.
  • Agent Forge:
    • runRegistry.test.ts has a new test that uses a real ws client on attachWebSocketServer. Each of the two mock steps has more than one text.delta and exactly one text.done, and the deltas add up to the text.done text.
    • abortableProvider.test.ts (new) covers a stream passed through, and a stream stopped while being read, which throws RunAbortedError.
  • reasoning.test.ts, agentStream.test.ts, guardrails, sub-agents, agent_await ([M4] Background sub-agents that need approval pause the lead at agent_await and resume #305), drift (M10c: compare a paused sub-agent's fingerprint fully on resume #300), remote usage ([M10b] Add remote sub-agent token usage to the lead's totals #297), the exporter and trace files ([M5a] lousho traces: file trace exporter, createAgent exporter, terminal viewer #286, [M5b] Agent Forge: persisted trace history in the Trace tab #301) and todo.updated ([N12] todo.updated stream event and useTodos() for React, Vue and Svelte #299) all pass unchanged.

Agent Forge, verified for real

I built the studio (npm run build:studio) and started node bin/lousho.js studio --port 4793 from an empty scratch directory. A script then drove it over HTTP and WS /agents/:id/stream. For the baseline, I ran the same script against a detached origin/main (a8711e1) studio on port 4794.

Flow origin/main this branch
PUT /agents/m9check (mock, current-date), POST /run "please use current-date" stopped, "This is a mock response." the same
WS event types run.start, step.start, text.delta, text.done, tool.start, tool.done, step.done, step.start, text.delta, text.done, step.done, run.done the same, except each text.delta is now 5 text.delta events ("This ", "is ", …) that join to the text.done text
Logs (WS log) trigger, llm, tool call, tool result, llm, finished (stop) identical
Spans over WS 8 8
Chat (POST /agents/m9check/message "hello there") transcript system/user/assistant/tool/assistant/user/assistant identical; the chat turn's one step came as 5 deltas + 1 text.done
GET /agents/m9check/traces 2 runs: 1 model call / 0 tools, 33 in / 6 out; 2 model calls / 1 tool, 37 in / 12 out identical
GET /runs/m9check/history [1, finished, stop, [], 39 tokens] identical
Approval (demo-approval agent): run → paused awaiting_approval with pendingApproval → POST /approve 202 → stopped "This is a mock response." as listed identical event sequence (… tool.start, approval.requested, step.done, run.done, run.start, tool.start, tool.done, …), except for the deltas

GET / served the client (200). Stop is covered by the existing "stop() then run() resumes from the last checkpoint" test and the new abortableProvider test.

Verification (on the merged head 37be264, after merging origin/main twice: a8711e1, then 59f2a75 / #303)

npx tsc --noEmit                                         exit 0
npm run lint                                             exit 0 (--max-warnings 0)
npm run build                                            exit 0
npm run build --workspace=packages/create-lousho-agent   exit 0
npm run test:types                                       exit 0
npm run docs:verify-snippets -- --skip-build             exit 0
npm run docs:llms:check                                  exit 0
npm run test:coverage                                    3580 passed, 1 failed: guardrails.test.ts "kills the underlying child process on timeout" (known flaky under load; 23/23 pass alone)
npx vitest run --coverage --config vitest.coverage.config.ts --coverage.reportOnFailure
                                                         246 files, 3581 passed, 7 skipped, 0 failed
npm run fallow                                           exit 0 - "0 above threshold · 5408 analyzed · maintainability 89.6 (good)"
npm run typecheck --workspace apps/agent-forge           exit 0
npm run typecheck:server --workspace apps/agent-forge    exit 0
npm run test --workspace apps/agent-forge -- --run       119 passed (14 files)
npm run test:server --workspace apps/agent-forge         133 passed (16 files)

Earlier full runs on d8d7415 also hit failures under load, and each was a different test. All of them pass when run alone, and none of them touches the run loop or the mock:

  • http.test.ts: timeout.
  • sandbox-wiring.test.ts: timeout.
  • NodeWorkspace.test.ts shell: timeout.
  • SubprocessSandbox.test.ts Docker integration: "no such container", because another agent is using the daemon.
  • backgroundSubagents.test.ts "reports a sub-agent paused for approval…": the test waits a fixed 10 ms. It sends without a listener, so it takes the unchanged generate() path.

Peer matrix, run locally (the CI peers job steps), for the provider and execution tests (src/providers src/execution src/createAgent src/context src/testing src/subagents):

  • ai@6.0.300 (ai@6 @ai-sdk/openai@3 @ai-sdk/anthropic@3): tsc 0, test:types 0, vitest 1168 passed.
  • ai@7.0.127 (ai@7 @ai-sdk/openai@4 @ai-sdk/anthropic@4): tsc 0, test:types 0, build 0, vitest 1170 passed.

After the matrix, npm ci restored the defaults (ai 4.3.19), and the lockfile is unchanged. I did not run pack-smoke: this change does not touch exports, bin, package.json or the build.

Live test

Skipped: the OpenRouter account is out of credit (/credits: total_usage 10.20 ≥ total_credits 10). Live test spend: before 0, after 0 (no calls). src/execution/sendStreams.live.test.ts is committed. It skips without OPENROUTER_API_KEY and runs only with npm run test:live. It records src/execution/__fixtures__/cassettes/send-streams.json. The replay test (sendStreams.replay.test.ts) is not included, because it needs that cassette. A checklist line was added to #260.

Docs

These pages were edited, with no heading changes:

  • docs/streaming.md: "Listening without iterating", and "Token streaming and providers".
  • docs/compaction.md, docs/hooks.md: the ctx.emit rule.
  • docs/guardrails.md: the output row says "a stream() run (one you iterate)".
  • docs/executor-api.md: a new streamModelCalls row.

There are no new pages. CHANGELOG has a BREAKING entry with the migration note, plus entries for the mock provider and Agent Forge. llms.txt / llms-full.txt are regenerated.

🤖 Generated with Claude Code

LinuxDevil and others added 2 commits October 2, 2026 21:15
…am()

RunEvents splits its one flag in two: `streamModelCalls` (stream each model
call when the provider can) and `iterated` (the caller iterates a stream()).
A run with listeners but no iteration streams its model calls; output
guardrails read `iterated`, so send() keeps checking only the final reply.
New ExecuteOptions.streamModelCalls (default true) restores whole steps.

MockLLMProvider.stream() now streams the step generate() returns (tool
calls, finish reason, usage, exact text): its old stream() dropped tool
calls and added a trailing space, which is what changed Agent Forge's
results in round 1. Forge's abortable provider checks stop() between
streamed chunks.

Closes #232

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[M9] send() and execute() with a listener stream model calls like stream()

1 participant