Skip to main content

traverse_runtime/
data_store.rs

1//! Portable state access for Traverse capabilities.
2//!
3//! Governed by spec `032-universal-data-access`.
4
5use serde::{Deserialize, Serialize};
6use serde_json::{Value, json};
7use std::collections::{BTreeMap, BTreeSet};
8use std::fs;
9use std::path::PathBuf;
10use traverse_contracts::CapabilityContract;
11
12const DATA_STORE_SPEC: &str = "032-universal-data-access";
13
14#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
15pub struct StateRecord {
16    pub key: String,
17    pub value: Value,
18    pub lamport_clock: u64,
19    pub writer_id: String,
20}
21
22#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
23pub struct MergeDecision {
24    pub key: String,
25    pub winning_writer_id: String,
26    pub winning_lamport_clock: u64,
27    pub resolution_rule: ConflictResolutionRule,
28}
29
30#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
31#[serde(rename_all = "snake_case")]
32pub enum ConflictResolutionRule {
33    OnlyLocal,
34    OnlyRemote,
35    HigherLamportClock,
36    WriterIdentityTieBreak,
37}
38
39#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
40pub struct SyncReport {
41    pub governing_spec: String,
42    pub decisions: Vec<MergeDecision>,
43}
44
45#[derive(Debug, Clone, PartialEq, Eq)]
46pub struct DataStoreError {
47    pub code: DataStoreErrorCode,
48    pub message: String,
49    pub details: Value,
50}
51
52#[derive(Debug, Clone, Copy, PartialEq, Eq)]
53pub enum DataStoreErrorCode {
54    SchemaValidationError,
55    NoStateSchemaDeclared,
56    LamportClockOverflow,
57    InvalidKey,
58    IoFailure,
59    SerializationFailure,
60    SyncFailure,
61}
62
63pub trait DataStore {
64    /// Reads a stored state record.
65    ///
66    /// # Errors
67    ///
68    /// Returns [`DataStoreError`] when the adapter cannot read the key.
69    fn read(&self, key: &str) -> Result<Option<StateRecord>, DataStoreError>;
70
71    /// Writes a stamped state record.
72    ///
73    /// # Errors
74    ///
75    /// Returns [`DataStoreError`] when the adapter cannot persist the record.
76    fn write(&mut self, record: StateRecord) -> Result<(), DataStoreError>;
77
78    /// Deletes a stored state record.
79    ///
80    /// # Errors
81    ///
82    /// Returns [`DataStoreError`] when the adapter cannot delete the key.
83    fn delete(&mut self, key: &str) -> Result<(), DataStoreError>;
84
85    /// Lists stored state keys.
86    ///
87    /// # Errors
88    ///
89    /// Returns [`DataStoreError`] when the adapter cannot enumerate keys.
90    fn list_keys(&self) -> Result<Vec<String>, DataStoreError>;
91}
92
93#[derive(Debug, Clone, PartialEq, Eq)]
94pub struct LamportClock {
95    writer_id: String,
96    value: u64,
97}
98
99impl LamportClock {
100    #[must_use]
101    pub fn new(writer_id: impl Into<String>) -> Self {
102        Self {
103            writer_id: writer_id.into(),
104            value: 0,
105        }
106    }
107
108    #[must_use]
109    pub fn with_value(writer_id: impl Into<String>, value: u64) -> Self {
110        Self {
111            writer_id: writer_id.into(),
112            value,
113        }
114    }
115
116    fn next(&mut self) -> Result<u64, DataStoreError> {
117        let next = self.value.checked_add(1).ok_or_else(|| {
118            data_store_error(
119                DataStoreErrorCode::LamportClockOverflow,
120                "lamport clock overflow",
121                json!({ "writer_id": self.writer_id }),
122            )
123        })?;
124        self.value = next;
125        Ok(next)
126    }
127}
128
129pub struct RuntimeDataStore<A> {
130    adapter: A,
131    clock: LamportClock,
132}
133
134impl<A: DataStore> RuntimeDataStore<A> {
135    #[must_use]
136    pub fn new(adapter: A, writer_id: impl Into<String>) -> Self {
137        Self {
138            adapter,
139            clock: LamportClock::new(writer_id),
140        }
141    }
142
143    #[must_use]
144    pub fn with_clock(adapter: A, clock: LamportClock) -> Self {
145        Self { adapter, clock }
146    }
147
148    /// Reads and validates a state value by key.
149    ///
150    /// # Errors
151    ///
152    /// Returns [`DataStoreError`] when the key is invalid, the adapter cannot
153    /// read the key, or the stored value violates the contract state schema.
154    pub fn read(
155        &self,
156        contract: &CapabilityContract,
157        key: &str,
158    ) -> Result<Option<Value>, DataStoreError> {
159        validate_key(key)?;
160        if contract.state_schema.is_none() {
161            return Ok(None);
162        }
163        self.adapter.read(key).and_then(|record| {
164            record
165                .map(|record| {
166                    validate_state_write(contract, key, &record.value)?;
167                    Ok(record.value)
168                })
169                .transpose()
170        })
171    }
172
173    /// Validates, stamps, and writes a state value for a capability contract.
174    ///
175    /// # Errors
176    ///
177    /// Returns [`DataStoreError`] when the key is invalid, no state schema is
178    /// declared, schema validation fails, the Lamport clock overflows, or the
179    /// adapter cannot persist the stamped record.
180    pub fn write(
181        &mut self,
182        contract: &CapabilityContract,
183        key: &str,
184        value: Value,
185    ) -> Result<StateRecord, DataStoreError> {
186        validate_state_write(contract, key, &value)?;
187        let record = StateRecord {
188            key: key.to_string(),
189            value,
190            lamport_clock: self.clock.next()?,
191            writer_id: self.clock.writer_id.clone(),
192        };
193        self.adapter.write(record.clone())?;
194        Ok(record)
195    }
196
197    /// Deletes a state value by key.
198    ///
199    /// # Errors
200    ///
201    /// Returns [`DataStoreError`] when the adapter cannot delete the key.
202    pub fn delete(&mut self, key: &str) -> Result<(), DataStoreError> {
203        self.adapter.delete(key)
204    }
205
206    /// Lists state keys.
207    ///
208    /// # Errors
209    ///
210    /// Returns [`DataStoreError`] when the adapter cannot enumerate keys.
211    pub fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
212        self.adapter.list_keys()
213    }
214
215    /// Triggers explicit sync after a reconnect event.
216    ///
217    /// # Errors
218    ///
219    /// Returns [`DataStoreError`] when either adapter cannot read, write, list,
220    /// or restore state during sync.
221    pub fn sync_on_reconnect(
222        &mut self,
223        remote: &mut dyn DataStore,
224    ) -> Result<SyncReport, DataStoreError> {
225        sync_adapters(&mut self.adapter, remote)
226    }
227
228    pub fn into_inner(self) -> A {
229        self.adapter
230    }
231}
232
233#[derive(Debug, Clone)]
234pub struct LocalFileDataStore {
235    root: PathBuf,
236}
237
238impl LocalFileDataStore {
239    /// Creates a local filesystem-backed data store rooted at `root`.
240    ///
241    /// # Errors
242    ///
243    /// Returns [`DataStoreError`] when the root directory cannot be created.
244    pub fn new(root: impl Into<PathBuf>) -> Result<Self, DataStoreError> {
245        let root = root.into();
246        fs::create_dir_all(&root).map_err(|error| io_error("create data store root", &error))?;
247        Ok(Self { root })
248    }
249
250    fn path_for_key(&self, key: &str) -> Result<PathBuf, DataStoreError> {
251        validate_key(key)?;
252        Ok(self.root.join(format!("{key}.json")))
253    }
254}
255
256impl DataStore for LocalFileDataStore {
257    fn read(&self, key: &str) -> Result<Option<StateRecord>, DataStoreError> {
258        let path = self.path_for_key(key)?;
259        if !path.exists() {
260            return Ok(None);
261        }
262        let text =
263            fs::read_to_string(&path).map_err(|error| io_error("read state record", &error))?;
264        serde_json::from_str(&text)
265            .map(Some)
266            .map_err(|error| serialization_error("deserialize state record", &error))
267    }
268
269    fn write(&mut self, record: StateRecord) -> Result<(), DataStoreError> {
270        let path = self.path_for_key(&record.key)?;
271        let text = serde_json::to_string_pretty(&record)
272            .map_err(|error| serialization_error("serialize state record", &error))?;
273        fs::write(path, text).map_err(|error| io_error("write state record", &error))
274    }
275
276    fn delete(&mut self, key: &str) -> Result<(), DataStoreError> {
277        let path = self.path_for_key(key)?;
278        match fs::remove_file(path) {
279            Ok(()) => Ok(()),
280            Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
281            Err(error) => Err(io_error("delete state record", &error)),
282        }
283    }
284
285    fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
286        let mut keys = Vec::new();
287        for entry in
288            fs::read_dir(&self.root).map_err(|error| io_error("list state keys", &error))?
289        {
290            let entry = entry.map_err(|error| io_error("read state key entry", &error))?;
291            let path = entry.path();
292            if path.extension().and_then(|extension| extension.to_str()) != Some("json") {
293                continue;
294            }
295            if let Some(key) = path.file_stem().and_then(|stem| stem.to_str()) {
296                keys.push(key.to_string());
297            }
298        }
299        keys.sort();
300        Ok(keys)
301    }
302}
303
304/// Validates a capability state write against the contract-declared state schema.
305///
306/// # Errors
307///
308/// Returns [`DataStoreError`] when the key is invalid, the contract does not
309/// declare a state schema, the key is not declared by the schema, or the value
310/// does not match the declared key schema.
311pub fn validate_state_write(
312    contract: &CapabilityContract,
313    key: &str,
314    value: &Value,
315) -> Result<(), DataStoreError> {
316    validate_key(key)?;
317    let schema = contract.state_schema.as_ref().ok_or_else(|| {
318        data_store_error(
319            DataStoreErrorCode::NoStateSchemaDeclared,
320            "no_state_schema_declared",
321            json!({ "capability_id": contract.id, "key": key }),
322        )
323    })?;
324    let property_schema = schema
325        .get("properties")
326        .and_then(Value::as_object)
327        .and_then(|properties| properties.get(key))
328        .ok_or_else(|| {
329            data_store_error(
330                DataStoreErrorCode::SchemaValidationError,
331                "schema_validation_error",
332                json!({ "key": key, "reason": "state key is not declared in schema" }),
333            )
334        })?;
335    let mut violations = Vec::new();
336    crate::validate_value_against_schema(value, property_schema, "$", &mut violations);
337    if violations.is_empty() {
338        Ok(())
339    } else {
340        Err(data_store_error(
341            DataStoreErrorCode::SchemaValidationError,
342            "schema_validation_error",
343            json!({ "key": key, "violations": violations }),
344        ))
345    }
346}
347
348fn sync_adapters(
349    local: &mut dyn DataStore,
350    remote: &mut dyn DataStore,
351) -> Result<SyncReport, DataStoreError> {
352    let keys = merged_keys(local.list_keys()?, remote.list_keys()?);
353    let mut decisions = Vec::new();
354    let mut snapshots = BTreeMap::new();
355
356    for key in keys {
357        let local_record = local.read(&key)?;
358        let remote_record = remote.read(&key)?;
359        snapshots.insert(key.clone(), local_record.clone());
360        let Some((winner, rule)) = merge_records(local_record.as_ref(), remote_record.as_ref())
361        else {
362            continue;
363        };
364        apply_winner(local, remote, &key, &winner).map_err(|error| {
365            rollback_local(local, &snapshots);
366            data_store_error(
367                DataStoreErrorCode::SyncFailure,
368                "sync failed; local state restored",
369                json!({ "key": key, "cause": error.message }),
370            )
371        })?;
372        decisions.push(MergeDecision {
373            key,
374            winning_writer_id: winner.writer_id,
375            winning_lamport_clock: winner.lamport_clock,
376            resolution_rule: rule,
377        });
378    }
379
380    Ok(SyncReport {
381        governing_spec: DATA_STORE_SPEC.to_string(),
382        decisions,
383    })
384}
385
386fn merged_keys(local_keys: Vec<String>, remote_keys: Vec<String>) -> Vec<String> {
387    local_keys
388        .into_iter()
389        .chain(remote_keys)
390        .collect::<BTreeSet<_>>()
391        .into_iter()
392        .collect()
393}
394
395fn merge_records(
396    local: Option<&StateRecord>,
397    remote: Option<&StateRecord>,
398) -> Option<(StateRecord, ConflictResolutionRule)> {
399    match (local, remote) {
400        (Some(record), None) => Some((record.clone(), ConflictResolutionRule::OnlyLocal)),
401        (None, Some(record)) => Some((record.clone(), ConflictResolutionRule::OnlyRemote)),
402        (Some(local), Some(remote)) => Some(select_conflict_winner(local, remote)),
403        (None, None) => None,
404    }
405}
406
407fn select_conflict_winner(
408    local: &StateRecord,
409    remote: &StateRecord,
410) -> (StateRecord, ConflictResolutionRule) {
411    if local.lamport_clock > remote.lamport_clock {
412        return (local.clone(), ConflictResolutionRule::HigherLamportClock);
413    }
414    if remote.lamport_clock > local.lamport_clock {
415        return (remote.clone(), ConflictResolutionRule::HigherLamportClock);
416    }
417    if local.writer_id >= remote.writer_id {
418        (
419            local.clone(),
420            ConflictResolutionRule::WriterIdentityTieBreak,
421        )
422    } else {
423        (
424            remote.clone(),
425            ConflictResolutionRule::WriterIdentityTieBreak,
426        )
427    }
428}
429
430fn apply_winner(
431    local: &mut dyn DataStore,
432    remote: &mut dyn DataStore,
433    key: &str,
434    winner: &StateRecord,
435) -> Result<(), DataStoreError> {
436    if local.read(key)?.as_ref() != Some(winner) {
437        local.write(winner.clone())?;
438    }
439    if remote.read(key)?.as_ref() != Some(winner) {
440        remote.write(winner.clone())?;
441    }
442    Ok(())
443}
444
445fn rollback_local(local: &mut dyn DataStore, snapshots: &BTreeMap<String, Option<StateRecord>>) {
446    for (key, snapshot) in snapshots {
447        let result = match snapshot {
448            Some(record) => local.write(record.clone()),
449            None => local.delete(key),
450        };
451        let _ignored = result.is_ok();
452    }
453}
454
455fn validate_key(key: &str) -> Result<(), DataStoreError> {
456    let valid = !key.is_empty()
457        && key
458            .chars()
459            .all(|character| character.is_ascii_alphanumeric() || matches!(character, '_' | '-'));
460    if valid {
461        Ok(())
462    } else {
463        Err(data_store_error(
464            DataStoreErrorCode::InvalidKey,
465            "state key must be non-empty and contain only ASCII letters, numbers, '_' or '-'",
466            json!({ "key": key }),
467        ))
468    }
469}
470
471fn data_store_error(code: DataStoreErrorCode, message: &str, details: Value) -> DataStoreError {
472    DataStoreError {
473        code,
474        message: message.to_string(),
475        details,
476    }
477}
478
479fn io_error(action: &str, error: &std::io::Error) -> DataStoreError {
480    data_store_error(
481        DataStoreErrorCode::IoFailure,
482        action,
483        json!({ "error": error.to_string() }),
484    )
485}
486
487fn serialization_error(action: &str, error: &serde_json::Error) -> DataStoreError {
488    data_store_error(
489        DataStoreErrorCode::SerializationFailure,
490        action,
491        json!({ "error": error.to_string() }),
492    )
493}
494
495#[cfg(test)]
496#[allow(clippy::expect_used)]
497mod tests {
498    use super::*;
499    use serde_json::json;
500    use std::cell::Cell;
501    use traverse_contracts::{
502        BinaryFormat, CapabilityContract, Condition, DependencyReference, Entrypoint,
503        EntrypointKind, EventReference, Execution, ExecutionConstraints, ExecutionTarget,
504        FilesystemAccess, HostApiAccess, IdReference, Lifecycle, NetworkAccess, Owner, Provenance,
505        ProvenanceSource, SchemaContainer, ServiceType, SideEffect, SideEffectKind,
506        ValidationEvidence,
507    };
508    use uuid::Uuid;
509
510    #[derive(Debug, Clone, Default)]
511    struct MemoryDataStore {
512        records: BTreeMap<String, StateRecord>,
513        fail_writes: Cell<bool>,
514    }
515
516    #[derive(Debug, Clone, Default)]
517    struct PhantomKeyStore;
518
519    impl DataStore for MemoryDataStore {
520        fn read(&self, key: &str) -> Result<Option<StateRecord>, DataStoreError> {
521            Ok(self.records.get(key).cloned())
522        }
523
524        fn write(&mut self, record: StateRecord) -> Result<(), DataStoreError> {
525            if self.fail_writes.get() {
526                return Err(data_store_error(
527                    DataStoreErrorCode::IoFailure,
528                    "forced write failure",
529                    json!({ "key": record.key }),
530                ));
531            }
532            self.records.insert(record.key.clone(), record);
533            Ok(())
534        }
535
536        fn delete(&mut self, key: &str) -> Result<(), DataStoreError> {
537            self.records.remove(key);
538            Ok(())
539        }
540
541        fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
542            Ok(self.records.keys().cloned().collect())
543        }
544    }
545
546    impl DataStore for PhantomKeyStore {
547        fn read(&self, _key: &str) -> Result<Option<StateRecord>, DataStoreError> {
548            Ok(None)
549        }
550
551        fn write(&mut self, _record: StateRecord) -> Result<(), DataStoreError> {
552            Ok(())
553        }
554
555        fn delete(&mut self, _key: &str) -> Result<(), DataStoreError> {
556            Ok(())
557        }
558
559        fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
560            Ok(vec!["phantom".to_string()])
561        }
562    }
563
564    #[test]
565    fn runtime_data_store_validates_writes_and_reads_from_local_file_adapter() {
566        let root = temp_root("valid");
567        let adapter = LocalFileDataStore::new(&root).expect("local adapter should initialize");
568        let mut store = RuntimeDataStore::new(adapter, "writer-a");
569        let contract = stateful_contract(Some(json!({
570            "type": "object",
571            "properties": {
572                "draft": {"type": "string"}
573            }
574        })));
575
576        let record = store
577            .write(&contract, "draft", json!("ready"))
578            .expect("valid state write should succeed");
579
580        assert_eq!(record.lamport_clock, 1);
581        assert_eq!(
582            store.read(&contract, "draft").expect("read should succeed"),
583            Some(json!("ready"))
584        );
585        assert_eq!(
586            store.list_keys().expect("list should succeed"),
587            vec!["draft".to_string()]
588        );
589        store.delete("draft").expect("delete should succeed");
590        assert_eq!(
591            store.read(&contract, "draft").expect("read should succeed"),
592            None
593        );
594    }
595
596    #[test]
597    fn runtime_data_store_rejects_missing_schema_bad_keys_and_schema_violations() {
598        let adapter = MemoryDataStore::default();
599        let mut store = RuntimeDataStore::new(adapter, "writer-a");
600        let no_schema = stateful_contract(None);
601        let schema = stateful_contract(Some(json!({
602            "type": "object",
603            "properties": {
604                "count": {"type": "integer"}
605            }
606        })));
607
608        let missing = store
609            .write(&no_schema, "count", json!(1))
610            .expect_err("missing state schema should fail");
611        assert_eq!(missing.code, DataStoreErrorCode::NoStateSchemaDeclared);
612
613        let invalid_key = store
614            .write(&schema, "bad.key", json!(1))
615            .expect_err("invalid key should fail");
616        assert_eq!(invalid_key.code, DataStoreErrorCode::InvalidKey);
617
618        let undeclared = store
619            .write(&schema, "other", json!(1))
620            .expect_err("undeclared state key should fail");
621        assert_eq!(undeclared.code, DataStoreErrorCode::SchemaValidationError);
622
623        let wrong_type = store
624            .write(&schema, "count", json!("one"))
625            .expect_err("wrong state type should fail");
626        assert_eq!(wrong_type.code, DataStoreErrorCode::SchemaValidationError);
627
628        let no_schema_read = store
629            .read(&no_schema, "count")
630            .expect("no-schema read should succeed");
631        assert_eq!(no_schema_read, None);
632
633        let bad_read_key = store
634            .read(&schema, "bad.key")
635            .expect_err("invalid read key should fail");
636        assert_eq!(bad_read_key.code, DataStoreErrorCode::InvalidKey);
637    }
638
639    #[test]
640    fn lamport_clock_overflow_is_rejected_before_adapter_write() {
641        let adapter = MemoryDataStore::default();
642        let clock = LamportClock::with_value("writer-a", u64::MAX);
643        let mut store = RuntimeDataStore::with_clock(adapter, clock);
644        let contract = stateful_contract(Some(json!({
645            "type": "object",
646            "properties": {
647                "draft": {"type": "string"}
648            }
649        })));
650
651        let error = store
652            .write(&contract, "draft", json!("ready"))
653            .expect_err("overflow should fail");
654
655        assert_eq!(error.code, DataStoreErrorCode::LamportClockOverflow);
656        assert!(store.into_inner().records.is_empty());
657    }
658
659    #[test]
660    fn runtime_data_store_validates_reads_before_returning_stored_values() {
661        let mut adapter = MemoryDataStore::default();
662        adapter
663            .write(record("count", "writer-a", 1, json!("not an integer")))
664            .expect("seed should succeed");
665        let store = RuntimeDataStore::new(adapter, "writer-a");
666        let contract = stateful_contract(Some(json!({
667            "type": "object",
668            "properties": {
669                "count": {"type": "integer"}
670            }
671        })));
672
673        let error = store
674            .read(&contract, "count")
675            .expect_err("invalid stored value should fail");
676
677        assert_eq!(error.code, DataStoreErrorCode::SchemaValidationError);
678    }
679
680    #[test]
681    fn reconnect_sync_merges_only_local_only_remote_clock_winner_and_writer_tie_breaks() {
682        let mut local = MemoryDataStore::default();
683        let mut remote = MemoryDataStore::default();
684        local
685            .write(record("local_only", "local-a", 1, json!("local")))
686            .expect("local write should succeed");
687        remote
688            .write(record("remote_only", "remote-a", 1, json!("remote")))
689            .expect("remote write should succeed");
690        local
691            .write(record("clock", "local-a", 2, json!("old")))
692            .expect("local write should succeed");
693        remote
694            .write(record("clock", "remote-a", 3, json!("new")))
695            .expect("remote write should succeed");
696        local
697            .write(record("tie", "writer-z", 4, json!("winner")))
698            .expect("local write should succeed");
699        remote
700            .write(record("tie", "writer-a", 4, json!("loser")))
701            .expect("remote write should succeed");
702
703        let report = sync_adapters(&mut local, &mut remote).expect("sync should succeed");
704
705        assert_eq!(report.governing_spec, "032-universal-data-access");
706        assert_eq!(report.decisions.len(), 4);
707        assert_eq!(
708            local.read("remote_only").expect("read should succeed"),
709            remote.read("remote_only").expect("read should succeed")
710        );
711        assert_eq!(
712            local.read("clock").expect("read should succeed"),
713            Some(record("clock", "remote-a", 3, json!("new")))
714        );
715        assert_eq!(
716            remote.read("tie").expect("read should succeed"),
717            Some(record("tie", "writer-z", 4, json!("winner")))
718        );
719        assert!(
720            report
721                .decisions
722                .iter()
723                .any(|decision| decision.resolution_rule
724                    == ConflictResolutionRule::WriterIdentityTieBreak)
725        );
726    }
727
728    #[test]
729    fn sync_failure_restores_local_snapshot() {
730        let mut local = MemoryDataStore::default();
731        let mut remote = MemoryDataStore::default();
732        local
733            .write(record("shared", "local-a", 2, json!("local")))
734            .expect("local write should succeed");
735        remote
736            .write(record("shared", "remote-a", 1, json!("remote")))
737            .expect("remote write should succeed");
738        remote.fail_writes.set(true);
739
740        let error = sync_adapters(&mut local, &mut remote).expect_err("sync should fail");
741
742        assert_eq!(error.code, DataStoreErrorCode::SyncFailure);
743        assert_eq!(
744            local.read("shared").expect("read should succeed"),
745            Some(record("shared", "local-a", 2, json!("local")))
746        );
747    }
748
749    #[test]
750    fn local_file_adapter_reports_bad_keys_and_bad_json() {
751        let root = temp_root("bad-json");
752        let adapter = LocalFileDataStore::new(&root).expect("local adapter should initialize");
753        let invalid = adapter
754            .read("bad.key")
755            .expect_err("invalid key should fail");
756        assert_eq!(invalid.code, DataStoreErrorCode::InvalidKey);
757
758        fs::write(root.join("broken.json"), "{").expect("bad json fixture should write");
759        let invalid_json = adapter
760            .read("broken")
761            .expect_err("invalid json should fail");
762        assert_eq!(invalid_json.code, DataStoreErrorCode::SerializationFailure);
763    }
764
765    #[test]
766    fn helper_paths_cover_remaining_datastore_branches() {
767        let mut local = RuntimeDataStore::new(MemoryDataStore::default(), "local-a");
768        let mut remote = MemoryDataStore::default();
769        remote
770            .write(record("remote_only", "remote-a", 1, json!("remote")))
771            .expect("remote seed should succeed");
772
773        let report = local
774            .sync_on_reconnect(&mut remote)
775            .expect("public reconnect sync should succeed");
776        assert_eq!(report.decisions.len(), 1);
777
778        assert!(merge_records(None, None).is_none());
779        let (_winner, rule) = select_conflict_winner(
780            &record("tie", "writer-a", 1, json!("local")),
781            &record("tie", "writer-z", 1, json!("remote")),
782        );
783        assert_eq!(rule, ConflictResolutionRule::WriterIdentityTieBreak);
784
785        let mut failing_local = MemoryDataStore::default();
786        failing_local.fail_writes.set(true);
787        let mut seeded_remote = MemoryDataStore::default();
788        seeded_remote
789            .write(record("missing_local", "remote-a", 1, json!("remote")))
790            .expect("remote seed should succeed");
791        let error =
792            sync_adapters(&mut failing_local, &mut seeded_remote).expect_err("sync should fail");
793        assert_eq!(error.code, DataStoreErrorCode::SyncFailure);
794        assert_eq!(
795            failing_local
796                .delete("missing_local")
797                .expect("delete should succeed"),
798            ()
799        );
800
801        let mut phantom_local = PhantomKeyStore;
802        let mut phantom_remote = PhantomKeyStore;
803        assert!(
804            sync_adapters(&mut phantom_local, &mut phantom_remote)
805                .expect("phantom sync should succeed")
806                .decisions
807                .is_empty()
808        );
809        phantom_local
810            .write(record("phantom", "writer-a", 1, json!("value")))
811            .expect("phantom write should succeed");
812        phantom_local
813            .delete("phantom")
814            .expect("phantom delete should succeed");
815
816        let root = temp_root("listing");
817        fs::create_dir_all(&root).expect("root should be created");
818        fs::write(root.join("skip.txt"), "not state").expect("non-json fixture should write");
819        let adapter = LocalFileDataStore::new(&root).expect("local adapter should initialize");
820        assert!(adapter.list_keys().expect("list should succeed").is_empty());
821        let mut delete_missing =
822            LocalFileDataStore::new(&root).expect("local adapter should initialize");
823        delete_missing
824            .delete("missing")
825            .expect("missing delete should succeed");
826        fs::create_dir(root.join("cant_delete.json")).expect("directory fixture should write");
827        let delete_failure = delete_missing
828            .delete("cant_delete")
829            .expect_err("directory delete should fail");
830        assert_eq!(delete_failure.code, DataStoreErrorCode::IoFailure);
831
832        let file_root = temp_root("file-root");
833        fs::write(&file_root, "not a directory").expect("file root fixture should write");
834        let io_failure = LocalFileDataStore::new(&file_root).expect_err("file root should fail");
835        assert_eq!(io_failure.code, DataStoreErrorCode::IoFailure);
836    }
837
838    fn record(key: &str, writer_id: &str, lamport_clock: u64, value: Value) -> StateRecord {
839        StateRecord {
840            key: key.to_string(),
841            value,
842            lamport_clock,
843            writer_id: writer_id.to_string(),
844        }
845    }
846
847    fn temp_root(name: &str) -> PathBuf {
848        std::env::temp_dir().join(format!("traverse-data-store-{name}-{}", Uuid::new_v4()))
849    }
850
851    fn stateful_contract(state_schema: Option<Value>) -> CapabilityContract {
852        CapabilityContract {
853            kind: "capability_contract".to_string(),
854            schema_version: "1.0.0".to_string(),
855            id: "stateful.example".to_string(),
856            namespace: "stateful".to_string(),
857            name: "example".to_string(),
858            version: "1.0.0".to_string(),
859            lifecycle: Lifecycle::Active,
860            owner: Owner {
861                team: "runtime".to_string(),
862                contact: "runtime@example.com".to_string(),
863            },
864            summary: "Stateful test capability".to_string(),
865            description: "Stateful test capability".to_string(),
866            inputs: SchemaContainer {
867                schema: json!({"type": "object"}),
868            },
869            outputs: SchemaContainer {
870                schema: json!({"type": "object"}),
871            },
872            preconditions: Vec::<Condition>::new(),
873            postconditions: Vec::<Condition>::new(),
874            side_effects: vec![SideEffect {
875                kind: SideEffectKind::StateChange,
876                description: "writes capability state".to_string(),
877            }],
878            emits: Vec::<EventReference>::new(),
879            consumes: Vec::<EventReference>::new(),
880            permissions: Vec::<IdReference>::new(),
881            execution: Execution {
882                binary_format: BinaryFormat::Wasm,
883                constraints: ExecutionConstraints {
884                    network_access: NetworkAccess::Forbidden,
885                    filesystem_access: FilesystemAccess::SandboxOnly,
886                    host_api_access: HostApiAccess::None,
887                },
888                entrypoint: Entrypoint {
889                    kind: EntrypointKind::WasiCommand,
890                    command: "run".to_string(),
891                },
892                preferred_targets: vec![ExecutionTarget::Local],
893            },
894            policies: Vec::<IdReference>::new(),
895            dependencies: Vec::<DependencyReference>::new(),
896            provenance: Provenance {
897                source: ProvenanceSource::Greenfield,
898                author: "Codex".to_string(),
899                created_at: "2026-04-19T00:00:00Z".to_string(),
900                spec_ref: Some("032-universal-data-access".to_string()),
901                adr_refs: Vec::new(),
902                exception_refs: Vec::new(),
903            },
904            evidence: Vec::<ValidationEvidence>::new(),
905            service_type: ServiceType::Stateful,
906            permitted_targets: vec![ExecutionTarget::Local],
907            event_trigger: None,
908            connector_requirements: Vec::new(),
909            state_schema,
910        }
911    }
912}