Skip to main content

kcode_k1_chat_persistence_store/
lib.rs

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}