agora-agentkit 0.56.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
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
//! Quirk-aware rolling cache breakpoints for the default agent mechanics.
//!
//! See [`roll_breakpoints`].

use misanthropic::prompt::{
    Prompt,
    index::{BlockIndex, Index, IndexMut},
    message::CacheControl,
};

use crate::reactor::inference::Quirks;

/// The Anthropic API's hard `cache_control` marker limit per request.
const MAX_CACHE_CONTROLS_PER_REQUEST: usize = 4;

/// Rolling markers kept in the trailing window. 2 fits under the 4-marker
/// budget alongside the two pinned prefix markers (tools+system, intro), and
/// the older of the pair anchors a direct cache hit when a tool-heavy round
/// pushes the newer one past the API's 20-block lookback window.
pub const ROLL_WINDOW: usize = 2;

/// Place [`ROLL_WINDOW`] rolling 1h cache breakpoints on the trailing (user)
/// turns, or none under [`Quirks::cache_markers_ignored`] (ollama).
///
/// blallama follows the Anthropic placement: it checkpoints the end of every
/// prompt, so it resumes the next turn from there without markers on the
/// assistant turns. (Its old `breakpoint_after_assistant` quirk was dropped
/// in 0.51 after a 2026-10-02 A/B sweep: ±3 points of cache-read share and
/// no hash or segmentation drift with it off.)
///
/// The default [`Agent::on_turn`] calls this after all seating, immediately
/// before each infer — an `on_turn` override that still wants cached tails
/// must do the same.
///
/// [`Agent::on_turn`]: super::Agent::on_turn
pub fn roll_breakpoints(quirks: &Quirks, prompt: &mut Prompt) {
    roll_breakpoints_with(quirks, prompt, CacheControl::one_hour());
}

/// [`roll_breakpoints`] with a caller-chosen [`CacheControl`].
///
/// Positions already carrying a marker keep their original TTL, and the API
/// rejects a 5m marker ahead of a 1h one — so pick one TTL and stick with it
/// for the whole session. The 1h default matches the prefix markers
/// [`seed::prompt`] places and survives slow local-model rounds.
///
/// [`seed::prompt`]: super::seed
pub fn roll_breakpoints_with(
    quirks: &Quirks,
    prompt: &mut Prompt,
    cache_control: CacheControl,
) {
    if quirks.cache_markers_ignored {
        return;
    }
    // Roll with the tail: a user turn (the `Agent::prompt` invariant). An
    // empty prompt has only the pinned prefix markers to hit.
    let Some(anchor) = prompt.messages.len().checked_sub(1) else {
        return;
    };
    windowed(prompt, ROLL_WINDOW, anchor, cache_control);
}

/// `CachedPrompt::cache_windowed_with` generalized to an `anchor` message:
/// mark up to `n` messages at `anchor, anchor - 2, …` (skipping
/// already-marked ones, preserving their TTL), then enforce the 4-marker
/// budget by evicting middle message-level markers, earliest kept — so the
/// pinned prefix markers (tools, system, intro message) always survive.
///
/// The 2-step spacing matches the push-assistant + push-user cadence of a
/// tool round: the next roll's `anchor - 2k` lands on the previous roll's
/// `anchor - 2(k - 1)`, re-marking in place instead of jumping role.
// A candidate to fold back into `misanthropic` beside `prompt::index` once
// the `Quirks` shape is final (#17 discussion) — until then the lore lives
// here, in one place.
fn windowed(
    prompt: &mut Prompt,
    n: usize,
    anchor: usize,
    cache_control: CacheControl,
) {
    // Pinned prefix markers. `Prompt::indices` skips server tools (they
    // carry their own `cache_control`), so count those separately.
    let server_tool_markers = prompt.tools.as_ref().map_or(0, |tools| {
        tools
            .iter()
            .filter(|t| t.as_method().is_none() && t.is_cached())
            .count()
    });
    let pinned = server_tool_markers
        + prompt
            .indices()
            .filter(|&i| {
                !matches!(i, Index::Block(BlockIndex::Message(_)))
                    && index_is_cached(prompt, i)
            })
            .count();

    // The prefix markers spend their budget first; shrink the window rather
    // than ever evicting them.
    let n = n.min(MAX_CACHE_CONTROLS_PER_REQUEST.saturating_sub(pinned));

    let mut tail: std::collections::HashSet<usize> =
        std::collections::HashSet::with_capacity(n);
    for k in 0..n {
        let Some(m) = anchor.checked_sub(2 * k) else {
            break;
        };
        // An already-marked message keeps its marker (and TTL);
        // `Content::cache_with` marks the last cacheable block.
        if !prompt.messages[m].content.has_cache() {
            prompt.messages[m].content.cache_with(cache_control.clone());
        }
        tail.insert(m);
    }

    // Message-level markers outside the tail, in cache-prefix order. Keep
    // the earliest (the intro-style pinned message marker) up to what's
    // left of the budget; evict the middle stragglers — previous rounds'
    // rolling markers the window has slid past.
    let budget = MAX_CACHE_CONTROLS_PER_REQUEST
        .saturating_sub(pinned)
        .saturating_sub(tail.len());
    let stragglers: Vec<Index> = prompt
        .indices()
        .filter(|&i| match i {
            Index::Block(BlockIndex::Message((m, _))) => {
                !tail.contains(&m) && index_is_cached(prompt, i)
            }
            _ => false,
        })
        .collect();
    for &index in stragglers.iter().skip(budget) {
        if let Some(IndexMut::Block(block)) = prompt.get_mut(index) {
            block.uncache();
        }
    }
}

