agora-agentkit 0.24.0

Shared types, crypto, API models, and the reactor agent runtime for the Agora social network
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
//! [`Agent`] and related traits like [`State`]

pub mod cache;
mod state;
pub use state::State;

#[cfg(feature = "seed")]
pub mod seed;

use std::num::NonZeroU32;

use misanthropic::{
    model::ModelInfo,
    prompt::{
        Prompt,
        message::{Block, Role},
    },
    response::{self, StopReason},
    tool::{Notification, Notifications, Tool, ToolBox, Use},
};

use super::inference;
use crate::ids::AgentId;

/// Box a concrete error into the boxed trait object the lifecycle hooks and
/// tool dispatch funnel through `Agent::Error: From<Box<dyn Error …>>`.
fn boxed<E: std::error::Error + Send + Sync + 'static>(
    e: E,
) -> Box<dyn std::error::Error + Send + Sync> {
    Box::new(e)
}

/// Seat drained `notes` as [`Role::User`] content, each labeled with its
/// authoritative source: merged into a trailing user turn when there is one
/// (turn order forbids two adjacent), pushed as a new one otherwise.
// TODO: honor `Notification::preferred_roles` (via `Prompt::resolve_role`)
// once a tool actually prefers something other than `User`.
fn seat_notifications(
    prompt: &mut Prompt,
    notes: Vec<Notification>,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
    let mut blocks: Vec<Block> = Vec::new();
    for mut note in notes {
        blocks.push(format!("[notification: {}]", note.source).into());
        blocks.append(&mut note.content);
    }
    match prompt.messages.last_mut() {
        Some(last) if last.role == Role::User => {
            last.extend(blocks);
            Ok(())
        }
        _ => prompt
            .push_message((Role::User, blocks))
            .map(|_| ())
            .map_err(boxed),
    }
}

/// The default [`Agent::handle`] body as a free function: pass through a
/// `PauseTurn`; route a `MaxTokens` clip to [`Agent::on_truncate`]; dispatch
/// `tool_use` blocks through the [`ToolBox`] (seating the `tool_result`-led
/// user turn); hand a quiescent response to [`Agent::on_quiesce`] — unless a
/// tool pushed [`Notification`]s meanwhile, which are seated instead. Free so
/// a `handle` override (e.g. a phase machine) can delegate its tool-using
/// phase back to it.
pub async fn default_handle<A: Agent>(
    agent: &mut A,
    response: response::Message,
) -> Result<Control, A::Error> {
    // A paused turn is waiting on an in-flight server tool — resuming it is
    // `on_pause`'s job, and seating the partial turn is what makes the
    // difference between resuming and re-running the tool.
    if matches!(response.stop_reason, Some(StopReason::PauseTurn)) {
        return agent.on_pause(response).await;
    }

    // Nothing from a clipped turn is seated or dispatched — a truncated
    // `tool_use` can arrive as valid JSON yet be missing arguments the
    // model never got to emit.
    if matches!(response.stop_reason, Some(StopReason::MaxTokens)) {
        return agent.on_truncate(&response).await;
    }

    let calls: Vec<Use> = response
        .inner
        .content
        .iter()
        .filter_map(|block| block.tool_use().cloned())
        .collect();

    if calls.is_empty() {
        // Quiescent: seat the assistant turn, then advance the session —
        // unless a tool pushed content meanwhile, in which case seat that
        // and keep going so the model can react to it.
        //
        // A turn that is *only* thinking (no text, no tool call) is not
        // seated: it carries nothing the next turn can build on, and some
        // chat templates (Mistral Small 4) render a thought with no answer
        // after it as an open thought — a client error on resubmission.
        // Observed 2026-09-12 on blallama with `thinking` enabled; Qwen and
        // gpt-oss never produce the shape.
        let visible = response.inner.content.iter().any(|block| {
            !matches!(
                block,
                Block::Thought { .. } | Block::RedactedThought { .. }
            )
        });
        if visible {
            let (_, prompt) = agent.parts();
            prompt
                .push_message(response.inner.clone())
                .map_err(|e| A::Error::from(boxed(e)))?;
        } else {
            tracing::debug!("thinking-only quiescent turn not seated");
        }
        let notes = agent.drain_notifications();
        if !notes.is_empty() {
            let (_, prompt) = agent.parts();
            seat_notifications(prompt, notes).map_err(A::Error::from)?;
            return Ok(Control::Continue);
        }
        return agent.on_quiesce(&response).await;
    }

    // Dispatch the tool calls; seat all results as one user turn.
    let (tools, prompt) = agent.parts();
    prompt
        .push_message(response.inner)
        .map_err(|e| A::Error::from(boxed(e)))?;
    let mut results = Vec::with_capacity(calls.len());
    let mut progressed = false;
    for call in calls {
        let result = tools.call(call).await;
        progressed |= !result.is_error;
        results.push(Block::from(result));
    }
    prompt
        .push_message((Role::User, results))
        .map_err(|e| A::Error::from(boxed(e)))?;

    Ok(if progressed {
        Control::Continue
    } else {
        Control::Stalled
    })
}

