oxedyne/fe2o3/fe2o3_o3db_sync/tests/dist_o3db_storage.rs
5.6 KiB, 9 runs
created by r1870400018:11531, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! Integration test for the O3db-backed [`Storage`] adapter. |
| 2 | //! |
| 3 | //! Boots a fresh [`O3db`] instance in a per-test directory, wraps it in |
| 4 | //! [`O3dbStorage`], and exercises the full [`Storage`] contract |
| 5 | //! (`put`, `get`, `delete`, `digests`) plus a few edge cases. The test |
| 6 | //! runs only with `--features dist`. |
| 7 | //! |
| 8 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 9 | //! Anthropic Claude |
| 10 | |
| 11 | #![cfg(feature = "dist")] |
| 12 | |
| 13 | use oxedyne_fe2o3_core::prelude::*; |
| 14 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 15 | use oxedyne_fe2o3_hash::{ |
| 16 | csum::ChecksumScheme, |
| 17 | hash::HashScheme, |
| 18 | }; |
| 19 | use oxedyne_fe2o3_jdat::prelude::*; |
| 20 | use oxedyne_fe2o3_o3db_sync::{ |
| 21 | data::core::RestSchemesInput, |
| 22 | dist::{ |
| 23 | o3db_storage::O3dbStorage, |
| 24 | record::{ |
| 25 | Record, |
| 26 | RecordId, |
| 27 | }, |
| 28 | storage::Storage, |
| 29 | }, |
| 30 | test::setup, |
| 31 | }; |
| 32 | |
| 33 | use std::{ |
| 34 | path::Path, |
| 35 | sync::Arc, |
| 36 | thread, |
| 37 | time::Duration, |
| 38 | }; |
| 39 | |
| 40 | |
| 41 | #[test] |
| 42 | fn o3db_storage_round_trip() -> Outcome<()> { |
| 43 | |
| 44 | let db_root = res!(Path::new("./test_db_dist_o3db_storage") |
| 45 | .canonicalize() |
| 46 | .or_else(|_| -> std::io::Result<_> { |
| 47 | ok!(std::fs::create_dir_all("./test_db_dist_o3db_storage")); |
| 48 | Path::new("./test_db_dist_o3db_storage").canonicalize() |
| 49 | })); |
| 50 | |
| 51 | let enckey = [0x33u8; 32]; |
| 52 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..])); |
| 53 | let crc32 = ChecksumScheme::new_crc32(); |
| 54 | let user = setup::Uid::default(); |
| 55 | |
| 56 | let schms_input = RestSchemesInput::new( |
| 57 | Some(aes_gcm.clone()), |
| 58 | None::<HashScheme>, |
| 59 | None::<HashScheme>, |
| 60 | Some(crc32.clone()), |
| 61 | ); |
| 62 | |
| 63 | let mut cfg = res!(setup::default_cfg()); |
| 64 | cfg.num_zones = 2; |
| 65 | cfg.num_cbots_per_zone = 2; |
| 66 | cfg.num_igbots_per_zone = 2; |
| 67 | cfg.data_file_max_bytes = 200_000; |
| 68 | // Keep zone directories inside the per-test db root so this test |
| 69 | // cannot inherit or contaminate state from sibling tests. |
| 70 | cfg.zone_overrides = mapdat!{ |
| 71 | 1u16 => mapdat!{ "dir" => "", "max_size" => 1_000_000u64 }, |
| 72 | }.get_map().unwrap(); |
| 73 | |
| 74 | let db = res!(setup::start_db( |
| 75 | db_root.clone(), |
| 76 | Some(cfg), |
| 77 | schms_input, |
| 78 | None, |
| 79 | true, // gc on |
| 80 | true, // wipe pre-existing |
| 81 | )); |
| 82 | |
| 83 | thread::sleep(Duration::from_secs(1)); |
| 84 | |
| 85 | let db = Arc::new(db); |
| 86 | let storage = O3dbStorage::new(Arc::clone(&db), user); |
| 87 | |
| 88 | // 1. Put three records across two tables. |
| 89 | let id_a = RecordId::from_bytes([0x01; 32]); |
| 90 | let rec_a = Record::new(id_a, "identity", b"alpha".to_vec()); |
| 91 | res!(storage.put(&rec_a)); |
| 92 | |
| 93 | let id_b = RecordId::from_bytes([0x02; 32]); |
| 94 | let rec_b = Record::new(id_b, "identity", b"beta".to_vec()); |
| 95 | res!(storage.put(&rec_b)); |
| 96 | |
| 97 | let id_c = RecordId::from_bytes([0x03; 32]); |
| 98 | let rec_c = Record::new(id_c, "escrow", b"gamma-with-a-longer-payload".to_vec()); |
| 99 | res!(storage.put(&rec_c)); |
| 100 | |
| 101 | thread::sleep(Duration::from_millis(500)); |
| 102 | |
| 103 | // 2. Fetch each record back and verify round-trip equality. |
| 104 | let fetched_a = res!(storage.get("identity", &id_a)); |
| 105 | assert_eq!(fetched_a.as_ref(), Some(&rec_a), |
| 106 | "identity record A did not round-trip"); |
| 107 | |
| 108 | let fetched_b = res!(storage.get("identity", &id_b)); |
| 109 | assert_eq!(fetched_b.as_ref(), Some(&rec_b), |
| 110 | "identity record B did not round-trip"); |
| 111 | |
| 112 | let fetched_c = res!(storage.get("escrow", &id_c)); |
| 113 | assert_eq!(fetched_c.as_ref(), Some(&rec_c), |
| 114 | "escrow record C did not round-trip"); |
| 115 | |
| 116 | // 3. A missing id returns None rather than erroring. |
| 117 | let missing_id = RecordId::from_bytes([0x99; 32]); |
| 118 | let fetched_missing = res!(storage.get("identity", &missing_id)); |
| 119 | assert!(fetched_missing.is_none(), |
| 120 | "missing record unexpectedly returned something"); |
| 121 | |
| 122 | // 4. digests("identity") enumerates exactly A and B. |
| 123 | let digests_identity = res!(storage.digests("identity")); |
| 124 | assert_eq!(digests_identity.len(), 2, |
| 125 | "expected 2 identity digests, got {}", digests_identity.len()); |
| 126 | let mut identity_ids: Vec<_> = digests_identity.iter() |
| 127 | .map(|d| d.id) |
| 128 | .collect(); |
| 129 | identity_ids.sort_by(|a, b| a.as_bytes().cmp(b.as_bytes())); |
| 130 | assert_eq!(identity_ids, vec![id_a, id_b]); |
| 131 | |
| 132 | // 5. digests("escrow") enumerates only C. |
| 133 | let digests_escrow = res!(storage.digests("escrow")); |
| 134 | assert_eq!(digests_escrow.len(), 1); |
| 135 | assert_eq!(digests_escrow[0].id, id_c); |
| 136 | |
| 137 | // 6. Content hashes differ between records with different payloads. |
| 138 | let digest_a = digests_identity.iter().find(|d| d.id == id_a) |
| 139 | .expect("digest for A"); |
| 140 | let digest_b = digests_identity.iter().find(|d| d.id == id_b) |
| 141 | .expect("digest for B"); |
| 142 | assert_ne!(digest_a.content, digest_b.content, |
| 143 | "different payloads should hash differently"); |
| 144 | |
| 145 | // 7. delete reports `true` for a present record and `false` for a |
| 146 | // subsequent attempt on the same key. |
| 147 | let was_present = res!(storage.delete("identity", &id_a)); |
| 148 | assert!(was_present, "first delete of A should return true"); |
| 149 | |
| 150 | // Give the delete tombstone time to land in cache / on disk. |
| 151 | thread::sleep(Duration::from_secs(2)); |
| 152 | |
| 153 | let after_delete = res!(storage.get("identity", &id_a)); |
| 154 | assert!(after_delete.is_none(), |
| 155 | "A should be gone after delete, but got {:?}", after_delete); |
| 156 | |
| 157 | let was_present_again = res!(storage.delete("identity", &id_a)); |
| 158 | assert!(!was_present_again, |
| 159 | "second delete of A should return false"); |
| 160 | |
| 161 | // 8. digests("identity") is now just B. |
| 162 | thread::sleep(Duration::from_millis(500)); |
| 163 | let digests_identity_after = res!(storage.digests("identity")); |
| 164 | assert_eq!(digests_identity_after.len(), 1); |
| 165 | assert_eq!(digests_identity_after[0].id, id_b); |
| 166 | |
| 167 | // 9. Graceful shutdown. Drop the adapter so the Arc refcount drops |
| 168 | // back to one, then reclaim the db and shut it down. Skipping |
| 169 | // shutdown leaves bot threads spinning and pollutes sibling |
| 170 | // tests' log streams. |
| 171 | drop(storage); |
| 172 | let db = match Arc::try_unwrap(db) { |
| 173 | Ok(db) => db, |
| 174 | Err(_) => return Err(err!( |
| 175 | "Could not reclaim db from Arc: extra references outlived storage."; |
| 176 | Bug)), |
| 177 | }; |
| 178 | res!(db.shutdown()); |
| 179 | |
| 180 | Ok(()) |
| 181 | } |