supercode-harness 0.5.149

The optional native Volter Harness agent and tool harness
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
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
//! Token spend: what each harness session used, read from the harness's own records (Claude Code's `message.usage`
//! on each response, Codex's cumulative `token_count` totals) and priced at the model's list rates
//! ([`crate::pricing::built_in`]). A model with no list price is unrecorded, never zero.
//!
//! [`session_spend`] answers one loaded session per UTC day and model (`harness.v1.sessions.usage`), from a moment
//! on when asked. [`machine_spend`] answers every Claude Code and Codex session on this machine with use in a window
//! (`harness.v1.sessions.spend`, `supercode sessions usage`): a subagent's use counts on its parent session. It reads
//! only transcripts modified inside the window, and of each only its tail back to the window's start.

use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::io::{Read, Seek, SeekFrom};
use std::path::Path;

use serde_json::{json, Map, Value};
use supercode_interchange::sidecar::{ms_to_rfc3339, rfc3339_to_ms};

use crate::pricing::{built_in, RecordedTokens};
use crate::HarnessHomes;

/// One response's (Claude Code) or one turn's (Codex) recorded use.
#[derive(Debug, Clone)]
pub struct UseRecord {
    /// When the harness recorded it.
    pub at: Option<String>,
    /// The same, in unix milliseconds.
    pub at_ms: Option<i64>,
    /// The model that answered.
    pub model: String,
    /// What it used.
    pub tokens: RecordedTokens,
    /// Whether it ran in fast mode (`usage.speed: "fast"`).
    pub fast: bool,
    /// Whether it ran US-only (`usage.inference_geo: "us"`).
    pub us_only: bool,
}

/// Use summed, and what it cost at list rates where its models are priced.
#[derive(Debug, Clone, Default)]
pub struct Spend {
    /// Tokens of every kind.
    pub tokens: RecordedTokens,
    /// List-rate cost of the priced part, in US dollars.
    pub cost_usd: f64,
    /// Whether any of it was priced.
    pub priced: bool,
    /// Whether any of it was not: its model has no list price.
    pub unpriced: bool,
}

impl Spend {
    /// Add one record, priced at its model's list rates.
    pub fn record(&mut self, record: &UseRecord) {
        self.tokens.add(&record.tokens);
        match built_in(&record.model) {
            Some(price) => {
                self.cost_usd +=
                    price.recorded_cost_usd(&record.tokens, record.fast, record.us_only);
                self.priced = true;
            }
            None => self.unpriced = true,
        }
    }

    /// Add another sum.
    pub fn add(&mut self, other: &Spend) {
        self.tokens.add(&other.tokens);
        self.cost_usd += other.cost_usd;
        self.priced |= other.priced;
        self.unpriced |= other.unpriced;
    }

    /// The fields every reading answers: each token kind (`cache_write_tokens` is both writes), `cost` only where some
    /// use was priced, and `cost_unrecorded` when some was not.
    pub fn fields(&self) -> Map<String, Value> {
        let t = &self.tokens;
        let mut fields = Map::new();
        fields.insert("input_tokens".into(), json!(t.input));
        fields.insert("output_tokens".into(), json!(t.output));
        fields.insert("cache_read_tokens".into(), json!(t.cache_read));
        fields.insert(
            "cache_write_tokens".into(),
            json!(t.cache_write_5m + t.cache_write_1h),
        );
        fields.insert("cache_write_5m_tokens".into(), json!(t.cache_write_5m));
        fields.insert("cache_write_1h_tokens".into(), json!(t.cache_write_1h));
        if self.priced {
            fields.insert(
                "cost".into(),
                json!({"amount": (self.cost_usd * 1e6).round() / 1e6, "currency": "USD", "source": "rate_table"}),
            );
        }
        if self.unpriced {
            fields.insert("cost_unrecorded".into(), json!(true));
        }
        fields
    }
}

