1use corium_core::{
4 Datom, EntityId,
5 encoding::{decode_value, encode_value},
6};
7use std::{
8 fs::{self, File, OpenOptions},
9 io::{self, Read, Write},
10 path::{Path, PathBuf},
11 sync::{Arc, RwLock},
12};
13use thiserror::Error;
14
15#[derive(Clone, Debug, Eq, PartialEq)]
17pub struct TxRecord {
18 pub t: u64,
20 pub tx_instant: i64,
22 pub datoms: Vec<Datom>,
24}
25
26#[derive(Debug, Error)]
28pub enum LogError {
29 #[error("log I/O failed: {0}")]
31 Io(#[from] io::Error),
32 #[error("corrupt transaction log")]
34 Corrupt,
35}
36
37pub trait TransactionLog: Send + Sync {
39 fn append(&self, record: &TxRecord) -> Result<(), LogError>;
44 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
49 fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
54 self.tx_range(0, None)
55 }
56}
57
58#[derive(Clone, Default)]
60pub struct MemoryLog(Arc<RwLock<Vec<TxRecord>>>);
61impl TransactionLog for MemoryLog {
62 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
63 let mut records = self.0.write().expect("poisoned log lock");
64 if records.last().map_or(1, |r| r.t + 1) != record.t {
65 return Err(LogError::Corrupt);
66 }
67 records.push(record.clone());
68 Ok(())
69 }
70 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
71 Ok(self
72 .0
73 .read()
74 .expect("poisoned log lock")
75 .iter()
76 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
77 .cloned()
78 .collect())
79 }
80}
81
82pub struct FileLog {
88 path: PathBuf,
89 next_t: RwLock<u64>,
90}
91impl FileLog {
92 pub fn open(path: impl AsRef<Path>) -> Result<Self, LogError> {
98 let path = path.as_ref().to_path_buf();
99 if let Some(parent) = path.parent() {
100 fs::create_dir_all(parent)?;
101 }
102 OpenOptions::new().create(true).append(true).open(&path)?;
103 let (records, durable_len) = read_records(&path)?;
104 if fs::metadata(&path)?.len() > durable_len {
105 let file = OpenOptions::new().write(true).open(&path)?;
106 file.set_len(durable_len)?;
107 file.sync_all()?;
108 }
109 Ok(Self {
110 path,
111 next_t: RwLock::new(records.last().map_or(1, |r| r.t + 1)),
112 })
113 }
114}
115impl TransactionLog for FileLog {
116 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
117 let mut next_t = self.next_t.write().expect("poisoned log lock");
118 if *next_t != record.t {
119 return Err(LogError::Corrupt);
120 }
121 let payload = encode_record(record);
122 let mut file = OpenOptions::new().append(true).open(&self.path)?;
123 file.write_all(
124 &u64::try_from(payload.len())
125 .map_err(|_| LogError::Corrupt)?
126 .to_be_bytes(),
127 )?;
128 file.write_all(&payload)?;
129 file.sync_all()?;
130 *next_t += 1;
131 Ok(())
132 }
133 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
134 let _guard = self.next_t.read().expect("poisoned log lock");
135 Ok(read_records(&self.path)?
136 .0
137 .into_iter()
138 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
139 .collect())
140 }
141}
142
143fn encode_record(record: &TxRecord) -> Vec<u8> {
144 let mut out = Vec::new();
145 out.extend_from_slice(&record.t.to_be_bytes());
146 out.extend_from_slice(&record.tx_instant.to_be_bytes());
147 out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
148 for d in &record.datoms {
149 out.extend_from_slice(&d.e.raw().to_be_bytes());
150 out.extend_from_slice(&d.a.raw().to_be_bytes());
151 out.extend_from_slice(&d.tx.raw().to_be_bytes());
152 out.push(u8::from(d.added));
153 let v = encode_value(&d.v);
154 out.extend_from_slice(&(v.len() as u64).to_be_bytes());
155 out.extend_from_slice(&v);
156 }
157 out
158}
159fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
160 fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
161 let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
162 *bytes = &bytes[n..];
163 Ok(value)
164 }
165 fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
166 Ok(u64::from_be_bytes(
167 take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
168 ))
169 }
170 let t = u64_be(&mut bytes)?;
171 let tx_instant = i64::from_be_bytes(
172 take(&mut bytes, 8)?
173 .try_into()
174 .map_err(|_| LogError::Corrupt)?,
175 );
176 let count = u64_be(&mut bytes)?;
177 let mut datoms = Vec::new();
178 for _ in 0..count {
179 let e = EntityId::from_raw(u64_be(&mut bytes)?);
180 let a = EntityId::from_raw(u64_be(&mut bytes)?);
181 let tx = EntityId::from_raw(u64_be(&mut bytes)?);
182 let added = take(&mut bytes, 1)?[0] != 0;
183 let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
184 let raw = take(&mut bytes, len)?;
185 let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
186 if used != len {
187 return Err(LogError::Corrupt);
188 }
189 datoms.push(Datom { e, a, v, tx, added });
190 }
191 if !bytes.is_empty() {
192 return Err(LogError::Corrupt);
193 }
194 Ok(TxRecord {
195 t,
196 tx_instant,
197 datoms,
198 })
199}
200fn read_records(path: &Path) -> Result<(Vec<TxRecord>, u64), LogError> {
206 let mut file = File::open(path)?;
207 let mut records = Vec::new();
208 let mut durable_len = 0_u64;
209 loop {
210 let mut len = [0; 8];
211 match file.read_exact(&mut len) {
212 Ok(()) => {}
213 Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
214 Err(e) => return Err(e.into()),
215 }
216 let len = usize::try_from(u64::from_be_bytes(len)).map_err(|_| LogError::Corrupt)?;
217 let mut payload = vec![0; len];
218 match file.read_exact(&mut payload) {
219 Ok(()) => {}
220 Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
221 Err(e) => return Err(e.into()),
222 }
223 records.push(decode_record(&payload)?);
224 durable_len += 8 + len as u64;
225 }
226 Ok((records, durable_len))
227}