use super::*;
#[derive(Debug)]
pub enum PruneResult {
NothingToPrune,
Pruned {
count: u64,
oldest_seq: u64,
newest_seq: u64,
},
}
pub async fn enforce_retention(store: &Store) -> Result<PruneResult> {
let retention_days = get_retention_days(store).await;
if retention_days == 0 {
return Ok(PruneResult::NothingToPrune);
}
let cutoff_ms = now_ms().saturating_sub(retention_days.saturating_mul(86_400_000));
let keys = store.scan_keys(EVENT_PREFIX).await?;
let mut expired: Vec<(String, u64)> = Vec::new();
for key in &keys {
let Some(seq) = key
.strip_prefix(EVENT_PREFIX)
.and_then(|s| s.parse::<u64>().ok())
else {
break;
};
let Ok(Some(bytes)) = store.get_raw_bytes(key).await else {
break;
};
let Ok(event) = serde_json::from_slice::<EnforcementEvent>(&bytes) else {
break;
};
if event.recorded_at_ms >= cutoff_ms {
break;
}
expired.push((key.clone(), seq));
}
if expired.is_empty() {
return Ok(PruneResult::NothingToPrune);
}
let oldest_seq = expired.first().expect("checked non-empty above").1;
let newest_seq = expired.last().expect("checked non-empty above").1;
let count = expired.len() as u64;
for (key, _) in &expired {
store.delete(key).await?;
}
let recorded = record_event(
store,
EnforcementEventType::RetentionPruned {
pruned_count: count,
oldest_pruned_seq: oldest_seq,
newest_pruned_seq: newest_seq,
},
SubjectKind::System,
"enforcement:retention".to_string(),
"system".to_string(),
None,
"retention_policy_enforced".to_string(),
None,
)
.await?;
if recorded.is_none() {
anyhow::bail!(
"pruned {count} enforcement events (seq {oldest_seq}-{newest_seq}) but could not \
record the RetentionPruned event — the deletion is unattested"
);
}
Ok(PruneResult::Pruned {
count,
oldest_seq,
newest_seq,
})
}
pub const STARTUP_GAP_THRESHOLD_MS: u64 = 24 * 60 * 60 * 1000;
pub async fn detect_startup_gap(
store: &Store,
gap_threshold_ms: u64,
) -> Result<Option<EnforcementEvent>> {
let keys = store.scan_keys(EVENT_PREFIX).await?;
let mut newest = None;
for (i, key) in keys.iter().enumerate().rev() {
if let Ok(Some(bytes)) = store.get_raw_bytes(key).await {
if let Ok(event) = serde_json::from_slice::<EnforcementEvent>(&bytes) {
newest = Some((i, event));
break;
}
}
}
let Some((newest_idx, newest)) = newest else {
return Ok(None);
};
if matches!(
newest.event_type,
EnforcementEventType::CleanShutdown { .. }
) {
return Ok(None);
}
let (cause, threshold_ms) = if has_terminator(store, &keys[..newest_idx]).await {
(GapCause::UncleanShutdown, 0)
} else {
(GapCause::Unknown, gap_threshold_ms)
};
let current = now_ms();
if current.saturating_sub(newest.recorded_at_ms) <= threshold_ms {
return Ok(None);
}
let writer = shared_writer(store).await?;
let event = writer
.lock()
.await
.detect_and_record_gap(store, newest.recorded_at_ms, current, cause)
.await?;
Ok(Some(event))
}
async fn has_terminator(store: &Store, keys: &[String]) -> bool {
for key in keys.iter().rev() {
if let Ok(Some(bytes)) = store.get_raw_bytes(key).await {
if let Ok(event) = serde_json::from_slice::<EnforcementEvent>(&bytes) {
if matches!(event.event_type, EnforcementEventType::CleanShutdown { .. }) {
return true;
}
}
}
}
false
}
pub async fn record_clean_shutdown(
store: &Store,
reason: &str,
) -> Result<Option<EnforcementEvent>> {
record_event(
store,
EnforcementEventType::CleanShutdown {
reason: reason.to_string(),
},
SubjectKind::System,
"enforcement:stream".to_string(),
"system".to_string(),
None,
"clean_shutdown".to_string(),
None,
)
.await
}
pub async fn scan_events_since(store: &Store, since_ms: u64) -> Result<Vec<EnforcementEvent>> {
let all = scan_enforcement_events(store, 0, u64::MAX).await?;
Ok(all
.into_iter()
.filter(|e| e.recorded_at_ms >= since_ms)
.collect())
}
pub async fn count_events_by_type(store: &Store, since_ms: u64) -> Result<EnforcementEventCounts> {
let events = scan_events_since(store, since_ms).await?;
Ok(aggregate_event_counts(&events))
}
pub const CONSULTATION_COALESCE_MS: u64 = 2_000;