Skip to content

spike(examples): serve the pi harness through Channels - #2441

Draft
cjol wants to merge 12 commits into
channels/11-telegramfrom
channels/12-pi-harness-spike
Draft

cjol wants to merge 12 commits into
channels/11-telegramfrom
channels/12-pi-harness-spike

Conversation

@cjol

@cjol cjol commented Oct 1, 2026 •

Copy link
Copy Markdown
Member

This PR is an experimental spike on top of the Channels stack. It serves the example PiHarness from examples/next/harnesses/pi through Channels, so the Web Channel browser client and npx agents tui talk to a pi-durable agent. It is a proof of concept, not meant to merge as is.

Why

  • Channels and the examples/next/harnesses both model a conversation with an agent, and the harness shape is the more stable of the two. Running one on the other shows where Channels' types should move to match it.
  • The bridge is typed against ConversationHarness, the part of the shared harness shape it uses (submit, abort, wait, messages, session().events()), rather than against PiHarness directly. PiHarness satisfies it structurally.

Code Changes

  • examples/next/pi-channels/src/harness-channels.ts maps the two models. An inbound message is a harness submission whose operation id is the event id and the turn id. Turn status follows pi's submission events. One pi run is one Channels response, fed from session.events(). The transcript and snapshot are projections of harness.messages(), with tool results folded into their call. cancel is harness.abort. Approvals and client tool results are ignored because pi-durable has neither yet.
  • src/server.ts is one Durable Object per room with PiHarness, Streams, Channels with the Web Channel, and WebSockets, behind a ChannelGateway. It imports the harness and Workers AI provider from ../harnesses/pi/src. It uses @cloudflare/vite-plugin 1.62, which detects a workers.dev behind Access for remote bindings and authenticates interactively or with CLOUDFLARE_ACCESS_CLIENT_ID / CLOUDFLARE_ACCESS_CLIENT_SECRET.
  • src/client.tsx is the minimal example's browser client. pnpm tui [room] runs the Pi TUI client.
  • The README records what the spike surfaced: harness-owned message ids need mapping back to surface ids, one pi run can answer several turns (steers), a run spans several assistant messages, events are harness specific, and the in-memory event watch does not resume an in-flight run after eviction.

cjol added 11 commits October 1, 2026 15:02
This PR removes the first-pass Channels implementation so the rest of the stack can add the turns model without carrying two designs. It is the base of the Channels-on-turns stack.

- The first pass modelled Channels as delivery adapters driven by a host, with fallback and fan-out composites on top. The turns model makes the agent own its turns and transcript, and every surface renders the same conversation.
- Keeping both would leave two overlapping APIs, two sets of docs and two examples while the new layers land.
- Removing first, in one reviewable layer, keeps each later PR additive.

Breaking for `agents/channels`:

- Removed the host, `fallback` and `fanout` composites, and the voice, AI SDK and TanStack AI helpers from the first pass.
- Slack, Telegram and email keep ingress only. Delivery comes back in later layers of this stack.

- Deleted `packages/agents/src/channels/host`, `fallback.ts`, `fanout.ts`, `voice.ts`, `tanstack-ai.ts`, the old `ai-sdk.ts`, and the live test suite with its snapshots.
- Removed `examples/channels`, `docs/agents/channels.md` and `design/channels.md`.
- Reduced the Slack, Telegram and email adapters and their tests to verified ingress and normalization.
- Dropped the removed entry points from `package.json`, `scripts/build.ts`, `AGENTS.md` and the README.
This PR adds the vocabulary and wire types for the Channels turns model, with no runtime behaviour yet.

## Why

- Later layers (core, gateway, Web Channel, Slack, Telegram) all speak the same protocol. Landing the types and terms first lets reviewers agree on the model before reading implementations.
- A glossary and sequence diagrams make the turn, response and surface terms concrete enough to review against.

## Public API Surface

Additive, exported from `agents/channels`:

| Symbol | Kind | Notes |
| --- | --- | --- |
| `TranscriptMessage`, `MessagePart`, `ToolPart`, `ToolState` | types | The transcript every surface renders |
| `ResponseChunk` | type | The streamed grammar of one response |
| `TurnStatus`, `TurnOutcome` | types | A turn's published state |
| `ConversationSnapshot` | type | Messages plus open turns, sent on join |
| `InboundEvent`, `EventOrigin`, `Participant`, `ClientTool` | types | What a surface sends the agent, and from whom |
| `ConversationAgent` | interface | What an agent implements to serve Channels |
| `Json` | type | JSON values carried by the protocol |