/// The use one Claude Code transcript line records, with its response's message id: Claude writes a response's
/// `usage` on each of its lines, so a reader keeps one per id (the last). A synthetic message records none.
/// `cache_creation_input_tokens` is split by TTL where Claude records the split (`usage.cache_creation`); a write
/// the split does not account for is at the default 5-minute TTL.
pub fn claude_use(record: &Value) -> Option<(String, UseRecord)> {
    let message = record.get("message")?;
    let id = message.get("id")?.as_str()?;
    let usage = message.get("usage")?;
    let model = message
        .get("model")
        .and_then(Value::as_str)
        .unwrap_or("unknown");
    if model == "<synthetic>" {
        return None;
    }
    let n = |value: Option<&Value>| value.and_then(Value::as_u64).unwrap_or(0);
    let written = n(usage.get("cache_creation_input_tokens"));
    let (write_5m, write_1h) = match usage.get("cache_creation") {
        Some(split) => {
            let write_5m = n(split.get("ephemeral_5m_input_tokens"));
            let write_1h = n(split.get("ephemeral_1h_input_tokens"));
            (
                write_5m + written.saturating_sub(write_5m + write_1h),
                write_1h,
            )
        }
        None => (written, 0),
    };
    let at = record
        .get("timestamp")
        .and_then(Value::as_str)
        .map(str::to_string);
    Some((
        id.to_string(),
        UseRecord {
            at_ms: at.as_deref().and_then(rfc3339_to_ms),
            at,
            model: model.to_string(),
            tokens: RecordedTokens {
                input: n(usage.get("input_tokens")),
                output: n(usage.get("output_tokens")),
                cache_read: n(usage.get("cache_read_input_tokens")),
                cache_write_5m: write_5m,
                cache_write_1h: write_1h,
            },
            fast: usage.get("speed").and_then(Value::as_str) == Some("fast"),
            us_only: usage.get("inference_geo").and_then(Value::as_str) == Some("us"),
        },
    ))
}

/// Codex's cumulative `token_count` totals, read in order: each change is one turn's use. Codex counts cached prompt
/// tokens inside `input_tokens`; a use separates them, as Claude's `usage` does. The model is the latest
/// `turn_context`'s.
#[derive(Debug, Clone)]
pub struct CodexUse {
    model: String,
    previous: [u64; 3],
}

impl CodexUse {
    /// A reader starting from no totals, with the session's recorded model until a `turn_context` names one.
    pub fn new(model: Option<String>) -> Self {
        CodexUse {
            model: model.unwrap_or_else(|| "unknown".into()),
            previous: [0; 3],
        }
    }

    /// The use `record` adds, when it is a changed total.
    pub fn next(&mut self, record: &Value) -> Option<UseRecord> {
        let payload = record.get("payload").unwrap_or(&Value::Null);
        if record.get("type").and_then(Value::as_str) == Some("turn_context") {
            if let Some(model) = payload.get("model").and_then(Value::as_str) {
                self.model = model.to_string();
            }
            return None;
        }
        if payload.get("type").and_then(Value::as_str) != Some("token_count") {
            return None;
        }
        let total = payload.get("info")?.get("total_token_usage")?;
        let n = |value: Option<&Value>| value.and_then(Value::as_u64).unwrap_or(0);
        let now = [
            n(total.get("input_tokens")),
            n(total.get("output_tokens")),
            n(total.get("cached_input_tokens")),
        ];
        if now == self.previous {
            return None;
        }
        let [input, output, cached] = [0, 1, 2].map(|i| now[i].saturating_sub(self.previous[i]));
        self.previous = now;
        let at = record
            .get("timestamp")
            .and_then(Value::as_str)
            .map(str::to_string);
        Some(UseRecord {
            at_ms: at.as_deref().and_then(rfc3339_to_ms),
            at,
            model: self.model.clone(),
            tokens: RecordedTokens {
                input: input.saturating_sub(cached),
                output,
                cache_read: cached,
                ..RecordedTokens::default()
            },
            fast: false,
            us_only: false,
        })
    }
}

/// One loaded session's use per UTC day and model, each record counted from `since_ms` on (all of it without one).
pub fn session_spend(
    source: &crate::SessionSource,
    model: Option<String>,
    raw: &[String],
    since_ms: Option<i64>,
) -> Vec<Value> {
    let mut days: BTreeMap<(String, String), Spend> = BTreeMap::new();
    let mut count = |record: &UseRecord| {
        if since_ms.is_some_and(|since| record.at_ms.is_none_or(|at| at < since)) {
            return;
        }
        let day = record
            .at
            .as_deref()
            .and_then(|at| at.get(..10))
            .unwrap_or("unknown")
            .to_string();
        days.entry((day, record.model.clone()))
            .or_default()
            .record(record);
    };
    let records = raw
        .iter()
        .filter_map(|line| serde_json::from_str::<Value>(line).ok());
    match source {
        crate::SessionSource::ClaudeCode => {
            let responses: BTreeMap<String, UseRecord> =
                records.filter_map(|record| claude_use(&record)).collect();
            responses.values().for_each(&mut count);
        }
        crate::SessionSource::Codex => {
            let mut codex = CodexUse::new(model);
            for record in records {
                if let Some(used) = codex.next(&record) {
                    count(&used);
                }
            }
        }
        _ => {}
    }
    days.into_iter()
        .map(|((day, model), spend)| {
            let mut row = Map::new();
            row.insert("at".into(), json!(format!("{day}T00:00:00Z")));
            row.insert("model".into(), json!(model));
            row.extend(spend.fields());
            Value::Object(row)
        })
        .collect()
}

