Skip to content

feat(examples): host pi-durable in the pi harness example - #2419

Merged
aron-cf merged 5 commits into
cloudflare:mainfrom
mattzcarey:feat/pi-durable-driver
Oct 1, 2026
Merged

aron-cf merged 5 commits into
cloudflare:mainfrom
mattzcarey:feat/pi-durable-driver

Conversation

@mattzcarey

@mattzcarey mattzcarey commented Sep 30, 2026 •

Copy link
Copy Markdown
Member

Moves examples/next/harnesses/pi from pi-agent-core's AgentHarness (vendored 0.84) to @earendil-works/pi-durable, pi's new durable harness, and replaces the Tasks driving with one Lifecycle job per session that PiHarness owns. Example-only: nothing is exported from agents.

export class PiAgent extends DurableObject<Env> {
  readonly harness = new PiHarness({
    models: createModels({ providers: [workersAI(this.env.AI)] }),
    model: { provider: "cloudflare-workers-ai", modelId: MODEL_ID },
    tools: createTools() // sleep, current_time
  });
  // App glue, not the harness: this app's socket protocol on session.events().
  readonly sockets = new PiSessionSockets(this.harness, (tag) =>
    this.ctx.getWebSockets(tag)
  );
  readonly webSockets = new WebSockets(this.sockets.options());
  readonly lifecycle = Lifecycle.install(this)
    .use(this.webSockets)
    .use(this.harness);

  async onStart() {
    await this.sockets.reattach(); // watches are in memory
  }
}

const { text } = await this.harness.prompt("What is 47 × 19?");
const side = await this.harness.sessions.create();
await side.submit("and now in hex", { whenBusy: "steer" });

How it fits together

flowchart LR
  submit["harness.submit()"] --> wake["1. push wake job pi-wake:session"]
  wake --> admit["2. pi: conversation.submit(requestId)\n(the one admission)"]
  admit --> kick["3. push wake job again (due now)"]
  kick --> step["onJob: pi has live tasks?"]
  step -->|yes| idle["waitForIdle in alarm work,\nreschedule +30s heartbeat"]
  idle -->|wait ends| kick
  step -->|long pi wait| later["reschedule to the deadline"]
  step -->|no| done["complete: no job, no alarm"]
  admit --> tables[("pi_* tables\n(session-store.ts)")]
  tables --> events["session.events()"]
  events -. app glue .-> ws["sockets.ts → WebSockets → view.ts"]
Loading
  • pi owns the run: transcript, inbox (steer / follow-up), generation and tool tasks, replay-safe vs unsafe tools, retries, abort. It stores all of it in its own schema.
  • session-store.ts: pi's SqliteDatabase facade over ctx.storage.sql (transactionSync, tables under pi_). pi's own storage conformance suite runs against it on a real Durable Object.
  • The wake is a Lifecycle job, not a primitive. One job per session (pi-wake:<session>), owned by PiHarness.onJob. Input goes to pi once, in submit(), after the job is pushed, so an eviction at any point leaves an alarm that restarts the object. The job waits while pi has live tasks in the session, reschedules itself as a heartbeat, and completes when there are none. It never admits input or replays model or tool work. pi resumes its own tasks on open. Long pi waits reschedule the job to the deadline, background tasks are re-checked every 30 s, and a wait caps at 10 minutes to stay under the alarm wall-time limit. A same-id push while the job is dispatching supersedes its outcome, which is what keeps a submit from being lost to a completing wake.
  • Harness surface: prompt, submit, abort, wait, messages, sessions (create, get, fork, list), session(id), and session.events(), which is pi's own AgentEvent stream (a snapshot, then one batch per commit). The harness picks no transport.
  • App glue (outside src/harness): sockets.ts + protocol.ts put one session per socket on WebSockets, and view.ts folds events for the UI and the tests. pi commits partial answers and tool output as it goes, so a late or reconnecting client (including after hibernation) gets the current state from a snapshot.
  • pi build: pins pi main (2bbfcca4) as vendored archives. npm 0.99.1 predates the inbox, events, ownership and subagents. vendor/pi-dev/pack.mjs rebuilds them.

Why not Tasks, the driver, or the state machine