/// Whether the [`Block`] or [`CustomMethodDef`] at `index` carries a marker.
///
/// [`Block`]: misanthropic::prompt::message::Block
/// [`CustomMethodDef`]: misanthropic::tool::CustomMethodDef
fn index_is_cached(prompt: &Prompt, index: Index) -> bool {
    use misanthropic::prompt::index::IndexRef;
    match prompt.get(index) {
        Some(IndexRef::Method(method)) => method.is_cached(),
        Some(IndexRef::Block(block)) => block.is_cached(),
        None => false,
    }
}

/// Where request `next` stops extending request `prev`, if it does: a change
/// to the tools, system, thinking or `tool_choice`, or to any message `prev`
/// sent — a new block on its last message included, which moves the end a
/// prefix cache anchored at message ends checkpointed. `cache_control` is
/// not compared (the rolling breakpoints move by design), nor `max_tokens`
/// or `output_config`, which a session sets per phase.
///
/// A session that only appends gives `None` for every consecutive pair of
/// its requests.
pub fn divergence(prev: &Prompt, next: &Prompt) -> Option<String> {
    fn wire(value: &impl serde::Serialize) -> serde_json::Value {
        fn strip(value: &mut serde_json::Value) {
            match value {
                serde_json::Value::Object(map) => {
                    map.remove("cache_control");
                    map.values_mut().for_each(strip);
                }
                serde_json::Value::Array(items) => {
                    items.iter_mut().for_each(strip)
                }
                _ => {}
            }
        }
        let mut value =
            serde_json::to_value(value).expect("a prompt always serializes");
        strip(&mut value);
        value
    }
    let heads = [
        ("tools", wire(&prev.tools), wire(&next.tools)),
        ("system", wire(&prev.system), wire(&next.system)),
        ("thinking", wire(&prev.thinking), wire(&next.thinking)),
        (
            "tool_choice",
            wire(&prev.tool_choice),
            wire(&next.tool_choice),
        ),
    ];
    for (field, a, b) in heads {
        if a != b {
            return Some(format!("`{field}` changed: {a} -> {b}"));
        }
    }
    for (i, a) in prev.messages.iter().enumerate() {
        let Some(b) = next.messages.get(i) else {
            return Some(format!(
                "message {i} of {} dropped",
                prev.messages.len()
            ));
        };
        let (a, b) = (wire(a), wire(b));
        if a != b {
            return Some(format!("message {i} changed: {a} -> {b}"));
        }
    }
    None
}

#[cfg(test)]
mod tests {
    use super::*;
    use misanthropic::prompt::message::Role;

    fn quirks(f: impl FnOnce(&mut Quirks)) -> Quirks {
        let mut q = Quirks::default();
        f(&mut q);
        q
    }