/// A window's start: a duration back from `now_ms` (`90s`, `10m`, `1h`, `2d`) or an RFC 3339 UTC moment.
pub fn parse_since(text: &str, now_ms: i64) -> Result<i64, String> {
    let text = text.trim();
    if let Some(at) = rfc3339_to_ms(text) {
        return Ok(at);
    }
    let (number, unit) = text.split_at(text.find(|c: char| !c.is_ascii_digit()).unwrap_or(text.len()));
    let amount: i64 = number
        .parse()
        .map_err(|_| format!("`{text}` is neither a duration (10m, 1h, 2d) nor an RFC 3339 moment"))?;
    let unit_ms = match unit {
        "s" => 1_000,
        "m" => 60_000,
        "h" => 3_600_000,
        "d" => 86_400_000,
        _ => {
            return Err(format!(
                "`{text}`: a duration's unit is s, m, h or d (10m, 1h, 2d)"
            ))
        }
    };
    Ok(now_ms - amount * unit_ms)
}

/// One session's use in a window, its subagents' included.
#[derive(Debug, Default)]
struct SessionSpend {
    by_model: BTreeMap<String, Spend>,
    cwd: Option<String>,
    subagents: BTreeSet<String>,
    last_ms: i64,
}

impl SessionSpend {
    fn record(&mut self, record: &UseRecord) {
        self.by_model
            .entry(record.model.clone())
            .or_default()
            .record(record);
        self.last_ms = self.last_ms.max(record.at_ms.unwrap_or_default());
    }

    fn total(&self) -> Spend {
        let mut total = Spend::default();
        for spend in self.by_model.values() {
            total.add(spend);
        }
        total
    }
}

/// Every Claude Code and Codex session on this machine with recorded use from `since_ms` to now, a subagent's use on
/// its parent. Sorted by `sort`: `cost` (the default) puts priced sessions first, most cost first, then the rest by
/// tokens; `tokens` sorts all by tokens. `harnesses` limits the reading (empty: both); `limit` keeps the first rows.
pub fn machine_spend(
    homes: &HarnessHomes,
    since_ms: i64,
    now_ms: i64,
    harnesses: &[String],
    sort: &str,
    limit: Option<usize>,
) -> Value {
    let wants = |harness: &str| harnesses.is_empty() || harnesses.iter().any(|h| h == harness);
    let mut sessions: BTreeMap<(String, String), SessionSpend> = BTreeMap::new();
    if wants(crate::HarnessId::CLAUDE_CODE) {
        claude_spend(&homes.claude_code, since_ms, &mut sessions);
    }
    if wants(crate::HarnessId::CODEX) {
        codex_spend(&homes.codex, since_ms, &mut sessions);
    }
    let names: HashMap<String, String> =
        crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
            .into_iter()
            .map(|session| (session.session_id, session.name))
            .collect();
    let machine = crate::mailbox::local_machine_name();
    let mut rows: Vec<(Spend, Value)> = sessions
        .into_iter()
        .map(|((harness, session_id), spend)| {
            let total = spend.total();
            let name = match harness.as_str() {
                crate::HarnessId::CODEX => Some(crate::mail_route::codex_name(&session_id)),
                _ => names.get(&session_id).cloned(),
            };
            let mut row = Map::new();
            row.insert("harness".into(), json!(harness));
            row.insert("session_id".into(), json!(session_id));
            row.insert(
                "address".into(),
                json!(crate::mailbox::MailAddress::new(&machine, &harness, &session_id)
                    .ok()
                    .map(|address| address.to_string())),
            );
            row.insert("name".into(), json!(name));
            row.insert("cwd".into(), json!(spend.cwd));
            row.insert("subagents".into(), json!(spend.subagents.len()));
            row.insert(
                "last_at".into(),
                json!((spend.last_ms > 0).then(|| ms_to_rfc3339(spend.last_ms))),
            );
            row.extend(total.fields());
            row.insert(
                "models".into(),
                Value::Array(
                    spend
                        .by_model
                        .iter()
                        .map(|(model, spend)| {
                            let mut entry = Map::new();
                            entry.insert("model".into(), json!(model));
                            entry.extend(spend.fields());
                            Value::Object(entry)
                        })
                        .collect(),
                ),
            );
            (total, Value::Object(row))
        })
        .collect();
    let by_tokens = sort == "tokens";
    rows.sort_by(|(a, _), (b, _)| {
        if by_tokens || a.priced != b.priced || !a.priced {
            (b.priced && !by_tokens)
                .cmp(&(a.priced && !by_tokens))
                .then(b.tokens.total().cmp(&a.tokens.total()))
        } else {
            b.cost_usd.total_cmp(&a.cost_usd)
        }
    });
    let mut total = Spend::default();
    for (spend, _) in &rows {
        total.add(spend);
    }
    let count = rows.len();
    let listed: Vec<Value> = rows
        .into_iter()
        .take(limit.unwrap_or(usize::MAX))
        .map(|(_, row)| row)
        .collect();
    let mut answer = Map::new();
    answer.insert("since".into(), json!(ms_to_rfc3339(since_ms)));
    answer.insert("until".into(), json!(ms_to_rfc3339(now_ms)));
    answer.insert("session_count".into(), json!(count));
    answer.insert("sessions".into(), Value::Array(listed));
    answer.insert("total".into(), Value::Object(total.fields()));
    Value::Object(answer)
}

