1use 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 fn read(&self, key: &str) -> Result<Option<StateRecord>, DataStoreError>;
70
71 fn write(&mut self, record: StateRecord) -> Result<(), DataStoreError>;
77
78 fn delete(&mut self, key: &str) -> Result<(), DataStoreError>;
84
85 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 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 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 pub fn delete(&mut self, key: &str) -> Result<(), DataStoreError> {
203 self.adapter.delete(key)
204 }
205
206 pub fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
212 self.adapter.list_keys()
213 }
214
215 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 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
304pub 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}