## Code Changes

- `protocol.ts` holds the types and is re-exported from the channels index.
- `docs/CONTEXT.md` defines the terms. `docs/diagrams.md` walks through the gateway and one turn on every surface.
… surfaces

This PR adds `Channels`, a Lifecycle capability that serves one conversation's surfaces from the agent that owns its turns.

## Why

- The agent owns the transcript and turns. Something has to take inbound events to the agent, publish turn and message updates, and stream responses to every surface that joined.
- Doing this inside a capability keeps the agent's code to `onEvent` and `snapshot`, and reuses `Streams` for durable, resumable responses instead of a channel-specific stream store.

## Public API Surface

Additive, exported from `agents/channels`:

| Symbol | Kind | Notes |
| --- | --- | --- |
| `Channels` | class | Lifecycle capability; `receive`, response and publish methods |
| `ChannelsOptions`, `ChannelsHost`, `ChannelState` | types | Configuration and per-channel state |
| `ConversationChannel` | interface | What a surface implementation provides |
| `GatewayEvent`, `ConversationUpdate`, `ResponseEnding`, `ResponseReadOptions` | types | Events and updates passed between layers |
| `ResponseWriter`, `responseWriter`, `ResponseGrammarError` | API | Writes chunks and enforces the response grammar |

## Code Changes

- `conversations.ts` implements `Channels`: it routes events to `onEvent`, fans updates out to each `ConversationChannel`, and opens responses on `Streams`.
- `response.ts` validates chunk order so a malformed response fails at the writer, not on a surface.
…agent

This PR adds `ChannelGateway`, the Worker entry point that authenticates channel webhooks and routes each event to the agent holding its conversation.

## Why

- Webhooks arrive at the Worker, not the agent. The Worker has to verify them, normalize them, pick a route, and call the right Durable Object.
- Keeping routing in the gateway lets the agent stay unaware of providers. The agent only sees `GatewayEvent`s and an `EventOrigin`.
- Routes default to the event's thread, with `defaultRoute` and `onRoute` for applications that group threads differently.

## Public API Surface

Additive, exported from `agents/channels`:

| Symbol | Kind | Notes |
| --- | --- | --- |
| `ChannelGateway` | class | `fetch` for webhooks, plus outbound `deliver` |
| `ChannelGatewayOptions` | type | `channels`, `agent(route)`, `defaultRoute`, `findUser`, `onRoute` |
| `GatewayAgent`, `ChannelRouteEvent` | types | The agent stub shape and the route observer event |

## Code Changes

- `gateway.ts` verifies, normalizes, routes and dispatches, with a stable `dispatchId` per Channel and event.
- `stream.ts` and `channel.ts` gain the stream collection the gateway uses for outbound delivery.
- `gateway.test.ts` covers routing, identity and dispatch.
This PR adds the Web Channel, which connects browsers to a conversation over WebSockets, and its browser client.

- Browsers need to join a conversation, receive the snapshot and live updates, and send messages, approvals and client tool results.
- The agent already owns its WebSocket connections through the `WebSockets` capability. The Web Channel shares those connections instead of opening its own, so `WebSockets` gains `use` for adding protocol handlers that can claim a message.
- The gateway resolves who is connecting from trusted request data and passes it in a header, so the Web Channel never trusts identity from the client.

Additive:

| Symbol | Entry point | Notes |
| --- | --- | --- |
| `WebChannel`, `WebChannelOptions` | `agents/channels/web` | Server side of the Web Channel |
| `ClientFrame`, `ServerFrame`, `parseClientFrame`, `WEB_IDENTITY_HEADER`, `WebIdentity` | `agents/channels/web` | Wire protocol |
| `WebChannelClient`, `WebChannelClientState`, `WebChannelActivity` | `agents/channels/web/client` | Browser client |
| `ChannelGatewayOptions.web`, `GatewayWebIdentity` | `agents/channels` | Resolve a WebSocket upgrade |
| `WebSockets.use` | `agents/websockets` | Add handlers that run before the configured ones |