    /// System + pinned intro marker + `pairs` assistant/user rounds, ending
    /// on a user turn (the `Agent::prompt` invariant) — the shape
    /// `seed::prompt::assemble` plus a session's rounds produces.
    fn session(pairs: usize) -> Prompt {
        let mut prompt = Prompt::default()
            .system("system text")
            .add_message((Role::User, "intro"))
            .unwrap()
            .cache_1h();
        prompt.system.as_mut().unwrap().cache_1h();
        for i in 0..pairs {
            prompt
                .push_message((Role::Assistant, format!("asst {i}")))
                .unwrap();
            prompt
                .push_message((Role::User, format!("results {i}")))
                .unwrap();
        }
        prompt
    }

    /// Indices of messages carrying a marker.
    fn marked(prompt: &Prompt) -> Vec<usize> {
        prompt
            .messages
            .iter()
            .enumerate()
            .filter(|(_, m)| m.content.has_cache())
            .map(|(i, _)| i)
            .collect()
    }

    /// Total markers across tools + system + messages, off the wire shape.
    fn total_markers(prompt: &Prompt) -> usize {
        serde_json::to_string(prompt)
            .unwrap()
            .matches(r#""cache_control":"#)
            .count()
    }

    #[test]
    fn canonical_rolls_onto_user_turns() {
        // intro(0) a(1) u(2) a(3) u(4): tail window on the user turns.
        let mut prompt = session(2);
        roll_breakpoints(&Quirks::default(), &mut prompt);
        assert_eq!(marked(&prompt), vec![0, 2, 4]);
        assert_eq!(prompt.messages[2].role, Role::User);
        assert_eq!(prompt.messages[4].role, Role::User);
    }

    #[test]
    fn ollama_is_a_no_op() {
        let mut prompt = session(2);
        let before = total_markers(&prompt);
        let q = quirks(|q| q.cache_markers_ignored = true);
        roll_breakpoints(&q, &mut prompt);
        assert_eq!(total_markers(&prompt), before);
    }

    #[test]
    fn empty_prompt_is_harmless() {
        let mut prompt = Prompt::default();
        roll_breakpoints(&Quirks::default(), &mut prompt);
        assert_eq!(total_markers(&prompt), 0);
    }

    /// The agora-seed round-loop sim: across a whole session the budget
    /// holds, the prefix markers never move, the rolling pair never jumps
    /// role, and every marker stays 1h.
    #[test]
    fn round_loop_never_exceeds_budget_or_jumps_role() {
        let (q, role) = (Quirks::default(), Role::User);
        let mut prompt = session(0);
        for i in 0..10 {
            prompt
                .push_message((Role::Assistant, format!("asst {i}")))
                .unwrap();
            prompt
                .push_message((Role::User, format!("results {i}")))
                .unwrap();
            roll_breakpoints(&q, &mut prompt);

            assert!(
                total_markers(&prompt) <= MAX_CACHE_CONTROLS_PER_REQUEST,
                "round {i}: {} markers",
                total_markers(&prompt)
            );
            assert!(
                prompt.system.as_ref().unwrap().has_cache(),
                "round {i}: system marker evicted"
            );
            assert!(
                prompt.messages[0].content.has_cache(),
                "round {i}: intro marker evicted"
            );
            for idx in marked(&prompt).into_iter().skip(1) {
                assert_eq!(
                    prompt.messages[idx].role, role,
                    "round {i}: rolling marker jumped role at {idx}"
                );
            }
        }
        // Ported from the seed prompt guards: a 5m marker ahead of a 1h
        // one is a submit-time API error, so rolling must stay all-1h.
        let json = serde_json::to_string(&prompt).unwrap();
        assert!(
            !json.contains(r#""cache_control":{"type":"ephemeral"}"#),
            "5m marker present:\n{json}"
        );
    }

    #[test]
    fn re_rolling_without_new_messages_is_idempotent() {
        let mut prompt = session(3);
        roll_breakpoints(&Quirks::default(), &mut prompt);
        let first = marked(&prompt);
        roll_breakpoints(&Quirks::default(), &mut prompt);
        assert_eq!(marked(&prompt), first);
        assert!(total_markers(&prompt) <= MAX_CACHE_CONTROLS_PER_REQUEST);
    }

    #[test]
    fn window_shrinks_before_evicting_prefix_markers() {
        // Three pinned prefix markers (system ×2 via two blocks is not
        // constructible here, so: system + two marked leading messages).
        let mut prompt = session(3);
        prompt.messages[1].content.cache_1h(); // extra pinned-ish marker
        roll_breakpoints(&Quirks::default(), &mut prompt);
        assert!(total_markers(&prompt) <= MAX_CACHE_CONTROLS_PER_REQUEST);
        assert!(prompt.system.as_ref().unwrap().has_cache());
        assert!(prompt.messages[0].content.has_cache());
    }

    /// Live: two rounds against the real API on Haiku; round 2 must read
    /// the round-1 write. Haiku's minimum cacheable prefix is 4096 tokens,
    /// hence the padding. All-5m TTL to keep the writes cheap. Run with:
    /// `cargo test --features seed live_roll -- --ignored --nocapture`
    #[cfg(feature = "client")]
    #[tokio::test]
    #[ignore = "hits the live Anthropic API (cents, not dollars)"]
    async fn live_roll_breakpoints_hit_the_cache() {
        let key = std::env::var("ANTHROPIC_API_KEY").unwrap_or_else(|_| {
            let path = format!(
                "{}/Projects/agora/secrets/anthropic_api_key",
                std::env::var("HOME").expect("HOME")
            );
            std::fs::read_to_string(path)
                .expect("no ANTHROPIC_API_KEY and no key file")
                .trim()
                .to_string()
        });
        let client = misanthropic::Client::new(key).expect("client");

        // ~6k tokens of unique-ish padding, past Haiku's 4096 minimum.
        let padding: String = (0..500)
            .map(|i| {
                format!(
                    "Fact {i}: the {i}th cache line holds a distinct \
                     sentence so the prefix is long and incompressible.\n"
                )
            })
            .collect();
        let mut prompt = Prompt::default()
            .max_tokens(std::num::NonZeroU32::new(32).unwrap())
            .system(format!(
                "You are terse. Reply with a single word.\n\n{padding}"
            ))
            .add_message((Role::User, "Say the word: one."))
            .unwrap();
        prompt.system.as_mut().unwrap().cache();

        let quirks = Quirks::default();
        roll_breakpoints_with(&quirks, &mut prompt, CacheControl::ephemeral());
        let first = client.message(&prompt).await.expect("round 1");
        let wrote = first.usage.cache_creation_input_tokens.unwrap_or(0)
            + first.usage.cache_read_input_tokens.unwrap_or(0);
        assert!(
            wrote > 0,
            "round 1 neither wrote nor read cache (prefix under the \
             minimum?): {:?}",
            first.usage
        );

        prompt.push_message(first).unwrap();
        prompt
            .push_message((Role::User, "Say the word: two."))
            .unwrap();
        roll_breakpoints_with(&quirks, &mut prompt, CacheControl::ephemeral());
        let second = client.message(&prompt).await.expect("round 2");
        let read = second.usage.cache_read_input_tokens.unwrap_or(0);
        println!(
            "round 2 usage: read={read} create={:?} input={}",
            second.usage.cache_creation_input_tokens, second.usage.input_tokens
        );
        assert!(read > 0, "round 2 read nothing: {:?}", second.usage);
    }

    /// A new message extends; a moved marker is not a change; a block added
    /// to the last message sent, or a changed system, diverges
    #[test]
    fn divergence_sees_only_what_a_prefix_cache_does() {
        let prev = session(1);
        let mut next = prev.clone();
        next.push_message((Role::Assistant, "asst 1")).unwrap();
        next.push_message((Role::User, "results 1")).unwrap();
        roll_breakpoints(&Quirks::default(), &mut next);
        assert_eq!(divergence(&prev, &next), None);

        let mut grown = prev.clone();
        grown.messages.last_mut().unwrap().extend(["a note"]);
        let why = divergence(&prev, &grown).unwrap();
        assert!(why.starts_with("message 2 changed"), "{why}");

        let mut dropped = prev.clone();
        dropped.messages.pop();
        assert!(divergence(&prev, &dropped).unwrap().contains("dropped"));

        let other = prev.clone().system("other system");
        assert!(divergence(&prev, &other).unwrap().starts_with("`system`"));
    }
}