use std::path::{Path, PathBuf};
use std::time::{SystemTime, UNIX_EPOCH};
use serde::{Deserialize, Serialize};
use tokio::sync::Mutex;
use uuid::Uuid;
use crate::{AppendResult, Client, ClientError};
use kindling_types::{Id, ObservationInput};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SpoolStatus {
pub pending_count: usize,
pub spool_path: PathBuf,
#[serde(default)]
pub last_flush_time_ms: Option<i64>,
#[serde(default)]
pub last_error: Option<String>,
pub replay_attempts: u64,
#[serde(default)]
pub dropped_count: u64,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
struct SpoolRuntime {
last_flush_time_ms: Option<i64>,
last_error: Option<String>,
replay_attempts: u64,
#[serde(default)]
dropped_count: u64,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct SpoolConfig {
pub spool_path: PathBuf,
pub max_bytes: Option<u64>,
pub max_age_ms: Option<i64>,
}
impl SpoolConfig {
pub fn new(spool_path: impl Into<PathBuf>) -> Self {
Self {
spool_path: spool_path.into(),
max_bytes: None,
max_age_ms: None,
}
}
pub fn with_max_bytes(mut self, max_bytes: u64) -> Self {
self.max_bytes = Some(max_bytes);
self
}
pub fn with_max_age_ms(mut self, max_age_ms: i64) -> Self {
self.max_age_ms = Some(max_age_ms);
self
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SpoolEntry {
pub input: ObservationInput,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub capsule_id: Option<Id>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub validate: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub spooled_at: Option<i64>,
}
#[derive(Debug)]
pub enum AppendOutcome {
Delivered(Box<AppendResult>),
Spooled,
}
#[derive(Debug, PartialEq, Eq)]
pub struct FlushReport {
pub replayed: usize,
pub remaining: usize,
}
#[derive(Debug, thiserror::Error)]
pub enum SpoolError {
#[error("spool io error: {0}")]
Io(#[from] std::io::Error),
#[error("spool serde error: {0}")]
Serde(#[from] serde_json::Error),
#[error("client error: {0}")]
Client(#[from] ClientError),
}
#[derive(Debug)]
pub struct SpooledClient {
client: Client,
spool_path: PathBuf,
max_bytes: Option<u64>,
max_age_ms: Option<i64>,
file_lock: Mutex<()>,
runtime: Mutex<SpoolRuntime>,
}
impl SpooledClient {
pub fn new(client: Client, spool_path: PathBuf) -> Self {
let runtime = load_runtime_sidecar(&spool_path);
Self {
client,
spool_path,
max_bytes: None,
max_age_ms: None,
file_lock: Mutex::new(()),
runtime: Mutex::new(runtime),
}
}
pub fn with_config(client: Client, config: SpoolConfig) -> Self {
let runtime = load_runtime_sidecar(&config.spool_path);
Self {
client,
spool_path: config.spool_path,
max_bytes: config.max_bytes,
max_age_ms: config.max_age_ms,
file_lock: Mutex::new(()),
runtime: Mutex::new(runtime),
}
}
pub fn client(&self) -> &Client {
&self.client
}
pub async fn append_observation(
&self,
mut input: ObservationInput,
capsule_id: Option<Id>,
validate: Option<bool>,
) -> Result<AppendOutcome, SpoolError> {
if input.id.is_none() {
input.id = Some(Uuid::new_v4().to_string());
}
if self.pending_count()? > 0 {
self.flush().await?;
}
match self
.client
.append_observation(input.clone(), capsule_id.clone(), validate)
.await
{
Ok(result) => Ok(AppendOutcome::Delivered(Box::new(result))),
Err(err) if is_connectivity_error(&err) => {
self.record_connectivity_error(&err).await;
let entry = SpoolEntry {
input,
capsule_id,
validate,
spooled_at: Some(now_ms()),
};
self.append_to_spool(&entry).await?;
Ok(AppendOutcome::Spooled)
}
Err(err) => Err(SpoolError::Client(err)),
}
}
pub async fn flush(&self) -> Result<FlushReport, SpoolError> {
let _guard = self.file_lock.lock().await;
let entries = read_spool(&self.spool_path)?;
let total = entries.len();
if total == 0 {
return Ok(FlushReport {
replayed: 0,
remaining: 0,
});
}
let mut replayed = 0usize;
let mut propagate: Option<ClientError> = None;
for (idx, entry) in entries.iter().enumerate() {
self.bump_replay_attempts().await;
match self
.client
.append_observation(
entry.input.clone(),
entry.capsule_id.clone(),
entry.validate,
)
.await
{
Ok(_) => replayed += 1,
Err(err) if is_connectivity_error(&err) => {
self.record_connectivity_error(&err).await;
break;
}
Err(err) => {
let _ = idx;
propagate = Some(err);
break;
}
}
}
let mut remainder: Vec<SpoolEntry> = entries[replayed..].to_vec();
let dropped = if self.max_bytes.is_some() || self.max_age_ms.is_some() {
let before = remainder.len();
remainder = trim_entries(remainder, self.max_bytes, self.max_age_ms, now_ms());
before - remainder.len()
} else {
0
};
rewrite_spool(&self.spool_path, &remainder)?;
if dropped > 0 {
self.bump_dropped(dropped as u64).await;
}
if let Some(err) = propagate {
return Err(SpoolError::Client(err));
}
if replayed > 0 {
self.record_successful_flush().await;
}
Ok(FlushReport {
replayed,
remaining: remainder.len(),
})
}
async fn bump_replay_attempts(&self) {
let mut rt = self.runtime.lock().await;
rt.replay_attempts = rt.replay_attempts.saturating_add(1);
persist_runtime_sidecar(&self.spool_path, &rt);
}
async fn bump_dropped(&self, n: u64) {
let mut rt = self.runtime.lock().await;
rt.dropped_count = rt.dropped_count.saturating_add(n);
persist_runtime_sidecar(&self.spool_path, &rt);
}
async fn record_connectivity_error(&self, err: &ClientError) {
let mut rt = self.runtime.lock().await;
rt.last_error = Some(err.to_string());
persist_runtime_sidecar(&self.spool_path, &rt);
}
async fn record_successful_flush(&self) {
let mut rt = self.runtime.lock().await;
rt.last_flush_time_ms = Some(now_ms());
rt.last_error = None;
persist_runtime_sidecar(&self.spool_path, &rt);
}
pub fn pending_count(&self) -> Result<usize, SpoolError> {
Ok(read_spool(&self.spool_path)?.len())
}
pub async fn spool_status(&self) -> Result<SpoolStatus, SpoolError> {
let runtime = self.runtime.lock().await;
Ok(SpoolStatus {
pending_count: self.pending_count()?,
spool_path: self.spool_path.clone(),
last_flush_time_ms: runtime.last_flush_time_ms,
last_error: runtime.last_error.clone(),
replay_attempts: runtime.replay_attempts,
dropped_count: runtime.dropped_count,
})
}
pub fn spool_status_from_path(
spool_path: impl Into<PathBuf>,
) -> Result<SpoolStatus, SpoolError> {
let spool_path = spool_path.into();
let pending_count = read_spool(&spool_path)?.len();
let runtime = load_runtime_sidecar(&spool_path);
Ok(SpoolStatus {
pending_count,
spool_path,
last_flush_time_ms: runtime.last_flush_time_ms,
last_error: runtime.last_error,
replay_attempts: runtime.replay_attempts,
dropped_count: runtime.dropped_count,
})
}
async fn append_to_spool(&self, entry: &SpoolEntry) -> Result<(), SpoolError> {
use std::io::Write;
let _guard = self.file_lock.lock().await;
let line = serde_json::to_string(entry)?;
let mut file = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&self.spool_path)?;
file.write_all(line.as_bytes())?;
file.write_all(b"\n")?;
file.flush()?;
Ok(())
}
}
fn status_sidecar_path(spool_path: &Path) -> PathBuf {
let name = format!(
"{}.status.json",
spool_path
.file_name()
.map(|n| n.to_string_lossy().into_owned())
.unwrap_or_else(|| "spool".to_string())
);
match spool_path.parent() {
Some(dir) => dir.join(name),
None => PathBuf::from(name),
}
}
fn load_runtime_sidecar(spool_path: &Path) -> SpoolRuntime {
let path = status_sidecar_path(spool_path);
let contents = match std::fs::read_to_string(&path) {
Ok(c) => c,
Err(_) => return SpoolRuntime::default(),
};
serde_json::from_str(&contents).unwrap_or_default()
}
fn persist_runtime_sidecar(spool_path: &Path, runtime: &SpoolRuntime) {
let _ = (|| -> Result<(), SpoolError> {
let path = status_sidecar_path(spool_path);
let line = serde_json::to_string(runtime)?;
let tmp = temp_sibling(&path);
std::fs::write(&tmp, format!("{line}\n"))?;
std::fs::rename(&tmp, &path)?;
Ok(())
})();
}
fn now_ms() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as i64
}
fn is_connectivity_error(err: &ClientError) -> bool {
matches!(err, ClientError::Unavailable(_) | ClientError::Http(_))
}
fn entry_size(entry: &SpoolEntry) -> u64 {
serde_json::to_string(entry)
.map(|s| s.len() as u64 + 1)
.unwrap_or(0)
}
fn trim_entries(
mut entries: Vec<SpoolEntry>,
max_bytes: Option<u64>,
max_age_ms: Option<i64>,
now: i64,
) -> Vec<SpoolEntry> {
if let Some(max_age) = max_age_ms {
let cutoff = now.saturating_sub(max_age);
let mut drop_n = 0;
for entry in &entries {
match entry.spooled_at {
Some(t) if t < cutoff => drop_n += 1,
_ => break,
}
}
entries.drain(0..drop_n);
}
if let Some(max_bytes) = max_bytes {
let sizes: Vec<u64> = entries.iter().map(entry_size).collect();
let mut total: u64 = sizes.iter().sum();
let mut drop_n = 0;
while total > max_bytes && entries.len() - drop_n > 1 {
total -= sizes[drop_n];
drop_n += 1;
}
entries.drain(0..drop_n);
}
entries
}
fn read_spool(path: &Path) -> Result<Vec<SpoolEntry>, SpoolError> {
let contents = match std::fs::read_to_string(path) {
Ok(c) => c,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(e) => return Err(SpoolError::Io(e)),
};
let lines: Vec<&str> = contents.split('\n').collect();
let mut entries = Vec::new();
let last_idx = lines.len().saturating_sub(1);
for (idx, raw) in lines.iter().enumerate() {
let line = raw.trim_end_matches('\r');
if line.is_empty() {
continue;
}
match serde_json::from_str::<SpoolEntry>(line) {
Ok(entry) => entries.push(entry),
Err(e) => {
if idx == last_idx {
break;
}
return Err(SpoolError::Serde(e));
}
}
}
Ok(entries)
}
fn rewrite_spool(path: &Path, entries: &[SpoolEntry]) -> Result<(), SpoolError> {
use std::io::Write;
if entries.is_empty() {
match std::fs::remove_file(path) {
Ok(()) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(e) => Err(SpoolError::Io(e)),
}
} else {
let tmp = temp_sibling(path);
{
let mut file = std::fs::File::create(&tmp)?;
for entry in entries {
let line = serde_json::to_string(entry)?;
file.write_all(line.as_bytes())?;
file.write_all(b"\n")?;
}
file.flush()?;
}
std::fs::rename(&tmp, path)?;
Ok(())
}
}
fn temp_sibling(path: &Path) -> PathBuf {
let mut name = path
.file_name()
.map(|n| n.to_os_string())
.unwrap_or_default();
name.push(format!(".tmp-{}", Uuid::new_v4()));
match path.parent() {
Some(dir) => dir.join(name),
None => PathBuf::from(name),
}
}
#[cfg(test)]
mod trim_tests {
use super::*;
use kindling_types::{ObservationKind, ScopeIds};
fn entry(content: &str, spooled_at: Option<i64>) -> SpoolEntry {
SpoolEntry {
input: ObservationInput {
id: Some(format!("id-{content}")),
kind: ObservationKind::Message,
content: content.to_string(),
provenance: None,
ts: None,
scope_ids: ScopeIds::default(),
redacted: None,
},
capsule_id: None,
validate: None,
spooled_at,
}
}
fn contents(entries: &[SpoolEntry]) -> Vec<String> {
entries.iter().map(|e| e.input.content.clone()).collect()
}
fn assert_is_suffix(original: &[SpoolEntry], kept: &[SpoolEntry]) {
let orig = contents(original);
let kept = contents(kept);
assert!(
orig.ends_with(&kept),
"kept {kept:?} is not a suffix of {orig:?}"
);
}
#[test]
fn no_caps_is_a_noop() {
let input = vec![entry("a", Some(1)), entry("b", Some(2))];
let kept = trim_entries(input.clone(), None, None, 1_000);
assert_eq!(contents(&kept), contents(&input));
}
#[test]
fn empty_spool_is_a_noop() {
let kept = trim_entries(Vec::new(), Some(10), Some(10), 1_000);
assert!(kept.is_empty());
}
#[test]
fn age_drops_oldest_prefix_only() {
let input = vec![
entry("a", Some(500)), entry("b", Some(800)), entry("c", Some(950)), entry("d", Some(990)), ];
let kept = trim_entries(input.clone(), None, Some(100), 1_000);
assert_eq!(contents(&kept), vec!["c", "d"]);
assert_is_suffix(&input, &kept);
}
#[test]
fn age_stops_at_first_unexpired_even_if_later_ones_are_old() {
let input = vec![
entry("a", Some(500)), entry("b", Some(950)), entry("c", Some(400)), ];
let kept = trim_entries(input.clone(), None, Some(100), 1_000);
assert_eq!(contents(&kept), vec!["b", "c"]);
assert_is_suffix(&input, &kept);
}
#[test]
fn legacy_entry_without_stamp_is_never_age_trimmed() {
let input = vec![
entry("a", None), entry("b", Some(100)), ];
let kept = trim_entries(input.clone(), None, Some(100), 1_000_000);
assert_eq!(contents(&kept), vec!["a", "b"]);
}
#[test]
fn bytes_drops_oldest_until_under_cap() {
let input = vec![entry("a", None), entry("b", None), entry("c", None)];
let each = entry_size(&input[0]);
let kept = trim_entries(input.clone(), Some(each * 2), None, 0);
assert_eq!(contents(&kept), vec!["b", "c"]);
assert_is_suffix(&input, &kept);
let total: u64 = kept.iter().map(entry_size).sum();
assert!(total <= each * 2);
}
#[test]
fn lone_entry_larger_than_cap_is_retained() {
let big = entry(&"x".repeat(4096), None);
let cap = entry_size(&big) / 2;
let kept = trim_entries(vec![big.clone()], Some(cap), None, 0);
assert_eq!(contents(&kept), vec![big.input.content]);
}
#[test]
fn bytes_keeps_at_least_one_even_when_all_oversize() {
let input = vec![
entry(&"x".repeat(1000), None),
entry(&"y".repeat(1000), None),
];
let kept = trim_entries(input.clone(), Some(10), None, 0);
assert_eq!(contents(&kept), vec!["y".repeat(1000)]);
}
#[test]
fn age_then_bytes_compose_into_a_suffix() {
let input = vec![
entry("a", Some(100)), entry("b", Some(200)), entry("c", Some(950)), entry("d", Some(960)), entry("e", Some(970)), ];
let each = entry_size(&input[0]);
let kept = trim_entries(input.clone(), Some(each * 2), Some(100), 1_000);
assert_eq!(contents(&kept), vec!["d", "e"]);
assert_is_suffix(&input, &kept);
}
}