kcode-k1-chat-persistence-store 0.3.0

Filesystem JSONL projection for K1 chat persistence records
Documentation
pub use kcode_k1_chat_persistence_records::{
    Batch, BoxId, CODEC_VERSION, Error, EventRecord, Record, SessionId, SessionLog, ToolCallId,
    TxId,
};
use kcode_k1_chat_persistence_records::{
    encode_line, locate_predecessor, strict_prefix, to_log, validate_session_records,
};
use std::fs::{self, File, OpenOptions};
use std::io::{self, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};

type Result<T, E = Error> = std::result::Result<T, E>;

pub struct Projection {
    root: PathBuf,
    sessions: PathBuf,
}

impl Projection {
    pub fn new(root: impl AsRef<Path>) -> Result<Self> {
        let root = root.as_ref().to_path_buf();
        fs::create_dir_all(&root)?;
        let sessions = root.join("sessions");
        fs::create_dir_all(&sessions)?;
        Ok(Self { root, sessions })
    }
    pub fn load(&self, session: SessionId) -> Result<SessionLog> {
        let records = self
            .read_strict(session)?
            .into_iter()
            .map(|value| value.1)
            .collect::<Vec<_>>();
        validate_session_records(session, &records)?;
        Ok(to_log(records))
    }
    pub fn apply(&mut self, txid: TxId, batch: &Batch, reconcile_first: bool) -> Result<()> {
        batch.encode()?;
        fs::create_dir_all(&self.sessions)?;
        let path = self.path(batch.session_id);
        if reconcile_first {
            let bytes = match fs::read(&path) {
                Ok(value) => value,
                Err(error) if error.kind() == io::ErrorKind::NotFound => Vec::new(),
                Err(error) => return Err(error.into()),
            };
            let end = locate_predecessor(&bytes, batch.predecessor.as_ref())?;
            let mut combined = strict_prefix(&bytes[..end])?
                .into_iter()
                .map(|value| value.1)
                .collect::<Vec<_>>();
            combined.extend(batch.records.clone());
            validate_session_records(batch.session_id, &combined)?;
            let mut file = OpenOptions::new()
                .create(true)
                .read(true)
                .write(true)
                .truncate(false)
                .open(&path)?;
            file.set_len(end as u64)?;
            file.seek(SeekFrom::Start(end as u64))?;
            append_lines(&mut file, txid, &batch.records)?;
            file.sync_all()?;
        } else {
            let current = self.read_strict(batch.session_id)?;
            let mut combined = current
                .iter()
                .map(|value| value.1.clone())
                .collect::<Vec<_>>();
            if combined.last() != batch.predecessor.as_ref() {
                return Err(Error::Invalid("predecessor mismatch"));
            }
            combined.extend(batch.records.clone());
            validate_session_records(batch.session_id, &combined)?;
            let mut file = OpenOptions::new().create(true).append(true).open(&path)?;
            append_lines(&mut file, txid, &batch.records)?;
            file.sync_all()?;
        }
        Ok(())
    }
    pub fn discard_all(&mut self) -> Result<()> {
        match fs::remove_dir_all(&self.sessions) {
            Ok(()) => {}
            Err(error) if error.kind() == io::ErrorKind::NotFound => {}
            Err(error) => return Err(error.into()),
        }
        File::open(&self.root)?.sync_all()?;
        Ok(())
    }
    fn path(&self, session: SessionId) -> PathBuf {
        self.sessions.join(format!("{}.jsonl", hex(session)))
    }
    fn read_strict(&self, session: SessionId) -> Result<Vec<(TxId, Record, usize)>> {
        let bytes = match fs::read(self.path(session)) {
            Ok(value) => value,
            Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()),
            Err(error) => return Err(error.into()),
        };
        strict_prefix(&bytes)
    }
}

fn append_lines(file: &mut File, txid: TxId, records: &[Record]) -> Result<()> {
    for record in records {
        file.write_all(&encode_line(txid, record)?)?;
    }
    Ok(())
}
fn hex(bytes: [u8; 12]) -> String {
    bytes.iter().map(|value| format!("{value:02x}")).collect()
}