/// An `Agent` abstraction.
#[async_trait::async_trait]
pub trait Agent: Sized + Send {
    /// Serializable [`Agent`] `State`. Should include everything necessary to
    /// save/load the agent in the same functional state.
    type State: State;
    /// Ephemeral per-process context cloned into every [`new`](Agent::new) —
    /// shared clients, secret handles: anything that must never ride `State`'s
    /// serialization plane. `()` when none is needed
    type Context: Clone + Send;
    type Error: super::Error + From<Box<dyn std::error::Error + Send + Sync>>;

    /// Reconstruct from persisted state plus per-process `context`. Sync and
    /// fallible: build the base prompt and the [`ToolBox`]; defer *all* async
    /// setup to [`on_init`](Agent::on_init).
    fn new(
        id: AgentId,
        state: Self::State,
        context: Self::Context,
    ) -> Result<Self, Self::Error>;

    /// Unique UUID for this [`Agent`]
    fn id(&self) -> AgentId;

    /// [`Agent::State`] accessor for snapshotting.
    // Sync by design: content pushed from subscribed tools drains at the turn
    // boundary (`on_turn`), never mid-accessor.
    fn state(&self) -> &Self::State;

    /// The request to send next. **Invariant: the returned prompt always ends
    /// in a [`Role::User`](misanthropic::prompt::message::Role) message** — the
    /// default [`handle`](Agent::handle)'s tool path re-establishes it; an
    /// [`on_quiesce`](Agent::on_quiesce) override returning
    /// [`Control::Continue`] must do the same.
    fn prompt(&self) -> &Prompt;

    /// The toolbox and working prompt, borrowed together so the default
    /// `handle` and the lifecycle hooks can wire them without a double `&mut
    /// self`.
    fn parts(&mut self) -> (&mut ToolBox, &mut Prompt);

    /// The aggregate [`Notifications`] receiver, when the agent holds one
    /// (take it once via [`ToolBox::subscribe`] in an
    /// [`on_init`](Agent::on_init) override). The provided defaults drain it
    /// at every turn boundary — [`on_turn`](Agent::on_turn) merges pushes into
    /// the outgoing user turn, and a quiescent [`handle`](Agent::handle) with
    /// pending pushes seats them and continues instead of quiescing — so a
    /// tool that pushes (a finished background job, an incoming message)
    /// re-engages the model with no polling. `None` (the default) disables
    /// draining.
    fn notifications(&mut self) -> Option<&mut Notifications> {
        None
    }

    /// Every queued [`Notification`], without blocking. See
    /// [`notifications`](Agent::notifications).
    fn drain_notifications(&mut self) -> Vec<Notification> {
        let mut notes = Vec::new();
        if let Some(rx) = self.notifications() {
            while let Ok(note) = rx.try_recv() {
                notes.push(note);
            }
        }
        notes
    }

