klieo-runlog 3.4.0

Tier 2 observability — RunLog aggregate + replay engine for klieo agents.
Documentation
//! `KvStore`-backed persistence for the [`Capture`] artifact (ADR-046). The
//! bucket holds one JSON blob per run, keyed by run id — the durable,
//! cross-process channel for embedded and distributed callers. Offline operator
//! export/import of the same serde blob goes through the `klieo-replay` CLI's
//! file path; both share [`Capture`]'s serialization.

use chrono::{DateTime, Utc};
use klieo_core::{KvStore, RunId};

use crate::capture::Capture;
use crate::error::RunLogError;

/// Every capture persist/load/gc operation keys off this one bucket; reads
/// against any other bucket miss all stored captures. Mirrors the ADR-045
/// checkpoint bucket scheme (one blob per run id).
pub const CAPTURE_BUCKET: &str = "klieo.captures";

/// Stores `capture` as the sole blob for its run id; a prior capture for the
/// same run is overwritten (last write wins, no history).
pub async fn persist_capture(kv: &dyn KvStore, capture: &Capture) -> Result<(), RunLogError> {
    let bytes = serde_json::to_vec(capture).map_err(|e| RunLogError::CaptureCodec {
        source: Box::new(e),
    })?;
    kv.put(CAPTURE_BUCKET, &capture.run_id.to_string(), bytes.into())
        .await
        .map_err(|e| RunLogError::CaptureBackend {
            source: Box::new(e),
        })?;
    Ok(())
}

/// `None` means no capture was ever persisted for `run_id`; a stored-but-
/// undecodable blob yields [`RunLogError::CaptureCodec`] rather than `None`, so
/// callers can distinguish "absent" from "corrupt".
pub async fn load_capture(kv: &dyn KvStore, run_id: RunId) -> Result<Option<Capture>, RunLogError> {
    let Some(entry) = kv
        .get(CAPTURE_BUCKET, &run_id.to_string())
        .await
        .map_err(|e| RunLogError::CaptureBackend {
            source: Box::new(e),
        })?
    else {
        return Ok(None);
    };
    let capture = serde_json::from_slice(&entry.value).map_err(|e| RunLogError::CaptureCodec {
        source: Box::new(e),
    })?;
    Ok(Some(capture))
}

/// Delete every stored capture whose `captured_at` is at or before `cutoff`,
/// returning the count removed. Enumerates keys one bounded page at a time then
/// fetches one blob per key, so neither the key list nor the payloads pull the
/// whole bucket into memory at once. A blob that fails to load is logged and
/// skipped, so one corrupt entry never aborts the sweep; only the key
/// enumeration propagates.
pub async fn gc_captures(kv: &dyn KvStore, cutoff: DateTime<Utc>) -> Result<u64, RunLogError> {
    gc_captures_paged(kv, cutoff, GC_KEY_PAGE).await
}

/// Sweep page size: large enough to amortize the per-page round trip, small
/// enough that one page of keys stays a bounded slice rather than the whole
/// bucket. The exact value is not load-bearing — any backend that cares
/// overrides `keys_paginated`. Matches `gc_checkpoints`' page size.
const GC_KEY_PAGE: usize = 256;

/// Page-walking core of [`gc_captures`], parameterized on page size so tests can
/// drive the multi-page cursor loop without seeding a full page of keys.
async fn gc_captures_paged(
    kv: &dyn KvStore,
    cutoff: DateTime<Utc>,
    page_size: usize,
) -> Result<u64, RunLogError> {
    let mut removed = 0u64;
    let mut cursor = None;
    loop {
        let page = kv
            .keys_paginated(CAPTURE_BUCKET, cursor, page_size)
            .await
            .map_err(|e| RunLogError::CaptureBackend {
                source: Box::new(e),
            })?;
        for key in &page.keys {
            let Some(capture) = load_one_for_gc(kv, key).await else {
                continue;
            };
            if capture.captured_at > cutoff {
                continue;
            }
            match kv.delete(CAPTURE_BUCKET, key).await {
                Ok(()) => removed += 1,
                Err(e) => tracing::warn!(
                    target: "klieo.capture.gc",
                    operation = "delete",
                    key = %key,
                    error = %e,
                    "capture gc delete failed; continuing"
                ),
            }
        }
        match page.next {
            Some(c) => cursor = Some(c),
            None => break,
        }
    }
    Ok(removed)
}

/// Best-effort single-blob load for [`gc_captures`]: returns `None` (after
/// logging) on a backend or codec failure so the sweep continues past a missing
/// or corrupt entry rather than aborting on it.
async fn load_one_for_gc(kv: &dyn KvStore, key: &str) -> Option<Capture> {
    let entry = match kv.get(CAPTURE_BUCKET, key).await {
        Ok(Some(entry)) => entry,
        Ok(None) => return None,
        Err(e) => {
            tracing::warn!(
                target: "klieo.capture.gc",
                operation = "read",
                key = %key,
                error = %e,
                "capture gc read failed; skipping"
            );
            return None;
        }
    };
    match serde_json::from_slice(&entry.value) {
        Ok(capture) => Some(capture),
        Err(e) => {
            tracing::warn!(
                target: "klieo.capture.gc",
                operation = "deserialize",
                key = %key,
                error = %e,
                "skipping undeserializable capture entry"
            );
            None
        }
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::types::{RunLog, RunStatus, Usage};
    use chrono::Utc;
    use klieo_core::test_utils::fake_kv;
    use klieo_core::RunId;

    fn capture_stamped(run_id: RunId, captured_at: chrono::DateTime<Utc>) -> Capture {
        let now = Utc::now();
        let run_log = RunLog {
            run_id,
            agent: "writer".into(),
            started_at: now,
            finished_at: Some(now),
            status: RunStatus::Completed,
            steps: Vec::new(),
            tokens: Usage::default(),
            cost_estimate: None,
        };
        let mut capture = Capture::from_runlog(run_log, Vec::new());
        capture.captured_at = captured_at;
        capture
    }

    #[tokio::test]
    async fn gc_captures_paged_reaps_eligible_across_pages() {
        let kv = fake_kv();
        let cutoff = Utc::now();
        let stale = cutoff - chrono::Duration::hours(1);
        let ids: Vec<RunId> = (0..5).map(|_| RunId::new()).collect();
        for id in &ids {
            persist_capture(kv.as_ref(), &capture_stamped(*id, stale))
                .await
                .unwrap();
        }

        // page_size 2 over 5 keys forces three pages + cursor advance.
        let removed = gc_captures_paged(kv.as_ref(), cutoff, 2).await.unwrap();

        assert_eq!(removed, 5);
        for id in &ids {
            assert!(load_capture(kv.as_ref(), *id).await.unwrap().is_none());
        }
    }

    #[tokio::test]
    async fn gc_captures_paged_keeps_captures_after_cutoff() {
        let kv = fake_kv();
        let cutoff = Utc::now();
        let fresh = cutoff + chrono::Duration::hours(1);
        let keep = RunId::new();
        persist_capture(kv.as_ref(), &capture_stamped(keep, fresh))
            .await
            .unwrap();

        let removed = gc_captures_paged(kv.as_ref(), cutoff, 2).await.unwrap();

        assert_eq!(removed, 0);
        assert!(load_capture(kv.as_ref(), keep).await.unwrap().is_some());
    }
}