1use kcode_k1_chat_chatend::{BoxContent, ChatBox, Chatend};
2pub use kcode_k1_chat_chatend::{BoxId, ToolCallId};
3pub use kcode_k1_transaction_id::TxId;
4use serde::{Deserialize, Serialize};
5use serde_json::Value;
6use std::fs::{self, File, OpenOptions};
7use std::io::{self, Seek, SeekFrom, Write};
8use std::path::{Path, PathBuf};
9
10pub type SessionId = [u8; 12];
11pub const CODEC_VERSION: u32 = 1;
12
13#[derive(Clone, Debug, PartialEq, Eq)]
14pub struct EventRecord {
15 pub after_box_id: u64,
16 pub event_index: u64,
17 pub connected_box_id: u64,
18 pub handler: String,
19 pub data: Value,
20}
21
22impl EventRecord {
23 pub fn new(
24 after_box_id: u64,
25 event_index: u64,
26 connected_box_id: u64,
27 handler: String,
28 data: Value,
29 ) -> Result<Self, Error> {
30 let value = Self {
31 after_box_id,
32 event_index,
33 connected_box_id,
34 handler,
35 data,
36 };
37 validate_event(&value)?;
38 Ok(value)
39 }
40}
41
42#[derive(Clone, Debug, PartialEq, Eq)]
43pub enum Record {
44 Box(ChatBox),
45 Event(EventRecord),
46}
47
48impl Record {
49 pub fn chat_box(value: ChatBox) -> Self {
50 Self::Box(value)
51 }
52
53 pub fn event(value: EventRecord) -> Self {
54 Self::Event(value)
55 }
56}
57
58#[derive(Clone, Debug, PartialEq, Eq)]
59pub struct Batch {
60 pub version: u32,
61 pub session_id: SessionId,
62 pub predecessor: Option<Record>,
63 pub records: Vec<Record>,
64}
65
66impl Batch {
67 pub fn new(
68 session_id: SessionId,
69 predecessor: Option<Record>,
70 records: Vec<Record>,
71 ) -> Result<Self, Error> {
72 let value = Self {
73 version: CODEC_VERSION,
74 session_id,
75 predecessor,
76 records,
77 };
78 validate_batch(&value)?;
79 Ok(value)
80 }
81
82 pub fn encode(&self) -> Result<Vec<u8>, Error> {
83 validate_batch(self)?;
84 Ok(serde_json::to_vec(&WireBatch::from(self))?)
85 }
86
87 pub fn decode(bytes: &[u8]) -> Result<Self, Error> {
88 let value = Self::try_from(serde_json::from_slice::<WireBatch>(bytes)?)?;
89 validate_batch(&value)?;
90 Ok(value)
91 }
92}
93
94#[derive(Clone, Debug, Default, PartialEq, Eq)]
95pub struct SessionLog {
96 pub boxes: Vec<ChatBox>,
97 pub events: Vec<EventRecord>,
98 pub records: Vec<Record>,
99}
100
101#[derive(Debug)]
102pub enum Error {
103 Io(io::Error),
104 Json(serde_json::Error),
105 Invalid(&'static str),
106}
107
108impl std::fmt::Display for Error {
109 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
110 match self {
111 Self::Io(e) => write!(f, "I/O error: {e}"),
112 Self::Json(e) => write!(f, "JSON error: {e}"),
113 Self::Invalid(e) => f.write_str(e),
114 }
115 }
116}
117
118impl std::error::Error for Error {}
119
120impl From<io::Error> for Error {
121 fn from(value: io::Error) -> Self {
122 Self::Io(value)
123 }
124}
125
126impl From<serde_json::Error> for Error {
127 fn from(value: serde_json::Error) -> Self {
128 Self::Json(value)
129 }
130}
131
132type Result<T, E = Error> = std::result::Result<T, E>;
133
134pub struct Projection {
135 root: PathBuf,
136 sessions: PathBuf,
137}
138
139impl Projection {
140 pub fn new(root: impl AsRef<Path>) -> Result<Self> {
141 let root = root.as_ref().to_path_buf();
142 fs::create_dir_all(&root)?;
143 let sessions = root.join("sessions");
144 fs::create_dir_all(&sessions)?;
145 Ok(Self { root, sessions })
146 }
147
148 pub fn load(&self, session: SessionId) -> Result<SessionLog> {
149 let records = self
150 .read_strict(session)?
151 .into_iter()
152 .map(|value| value.1)
153 .collect::<Vec<_>>();
154 validate_session_records(session, &records)?;
155 Ok(to_log(records))
156 }
157
158 pub fn apply(&mut self, txid: TxId, batch: &Batch, reconcile_first: bool) -> Result<()> {
159 validate_batch(batch)?;
160 fs::create_dir_all(&self.sessions)?;
161 let path = self.path(batch.session_id);
162
163 if reconcile_first {
164 let bytes = match fs::read(&path) {
165 Ok(value) => value,
166 Err(error) if error.kind() == io::ErrorKind::NotFound => Vec::new(),
167 Err(error) => return Err(error.into()),
168 };
169 let end = locate_predecessor(&bytes, batch.predecessor.as_ref())?;
170 let mut combined = strict_prefix(&bytes[..end])?
171 .into_iter()
172 .map(|value| value.1)
173 .collect::<Vec<_>>();
174 combined.extend(batch.records.clone());
175 validate_session_records(batch.session_id, &combined)?;
176
177 let mut file = OpenOptions::new()
178 .create(true)
179 .read(true)
180 .write(true)
181 .truncate(false)
182 .open(&path)?;
183 file.set_len(end as u64)?;
184 file.seek(SeekFrom::Start(end as u64))?;
185 append_lines(&mut file, txid, &batch.records)?;
186 file.sync_all()?;
187 } else {
188 let current = self.read_strict(batch.session_id)?;
189 let mut combined = current
190 .iter()
191 .map(|value| value.1.clone())
192 .collect::<Vec<_>>();
193 if combined.last() != batch.predecessor.as_ref() {
194 return Err(Error::Invalid("predecessor mismatch"));
195 }
196 combined.extend(batch.records.clone());
197 validate_session_records(batch.session_id, &combined)?;
198
199 let mut file = OpenOptions::new().create(true).append(true).open(&path)?;
200 append_lines(&mut file, txid, &batch.records)?;
201 file.sync_all()?;
202 }
203
204 Ok(())
205 }
206
207 pub fn discard_all(&mut self) -> Result<()> {
208 match fs::remove_dir_all(&self.sessions) {
209 Ok(()) => {}
210 Err(error) if error.kind() == io::ErrorKind::NotFound => {}
211 Err(error) => return Err(error.into()),
212 }
213 File::open(&self.root)?.sync_all()?;
214 Ok(())
215 }
216
217 fn path(&self, session: SessionId) -> PathBuf {
218 self.sessions.join(format!("{}.jsonl", hex(session)))
219 }
220
221 fn read_strict(&self, session: SessionId) -> Result<Vec<(TxId, Record, usize)>> {
222 let bytes = match fs::read(self.path(session)) {
223 Ok(value) => value,
224 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()),
225 Err(error) => return Err(error.into()),
226 };
227 strict_prefix(&bytes)
228 }
229}
230
231fn append_lines(file: &mut File, txid: TxId, records: &[Record]) -> Result<()> {
232 for record in records {
233 let mut bytes = serde_json::to_vec(&WireLine {
234 txid: hex(*txid.as_bytes()),
235 record: WireRecord::from(record),
236 })?;
237 bytes.push(b'\n');
238 file.write_all(&bytes)?;
239 }
240 Ok(())
241}
242
243fn strict_prefix(bytes: &[u8]) -> Result<Vec<(TxId, Record, usize)>> {
244 if !bytes.is_empty() && bytes.last() != Some(&b'\n') {
245 return Err(Error::Invalid("line is not newline terminated"));
246 }
247
248 let mut records = Vec::new();
249 let mut start = 0;
250 while start < bytes.len() {
251 let relative = bytes[start..]
252 .iter()
253 .position(|byte| *byte == b'\n')
254 .ok_or(Error::Invalid("line is not newline terminated"))?;
255 let end = start + relative + 1;
256 let wire: WireLine = serde_json::from_slice(&bytes[start..end - 1])?;
257 records.push((
258 TxId::from_bytes(parse_hex(&wire.txid)?),
259 Record::try_from(wire.record)?,
260 end,
261 ));
262 start = end;
263 }
264
265 validate_records(
266 &records
267 .iter()
268 .map(|value| value.1.clone())
269 .collect::<Vec<_>>(),
270 )?;
271 Ok(records)
272}
273
274fn locate_predecessor(bytes: &[u8], wanted: Option<&Record>) -> Result<usize> {
275 let Some(wanted) = wanted else {
276 return Ok(0);
277 };
278
279 let mut start = 0;
280 let mut found = None;
281 let mut prefix = Vec::new();
282 while start < bytes.len() {
283 let Some(relative) = bytes[start..].iter().position(|byte| *byte == b'\n') else {
284 break;
285 };
286 let end = start + relative + 1;
287 let parsed = serde_json::from_slice::<WireLine>(&bytes[start..end - 1])
288 .map_err(Error::from)
289 .and_then(|wire| Record::try_from(wire.record));
290 let record = match parsed {
291 Ok(value) => value,
292 Err(_) if found.is_some() => break,
293 Err(error) => return Err(error),
294 };
295
296 if same_identity(&record, wanted) {
297 if &record != wanted {
298 return Err(Error::Invalid("predecessor identity mismatch"));
299 }
300 if found.is_some() {
301 return Err(Error::Invalid("duplicate predecessor identity"));
302 }
303 found = Some(end);
304 }
305 if found.is_none() {
306 prefix.push(record);
307 validate_records(&prefix)?;
308 }
309 start = end;
310 }
311
312 found.ok_or(Error::Invalid("predecessor absent"))
313}
314
315fn same_identity(left: &Record, right: &Record) -> bool {
316 match (left, right) {
317 (Record::Box(left), Record::Box(right)) => left.id() == right.id(),
318 (Record::Event(left), Record::Event(right)) => {
319 (left.after_box_id, left.event_index, left.connected_box_id)
320 == (
321 right.after_box_id,
322 right.event_index,
323 right.connected_box_id,
324 )
325 }
326 _ => false,
327 }
328}
329
330fn validate_batch(batch: &Batch) -> Result<()> {
331 if batch.version != CODEC_VERSION {
332 return Err(Error::Invalid("unsupported version"));
333 }
334 if batch.records.is_empty() {
335 return Err(Error::Invalid("empty batch"));
336 }
337 if let Some(record) = &batch.predecessor {
338 validate_record_session(batch.session_id, record)?;
339 }
340 for record in &batch.records {
341 validate_record_session(batch.session_id, record)?;
342 }
343 validate_suffix(batch.predecessor.as_ref(), &batch.records)?;
344 if batch.predecessor.is_none() {
345 validate_session_records(batch.session_id, &batch.records)?;
346 }
347 Ok(())
348}
349
350fn validate_suffix(predecessor: Option<&Record>, suffix: &[Record]) -> Result<()> {
351 let (mut latest, mut event_index) = match predecessor {
352 None => (0, 0),
353 Some(Record::Box(value)) => (value.id().get(), 0),
354 Some(Record::Event(value)) => {
355 validate_event(value)?;
356 (value.after_box_id, value.event_index)
357 }
358 };
359
360 for record in suffix {
361 match record {
362 Record::Box(value) => {
363 if value.id().get()
364 != latest
365 .checked_add(1)
366 .ok_or(Error::Invalid("box ID overflow"))?
367 {
368 return Err(Error::Invalid("noncontiguous box ID"));
369 }
370 latest = value.id().get();
371 event_index = 0;
372 }
373 Record::Event(value) => {
374 validate_event(value)?;
375 if value.after_box_id != latest
376 || value.event_index
377 != event_index
378 .checked_add(1)
379 .ok_or(Error::Invalid("event index overflow"))?
380 || value.connected_box_id > latest
381 {
382 return Err(Error::Invalid("invalid event order or association"));
383 }
384 event_index = value.event_index;
385 }
386 }
387 }
388 Ok(())
389}
390
391fn validate_event(value: &EventRecord) -> Result<()> {
392 if value.handler.is_empty() {
393 Err(Error::Invalid("empty event handler"))
394 } else {
395 Ok(())
396 }
397}
398
399fn validate_record_session(session: SessionId, record: &Record) -> Result<()> {
400 let id = match record {
401 Record::Box(value) => match value.content() {
402 BoxContent::KtoolCall { tool_call_id, .. }
403 | BoxContent::KtoolReturn { tool_call_id, .. } => Some(tool_call_id),
404 _ => None,
405 },
406 Record::Event(_) => None,
407 };
408
409 if id.is_some_and(|value| value.session() != session) {
410 Err(Error::Invalid("ToolCallId belongs to another session"))
411 } else {
412 Ok(())
413 }
414}
415
416fn validate_session_records(session: SessionId, records: &[Record]) -> Result<()> {
417 validate_records(records)?;
418 for record in records {
419 validate_record_session(session, record)?;
420 }
421 Ok(())
422}
423
424fn validate_records(records: &[Record]) -> Result<()> {
425 validate_suffix(None, records)?;
426 let boxes = records
427 .iter()
428 .filter_map(|record| match record {
429 Record::Box(value) => Some(value.clone()),
430 Record::Event(_) => None,
431 })
432 .collect();
433 Chatend::recover(boxes).map_err(|_| Error::Invalid("invalid Chatend recovery"))?;
434 Ok(())
435}
436
437fn to_log(records: Vec<Record>) -> SessionLog {
438 let boxes = records
439 .iter()
440 .filter_map(|record| match record {
441 Record::Box(value) => Some(value.clone()),
442 Record::Event(_) => None,
443 })
444 .collect();
445 let events = records
446 .iter()
447 .filter_map(|record| match record {
448 Record::Event(value) => Some(value.clone()),
449 Record::Box(_) => None,
450 })
451 .collect();
452 SessionLog {
453 boxes,
454 events,
455 records,
456 }
457}
458
459fn hex(bytes: [u8; 12]) -> String {
460 bytes.iter().map(|value| format!("{value:02x}")).collect()
461}
462
463fn parse_hex(value: &str) -> Result<[u8; 12]> {
464 if value.len() != 24
465 || !value
466 .bytes()
467 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
468 {
469 return Err(Error::Invalid("invalid lowercase 24-hex value"));
470 }
471
472 let mut output = [0; 12];
473 for (index, slot) in output.iter_mut().enumerate() {
474 *slot = u8::from_str_radix(&value[index * 2..index * 2 + 2], 16)
475 .map_err(|_| Error::Invalid("invalid hex"))?;
476 }
477 Ok(output)
478}
479
480#[derive(Serialize, Deserialize)]
481#[serde(deny_unknown_fields)]
482struct WireBatch {
483 version: u32,
484 session_id: String,
485 predecessor: Option<WireRecord>,
486 records: Vec<WireRecord>,
487}
488
489impl From<&Batch> for WireBatch {
490 fn from(value: &Batch) -> Self {
491 Self {
492 version: value.version,
493 session_id: hex(value.session_id),
494 predecessor: value.predecessor.as_ref().map(WireRecord::from),
495 records: value.records.iter().map(WireRecord::from).collect(),
496 }
497 }
498}
499
500impl TryFrom<WireBatch> for Batch {
501 type Error = Error;
502
503 fn try_from(value: WireBatch) -> Result<Self> {
504 Ok(Self {
505 version: value.version,
506 session_id: parse_hex(&value.session_id)?,
507 predecessor: value.predecessor.map(Record::try_from).transpose()?,
508 records: value
509 .records
510 .into_iter()
511 .map(Record::try_from)
512 .collect::<Result<_>>()?,
513 })
514 }
515}
516
517#[derive(Serialize, Deserialize)]
518#[serde(deny_unknown_fields)]
519struct WireLine {
520 txid: String,
521 record: WireRecord,
522}
523
524#[derive(Serialize, Deserialize)]
525#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
526enum WireRecord {
527 Box {
528 id: u64,
529 content: WireContent,
530 },
531 Event {
532 after_box_id: u64,
533 event_index: u64,
534 connected_box_id: u64,
535 handler: String,
536 data: Value,
537 },
538}
539
540impl From<&Record> for WireRecord {
541 fn from(value: &Record) -> Self {
542 match value {
543 Record::Box(value) => Self::Box {
544 id: value.id().get(),
545 content: WireContent::from(value.content()),
546 },
547 Record::Event(value) => Self::Event {
548 after_box_id: value.after_box_id,
549 event_index: value.event_index,
550 connected_box_id: value.connected_box_id,
551 handler: value.handler.clone(),
552 data: value.data.clone(),
553 },
554 }
555 }
556}
557
558impl TryFrom<WireRecord> for Record {
559 type Error = Error;
560
561 fn try_from(value: WireRecord) -> Result<Self> {
562 Ok(match value {
563 WireRecord::Box { id, content } => {
564 Self::Box(ChatBox::new(BoxId::new(id), BoxContent::try_from(content)?))
565 }
566 WireRecord::Event {
567 after_box_id,
568 event_index,
569 connected_box_id,
570 handler,
571 data,
572 } => Self::Event(EventRecord::new(
573 after_box_id,
574 event_index,
575 connected_box_id,
576 handler,
577 data,
578 )?),
579 })
580 }
581}
582
583#[derive(Serialize, Deserialize)]
584#[serde(tag = "type", rename_all = "snake_case", deny_unknown_fields)]
585enum WireContent {
586 System {
587 text: String,
588 },
589 User {
590 text: String,
591 },
592 Kennedy {
593 text: String,
594 },
595 Attachment,
596 KtoolCall {
597 session: String,
598 sequence: u64,
599 name: String,
600 arguments: String,
601 },
602 KtoolReturn {
603 session: String,
604 sequence: u64,
605 originating_call: u64,
606 result: WireResult,
607 },
608}
609
610impl From<&BoxContent> for WireContent {
611 fn from(value: &BoxContent) -> Self {
612 match value {
613 BoxContent::System(text) => Self::System { text: text.clone() },
614 BoxContent::User(text) => Self::User { text: text.clone() },
615 BoxContent::Kennedy { text } => Self::Kennedy { text: text.clone() },
616 BoxContent::Attachment => Self::Attachment,
617 BoxContent::KtoolCall {
618 tool_call_id,
619 name,
620 arguments,
621 } => Self::KtoolCall {
622 session: hex(tool_call_id.session()),
623 sequence: tool_call_id.sequence(),
624 name: name.clone(),
625 arguments: arguments.clone(),
626 },
627 BoxContent::KtoolReturn {
628 tool_call_id,
629 originating_call,
630 result,
631 } => Self::KtoolReturn {
632 session: hex(tool_call_id.session()),
633 sequence: tool_call_id.sequence(),
634 originating_call: originating_call.get(),
635 result: WireResult::from(result),
636 },
637 }
638 }
639}
640
641impl TryFrom<WireContent> for BoxContent {
642 type Error = Error;
643
644 fn try_from(value: WireContent) -> Result<Self> {
645 Ok(match value {
646 WireContent::System { text } => Self::System(text),
647 WireContent::User { text } => Self::User(text),
648 WireContent::Kennedy { text } => Self::Kennedy { text },
649 WireContent::Attachment => Self::Attachment,
650 WireContent::KtoolCall {
651 session,
652 sequence,
653 name,
654 arguments,
655 } => Self::KtoolCall {
656 tool_call_id: ToolCallId::new(parse_hex(&session)?, sequence),
657 name,
658 arguments,
659 },
660 WireContent::KtoolReturn {
661 session,
662 sequence,
663 originating_call,
664 result,
665 } => Self::KtoolReturn {
666 tool_call_id: ToolCallId::new(parse_hex(&session)?, sequence),
667 originating_call: BoxId::new(originating_call),
668 result: result.into(),
669 },
670 })
671 }
672}
673
674#[derive(Serialize, Deserialize)]
675#[serde(tag = "status", rename_all = "snake_case", deny_unknown_fields)]
676enum WireResult {
677 Ok { value: String },
678 Err { value: String },
679}
680
681impl From<&std::result::Result<String, String>> for WireResult {
682 fn from(value: &std::result::Result<String, String>) -> Self {
683 match value {
684 Ok(value) => Self::Ok {
685 value: value.clone(),
686 },
687 Err(value) => Self::Err {
688 value: value.clone(),
689 },
690 }
691 }
692}
693
694impl From<WireResult> for std::result::Result<String, String> {
695 fn from(value: WireResult) -> Self {
696 match value {
697 WireResult::Ok { value } => Ok(value),
698 WireResult::Err { value } => Err(value),
699 }
700 }
701}
702
703#[cfg(test)]
704mod tests {
705 use super::*;
706 use serde_json::json;
707 use std::sync::atomic::{AtomicU64, Ordering};
708
709 fn boxed(id: u64, content: BoxContent) -> Record {
710 Record::Box(ChatBox::new(BoxId::new(id), content))
711 }
712
713 fn event(after: u64, index: u64) -> Record {
714 Record::Event(
715 EventRecord::new(
716 after,
717 index,
718 after,
719 "handler".into(),
720 json!({"z":[true,null,{"a":1}]}),
721 )
722 .unwrap(),
723 )
724 }
725
726 fn directory() -> PathBuf {
727 static NEXT: AtomicU64 = AtomicU64::new(0);
728 let path = std::env::temp_dir().join(format!(
729 "k1-persistence-{}-{}",
730 std::process::id(),
731 NEXT.fetch_add(1, Ordering::Relaxed)
732 ));
733 let _ = fs::remove_dir_all(&path);
734 path
735 }
736
737 fn tx(value: u8) -> TxId {
738 TxId::from_bytes([value; 12])
739 }
740
741 #[test]
742 fn codec_order_and_session_validation() {
743 let tool = ToolCallId::new([1; 12], 7);
744 let records = vec![
745 event(0, 1),
746 boxed(1, BoxContent::System("s".into())),
747 boxed(2, BoxContent::User("u".into())),
748 boxed(3, BoxContent::Kennedy { text: "k".into() }),
749 boxed(4, BoxContent::Attachment),
750 boxed(
751 5,
752 BoxContent::KtoolCall {
753 tool_call_id: tool,
754 name: "n".into(),
755 arguments: "{}".into(),
756 },
757 ),
758 boxed(
759 6,
760 BoxContent::KtoolReturn {
761 tool_call_id: tool,
762 originating_call: BoxId::new(5),
763 result: Err("e".into()),
764 },
765 ),
766 event(6, 1),
767 ];
768 let batch = Batch::new([1; 12], None, records).unwrap();
769 assert_eq!(Batch::decode(&batch.encode().unwrap()).unwrap(), batch);
770
771 let wrong = boxed(
772 1,
773 BoxContent::KtoolCall {
774 tool_call_id: ToolCallId::new([2; 12], 1),
775 name: "n".into(),
776 arguments: "{}".into(),
777 },
778 );
779 assert!(Batch::new([1; 12], None, vec![wrong]).is_err());
780 assert!(
781 Batch::new(
782 [2; 12],
783 None,
784 vec![boxed(
785 1,
786 BoxContent::KtoolCall {
787 tool_call_id: ToolCallId::new([2; 12], 1),
788 name: "n".into(),
789 arguments: "{}".into(),
790 },
791 )],
792 )
793 .is_ok()
794 );
795 assert!(Batch::new([0; 12], None, vec![boxed(2, BoxContent::Attachment)]).is_err());
796 assert!(EventRecord::new(0, 1, 0, String::new(), json!(null)).is_err());
797 }
798
799 #[test]
800 fn append_load_reconcile_and_discard() {
801 let root = directory();
802 let mut projection = Projection::new(&root).unwrap();
803 let first = Batch::new(
804 [4; 12],
805 None,
806 vec![boxed(1, BoxContent::System("a".into()))],
807 )
808 .unwrap();
809 projection.apply(tx(1), &first, false).unwrap();
810
811 let next = Batch::new(
812 [4; 12],
813 Some(first.records[0].clone()),
814 vec![event(1, 1), boxed(2, BoxContent::User("b".into()))],
815 )
816 .unwrap();
817 projection.apply(tx(2), &next, false).unwrap();
818 projection.apply(tx(2), &next, true).unwrap();
819 assert_eq!(projection.load([4; 12]).unwrap().records.len(), 3);
820
821 let mut file = OpenOptions::new()
822 .append(true)
823 .open(projection.path([4; 12]))
824 .unwrap();
825 file.write_all(b"{partial").unwrap();
826 projection.apply(tx(2), &next, true).unwrap();
827 assert_eq!(projection.load([4; 12]).unwrap().records.len(), 3);
828 assert!(projection.load([5; 12]).unwrap().records.is_empty());
829
830 projection.discard_all().unwrap();
831 assert!(projection.load([4; 12]).unwrap().records.is_empty());
832 fs::remove_dir_all(root).unwrap();
833 }
834
835 #[test]
836 fn malformed_lines_and_predecessors_fail_closed() {
837 let root = directory();
838 let mut projection = Projection::new(&root).unwrap();
839 fs::write(projection.path([1; 12]), b"{}\n").unwrap();
840 assert!(projection.load([1; 12]).is_err());
841 fs::write(projection.path([2; 12]), b"{}").unwrap();
842 assert!(projection.load([2; 12]).is_err());
843
844 let missing = Batch::new(
845 [3; 12],
846 Some(boxed(9, BoxContent::Attachment)),
847 vec![boxed(10, BoxContent::Attachment)],
848 )
849 .unwrap();
850 assert!(projection.apply(tx(3), &missing, true).is_err());
851 fs::remove_dir_all(root).unwrap();
852 }
853}