/// How far before a window's start a tail is read: records are mostly, not strictly, in time order (a response's
/// lines, a subagent's result), so the read stops only at a record this much older than the window, and what it
/// reads in between is still counted by each record's own time. The limit: a record inside the window that sits in
/// the file before one more than ten minutes older than the window's start is not read, so not counted.
const OUT_OF_ORDER_MS: i64 = 10 * 60 * 1000;

/// Whether `path` was written at or after `since_ms`.
fn modified_since(path: &Path, since_ms: i64) -> bool {
    std::fs::metadata(path)
        .and_then(|meta| meta.modified())
        .ok()
        .and_then(|at| at.duration_since(std::time::UNIX_EPOCH).ok())
        .is_some_and(|at| at.as_millis() as i64 >= since_ms)
}

fn record_ms(record: &Value) -> Option<i64> {
    record
        .get("timestamp")
        .and_then(Value::as_str)
        .and_then(rfc3339_to_ms)
}

/// Claude Code's transcripts under `projects`: `<project>/<session>.jsonl`, and a subagent's
/// `<project>/<session>/subagents/*.jsonl`, counted on `<session>`.
fn claude_spend(
    projects: &Path,
    since_ms: i64,
    sessions: &mut BTreeMap<(String, String), SessionSpend>,
) {
    let Ok(projects) = std::fs::read_dir(projects) else {
        return;
    };
    for project in projects.flatten() {
        let Ok(entries) = std::fs::read_dir(project.path()) else {
            continue;
        };
        for entry in entries.flatten() {
            let path = entry.path();
            if path.is_dir() {
                let Some(parent) = path.file_name().and_then(|name| name.to_str()) else {
                    continue;
                };
                let Ok(subagents) = std::fs::read_dir(path.join("subagents")) else {
                    continue;
                };
                for subagent in subagents.flatten() {
                    let file = subagent.path();
                    if file.extension().is_some_and(|ext| ext == "jsonl")
                        && modified_since(&file, since_ms)
                    {
                        let agent = file
                            .file_stem()
                            .map(|stem| stem.to_string_lossy().into_owned());
                        claude_file(&file, since_ms, parent, agent, sessions);
                    }
                }
            } else if path.extension().is_some_and(|ext| ext == "jsonl")
                && modified_since(&path, since_ms)
            {
                if let Some(session) = path.file_stem().and_then(|stem| stem.to_str()) {
                    claude_file(&path, since_ms, session, None, sessions);
                }
            }
        }
    }
}

fn claude_file(
    path: &Path,
    since_ms: i64,
    session_id: &str,
    subagent: Option<String>,
    sessions: &mut BTreeMap<(String, String), SessionSpend>,
) {
    let records = tail_records(path, |record| {
        record_ms(record).is_some_and(|at| at < since_ms - OUT_OF_ORDER_MS)
    });
    let cwd = records
        .iter()
        .rev()
        .find_map(|record| record.get("cwd").and_then(Value::as_str))
        .map(str::to_string);
    let responses: BTreeMap<String, UseRecord> = records
        .iter()
        .filter_map(claude_use)
        .filter(|(_, used)| used.at_ms.is_some_and(|at| at >= since_ms))
        .collect();
    if responses.is_empty() {
        return;
    }
    let session = sessions
        .entry((crate::HarnessId::CLAUDE_CODE.to_string(), session_id.to_string()))
        .or_default();
    for used in responses.values() {
        session.record(used);
    }
    // A session's folder is its own transcript's; a subagent's stands in until that is read.
    if subagent.is_none() || session.cwd.is_none() {
        session.cwd = cwd.or(session.cwd.take());
    }
    if let Some(agent) = subagent {
        session.subagents.insert(agent);
    }
}

