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}