xagent-pi 0.2.3

Self-contained local brain (chat UI + API + SSE) for the Pi agent, tunneled into xagent-service.
//! Disk persistence for sessions + messages. One JSON file per session.

use anyhow::Result;
use serde::{Deserialize, Serialize};
use std::path::PathBuf;
#[cfg(test)]
use std::path::Path;

#[derive(Clone, Default, Serialize, Deserialize)]
pub struct Msg {
    pub role: String,
    pub text: String,
    pub ts: u64,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub call_id: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub tool_name: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub input: Option<serde_json::Value>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub result: Option<serde_json::Value>,
    #[serde(default)]
    pub is_error: bool,
}

#[derive(Clone, Serialize, Deserialize)]
pub struct SessionMeta {
    pub id: String,
    #[serde(default = "default_project_id")]
    pub project_id: String,
    pub title: String,
    pub created_at: u64,
    pub updated_at: u64,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub worktree: Option<WorktreeMeta>,
}

#[derive(Clone, Serialize, Deserialize)]
pub struct WorktreeMeta {
    pub base_path: String,
    pub worktree_path: String,
    pub branch: String,
    pub name: String,
    pub created_at: u64,
}

fn default_project_id() -> String {
    "default".to_string()
}

#[derive(Serialize, Deserialize)]
struct SessionFile {
    meta: SessionMeta,
    messages: Vec<Msg>,
}

#[derive(Clone)]
pub struct Store {
    dir: PathBuf,
}

pub fn now_ms() -> u64 {
    use std::time::{SystemTime, UNIX_EPOCH};
    SystemTime::now()
        .duration_since(UNIX_EPOCH)
        .map(|d| d.as_millis() as u64)
        .unwrap_or(0)
}

impl Store {
    pub fn new(dir: PathBuf) -> Result<Self> {
        std::fs::create_dir_all(&dir)?;
        Ok(Self { dir })
    }

    fn path(&self, id: &str) -> PathBuf {
        self.dir.join(format!("{id}.json"))
    }

    pub fn list(&self) -> Vec<SessionMeta> {
        let mut out = Vec::new();
        if let Ok(rd) = std::fs::read_dir(&self.dir) {
            for entry in rd.flatten() {
                let p = entry.path();
                if p.extension().and_then(|e| e.to_str()) != Some("json") {
                    continue;
                }
                if let Ok(bytes) = std::fs::read(&p) {
                    if let Ok(sf) = serde_json::from_slice::<SessionFile>(&bytes) {
                        out.push(sf.meta);
                    }
                }
            }
        }
        out.sort_by(|a, b| b.updated_at.cmp(&a.updated_at));
        out
    }

    pub fn load(&self, id: &str) -> Option<(SessionMeta, Vec<Msg>)> {
        let bytes = std::fs::read(self.path(id)).ok()?;
        let sf: SessionFile = serde_json::from_slice(&bytes).ok()?;
        Some((sf.meta, sf.messages))
    }

    pub fn save(&self, meta: &SessionMeta, messages: &[Msg]) -> Result<()> {
        let sf = SessionFile {
            meta: meta.clone(),
            messages: messages.to_vec(),
        };
        let bytes = serde_json::to_vec_pretty(&sf)?;
        let final_path = self.path(&meta.id);
        let tmp = final_path.with_extension("json.tmp");
        std::fs::write(&tmp, bytes)?;
        std::fs::rename(&tmp, &final_path)?;
        Ok(())
    }

    pub fn delete(&self, id: &str) -> Result<()> {
        let _ = std::fs::remove_file(self.path(id));
        Ok(())
    }