/// Codex's rollouts under `sessions` (`YYYY/MM/DD/rollout-*.jsonl`); a subagent thread's use counts on its parent.
fn codex_spend(
    root: &Path,
    since_ms: i64,
    sessions: &mut BTreeMap<(String, String), SessionSpend>,
) {
    let mut stack = vec![(root.to_path_buf(), 0)];
    while let Some((directory, depth)) = stack.pop() {
        let Ok(entries) = std::fs::read_dir(&directory) else {
            continue;
        };
        for entry in entries.flatten() {
            let path = entry.path();
            if path.is_dir() {
                if depth < 4 {
                    stack.push((path, depth + 1));
                }
                continue;
            }
            let is_rollout = path
                .file_name()
                .and_then(|name| name.to_str())
                .is_some_and(|name| name.starts_with("rollout-") && name.ends_with(".jsonl"));
            if is_rollout && modified_since(&path, since_ms) {
                codex_file(&path, since_ms, sessions);
            }
        }
    }
}

fn codex_file(path: &Path, since_ms: i64, sessions: &mut BTreeMap<(String, String), SessionSpend>) {
    let Some((thread, parent)) = crate::codex_peer::rollout_session(path) else {
        return;
    };
    // Read back to a total recorded before the window (each later total is a delta from it) and the model then.
    let (mut total_before, mut model_before) = (false, false);
    let records = tail_records(path, |record| {
        if record_ms(record).is_some_and(|at| at < since_ms - OUT_OF_ORDER_MS) {
            let kind = record.get("type").and_then(Value::as_str);
            total_before |= record.pointer("/payload/type").and_then(Value::as_str) == Some("token_count");
            model_before |= kind == Some("turn_context");
        }
        total_before && model_before
    });
    let mut codex = CodexUse::new(None);
    let used: Vec<UseRecord> = records
        .iter()
        .filter_map(|record| codex.next(record))
        .filter(|used| used.at_ms.is_some_and(|at| at >= since_ms))
        .collect();
    if used.is_empty() {
        return;
    }
    let session_id = parent.clone().unwrap_or_else(|| thread.clone());
    let session = sessions
        .entry((crate::HarnessId::CODEX.to_string(), session_id))
        .or_default();
    for record in &used {
        session.record(record);
    }
    if parent.is_some() {
        session.subagents.insert(thread);
    }
    if session.cwd.is_none() || parent.is_none() {
        if let Some(cwd) = crate::codex_peer::rollout_cwd(path) {
            session.cwd = Some(cwd.to_string_lossy().into_owned());
        }
    }
}

/// The records of a JSONL file from the newest back to the first one `start` accepts (included), oldest first; the
/// whole file when none does. Reads back from the end in blocks, so a long transcript costs only its tail. `start`
/// sees each record newest first.
fn tail_records(path: &Path, mut start: impl FnMut(&Value) -> bool) -> Vec<Value> {
    const BLOCK: u64 = 256 * 1024;
    let Ok(mut file) = std::fs::File::open(path) else {
        return Vec::new();
    };
    let Ok(length) = file.metadata().map(|meta| meta.len()) else {
        return Vec::new();
    };
    let mut newest_first = Vec::new();
    // The bytes before the first newline of the block read last: a line whose start is further back.
    let mut carry: Vec<u8> = Vec::new();
    let mut end = length;
    while end > 0 {
        let begin = end.saturating_sub(BLOCK);
        let mut block = vec![0u8; (end - begin) as usize];
        if file.seek(SeekFrom::Start(begin)).is_err() || file.read_exact(&mut block).is_err() {
            break;
        }
        block.extend_from_slice(&carry);
        let first_newline = if begin == 0 {
            None
        } else {
            match block.iter().position(|byte| *byte == b'\n') {
                Some(at) => Some(at),
                None => {
                    carry = block;
                    end = begin;
                    continue;
                }
            }
        };
        let (head, lines) = match first_newline {
            Some(at) => (block[..at].to_vec(), &block[at + 1..]),
            None => (Vec::new(), &block[..]),
        };
        let mut reached = false;
        for line in lines.split(|byte| *byte == b'\n').rev() {
            let Ok(record) = serde_json::from_slice::<Value>(line) else {
                continue;
            };
            reached = start(&record);
            newest_first.push(record);
            if reached {
                break;
            }
        }
        if reached {
            break;
        }
        carry = head;
        end = begin;
    }
    newest_first.reverse();
    newest_first
}