Repository navigation
feat(agents): agents/driver, and an experimental ThinkHarness scored against Think's tests - #2396
mattzcarey wants to merge 7 commits into
Conversation
Extract HarnessDriver and DurableToolRuns from the harness-driver branch into their own agents/driver entrypoint, with Lifecycle pushSync/cancelSync for enqueueing a scope job in the same transaction as its submission. Co-authored-by: Matt Carey <mcarey@cloudflare.com>
The host now owns the driver and passes it to each harness:
readonly driver = new Driver();
readonly harness = new PiHarness({ driver: this.driver });
lifecycle.use(this.driver).use(this.harness)
A harness calls driver.register(id, runtime, { settle, fail }) and keeps
the returned handle (submit, wake, defer, cancel, pending, waitForIdle).
Several runtimes share one driver; their queues are keyed by runtime id.
Also:
- waiting without notBefore parks the scope with no job until wake(),
for open-ended waits such as a human approval
- a wake() during an in-flight drive is no longer lost when that drive
returns waiting
- rename HarnessDriver* to Driver*, driverId to runtimeId, and the
submissions table to cf_agents_driver_submissions
- drop cancelTools (runtime.cancel runs at the same point) and the
pre-release ALTER TABLE column migrations
- docs/agents/driver.md and a changeset
Replace the four-method runtime (inspect, admit, drive, cancel) and its three status vocabularies with one method: step(operation, signal) -> continue | sleep until | park | done plus an optional stop(operation), and an onFail registration hook in place of settle/fail. A step reads the harness's durable records, so the old inspect and admit fold into it, and a finished operation answers done again. The handle is submit, wake, stop, pending and waitForIdle. defer is gone (return sleep), and cancel is renamed stop. Also: - stop() aborts a running step and waits up to 5s for it to unwind before calling the runtime's stop, so the two never run at once - an aborted step's error is not counted toward onFail, and a step's error is never charged to the next operation in the queue - a runtime stop that throws is retried by the loop without stepping - malformed step answers count as a throwing step - the submissions table drops stream_id and uses running/started_at/ stop_requested
Think's turn loop as a Driver runtime. Each chat is one queue of turns; one step of a turn does one thing, chosen from the transcript: call the model once, run one tool call, or settle one tool part. - model steps stream AI SDK UI chunks through an onChunk hook; the tools the model sees have no execute, so the harness runs each tool call itself, one per step - needsApproval tools and client tools park the turn until answer() or resolveTool() - a tool call records that it started, so after an eviction a tool that never finished is reported as interrupted (recovery: never, the default) or run again (recovery: safe) - stop() settles open tool parts and ends the turn as stopped; the next queued turn in the chat runs - model errors end the turn as error instead of being retried Experimental and undocumented until @cloudflare/think's test suite passes on it.
- src/harness/think.ts: an internal, experimental Think that extends Agent and installs Sessions, a Driver and a ThinkHarness. It keeps Think's override hooks (getModel, getSystemPrompt, getTools, beforeTurn, beforeStep, beforeToolCall, afterToolCall, onStepFinish, onChunk, onChatResponse, onChatError), chat(), saveMessages(), messages, and the cf_agent_chat WebSocket protocol. Any other method the real Think has throws 'Think.<name> is not supported by the harness-backed Think yet', so a failing test names what is missing. - src/tests/vitest.harness.config.ts: the same workers suite, with the test agents' ../think imports pointed at src/harness/compat.ts. - scripts/harness-compat.ts (pnpm test:harness): runs it and writes harness-compat.md and harness-compat.json. The denominator is every test the normal suite collects. --check fails when a test that passed in harness-compat.json fails now. Baseline: 378 of 1192 tests pass. The harness also gains beforeStep, onModelChunk and the AI SDK step result on onStepFinish, and ends a turn at once when the host's hooks, model or tools throw, keeping the driver's retries for its own storage failures.
…d Think - Sessions gets Think's reservedMetadataKeys, so messages are written the same way into the same rows - scheduler jobs the previous engine queued for in-flight recovery (_chatRecoveryRetry, _chatRecoveryContinue, _cfRetryMessengerRecoveryDelivery) complete as no-ops instead of failing on retry: in-flight work may die in the move - src/harness/README.md lists the durable state that has to carry over, and where each piece lives
🦋 Changeset detectedLatest commit: 097a2b0 The changes in this PR will be included in the next version bump. This PR includes changesets to release 2 packages
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 |
🟢 agents import sizes: 2 entry points changed, no growth
Changed exports (4)
How this worksEach 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 |
| this.lifecycle.jobs.list(); | ||
| const result = this.lifecycle.storage.transactionSync(() => { | ||
| const enqueued = store.enqueue(scope, id, input); | ||
| this.#pushJobSync(registration.id, scope, Date.now()); |
There was a problem hiding this comment.
🔴 New submissions interrupt sleeping operations
When a scope already has a sleeping head, submit replaces its scheduled job with an immediate one. The head runs before its requested wake time, defeating sleeps and retry backoff.
Learn more
Each scope has one job for its oldest operation. A sleep or failed step schedules that job at a future time. Enqueuing a second operation currently replaces that job with a due-now job, although the new operation cannot run before the head. This also defeats exponential retry delays when submissions arrive during an outage.
Example: Operation A returns sleep until 12:00 and remains the head. At 11:00, operation B is submitted in the same scope. The replaced job fires immediately and steps A at 11:00 instead of 12:00.
Recommended fix: In #submit, only create or move the scope job to now if the enqueued operation becomes the head and no earlier job is already scheduled. Preserve the existing head's job time when adding work behind it, while still recovering a missing job.
Was this helpful? React with 👍 or 👎 to provide feedback.
| const outcome = this.#afterStopFailure(registration, current, error); | ||
| await this.#applyOutcome(registration, current.scope, outcome); |
There was a problem hiding this comment.
🔴 Queued stop failure delays the head
When stopping a queued operation fails, stop applies that operation's backoff to the shared scope job. The unrelated head waits for the queued operation's retry instead of following its own schedule.
Learn more
A scope has one job, representing only its oldest operation. The success path already checks wasHead before updating that job. The failure path does not. A failed stop of a later operation therefore replaces the head's due time with the later operation's backoff; retrying the later stop through the scope job is impossible until it becomes head.
Example: A is the head and due now, while B is queued behind it. stop(B) invokes a stop callback that throws. The shared job moves to the retry time for B, and A does not step when due.
Recommended fix: Apply stop-failure backoff to the scope job only if the stopped operation was the head. Preserve a durable stop request for queued operations and retry their stop when they reach the head, or add an independent retry job for queued stops.
Was this helpful? React with 👍 or 👎 to provide feedback.
| let inspection = await this.#runtime.inspect(runId, run.owner); | ||
| if (inspection.status === "not-started") { | ||
| await this.#runtime.start(runId, run.input, run.owner); | ||
| this.#markRunning(runId); | ||
| run = this.#get(runId) ?? run; | ||
| inspection = await this.#runtime.inspect(runId, run.owner); | ||
| } |
There was a problem hiding this comment.
🔴 Cancelled tool runs can restart or complete
If cancel runs while onJob awaits inspection, onJob continues using the stale active run. It can start the cancelled tool or overwrite its cancelled status with a result.
Learn more
The job reads an active run and then awaits external runtime calls. During those awaits, cancel can call the runtime's cancellation callback, persist cancelled, and remove the job. The resumed job still uses its original run object. It can call start for a cancelled run, and #settle updates rows without requiring an active status, so a terminal inspection can replace cancelled with completed or failed.
Example: inspect pauses with a not-started response pending. Another request cancels the run. When inspect returns, the job calls start despite the persisted cancellation.
Recommended fix: Recheck the persisted run after each awaited runtime call before starting or settling. Make #settle conditional on status IN ('pending', 'running'), and avoid waking/notifying from a stale outcome.
Was this helpful? React with 👍 or 👎 to provide feedback.
| const runs = this.#active().filter( | ||
| (run) => | ||
| run.owner.operationId === operationId && | ||
| run.owner.cancellation === "with-parent" | ||
| ); |
There was a problem hiding this comment.
🟡 Tool cancellation crosses runtime boundaries
When two runtimes use the same operation ID, cancelByOperation cancels attached runs owned by both. The filter ignores owner.runtimeId and owner.scope, so an unrelated operation loses its tools.
Learn more
The driver permits the same operation ID in different runtime registrations, as shown by its multiple-runtime queue test. Tool owners also record the runtime and scope. Matching on operation ID alone conflates two unrelated owners, so a parent cancelling its attached tools cancels another parent's attached tools too.
Example: Runtimes pi and think each submit op-1 and start attached tool runs. Calling cancelByOperation('op-1') for pi cancels think's tool run as well.
Recommended fix: Require a full owner identity, including runtime ID and scope, in cancelByOperation, or accept a DurableToolOwner and match all parent identity fields.
Was this helpful? React with 👍 or 👎 to provide feedback.
| const maxSteps = config.maxSteps ?? this.#options.maxSteps ?? 10; | ||
| if (turn.step >= maxSteps) { | ||
| return this.#finish(context, turn, "completed"); | ||
| } |
There was a problem hiding this comment.
🟡 Turn stop conditions never take effect
When beforeTurn returns stopWhen, the harness still controls termination only through maxSteps. A requested stop condition cannot end the turn, so additional model calls run.
Learn more
The per-turn configuration type exposes stopWhen, but the harness never reads it. Each streamText call is intentionally limited to one model step, leaving the harness responsible for evaluating the caller's condition across calls. The only configurable termination currently checked here is maxSteps.
Example: A caller supplies stopWhen that returns true after the first tool result. The harness calls the model again after settling that tool until it reaches a no-tool response or maxSteps.
Recommended fix: Track the completed step results needed by StopCondition across driver steps and evaluate config.stopWhen at the same boundary as maxSteps, or remove the exposed option until it is supported.
Was this helpful? React with 👍 or 👎 to provide feedback.
| abort(); | ||
| return; | ||
| } | ||
| signal?.addEventListener("abort", abort, { once: true }); |
There was a problem hiding this comment.
🟡 Completed tool waits retain abort listeners
When wait receives a long-lived signal, it adds an abort listener without removing it after settlement. Repeated completed waits retain callbacks and waiter state until that signal aborts or is collected.
Learn more
A pending wait registers an abort callback on its caller's signal. #notify resolves or rejects and removes waiters from the coordinator map, but it never removes those signals' callbacks. The signal retains each callback's closure, including its waiter set, for its lifetime.
Example: A request controller is reused while waiting for 1,000 tool runs. All 1,000 waits settle, but the controller still holds 1,000 abort listeners.
Recommended fix: Store the signal and abort handler with each waiter, remove the handler in #notify and on abort, and avoid adding one if the run settled during waiter registration.
Was this helpful? React with 👍 or 👎 to provide feedback.
|
|
||
| A runtime has one required method, `step`, and an optional `stop`. Pass | ||
| the driver an object that calls your private methods, so `step` and `stop` | ||
| don't become part of the harness's public API, and keep the handle you get |
agents
@cloudflare/ai-chat
@cloudflare/codemode
hono-agents
@cloudflare/shell
@cloudflare/think
@cloudflare/voice
@cloudflare/worker-bundler
commit: |
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
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
|
Closing in favour of #2512, which builds ThinkHarness on the Lifecycle job queue directly (no driver) with Sessions and Streams. The compat scoreboard runner from this PR carried over. |
…cycle capability (#2512) * feat(agents): add experimental ThinkHarness, Think's agent loop as a Lifecycle capability Sessions for transcripts, Streams for in-flight model output, and one Lifecycle wake job per session. The harness runs each server tool call itself, with a per-tool recovery policy. ThinkChat serves a session over Think's useAgentChat WebSocket protocol. * fix(agents): ThinkHarness refuses to steer instead of queueing as a follow-up * test(think): score Think's suite against a Think built on ThinkHarness pnpm test:harness runs Think's workers suite with the test agents' ../think import pointed at an internal Think built on agents/harness/think, and records the score in harness-compat.md (427 of 1225 today). test:harness:check fails when a recorded pass regresses. The compat runner and resolver come from #2396. * refactor(agents): ThinkHarness owns its Streams; tools carry their own recovery The harness creates its in-flight Streams internally (only it reads them), so hosts install Sessions and the harness only. A tool's interrupted-call policy moves from recovery.tools to an optional recovery field on the tool; recovery keeps only the budget (maxAttempts, backoffMs, stallTimeoutMs). * docs(examples): add examples/next/harnesses/think Think's agent loop on a plain Durable Object with ThinkHarness, served to the stock useAgentChat by ThinkChat. Three tools show the three ways a call runs: on the server and rerun after an eviction, on the server after approval, and in the browser. * refactor(agents): ThinkHarness owns its Sessions The harness writes every message and its operations reference them, so it owns the Sessions capability as it does Streams; hosts install the harness only. session.transcript is the Session handle for branches, search, compaction and direct writes, and direct writes reach subscribe() and ThinkChat clients through the Sessions change feed. reservedMetadataKeys passes through to the transcript. * fix(examples): dedupe React in the think harness example The agents workspace link brought a second React instance into the bundle, so useAgent's hooks read a null dispatcher and the page crashed on load. * refactor(agents): branches, search and compact on ThinkSession; rename internal records ThinkSession gets branches(messageId), search(query) and compact(), and the session.transcript getter goes; configureSession still hands over the raw Sessions handle. The internal operation tables are OperationRecords in records.ts, so they do not share a name with agents/harness/store. * chore: ThinkHarness ships as a patch * fix(agents): address review of ThinkHarness - onStart keeps a wake already set for later (recovery backoff, the alarm memory-limit breaker's delay) instead of replacing it with one due now; a sealed breaker settles the running operation as out_of_memory. - Streamed chunks wait at most 100ms in memory before they are written, so output a client saw survives an eviction. - Client input may only add user messages and cannot rewrite a stored message by reusing its id; a client tool cannot replace a server tool. - ThinkChat sends one start and one finish per operation, so a multi-step turn is one message for every client, and answers a STREAM_PENDING probe with STREAM_RESUME_NONE when the work ends without streaming. - Streams rechecks the v1 legacy table before use, so a second Streams instance dropping it cannot break the harness's appends. - Example: one chat per browser, busy while a message is submitted, themed error blocks, and a composition snippet in the README. - Scoreboard: the headline counts only tests that construct Think. * fix(agents): address second ThinkHarness review - Deletions and compactions on a session's transcript emit a new transcript event; ThinkChat re-sends the transcript instead of clearing the chat, and compact() no longer reports itself as a reset. - The change-feed bridge skips the harness's own writes by message id, so a direct write landing during a harness commit still reaches listeners. - A sealed memory-limit breaker only marks running operations; the next drive settles them unanswered (out_of_memory) through the usual path, keeping the partial answer and notifying clients, and never reruns them. - A client tool result answers only a client tool (not_client_tool otherwise), so a server tool's result comes from running it. * fix(agents): migrate ThinkHarness operations tables that predate abandon_reason CREATE TABLE IF NOT EXISTS leaves an existing table as it was, so a sealed memory-limit breaker failed to mark running operations on objects whose table was created before the column.
Adds
agents/driver, a durable queue-and-step loop for agent harnesses, and rebuilds Think's turn loop on it as an experimentalagents/harness/think.@cloudflare/thinknow runs its own test suite against a Think built on that harness and records the score, so Think can move onto it once every test passes.agents/driverOne
Driverper object. The host installs it, and each harness registers a runtime on it and keeps the handle:flowchart LR submit["submit(scope, input)"] --> queue[("FIFO per scope")] queue --> job["one singleflight Lifecycle job per scope"] job --> step["runtime.step(op, signal)"] step -->|continue| job step -->|sleep until| job step -->|park| idle["no job, no alarm"] idle -->|"wake(scope)"| job step -->|done| next["remove, step the next op in the scope"] step -->|throws| retry["backoff, onFail after maxAttempts"]stepdoes one bounded piece of work. It may run again for the same state after an eviction, so it reads the runtime's durable records.trackAlarmWork, with a heartbeat job, so a long model call doesn't block the alarm.stop(id)aborts a running step, waits up to 5s for it to unwind, then calls the runtime'sstop. An aborted step's error never counts towardonFailand is never charged to the next op.wake()that lands while a step is running is kept, even if that step answersparkorsleep.DurableToolRuns, for tool work that outlives one step, comes along unchanged.Based on @aron-cf's
harness-driverbranch. The first commit is his driver code as it was.agents/harness/think(experimental, undocumented)Think's turn loop as a driver runtime. Each chat is one queue of turns, and each step reads the transcript and does exactly one thing:
flowchart TD s["step(turn)"] --> t{"this turn's assistant message"} t -->|"none yet, or every tool call has a result"| m["call the model once"] t -->|"approval-requested, or a client tool"| p["park"] t -->|"approved, or input-available"| r["run one tool call"] t -->|"denied"| d["tool-output-denied"] t -->|"last model step made no tool calls"| f["end the turn"]execute, so the harness runs each call itself. That allows a recovery policy per tool. If the object was evicted during a call, the call is reported to the model as interrupted (recovery: "never", the default) or run again (recovery: "safe").beforeTurn,beforeStep,beforeToolCall,afterToolCall,onChunk(UI chunks),onModelChunk,onStepFinish,onTurnEnd.Scoring Think against it
packages/think/src/harness/think.tsis an internalThink extends Agentthat installs Sessions, a Driver and a ThinkHarness. It keeps Think's override hooks,chat(),saveMessages(),messages, and thecf_agent_chatWebSocket protocol. Every other method the real Think has throwsThink.<name> is not supported by the harness-backed Think yet, so each failing test names what is missing.pnpm test:harnessruns the unchanged workers suite with the test agents'../thinkimport pointed at that class. It writesharness-compat.mdandharness-compat.json. The denominator is every test the normal suite collects.pnpm test:harness:checkfails when a test recorded as passing now fails.Today: 378 of 1196 pass. The largest gaps are submissions (94 tests),
runTurn(54), agent tools (about 40), actions (about 27) and chat recovery internals (25).Storage
Sessionscapability, session and tables as Think. Agent schedules, state and MCP servers belong to the Agent base class, which is unchanged.think_config, workspace files, messenger subscriptions) is listed inpackages/think/src/harness/README.md. Each will be ported onto the tables Think already wrote.Not in this PR
DurableToolRunsto thestepstyle.test:harness:checkin CI.