- `web/channel.ts` serves the protocol on shared connections. `web/client.ts` keeps transcript and turn state with `apply-chunks.ts`.
- The gateway upgrades `/channels/<route>` by default.
- A `ChannelsHarnessObject` in the worker tests drives both sides end to end.
This PR adds `agents/channels/ai-sdk`, which connects an AI SDK agent to Channels.

- An AI SDK agent speaks `UIMessage` and `UIMessageChunk`. Channels speaks `TranscriptMessage` and `ResponseChunk`. Without conversions, every agent would write the same mapping, including tool states and approvals.
- This is split from the example so the package change can be reviewed on its own.

Additive, from `agents/channels/ai-sdk`:

| Symbol | Kind | Notes |
| --- | --- | --- |
| `toResponseChunks` | function | Converts an AI SDK stream to response chunks |
| `toTranscriptMessage`, `toUIMessage` | functions | Message conversions in each direction |
| `awaitsInput`, `answerToolCall` | functions | Approvals and client tool results |
| `createSendMessageTool` | function | AI SDK tool that sends to a surface through the gateway |
| `AiSdkConversionOptions`, `CreateSendMessageToolOptions` | types | Options |

- `ai-sdk-turns.ts` holds the conversions. `ai-sdk.ts` adds the send tool and re-exports them.
- New entry point in `package.json` and `scripts/build.ts`, and a minor changeset for `agents`.
This PR adds `examples/channels-minimal-agent`, one AI SDK agent served to browsers through the Web Channel.

## Why

- The stack needs a small runnable agent that shows the turns model end to end: one conversation per object, one turn at a time, and queued messages.
- It uses the AI SDK connectors from the previous PR, Workers AI, and `Sessions`, `Streams`, `Tasks` and `WebSockets` for state.

## Code Changes

- `src/server.ts` implements the agent and the gateway. A client tool (`getLocation`) runs in the browser of the participant who asked, and an approval tool (`forgetEverything`) waits for consent.
- `src/client.tsx` is a Kumo chat UI on `WebChannelClient`.
- Adds the example to the root README and the lockfile.
This PR adds `WebChannelChatTransport`, an AI SDK `ChatTransport` over the Web Channel, and uses it for an `@ai-sdk/tui` client in the minimal example.

- AI SDK UIs expect a `ChatTransport`. A transport over `WebChannelClient` lets them join a Channels conversation without knowing the protocol.
- The example's terminal client shows the same conversation as the browser, which demonstrates that surfaces share turns.

Additive, from `agents/channels/web/ai-sdk`:

| Symbol | Kind | Notes |
| --- | --- | --- |
| `WebChannelChatTransport` | class | `ChatTransport<UIMessage>` over a `WebChannelClient` |
| `WebChannelChatTransportOptions`, `WebChannelSendOptions` | types | Options |
| `toUIMessageChunk` | function | From `agents/channels/ai-sdk`, response chunk to `UIMessageChunk` |

- `web/ai-sdk.ts` maps transport calls to Web Channel events and response chunks to `UIMessageChunk`s with `toUIMessageChunk`, added to `ai-sdk-turns.ts` and exported from `agents/channels/ai-sdk`.
- `examples/channels-minimal-agent/src/tui.ts` and a `tui` script run the terminal client.
This PR adds `npx agents tui <url>`, a terminal client for any conversation's Web Channel, built on `@earendil-works/pi-tui`.

## Why

- Developers need to talk to a Channels agent from a terminal without writing a client. The `@ai-sdk/tui` client in the example is tied to that example.
- The client covers what the Web Channel offers: live transcript and turns from every surface, approvals, client tool results and cancellation.
- Its dependencies are bundled into `dist/cli.js`, so installing `agents` adds none.

## Public API Surface

Additive CLI:

| Command or flag | Notes |
| --- | --- |
| `agents tui <url>` | Connects to a Web Channel URL |
| `--header`, `-H` | Extra headers on the WebSocket upgrade |
| `--as` | Sets the `as` query parameter, which demo agents read as the participant |
| `CF_ACCESS_CLIENT_ID`, `CF_ACCESS_CLIENT_SECRET` | Sent as an Access service token |

## Code Changes

- `channels/tui/` holds argument parsing, the socket, the view model, and rendering components with Markdown highlighting and collapsible reasoning and tool blocks.
- `cli/index.ts` dispatches `tui`. `scripts/build.ts` bundles the CLI.
- The minimal example gains a `tui2` script and a README entry. Patch changeset for `agents`.
This PR makes Slack a conversation surface: a Slack thread that joins a conversation sees every turn, from any surface.