pi-durable is already the replay authority, so the SDK only needs to wake the object. Tasks journals steps for replay (redundant here). The driver (#2396) worked, but with pi owning admission and replay it reduced to a heartbeat job per session, which is what Lifecycle jobs already are. The state machine (#2338) drives pi in passes, which pi-durable's scheduler removes the need for. Full reasoning is in NOTES.md.

Open issues (tracked in NOTES.md)

  • pi sleeps with setTimeout, which does not keep an object alive. Generation retry/poll deadlines are read from pi's LiveDoc and handed to the alarm; custom task sleeps are invisible.
  • Upstream asks for pi:
    • scheduler-owned sleeps with inspect().nextWakeAt, or an injectable timer
    • submission lookup by request id
    • conversation listing
    • a table-prefix option (today it is a SQL rewrite)
  • Not done yet:
    • compaction (not upstream yet)
    • subagent tools
    • ExecutionEnv for pi's coding tools

Tests

  • session-store.test.ts: pi's storage conformance suite.
  • harness.test.ts: faux-provider turns, follow-ups while busy, abort, sessions, and a crash mid-tool-call recovered by the alarm (a safe tool reruns, an unsafe one is reported to the model as interrupted).
  • sockets.test.ts: real WebSockets, joining mid-run, and a socket that outlives eviction.

@changeset-bot

changeset-bot Bot commented Sep 30, 2026 •

Copy link
Copy Markdown

⚠️ No Changeset found

Latest commit: 36dddfb

Merging this PR will not cause a version bump for any packages. If these changes should not result in a new version, you're good to go. If these changes should result in a version bump, you need to add a changeset.

This PR includes no changesets

When changesets are added to this PR, you'll see the packages that this PR includes changesets for and the associated semver types

Click here to learn what changesets are, and how to add one.

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

@agent-think

agent-think Bot commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

🔴 agents import sizes: 1 entry point over threshold

Entry point Exports Largest gzip change Size now
🔴 agents/vite 1 resized +65.2 KiB (+18.42%) 419 KiB
Changed exports (1)
Import Gzip change Size now
🔴 agents/vite#default +65.2 KiB (+18.42%) 419 KiB
How this works

Each runtime export is bundled on its own, minified, and gzipped. Changes smaller than 100 B, or smaller than 1% and 1 KiB, are ignored. Growth over 10% or 5 KiB is marked 🔴. This report is informational and does not fail CI. The workflow artifact contains every measurement.

Compared 040458ed → 5b54e569 · workflow run · reported by agent-think[bot]

@mattzcarey
mattzcarey force-pushed the feat/pi-durable-driver branch from 5b54e56 to 87a8503 Compare September 30, 2026 16:05
@mattzcarey
mattzcarey changed the base branch from main to feat/lifecycle-jobs-sync September 30, 2026 16:05
Move examples/next/harnesses/pi from pi-agent-core's AgentHarness (0.84) to
@earendil-works/pi-durable from pi main (2bbfcca4). PiHarness keeps the
harness interface (prompt, submit, sessions, session(id), webSockets) and is
now a driver runtime: each submission is one driver operation whose step
admits the input into pi by request id and waits for pi to settle it.

- src/driver: the Driver from cloudflare#2396, example-local
- src/harness/session-store.ts: pi's SqliteDatabase on Durable Object SQLite,
  tables under pi_, checked with pi's storage conformance suite
- the wire is pi's own agent events; one reducer for client and tests
- NOTES.md records the Tasks/driver/state-machine decision and what was hard
The harness exposes session.events() and stops owning a transport. The
one-session-per-socket protocol (sockets.ts, protocol.ts) and the UI
reducer (view.ts) are app glue on the harness's public API; the host
re-watches sockets from its own onStart.
submit() admits into pi once, after the session's wake has a job. The
wake's step waits while pi has live tasks in the session and parks when
it has none, so it never admits or replays. The copied driver writes its
row and job atomically with jobs.pushSync (cloudflare#2420).
…ake job

pi owns admission and replay, so the driver reduced to a heartbeat job per
session. PiHarness now owns that job directly (onJob), and the copied driver
is gone. The example no longer needs jobs.pushSync.
@mattzcarey
mattzcarey force-pushed the feat/pi-durable-driver branch from 87a8503 to 36dddfb Compare September 30, 2026 16:48
@mattzcarey mattzcarey changed the title feat(examples): host pi-durable on a driver in the pi harness example feat(examples): host pi-durable in the pi harness example Sep 30, 2026
@mattzcarey
mattzcarey changed the base branch from feat/lifecycle-jobs-sync to main September 30, 2026 16:48
@aron-cf
aron-cf marked this pull request as ready for review October 1, 2026 08:55
@aron-cf
aron-cf merged commit 3c6a1d0 into cloudflare:main Oct 1, 2026
3 of 5 checks passed

@devin-ai-integration devin-ai-integration Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Devin Review found 6 potential issues.

Devin Review

Comment on lines +201 to +204
const sessions = new Set(
inspection.tasks.map((task) => String(task.record.conversationId))
);
for (const session of sessions) await this.#wake(session);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔴 Background subagents lose their wake alarm

When a turn starts a background subagent after startup, onStart never schedules that child's wake. The parent's wake only tracks its own conversation, so the child can stop progressing after eviction.

Learn more

The wake job is the Durable Object alarm that resumes pi's in-memory scheduler after eviction. Startup creates jobs for conversations that already have tasks, but pi can create a background child conversation later. The parent's #wakeStep ignores tasks whose conversation ID differs from its own. When the parent settles, its job completes while the child has no job; after eviction its scheduler cannot resume.

Example: A root turn launches a background child and returns immediately. Startup saw only the root; the root's wake completes, and the child has no alarm even though its task remains live. The child resumes only when another request starts the object.

Recommended fix: Register or refresh a per-conversation wake when pi creates tasks or child conversations, including those created after startup. Alternatively, keep a global wake while any live task exists; test a background child that outlives its parent across eviction.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Comment on lines +430 to +433
const wakeAt = await this.#longWait(pi, conversation.id, BG);
if (wakeAt !== undefined) return { rescheduleAt: wakeAt };
// Background tasks are outside the conversation's idle wait.
if (tasks.every((task) => task.record.background)) return heartbeat;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔴 Long retries stall concurrent background work

When a generation retry is over 60 seconds away, #longWait parks the session's wake until that deadline. Other background tasks in that session lose their 30-second heartbeat and can stall after eviction.

Learn more

A wake job is scheduled per session and is responsible for every live task in that session. #longWait reads only the generation's retry or deferred-poll deadline, while the later background-task check would normally reschedule the job in 30 seconds. If background work overlaps a long generation wait, returning the generation deadline skips that check. An eviction can then leave the background task without an alarm until the generation deadline.

Example: A background task is running in session 1 while its generation retries in five minutes. The wake schedules its next alarm five minutes away; after eviction, the background task does not restart at the usual 30-second check.

Recommended fix: Choose the earliest necessary wake across all live session tasks. Keep the 30-second heartbeat whenever background tasks need polling, even if the generation's own retry can wait longer.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Comment on lines +441 to +444
// The wait runs past this dispatch, inside the alarm's work, so the
// object stays alive for it. The heartbeat covers an eviction.
this.lifecycle.trackAlarmWork(wait);
return heartbeat;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔍 Verify the wake's lifetime guarantee

trackAlarmWork(wait) records a promise for memory-limit accounting; it does not await the promise. The documented claim that this keeps an alarm invocation open needs verification.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

{
"name": "@cloudflare/agents-next-pi-harness-example",
"description": "Experimental: pi AgentHarness composed as an example-local Lifecycle capability",
"description": "Experimental: pi-durable hosted on a Durable Object as an example-local PiHarness on a driver",

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔍 Package description names the obsolete driver

The package description says the harness runs on a driver, but this change uses Lifecycle jobs instead.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Comment on lines +17 to +23
function sessionFromRequest(request: Request, fallback: PiSessionId): string {
const session = new URL(request.url).searchParams.get(SESSION_QUERY);
if (session === null || session === "") return fallback;
if (!/^[1-9][0-9]{0,15}$/.test(session)) {
throw new Error(`Invalid pi session ${JSON.stringify(session)}`);
}
return session;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟥 Socket clients can control other sessions

A client can select any numeric session ID, including root 1, without an ownership check. The socket sends that session's transcript and accepts submit, abort, and reset commands, exposing and modifying other conversations.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Comment on lines +134 to +136
let message: PiClientMessage;
try {
message = JSON.parse(raw) as PiClientMessage;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟨 Socket commands lack runtime validation

Parsed socket JSON is treated as a valid command without checking its fields or size. Malformed or oversized submissions reach pi, allowing clients to trigger errors or excessive storage and model work.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

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.

2 participants