    #[cfg(test)]
    pub fn dir(&self) -> &Path {
        &self.dir
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    fn temp_dir() -> PathBuf {
        let d = std::env::temp_dir().join(format!(
            "xagent-pi-store-test-{}-{}",
            std::process::id(),
            uuid::Uuid::new_v4()
        ));
        std::fs::create_dir_all(&d).unwrap();
        d
    }

    fn meta(id: &str, updated_at: u64) -> SessionMeta {
        SessionMeta {
            id: id.to_string(),
            project_id: default_project_id(),
            title: format!("session {id}"),
            created_at: 1,
            updated_at,
            worktree: None,
        }
    }

    fn msgs() -> Vec<Msg> {
        vec![
            Msg {
                role: "user".into(),
                text: "hi".into(),
                ts: 10,
                ..Msg::default()
            },
            Msg {
                role: "assistant".into(),
                text: "hello".into(),
                ts: 11,
                ..Msg::default()
            },
        ]
    }

    #[test]
    fn new_creates_directory() {
        let dir = std::env::temp_dir().join(format!(
            "xagent-pi-store-mkdir-{}",
            uuid::Uuid::new_v4()
        ));
        let store = Store::new(dir.clone()).unwrap();
        assert!(store.dir().is_dir());
        std::fs::remove_dir_all(&dir).unwrap();
    }

    #[test]
    fn save_then_load_roundtrip() {
        let dir = temp_dir();
        let store = Store::new(dir.clone()).unwrap();
        let m = meta("abc", 100);
        let msgs = msgs();
        store.save(&m, &msgs).unwrap();

        let (loaded_meta, loaded_msgs) = store.load("abc").expect("session should exist");
        assert_eq!(loaded_meta.id, m.id);
        assert_eq!(loaded_meta.title, m.title);
        assert_eq!(loaded_meta.updated_at, m.updated_at);
        assert_eq!(loaded_msgs.len(), 2);
        assert_eq!(loaded_msgs[0].role, "user");
        assert_eq!(loaded_msgs[0].text, "hi");
        assert_eq!(loaded_msgs[1].text, "hello");
        std::fs::remove_dir_all(&dir).unwrap();
    }

    #[test]
    fn load_missing_returns_none() {
        let dir = temp_dir();
        let store = Store::new(dir.clone()).unwrap();
        assert!(store.load("nope").is_none());
        std::fs::remove_dir_all(&dir).unwrap();
    }

    #[test]
    fn list_sorts_by_updated_at_desc() {
        let dir = temp_dir();
        let store = Store::new(dir.clone()).unwrap();
        store.save(&meta("old", 100), &[]).unwrap();
        store.save(&meta("new", 300), &[]).unwrap();
        store.save(&meta("mid", 200), &[]).unwrap();

        let list = store.list();
        let ids: Vec<&str> = list.iter().map(|m| m.id.as_str()).collect();
        assert_eq!(ids, vec!["new", "mid", "old"]);
        std::fs::remove_dir_all(&dir).unwrap();
    }

    #[test]
    fn list_ignores_non_json_and_corrupt_files() {
        let dir = temp_dir();
        let store = Store::new(dir.clone()).unwrap();
        store.save(&meta("good", 50), &[]).unwrap();
        std::fs::write(dir.join("not-json.json"), "{{{{ not json").unwrap();
        std::fs::write(dir.join("notes.txt"), "irrelevant").unwrap();

        let list = store.list();
        assert_eq!(list.len(), 1);
        assert_eq!(list[0].id, "good");
        std::fs::remove_dir_all(&dir).unwrap();
    }

    #[test]
    fn delete_removes_session_file() {
        let dir = temp_dir();
        let store = Store::new(dir.clone()).unwrap();
        store.save(&meta("gone", 1), &[]).unwrap();
        assert!(store.load("gone").is_some());
        store.delete("gone").unwrap();
        assert!(store.load("gone").is_none());
        // deleting a missing id is not an error
        store.delete("never-existed").unwrap();
        std::fs::remove_dir_all(&dir).unwrap();
    }

    #[test]
    fn save_leaves_no_temp_file() {
        let dir = temp_dir();
        let store = Store::new(dir.clone()).unwrap();
        store.save(&meta("atomic", 7), &[]).unwrap();
        let leftovers: Vec<_> = std::fs::read_dir(&dir)
            .unwrap()
            .flatten()
            .filter(|e| e.path().extension().and_then(|x| x.to_str()) == Some("tmp"))
            .collect();
        assert!(leftovers.is_empty(), "tmp files left after save: {leftovers:?}");
        std::fs::remove_dir_all(&dir).unwrap();
    }

    #[test]
    fn overwrite_updates_existing_session() {
        let dir = temp_dir();
        let store = Store::new(dir.clone()).unwrap();
        store.save(&meta("same", 1), &msgs()).unwrap();
        let extra = vec![Msg {
            role: "assistant".into(),
            text: "bye".into(),
            ts: 12,
            ..Msg::default()
        }];
        store.save(&meta("same", 2), &extra).unwrap();
        let (_, loaded) = store.load("same").unwrap();
        assert_eq!(loaded.len(), 1);
        assert_eq!(loaded[0].text, "bye");
        std::fs::remove_dir_all(&dir).unwrap();
    }
}