use std::collections::HashMap;
use std::path::{Path, PathBuf};
use serde::{Deserialize, Serialize};
use super::destination::LogDestination;
use super::error::DrainError;
pub const MANIFEST_FILENAME: &str = ".drain-manifest.json";
pub const MANIFEST_VERSION: u32 = 1;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ManifestEntry {
pub relative_file: String,
pub size: u64,
pub mtime_unix: i64,
pub sha256: String,
pub uploaded_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DrainManifest {
pub version: u32,
pub entries: Vec<ManifestEntry>,
}
impl Default for DrainManifest {
fn default() -> Self {
Self {
version: MANIFEST_VERSION,
entries: Vec::new(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StatDecision {
SkipUnchanged,
NeedsHash,
}
impl DrainManifest {
fn index(&self) -> HashMap<&str, &ManifestEntry> {
self.entries
.iter()
.map(|e| (e.relative_file.as_str(), e))
.collect()
}
pub fn decide(&self, relative_file: &str, size: u64, mtime_unix: i64) -> StatDecision {
match self.index().get(relative_file) {
Some(entry) if entry.size == size && entry.mtime_unix == mtime_unix => {
StatDecision::SkipUnchanged
}
_ => StatDecision::NeedsHash,
}
}
pub fn digest_matches(&self, relative_file: &str, sha256: &str) -> bool {
self.index()
.get(relative_file)
.is_some_and(|entry| entry.sha256 == sha256)
}
pub fn record(&mut self, entry: ManifestEntry) {
match self
.entries
.iter_mut()
.find(|e| e.relative_file == entry.relative_file)
{
Some(existing) => *existing = entry,
None => self.entries.push(entry),
}
self.entries
.sort_by(|a, b| a.relative_file.cmp(&b.relative_file));
}
pub async fn load(
dest: &dyn LogDestination,
state_dir: &Path,
remote_key: &str,
cache_key: &str,
) -> Result<Self, DrainError> {
let cache_path = cache_file(state_dir, cache_key);
if let Some(raw) = dest.get(remote_key).await? {
match Self::decode(&raw) {
Ok(remote) => {
write_cache(&cache_path, &remote);
return Ok(remote);
}
Err(reason) => {
tracing::warn!(
key = %remote_key,
%reason,
"log-drain manifest is undecodable; treating as absent and re-uploading"
);
}
}
}
Ok(read_cache(&cache_path).unwrap_or_default())
}
pub async fn save(
&self,
dest: &dyn LogDestination,
state_dir: &Path,
remote_key: &str,
cache_key: &str,
) -> Result<(), DrainError> {
let body = serde_json::to_vec_pretty(self).map_err(|e| DrainError::Manifest {
key: remote_key.to_string(),
reason: format!("could not serialise: {e}"),
})?;
dest.put(
remote_key,
bytes::Bytes::from(body),
super::destination::PutMeta {
content_type: Some("application/json".to_string()),
content_encoding: None,
},
)
.await?;
write_cache(&cache_file(state_dir, cache_key), self);
Ok(())
}
fn decode(raw: &[u8]) -> Result<Self, String> {
let parsed: Self = serde_json::from_slice(raw).map_err(|e| e.to_string())?;
if parsed.version != MANIFEST_VERSION {
return Err(format!(
"schema version {} is not the supported {MANIFEST_VERSION}",
parsed.version
));
}
Ok(parsed)
}
}
fn cache_file(state_dir: &Path, cache_key: &str) -> PathBuf {
state_dir
.join("log-drain")
.join(cache_key)
.join("manifest.json")
}
fn read_cache(path: &Path) -> Option<DrainManifest> {
let raw = std::fs::read(path).ok()?;
DrainManifest::decode(&raw).ok()
}
fn write_cache(path: &Path, manifest: &DrainManifest) {
let Some(parent) = path.parent() else { return };
if let Err(e) = std::fs::create_dir_all(parent) {
tracing::warn!(path = %parent.display(), error = %e, "log-drain cache dir unwritable");
return;
}
match serde_json::to_vec_pretty(manifest) {
Ok(body) => {
if let Err(e) = std::fs::write(path, body) {
tracing::warn!(path = %path.display(), error = %e, "log-drain cache write failed");
}
}
Err(e) => tracing::warn!(error = %e, "log-drain cache serialisation failed"),
}
}