#[cfg(test)]
mod tests {
    use super::*;
    use std::sync::atomic::{AtomicU64, Ordering};
    const FIRST_RECORD: &str = r#"{"kind":"box","id":1,"type":"System Message","contents":"s","hidden_type":"","hidden_contents":""}"#;
    fn batch(session: u8, predecessor: &str, records: &str) -> Batch {
        let session_id = format!("{session:02x}").repeat(12);
        Batch::decode(format!(r#"{{"version":2,"session_id":"{session_id}","predecessor":{predecessor},"records":[{records}]}}"#).as_bytes()).unwrap()
    }
    fn directory() -> PathBuf {
        static NEXT: AtomicU64 = AtomicU64::new(0);
        let path = std::env::temp_dir().join(format!(
            "k1-persistence-store-{}-{}",
            std::process::id(),
            NEXT.fetch_add(1, Ordering::Relaxed)
        ));
        let _ = fs::remove_dir_all(&path);
        path
    }
    fn tx(value: u8) -> TxId {
        TxId::from_bytes([value; 12])
    }
    #[test]
    fn records_v2_facade_exact_bytes_and_open_fields() {
        fn assert_type<T>() {}
        assert_type::<Batch>();
        assert_type::<EventRecord>();
        assert_type::<Record>();
        assert_type::<SessionId>();
        assert_type::<SessionLog>();
        assert_type::<TxId>();
        assert_type::<BoxId>();
        assert_type::<ToolCallId>();
        assert_type::<Error>();
        assert_eq!(CODEC_VERSION, 2);
        let records = format!(
            "{FIRST_RECORD},{{\"kind\":\"box\",\"id\":2,\"type\":\"Future Kind\",\"contents\":\"opaque\\ncontents\",\"hidden_type\":\"future/v9\",\"hidden_contents\":\"hidden bytes\"}}"
        );
        let value = batch(1, "null", &records);
        let expected_batch = br#"{"version":2,"session_id":"010101010101010101010101","predecessor":null,"records":[{"kind":"box","id":1,"type":"System Message","contents":"s","hidden_type":"","hidden_contents":""},{"kind":"box","id":2,"type":"Future Kind","contents":"opaque\ncontents","hidden_type":"future/v9","hidden_contents":"hidden bytes"}]}"#;
        assert_eq!(value.encode().unwrap(), expected_batch);
        let root = directory();
        let mut projection = Projection::new(&root).unwrap();
        projection.apply(tx(0xab), &value, false).unwrap();
        let expected_line = b"{\"txid\":\"abababababababababababab\",\"record\":{\"kind\":\"box\",\"id\":1,\"type\":\"System Message\",\"contents\":\"s\",\"hidden_type\":\"\",\"hidden_contents\":\"\"}}\n{\"txid\":\"abababababababababababab\",\"record\":{\"kind\":\"box\",\"id\":2,\"type\":\"Future Kind\",\"contents\":\"opaque\\ncontents\",\"hidden_type\":\"future/v9\",\"hidden_contents\":\"hidden bytes\"}}\n";
        assert_eq!(fs::read(projection.path([1; 12])).unwrap(), expected_line);
        fs::remove_dir_all(root).unwrap();
    }
    #[test]
    fn append_load_reconcile_and_discard() {
        let root = directory();
        let mut projection = Projection::new(&root).unwrap();
        let first = batch(4, "null", FIRST_RECORD);
        projection.apply(tx(1), &first, false).unwrap();
        let records = r#"{"kind":"event","after_box_id":1,"event_index":1,"connected_box_id":1,"handler":"h","data":null},{"kind":"box","id":2,"type":"User Message","contents":"u","hidden_type":"","hidden_contents":""}"#;
        let next = batch(4, FIRST_RECORD, records);
        projection.apply(tx(2), &next, false).unwrap();
        projection.apply(tx(2), &next, true).unwrap();
        assert_eq!(projection.load([4; 12]).unwrap().records.len(), 3);
        let mut file = OpenOptions::new()
            .append(true)
            .open(projection.path([4; 12]))
            .unwrap();
        file.write_all(b"{partial").unwrap();
        projection.apply(tx(2), &next, true).unwrap();
        assert_eq!(projection.load([4; 12]).unwrap().records.len(), 3);
        assert!(projection.load([5; 12]).unwrap().records.is_empty());
        projection.discard_all().unwrap();
        assert!(projection.load([4; 12]).unwrap().records.is_empty());
        fs::remove_dir_all(root).unwrap();
    }
    #[test]
    fn malformed_lines_and_predecessors_fail_closed() {
        let root = directory();
        let mut projection = Projection::new(&root).unwrap();
        fs::write(projection.path([1; 12]), b"{}\n").unwrap();
        assert!(projection.load([1; 12]).is_err());
        fs::write(projection.path([2; 12]), b"{}").unwrap();
        assert!(projection.load([2; 12]).is_err());
        let missing = batch(
            3,
            r#"{"kind":"box","id":9,"type":"Attachment","contents":"","hidden_type":"","hidden_contents":""}"#,
            r#"{"kind":"box","id":10,"type":"Attachment","contents":"","hidden_type":"","hidden_contents":""}"#,
        );
        assert!(projection.apply(tx(3), &missing, true).is_err());
        fs::remove_dir_all(root).unwrap();
    }
}