1use super::{
10 DataStore, DataStoreError, DataStoreErrorCode, LocalDataClassification, StateRecord,
11 data_store_error, digest_for_record, serialization_error, validate_key,
12};
13use serde::{Deserialize, Serialize};
14use serde_json::json;
15use std::cell::RefCell;
16
17const REMOTE_DATA_STORE_FORMAT: &str = "remote-datastore/1";
18const REMOTE_CONSISTENCY_MODE: &str = "read-your-write";
19
20#[derive(Debug, Clone, PartialEq, Eq)]
23pub struct RemoteVersionToken(pub String);
24
25#[derive(Debug, Clone, PartialEq, Eq)]
27pub struct RemoteObject {
28 pub bytes: Vec<u8>,
29 pub version: RemoteVersionToken,
30}
31
32#[derive(Debug, Clone, Copy, PartialEq, Eq)]
35pub enum RemoteBackendFailure {
36 Unavailable,
37 Timeout,
38 Unauthorized,
39 ScopeDenied,
40 BackendFailed,
41}
42
43#[derive(Debug, Clone, Copy, PartialEq, Eq)]
45pub enum RemoteWriteOutcome {
46 Acknowledged { retry_count: u32 },
48 Conflict,
50 Unknown { retry_count: u32 },
53 Failed(RemoteBackendFailure),
55}
56
57pub trait RemoteDataStoreBackend {
61 fn get(&self, key: &str) -> Result<Option<RemoteObject>, RemoteBackendFailure>;
68
69 fn put_conditional(
72 &mut self,
73 key: &str,
74 bytes: Vec<u8>,
75 expected_version: Option<RemoteVersionToken>,
76 ) -> RemoteWriteOutcome;
77
78 fn delete_conditional(
81 &mut self,
82 key: &str,
83 expected_version: Option<RemoteVersionToken>,
84 ) -> RemoteWriteOutcome;
85}
86
87#[derive(Debug, Clone, PartialEq, Eq)]
91pub struct RemoteOperationEvidence {
92 pub operation: &'static str,
93 pub outcome: &'static str,
94 pub classification: LocalDataClassification,
95 pub consistency_mode: &'static str,
96 pub retry_count: u32,
97 pub code: Option<&'static str>,
98}
99
100#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
101struct RemoteDataStoreEnvelope {
102 format: String,
103 classification: LocalDataClassification,
104 digest: String,
105 record: StateRecord,
106}
107
108pub struct RemoteKeyValueDataStore<B: RemoteDataStoreBackend> {
112 backend: B,
113 classification: LocalDataClassification,
114 last_evidence: RefCell<Option<RemoteOperationEvidence>>,
115}
116
117impl<B: RemoteDataStoreBackend> RemoteKeyValueDataStore<B> {
118 #[must_use]
119 pub fn new(backend: B, classification: LocalDataClassification) -> Self {
120 Self {
121 backend,
122 classification,
123 last_evidence: RefCell::new(None),
124 }
125 }
126
127 #[must_use]
129 pub fn last_evidence(&self) -> Option<RemoteOperationEvidence> {
130 self.last_evidence.borrow().clone()
131 }
132
133 fn record_evidence(
134 &self,
135 operation: &'static str,
136 outcome: &'static str,
137 retry_count: u32,
138 code: Option<&'static str>,
139 ) {
140 *self.last_evidence.borrow_mut() = Some(RemoteOperationEvidence {
141 operation,
142 outcome,
143 classification: self.classification,
144 consistency_mode: REMOTE_CONSISTENCY_MODE,
145 retry_count,
146 code,
147 });
148 }
149
150 fn decode_envelope(&self, bytes: &[u8]) -> Result<StateRecord, DataStoreError> {
151 let envelope: RemoteDataStoreEnvelope = serde_json::from_slice(bytes)
152 .map_err(|_| remote_integrity_error("malformed_envelope"))?;
153 if envelope.format != REMOTE_DATA_STORE_FORMAT {
154 return Err(remote_integrity_error("unknown_format_version"));
155 }
156 if envelope.classification != self.classification {
157 return Err(remote_integrity_error("classification_mismatch"));
158 }
159 let expected = digest_for_record(&envelope.record)?;
160 if expected != envelope.digest {
161 return Err(remote_integrity_error("digest_mismatch"));
162 }
163 Ok(envelope.record)
164 }
165
166 fn current_version(
167 &self,
168 operation: &'static str,
169 key: &str,
170 ) -> Result<Option<RemoteVersionToken>, DataStoreError> {
171 match self.backend.get(key) {
172 Ok(object) => Ok(object.map(|object| object.version)),
173 Err(failure) => {
174 let code = remote_failure_code(failure);
175 self.record_evidence(operation, "failed", 0, Some(code));
176 Err(remote_backend_error(failure))
177 }
178 }
179 }
180
181 fn handle_write_outcome(
182 &self,
183 operation: &'static str,
184 outcome: RemoteWriteOutcome,
185 ) -> Result<(), DataStoreError> {
186 match outcome {
187 RemoteWriteOutcome::Acknowledged { retry_count } => {
188 self.record_evidence(operation, "acknowledged", retry_count, None);
189 Ok(())
190 }
191 RemoteWriteOutcome::Conflict => {
192 self.record_evidence(operation, "conflict", 0, Some("remote_conflict"));
193 Err(remote_conflict_error())
194 }
195 RemoteWriteOutcome::Unknown { retry_count } => {
196 self.record_evidence(
197 operation,
198 "outcome_unknown",
199 retry_count,
200 Some("remote_outcome_unknown"),
201 );
202 Err(remote_outcome_unknown_error(retry_count))
203 }
204 RemoteWriteOutcome::Failed(failure) => {
205 let code = remote_failure_code(failure);
206 self.record_evidence(operation, "failed", 0, Some(code));
207 Err(remote_backend_error(failure))
208 }
209 }
210 }
211}
212
213impl<B: RemoteDataStoreBackend> DataStore for RemoteKeyValueDataStore<B> {
214 fn read(&self, key: &str) -> Result<Option<StateRecord>, DataStoreError> {
215 validate_key(key)?;
216 match self.backend.get(key) {
217 Ok(None) => {
218 self.record_evidence("read", "not_found", 0, None);
219 Ok(None)
220 }
221 Ok(Some(object)) => match self.decode_envelope(&object.bytes) {
222 Ok(record) => {
223 self.record_evidence("read", "acknowledged", 0, None);
224 Ok(Some(record))
225 }
226 Err(error) => {
227 self.record_evidence(
228 "read",
229 "integrity_failed",
230 0,
231 Some("remote_integrity_failed"),
232 );
233 Err(error)
234 }
235 },
236 Err(failure) => {
237 let code = remote_failure_code(failure);
238 self.record_evidence("read", "failed", 0, Some(code));
239 Err(remote_backend_error(failure))
240 }
241 }
242 }
243
244 fn write(&mut self, record: StateRecord) -> Result<(), DataStoreError> {
245 validate_key(&record.key)?;
246 let expected_version = self.current_version("write", &record.key)?;
247 let digest = digest_for_record(&record)?;
248 let envelope = RemoteDataStoreEnvelope {
249 format: REMOTE_DATA_STORE_FORMAT.to_string(),
250 classification: self.classification,
251 digest,
252 record: record.clone(),
253 };
254 let bytes = serde_json::to_vec(&envelope)
255 .map_err(|error| serialization_error("serialize remote envelope", &error))?;
256 let outcome = self
257 .backend
258 .put_conditional(&record.key, bytes, expected_version);
259 self.handle_write_outcome("write", outcome)
260 }
261
262 fn delete(&mut self, key: &str) -> Result<(), DataStoreError> {
263 validate_key(key)?;
264 let expected_version = self.current_version("delete", key)?;
265 let outcome = self.backend.delete_conditional(key, expected_version);
266 self.handle_write_outcome("delete", outcome)
267 }
268
269 fn list_keys(&self) -> Result<Vec<String>, DataStoreError> {
270 self.record_evidence("list_keys", "unsupported", 0, Some("remote_backend_failed"));
271 Err(data_store_error(
272 DataStoreErrorCode::RemoteBackendFailed,
273 "remote_backend_failed",
274 json!({ "reason": "remote_scan_not_supported_in_v1" }),
275 ))
276 }
277}
278
279fn remote_failure_code(failure: RemoteBackendFailure) -> &'static str {
280 match failure {
281 RemoteBackendFailure::Unavailable => "remote_unavailable",
282 RemoteBackendFailure::Timeout => "remote_timeout",
283 RemoteBackendFailure::Unauthorized => "remote_unauthorized",
284 RemoteBackendFailure::ScopeDenied => "remote_scope_denied",
285 RemoteBackendFailure::BackendFailed => "remote_backend_failed",
286 }
287}
288
289fn remote_backend_error(failure: RemoteBackendFailure) -> DataStoreError {
290 let code = match failure {
291 RemoteBackendFailure::Unavailable => DataStoreErrorCode::RemoteUnavailable,
292 RemoteBackendFailure::Timeout => DataStoreErrorCode::RemoteTimeout,
293 RemoteBackendFailure::Unauthorized => DataStoreErrorCode::RemoteUnauthorized,
294 RemoteBackendFailure::ScopeDenied => DataStoreErrorCode::RemoteScopeDenied,
295 RemoteBackendFailure::BackendFailed => DataStoreErrorCode::RemoteBackendFailed,
296 };
297 data_store_error(code, remote_failure_code(failure), json!({}))
298}
299
300fn remote_conflict_error() -> DataStoreError {
301 data_store_error(
302 DataStoreErrorCode::RemoteConflict,
303 "remote_conflict",
304 json!({}),
305 )
306}
307
308fn remote_outcome_unknown_error(retry_count: u32) -> DataStoreError {
309 data_store_error(
310 DataStoreErrorCode::RemoteOutcomeUnknown,
311 "remote_outcome_unknown",
312 json!({ "retry_count": retry_count }),
313 )
314}
315
316fn remote_integrity_error(reason: &str) -> DataStoreError {
317 data_store_error(
318 DataStoreErrorCode::RemoteIntegrityFailed,
319 "remote_integrity_failed",
320 json!({ "reason": reason }),
321 )
322}
323
324#[cfg(test)]
325#[allow(clippy::expect_used, clippy::unwrap_used)]
326mod tests {
327 use super::*;
328 use std::collections::BTreeMap;
329
330 #[derive(Debug, Clone)]
331 struct StoredObject {
332 bytes: Vec<u8>,
333 version: u64,
334 }
335
336 #[derive(Default)]
343 struct FakeS3Backend {
344 objects: BTreeMap<String, StoredObject>,
345 next_version: u64,
346 fail_get: Option<RemoteBackendFailure>,
347 fail_put: Option<RemoteBackendFailure>,
348 fail_delete: Option<RemoteBackendFailure>,
349 force_conflict: bool,
350 force_unknown_retry_count: Option<u32>,
351 }
352
353 impl FakeS3Backend {
354 fn version_token(version: u64) -> RemoteVersionToken {
355 RemoteVersionToken(format!("v{version}"))
356 }
357
358 fn current_token(&self, key: &str) -> Option<RemoteVersionToken> {
359 self.objects
360 .get(key)
361 .map(|object| Self::version_token(object.version))
362 }
363 }
364
365 impl RemoteDataStoreBackend for FakeS3Backend {
366 fn get(&self, key: &str) -> Result<Option<RemoteObject>, RemoteBackendFailure> {
367 if let Some(failure) = self.fail_get {
368 return Err(failure);
369 }
370 Ok(self.objects.get(key).map(|object| RemoteObject {
371 bytes: object.bytes.clone(),
372 version: Self::version_token(object.version),
373 }))
374 }
375
376 fn put_conditional(
377 &mut self,
378 key: &str,
379 bytes: Vec<u8>,
380 expected_version: Option<RemoteVersionToken>,
381 ) -> RemoteWriteOutcome {
382 if let Some(failure) = self.fail_put {
383 return RemoteWriteOutcome::Failed(failure);
384 }
385 if self.force_conflict || expected_version != self.current_token(key) {
386 self.force_conflict = false;
387 return RemoteWriteOutcome::Conflict;
388 }
389 if let Some(retry_count) = self.force_unknown_retry_count.take() {
390 return RemoteWriteOutcome::Unknown { retry_count };
391 }
392 self.next_version += 1;
393 let version = self.next_version;
394 self.objects
395 .insert(key.to_string(), StoredObject { bytes, version });
396 RemoteWriteOutcome::Acknowledged { retry_count: 0 }
397 }
398
399 fn delete_conditional(
400 &mut self,
401 key: &str,
402 expected_version: Option<RemoteVersionToken>,
403 ) -> RemoteWriteOutcome {
404 if let Some(failure) = self.fail_delete {
405 return RemoteWriteOutcome::Failed(failure);
406 }
407 if self.force_conflict || expected_version != self.current_token(key) {
408 self.force_conflict = false;
409 return RemoteWriteOutcome::Conflict;
410 }
411 self.objects.remove(key);
412 RemoteWriteOutcome::Acknowledged { retry_count: 0 }
413 }
414 }
415
416 fn record(key: &str, value: &str) -> StateRecord {
417 StateRecord {
418 key: key.to_string(),
419 value: json!({ "v": value }),
420 lamport_clock: 1,
421 writer_id: "writer-a".to_string(),
422 }
423 }
424
425 fn store_with_backend(
426 backend: FakeS3Backend,
427 classification: LocalDataClassification,
428 ) -> RemoteKeyValueDataStore<FakeS3Backend> {
429 RemoteKeyValueDataStore::new(backend, classification)
430 }
431
432 #[test]
433 fn write_then_read_your_write_round_trips() {
434 let mut store =
435 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
436 store.write(record("alpha", "one")).expect("write succeeds");
437 let read = store.read("alpha").expect("read succeeds");
438 assert_eq!(read, Some(record("alpha", "one")));
439 let evidence = store.last_evidence().expect("evidence recorded");
440 assert_eq!(evidence.operation, "read");
441 assert_eq!(evidence.outcome, "acknowledged");
442 assert_eq!(evidence.consistency_mode, "read-your-write");
443 }
444
445 #[test]
446 fn read_of_missing_key_returns_none_without_error() {
447 let store = store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
448 assert_eq!(store.read("missing").expect("read succeeds"), None);
449 assert_eq!(store.last_evidence().unwrap().outcome, "not_found");
450 }
451
452 #[test]
453 fn ambiguous_write_outcome_is_reported_as_outcome_unknown_not_success() {
454 let backend = FakeS3Backend {
455 force_unknown_retry_count: Some(2),
456 ..Default::default()
457 };
458 let mut store = store_with_backend(backend, LocalDataClassification::Public);
459 let error = store
460 .write(record("alpha", "one"))
461 .expect_err("ambiguous outcome fails");
462 assert_eq!(error.code, DataStoreErrorCode::RemoteOutcomeUnknown);
463 assert_eq!(error.message, "remote_outcome_unknown");
464 let evidence = store.last_evidence().unwrap();
465 assert_eq!(evidence.outcome, "outcome_unknown");
466 assert_eq!(evidence.retry_count, 2);
467 assert_eq!(evidence.code, Some("remote_outcome_unknown"));
468 assert_eq!(store.read("alpha").expect("read succeeds"), None);
469 }
470
471 #[test]
472 fn concurrent_write_conflict_never_overwrites_or_retries() {
473 let mut store =
474 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
475 store
476 .write(record("alpha", "one"))
477 .expect("first write succeeds");
478 store.backend.force_conflict = true;
479 let error = store
480 .write(record("alpha", "two"))
481 .expect_err("conflicting write fails closed");
482 assert_eq!(error.code, DataStoreErrorCode::RemoteConflict);
483 assert_eq!(
484 store.read("alpha").expect("read succeeds"),
485 Some(record("alpha", "one")),
486 "a conflicting write must never overwrite the existing value"
487 );
488 }
489
490 #[test]
491 fn delete_conflict_is_reported_and_leaves_the_record_intact() {
492 let mut store =
493 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
494 store.write(record("alpha", "one")).expect("write succeeds");
495 store.backend.force_conflict = true;
496 let error = store
497 .delete("alpha")
498 .expect_err("conflicting delete fails closed");
499 assert_eq!(error.code, DataStoreErrorCode::RemoteConflict);
500 assert!(store.read("alpha").expect("read succeeds").is_some());
501 }
502
503 #[test]
504 fn delete_removes_the_record_and_subsequent_read_returns_none() {
505 let mut store =
506 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
507 store.write(record("alpha", "one")).expect("write succeeds");
508 store.delete("alpha").expect("delete succeeds");
509 let evidence = store.last_evidence().expect("evidence recorded");
510 assert_eq!(evidence.operation, "delete");
511 assert_eq!(evidence.outcome, "acknowledged");
512 assert_eq!(store.read("alpha").expect("read succeeds"), None);
513 }
514
515 #[test]
516 fn denied_scope_and_unauthorized_map_to_stable_codes() {
517 let scope_denied = FakeS3Backend {
518 fail_get: Some(RemoteBackendFailure::ScopeDenied),
519 ..Default::default()
520 };
521 let scope_store = store_with_backend(scope_denied, LocalDataClassification::Public);
522 let error = scope_store.read("alpha").expect_err("denied scope fails");
523 assert_eq!(error.code, DataStoreErrorCode::RemoteScopeDenied);
524 assert_eq!(error.message, "remote_scope_denied");
525
526 let unauthorized = FakeS3Backend {
527 fail_get: Some(RemoteBackendFailure::Unauthorized),
528 ..Default::default()
529 };
530 let auth_store = store_with_backend(unauthorized, LocalDataClassification::Public);
531 let error = auth_store.read("alpha").expect_err("unauthorized fails");
532 assert_eq!(error.code, DataStoreErrorCode::RemoteUnauthorized);
533 assert_eq!(error.message, "remote_unauthorized");
534 }
535
536 #[test]
537 fn outage_and_timeout_map_to_stable_codes() {
538 let unavailable = FakeS3Backend {
539 fail_get: Some(RemoteBackendFailure::Unavailable),
540 ..Default::default()
541 };
542 let store = store_with_backend(unavailable, LocalDataClassification::Public);
543 assert_eq!(
544 store.read("alpha").expect_err("outage fails").code,
545 DataStoreErrorCode::RemoteUnavailable
546 );
547
548 let timeout = FakeS3Backend {
549 fail_get: Some(RemoteBackendFailure::Timeout),
550 ..Default::default()
551 };
552 let store = store_with_backend(timeout, LocalDataClassification::Public);
553 assert_eq!(
554 store.read("alpha").expect_err("timeout fails").code,
555 DataStoreErrorCode::RemoteTimeout
556 );
557 }
558
559 #[test]
560 fn backend_failure_during_write_precondition_lookup_is_reported() {
561 let backend = FakeS3Backend {
562 fail_get: Some(RemoteBackendFailure::BackendFailed),
563 ..Default::default()
564 };
565 let mut store = store_with_backend(backend, LocalDataClassification::Public);
566 let error = store
567 .write(record("alpha", "one"))
568 .expect_err("precondition lookup failure surfaces");
569 assert_eq!(error.code, DataStoreErrorCode::RemoteBackendFailed);
570 }
571
572 #[test]
573 fn backend_failure_during_put_after_precondition_check_is_reported() {
574 let mut store =
575 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
576 store.backend.fail_put = Some(RemoteBackendFailure::BackendFailed);
577 let error = store
578 .write(record("alpha", "one"))
579 .expect_err("backend failure during put surfaces");
580 assert_eq!(error.code, DataStoreErrorCode::RemoteBackendFailed);
581 let evidence = store.last_evidence().unwrap();
582 assert_eq!(evidence.operation, "write");
583 assert_eq!(evidence.outcome, "failed");
584 }
585
586 #[test]
587 fn backend_failure_during_delete_after_precondition_check_is_reported() {
588 let mut store =
589 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
590 store.write(record("alpha", "one")).expect("write succeeds");
591 store.backend.fail_delete = Some(RemoteBackendFailure::Unavailable);
592 let error = store
593 .delete("alpha")
594 .expect_err("backend failure during delete surfaces");
595 assert_eq!(error.code, DataStoreErrorCode::RemoteUnavailable);
596 }
597
598 #[test]
599 fn malformed_remote_bytes_fail_closed() {
600 let mut store =
601 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
602 store.write(record("alpha", "one")).expect("write succeeds");
603 if let Some(object) = store.backend.objects.get_mut("alpha") {
604 object.bytes = b"not json".to_vec();
605 }
606 let error = store
607 .read("alpha")
608 .expect_err("malformed bytes fail closed");
609 assert_eq!(error.code, DataStoreErrorCode::RemoteIntegrityFailed);
610 }
611
612 #[test]
613 fn unknown_envelope_format_fails_closed() {
614 let mut store =
615 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
616 store.write(record("alpha", "one")).expect("write succeeds");
617 let envelope = RemoteDataStoreEnvelope {
618 format: "remote-datastore/999".to_string(),
619 classification: LocalDataClassification::Public,
620 digest: "sha256:00".to_string(),
621 record: record("alpha", "one"),
622 };
623 let bytes = serde_json::to_vec(&envelope).expect("serialize envelope");
624 if let Some(object) = store.backend.objects.get_mut("alpha") {
625 object.bytes = bytes;
626 }
627 let error = store
628 .read("alpha")
629 .expect_err("unknown format fails closed");
630 assert_eq!(error.code, DataStoreErrorCode::RemoteIntegrityFailed);
631 }
632
633 #[test]
634 fn tampered_remote_bytes_fail_closed_with_integrity_error() {
635 let mut store =
636 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
637 store.write(record("alpha", "one")).expect("write succeeds");
638 let envelope = RemoteDataStoreEnvelope {
639 format: REMOTE_DATA_STORE_FORMAT.to_string(),
640 classification: LocalDataClassification::Public,
641 digest: "sha256:0000000000000000000000000000000000000000000000000000000000000000"
642 .to_string(),
643 record: record("alpha", "tampered"),
644 };
645 let bytes = serde_json::to_vec(&envelope).expect("serialize envelope");
646 if let Some(object) = store.backend.objects.get_mut("alpha") {
647 object.bytes = bytes;
648 }
649 let error = store.read("alpha").expect_err("tampered bytes fail closed");
650 assert_eq!(error.code, DataStoreErrorCode::RemoteIntegrityFailed);
651 }
652
653 #[test]
654 fn classification_mismatch_fails_closed() {
655 let mut store =
656 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Private);
657 store.write(record("alpha", "one")).expect("write succeeds");
658 let backend = std::mem::take(&mut store.backend);
659 let mismatched = store_with_backend(backend, LocalDataClassification::Public);
660 let error = mismatched
661 .read("alpha")
662 .expect_err("classification mismatch fails closed");
663 assert_eq!(error.code, DataStoreErrorCode::RemoteIntegrityFailed);
664 }
665
666 #[test]
667 fn list_keys_is_explicitly_unsupported_in_v1() {
668 let store = store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
669 let error = store.list_keys().expect_err("scan is out of scope for v1");
670 assert_eq!(error.code, DataStoreErrorCode::RemoteBackendFailed);
671 }
672
673 #[test]
674 fn evidence_never_contains_keys_values_or_credentials() {
675 let mut store =
676 store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
677 store
678 .write(record("super-secret-key", "super-secret-value"))
679 .expect("write succeeds");
680 let evidence = store.last_evidence().expect("evidence recorded");
681 let debug = format!("{evidence:?}");
682 assert!(!debug.contains("super-secret-key"));
683 assert!(!debug.contains("super-secret-value"));
684 }
685
686 #[test]
687 fn invalid_key_is_rejected_before_touching_the_backend() {
688 let store = store_with_backend(FakeS3Backend::default(), LocalDataClassification::Public);
689 let error = store.read("has a space").expect_err("invalid key rejected");
690 assert_eq!(error.code, DataStoreErrorCode::InvalidKey);
691 }
692}