1use alloc::vec::Vec;
29
30use xxhash_rust::xxh3::xxh3_64;
31
32use crate::error::Error;
33
34const U32_BYTES: usize = core::mem::size_of::<u32>();
36const U64_BYTES: usize = core::mem::size_of::<u64>();
38const F32_BYTES: usize = core::mem::size_of::<f32>();
40
41const HEADER: usize = U32_BYTES + U32_BYTES;
43
44#[derive(Clone, Copy, Debug, PartialEq, Eq)]
46#[cfg_attr(feature = "serde", derive(serde::Serialize))]
47pub struct JournalEntry<'a> {
48 pub op: u8,
50 pub payload: &'a [u8],
52}
53
54#[derive(Clone, Debug, PartialEq, Eq)]
56pub struct JournalScan<'a> {
57 pub entries: Vec<JournalEntry<'a>>,
59 pub truncated_tail: bool,
62}
63
64fn body_checksum(body: &[u8]) -> u32 {
67 xxh3_64(body) as u32
68}
69
70pub fn encode_entry(out: &mut Vec<u8>, op: u8, payload: &[u8]) {
73 let len = 1 + payload.len();
74 let len32 = u32::try_from(len).expect("journal payload fits u32 by construction");
75 out.reserve(HEADER + len);
76 out.extend_from_slice(&len32.to_le_bytes());
77 let check_pos = out.len();
80 out.extend_from_slice(&[0u8; U32_BYTES]);
81 out.push(op);
82 out.extend_from_slice(payload);
83 let check = body_checksum(&out[check_pos + U32_BYTES..]);
84 out[check_pos..check_pos + U32_BYTES].copy_from_slice(&check.to_le_bytes());
85}
86
87#[derive(Clone, Debug, PartialEq)]
94#[cfg_attr(feature = "serde", derive(serde::Serialize))]
95pub enum Op<'a> {
96 Remember {
98 now: u64,
100 valid_from: u64,
103 entity: Option<&'a str>,
105 text: &'a str,
107 tags: Vec<&'a str>,
109 links: Vec<(&'a str, &'a str)>,
111 vector: Vec<f32>,
115 metadata: Vec<(&'a str, &'a str)>,
119 revises: crate::id::FactId,
121 assigned: crate::id::FactId,
124 },
125 Forget {
127 now: u64,
129 fact: crate::id::FactId,
131 },
132 Link {
134 now: u64,
136 src: &'a str,
138 rel: &'a str,
140 dst: &'a str,
142 provenance: crate::id::FactId,
144 },
145 Maintain {
147 now: u64,
149 },
150}
151
152fn put_str(out: &mut Vec<u8>, s: &str) {
154 out.extend_from_slice(&(s.len() as u32).to_le_bytes());
155 out.extend_from_slice(s.as_bytes());
156}
157
158fn take_str<'a>(bytes: &'a [u8], at: &mut usize) -> Result<&'a str, Error> {
160 let len = take_u32(bytes, at)? as usize;
161 let end = at
162 .checked_add(len)
163 .filter(|&e| e <= bytes.len())
164 .ok_or(Error::Corrupt("journal string overruns its record"))?;
165 let s = core::str::from_utf8(&bytes[*at..end])
166 .map_err(|_| Error::Corrupt("journal string is not UTF-8"))?;
167 *at = end;
168 Ok(s)
169}
170
171fn take_u32(bytes: &[u8], at: &mut usize) -> Result<u32, Error> {
173 let end = *at + U32_BYTES;
174 if end > bytes.len() {
175 return Err(Error::Corrupt("journal record truncated inside a field"));
176 }
177 let v = u32::from_le_bytes(bytes[*at..end].try_into().unwrap());
178 *at = end;
179 Ok(v)
180}
181
182fn take_u64(bytes: &[u8], at: &mut usize) -> Result<u64, Error> {
184 let end = *at + U64_BYTES;
185 if end > bytes.len() {
186 return Err(Error::Corrupt("journal record truncated inside a field"));
187 }
188 let v = u64::from_le_bytes(bytes[*at..end].try_into().unwrap());
189 *at = end;
190 Ok(v)
191}
192
193fn take_vec_f32(bytes: &[u8], at: &mut usize) -> Result<Vec<f32>, Error> {
195 let count = take_u32(bytes, at)? as usize;
196 let end = *at as u64 + count as u64 * F32_BYTES as u64;
202 if end > bytes.len() as u64 {
203 return Err(Error::Corrupt("journal vector overruns its record"));
204 }
205 let end = end as usize;
206 let mut v = Vec::with_capacity(count);
207 let mut p = *at;
208 while p < end {
209 v.push(f32::from_le_bytes(
210 bytes[p..p + F32_BYTES].try_into().unwrap(),
211 ));
212 p += F32_BYTES;
213 }
214 *at = end;
215 Ok(v)
216}
217
218impl<'a> Op<'a> {
219 pub fn encode(&self, out: &mut Vec<u8>) {
222 let mut payload = Vec::new();
223 let op = match self {
224 Op::Remember {
225 now,
226 valid_from,
227 entity,
228 text,
229 tags,
230 links,
231 vector,
232 metadata,
233 revises,
234 assigned,
235 } => {
236 payload.extend_from_slice(&now.to_le_bytes());
237 payload.extend_from_slice(&valid_from.to_le_bytes());
238 payload.extend_from_slice(&revises.0.to_le_bytes());
239 payload.extend_from_slice(&assigned.0.to_le_bytes());
240 match entity {
241 Some(name) => {
242 payload.push(1);
243 put_str(&mut payload, name);
244 }
245 None => payload.push(0),
246 }
247 put_str(&mut payload, text);
248 payload.push(tags.len() as u8);
249 for tag in tags {
250 put_str(&mut payload, tag);
251 }
252 payload.push(links.len() as u8);
253 for (rel, dst) in links {
254 put_str(&mut payload, rel);
255 put_str(&mut payload, dst);
256 }
257 payload.extend_from_slice(&(vector.len() as u32).to_le_bytes());
258 for &x in vector {
259 payload.extend_from_slice(&x.to_le_bytes());
260 }
261 payload.extend_from_slice(&(metadata.len() as u32).to_le_bytes());
262 for (k, v) in metadata {
263 put_str(&mut payload, k);
264 put_str(&mut payload, v);
265 }
266 if revises.is_none() { 1 } else { 2 }
267 }
268 Op::Forget { now, fact } => {
269 payload.extend_from_slice(&now.to_le_bytes());
270 payload.extend_from_slice(&fact.0.to_le_bytes());
271 3
272 }
273 Op::Link {
274 now,
275 src,
276 rel,
277 dst,
278 provenance,
279 } => {
280 payload.extend_from_slice(&now.to_le_bytes());
281 payload.extend_from_slice(&provenance.0.to_le_bytes());
282 put_str(&mut payload, src);
283 put_str(&mut payload, rel);
284 put_str(&mut payload, dst);
285 4
286 }
287 Op::Maintain { now } => {
288 payload.extend_from_slice(&now.to_le_bytes());
289 5
290 }
291 };
292 encode_entry(out, op, &payload);
293 }
294
295 pub fn decode(op: u8, payload: &'a [u8]) -> Result<Op<'a>, Error> {
300 use crate::id::FactId;
301 let at = &mut 0usize;
302 let decoded = match op {
303 1 | 2 => {
304 let now = take_u64(payload, at)?;
305 let valid_from = take_u64(payload, at)?;
306 let revises = FactId(take_u32(payload, at)?);
307 let assigned = FactId(take_u32(payload, at)?);
308 if (op == 2) == revises.is_none() {
309 return Err(Error::Corrupt("journal revises field disagrees with op"));
310 }
311 let entity = match payload.get(*at) {
312 Some(0) => {
313 *at += 1;
314 None
315 }
316 Some(1) => {
317 *at += 1;
318 Some(take_str(payload, at)?)
319 }
320 _ => return Err(Error::Corrupt("journal entity flag is invalid")),
321 };
322 let text = take_str(payload, at)?;
323 let tag_cnt = *payload
324 .get(*at)
325 .ok_or(Error::Corrupt("journal record truncated inside a field"))?;
326 *at += 1;
327 let mut tags = Vec::with_capacity(tag_cnt as usize);
328 for _ in 0..tag_cnt {
329 tags.push(take_str(payload, at)?);
330 }
331 let link_cnt = *payload
332 .get(*at)
333 .ok_or(Error::Corrupt("journal record truncated inside a field"))?;
334 *at += 1;
335 let mut links = Vec::with_capacity(link_cnt as usize);
336 for _ in 0..link_cnt {
337 let rel = take_str(payload, at)?;
338 let dst = take_str(payload, at)?;
339 links.push((rel, dst));
340 }
341 let vector = take_vec_f32(payload, at)?;
342 let meta_cnt = take_u32(payload, at)?;
343 let mut metadata = Vec::new();
344 for _ in 0..meta_cnt {
345 let k = take_str(payload, at)?;
346 let v = take_str(payload, at)?;
347 metadata.push((k, v));
348 }
349 Op::Remember {
350 now,
351 valid_from,
352 entity,
353 text,
354 tags,
355 links,
356 vector,
357 metadata,
358 revises,
359 assigned,
360 }
361 }
362 3 => Op::Forget {
363 now: take_u64(payload, at)?,
364 fact: FactId(take_u32(payload, at)?),
365 },
366 4 => {
367 let now = take_u64(payload, at)?;
368 let provenance = FactId(take_u32(payload, at)?);
369 let src = take_str(payload, at)?;
370 let rel = take_str(payload, at)?;
371 let dst = take_str(payload, at)?;
372 Op::Link {
373 now,
374 src,
375 rel,
376 dst,
377 provenance,
378 }
379 }
380 5 => Op::Maintain {
381 now: take_u64(payload, at)?,
382 },
383 _ => return Err(Error::Corrupt("unknown journal op")),
384 };
385 if *at != payload.len() {
386 return Err(Error::Corrupt("journal record has trailing bytes"));
387 }
388 Ok(decoded)
389 }
390}
391
392pub fn scan(journal: &[u8]) -> Result<JournalScan<'_>, Error> {
395 let mut entries = Vec::new();
396 let mut pos = 0usize;
397 while pos < journal.len() {
398 let rest = &journal[pos..];
399 if rest.len() < HEADER {
400 return Ok(JournalScan {
401 entries,
402 truncated_tail: true,
403 });
404 }
405 let len = u32::from_le_bytes(rest[..U32_BYTES].try_into().unwrap()) as usize;
406 if len == 0 {
407 return Err(Error::Corrupt("journal record with zero length"));
408 }
409 let Some(body) = HEADER
415 .checked_add(len)
416 .and_then(|end| rest.get(HEADER..end))
417 else {
418 return Ok(JournalScan {
419 entries,
420 truncated_tail: true,
421 });
422 };
423 let want = u32::from_le_bytes(rest[U32_BYTES..HEADER].try_into().unwrap());
424 if body_checksum(body) != want {
425 if pos + HEADER + len == journal.len() {
426 return Ok(JournalScan {
427 entries,
428 truncated_tail: true,
429 });
430 }
431 return Err(Error::Corrupt("journal checksum mismatch mid-stream"));
432 }
433 entries.push(JournalEntry {
434 op: body[0],
435 payload: &body[1..],
436 });
437 pos += HEADER + len;
438 }
439 Ok(JournalScan {
440 entries,
441 truncated_tail: false,
442 })
443}