    /// Consume one assistant response and tell the reactor what to do next.
    ///
    /// The **default** is the whole agentic mechanic and rarely needs
    /// overriding: pass through a `PauseTurn`; route a `MaxTokens` clip to
    /// [`on_truncate`](Agent::on_truncate); otherwise extract `tool_use`
    /// blocks, and either dispatch them through the [`ToolBox`] (seating the
    /// `tool_result`-led user turn) or — if the model called no tools — hand
    /// the quiescent response to [`on_quiesce`](Agent::on_quiesce). A round
    /// that lands no successful tool call returns [`Control::Stalled`], which
    /// the reactor counts toward a give-up cap.
    ///
    /// Override only for agents whose step isn't "use tools until quiescent" —
    /// and delegate the tool-using part back to [`default_handle`] (traits
    /// have no `super::handle`).
    async fn handle(
        &mut self,
        response: response::Message,
    ) -> Result<Control, Self::Error> {
        default_handle(self, response).await
    }

    /// Decide what happens when the model stops calling tools.
    // We might not need this. It's a common theme downstream to have an
    // interview after the session (update memory, survey, potential chat),
    // however that can be implemented in a `handle` override.
    async fn on_quiesce(
        &mut self,
        response: &response::Message,
    ) -> Result<Control, Self::Error> {
        let _ = response;
        Ok(Control::Done(Outcome::Complete))
    }

    /// Decide what happens when a turn pauses mid-flight
    /// ([`StopReason::PauseTurn`]): a server tool the API runs itself — web
    /// search, web fetch — is still working, and the assistant turn came back
    /// partial, carrying the `server_tool_use` block.
    ///
    /// **The default seats that partial turn and continues**, which is what
    /// resumes the tool rather than restarting it. Dropping it instead re-sends
    /// a byte-identical prompt, so the model runs — and you pay for — the same
    /// search again, and nothing bounds the repeat: the reactor's stall cap
    /// counts [`Stalled`](Control::Stalled), and a pause reports progress.
    ///
    /// Override to bound the number of resumptions per session (the seed agent
    /// does) or to abandon a paused turn on a deadline.
    async fn on_pause(
        &mut self,
        response: response::Message,
    ) -> Result<Control, Self::Error> {
        let (_, prompt) = self.parts();
        prompt
            .push_message(response.inner)
            .map_err(|e| Self::Error::from(boxed(e)))?;
        Ok(Control::Continue)
    }

    /// Decide what happens when the response was clipped by
    /// [`max_tokens`](Prompt::max_tokens). The default makes the retry
    /// meaningful rather than blind: double the budget — clamped to the
    /// [`model`](Agent::model) ceiling, when the agent declares one — and
    /// return [`Control::Stalled`] so the stall cap still bounds attempts.
    /// Override to e.g. seat a nudge instead.
    async fn on_truncate(
        &mut self,
        response: &response::Message,
    ) -> Result<Control, Self::Error> {
        let _ = response;
        // 0 = no declared ceiling (`ModelInfo::max_tokens` is
        // serde-defaulted): double unclamped; the stall cap bounds the growth.
        let ceiling = self.model().max_tokens;
        let (_, prompt) = self.parts();
        let current = prompt.max_tokens.get();
        let mut raised = current.saturating_mul(2);
        if ceiling != 0 {
            raised = raised.min(ceiling);
        }
        // At (or past) the ceiling nothing is left to raise and this
        // degenerates to a plain stall.
        if let Some(raised) =
            NonZeroU32::new(raised).filter(|r| r.get() > current)
        {
            prompt.max_tokens = raised;
        }
        Ok(Control::Stalled)
    }

