Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion book/src/ch02-what-is-an-agent-loop.md
Original file line number Diff line number Diff line change
Expand Up @@ -160,7 +160,7 @@ For loss-free transcript reconstruction (persistence, replication, audit), regis
- **`LoopObserver`** sees a stream of `AgentEvent`s. Content arrives as deltas — partial fragments that don't carry their parent-`Item` identity — interleaved with lifecycle and telemetry events. Useful for UIs and logging, but a consumer cannot reassemble the canonical transcript from this stream alone.
- **`TranscriptObserver`** fires exactly once per `Item` appended, receiving a session-addressed `TranscriptEvent` with the fully-formed `Item` ready to persist. Calls happen synchronously from the driver, in transcript order, at the single mutation point that owns the transcript — so what the observer sees is what the loop will send to the model on the next turn.

Mutator-driven rewrites (compaction, redaction, repair) do not fire `on_transcript_event`; they are signalled separately by `AgentEvent::MutationFinished`, which a persistence layer can use to snapshot the post-mutation state.
Mutator-driven rewrites (compaction, redaction, repair) replace history rather than appending to it, so they do not fire `on_transcript_event`. They fire the trait's other required method, `on_transcript_rewrite`, with the complete canonical transcript. Between the two, a persistence consumer cannot diverge from the driver.

## The three-layer model

