use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use std::path::{Path, PathBuf};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LocalOutputObservation {
pub path: PathBuf,
pub dataset_uri: String,
pub dataset_id: String,
pub kind: String,
pub pipeline: String,
pub row: String,
pub run_id: String,
pub pre_existing: bool,
pub replaced: bool,
pub retention_days: Option<u32>,
pub observed_at: DateTime<Utc>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct LocalOutputRecord {
pub id: String,
pub path: String,
pub dataset_uri: String,
pub dataset_id: String,
pub kind: String,
pub pipeline: String,
pub row: String,
pub run_id: String,
pub pre_existing: bool,
#[serde(default)]
pub replaced: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub retention_days: Option<u32>,
pub first_written_at: DateTime<Utc>,
pub last_written_at: DateTime<Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub deleted_at: Option<DateTime<Utc>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub deleted_bytes: Option<u64>,
}
impl LocalOutputRecord {
pub fn new(obs: &LocalOutputObservation) -> Self {
Self {
id: output_id(&obs.path),
path: obs.path.to_string_lossy().to_string(),
dataset_uri: obs.dataset_uri.clone(),
dataset_id: obs.dataset_id.clone(),
kind: obs.kind.clone(),
pipeline: obs.pipeline.clone(),
row: obs.row.clone(),
run_id: obs.run_id.clone(),
pre_existing: obs.pre_existing,
replaced: obs.pre_existing && obs.replaced,
retention_days: obs.retention_days,
first_written_at: obs.observed_at,
last_written_at: obs.observed_at,
deleted_at: None,
deleted_bytes: None,
}
}
pub fn observe(&mut self, obs: &LocalOutputObservation) {
self.dataset_uri = obs.dataset_uri.clone();
self.dataset_id = obs.dataset_id.clone();
self.kind = obs.kind.clone();
self.pipeline = obs.pipeline.clone();
self.row = obs.row.clone();
self.run_id = obs.run_id.clone();
self.retention_days = obs.retention_days;
if obs.observed_at > self.last_written_at {
self.last_written_at = obs.observed_at;
}
if obs.observed_at < self.first_written_at {
self.first_written_at = obs.observed_at;
}
if self.pre_existing && obs.replaced {
self.replaced = true;
}
self.deleted_at = None;
self.deleted_bytes = None;
}
pub fn state(&self) -> LocalOutputState {
if self.deleted_at.is_some() {
LocalOutputState::Expired
} else if self.pre_existing && self.replaced {
LocalOutputState::Replaced
} else if self.pre_existing {
LocalOutputState::External
} else {
LocalOutputState::Present
}
}
pub fn age_secs(&self, now: DateTime<Utc>) -> u64 {
now.signed_duration_since(self.last_written_at)
.num_seconds()
.max(0) as u64
}
pub fn effective_retention_days(&self, default_days: u32) -> Option<u32> {
match self.retention_days.unwrap_or(default_days) {
0 => None,
days => Some(days),
}
}
pub fn is_expired_by_age(&self, default_days: u32, now: DateTime<Utc>) -> bool {
match self.effective_retention_days(default_days) {
None => false,
Some(days) => self.age_secs(now) >= u64::from(days) * 86_400,
}
}
pub fn fs_path(&self) -> &Path {
Path::new(&self.path)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum LocalOutputState {
Present,
Expired,
External,
Replaced,
}
impl LocalOutputState {
pub fn as_str(self) -> &'static str {
match self {
Self::Present => "present",
Self::Expired => "expired",
Self::External => "external",
Self::Replaced => "replaced",
}
}
}
pub fn output_id(path: &Path) -> String {
use sha2::{Digest, Sha256};
let digest = Sha256::digest(path.to_string_lossy().as_bytes());
crate::serve::history::catalog::hex_prefix(&digest, 16)
}
#[derive(Debug, Clone, Default)]
pub struct LocalOutputFilter {
pub dataset_id: Option<String>,
pub pipeline: Option<String>,
pub include_deleted: bool,
pub limit: usize,
}
pub fn matches(rec: &LocalOutputRecord, filter: &LocalOutputFilter) -> bool {
if !filter.include_deleted && rec.deleted_at.is_some() {
return false;
}
if let Some(id) = &filter.dataset_id
&& rec.dataset_id != *id
{
return false;
}
if let Some(pipeline) = &filter.pipeline
&& rec.pipeline != *pipeline
{
return false;
}
true
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SweepScope {
Output(String),
Dataset(String),
Run(String),
OlderThanDays(u32),
Expired,
All,
}
impl SweepScope {
pub fn requires_confirmation(&self) -> bool {
matches!(self, Self::All | Self::OlderThanDays(0))
}
pub fn label(&self) -> &'static str {
match self {
Self::Output(_) => "output",
Self::Dataset(_) => "dataset",
Self::Run(_) => "run",
Self::OlderThanDays(_) => "older_than",
Self::Expired => "expired",
Self::All => "all",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SkipReason {
PreExisting,
AlreadyDeleted,
NotOnDisk,
InFlight,
DeleteFailed,
}
impl SkipReason {
pub fn as_str(self) -> &'static str {
match self {
Self::PreExisting => "pre_existing",
Self::AlreadyDeleted => "already_deleted",
Self::NotOnDisk => "not_on_disk",
Self::InFlight => "in_flight",
Self::DeleteFailed => "delete_failed",
}
}
pub fn marks_expired(self) -> bool {
matches!(self, Self::NotOnDisk)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SweepOutcome {
pub id: String,
pub path: String,
pub dataset_uri: String,
pub deleted: bool,
pub bytes: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub skipped: Option<SkipReason>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct SweepReport {
pub dry_run: bool,
pub scope: String,
pub deleted: usize,
pub bytes: u64,
pub skipped: usize,
pub outputs: Vec<SweepOutcome>,
}
impl SweepReport {
pub fn push(&mut self, outcome: SweepOutcome) {
if outcome.deleted {
self.deleted += 1;
self.bytes += outcome.bytes;
} else {
self.skipped += 1;
}
self.outputs.push(outcome);
}
pub fn skipped_for(&self, reason: SkipReason) -> usize {
self.outputs
.iter()
.filter(|o| o.skipped == Some(reason))
.count()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn ts(s: &str) -> DateTime<Utc> {
DateTime::parse_from_rfc3339(s).unwrap().with_timezone(&Utc)
}
fn obs(path: &str, at: &str) -> LocalOutputObservation {
LocalOutputObservation {
path: PathBuf::from(path),
dataset_uri: format!("file://{path}"),
dataset_id: "abc123".into(),
kind: "jsonl".into(),
pipeline: "p".into(),
row: "default".into(),
run_id: "r1".into(),
pre_existing: false,
replaced: false,
retention_days: None,
observed_at: ts(at),
}
}
#[test]
fn output_id_is_stable_and_path_specific() {
let a = output_id(Path::new("/tmp/out.jsonl"));
assert_eq!(a, output_id(Path::new("/tmp/out.jsonl")));
assert_ne!(a, output_id(Path::new("/tmp/other.jsonl")));
assert_eq!(a.len(), 16);
assert!(a.chars().all(|c| c.is_ascii_hexdigit()));
}
#[test]
fn new_row_starts_present() {
let rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z"));
assert_eq!(rec.state(), LocalOutputState::Present);
assert_eq!(rec.first_written_at, rec.last_written_at);
assert!(rec.deleted_at.is_none());
}
#[test]
fn observe_refreshes_last_written_but_keeps_first() {
let mut rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z"));
let mut second = obs("/tmp/a.jsonl", "2026-08-05T00:00:00Z");
second.run_id = "r2".into();
rec.observe(&second);
assert_eq!(rec.first_written_at, ts("2026-08-01T00:00:00Z"));
assert_eq!(rec.last_written_at, ts("2026-08-05T00:00:00Z"));
assert_eq!(rec.run_id, "r2");
}
#[test]
fn observe_never_reclassifies_a_faucet_created_file_as_external() {
let mut rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z"));
let mut second = obs("/tmp/a.jsonl", "2026-08-05T00:00:00Z");
second.pre_existing = true;
rec.observe(&second);
assert!(!rec.pre_existing);
assert_eq!(rec.state(), LocalOutputState::Present);
}
#[test]
fn observe_never_backdates_first_written() {
let mut rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2026-08-05T00:00:00Z"));
rec.observe(&obs("/tmp/a.jsonl", "2026-08-09T00:00:00Z"));
assert_eq!(rec.first_written_at, ts("2026-08-05T00:00:00Z"));
rec.observe(&obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z"));
assert_eq!(rec.first_written_at, ts("2026-08-01T00:00:00Z"));
assert_eq!(
rec.last_written_at,
ts("2026-08-09T00:00:00Z"),
"an older observation must not roll `last_written_at` backwards"
);
}
#[test]
fn a_rewrite_un_expires_the_row() {
let mut rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z"));
rec.deleted_at = Some(ts("2026-08-09T00:00:00Z"));
rec.deleted_bytes = Some(42);
assert_eq!(rec.state(), LocalOutputState::Expired);
rec.observe(&obs("/tmp/a.jsonl", "2026-08-10T00:00:00Z"));
assert_eq!(rec.state(), LocalOutputState::Present);
assert!(rec.deleted_bytes.is_none());
}
#[test]
fn a_truncated_pre_existing_file_reads_as_replaced_and_a_later_truncation_upgrades() {
let mut o = obs("/tmp/theirs.jsonl", "2026-08-01T00:00:00Z");
o.pre_existing = true;
o.replaced = true;
assert_eq!(
LocalOutputRecord::new(&o).state(),
LocalOutputState::Replaced
);
let mut first = obs("/tmp/theirs.csv", "2026-08-01T00:00:00Z");
first.pre_existing = true;
let mut rec = LocalOutputRecord::new(&first);
assert_eq!(rec.state(), LocalOutputState::External);
let mut second = obs("/tmp/theirs.csv", "2026-08-02T00:00:00Z");
second.pre_existing = true;
second.replaced = true;
rec.observe(&second);
assert_eq!(rec.state(), LocalOutputState::Replaced);
assert!(rec.pre_existing, "still never collectable");
let mut ours = obs("/tmp/ours.jsonl", "2026-08-01T00:00:00Z");
ours.replaced = true;
assert_eq!(
LocalOutputRecord::new(&ours).state(),
LocalOutputState::Present
);
assert_eq!(LocalOutputState::Replaced.as_str(), "replaced");
}
#[test]
fn pre_existing_row_reads_as_external() {
let mut o = obs("/tmp/theirs.jsonl", "2026-08-01T00:00:00Z");
o.pre_existing = true;
assert_eq!(
LocalOutputRecord::new(&o).state(),
LocalOutputState::External
);
}
#[test]
fn expiry_uses_the_default_window_when_the_row_has_no_override() {
let rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z"));
assert!(!rec.is_expired_by_age(7, ts("2026-08-07T23:59:59Z")));
assert!(rec.is_expired_by_age(7, ts("2026-08-08T00:00:00Z")));
}
#[test]
fn a_per_pipeline_override_wins_over_the_default() {
let mut o = obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z");
o.retention_days = Some(1);
let rec = LocalOutputRecord::new(&o);
assert_eq!(rec.effective_retention_days(7), Some(1));
assert!(rec.is_expired_by_age(7, ts("2026-08-02T00:00:01Z")));
}
#[test]
fn zero_days_means_keep_forever_at_either_level() {
let mut o = obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z");
o.retention_days = Some(0);
let rec = LocalOutputRecord::new(&o);
assert_eq!(rec.effective_retention_days(7), None);
assert!(!rec.is_expired_by_age(7, ts("2030-01-01T00:00:00Z")));
let rec = LocalOutputRecord::new(&obs("/tmp/b.jsonl", "2026-08-01T00:00:00Z"));
assert_eq!(rec.effective_retention_days(0), None);
assert!(!rec.is_expired_by_age(0, ts("2030-01-01T00:00:00Z")));
}
#[test]
fn age_of_a_future_timestamp_is_zero_not_negative() {
let rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2030-01-01T00:00:00Z"));
assert_eq!(rec.age_secs(ts("2026-08-01T00:00:00Z")), 0);
assert!(!rec.is_expired_by_age(7, ts("2026-08-01T00:00:00Z")));
}
#[test]
fn filter_hides_deleted_rows_unless_asked() {
let mut rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z"));
rec.deleted_at = Some(ts("2026-08-09T00:00:00Z"));
assert!(!matches(&rec, &LocalOutputFilter::default()));
assert!(matches(
&rec,
&LocalOutputFilter {
include_deleted: true,
..Default::default()
}
));
}
#[test]
fn filter_narrows_by_dataset_and_pipeline() {
let rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z"));
assert!(matches(
&rec,
&LocalOutputFilter {
dataset_id: Some("abc123".into()),
..Default::default()
}
));
assert!(!matches(
&rec,
&LocalOutputFilter {
dataset_id: Some("other".into()),
..Default::default()
}
));
assert!(matches(
&rec,
&LocalOutputFilter {
pipeline: Some("p".into()),
..Default::default()
}
));
assert!(!matches(
&rec,
&LocalOutputFilter {
pipeline: Some("q".into()),
..Default::default()
}
));
}
#[test]
fn only_clean_all_requires_confirmation() {
assert!(SweepScope::All.requires_confirmation());
assert!(SweepScope::OlderThanDays(0).requires_confirmation());
assert!(!SweepScope::Expired.requires_confirmation());
assert!(!SweepScope::OlderThanDays(3).requires_confirmation());
assert!(!SweepScope::Output("x".into()).requires_confirmation());
assert!(!SweepScope::Dataset("d".into()).requires_confirmation());
assert!(!SweepScope::Run("r".into()).requires_confirmation());
}
#[test]
fn scope_labels_are_distinct() {
let labels = [
SweepScope::Output("x".into()).label(),
SweepScope::Dataset("d".into()).label(),
SweepScope::Run("r".into()).label(),
SweepScope::OlderThanDays(1).label(),
SweepScope::Expired.label(),
SweepScope::All.label(),
];
let unique: std::collections::BTreeSet<_> = labels.iter().collect();
assert_eq!(unique.len(), labels.len());
}
#[test]
fn only_a_missing_file_marks_the_row_expired() {
assert!(SkipReason::NotOnDisk.marks_expired());
for r in [
SkipReason::PreExisting,
SkipReason::AlreadyDeleted,
SkipReason::InFlight,
SkipReason::DeleteFailed,
] {
assert!(!r.marks_expired(), "{}", r.as_str());
}
}
#[test]
fn report_totals_track_deleted_and_skipped() {
let mut rep = SweepReport::default();
rep.push(SweepOutcome {
id: "a".into(),
path: "/tmp/a".into(),
dataset_uri: "file:///tmp/a".into(),
deleted: true,
bytes: 100,
skipped: None,
error: None,
});
rep.push(SweepOutcome {
id: "b".into(),
path: "/tmp/b".into(),
dataset_uri: "file:///tmp/b".into(),
deleted: false,
bytes: 0,
skipped: Some(SkipReason::PreExisting),
error: None,
});
assert_eq!((rep.deleted, rep.bytes, rep.skipped), (1, 100, 1));
assert_eq!(rep.skipped_for(SkipReason::PreExisting), 1);
assert_eq!(rep.skipped_for(SkipReason::NotOnDisk), 0);
}
#[test]
fn state_and_reason_strings_are_stable() {
assert_eq!(LocalOutputState::Present.as_str(), "present");
assert_eq!(LocalOutputState::Expired.as_str(), "expired");
assert_eq!(LocalOutputState::External.as_str(), "external");
assert_eq!(SkipReason::PreExisting.as_str(), "pre_existing");
assert_eq!(SkipReason::NotOnDisk.as_str(), "not_on_disk");
}
#[test]
fn record_round_trips_through_json() {
let rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z"));
let back: LocalOutputRecord =
serde_json::from_str(&serde_json::to_string(&rec).unwrap()).unwrap();
assert_eq!(rec, back);
}
#[test]
fn fs_path_matches_the_stored_string() {
let rec = LocalOutputRecord::new(&obs("/tmp/a.jsonl", "2026-08-01T00:00:00Z"));
assert_eq!(rec.fs_path(), Path::new("/tmp/a.jsonl"));
}
}