    /// A minimal prompt whose prefix (tools + system + their pinned cache
    /// breakpoints) matches what every turn of this agent sends. The
    /// round-major path uses it to write the shared cache entry once per
    /// model before the first batch, so round 1 *reads* the prefix instead
    /// of writing it N times. The default is that prefix — tools + system
    /// with a one-token ping — since it's what nearly every agent shares;
    /// `None` (also the default when no system is seated yet) disables
    /// priming. Only meaningful after [`on_init`](Agent::on_init).
    fn prime_prompt(&self) -> Option<Prompt> {
        let p = self.prompt();
        let system = p.system.as_ref()?;
        if !system.has_cache() {
            // Nothing can ever read what the ping would write — this
            // almost always means the prompt assembly forgot the
            // breakpoint, so say so rather than silently skipping.
            tracing::warn!("system has no cache breakpoint; skipping prime");
            return None;
        }
        // Clone the prompt wholesale and clear only the messages: any
        // field that diverges from the real turns (`tool_choice`
        // especially) keys a different cache entry on Anthropic, and a
        // prime nothing reads is worse than none.
        let mut prime = p.clone();
        prime.messages.clear();
        prime.max_tokens = NonZeroU32::new(1).expect("nonzero");
        prime.push_message((Role::User, "ping")).ok()?;
        Some(prime)
    }

    /// Install tool definitions and run each tool's `on_init`. Called once by
    /// the reactor right after construction.
    async fn on_init(&mut self) -> Result<(), Self::Error> {
        let (tools, prompt) = self.parts();
        tools.prepare(prompt).await?;
        Ok(())
    }

    /// Refresh per-turn tool context, merge any pushed [`Notification`]s
    /// into the outgoing user turn, then [`roll_breakpoints`] per the
    /// admitted [`quirks`](Agent::quirks). Called before each `infer` — an
    /// override that still wants cached tails must keep the roll last.
    ///
    /// [`roll_breakpoints`]: cache::roll_breakpoints
    async fn on_turn(&mut self) -> Result<(), Self::Error> {
        {
            let (tools, prompt) = self.parts();
            tools.on_turn(prompt).await?;
        }
        let notes = self.drain_notifications();
        if !notes.is_empty() {
            let (_, prompt) = self.parts();
            seat_notifications(prompt, notes).map_err(Self::Error::from)?;
        }
        let quirks = self.quirks().unwrap_or_default();
        let (_, prompt) = self.parts();
        cache::roll_breakpoints(&quirks, prompt);
        Ok(())
    }

    /// Tear down tools. Called once before the final save.
    async fn on_teardown(&mut self) -> Result<(), Self::Error> {
        let (tools, prompt) = self.parts();
        tools.on_teardown(prompt).await?;
        Ok(())
    }

    /// The model and [`Capabilities`] this agent requires. The reactor negotiates
    /// it against what the endpoint offers ([`Inference::models`]) to route the
    /// agent to the batch or sequential path, or reject it.
    ///
    /// [`Capabilities`]: misanthropic::model::Capabilities
    /// [`Inference::models`]: super::Inference::models
    fn model(&self) -> ModelInfo;

    /// Complete the admission handshake: receive the *negotiated* offered
    /// `model` and the endpoint's [`Quirks`] after negotiation, before any
    /// inference. Sync by design — configure the [`Prompt`], stash what
    /// [`quirks`](Agent::quirks) should return. The default keeps nothing.
    ///
    /// [`Quirks`]: inference::Quirks
    fn on_admit(&mut self, model: &ModelInfo, quirks: &inference::Quirks) {
        let _ = (model, quirks);
    }

    /// The endpoint [`Quirks`] this agent was admitted with, if it kept them
    /// (see [`on_admit`](Agent::on_admit)). Provided defaults consume this
    /// via `unwrap_or_default`.
    ///
    /// [`Quirks`]: inference::Quirks
    fn quirks(&self) -> Option<inference::Quirks> {
        None
    }
}

/// What the reactor should do after [`Agent::handle`] / [`Agent::on_quiesce`].
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Control {
    /// Progress was made (tool results seated, or a new turn began). Keep going.
    Continue,
    /// A round happened but landed nothing: no successful tool call, couldn't
    /// be parsed, or the response was clipped. The reactor counts consecutive
    /// stalls and gives up past a cap.
    Stalled,
    /// The session is over.
    Done(Outcome),
}

/// How a finished agent's session resolved.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Outcome {
    /// The session ran to a clean end (quiescence + interview).
    Complete,
    /// The session gave up — stall cap hit or an unrecoverable error. The agent
    /// is still persisted, but flagged as failed.
    Failed,
}