Skip to content

feat(pi): stream pi durable conversations to clients - #323

Merged
eersnington merged 8 commits into
stack/feat-pi-add-pidurable-actor-zlmoklxmfrom
stack/feat-pi-stream-pi-durable-conversations-to-clients-rvqykpzv
Oct 5, 2026
Merged

eersnington merged 8 commits into
stack/feat-pi-add-pidurable-actor-zlmoklxmfrom
stack/feat-pi-stream-pi-durable-conversations-to-clients-rvqykpzv

Conversation

@eersnington

@eersnington eersnington commented Oct 5, 2026 •

Copy link
Copy Markdown
Member

Clients can follow a conversation live. Each watch action returns the current value and then streams frames to the calling connection only.

import type { PiEventsFrame } from "@rivet-dev/pi/durable";

const chat = client.agent.getOrCreate(["chat-7"]).connect();
let view;
chat.on("pi.events", (frame) => {
	// RivetKit types event payloads as unknown on the client.
	const { seq, events } = frame as PiEventsFrame;
	view = applyEvents(view, events); // seq 0 replaces everything, as after a wake
});
const { snapshot } = await chat.conversation.watchEvents(root.id); // the starting point, including a partial answer
view = snapshot;
Action Returns Then sends
conversation.watchEvents(cid) { snapshot } pi.events batches of agent events
conversation.watch(cid) { value }: the conversation view pi.view Chord ops
harness.watchTaskGraph() { value } pi.taskGraph Chord ops
harness.watchDoc(kind, cid) { value }, or undefined before the document exists pi.doc with the whole value

Each has an unwatch* counterpart.

sequenceDiagram
  participant C as Client
  participant A as Actor
  C->>A: watchEvents(1)
  A-->>C: { snapshot } (frame 0)
  A->>C: pi.events seq 1, 2, 3 ...
  Note over A: actor sleeps, connection hibernates
  A->>C: on wake: pi.events seq 0 with a fresh snapshot
  A->>C: seq 1, 2 ...
Loading
  • One Pi watch per connection, conversation, and kind. Pi takes each snapshot atomically when the watch attaches, so a late joiner's snapshot matches its stream.
  • seq counts frames per watch. A client that sees a gap calls the watch action again. Nothing is replayed, as in Pi Durable and Cloudflare.
  • Watches are stored in pi_durable_watch. On wake, connections that are still open get their watches back with a seq 0 snapshot. This needs fix(rivetkit): stop hibernating actor connections by default rivet#5835, which delivers events sent from onWake.
  • Watch changes (watch, unwatch, disconnect, close, re-attach) run one at a time, so overlapping calls leave exactly one stream or none.
  • onDisconnect ends every watch of the connection.
  • A document watch can start before the document exists. When a commit creates it, the watch attaches and sends its value as the first frame. If that attach fails, it is logged and the next write retries.
  • A watch over a stateless HTTP call returns the starting value only. Its watch ends with the request, so use connect() for frames.

This is part 4 of 11 in a stack:

@eersnington
eersnington force-pushed the stack/feat-pi-add-pidurable-actor-zlmoklxm branch from d19db6a to a8731a7 Compare October 5, 2026 18:53
@eersnington
eersnington force-pushed the stack/feat-pi-stream-pi-durable-conversations-to-clients-rvqykpzv branch from 0bfc39b to 9657e88 Compare October 5, 2026 18:53
@eersnington
eersnington merged commit 93bf3de into stack/feat-pi-add-pidurable-actor-zlmoklxm Oct 5, 2026
0 of 3 checks passed
@eersnington
eersnington deleted the stack/feat-pi-stream-pi-durable-conversations-to-clients-rvqykpzv branch October 5, 2026 19:15
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