Expand Down
8 changes: 7 additions & 1 deletion book/src/ch06-driving-the-loop.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ pub struct LoopDriver<S: ModelSession> {
mutators: Vec<Arc<dyn LoopMutator>>,
observers: Vec<Arc<dyn LoopObserver>>,
transcript_observers: Vec<Arc<dyn TranscriptObserver>>,
deliveries: Vec<Arc<dyn NativeDelivery>>,
transcript: Vec<Item>,
pending_input: Vec<Item>,
pending_approvals: BTreeMap<ToolCallId, PendingApprovalToolCall>,
Expand All @@ -36,6 +37,7 @@ impl<S: ModelSession> LoopDriver<S> {
-> Result<(), LoopError>;
pub fn set_next_turn_cache(&mut self, cache: PromptCacheRequest) -> Result<(), LoopError>;
pub fn snapshot(&self) -> LoopSnapshot;
pub fn take_delivery_errors(&mut self) -> DeliveryFailures;
}
```

Expand Down Expand Up @@ -282,7 +284,9 @@ The full event taxonomy (`AgentEvent` is `#[non_exhaustive]` — keep a wildcard

Observers are called inline, synchronously, in registration order. The loop task blocks briefly for each observer call. This is acceptable because observers should be fast — write to stderr, increment a counter, append to a buffer. Expensive processing should happen asynchronously behind a channel adapter.

For loss-free transcript reconstruction (persistence, replication, audit), the driver also fans out to a separate `TranscriptObserver` channel that fires once per `Item` appended — a session-addressed `TranscriptEvent` carrying the item, in transcript order. `LoopObserver` alone is not sufficient for this — content deltas span partial parts and historically tool results were appended without an event at all. Mutator-driven rewrites do **not** fire `on_transcript_event`; those are signaled by `AgentEvent::MutationFinished`. Register via `AgentBuilder::transcript_observer`.
For loss-free transcript reconstruction (persistence, replication, audit), the driver also fans out to a separate `TranscriptObserver` channel that fires once per `Item` appended — a session-addressed `TranscriptEvent` carrying the item, in transcript order. `LoopObserver` alone is not sufficient for this — content deltas span partial parts and historically tool results were appended without an event at all. Mutator-driven rewrites do **not** fire `on_transcript_event`; they fire the trait's other required method, `on_transcript_rewrite`, with the complete canonical transcript. Register via `AgentBuilder::transcript_observer`.

Observers are synchronous and infallible, which is the right shape for telemetry but the wrong one for a host that has to *await* its own delivery of a fact at the moment it happens. That is what `AgentBuilder::delivery` is for: a `NativeDelivery` target is awaited at the emission site of each `NativeFact` — streaming progress as the model produces it, the read-only pre-commit `BeforeFinish`, and the post-commit `TurnFinished`. Delivery is read-only and cannot drive the loop, and a `DeliveryError` is diagnostics only: it never fails the operation, never rolls a commit back and never produces a second terminal event. Drain the diagnostics with `LoopDriver::take_delivery_errors`.

## Building the agent

Expand All @@ -299,6 +303,7 @@ let agent = Agent::builder()
.compaction(config) // default: none
.observer(reporter) // default: none
.transcript_observer(persistence) // default: none
.delivery(awaited_target) // default: none
.transcript(vec![system_item]) // default: empty
.input(vec![first_user_turn]) // default: empty (one-shot opener)
.build()?;
Expand All @@ -318,6 +323,7 @@ The builder validates that a model adapter is set. Everything else has sensible
| `compaction` | `None` | Transcript grows without bounds |
| `observers` | `[]` | No event reporting |
| `transcript_observers` | `[]` | No transcript persistence hook |
| `deliveries` | `[]` | No awaited fact delivery |

`Agent::start()` consumes the agent and returns a `LoopDriver` with the supplied transcript loaded passively. The first call to `next()` yields `AwaitingInput`; the host supplies the first user turn via `InputRequest::submit`, and the driver dispatches the model on the next `next()`. The agent's immutable configuration (adapter, tool sources, permissions) is moved into the driver. Multiple drivers can be created from the same `Agent` type by cloning it first.

Expand Down
16 changes: 14 additions & 2 deletions book/src/ch16-compaction.md
Original file line number Diff line number Diff line change
Expand Up @@ -261,19 +261,31 @@ This is why caching is configured separately from compaction in agentkit. Compac

## Loop integration

Compactors register as `LoopMutator`s. The loop runs every registered mutator at each `MutationPoint` — `AfterToolResult` (between tool results and the next inference call) and `AfterTurnEnded` (after the assistant final, interrupt, or cancellation). The trigger decides which points are relevant.
Compactors register as `LoopMutator`s. The loop runs every registered mutator at each `MutationPoint`:

- `TurnStarted` — once when the driver creates a logical turn, after queued input is appended and before any other point or inference. This is the only point that also runs for turns that never dispatch an inference.
- `AfterToolResult` — between tool results and the next inference call.
- `AfterTurnEnded` — before the first inference call of a new turn.

The trigger decides which points are relevant. A trigger that ignores `point` (such as `item_count_trigger`) now also gets the `TurnStarted` opportunity, so a transcript over the threshold compacts a little earlier in the turn. Filter on `point` if that matters.

The chain is transactional. Mutators edit a candidate copy of the transcript; the loop validates it, re-checks cancellation, and only then assigns it to the live transcript in one synchronous step. A mutator that errors, produces a protocol-invalid transcript, cancels, or whose future is dropped leaves the live transcript exactly as it was.

When a compactor fires:

1. The compactor emits `AgentEvent::MutationStarted { mutator, point, .. }` with a stable label it chose
2. The strategy pipeline transforms the transcript through the cursor
3. The loop validates transcript invariants (tool_use ↔ tool_result pairing) and hard-fails with `LoopError::Mutator` on a protocol violation
4. The compactor emits `AgentEvent::MutationFinished { mutator, dirty, metadata, .. }`
5. If the committed transcript actually differs from the live one, the loop delivers one `TranscriptObserver::on_transcript_rewrite` with the complete canonical transcript

```text
Turn lifecycle with a registered compactor:

next() → merge pending input
next() → open turn, merge pending input
│
▼
run mutators at TurnStarted
│
▼
begin model turn
Expand Down
2 changes: 1 addition & 1 deletion book/src/ch19-reporting.md
Original file line number Diff line number Diff line change
Expand Up @@ -269,7 +269,7 @@ impl LoopObserver for AuditLogger {
| Usage | `UsageUpdated` |
| Diagnostic | `Warning` |

For loss-free transcript reconstruction, register a `TranscriptObserver` alongside `LoopObserver`. It fires a session-addressed `TranscriptEvent` once per `Item` appended, in transcript order — including the synthetic placeholder and the eventual real result for background-detached tools, correlated by `call_id` through the matching `ToolExecutionProgress` (placeholder) and `ToolResultReceived` (terminal result) events.
For loss-free transcript reconstruction, register a `TranscriptObserver` alongside `LoopObserver`. It fires a session-addressed `TranscriptEvent` once per `Item` appended, in transcript order — including the synthetic placeholder and the eventual real result for background-detached tools, correlated by `call_id` through the matching `ToolExecutionProgress` (placeholder) and `ToolResultReceived` (terminal result) events. Its other required method, `on_transcript_rewrite`, carries the complete canonical transcript when a mutator rewrote history instead of appending to it.

### Event timeline for a typical turn

Expand Down
33 changes: 24 additions & 9 deletions book/src/session-persistence.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,21 +6,24 @@ The agent loop has no built-in storage backend. Persistence is intentionally a h

| Primitive | Purpose |
| ------------------------------------------------ | ----------------------------------------------------------------------------------------------------------------------- |
| `AgentBuilder::transcript(items)` | Restore prior transcript before the loop starts. |
| `TranscriptObserver::on_transcript_event(event)` | Mirror every newly-appended item to durable storage as the loop runs (`TranscriptEvent` carries `session_id` + `item`). |
| `LoopDriver::snapshot() -> LoopSnapshot` | Read-only point-in-time view of `transcript` and `pending_input` for ad-hoc dumps, audit, or full-state checkpoints. |
| `AgentBuilder::transcript(items)` | Restore prior transcript before the loop starts. |
| `TranscriptObserver::on_transcript_event(event)` | Mirror every newly-appended item to durable storage as the loop runs (`TranscriptEvent` carries `session_id` + `item`). |
| `TranscriptObserver::on_transcript_rewrite(event)` | Replace the stored transcript when a mutator rewrote history (`TranscriptRewriteEvent` carries `session_id` + `items`). |
| `LoopDriver::snapshot() -> LoopSnapshot` | Read-only point-in-time view of `transcript` and `pending_input` for ad-hoc dumps, audit, or full-state checkpoints. |

That is the whole protocol. Any storage backend — in-memory map, sqlite, Postgres, S3, Redis — implements the same shape:

1. **On startup**: load the prior `Vec<Item>` for the session id (or empty for a fresh session) and pass it to `AgentBuilder::transcript`.
2. **During the run**: register a `TranscriptObserver` that appends each `Item` to durable storage.
2. **During the run**: register a `TranscriptObserver` that appends each `Item` to durable storage and replaces it on rewrite.
3. **On shutdown** (graceful or not): nothing required — the observer has already persisted every appended item.

## Two important guarantees

**Append-only ordering.** `on_transcript_event` is called synchronously by the loop, in the exact order items land in the transcript. The observer is the single mutation point — every push to the transcript funnels through it. This means a strictly monotonic `seq` column on a sqlite `items` table reproduces the transcript byte-for-byte on reload.

**Mutators are out-of-band.** Mutator-driven transcript rewrites (compaction, redaction, repair) do **not** fire `on_transcript_event`. They are signalled via `AgentEvent::MutationFinished { dirty: true, .. }`, observable through a `LoopObserver`. A mutation-aware persistor subscribes to both channels and replaces the stored transcript when it sees a dirty mutation finish. An agent without mutators (most coding agents that rely on the provider's prompt cache plus a long context window) can ignore this.
**Mutator rewrites come through the second method.** Mutator-driven transcript rewrites (compaction, redaction, repair) replace history rather than appending to it, so they do **not** fire `on_transcript_event`. They fire `on_transcript_rewrite` with the complete canonical transcript — once per mutation point whose mutator chain produced a transcript that actually differs from the live one. Both methods are required, so an append-only consumer has to decide explicitly (a one-line no-op body) rather than silently diverging from the driver after a compaction pass.

**Best-effort post-commit observation, not a commit gate.** Both methods are synchronous and infallible, and the loop has already committed by the time they run. A failed write cannot roll the turn back; a host that needs acknowledged durable publication owns that linearization point itself.

## A complete sqlite implementation

Expand Down Expand Up @@ -72,6 +75,18 @@ let agent = Agent::builder()
.build()?;
```

```rust,ignore
impl TranscriptObserver for SqliteTranscriptObserver {
fn on_transcript_event(&self, event: TranscriptEvent<'_>) {
// append one item
}

fn on_transcript_rewrite(&self, event: TranscriptRewriteEvent<'_>) {
// replace the stored rows with event.items
}
}
```

That is the entire round-trip. Run the example twice with the same `--session` flag and the second run resumes mid-conversation — the first `next()` call returns `AwaitingInput` because the transcript is loaded but no input is queued, and the host supplies the next user message in response.

## Choosing a backend
Expand All @@ -93,10 +108,10 @@ The integration test crate exercises the round-trip pattern internally; see `cra

## Mutation-aware persistence

If your agent registers any `LoopMutator`s (compaction, redaction, repair), the persistence flow extends:
If your agent registers any `LoopMutator`s (compaction, redaction, repair), the second observer method carries it:

1. `TranscriptObserver::on_transcript_event` continues to mirror new items as they arrive.
2. A `LoopObserver` subscribes to `AgentEvent::MutationFinished { dirty: true, .. }` and uses it as a signal to replace the stored transcript.
3. After a dirty mutation finish, call `LoopDriver::snapshot()` from the host's main task and replace the persisted transcript with `snapshot.transcript`. Subsequent `on_transcript_event` calls resume appending from the new tail.
2. `TranscriptObserver::on_transcript_rewrite` delivers the complete canonical transcript after a mutator chain committed a change. Replace the stored rows with `event.items`; subsequent `on_transcript_event` calls resume appending from the new tail.
3. A mutation that errored, produced a protocol-invalid transcript, was cancelled, or whose future was dropped leaves the live transcript untouched and notifies nothing — there is no partial rewrite to reconcile. A mutator that writes the same value back is not a change and notifies nothing either.

The two channels exist precisely so persistence can stay simple in the no-mutator case (one observer) without sacrificing correctness when mutators are wired in (one observer plus one event listener).
`AgentEvent::MutationStarted` / `MutationFinished` remain the mutator's own telemetry (which mutator ran, why, how much it replaced). They are not the persistence signal: a mutator chooses its own `dirty` label, while `on_transcript_rewrite` fires on the driver's own value comparison.
30 changes: 30 additions & 0 deletions crates/agentkit-loop/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,36 @@ This crate provides:

Use it as the central coordinator between model providers, tool execution, and application UI or control flow.

## Interception hooks

Besides `mutate`, a `LoopMutator` registered with `AgentBuilder::mutator` can intercept values before the loop consumes or commits them:

- `on_session_start` edits session options before the model adapter starts the session.
- `on_model_request` edits a single inference request (transcript, tools, cache, metadata). These edits are not persisted.
- `on_model_response` edits complete model output before it is committed, returned, or sent to tools. Content may change — including tool-call arguments, which reach the executor. Item identity, accounting and tool-call linkage may not. `payload.disposition` says whether the loop will dispatch this response's tool calls (`ContinueWithTools`) or take the normal finish branch (`FinishTurnCandidate`), computed from the loop's own branch predicate rather than the response content.

All hooks default to no-ops and run in registration order. Read-only notifications belong in `LoopObserver`; individual tool interception belongs at `ToolExecutor`.

## Transcript mutation

`mutate` runs at every `MutationPoint`: `TurnStarted` (once per logical turn the driver creates, with queued input already appended, and the only point that also runs for turns that never dispatch an inference), `AfterToolResult`, and `AfterTurnEnded`.

The chain is transactional. Mutators edit a candidate copy; the loop validates invariants, re-checks cancellation, then assigns the result to the live transcript in one synchronous step. A mutator that errors, produces a protocol-invalid transcript, cancels, or whose future is dropped leaves the live transcript untouched. A committed change is published to every `TranscriptObserver` as one `on_transcript_rewrite` carrying the complete canonical transcript; writing the same value back is not a change and publishes nothing.

## Awaited delivery

`LoopObserver` is synchronous and infallible. A host that must *await* its own delivery of a fact at the moment it happens registers a `NativeDelivery` with `AgentBuilder::delivery`. The driver awaits each target at the fact's emission site:

- `NativeFact::Progress` — model deltas, usage, tool calls, attempt supersession, and the loop-authored background-detach placeholder, delivered as the driver consumes them rather than buffered until `next()` returns.
- `NativeFact::BeforeFinish` — the logical turn is about to finish, before its own terminal items are appended, so a consumer sees terminal output and cancellation partials before they commit.
- `NativeFact::TurnFinished` — the turn finished and its items are committed.

The terminal pair is delivered exactly once per logical turn — including turns that end through cancellation, an error, a failed cleanup, or `retire_interrupted_turn` — for uninterrupted calls and for cooperative cancellation followed by retirement. A hard abort breaks that: `BeforeFinish` is awaited before anything commits, so dropping the `next()` future inside it loses the turn's terminal output candidate and leaves the turn active, and a later `retire_interrupted_turn` delivers a second prefinish for the same turn with a cancelled result. Treat prefinish as at-least-once if you drop driver futures.

`HookCtx::cancellation` is `None` for the terminal facts, and also `None` for `Progress` when the agent was built without `AgentBuilder::cancellation` — match on the `NativeFact` variant rather than on the handle's presence.

Delivery is read-only; a `DeliveryError` is diagnostics only — it never fails the operation, never stops the stream, never rolls a commit back, never replays and never produces a second terminal event, and later facts for the same turn still arrive. Drain it with `LoopDriver::take_delivery_errors`, which resets both the retained failures and the `dropped` count of what the bounded buffer discarded.

## Quick start

```rust,no_run
Expand Down
Loading
Loading