- After the cleanup, Slack only had ingress. A thread should show the agent's answers as they stream, ask for approvals with buttons, and say when a turn fails.
- Turns started elsewhere are quoted above the answer, so a thread can follow a conversation that began in the browser.
- `ConversationSurfaces` holds the per-surface lifecycle so Telegram can reuse it in the next PR.

No new exports. `slack()` from `agents/channels/slack` now renders turns, approvals and failures on threads that joined a conversation, in addition to ingress.

- `surfaces.ts` (internal) tracks joined surfaces and drives a renderer per turn.
- `adapters/slack.ts` streams with `chat.startStream`, `appendStream` and `stopStream`, and posts and updates approval messages.
- `slack-surfaces.test.ts` runs Slack against the Channels harness.
This PR makes Telegram a conversation surface, on the same `ConversationSurfaces` lifecycle as Slack.

## Why

- After the cleanup, Telegram only had ingress. A chat that joins a conversation should see each turn stream in, answer approvals, and learn when a turn fails.
- Sharing the surface lifecycle with Slack keeps turn handling in one place. Only rendering is Telegram specific.

## Public API Surface

No new exports. `telegram()` from `agents/channels/telegram` now renders turns, approvals and failures on chats that joined a conversation, in addition to ingress.

## Code Changes

- `adapters/telegram.ts` streams turns with `sendMessageDraft` and finishes them as messages. Approvals ask the chat to reply YES or NO.
- `channels-exports.test-d.ts` checks the public exports across the stack.
@cjol
cjol added this pull request to stack #2440 October 1, 2026 14:21
@changeset-bot

changeset-bot Bot commented Oct 1, 2026 •

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: dae0d97

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 2 packages
Name Type
agents Minor
@cloudflare/agent-think Patch

Not sure what this means? Click here to learn what changesets are.

Click here if you're a maintainer who wants to add another changeset to this PR

@agent-think

agent-think Bot commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

✅ agents import sizes: no significant changes (2ff75490 → 1c7f7e92, workflow run)

@cjol
cjol force-pushed the channels/12-pi-harness-spike branch from 1c7f7e9 to 5814d31 Compare October 1, 2026 14:40
This PR is an experimental spike on top of the Channels stack. It serves the example `PiHarness` from `examples/next/harnesses/pi` through Channels, so the Web Channel browser client and `npx agents tui` talk to a pi-durable agent. It is a proof of concept, not meant to merge as is.

## Why

- Channels and the `examples/next/harnesses` both model a conversation with an agent, and the harness shape is the more stable of the two. Running one on the other shows where Channels' types should move to match it.
- The bridge is typed against `ConversationHarness`, the part of the shared harness shape it uses (`submit`, `abort`, `wait`, `messages`, `session().events()`), rather than against `PiHarness` directly. `PiHarness` satisfies it structurally.

## Code Changes

- `examples/next/pi-channels/src/harness-channels.ts` maps the two models. An inbound message is a harness submission whose operation id is the event id and the turn id. Turn status follows pi's `submission` events. One pi run is one Channels response, fed from `session.events()`. The transcript and snapshot are projections of `harness.messages()`, with tool results folded into their call. `cancel` is `harness.abort`. Approvals and client tool results are ignored because pi-durable has neither yet.
- `src/server.ts` is one Durable Object per room with `PiHarness`, `Streams`, `Channels` with the Web Channel, and `WebSockets`, behind a `ChannelGateway`. It imports the harness and Workers AI provider from `../harnesses/pi/src`. It uses `@cloudflare/vite-plugin` 1.62, which detects a `workers.dev` behind Access for remote bindings and authenticates interactively or with `CLOUDFLARE_ACCESS_CLIENT_ID` / `CLOUDFLARE_ACCESS_CLIENT_SECRET`.
- `src/client.tsx` is the minimal example's browser client. `pnpm tui [room]` runs the Pi TUI client.
- The README records what the spike surfaced: harness-owned message ids need mapping back to surface ids, one pi run can answer several turns (steers), a run spans several assistant messages, events are harness specific, and the in-memory event watch does not resume an in-flight run after eviction.
@cjol
cjol force-pushed the channels/12-pi-harness-spike branch from 5814d31 to dae0d97 Compare October 1, 2026 14:40
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.

1 participant