Skip to main content

traverse_runtime/
data_store.rs

1//! Portable state access for Traverse capabilities.
2//!
3//! General operations are governed by spec `032-universal-data-access`; the
4//! local-file adapter durability boundary is governed by spec
5//! `518-durable-local-datastore`. Retention prune and verified backup/restore
6//! are governed by spec `083-datastore-retention-backup`.
7
8#[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
51/// The synchronization report is governed by the approved protocol, rather
52/// than the earlier generic `DataStore` surface that supplied its merge helper.
53const 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/// Classification recorded with each locally durable record.
128#[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    /// Host-supplied RFC3339 instant stamped at write time for age-based prune.
151    /// Absent on legacy envelopes; age prune treats missing stamps as retained.
152    #[serde(default, skip_serializing_if = "Option::is_none")]
153    retained_at: Option<String>,
154}
155
156/// Version-two wrapper keeps the durable payload self-identifying and
157/// independently integrity-protected. The payload remains opaque to callers.
158#[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/// Stable, secret-free key-provider failure codes.
225#[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/// A key-provider failure that never includes key material or provider internals.
234#[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
266/// Host-owned source of AES-256 keys.
267pub trait KeyProvider: Send + Sync {
268    /// Returns the key id used for new private writes.
269    ///
270    /// # Errors
271    ///
272    /// Returns a secret-free [`KeyProviderError`] when no write key is available.
273    fn active_key_id(&self) -> Result<String, KeyProviderError>;
274
275    /// Returns key material for an envelope key id.
276    ///
277    /// # Errors
278    ///
279    /// Returns a secret-free [`KeyProviderError`] for missing, expired, or
280    /// unavailable keys.
281    fn key_for(&self, key_id: &str) -> Result<Zeroizing<[u8; 32]>, KeyProviderError>;
282}
283
284/// In-memory provider for tests and host wiring.
285#[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    /// Reads a stored state record.
343    ///
344    /// # Errors
345    ///
346    /// Returns [`DataStoreError`] when the adapter cannot read the key.
347    fn read(&self, key: &str) -> Result<Option<StateRecord>, DataStoreError>;
348
349    /// Writes a stamped state record.
350    ///
351    /// # Errors
352    ///
353    /// Returns [`DataStoreError`] when the adapter cannot persist the record.
354    fn write(&mut self, record: StateRecord) -> Result<(), DataStoreError>;
355
356    /// Deletes a stored state record.
357    ///
358    /// # Errors
359    ///
360    /// Returns [`DataStoreError`] when the adapter cannot delete the key.
361    fn delete(&mut self, key: &str) -> Result<(), DataStoreError>;
362
363    /// Lists stored state keys.
364    ///
365    /// # Errors
366    ///
367    /// Returns [`DataStoreError`] when the adapter cannot enumerate keys.
368    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    /// Reads and validates a state value by key.
427    ///
428    /// # Errors
429    ///
430    /// Returns [`DataStoreError`] when the key is invalid, the adapter cannot
431    /// read the key, or the stored value violates the contract state schema.
432    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    /// Validates, stamps, and writes a state value for a capability contract.
452    ///
453    /// # Errors
454    ///
455    /// Returns [`DataStoreError`] when the key is invalid, no state schema is
456    /// declared, schema validation fails, the Lamport clock overflows, or the
457    /// adapter cannot persist the stamped record.
458    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    /// Deletes a state value by key.
476    ///
477    /// # Errors
478    ///
479    /// Returns [`DataStoreError`] when the adapter cannot delete the key.
480    pub fn delete(&mut self, key: &str) -> Result<(), DataStoreError> {
481        self.adapter.delete(key)
482    }
483
484    /// Lists state keys.
485    ///
486    /// # Errors
487    ///
488    /// Returns [`DataStoreError`] when the adapter cannot enumerate keys.
489    pub fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
490        self.adapter.list_keys()
491    }
492
493    /// Triggers explicit sync after a reconnect event.
494    ///
495    /// # Errors
496    ///
497    /// Returns [`DataStoreError`] when either adapter cannot read, write, list,
498    /// or restore state during sync.
499    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    /// Host-supplied retained-at stamp applied to subsequent writes (no OS clock).
517    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    /// Creates a local filesystem-backed data store rooted at `root`.
540    ///
541    /// # Errors
542    ///
543    /// Returns [`DataStoreError`] when the root directory cannot be created or
544    /// another process owns the store.
545    pub fn new(root: impl Into<PathBuf>) -> Result<Self, DataStoreError> {
546        Self::with_classification(root, LocalDataClassification::Private)
547    }
548
549    /// Creates a local filesystem-backed data store with explicit persisted
550    /// record classification.
551    ///
552    /// # Errors
553    ///
554    /// Returns [`DataStoreError`] when the root directory cannot be created or
555    /// another process owns the store.
556    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    /// Configures the host-owned key provider used for private records.
581    #[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    /// Sets the host-supplied retained-at stamp used by subsequent writes.
588    ///
589    /// Traverse does not read the OS clock; hosts must supply RFC3339 instants
590    /// when age-based retention is desired.
591    pub fn set_write_retained_at(&mut self, retained_at: Option<String>) {
592        self.write_retained_at = retained_at;
593    }
594
595    /// Returns the store root path for host-owned maintenance construction.
596    #[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    // Host-local UUID entropy avoids aes-gcm's getrandom Generate path so
1010    // wasm32 `--no-default-features` checks do not require a getrandom backend.
1011    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
1061/// Validates a capability state write against the contract-declared state schema.
1062///
1063/// # Errors
1064///
1065/// Returns [`DataStoreError`] when the key is invalid, the contract does not
1066/// declare a state schema, the key is not declared by the schema, or the value
1067/// does not match the declared key schema.
1068pub 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}