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();
}
}