use chrono::{DateTime, Utc};
use klieo_core::{KvStore, RunId};
use crate::capture::Capture;
use crate::error::RunLogError;
pub const CAPTURE_BUCKET: &str = "klieo.captures";
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(())
}
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))
}
pub async fn gc_captures(kv: &dyn KvStore, cutoff: DateTime<Utc>) -> Result<u64, RunLogError> {
gc_captures_paged(kv, cutoff, GC_KEY_PAGE).await
}
const GC_KEY_PAGE: usize = 256;
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)
}
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();
}
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());
}
}