1#[path = "data_store_coordinator.rs"]
9mod coordinator;
10#[path = "data_store_hosted_sync.rs"]
11mod hosted_sync;
12#[path = "data_store_maintenance.rs"]
13mod maintenance;
14#[path = "data_store_remote.rs"]
15mod remote;
16pub use coordinator::{DataStoreCoordinator, DataStoreCoordinatorError};
17pub use hosted_sync::{
18 AblyEdgeError, AblyHistoryBatch, AblyHostedSyncTransport, AblyRealtimeEdge,
19 EncryptedSyncOperation, HostedSyncConnectionState, HostedSyncCredential,
20 HostedSyncDegradedReason, HostedSyncError, HostedSyncErrorCode, HostedSyncLineageEvidence,
21 HostedSyncObservation, HostedSyncPublishReceipt, HostedSyncReplayResult, HostedSyncTransport,
22 InMemoryAblyEdge, InMemoryHostedSyncTransport, SyncScopeId, run_hosted_sync_conformance,
23};
24pub use maintenance::{
25 BackupManifest, BackupRecordIndexEntry, DataStoreMaintenance, DataStoreMigration,
26 LocalFileDataStoreMaintenance, MaintenanceError, MaintenanceErrorCode, MaintenanceEvidence,
27 MigrationError, MigrationErrorCode, MigrationReport, RetentionPolicy,
28};
29pub use remote::{
30 RemoteBackendFailure, RemoteDataStoreBackend, RemoteKeyValueDataStore, RemoteObject,
31 RemoteOperationEvidence, RemoteVersionToken, RemoteWriteOutcome,
32};
33
34#[cfg(feature = "datastore-encryption")]
35use aes_gcm::aead::consts::U12;
36#[cfg(feature = "datastore-encryption")]
37use aes_gcm::aead::{Aead, KeyInit, Payload};
38#[cfg(feature = "datastore-encryption")]
39use aes_gcm::{Aes256Gcm, Nonce};
40use serde::{Deserialize, Serialize};
41use serde_json::{Value, json};
42use sha2::{Digest, Sha256};
43use std::collections::{BTreeMap, BTreeSet};
44use std::fs::{self, File, OpenOptions, TryLockError};
45use std::io::Write;
46use std::path::PathBuf;
47use std::sync::Arc;
48use traverse_contracts::CapabilityContract;
49use zeroize::Zeroizing;
50
51const DATA_STORE_SPEC: &str = "089-datastore-synchronization";
54const LOCAL_DATA_STORE_FORMAT: &str = "local-datastore/1";
55const LOCAL_DATA_STORE_V2_FORMAT: &str = "local-datastore/2";
56const LOCAL_DATA_STORE_V2_FORMAT_VERSION: u32 = 2;
57const LOCAL_DATA_STORE_V2_SCHEMA_VERSION: u32 = 1;
58const LOCAL_DATA_STORE_LOCK_FILE: &str = ".traverse-datastore.lock";
59const HEXADECIMAL_DIGITS: &[u8; 16] = b"0123456789abcdef";
60
61#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
62pub struct StateRecord {
63 pub key: String,
64 pub value: Value,
65 pub lamport_clock: u64,
66 pub writer_id: String,
67}
68
69#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
70pub struct MergeDecision {
71 pub key: String,
72 pub winning_writer_id: String,
73 pub winning_lamport_clock: u64,
74 pub resolution_rule: ConflictResolutionRule,
75}
76
77#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
78#[serde(rename_all = "snake_case")]
79pub enum ConflictResolutionRule {
80 OnlyLocal,
81 OnlyRemote,
82 HigherLamportClock,
83 WriterIdentityTieBreak,
84}
85
86#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
87pub struct SyncReport {
88 pub governing_spec: String,
89 pub decisions: Vec<MergeDecision>,
90}
91
92#[derive(Debug, Clone, PartialEq, Eq)]
93pub struct DataStoreError {
94 pub code: DataStoreErrorCode,
95 pub message: String,
96 pub details: Value,
97}
98
99#[derive(Debug, Clone, Copy, PartialEq, Eq)]
100pub enum DataStoreErrorCode {
101 SchemaValidationError,
102 NoStateSchemaDeclared,
103 LamportClockOverflow,
104 InvalidKey,
105 IoFailure,
106 SerializationFailure,
107 SyncFailure,
108 IntegrityCheckFailed,
109 StoreLocked,
110 DurabilityCommitFailed,
111 KeyProviderRequired,
112 KeyNotFound,
113 KeyExpired,
114 KeyProviderFailure,
115 CryptoFailure,
116 ClassificationChangeNotAllowed,
117 RemoteConflict,
118 RemoteUnavailable,
119 RemoteTimeout,
120 RemoteOutcomeUnknown,
121 RemoteUnauthorized,
122 RemoteScopeDenied,
123 RemoteIntegrityFailed,
124 RemoteBackendFailed,
125}
126
127#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
129#[serde(rename_all = "snake_case")]
130pub enum LocalDataClassification {
131 Public,
132 Private,
133}
134
135#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
136struct LocalDataStoreEnvelope {
137 format: String,
138 classification: LocalDataClassification,
139 digest: String,
140 #[serde(default, skip_serializing_if = "Option::is_none")]
141 record: Option<StateRecord>,
142 #[serde(default, skip_serializing_if = "Option::is_none")]
143 record_key: Option<String>,
144 #[serde(default, skip_serializing_if = "Option::is_none")]
145 key_id: Option<String>,
146 #[serde(default, skip_serializing_if = "Option::is_none")]
147 nonce: Option<String>,
148 #[serde(default, skip_serializing_if = "Option::is_none")]
149 ciphertext: Option<String>,
150 #[serde(default, skip_serializing_if = "Option::is_none")]
153 retained_at: Option<String>,
154}
155
156#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
159struct LocalDataStoreV2Envelope {
160 format: String,
161 format_version: u32,
162 schema_version: u32,
163 payload_integrity: String,
164 integrity: LocalDataStoreIntegrity,
165 encryption_disclosure: String,
166 payload: LocalDataStoreEnvelope,
167}
168
169#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
170struct LocalDataStoreIntegrity {
171 algorithm: String,
172 content_digest: String,
173}
174
175fn v2_envelope(
176 payload: LocalDataStoreEnvelope,
177) -> Result<LocalDataStoreV2Envelope, DataStoreError> {
178 let payload_bytes = serde_json::to_vec(&payload)
179 .map_err(|error| serialization_error("serialize v2 data store payload", &error))?;
180 let payload_integrity = digest_bytes(&payload_bytes);
181 Ok(LocalDataStoreV2Envelope {
182 format: LOCAL_DATA_STORE_V2_FORMAT.to_string(),
183 format_version: LOCAL_DATA_STORE_V2_FORMAT_VERSION,
184 schema_version: LOCAL_DATA_STORE_V2_SCHEMA_VERSION,
185 integrity: LocalDataStoreIntegrity {
186 algorithm: "sha256".to_string(),
187 content_digest: payload_integrity.clone(),
188 },
189 payload_integrity,
190 encryption_disclosure: match payload.classification {
191 LocalDataClassification::Public => "not_enabled".to_string(),
192 LocalDataClassification::Private => "host_managed_opaque".to_string(),
193 },
194 payload,
195 })
196}
197
198fn decode_v2_envelope(value: Value) -> Result<LocalDataStoreEnvelope, DataStoreError> {
199 let envelope: LocalDataStoreV2Envelope =
200 serde_json::from_value(value).map_err(|_| integrity_error("malformed_envelope"))?;
201 if envelope.format != LOCAL_DATA_STORE_V2_FORMAT
202 || envelope.format_version != LOCAL_DATA_STORE_V2_FORMAT_VERSION
203 || envelope.schema_version != LOCAL_DATA_STORE_V2_SCHEMA_VERSION
204 || envelope.integrity.algorithm != "sha256"
205 || envelope.integrity.content_digest != envelope.payload_integrity
206 || !matches!(
207 envelope.encryption_disclosure.as_str(),
208 "not_enabled" | "host_managed_opaque"
209 )
210 {
211 return Err(integrity_error("unknown_format_version"));
212 }
213 let payload_bytes = serde_json::to_vec(&envelope.payload)
214 .map_err(|error| serialization_error("serialize v2 payload for verification", &error))?;
215 if digest_bytes(&payload_bytes) != envelope.payload_integrity {
216 return Err(integrity_error("digest_mismatch"));
217 }
218 if envelope.payload.format != LOCAL_DATA_STORE_FORMAT {
219 return Err(integrity_error("unknown_format_version"));
220 }
221 Ok(envelope.payload)
222}
223
224#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
226#[serde(rename_all = "snake_case")]
227pub enum KeyProviderErrorCode {
228 MissingKey,
229 ExpiredKeyId,
230 ProviderFailure,
231}
232
233#[derive(Debug, Clone, PartialEq, Eq)]
235pub struct KeyProviderError {
236 pub code: KeyProviderErrorCode,
237 pub message: String,
238}
239
240impl KeyProviderError {
241 #[must_use]
242 pub fn missing_key() -> Self {
243 Self {
244 code: KeyProviderErrorCode::MissingKey,
245 message: "key_not_found".to_string(),
246 }
247 }
248
249 #[must_use]
250 pub fn expired_key_id() -> Self {
251 Self {
252 code: KeyProviderErrorCode::ExpiredKeyId,
253 message: "key_expired".to_string(),
254 }
255 }
256
257 #[must_use]
258 pub fn provider_failure() -> Self {
259 Self {
260 code: KeyProviderErrorCode::ProviderFailure,
261 message: "key_provider_failed".to_string(),
262 }
263 }
264}
265
266pub trait KeyProvider: Send + Sync {
268 fn active_key_id(&self) -> Result<String, KeyProviderError>;
274
275 fn key_for(&self, key_id: &str) -> Result<Zeroizing<[u8; 32]>, KeyProviderError>;
282}
283
284#[derive(Clone)]
286pub struct InMemoryKeyProvider {
287 active_key_id: String,
288 keys: BTreeMap<String, Zeroizing<[u8; 32]>>,
289 expired_key_ids: BTreeSet<String>,
290}
291
292impl InMemoryKeyProvider {
293 #[must_use]
294 pub fn new(active_key_id: impl Into<String>, key: [u8; 32]) -> Self {
295 let active_key_id = active_key_id.into();
296 let mut keys = BTreeMap::new();
297 keys.insert(active_key_id.clone(), Zeroizing::new(key));
298 Self {
299 active_key_id,
300 keys,
301 expired_key_ids: BTreeSet::new(),
302 }
303 }
304
305 #[must_use]
306 pub fn with_read_key(mut self, key_id: impl Into<String>, key: [u8; 32]) -> Self {
307 self.keys.insert(key_id.into(), Zeroizing::new(key));
308 self
309 }
310
311 #[must_use]
312 pub fn with_expired_key_id(mut self, key_id: impl Into<String>) -> Self {
313 self.expired_key_ids.insert(key_id.into());
314 self
315 }
316}
317
318impl KeyProvider for InMemoryKeyProvider {
319 fn active_key_id(&self) -> Result<String, KeyProviderError> {
320 if self.expired_key_ids.contains(&self.active_key_id) {
321 return Err(KeyProviderError::expired_key_id());
322 }
323 if self.keys.contains_key(&self.active_key_id) {
324 Ok(self.active_key_id.clone())
325 } else {
326 Err(KeyProviderError::missing_key())
327 }
328 }
329
330 fn key_for(&self, key_id: &str) -> Result<Zeroizing<[u8; 32]>, KeyProviderError> {
331 if self.expired_key_ids.contains(key_id) {
332 return Err(KeyProviderError::expired_key_id());
333 }
334 self.keys
335 .get(key_id)
336 .cloned()
337 .ok_or_else(KeyProviderError::missing_key)
338 }
339}
340
341pub trait DataStore {
342 fn read(&self, key: &str) -> Result<Option<StateRecord>, DataStoreError>;
348
349 fn write(&mut self, record: StateRecord) -> Result<(), DataStoreError>;
355
356 fn delete(&mut self, key: &str) -> Result<(), DataStoreError>;
362
363 fn list_keys(&self) -> Result<Vec<String>, DataStoreError>;
369}
370
371#[derive(Debug, Clone, PartialEq, Eq)]
372pub struct LamportClock {
373 writer_id: String,
374 value: u64,
375}
376
377impl LamportClock {
378 #[must_use]
379 pub fn new(writer_id: impl Into<String>) -> Self {
380 Self {
381 writer_id: writer_id.into(),
382 value: 0,
383 }
384 }
385
386 #[must_use]
387 pub fn with_value(writer_id: impl Into<String>, value: u64) -> Self {
388 Self {
389 writer_id: writer_id.into(),
390 value,
391 }
392 }
393
394 fn next(&mut self) -> Result<u64, DataStoreError> {
395 let next = self.value.checked_add(1).ok_or_else(|| {
396 data_store_error(
397 DataStoreErrorCode::LamportClockOverflow,
398 "lamport clock overflow",
399 json!({ "writer_id": self.writer_id }),
400 )
401 })?;
402 self.value = next;
403 Ok(next)
404 }
405}
406
407pub struct RuntimeDataStore<A> {
408 adapter: A,
409 clock: LamportClock,
410}
411
412impl<A: DataStore> RuntimeDataStore<A> {
413 #[must_use]
414 pub fn new(adapter: A, writer_id: impl Into<String>) -> Self {
415 Self {
416 adapter,
417 clock: LamportClock::new(writer_id),
418 }
419 }
420
421 #[must_use]
422 pub fn with_clock(adapter: A, clock: LamportClock) -> Self {
423 Self { adapter, clock }
424 }
425
426 pub fn read(
433 &self,
434 contract: &CapabilityContract,
435 key: &str,
436 ) -> Result<Option<Value>, DataStoreError> {
437 validate_key(key)?;
438 if contract.state_schema.is_none() {
439 return Ok(None);
440 }
441 self.adapter.read(key).and_then(|record| {
442 record
443 .map(|record| {
444 validate_state_write(contract, key, &record.value)?;
445 Ok(record.value)
446 })
447 .transpose()
448 })
449 }
450
451 pub fn write(
459 &mut self,
460 contract: &CapabilityContract,
461 key: &str,
462 value: Value,
463 ) -> Result<StateRecord, DataStoreError> {
464 validate_state_write(contract, key, &value)?;
465 let record = StateRecord {
466 key: key.to_string(),
467 value,
468 lamport_clock: self.clock.next()?,
469 writer_id: self.clock.writer_id.clone(),
470 };
471 self.adapter.write(record.clone())?;
472 Ok(record)
473 }
474
475 pub fn delete(&mut self, key: &str) -> Result<(), DataStoreError> {
481 self.adapter.delete(key)
482 }
483
484 pub fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
490 self.adapter.list_keys()
491 }
492
493 pub fn sync_on_reconnect(
500 &mut self,
501 remote: &mut dyn DataStore,
502 ) -> Result<SyncReport, DataStoreError> {
503 sync_adapters(&mut self.adapter, remote)
504 }
505
506 pub fn into_inner(self) -> A {
507 self.adapter
508 }
509}
510
511pub struct LocalFileDataStore {
512 root: PathBuf,
513 classification: LocalDataClassification,
514 lock_file: File,
515 key_provider: Option<Arc<dyn KeyProvider>>,
516 write_retained_at: Option<String>,
518}
519
520impl std::fmt::Debug for LocalFileDataStore {
521 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
522 formatter
523 .debug_struct("LocalFileDataStore")
524 .field("root", &self.root)
525 .field("classification", &self.classification)
526 .field("key_provider_configured", &self.key_provider.is_some())
527 .field("write_retained_at", &self.write_retained_at)
528 .finish_non_exhaustive()
529 }
530}
531
532impl Drop for LocalFileDataStore {
533 fn drop(&mut self) {
534 let _ = self.lock_file.unlock();
535 }
536}
537
538impl LocalFileDataStore {
539 pub fn new(root: impl Into<PathBuf>) -> Result<Self, DataStoreError> {
546 Self::with_classification(root, LocalDataClassification::Private)
547 }
548
549 pub fn with_classification(
557 root: impl Into<PathBuf>,
558 classification: LocalDataClassification,
559 ) -> Result<Self, DataStoreError> {
560 let root = root.into();
561 fs::create_dir_all(&root).map_err(|error| io_error("create data store root", &error))?;
562 let lock_path = root.join(LOCAL_DATA_STORE_LOCK_FILE);
563 let lock_file = OpenOptions::new()
564 .read(true)
565 .write(true)
566 .create(true)
567 .truncate(false)
568 .open(&lock_path)
569 .map_err(|error| io_error("open data store lock", &error))?;
570 lock_file.try_lock().map_err(lock_error)?;
571 Ok(Self {
572 root,
573 classification,
574 lock_file,
575 key_provider: None,
576 write_retained_at: None,
577 })
578 }
579
580 #[must_use]
582 pub fn with_key_provider(mut self, key_provider: Arc<dyn KeyProvider>) -> Self {
583 self.key_provider = Some(key_provider);
584 self
585 }
586
587 pub fn set_write_retained_at(&mut self, retained_at: Option<String>) {
592 self.write_retained_at = retained_at;
593 }
594
595 #[must_use]
597 pub fn root(&self) -> &PathBuf {
598 &self.root
599 }
600
601 fn path_for_key(&self, key: &str) -> Result<PathBuf, DataStoreError> {
602 validate_key(key)?;
603 Ok(self.root.join(format!("{key}.json")))
604 }
605
606 fn temporary_path_for_key(&self, key: &str) -> PathBuf {
607 self.root.join(format!(".{key}.{}.tmp", std::process::id()))
608 }
609
610 fn sync_root(&self) -> Result<(), DataStoreError> {
611 File::open(&self.root)
612 .and_then(|directory| directory.sync_all())
613 .map_err(|error| durability_error(&self.root, "parent_directory", &error))
614 }
615}
616
617impl DataStore for LocalFileDataStore {
618 fn read(&self, key: &str) -> Result<Option<StateRecord>, DataStoreError> {
619 let path = self.path_for_key(key)?;
620 if !path.exists() {
621 return Ok(None);
622 }
623 let text =
624 fs::read_to_string(&path).map_err(|error| io_error("read state record", &error))?;
625 let value: Value =
626 serde_json::from_str(&text).map_err(|_| integrity_error("malformed_envelope"))?;
627 if value.get("format").is_none() {
628 return Err(integrity_error("legacy_unverified"));
629 }
630 let envelope: LocalDataStoreEnvelope = match value.get("format").and_then(Value::as_str) {
631 Some(LOCAL_DATA_STORE_FORMAT) => {
632 serde_json::from_value(value).map_err(|_| integrity_error("malformed_envelope"))?
633 }
634 Some(LOCAL_DATA_STORE_V2_FORMAT) => decode_v2_envelope(value)?,
635 _ => return Err(integrity_error("unknown_format_version")),
636 };
637 match envelope.classification {
638 LocalDataClassification::Public => {
639 let record = envelope
640 .record
641 .ok_or_else(|| integrity_error("malformed_envelope"))?;
642 if envelope.key_id.is_some()
643 || envelope.nonce.is_some()
644 || envelope.ciphertext.is_some()
645 || envelope.record_key.is_some()
646 {
647 return Err(integrity_error("malformed_envelope"));
648 }
649 let expected_digest = digest_for_record(&record)?;
650 if envelope.digest != expected_digest {
651 return Err(integrity_error("digest_mismatch"));
652 }
653 Ok(Some(record))
654 }
655 LocalDataClassification::Private => self.decrypt_private_envelope(key, envelope),
656 }
657 }
658
659 fn write(&mut self, record: StateRecord) -> Result<(), DataStoreError> {
660 let path = self.path_for_key(&record.key)?;
661 let write_v2 = self.ensure_classification_unchanged(&path)?;
662 let record_key = record.key.clone();
663 let envelope = match self.classification {
664 LocalDataClassification::Public => LocalDataStoreEnvelope {
665 format: LOCAL_DATA_STORE_FORMAT.to_string(),
666 classification: self.classification,
667 digest: digest_for_record(&record)?,
668 record: Some(record),
669 record_key: None,
670 key_id: None,
671 nonce: None,
672 ciphertext: None,
673 retained_at: self.write_retained_at.clone(),
674 },
675 LocalDataClassification::Private => self.encrypt_private_record(record)?,
676 };
677 let text = if write_v2 {
678 serde_json::to_vec(&v2_envelope(envelope)?)
679 } else {
680 serde_json::to_vec(&envelope)
681 }
682 .map_err(|error| serialization_error("serialize state record envelope", &error))?;
683 let temporary_path = self.temporary_path_for_key(&record_key);
684 let write_result = (|| {
685 let mut temporary_file = File::create(&temporary_path)
686 .map_err(|error| io_error("create temporary state record", &error))?;
687 temporary_file
688 .write_all(&text)
689 .map_err(|error| io_error("write temporary state record", &error))?;
690 temporary_file
691 .sync_all()
692 .map_err(|error| durability_error(&self.root, "temporary_file", &error))?;
693 fs::rename(&temporary_path, &path)
694 .map_err(|error| io_error("atomically commit state record", &error))?;
695 self.sync_root()
696 })();
697 discard_temporary_record(&temporary_path);
698 write_result
699 }
700
701 fn delete(&mut self, key: &str) -> Result<(), DataStoreError> {
702 let path = self.path_for_key(key)?;
703 match fs::remove_file(path) {
704 Ok(()) => self.sync_root(),
705 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
706 Err(error) => Err(io_error("delete state record", &error)),
707 }
708 }
709
710 fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
711 let mut keys = Vec::new();
712 for entry in
713 fs::read_dir(&self.root).map_err(|error| io_error("list state keys", &error))?
714 {
715 let entry = entry.map_err(|error| io_error("read state key entry", &error))?;
716 let path = entry.path();
717 if path.extension().and_then(|extension| extension.to_str()) != Some("json") {
718 continue;
719 }
720 if let Some(key) = path.file_stem().and_then(|stem| stem.to_str()) {
721 keys.push(key.to_string());
722 }
723 }
724 keys.sort();
725 Ok(keys)
726 }
727}
728
729impl LocalFileDataStore {
730 #[cfg(feature = "datastore-encryption")]
731 fn encrypt_private_record(
732 &self,
733 record: StateRecord,
734 ) -> Result<LocalDataStoreEnvelope, DataStoreError> {
735 let provider = self.required_key_provider()?;
736 let key_id = provider.active_key_id().map_err(map_key_provider_error)?;
737 let key = provider.key_for(&key_id).map_err(map_key_provider_error)?;
738 let plaintext = Zeroizing::new(
739 serde_json::to_vec(&record)
740 .map_err(|error| serialization_error("serialize private state record", &error))?,
741 );
742 let cipher = Aes256Gcm::new_from_slice(key.as_ref()).map_err(|_| crypto_error())?;
743 let nonce = fresh_aes_nonce();
744 let aad = private_record_aad(&key_id, &record.key);
745 let ciphertext = cipher
746 .encrypt(
747 &nonce,
748 Payload {
749 msg: plaintext.as_ref(),
750 aad: &aad,
751 },
752 )
753 .map_err(|_| crypto_error())?;
754 let nonce = hex_encode(&nonce);
755 let ciphertext = hex_encode(&ciphertext);
756 let digest = digest_for_private_envelope(&key_id, &record.key, &nonce, &ciphertext);
757 Ok(LocalDataStoreEnvelope {
758 format: LOCAL_DATA_STORE_FORMAT.to_string(),
759 classification: LocalDataClassification::Private,
760 digest,
761 record: None,
762 record_key: Some(record.key),
763 key_id: Some(key_id),
764 nonce: Some(nonce),
765 ciphertext: Some(ciphertext),
766 retained_at: self.write_retained_at.clone(),
767 })
768 }
769
770 #[cfg(not(feature = "datastore-encryption"))]
771 fn encrypt_private_record(
772 &self,
773 _record: StateRecord,
774 ) -> Result<LocalDataStoreEnvelope, DataStoreError> {
775 Err(encryption_feature_disabled())
776 }
777
778 #[cfg(feature = "datastore-encryption")]
779 fn decrypt_private_envelope(
780 &self,
781 requested_key: &str,
782 envelope: LocalDataStoreEnvelope,
783 ) -> Result<Option<StateRecord>, DataStoreError> {
784 if envelope.record.is_some() {
785 return Err(integrity_error("malformed_envelope"));
786 }
787 let record_key = envelope
788 .record_key
789 .ok_or_else(|| integrity_error("malformed_envelope"))?;
790 let key_id = envelope
791 .key_id
792 .ok_or_else(|| integrity_error("malformed_envelope"))?;
793 let nonce_hex = envelope
794 .nonce
795 .ok_or_else(|| integrity_error("malformed_envelope"))?;
796 let ciphertext_hex = envelope
797 .ciphertext
798 .ok_or_else(|| integrity_error("malformed_envelope"))?;
799 if record_key != requested_key {
800 return Err(integrity_error("record_key_mismatch"));
801 }
802 let expected_digest =
803 digest_for_private_envelope(&key_id, &record_key, &nonce_hex, &ciphertext_hex);
804 if envelope.digest != expected_digest {
805 return Err(integrity_error("digest_mismatch"));
806 }
807 let provider = self.required_key_provider()?;
808 let key = provider.key_for(&key_id).map_err(map_key_provider_error)?;
809 let nonce_bytes = hex_decode(&nonce_hex)?;
810 let nonce = Nonce::try_from(nonce_bytes.as_slice())
811 .map_err(|_| integrity_error("invalid_nonce"))?;
812 let ciphertext = hex_decode(&ciphertext_hex)?;
813 let cipher = Aes256Gcm::new_from_slice(key.as_ref()).map_err(|_| crypto_error())?;
814 let aad = private_record_aad(&key_id, &record_key);
815 let plaintext = Zeroizing::new(
816 cipher
817 .decrypt(
818 &nonce,
819 Payload {
820 msg: &ciphertext,
821 aad: &aad,
822 },
823 )
824 .map_err(|_| integrity_error("authentication_failed"))?,
825 );
826 let record: StateRecord = serde_json::from_slice(&plaintext)
827 .map_err(|_| integrity_error("malformed_plaintext"))?;
828 if record.key != record_key {
829 return Err(integrity_error("record_key_mismatch"));
830 }
831 Ok(Some(record))
832 }
833
834 #[cfg(not(feature = "datastore-encryption"))]
835 fn decrypt_private_envelope(
836 &self,
837 _requested_key: &str,
838 _envelope: LocalDataStoreEnvelope,
839 ) -> Result<Option<StateRecord>, DataStoreError> {
840 Err(encryption_feature_disabled())
841 }
842
843 #[cfg(feature = "datastore-encryption")]
844 fn required_key_provider(&self) -> Result<&dyn KeyProvider, DataStoreError> {
845 self.key_provider.as_deref().ok_or_else(|| {
846 data_store_error(
847 DataStoreErrorCode::KeyProviderRequired,
848 "key_provider_required",
849 json!({ "reason": "private_record" }),
850 )
851 })
852 }
853
854 fn ensure_classification_unchanged(&self, path: &PathBuf) -> Result<bool, DataStoreError> {
855 if !path.exists() {
856 return Ok(false);
857 }
858 let bytes = fs::read(path).map_err(|error| io_error("read existing envelope", &error))?;
859 let value: Value =
860 serde_json::from_slice(&bytes).map_err(|_| integrity_error("malformed_envelope"))?;
861 let write_v2 =
862 value.get("format").and_then(Value::as_str) == Some(LOCAL_DATA_STORE_V2_FORMAT);
863 let envelope = if write_v2 {
864 decode_v2_envelope(value)?
865 } else {
866 serde_json::from_value(value).map_err(|_| integrity_error("malformed_envelope"))?
867 };
868 if envelope.format != LOCAL_DATA_STORE_FORMAT {
869 return Err(integrity_error("unknown_format_version"));
870 }
871 if envelope.classification != self.classification {
872 return Err(data_store_error(
873 DataStoreErrorCode::ClassificationChangeNotAllowed,
874 "classification_change_not_allowed",
875 json!({ "reason": "delete_before_reclassifying" }),
876 ));
877 }
878 Ok(write_v2)
879 }
880}
881
882fn digest_bytes(bytes: &[u8]) -> String {
883 format!("sha256:{}", hex_encode(&Sha256::digest(bytes)))
884}
885
886fn digest_for_record(record: &StateRecord) -> Result<String, DataStoreError> {
887 let canonical = serde_json::to_vec(record)
888 .map_err(|error| serialization_error("serialize canonical state record", &error))?;
889 let digest = Sha256::digest(canonical);
890 let mut hexadecimal = String::with_capacity(digest.len() * 2);
891 for byte in digest {
892 hexadecimal.push(char::from(HEXADECIMAL_DIGITS[usize::from(byte >> 4)]));
893 hexadecimal.push(char::from(HEXADECIMAL_DIGITS[usize::from(byte & 0x0f)]));
894 }
895 Ok(format!("sha256:{hexadecimal}"))
896}
897
898fn digest_for_private_envelope(
899 key_id: &str,
900 record_key: &str,
901 nonce: &str,
902 ciphertext: &str,
903) -> String {
904 let mut hasher = Sha256::new();
905 update_length_prefixed(&mut hasher, key_id.as_bytes());
906 update_length_prefixed(&mut hasher, record_key.as_bytes());
907 update_length_prefixed(&mut hasher, b"private");
908 update_length_prefixed(&mut hasher, nonce.as_bytes());
909 update_length_prefixed(&mut hasher, ciphertext.as_bytes());
910 format!("sha256:{}", hex_encode(&hasher.finalize()))
911}
912
913#[cfg(feature = "datastore-encryption")]
914fn private_record_aad(key_id: &str, record_key: &str) -> Vec<u8> {
915 let mut aad = Vec::new();
916 append_length_prefixed(&mut aad, key_id.as_bytes());
917 append_length_prefixed(&mut aad, record_key.as_bytes());
918 append_length_prefixed(&mut aad, b"private");
919 aad
920}
921
922fn update_length_prefixed(hasher: &mut Sha256, value: &[u8]) {
923 hasher.update(value.len().to_le_bytes());
924 hasher.update(value);
925}
926
927#[cfg(feature = "datastore-encryption")]
928fn append_length_prefixed(output: &mut Vec<u8>, value: &[u8]) {
929 output.extend_from_slice(&value.len().to_le_bytes());
930 output.extend_from_slice(value);
931}
932
933fn hex_encode(bytes: &[u8]) -> String {
934 let mut hexadecimal = String::with_capacity(bytes.len() * 2);
935 for byte in bytes {
936 hexadecimal.push(char::from(HEXADECIMAL_DIGITS[usize::from(byte >> 4)]));
937 hexadecimal.push(char::from(HEXADECIMAL_DIGITS[usize::from(byte & 0x0f)]));
938 }
939 hexadecimal
940}
941
942fn hex_decode(value: &str) -> Result<Vec<u8>, DataStoreError> {
943 if !value.len().is_multiple_of(2) {
944 return Err(integrity_error("invalid_hex"));
945 }
946 value
947 .as_bytes()
948 .chunks_exact(2)
949 .map(|pair| {
950 let high = decode_hex_digit(pair[0]).ok_or_else(|| integrity_error("invalid_hex"))?;
951 let low = decode_hex_digit(pair[1]).ok_or_else(|| integrity_error("invalid_hex"))?;
952 Ok((high << 4) | low)
953 })
954 .collect()
955}
956
957fn decode_hex_digit(value: u8) -> Option<u8> {
958 match value {
959 b'0'..=b'9' => Some(value - b'0'),
960 b'a'..=b'f' => Some(value - b'a' + 10),
961 _ => None,
962 }
963}
964
965#[cfg(feature = "datastore-encryption")]
966fn map_key_provider_error(error: KeyProviderError) -> DataStoreError {
967 let code = match error.code {
968 KeyProviderErrorCode::MissingKey => DataStoreErrorCode::KeyNotFound,
969 KeyProviderErrorCode::ExpiredKeyId => DataStoreErrorCode::KeyExpired,
970 KeyProviderErrorCode::ProviderFailure => DataStoreErrorCode::KeyProviderFailure,
971 };
972 let message = error.message;
973 data_store_error(
974 code,
975 &message,
976 json!({ "provider_error_code": key_provider_error_code(error.code) }),
977 )
978}
979
980#[cfg(feature = "datastore-encryption")]
981fn key_provider_error_code(code: KeyProviderErrorCode) -> &'static str {
982 match code {
983 KeyProviderErrorCode::MissingKey => "missing_key",
984 KeyProviderErrorCode::ExpiredKeyId => "expired_key_id",
985 KeyProviderErrorCode::ProviderFailure => "provider_failure",
986 }
987}
988
989#[cfg(feature = "datastore-encryption")]
990fn crypto_error() -> DataStoreError {
991 data_store_error(
992 DataStoreErrorCode::CryptoFailure,
993 "crypto_failed",
994 json!({ "reason": "encryption_failed" }),
995 )
996}
997
998#[cfg(not(feature = "datastore-encryption"))]
999fn encryption_feature_disabled() -> DataStoreError {
1000 data_store_error(
1001 DataStoreErrorCode::KeyProviderRequired,
1002 "key_provider_required",
1003 json!({ "reason": "datastore_encryption_feature_disabled" }),
1004 )
1005}
1006
1007#[cfg(feature = "datastore-encryption")]
1008fn fresh_aes_nonce() -> Nonce<U12> {
1009 let entropy = *uuid::Uuid::new_v4().as_bytes();
1012 let mut nonce_bytes = [0_u8; 12];
1013 nonce_bytes.copy_from_slice(&entropy[..12]);
1014 Nonce::<U12>::from(nonce_bytes)
1015}
1016
1017fn lock_error(error: TryLockError) -> DataStoreError {
1018 match error {
1019 TryLockError::WouldBlock => data_store_error(
1020 DataStoreErrorCode::StoreLocked,
1021 "store_locked",
1022 json!({ "reason": "exclusive_owner_active" }),
1023 ),
1024 TryLockError::Error(error) if error.kind() == std::io::ErrorKind::Unsupported => {
1025 data_store_error(
1026 DataStoreErrorCode::IoFailure,
1027 "storage_io_failed",
1028 json!({ "operation": "acquire_lock", "reason": "locking_unsupported" }),
1029 )
1030 }
1031 TryLockError::Error(_) => data_store_error(
1032 DataStoreErrorCode::IoFailure,
1033 "storage_io_failed",
1034 json!({ "operation": "acquire_lock", "reason": "lock_acquisition_failed" }),
1035 ),
1036 }
1037}
1038
1039fn durability_error(root: &PathBuf, stage: &str, error: &std::io::Error) -> DataStoreError {
1040 data_store_error(
1041 DataStoreErrorCode::DurabilityCommitFailed,
1042 "durability_commit_failed",
1043 json!({ "root": root, "stage": stage, "reason": error.to_string() }),
1044 )
1045}
1046
1047fn discard_temporary_record(path: &PathBuf) {
1048 if path.exists() {
1049 let _ignored = fs::remove_file(path).is_ok();
1050 }
1051}
1052
1053fn integrity_error(reason: &str) -> DataStoreError {
1054 data_store_error(
1055 DataStoreErrorCode::IntegrityCheckFailed,
1056 "integrity_check_failed",
1057 json!({ "reason": reason }),
1058 )
1059}
1060
1061pub fn validate_state_write(
1069 contract: &CapabilityContract,
1070 key: &str,
1071 value: &Value,
1072) -> Result<(), DataStoreError> {
1073 validate_key(key)?;
1074 let schema = contract.state_schema.as_ref().ok_or_else(|| {
1075 data_store_error(
1076 DataStoreErrorCode::NoStateSchemaDeclared,
1077 "no_state_schema_declared",
1078 json!({ "capability_id": contract.id, "key": key }),
1079 )
1080 })?;
1081 let property_schema = schema
1082 .get("properties")
1083 .and_then(Value::as_object)
1084 .and_then(|properties| properties.get(key))
1085 .ok_or_else(|| {
1086 data_store_error(
1087 DataStoreErrorCode::SchemaValidationError,
1088 "schema_validation_error",
1089 json!({ "key": key, "reason": "state key is not declared in schema" }),
1090 )
1091 })?;
1092 let mut violations = Vec::new();
1093 crate::validate_value_against_schema(value, property_schema, "$", &mut violations);
1094 if violations.is_empty() {
1095 Ok(())
1096 } else {
1097 Err(data_store_error(
1098 DataStoreErrorCode::SchemaValidationError,
1099 "schema_validation_error",
1100 json!({ "key": key, "violations": violations }),
1101 ))
1102 }
1103}
1104
1105fn sync_adapters(
1106 local: &mut dyn DataStore,
1107 remote: &mut dyn DataStore,
1108) -> Result<SyncReport, DataStoreError> {
1109 let keys = merged_keys(local.list_keys()?, remote.list_keys()?);
1110 let mut decisions = Vec::new();
1111 let mut snapshots = BTreeMap::new();
1112
1113 for key in keys {
1114 let local_record = local.read(&key)?;
1115 let remote_record = remote.read(&key)?;
1116 snapshots.insert(key.clone(), local_record.clone());
1117 let Some((winner, rule)) = merge_records(local_record.as_ref(), remote_record.as_ref())
1118 else {
1119 continue;
1120 };
1121 apply_winner(local, remote, &key, &winner).map_err(|error| {
1122 rollback_local(local, &snapshots);
1123 data_store_error(
1124 DataStoreErrorCode::SyncFailure,
1125 "sync failed; local state restored",
1126 json!({ "key": key, "cause": error.message }),
1127 )
1128 })?;
1129 decisions.push(MergeDecision {
1130 key,
1131 winning_writer_id: winner.writer_id,
1132 winning_lamport_clock: winner.lamport_clock,
1133 resolution_rule: rule,
1134 });
1135 }
1136
1137 Ok(SyncReport {
1138 governing_spec: DATA_STORE_SPEC.to_string(),
1139 decisions,
1140 })
1141}
1142
1143fn merged_keys(local_keys: Vec<String>, remote_keys: Vec<String>) -> Vec<String> {
1144 local_keys
1145 .into_iter()
1146 .chain(remote_keys)
1147 .collect::<BTreeSet<_>>()
1148 .into_iter()
1149 .collect()
1150}
1151
1152fn merge_records(
1153 local: Option<&StateRecord>,
1154 remote: Option<&StateRecord>,
1155) -> Option<(StateRecord, ConflictResolutionRule)> {
1156 match (local, remote) {
1157 (Some(record), None) => Some((record.clone(), ConflictResolutionRule::OnlyLocal)),
1158 (None, Some(record)) => Some((record.clone(), ConflictResolutionRule::OnlyRemote)),
1159 (Some(local), Some(remote)) => Some(select_conflict_winner(local, remote)),
1160 (None, None) => None,
1161 }
1162}
1163
1164fn select_conflict_winner(
1165 local: &StateRecord,
1166 remote: &StateRecord,
1167) -> (StateRecord, ConflictResolutionRule) {
1168 if local.lamport_clock > remote.lamport_clock {
1169 return (local.clone(), ConflictResolutionRule::HigherLamportClock);
1170 }
1171 if remote.lamport_clock > local.lamport_clock {
1172 return (remote.clone(), ConflictResolutionRule::HigherLamportClock);
1173 }
1174 if local.writer_id >= remote.writer_id {
1175 (
1176 local.clone(),
1177 ConflictResolutionRule::WriterIdentityTieBreak,
1178 )
1179 } else {
1180 (
1181 remote.clone(),
1182 ConflictResolutionRule::WriterIdentityTieBreak,
1183 )
1184 }
1185}
1186
1187fn apply_winner(
1188 local: &mut dyn DataStore,
1189 remote: &mut dyn DataStore,
1190 key: &str,
1191 winner: &StateRecord,
1192) -> Result<(), DataStoreError> {
1193 if local.read(key)?.as_ref() != Some(winner) {
1194 local.write(winner.clone())?;
1195 }
1196 if remote.read(key)?.as_ref() != Some(winner) {
1197 remote.write(winner.clone())?;
1198 }
1199 Ok(())
1200}
1201
1202fn rollback_local(local: &mut dyn DataStore, snapshots: &BTreeMap<String, Option<StateRecord>>) {
1203 for (key, snapshot) in snapshots {
1204 let result = match snapshot {
1205 Some(record) => local.write(record.clone()),
1206 None => local.delete(key),
1207 };
1208 let _ignored = result.is_ok();
1209 }
1210}
1211
1212fn validate_key(key: &str) -> Result<(), DataStoreError> {
1213 let valid = !key.is_empty()
1214 && key
1215 .chars()
1216 .all(|character| character.is_ascii_alphanumeric() || matches!(character, '_' | '-'));
1217 if valid {
1218 Ok(())
1219 } else {
1220 Err(data_store_error(
1221 DataStoreErrorCode::InvalidKey,
1222 "state key must be non-empty and contain only ASCII letters, numbers, '_' or '-'",
1223 json!({ "key": key }),
1224 ))
1225 }
1226}
1227
1228fn data_store_error(code: DataStoreErrorCode, message: &str, details: Value) -> DataStoreError {
1229 DataStoreError {
1230 code,
1231 message: message.to_string(),
1232 details,
1233 }
1234}
1235
1236fn io_error(action: &str, error: &std::io::Error) -> DataStoreError {
1237 data_store_error(
1238 DataStoreErrorCode::IoFailure,
1239 "storage_io_failed",
1240 json!({ "action": action, "reason": error.to_string() }),
1241 )
1242}
1243
1244fn serialization_error(action: &str, error: &serde_json::Error) -> DataStoreError {
1245 data_store_error(
1246 DataStoreErrorCode::SerializationFailure,
1247 action,
1248 json!({ "error": error.to_string() }),
1249 )
1250}
1251
1252#[cfg(test)]
1253#[allow(clippy::expect_used)]
1254mod tests {
1255 use super::*;
1256 use serde_json::json;
1257 use std::cell::Cell;
1258 use std::path::Path;
1259 use std::process::{Child, Command, Stdio};
1260 use std::thread;
1261 use std::time::Duration;
1262 use traverse_contracts::{
1263 BinaryFormat, CapabilityContract, Condition, DependencyReference, Entrypoint,
1264 EntrypointKind, EventReference, Execution, ExecutionConstraints, ExecutionTarget,
1265 FilesystemAccess, HostApiAccess, IdReference, Lifecycle, NetworkAccess, Owner, Provenance,
1266 ProvenanceSource, SchemaContainer, ServiceType, SideEffect, SideEffectKind,
1267 ValidationEvidence,
1268 };
1269 use uuid::Uuid;
1270
1271 #[derive(Debug, Clone, Default)]
1272 struct MemoryDataStore {
1273 records: BTreeMap<String, StateRecord>,
1274 fail_writes: Cell<bool>,
1275 }
1276
1277 #[derive(Debug, Clone, Default)]
1278 struct PhantomKeyStore;
1279
1280 struct FailingKeyProvider;
1281
1282 impl KeyProvider for FailingKeyProvider {
1283 fn active_key_id(&self) -> Result<String, KeyProviderError> {
1284 Err(KeyProviderError::provider_failure())
1285 }
1286
1287 fn key_for(&self, _key_id: &str) -> Result<Zeroizing<[u8; 32]>, KeyProviderError> {
1288 Err(KeyProviderError::provider_failure())
1289 }
1290 }
1291
1292 impl DataStore for MemoryDataStore {
1293 fn read(&self, key: &str) -> Result<Option<StateRecord>, DataStoreError> {
1294 Ok(self.records.get(key).cloned())
1295 }
1296
1297 fn write(&mut self, record: StateRecord) -> Result<(), DataStoreError> {
1298 if self.fail_writes.get() {
1299 return Err(data_store_error(
1300 DataStoreErrorCode::IoFailure,
1301 "forced write failure",
1302 json!({ "key": record.key }),
1303 ));
1304 }
1305 self.records.insert(record.key.clone(), record);
1306 Ok(())
1307 }
1308
1309 fn delete(&mut self, key: &str) -> Result<(), DataStoreError> {
1310 self.records.remove(key);
1311 Ok(())
1312 }
1313
1314 fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
1315 Ok(self.records.keys().cloned().collect())
1316 }
1317 }
1318
1319 impl DataStore for PhantomKeyStore {
1320 fn read(&self, _key: &str) -> Result<Option<StateRecord>, DataStoreError> {
1321 Ok(None)
1322 }
1323
1324 fn write(&mut self, _record: StateRecord) -> Result<(), DataStoreError> {
1325 Ok(())
1326 }
1327
1328 fn delete(&mut self, _key: &str) -> Result<(), DataStoreError> {
1329 Ok(())
1330 }
1331
1332 fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
1333 Ok(vec!["phantom".to_string()])
1334 }
1335 }
1336
1337 #[test]
1338 fn runtime_data_store_validates_writes_and_reads_from_local_file_adapter() {
1339 let root = temp_root("valid");
1340 let adapter = public_store(&root);
1341 let mut store = RuntimeDataStore::new(adapter, "writer-a");
1342 let contract = stateful_contract(Some(json!({
1343 "type": "object",
1344 "properties": {
1345 "draft": {"type": "string"}
1346 }
1347 })));
1348
1349 let record = store
1350 .write(&contract, "draft", json!("ready"))
1351 .expect("valid state write should succeed");
1352
1353 assert_eq!(record.lamport_clock, 1);
1354 assert_eq!(
1355 store.read(&contract, "draft").expect("read should succeed"),
1356 Some(json!("ready"))
1357 );
1358 assert_eq!(
1359 store.list_keys().expect("list should succeed"),
1360 vec!["draft".to_string()]
1361 );
1362 store.delete("draft").expect("delete should succeed");
1363 assert_eq!(
1364 store.read(&contract, "draft").expect("read should succeed"),
1365 None
1366 );
1367 }
1368
1369 #[test]
1370 fn runtime_data_store_rejects_missing_schema_bad_keys_and_schema_violations() {
1371 let adapter = MemoryDataStore::default();
1372 let mut store = RuntimeDataStore::new(adapter, "writer-a");
1373 let no_schema = stateful_contract(None);
1374 let schema = stateful_contract(Some(json!({
1375 "type": "object",
1376 "properties": {
1377 "count": {"type": "integer"}
1378 }
1379 })));
1380
1381 let missing = store
1382 .write(&no_schema, "count", json!(1))
1383 .expect_err("missing state schema should fail");
1384 assert_eq!(missing.code, DataStoreErrorCode::NoStateSchemaDeclared);
1385
1386 let invalid_key = store
1387 .write(&schema, "bad.key", json!(1))
1388 .expect_err("invalid key should fail");
1389 assert_eq!(invalid_key.code, DataStoreErrorCode::InvalidKey);
1390
1391 let undeclared = store
1392 .write(&schema, "other", json!(1))
1393 .expect_err("undeclared state key should fail");
1394 assert_eq!(undeclared.code, DataStoreErrorCode::SchemaValidationError);
1395
1396 let wrong_type = store
1397 .write(&schema, "count", json!("one"))
1398 .expect_err("wrong state type should fail");
1399 assert_eq!(wrong_type.code, DataStoreErrorCode::SchemaValidationError);
1400
1401 let no_schema_read = store
1402 .read(&no_schema, "count")
1403 .expect("no-schema read should succeed");
1404 assert_eq!(no_schema_read, None);
1405
1406 let bad_read_key = store
1407 .read(&schema, "bad.key")
1408 .expect_err("invalid read key should fail");
1409 assert_eq!(bad_read_key.code, DataStoreErrorCode::InvalidKey);
1410 }
1411
1412 #[test]
1413 fn lamport_clock_overflow_is_rejected_before_adapter_write() {
1414 let adapter = MemoryDataStore::default();
1415 let clock = LamportClock::with_value("writer-a", u64::MAX);
1416 let mut store = RuntimeDataStore::with_clock(adapter, clock);
1417 let contract = stateful_contract(Some(json!({
1418 "type": "object",
1419 "properties": {
1420 "draft": {"type": "string"}
1421 }
1422 })));
1423
1424 let error = store
1425 .write(&contract, "draft", json!("ready"))
1426 .expect_err("overflow should fail");
1427
1428 assert_eq!(error.code, DataStoreErrorCode::LamportClockOverflow);
1429 assert!(store.into_inner().records.is_empty());
1430 }
1431
1432 #[test]
1433 fn runtime_data_store_validates_reads_before_returning_stored_values() {
1434 let mut adapter = MemoryDataStore::default();
1435 adapter
1436 .write(record("count", "writer-a", 1, json!("not an integer")))
1437 .expect("seed should succeed");
1438 let store = RuntimeDataStore::new(adapter, "writer-a");
1439 let contract = stateful_contract(Some(json!({
1440 "type": "object",
1441 "properties": {
1442 "count": {"type": "integer"}
1443 }
1444 })));
1445
1446 let error = store
1447 .read(&contract, "count")
1448 .expect_err("invalid stored value should fail");
1449
1450 assert_eq!(error.code, DataStoreErrorCode::SchemaValidationError);
1451 }
1452
1453 #[test]
1454 fn reconnect_sync_merges_only_local_only_remote_clock_winner_and_writer_tie_breaks() {
1455 let mut local = MemoryDataStore::default();
1456 let mut remote = MemoryDataStore::default();
1457 local
1458 .write(record("local_only", "local-a", 1, json!("local")))
1459 .expect("local write should succeed");
1460 remote
1461 .write(record("remote_only", "remote-a", 1, json!("remote")))
1462 .expect("remote write should succeed");
1463 local
1464 .write(record("clock", "local-a", 2, json!("old")))
1465 .expect("local write should succeed");
1466 remote
1467 .write(record("clock", "remote-a", 3, json!("new")))
1468 .expect("remote write should succeed");
1469 local
1470 .write(record("tie", "writer-z", 4, json!("winner")))
1471 .expect("local write should succeed");
1472 remote
1473 .write(record("tie", "writer-a", 4, json!("loser")))
1474 .expect("remote write should succeed");
1475
1476 let report = sync_adapters(&mut local, &mut remote).expect("sync should succeed");
1477
1478 assert_eq!(report.governing_spec, "089-datastore-synchronization");
1479 assert_eq!(report.decisions.len(), 4);
1480 assert_eq!(
1481 local.read("remote_only").expect("read should succeed"),
1482 remote.read("remote_only").expect("read should succeed")
1483 );
1484 assert_eq!(
1485 local.read("clock").expect("read should succeed"),
1486 Some(record("clock", "remote-a", 3, json!("new")))
1487 );
1488 assert_eq!(
1489 remote.read("tie").expect("read should succeed"),
1490 Some(record("tie", "writer-z", 4, json!("winner")))
1491 );
1492 assert!(
1493 report
1494 .decisions
1495 .iter()
1496 .any(|decision| decision.resolution_rule
1497 == ConflictResolutionRule::WriterIdentityTieBreak)
1498 );
1499 }
1500
1501 #[test]
1502 fn sync_failure_restores_local_snapshot() {
1503 let mut local = MemoryDataStore::default();
1504 let mut remote = MemoryDataStore::default();
1505 local
1506 .write(record("shared", "local-a", 2, json!("local")))
1507 .expect("local write should succeed");
1508 remote
1509 .write(record("shared", "remote-a", 1, json!("remote")))
1510 .expect("remote write should succeed");
1511 remote.fail_writes.set(true);
1512
1513 let error = sync_adapters(&mut local, &mut remote).expect_err("sync should fail");
1514
1515 assert_eq!(error.code, DataStoreErrorCode::SyncFailure);
1516 assert_eq!(
1517 local.read("shared").expect("read should succeed"),
1518 Some(record("shared", "local-a", 2, json!("local")))
1519 );
1520 }
1521
1522 #[test]
1523 fn local_file_adapter_reports_bad_keys_and_bad_json() {
1524 let root = temp_root("bad-json");
1525 let adapter = LocalFileDataStore::new(&root).expect("local adapter should initialize");
1526 let invalid = adapter
1527 .read("bad.key")
1528 .expect_err("invalid key should fail");
1529 assert_eq!(invalid.code, DataStoreErrorCode::InvalidKey);
1530
1531 fs::write(root.join("broken.json"), "{").expect("bad json fixture should write");
1532 let invalid_json = adapter
1533 .read("broken")
1534 .expect_err("invalid json should fail");
1535 assert_eq!(invalid_json.code, DataStoreErrorCode::IntegrityCheckFailed);
1536 assert_eq!(invalid_json.message, "integrity_check_failed");
1537 }
1538
1539 #[test]
1540 fn helper_paths_cover_remaining_datastore_branches() {
1541 let mut local = RuntimeDataStore::new(MemoryDataStore::default(), "local-a");
1542 let mut remote = MemoryDataStore::default();
1543 remote
1544 .write(record("remote_only", "remote-a", 1, json!("remote")))
1545 .expect("remote seed should succeed");
1546
1547 let report = local
1548 .sync_on_reconnect(&mut remote)
1549 .expect("public reconnect sync should succeed");
1550 assert_eq!(report.decisions.len(), 1);
1551
1552 assert!(merge_records(None, None).is_none());
1553 let (_winner, rule) = select_conflict_winner(
1554 &record("tie", "writer-a", 1, json!("local")),
1555 &record("tie", "writer-z", 1, json!("remote")),
1556 );
1557 assert_eq!(rule, ConflictResolutionRule::WriterIdentityTieBreak);
1558
1559 let mut failing_local = MemoryDataStore::default();
1560 failing_local.fail_writes.set(true);
1561 let mut seeded_remote = MemoryDataStore::default();
1562 seeded_remote
1563 .write(record("missing_local", "remote-a", 1, json!("remote")))
1564 .expect("remote seed should succeed");
1565 let error =
1566 sync_adapters(&mut failing_local, &mut seeded_remote).expect_err("sync should fail");
1567 assert_eq!(error.code, DataStoreErrorCode::SyncFailure);
1568 assert_eq!(
1569 failing_local
1570 .delete("missing_local")
1571 .expect("delete should succeed"),
1572 ()
1573 );
1574
1575 let mut phantom_local = PhantomKeyStore;
1576 let mut phantom_remote = PhantomKeyStore;
1577 assert!(
1578 sync_adapters(&mut phantom_local, &mut phantom_remote)
1579 .expect("phantom sync should succeed")
1580 .decisions
1581 .is_empty()
1582 );
1583 phantom_local
1584 .write(record("phantom", "writer-a", 1, json!("value")))
1585 .expect("phantom write should succeed");
1586 phantom_local
1587 .delete("phantom")
1588 .expect("phantom delete should succeed");
1589
1590 let root = temp_root("listing");
1591 fs::create_dir_all(&root).expect("root should be created");
1592 fs::write(root.join("skip.txt"), "not state").expect("non-json fixture should write");
1593 let adapter = LocalFileDataStore::new(&root).expect("local adapter should initialize");
1594 assert!(adapter.list_keys().expect("list should succeed").is_empty());
1595 drop(adapter);
1596 let mut delete_missing =
1597 LocalFileDataStore::new(&root).expect("local adapter should initialize");
1598 delete_missing
1599 .delete("missing")
1600 .expect("missing delete should succeed");
1601 fs::create_dir(root.join("cant_delete.json")).expect("directory fixture should write");
1602 let delete_failure = delete_missing
1603 .delete("cant_delete")
1604 .expect_err("directory delete should fail");
1605 assert_eq!(delete_failure.code, DataStoreErrorCode::IoFailure);
1606
1607 let file_root = temp_root("file-root");
1608 fs::write(&file_root, "not a directory").expect("file root fixture should write");
1609 let io_failure = LocalFileDataStore::new(&file_root).expect_err("file root should fail");
1610 assert_eq!(io_failure.code, DataStoreErrorCode::IoFailure);
1611 }
1612
1613 #[test]
1614 fn local_file_adapter_writes_integrity_envelope_and_reopens() {
1615 let root = temp_root("integrity-envelope");
1616 let record = record("draft", "writer-a", 1, json!("ready"));
1617 let mut adapter =
1618 LocalFileDataStore::with_classification(&root, LocalDataClassification::Public)
1619 .expect("local adapter should initialize");
1620 adapter.write(record.clone()).expect("write should succeed");
1621
1622 let envelope: LocalDataStoreEnvelope = serde_json::from_slice(
1623 &fs::read(root.join("draft.json")).expect("envelope should be present"),
1624 )
1625 .expect("envelope should deserialize");
1626 assert_eq!(envelope.format, LOCAL_DATA_STORE_FORMAT);
1627 assert_eq!(envelope.classification, LocalDataClassification::Public);
1628 assert_eq!(
1629 envelope.digest,
1630 digest_for_record(&record).expect("digest should compute")
1631 );
1632
1633 drop(adapter);
1634 let reopened = LocalFileDataStore::new(&root).expect("reopen should acquire lock");
1635 assert_eq!(
1636 reopened.read("draft").expect("read should succeed"),
1637 Some(record)
1638 );
1639 }
1640
1641 #[test]
1642 fn private_records_encrypt_reopen_and_require_provider() {
1643 let root = temp_root("private-encryption");
1644 let provider: Arc<dyn KeyProvider> = Arc::new(InMemoryKeyProvider::new("key-1", [7; 32]));
1645 let original = record("secret", "writer-a", 1, json!("classified value"));
1646 let mut store = LocalFileDataStore::new(&root)
1647 .expect("private store should open")
1648 .with_key_provider(Arc::clone(&provider));
1649 store.write(original.clone()).expect("private write");
1650
1651 let bytes = fs::read(root.join("secret.json")).expect("private envelope");
1652 let text = String::from_utf8(bytes.clone()).expect("json utf8");
1653 let envelope: LocalDataStoreEnvelope =
1654 serde_json::from_slice(&bytes).expect("private envelope shape");
1655 assert_eq!(envelope.classification, LocalDataClassification::Private);
1656 assert!(envelope.record.is_none());
1657 assert_eq!(envelope.record_key.as_deref(), Some("secret"));
1658 assert_eq!(envelope.key_id.as_deref(), Some("key-1"));
1659 assert!(envelope.nonce.is_some());
1660 assert!(envelope.ciphertext.is_some());
1661 assert!(!text.contains("classified value"));
1662 assert!(!text.contains(&hex_encode(&[7; 32])));
1663
1664 let first_nonce = envelope.nonce;
1665 store.write(original.clone()).expect("second private write");
1666 let second: LocalDataStoreEnvelope =
1667 serde_json::from_slice(&fs::read(root.join("secret.json")).expect("second envelope"))
1668 .expect("second envelope shape");
1669 assert_ne!(first_nonce, second.nonce);
1670 drop(store);
1671
1672 let reopened = LocalFileDataStore::new(&root)
1673 .expect("private store should reopen")
1674 .with_key_provider(Arc::clone(&provider));
1675 assert_eq!(
1676 reopened.read("secret").expect("private read"),
1677 Some(original)
1678 );
1679 drop(reopened);
1680
1681 let no_provider = LocalFileDataStore::new(&root).expect("store without provider opens");
1682 let read_error = no_provider
1683 .read("secret")
1684 .expect_err("private read must require provider");
1685 assert_eq!(read_error.code, DataStoreErrorCode::KeyProviderRequired);
1686 }
1687
1688 #[test]
1689 fn private_write_without_provider_fails_before_commit_and_public_crud_works() {
1690 let private_root = temp_root("private-provider-required");
1691 let mut private =
1692 LocalFileDataStore::new(&private_root).expect("private store should open");
1693 let error = private
1694 .write(record("secret", "writer-a", 1, json!("value")))
1695 .expect_err("private write must fail");
1696 assert_eq!(error.code, DataStoreErrorCode::KeyProviderRequired);
1697 assert!(!private_root.join("secret.json").exists());
1698
1699 let public_root = temp_root("public-without-provider");
1700 let mut public = public_store(&public_root);
1701 let public_record = record("cache", "writer-a", 1, json!("visible"));
1702 public.write(public_record.clone()).expect("public write");
1703 assert_eq!(
1704 public.read("cache").expect("public read"),
1705 Some(public_record)
1706 );
1707 public.delete("cache").expect("public delete");
1708 assert!(public.list_keys().expect("public list").is_empty());
1709 }
1710
1711 #[test]
1712 fn private_authentication_and_key_provider_failures_are_stable() {
1713 let root = temp_root("private-authentication");
1714 let key = [9; 32];
1715 let provider: Arc<dyn KeyProvider> =
1716 Arc::new(InMemoryKeyProvider::new("key-a", key).with_read_key("key-b", key));
1717 let mut store = LocalFileDataStore::new(&root)
1718 .expect("private store should open")
1719 .with_key_provider(Arc::clone(&provider));
1720 store
1721 .write(record("secret", "writer-a", 1, json!("value")))
1722 .expect("private write");
1723
1724 let path = root.join("secret.json");
1725 let mut envelope: LocalDataStoreEnvelope =
1726 serde_json::from_slice(&fs::read(&path).expect("envelope")).expect("shape");
1727 let original_envelope = envelope.clone();
1728 let mut ciphertext =
1729 hex_decode(envelope.ciphertext.as_deref().expect("ciphertext")).expect("valid hex");
1730 ciphertext[0] ^= 1;
1731 envelope.ciphertext = Some(hex_encode(&ciphertext));
1732 envelope.digest = digest_for_private_envelope(
1733 envelope.key_id.as_deref().expect("key id"),
1734 envelope.record_key.as_deref().expect("record key"),
1735 envelope.nonce.as_deref().expect("nonce"),
1736 envelope.ciphertext.as_deref().expect("ciphertext"),
1737 );
1738 fs::write(&path, serde_json::to_vec(&envelope).expect("serialize")).expect("tamper");
1739 let ciphertext_authentication = store
1740 .read("secret")
1741 .expect_err("ciphertext tampering must fail");
1742 assert_eq!(
1743 ciphertext_authentication.code,
1744 DataStoreErrorCode::IntegrityCheckFailed
1745 );
1746 assert_eq!(
1747 ciphertext_authentication.details["reason"],
1748 "authentication_failed"
1749 );
1750
1751 envelope = original_envelope;
1752 envelope.key_id = Some("key-b".to_string());
1753 envelope.digest = digest_for_private_envelope(
1754 "key-b",
1755 envelope.record_key.as_deref().expect("record key"),
1756 envelope.nonce.as_deref().expect("nonce"),
1757 envelope.ciphertext.as_deref().expect("ciphertext"),
1758 );
1759 fs::write(&path, serde_json::to_vec(&envelope).expect("serialize")).expect("tamper");
1760 let authentication = store.read("secret").expect_err("AAD tampering must fail");
1761 assert_eq!(
1762 authentication.code,
1763 DataStoreErrorCode::IntegrityCheckFailed
1764 );
1765 assert_eq!(authentication.details["reason"], "authentication_failed");
1766 drop(store);
1767
1768 let missing_provider: Arc<dyn KeyProvider> =
1769 Arc::new(InMemoryKeyProvider::new("other", [3; 32]));
1770 let missing = LocalFileDataStore::new(&root)
1771 .expect("reopen")
1772 .with_key_provider(missing_provider)
1773 .read("secret")
1774 .expect_err("missing key must fail");
1775 assert_eq!(missing.code, DataStoreErrorCode::KeyNotFound);
1776 assert_eq!(missing.message, "key_not_found");
1777
1778 let expired_provider: Arc<dyn KeyProvider> =
1779 Arc::new(InMemoryKeyProvider::new("key-b", key).with_expired_key_id("key-b"));
1780 let expired = LocalFileDataStore::new(&root)
1781 .expect("reopen")
1782 .with_key_provider(expired_provider)
1783 .read("secret")
1784 .expect_err("expired key must fail");
1785 assert_eq!(expired.code, DataStoreErrorCode::KeyExpired);
1786
1787 let failing_provider: Arc<dyn KeyProvider> = Arc::new(FailingKeyProvider);
1788 let failure_root = temp_root("provider-failure");
1789 let provider_failure = LocalFileDataStore::new(&failure_root)
1790 .expect("open")
1791 .with_key_provider(failing_provider)
1792 .write(record("secret", "writer-a", 1, json!("value")))
1793 .expect_err("provider failure must fail");
1794 assert_eq!(
1795 provider_failure.code,
1796 DataStoreErrorCode::KeyProviderFailure
1797 );
1798 assert_eq!(provider_failure.message, "key_provider_failed");
1799 }
1800
1801 #[test]
1802 fn classification_change_requires_delete_before_write() {
1803 let root = temp_root("classification-immutable");
1804 let mut public = public_store(&root);
1805 public
1806 .write(record("shared", "writer-a", 1, json!("public")))
1807 .expect("public write");
1808 drop(public);
1809
1810 let provider: Arc<dyn KeyProvider> = Arc::new(InMemoryKeyProvider::new("key-1", [4; 32]));
1811 let mut private = LocalFileDataStore::new(&root)
1812 .expect("private reopen")
1813 .with_key_provider(provider);
1814 let rejected = private
1815 .write(record("shared", "writer-a", 2, json!("private")))
1816 .expect_err("in-place reclassification must fail");
1817 assert_eq!(
1818 rejected.code,
1819 DataStoreErrorCode::ClassificationChangeNotAllowed
1820 );
1821 private.delete("shared").expect("delete old classification");
1822 private
1823 .write(record("shared", "writer-a", 2, json!("private")))
1824 .expect("write after delete");
1825 }
1826
1827 #[test]
1828 #[allow(clippy::too_many_lines)]
1829 fn private_envelope_corruption_and_error_helpers_fail_closed() {
1830 let root = temp_root("private-corruption-branches");
1831 let key = [23; 32];
1832 let provider: Arc<dyn KeyProvider> = Arc::new(InMemoryKeyProvider::new("key-1", key));
1833 let mut store = LocalFileDataStore::new(&root)
1834 .expect("open")
1835 .with_key_provider(Arc::clone(&provider));
1836 store
1837 .write(record("secret", "writer-a", 1, json!("value")))
1838 .expect("write");
1839 let path = root.join("secret.json");
1840 let original: LocalDataStoreEnvelope =
1841 serde_json::from_slice(&fs::read(&path).expect("read")).expect("shape");
1842
1843 let mut plaintext_leak = original.clone();
1844 plaintext_leak.record = Some(record("secret", "writer-a", 1, json!("value")));
1845 write_envelope_fixture(&path, &plaintext_leak);
1846 assert_integrity_reason(&store, "malformed_envelope");
1847
1848 let mut wrong_key = original.clone();
1849 wrong_key.record_key = Some("other".to_string());
1850 wrong_key.digest = digest_for_private_envelope(
1851 wrong_key.key_id.as_deref().expect("key id"),
1852 "other",
1853 wrong_key.nonce.as_deref().expect("nonce"),
1854 wrong_key.ciphertext.as_deref().expect("ciphertext"),
1855 );
1856 write_envelope_fixture(&path, &wrong_key);
1857 assert_integrity_reason(&store, "record_key_mismatch");
1858
1859 let mut bad_digest = original.clone();
1860 bad_digest.digest = "sha256:deadbeef".to_string();
1861 write_envelope_fixture(&path, &bad_digest);
1862 assert_integrity_reason(&store, "digest_mismatch");
1863
1864 let mut odd_nonce = original.clone();
1865 odd_nonce.nonce = Some("0".to_string());
1866 odd_nonce.digest = digest_for_private_envelope(
1867 odd_nonce.key_id.as_deref().expect("key id"),
1868 odd_nonce.record_key.as_deref().expect("record key"),
1869 "0",
1870 odd_nonce.ciphertext.as_deref().expect("ciphertext"),
1871 );
1872 write_envelope_fixture(&path, &odd_nonce);
1873 assert_integrity_reason(&store, "invalid_hex");
1874
1875 let mut invalid_nonce = original.clone();
1876 invalid_nonce.nonce = Some("00".to_string());
1877 invalid_nonce.digest = digest_for_private_envelope(
1878 invalid_nonce.key_id.as_deref().expect("key id"),
1879 invalid_nonce.record_key.as_deref().expect("record key"),
1880 "00",
1881 invalid_nonce.ciphertext.as_deref().expect("ciphertext"),
1882 );
1883 write_envelope_fixture(&path, &invalid_nonce);
1884 assert_integrity_reason(&store, "invalid_nonce");
1885
1886 let mut invalid_ciphertext = original.clone();
1887 invalid_ciphertext.ciphertext = Some("gg".to_string());
1888 invalid_ciphertext.digest = digest_for_private_envelope(
1889 invalid_ciphertext.key_id.as_deref().expect("key id"),
1890 invalid_ciphertext
1891 .record_key
1892 .as_deref()
1893 .expect("record key"),
1894 invalid_ciphertext.nonce.as_deref().expect("nonce"),
1895 "gg",
1896 );
1897 write_envelope_fixture(&path, &invalid_ciphertext);
1898 assert_integrity_reason(&store, "invalid_hex");
1899
1900 let malformed_plaintext = encrypted_fixture("secret", "key-1", key, b"not-json".as_slice());
1901 write_envelope_fixture(&path, &malformed_plaintext);
1902 assert_integrity_reason(&store, "malformed_plaintext");
1903
1904 let mismatched_record =
1905 serde_json::to_vec(&record("other", "writer-a", 1, json!("value"))).expect("serialize");
1906 let mismatched_plaintext = encrypted_fixture("secret", "key-1", key, &mismatched_record);
1907 write_envelope_fixture(&path, &mismatched_plaintext);
1908 assert_integrity_reason(&store, "record_key_mismatch");
1909
1910 let mut unknown_format = original;
1911 unknown_format.format = "local-datastore/unknown".to_string();
1912 write_envelope_fixture(&path, &unknown_format);
1913 let unknown = store
1914 .write(record("secret", "writer-a", 2, json!("new")))
1915 .expect_err("unknown existing format");
1916 assert_eq!(unknown.details["reason"], "unknown_format_version");
1917
1918 assert_eq!(crypto_error().code, DataStoreErrorCode::CryptoFailure);
1919 assert_eq!(
1920 FailingKeyProvider
1921 .key_for("key-1")
1922 .expect_err("provider failure")
1923 .code,
1924 KeyProviderErrorCode::ProviderFailure
1925 );
1926 assert_eq!(KeyProviderError::missing_key().message, "key_not_found");
1927 for code in [
1928 KeyProviderErrorCode::MissingKey,
1929 KeyProviderErrorCode::ExpiredKeyId,
1930 KeyProviderErrorCode::ProviderFailure,
1931 ] {
1932 let encoded = serde_json::to_vec(&code).expect("serialize provider code");
1933 let decoded: KeyProviderErrorCode =
1934 serde_json::from_slice(&encoded).expect("deserialize provider code");
1935 assert_eq!(decoded, code);
1936 }
1937
1938 let mut missing_active = InMemoryKeyProvider::new("missing", [1; 32]);
1939 missing_active.keys.clear();
1940 assert_eq!(
1941 missing_active
1942 .active_key_id()
1943 .expect_err("missing active key")
1944 .code,
1945 KeyProviderErrorCode::MissingKey
1946 );
1947 let expired_active =
1948 InMemoryKeyProvider::new("expired", [1; 32]).with_expired_key_id("expired");
1949 assert_eq!(
1950 expired_active
1951 .active_key_id()
1952 .expect_err("expired active key")
1953 .code,
1954 KeyProviderErrorCode::ExpiredKeyId
1955 );
1956 assert!(format!("{store:?}").contains("key_provider_configured"));
1957
1958 let public_root = temp_root("public-extra-encryption-fields");
1959 let mut public = public_store(&public_root);
1960 public
1961 .write(record("cache", "writer-a", 1, json!("value")))
1962 .expect("public write");
1963 let public_path = public_root.join("cache.json");
1964 let mut public_envelope: LocalDataStoreEnvelope =
1965 serde_json::from_slice(&fs::read(&public_path).expect("read")).expect("shape");
1966 public_envelope.key_id = Some("unexpected".to_string());
1967 write_envelope_fixture(&public_path, &public_envelope);
1968 let malformed_public = public.read("cache").expect_err("extra private metadata");
1969 assert_eq!(malformed_public.details["reason"], "malformed_envelope");
1970 }
1971
1972 #[test]
1973 fn local_file_adapter_rejects_tampered_and_legacy_records() {
1974 let root = temp_root("tampered");
1975 let mut adapter = public_store(&root);
1976 adapter
1977 .write(record("draft", "writer-a", 1, json!("ready")))
1978 .expect("write should succeed");
1979
1980 let mut envelope: Value = serde_json::from_slice(
1981 &fs::read(root.join("draft.json")).expect("envelope should be present"),
1982 )
1983 .expect("fixture should deserialize");
1984 envelope["record"]["value"] = json!("tampered");
1985 fs::write(
1986 root.join("draft.json"),
1987 serde_json::to_vec(&envelope).expect("fixture should serialize"),
1988 )
1989 .expect("tampered fixture should write");
1990 let tampered = adapter.read("draft").expect_err("tampering must fail");
1991 assert_eq!(tampered.code, DataStoreErrorCode::IntegrityCheckFailed);
1992 assert_eq!(tampered.details["reason"], "digest_mismatch");
1993
1994 fs::write(
1995 root.join("legacy.json"),
1996 serde_json::to_vec(&record("legacy", "writer-a", 1, json!("old")))
1997 .expect("legacy fixture should serialize"),
1998 )
1999 .expect("legacy fixture should write");
2000 let legacy = adapter.read("legacy").expect_err("legacy must fail closed");
2001 assert_eq!(legacy.code, DataStoreErrorCode::IntegrityCheckFailed);
2002 assert_eq!(legacy.details["reason"], "legacy_unverified");
2003 }
2004
2005 #[test]
2006 fn local_file_adapter_ignores_temporary_records_and_rejects_second_owner() {
2007 let root = temp_root("temporary-and-lock");
2008 let mut adapter = public_store(&root);
2009 adapter
2010 .write(record("draft", "writer-a", 1, json!("committed")))
2011 .expect("write should succeed");
2012 fs::write(root.join(".draft.temporary.tmp"), "incomplete")
2013 .expect("temporary fixture should write");
2014
2015 assert_eq!(
2016 adapter.list_keys().expect("listing should succeed"),
2017 vec!["draft".to_string()]
2018 );
2019 assert_eq!(
2020 adapter
2021 .read("draft")
2022 .expect("committed read should succeed"),
2023 Some(record("draft", "writer-a", 1, json!("committed")))
2024 );
2025 let second_owner = LocalFileDataStore::new(&root).expect_err("second owner must fail");
2026 assert_eq!(second_owner.code, DataStoreErrorCode::StoreLocked);
2027 assert_eq!(second_owner.message, "store_locked");
2028 assert_eq!(
2029 second_owner.details,
2030 json!({ "reason": "exclusive_owner_active" })
2031 );
2032 }
2033
2034 #[test]
2035 fn local_file_adapter_lock_child() -> Result<(), String> {
2036 let Ok(root) = std::env::var("TRAVERSE_DATA_STORE_LOCK_CHILD_ROOT") else {
2037 return Ok(());
2038 };
2039 let root = PathBuf::from(root);
2040 let _adapter = LocalFileDataStore::new(&root).expect("child should acquire lock");
2041 fs::write(lock_child_ready_path(&root), "ready").expect("child should signal readiness");
2042 wait_for_lock_child_release(&root, 500)
2043 }
2044
2045 #[test]
2046 fn local_file_adapter_rejects_cross_process_owner_and_recovers_after_exit() {
2047 let root = temp_root("cross-process-lock");
2048 let mut initial_owner = public_store(&root);
2049 let committed = record("draft", "writer-a", 1, json!("committed"));
2050 initial_owner
2051 .write(committed.clone())
2052 .expect("initial write should succeed");
2053 drop(initial_owner);
2054
2055 let mut child = start_lock_child(&root);
2056 wait_for_lock_child(&root, 500).expect("lock child should become ready");
2057 let blocked = LocalFileDataStore::new(&root).expect_err("second process must be blocked");
2058 assert_eq!(blocked.code, DataStoreErrorCode::StoreLocked);
2059 assert_eq!(
2060 blocked.details,
2061 json!({ "reason": "exclusive_owner_active" })
2062 );
2063
2064 fs::write(lock_child_release_path(&root), "release").expect("parent should release child");
2065 assert!(child.wait().expect("child should exit").success());
2066 let reopened = LocalFileDataStore::new(&root).expect("released lock should reopen");
2067 assert_eq!(
2068 reopened
2069 .read("draft")
2070 .expect("committed record should remain readable"),
2071 Some(committed)
2072 );
2073 }
2074
2075 #[test]
2076 fn local_file_adapter_recovers_after_lock_owner_crash() {
2077 let root = temp_root("owner-crash-lock");
2078 let mut initial_owner = public_store(&root);
2079 initial_owner
2080 .write(record("draft", "writer-a", 1, json!("committed")))
2081 .expect("initial write should succeed");
2082 drop(initial_owner);
2083
2084 let mut child = start_lock_child(&root);
2085 wait_for_lock_child(&root, 500).expect("lock child should become ready");
2086 child.kill().expect("parent should terminate child");
2087 child.wait().expect("terminated child should exit");
2088
2089 let reopened = LocalFileDataStore::new(&root).expect("crashed owner lock should release");
2090 assert_eq!(
2091 reopened
2092 .read("draft")
2093 .expect("committed record should remain readable"),
2094 Some(record("draft", "writer-a", 1, json!("committed")))
2095 );
2096 }
2097
2098 #[test]
2099 fn local_file_adapter_reports_unknown_version_and_helper_failures_stably() {
2100 let root = temp_root("helper-failures");
2101 let mut adapter = public_store(&root);
2102 adapter
2103 .write(record("draft", "writer-a", 1, json!("ready")))
2104 .expect("write should succeed");
2105
2106 let mut envelope: Value = serde_json::from_slice(
2107 &fs::read(root.join("draft.json")).expect("envelope should be present"),
2108 )
2109 .expect("fixture should deserialize");
2110 envelope["format"] = json!("local-datastore/unsupported");
2111 fs::write(
2112 root.join("draft.json"),
2113 serde_json::to_vec(&envelope).expect("fixture should serialize"),
2114 )
2115 .expect("unknown-version fixture should write");
2116 let unknown_version = adapter
2117 .read("draft")
2118 .expect_err("unknown format version must fail");
2119 assert_eq!(
2120 unknown_version.code,
2121 DataStoreErrorCode::IntegrityCheckFailed
2122 );
2123 assert_eq!(unknown_version.details["reason"], "unknown_format_version");
2124
2125 let lock_io = lock_error(TryLockError::Error(std::io::Error::other(
2126 "lock device failure",
2127 )));
2128 assert_eq!(lock_io.code, DataStoreErrorCode::IoFailure);
2129 assert_eq!(lock_io.message, "storage_io_failed");
2130 assert_eq!(
2131 lock_io.details,
2132 json!({ "operation": "acquire_lock", "reason": "lock_acquisition_failed" })
2133 );
2134
2135 let unsupported_lock = lock_error(TryLockError::Error(std::io::Error::new(
2136 std::io::ErrorKind::Unsupported,
2137 "locking unavailable",
2138 )));
2139 assert_eq!(unsupported_lock.code, DataStoreErrorCode::IoFailure);
2140 assert_eq!(
2141 unsupported_lock.details,
2142 json!({ "operation": "acquire_lock", "reason": "locking_unsupported" })
2143 );
2144
2145 let durability = durability_error(
2146 &root,
2147 "temporary_file",
2148 &std::io::Error::other("sync failure"),
2149 );
2150 assert_eq!(durability.code, DataStoreErrorCode::DurabilityCommitFailed);
2151 assert_eq!(durability.message, "durability_commit_failed");
2152
2153 let temporary = root.join(".orphan.tmp");
2154 fs::write(&temporary, "incomplete").expect("temporary fixture should write");
2155 discard_temporary_record(&temporary);
2156 assert!(!temporary.exists());
2157 discard_temporary_record(&temporary);
2158
2159 let parse_error = serde_json::from_str::<Value>("{")
2160 .expect_err("invalid fixture must produce a serialization error");
2161 let serialization = serialization_error("deserialize fixture", &parse_error);
2162 assert_eq!(serialization.code, DataStoreErrorCode::SerializationFailure);
2163 }
2164
2165 fn record(key: &str, writer_id: &str, lamport_clock: u64, value: Value) -> StateRecord {
2166 StateRecord {
2167 key: key.to_string(),
2168 value,
2169 lamport_clock,
2170 writer_id: writer_id.to_string(),
2171 }
2172 }
2173
2174 fn temp_root(name: &str) -> PathBuf {
2175 std::env::temp_dir().join(format!("traverse-data-store-{name}-{}", Uuid::new_v4()))
2176 }
2177
2178 fn public_store(root: &Path) -> LocalFileDataStore {
2179 LocalFileDataStore::with_classification(root, LocalDataClassification::Public)
2180 .expect("public local adapter should initialize")
2181 }
2182
2183 fn write_envelope_fixture(path: &Path, envelope: &LocalDataStoreEnvelope) {
2184 fs::write(
2185 path,
2186 serde_json::to_vec(envelope).expect("serialize envelope"),
2187 )
2188 .expect("write envelope");
2189 }
2190
2191 fn assert_integrity_reason(store: &LocalFileDataStore, reason: &str) {
2192 let error = store.read("secret").expect_err("corruption must fail");
2193 assert_eq!(error.code, DataStoreErrorCode::IntegrityCheckFailed);
2194 assert_eq!(error.details["reason"], reason);
2195 }
2196
2197 #[cfg(feature = "datastore-encryption")]
2198 fn encrypted_fixture(
2199 record_key: &str,
2200 key_id: &str,
2201 key: [u8; 32],
2202 plaintext: &[u8],
2203 ) -> LocalDataStoreEnvelope {
2204 let cipher = Aes256Gcm::new_from_slice(&key).expect("valid AES-256 key");
2205 let nonce = fresh_aes_nonce();
2206 let ciphertext = cipher
2207 .encrypt(
2208 &nonce,
2209 Payload {
2210 msg: plaintext,
2211 aad: &private_record_aad(key_id, record_key),
2212 },
2213 )
2214 .expect("fixture encryption");
2215 let nonce = hex_encode(&nonce);
2216 let ciphertext = hex_encode(&ciphertext);
2217 LocalDataStoreEnvelope {
2218 format: LOCAL_DATA_STORE_FORMAT.to_string(),
2219 classification: LocalDataClassification::Private,
2220 digest: digest_for_private_envelope(key_id, record_key, &nonce, &ciphertext),
2221 record: None,
2222 record_key: Some(record_key.to_string()),
2223 key_id: Some(key_id.to_string()),
2224 nonce: Some(nonce),
2225 ciphertext: Some(ciphertext),
2226 retained_at: None,
2227 }
2228 }
2229
2230 fn lock_child_ready_path(root: &Path) -> PathBuf {
2231 root.join(".lock-child-ready")
2232 }
2233
2234 fn lock_child_release_path(root: &Path) -> PathBuf {
2235 root.join(".lock-child-release")
2236 }
2237
2238 fn start_lock_child(root: &Path) -> Child {
2239 Command::new(std::env::current_exe().expect("test binary path should resolve"))
2240 .args([
2241 "--exact",
2242 "data_store::tests::local_file_adapter_lock_child",
2243 "--nocapture",
2244 ])
2245 .env("TRAVERSE_DATA_STORE_LOCK_CHILD_ROOT", root)
2246 .stdout(Stdio::null())
2247 .stderr(Stdio::null())
2248 .spawn()
2249 .expect("lock child should start")
2250 }
2251
2252 #[test]
2253 fn lock_child_waits_report_bounded_timeouts() {
2254 let root = temp_root("lock-child-timeout");
2255 assert_eq!(
2256 wait_for_lock_child(&root, 0),
2257 Err("lock child did not become ready")
2258 );
2259 assert_eq!(
2260 wait_for_lock_child_release(&root, 0),
2261 Err("parent did not release the lock child".to_string())
2262 );
2263 }
2264
2265 fn wait_for_lock_child(root: &Path, attempts: usize) -> Result<(), &'static str> {
2266 if wait_for_path(&lock_child_ready_path(root), attempts) {
2267 Ok(())
2268 } else {
2269 Err("lock child did not become ready")
2270 }
2271 }
2272
2273 fn wait_for_lock_child_release(root: &Path, attempts: usize) -> Result<(), String> {
2274 if wait_for_path(&lock_child_release_path(root), attempts) {
2275 Ok(())
2276 } else {
2277 Err("parent did not release the lock child".to_string())
2278 }
2279 }
2280
2281 fn wait_for_path(path: &Path, attempts: usize) -> bool {
2282 for _ in 0..attempts {
2283 if path.exists() {
2284 return true;
2285 }
2286 thread::sleep(Duration::from_millis(10));
2287 }
2288 false
2289 }
2290
2291 #[test]
2292 fn v2_envelope_rejects_tampering_unknown_fields_and_invalid_inner_format() {
2293 let private = LocalDataStoreEnvelope {
2294 format: LOCAL_DATA_STORE_FORMAT.to_string(),
2295 classification: LocalDataClassification::Private,
2296 digest: "sha256:opaque".to_string(),
2297 record: None,
2298 record_key: Some("record".to_string()),
2299 key_id: Some("host-key".to_string()),
2300 nonce: Some("000000000000000000000000".to_string()),
2301 ciphertext: Some("00000000000000000000000000000000".to_string()),
2302 retained_at: None,
2303 };
2304 let wrapped = v2_envelope(private).expect("wrap private payload");
2305 assert_eq!(wrapped.encryption_disclosure, "host_managed_opaque");
2306
2307 let mut tampered = serde_json::to_value(&wrapped).expect("json");
2308 tampered["payload_integrity"] = Value::String("sha256:bad".to_string());
2309 assert!(decode_v2_envelope(tampered).is_err());
2310
2311 let mut payload_tampered = serde_json::to_value(&wrapped).expect("json");
2312 payload_tampered["payload_integrity"] = Value::String("sha256:bad".to_string());
2313 payload_tampered["integrity"]["content_digest"] = Value::String("sha256:bad".to_string());
2314 assert!(decode_v2_envelope(payload_tampered).is_err());
2315
2316 let mut unknown = serde_json::to_value(&wrapped).expect("json");
2317 unknown["format_version"] = json!(99);
2318 assert!(decode_v2_envelope(unknown).is_err());
2319
2320 let mut invalid_inner = wrapped;
2321 invalid_inner.payload.format = "unknown/9".to_string();
2322 let bytes = serde_json::to_vec(&invalid_inner.payload).expect("payload");
2323 invalid_inner.payload_integrity = digest_bytes(&bytes);
2324 invalid_inner.integrity.content_digest = invalid_inner.payload_integrity.clone();
2325 assert!(decode_v2_envelope(serde_json::to_value(invalid_inner).expect("json")).is_err());
2326 }
2327
2328 fn stateful_contract(state_schema: Option<Value>) -> CapabilityContract {
2329 CapabilityContract {
2330 kind: "capability_contract".to_string(),
2331 schema_version: "1.0.0".to_string(),
2332 id: "stateful.example".to_string(),
2333 namespace: "stateful".to_string(),
2334 name: "example".to_string(),
2335 version: "1.0.0".to_string(),
2336 lifecycle: Lifecycle::Active,
2337 owner: Owner {
2338 team: "runtime".to_string(),
2339 contact: "runtime@example.com".to_string(),
2340 },
2341 summary: "Stateful test capability".to_string(),
2342 description: "Stateful test capability".to_string(),
2343 inputs: SchemaContainer {
2344 schema: json!({"type": "object"}),
2345 },
2346 outputs: SchemaContainer {
2347 schema: json!({"type": "object"}),
2348 },
2349 preconditions: Vec::<Condition>::new(),
2350 postconditions: Vec::<Condition>::new(),
2351 side_effects: vec![SideEffect {
2352 kind: SideEffectKind::StateChange,
2353 description: "writes capability state".to_string(),
2354 }],
2355 emits: Vec::<EventReference>::new(),
2356 consumes: Vec::<EventReference>::new(),
2357 permissions: Vec::<IdReference>::new(),
2358 execution: Execution {
2359 binary_format: BinaryFormat::Wasm,
2360 constraints: ExecutionConstraints {
2361 network_access: NetworkAccess::Forbidden,
2362 filesystem_access: FilesystemAccess::SandboxOnly,
2363 host_api_access: HostApiAccess::None,
2364 },
2365 entrypoint: Entrypoint {
2366 kind: EntrypointKind::WasiCommand,
2367 command: "run".to_string(),
2368 },
2369 preferred_targets: vec![ExecutionTarget::Local],
2370 },
2371 policies: Vec::<IdReference>::new(),
2372 dependencies: Vec::<DependencyReference>::new(),
2373 provenance: Provenance {
2374 source: ProvenanceSource::Greenfield,
2375 author: "Codex".to_string(),
2376 created_at: "2026-04-19T00:00:00Z".to_string(),
2377 spec_ref: Some("032-universal-data-access".to_string()),
2378 adr_refs: Vec::new(),
2379 exception_refs: Vec::new(),
2380 },
2381 evidence: Vec::<ValidationEvidence>::new(),
2382 service_type: ServiceType::Stateful,
2383 permitted_targets: vec![ExecutionTarget::Local],
2384 event_trigger: None,
2385 connector_requirements: Vec::new(),
2386 state_schema,
2387 use_cases: Vec::new(),
2388 risk: traverse_contracts::default_risk_metadata(),
2389 }
2390 }
2391}