use crate::error::{BallError, Result};
use crate::plugin::SyncReport;
use crate::store::Store;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use std::fs;
use std::path::PathBuf;
const PENDING_DIR: &str = "pending-sync";
const SYNC_EVENT: &str = "sync";
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StagedSync {
pub event: String,
pub plugin: String,
pub staged_at: DateTime<Utc>,
pub report: SyncReport,
}
#[derive(Debug)]
pub struct StagedEntry {
pub id: String,
pub path: PathBuf,
pub data: StagedSync,
}
fn pending_root(store: &Store) -> PathBuf {
store.local_dir().join(PENDING_DIR)
}
fn pending_dir(store: &Store, event: &str) -> PathBuf {
pending_root(store).join(event)
}
fn make_stage_id(plugin: &str, ts: DateTime<Utc>) -> String {
crate::task::Task::generate_id(plugin, ts, 4)
}
pub fn stage_sync(store: &Store, plugin: &str, report: &SyncReport) -> Result<String> {
stage_sync_at(store, plugin, report, Utc::now())
}
fn stage_sync_at(
store: &Store,
plugin: &str,
report: &SyncReport,
start: DateTime<Utc>,
) -> Result<String> {
let dir = pending_dir(store, SYNC_EVENT);
fs::create_dir_all(&dir)?;
let mut now = start;
let mut id = make_stage_id(plugin, now);
let mut path = dir.join(format!("{id}.json"));
let mut tries = 0;
while path.exists() {
tries += 1;
if tries > 1000 {
return Err(BallError::Other(
"could not allocate unique stage id".into(),
));
}
now += chrono::Duration::milliseconds(1);
id = make_stage_id(plugin, now);
path = dir.join(format!("{id}.json"));
}
let payload = StagedSync {
event: SYNC_EVENT.into(),
plugin: plugin.into(),
staged_at: now,
report: report.clone(),
};
fs::write(&path, serde_json::to_vec_pretty(&payload)?)?;
Ok(id)
}
pub fn load_staged(store: &Store, id: &str) -> Result<StagedEntry> {
let root = pending_root(store);
if !root.exists() {
return Err(BallError::Other(format!("no staged entry: {id}")));
}
for entry in fs::read_dir(&root)? {
let event_dir = entry?.path();
if !event_dir.is_dir() {
continue;
}
let candidate = event_dir.join(format!("{id}.json"));
if candidate.exists() {
let bytes = fs::read(&candidate)?;
let data: StagedSync = serde_json::from_slice(&bytes)?;
return Ok(StagedEntry { id: id.into(), path: candidate, data });
}
}
Err(BallError::Other(format!("no staged entry: {id}")))
}
pub fn list_staged(store: &Store) -> Result<Vec<StagedEntry>> {
let root = pending_root(store);
if !root.exists() {
return Ok(Vec::new());
}
let mut out = Vec::new();
for entry in fs::read_dir(&root)? {
let event_dir = entry?.path();
if !event_dir.is_dir() {
continue;
}
for f in fs::read_dir(&event_dir)? {
let p = f?.path();
if p.extension().and_then(|s| s.to_str()) != Some("json") {
continue;
}
let id = p
.file_stem()
.and_then(|s| s.to_str())
.unwrap_or_default()
.to_string();
let bytes = fs::read(&p)?;
let data: StagedSync = serde_json::from_slice(&bytes)?;
out.push(StagedEntry { id, path: p, data });
}
}
out.sort_by(|a, b| a.id.cmp(&b.id));
Ok(out)
}
pub fn discard_staged(store: &Store, id: &str) -> Result<(String, String)> {
let entry = load_staged(store, id)?;
fs::remove_file(&entry.path)?;
Ok((entry.data.event, entry.data.plugin))
}
#[cfg(test)]
#[path = "human_gate_tests.rs"]
mod tests;