Skip to main content

kcode_k1_chat_persistence_store/
lib.rs

1use kcode_k1_chat_persistence_jsonl::{encode_line, locate_predecessor, strict_prefix};
2pub use kcode_k1_chat_persistence_records::{
3    Batch, BoxId, CODEC_VERSION, Error, EventRecord, Record, SessionId, SessionLog, ToolCallId,
4    TxId,
5};
6use kcode_k1_chat_persistence_records::{to_log, validate_session_records};
7use std::fs::{self, File, OpenOptions};
8use std::io::{self, Seek, SeekFrom, Write};
9use std::path::{Path, PathBuf};
10
11type Result<T, E = Error> = std::result::Result<T, E>;
12
13pub struct Projection {
14    root: PathBuf,
15    sessions: PathBuf,
16}
17
18impl Projection {
19    pub fn new(root: impl AsRef<Path>) -> Result<Self> {
20        let root = root.as_ref().to_path_buf();
21        fs::create_dir_all(&root)?;
22        let sessions = root.join("sessions");
23        fs::create_dir_all(&sessions)?;
24        Ok(Self { root, sessions })
25    }
26    pub fn load(&self, session: SessionId) -> Result<SessionLog> {
27        let records = self
28            .read_strict(session)?
29            .into_iter()
30            .map(|value| value.1)
31            .collect::<Vec<_>>();
32        validate_session_records(session, &records)?;
33        Ok(to_log(records))
34    }
35    pub fn apply(&mut self, txid: TxId, batch: &Batch, reconcile_first: bool) -> Result<()> {
36        batch.encode()?;
37        fs::create_dir_all(&self.sessions)?;
38        let path = self.path(batch.session_id);
39        if reconcile_first {
40            let bytes = match fs::read(&path) {
41                Ok(value) => value,
42                Err(error) if error.kind() == io::ErrorKind::NotFound => Vec::new(),
43                Err(error) => return Err(error.into()),
44            };
45            let end = locate_predecessor(&bytes, batch.predecessor.as_ref())?;
46            let mut combined = strict_prefix(&bytes[..end])?
47                .into_iter()
48                .map(|value| value.1)
49                .collect::<Vec<_>>();
50            combined.extend(batch.records.clone());
51            validate_session_records(batch.session_id, &combined)?;
52            let mut file = OpenOptions::new()
53                .create(true)
54                .read(true)
55                .write(true)
56                .truncate(false)
57                .open(&path)?;
58            file.set_len(end as u64)?;
59            file.seek(SeekFrom::Start(end as u64))?;
60            append_lines(&mut file, txid, &batch.records)?;
61            file.sync_all()?;
62        } else {
63            let current = self.read_strict(batch.session_id)?;
64            let mut combined = current
65                .iter()
66                .map(|value| value.1.clone())
67                .collect::<Vec<_>>();
68            if combined.last() != batch.predecessor.as_ref() {
69                return Err(Error::Invalid("predecessor mismatch"));
70            }
71            combined.extend(batch.records.clone());
72            validate_session_records(batch.session_id, &combined)?;
73            let mut file = OpenOptions::new().create(true).append(true).open(&path)?;
74            append_lines(&mut file, txid, &batch.records)?;
75            file.sync_all()?;
76        }
77        Ok(())
78    }
79    pub fn discard_all(&mut self) -> Result<()> {
80        match fs::remove_dir_all(&self.sessions) {
81            Ok(()) => {}
82            Err(error) if error.kind() == io::ErrorKind::NotFound => {}
83            Err(error) => return Err(error.into()),
84        }
85        File::open(&self.root)?.sync_all()?;
86        Ok(())
87    }
88    fn path(&self, session: SessionId) -> PathBuf {
89        self.sessions.join(format!("{}.jsonl", hex(session)))
90    }
91    fn read_strict(&self, session: SessionId) -> Result<Vec<(TxId, Record, usize)>> {
92        let bytes = match fs::read(self.path(session)) {
93            Ok(value) => value,
94            Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()),
95            Err(error) => return Err(error.into()),
96        };
97        strict_prefix(&bytes)
98    }
99}
100
101fn append_lines(file: &mut File, txid: TxId, records: &[Record]) -> Result<()> {
102    for record in records {
103        file.write_all(&encode_line(txid, record)?)?;
104    }
105    Ok(())
106}
107fn hex(bytes: [u8; 12]) -> String {
108    bytes.iter().map(|value| format!("{value:02x}")).collect()
109}
110
111#[cfg(test)]
112mod tests {
113    use super::*;
114    use kcode_k1_chat_chatend::AGENT_RESPONSE_TYPE;
115    use std::sync::atomic::{AtomicU64, Ordering};
116    const FIRST_RECORD: &str = r#"{"kind":"box","id":1,"type":"System Message","contents":"s","hidden_type":"","hidden_contents":""}"#;
117    fn batch(session: u8, predecessor: &str, records: &str) -> Batch {
118        let session_id = format!("{session:02x}").repeat(12);
119        Batch::decode(format!(r#"{{"version":2,"session_id":"{session_id}","predecessor":{predecessor},"records":[{records}]}}"#).as_bytes()).unwrap()
120    }
121    fn directory() -> PathBuf {
122        static NEXT: AtomicU64 = AtomicU64::new(0);
123        let path = std::env::temp_dir().join(format!(
124            "k1-persistence-store-{}-{}",
125            std::process::id(),
126            NEXT.fetch_add(1, Ordering::Relaxed)
127        ));
128        let _ = fs::remove_dir_all(&path);
129        path
130    }
131    fn tx(value: u8) -> TxId {
132        TxId::from_bytes([value; 12])
133    }
134    #[test]
135    fn records_v2_facade_exact_bytes_and_open_fields() {
136        fn assert_type<T>() {}
137        assert_type::<Batch>();
138        assert_type::<EventRecord>();
139        assert_type::<Record>();
140        assert_type::<SessionId>();
141        assert_type::<SessionLog>();
142        assert_type::<TxId>();
143        assert_type::<BoxId>();
144        assert_type::<ToolCallId>();
145        assert_type::<Error>();
146        assert_eq!(CODEC_VERSION, 2);
147        let records = format!(
148            "{FIRST_RECORD},{{\"kind\":\"box\",\"id\":2,\"type\":\"Future Kind\",\"contents\":\"opaque\\ncontents\",\"hidden_type\":\"future/v9\",\"hidden_contents\":\"hidden bytes\"}}"
149        );
150        let value = batch(1, "null", &records);
151        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"}]}"#;
152        assert_eq!(value.encode().unwrap(), expected_batch);
153        let root = directory();
154        let mut projection = Projection::new(&root).unwrap();
155        projection.apply(tx(0xab), &value, false).unwrap();
156        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";
157        assert_eq!(fs::read(projection.path([1; 12])).unwrap(), expected_line);
158        fs::remove_dir_all(root).unwrap();
159    }
160    #[test]
161    fn empty_agent_response_hidden_metadata_round_trips_generic_jsonl() {
162        let record = format!(
163            r#"{{"kind":"box","id":1,"type":"{AGENT_RESPONSE_TYPE}","contents":"","hidden_type":"reasoning/v1","hidden_contents":"opaque reasoning"}}"#
164        );
165        let value = batch(6, "null", &record);
166        let root = directory();
167        let mut projection = Projection::new(&root).unwrap();
168        projection.apply(tx(6), &value, false).unwrap();
169        let expected_line =
170            format!("{{\"txid\":\"060606060606060606060606\",\"record\":{record}}}\n");
171        assert_eq!(
172            fs::read(projection.path([6; 12])).unwrap(),
173            expected_line.as_bytes()
174        );
175        assert_eq!(projection.load([6; 12]).unwrap().records, value.records);
176        fs::remove_dir_all(root).unwrap();
177    }
178    #[test]
179    fn append_load_reconcile_and_discard() {
180        let root = directory();
181        let mut projection = Projection::new(&root).unwrap();
182        let first = batch(4, "null", FIRST_RECORD);
183        projection.apply(tx(1), &first, false).unwrap();
184        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":""}"#;
185        let next = batch(4, FIRST_RECORD, records);
186        projection.apply(tx(2), &next, false).unwrap();
187        projection.apply(tx(2), &next, true).unwrap();
188        assert_eq!(projection.load([4; 12]).unwrap().records.len(), 3);
189        let mut file = OpenOptions::new()
190            .append(true)
191            .open(projection.path([4; 12]))
192            .unwrap();
193        file.write_all(b"{partial").unwrap();
194        projection.apply(tx(2), &next, true).unwrap();
195        assert_eq!(projection.load([4; 12]).unwrap().records.len(), 3);
196        assert!(projection.load([5; 12]).unwrap().records.is_empty());
197        projection.discard_all().unwrap();
198        assert!(projection.load([4; 12]).unwrap().records.is_empty());
199        fs::remove_dir_all(root).unwrap();
200    }
201    #[test]
202    fn malformed_lines_and_predecessors_fail_closed() {
203        let root = directory();
204        let mut projection = Projection::new(&root).unwrap();
205        fs::write(projection.path([1; 12]), b"{}\n").unwrap();
206        assert!(projection.load([1; 12]).is_err());
207        fs::write(projection.path([2; 12]), b"{}").unwrap();
208        assert!(projection.load([2; 12]).is_err());
209        let missing = batch(
210            3,
211            r#"{"kind":"box","id":9,"type":"Attachment","contents":"","hidden_type":"","hidden_contents":""}"#,
212            r#"{"kind":"box","id":10,"type":"Attachment","contents":"","hidden_type":"","hidden_contents":""}"#,
213        );
214        assert!(projection.apply(tx(3), &missing, true).is_err());
215        fs::remove_dir_all(root).unwrap();
216    }
217}