1//! The native adapters held to the same laws as the in-memory one (`whiskers-conformance`), plus what 2//! only real threads and real files can show: that the atomicity the contracts promise holds when 3//! requests really do arrive at once, and that a restart really does find the state again. 4 5use std::path::{Path, PathBuf}; 6use std::sync::Arc; 7use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; 8 9use whiskers_conformance::suite::fixture::*; 10use whiskers_conformance::suite::{self}; 11use whiskers_core::{ChatError, ChatState, ChatStore, JsonChat, TokenLimit}; 12use whiskers_ports::{ 13 Allowance, Hub, Millis, ReserveError, SafePictureId, SecretName, Settlement, PictureShelf, Put, run_ready, 14}; 15use whiskersd::allowance::FileAllowance; 16use whiskersd::clock::SystemClock; 17use whiskersd::hubs::{FileChatHub, FileHouseholdHub, FileMemoryHub}; 18use whiskersd::journal::FileJournal; 19use whiskersd::judge::JevJudge; 20use whiskersd::pictures::DirShelf; 21use whiskersd::secrets::EnvSecrets; 22use whiskersd::speak::ElevenLabsVoice; 23 24static COUNTER: AtomicUsize = AtomicUsize::new(0); 25 26/// A directory of its own, empty. 27fn scratch(what: &str) -> PathBuf { 28 let d = std::env::temp_dir().join(format!("whiskers-conformance-{what}-{}-{}", std::process::id(), COUNTER.fetch_add(1, Ordering::SeqCst))); 29 let _ = std::fs::remove_dir_all(&d); 30 std::fs::create_dir_all(&d).unwrap(); 31 d 32} 33 34fn memory_hub_at(dir: &Path) -> FileMemoryHub { 35 FileMemoryHub::open(&dir.join("memory-hub.json")).expect("the memory hub opens") 36} 37 38fn allowance_at(dir: &Path) -> FileAllowance { 39 FileAllowance::open(Some(dir.join("tokens.json")), Some(dir.join("speak.json"))) 40} 41 42/// The native adapters' harness: a fixture for every capability there is a law for, each on a directory of 43/// its own, a restart being the same files opened again. 44struct NativeFixtures; 45 46impl MemoryFixture for NativeFixtures { 47 type Port = FileMemoryHub; 48 fn fresh(&self) -> FileMemoryHub { 49 memory_hub_at(&scratch("memory")) 50 } 51 fn restart(&self, hub: FileMemoryHub) -> FileMemoryHub { 52 memory_hub_at(hub.path().parent().unwrap()) 53 } 54} 55 56impl HouseholdFixture for NativeFixtures { 57 type Port = FileHouseholdHub; 58 fn fresh(&self) -> FileHouseholdHub { 59 FileHouseholdHub::open(&scratch("household").join("h.json")) 60 } 61 fn restart(&self, hub: FileHouseholdHub) -> FileHouseholdHub { 62 FileHouseholdHub::open(hub.path()) 63 } 64} 65 66impl ChatFixture for NativeFixtures { 67 type Port = FileChatHub; 68 fn fresh(&self) -> FileChatHub { 69 FileChatHub::open(&scratch("chat").join("c.json")) 70 } 71 fn restart(&self, hub: FileChatHub) -> FileChatHub { 72 FileChatHub::open(hub.path()) 73 } 74} 75 76/// The same hub with a store that can be told to refuse writes, as a full disk would. 77impl ChatFaultFixture for NativeFixtures { 78 type Port = FileChatHub<RefusingStore>; 79 fn fresh_refusing(&self) -> (FileChatHub<RefusingStore>, suite::hubs::WriteSwitch) { 80 let path = scratch("chat-broken").join("c.json"); 81 let refusing = Arc::new(AtomicBool::new(false)); 82 let store = RefusingStore { inner: JsonChat::open(&path), refusing: refusing.clone() }; 83 let hub = FileChatHub::with_store(&path, ChatState::default(), store); 84 (hub, Box::new(move |refuse| refusing.store(refuse, Ordering::SeqCst))) 85 } 86 /// A restart reads the file again; the storage that comes back works. 87 fn restart_refusing(&self, hub: FileChatHub<RefusingStore>) -> FileChatHub<RefusingStore> { 88 let store = RefusingStore { inner: JsonChat::open(hub.path()), refusing: Arc::new(AtomicBool::new(false)) }; 89 let state = store.load().expect("the file the hub kept can be read"); 90 FileChatHub::with_store(hub.path(), state, store) 91 } 92} 93 94impl JournalFixture for NativeFixtures { 95 type Port = FileJournal; 96 fn fresh(&self) -> FileJournal { 97 FileJournal::open(&scratch("journal")) 98 } 99 fn restart(&self, journal: FileJournal) -> FileJournal { 100 FileJournal::open(journal.dir()) 101 } 102} 103 104impl PicturesFixture for NativeFixtures { 105 type Port = DirShelf; 106 fn fresh(&self) -> DirShelf { 107 DirShelf::open(&scratch("pictures")) 108 } 109 fn restart(&self, shelf: DirShelf) -> DirShelf { 110 DirShelf::open(shelf.dir()) 111 } 112} 113 114impl AllowanceFixture for NativeFixtures { 115 type Port = FileAllowance; 116 fn fresh(&self) -> FileAllowance { 117 allowance_at(&scratch("allowance")) 118 } 119 fn restart(&self, allowance: FileAllowance) -> FileAllowance { 120 allowance.reopened() 121 } 122} 123 124fn no_secrets() -> EnvSecrets { 125 EnvSecrets::from_fn(|_| None) 126} 127 128impl JudgeFixture for NativeFixtures { 129 type Port = JevJudge; 130 fn without_credentials(&self) -> JevJudge { 131 JevJudge::start(&no_secrets()) 132 } 133} 134 135impl VoiceFixture for NativeFixtures { 136 type Port = ElevenLabsVoice<EnvSecrets>; 137 fn without_credentials(&self) -> ElevenLabsVoice<EnvSecrets> { 138 ElevenLabsVoice::new(no_secrets()) 139 } 140} 141 142impl ClockFixture for NativeFixtures { 143 type Port = SystemClock; 144 fn clock(&self) -> SystemClock { 145 SystemClock 146 } 147} 148 149/// The native adapters are held to every law there is: if a capability gains a law and a fixture trait, this 150/// stops compiling until the native harness has the fixture. 151const _EVERY_LAW_IS_RUN: fn(NativeFixtures) -> Suite<NativeFixtures> = Suite::complete; 152 153#[test] 154fn the_file_memory_hub_obeys_its_contract() { 155 Suite::new(NativeFixtures).memory(); 156} 157 158#[test] 159fn the_file_household_hub_obeys_its_contract() { 160 Suite::new(NativeFixtures).household(); 161} 162 163#[test] 164fn the_file_chat_hub_obeys_its_contract() { 165 Suite::new(NativeFixtures).chat().chat_write_faults(); 166} 167 168/// A chat store that is the real file until told to refuse writes. 169struct RefusingStore { 170 inner: JsonChat, 171 refusing: Arc<AtomicBool>, 172} 173 174impl ChatStore for RefusingStore { 175 fn load(&self) -> Result<ChatState, ChatError> { 176 self.inner.load() 177 } 178 179 fn save(&mut self, state: &ChatState) -> Result<(), ChatError> { 180 if self.refusing.load(Ordering::SeqCst) { 181 return Err(ChatError("no space left on device".into())); 182 } 183 self.inner.save(state) 184 } 185} 186 187#[test] 188fn the_file_journal_obeys_its_contract() { 189 Suite::new(NativeFixtures).journal(); 190} 191 192#[test] 193fn the_directory_of_pictures_obeys_its_contract() { 194 Suite::new(NativeFixtures).pictures(); 195} 196 197#[test] 198fn the_file_allowance_obeys_its_contract() { 199 Suite::new(NativeFixtures).allowance(); 200} 201 202#[test] 203fn environment_secrets_tell_absent_from_empty() { 204 suite::secrets::secrets(&|entries: &[(SecretName, &str)]| { 205 let set: Vec<(String, String)> = entries.iter().map(|(n, v)| (n.operator_name().to_owned(), (*v).to_owned())).collect(); 206 EnvSecrets::from_fn(move |name| set.iter().find(|(n, _)| n == name).map(|(_, v)| v.clone())) 207 }); 208} 209 210#[test] 211fn the_system_clock_reads_a_plausible_time() { 212 Suite::new(NativeFixtures).clock(); 213} 214 215#[test] 216fn without_credentials_the_judge_and_the_voice_say_so() { 217 Suite::new(NativeFixtures).judge().voice(); 218} 219 220// What only real threads show. 221 222fn at(ms: u64) -> Millis { 223 Millis::new(ms) 224} 225 226#[test] 227fn eight_lines_asked_at_once_cannot_pass_a_cap_with_room_for_one() { 228 let a = Arc::new(allowance_at(&scratch("race-voice"))); 229 let threads: Vec<_> = (0..8) 230 .map(|_| { 231 let a = a.clone(); 232 std::thread::spawn(move || run_ready(a.reserve_voice(at(1000), 6, 10)).map(|r| run_ready(a.settle_voice(r, Settlement::Spoken)).unwrap())) 233 }) 234 .collect(); 235 let results: Vec<_> = threads.into_iter().map(|t| t.join().unwrap()).collect(); 236 assert_eq!(results.iter().filter(|r| r.is_ok()).count(), 1, "exactly one line fits"); 237 assert!(results.iter().filter_map(|r| r.as_ref().err()).all(|e| matches!(e, ReserveError::DailyCap { .. }))); 238 assert_eq!(run_ready(a.voice_spend(at(1001))).unwrap().spent_today, 6); 239} 240 241#[test] 242fn merges_that_arrive_together_all_land_and_no_fact_is_numbered_twice() { 243 let hub = Arc::new(memory_hub_at(&scratch("race-memory"))); 244 let threads: Vec<_> = (0..8) 245 .map(|i| { 246 let hub = hub.clone(); 247 std::thread::spawn(move || { 248 let doc = suite::build::snapshot(vec![suite::build::fact(&format!("g{i}"), &format!("fact number {i}"))], &[]); 249 run_ready(hub.merge(doc)).unwrap(); 250 }) 251 }) 252 .collect(); 253 for t in threads { 254 t.join().unwrap(); 255 } 256 let doc = run_ready(hub.current()).unwrap(); 257 assert_eq!(doc.facts.len(), 8, "no merge lost another's work"); 258 let ids: std::collections::BTreeSet<u64> = doc.facts.iter().map(|f| f.id).collect(); 259 assert_eq!(ids.len(), 8, "and no id was handed out twice"); 260} 261 262#[test] 263fn eight_puts_of_one_picture_store_it_once_and_it_is_never_seen_half_written() { 264 let shelf = Arc::new(DirShelf::open(&scratch("race-pictures"))); 265 let id = SafePictureId::parse("a.jpg").unwrap(); 266 let body = vec![7u8; 1 << 20]; 267 let threads: Vec<_> = (0..8) 268 .map(|_| { 269 let (shelf, id, body) = (shelf.clone(), id.clone(), body.clone()); 270 std::thread::spawn(move || { 271 let put = run_ready(shelf.put(&id, &body)).unwrap(); 272 // Whatever another thread is doing, a picture that is there is all there. 273 if let Some(bytes) = run_ready(shelf.get(&id)).unwrap() { 274 assert_eq!(bytes.len(), body.len()); 275 } 276 put 277 }) 278 }) 279 .collect(); 280 let puts: Vec<Put> = threads.into_iter().map(|t| t.join().unwrap()).collect(); 281 assert_eq!(puts.iter().filter(|p| **p == Put::Stored).count(), 1); 282 assert_eq!(puts.iter().filter(|p| **p == Put::AlreadyKept).count(), 7); 283} 284 285#[test] 286fn the_ledger_and_the_spend_are_found_again_by_the_next_process() { 287 let dir = scratch("restart-files"); 288 let a = allowance_at(&dir); 289 run_ready(a.charge_thinking(at(10), whiskers_ports::Tokens::new(700))).unwrap(); 290 let r = run_ready(a.reserve_voice(at(10), 4, 10)).unwrap(); 291 run_ready(a.settle_voice(r, Settlement::Spoken)).unwrap(); 292 assert!(std::fs::read_to_string(dir.join("tokens.json")).unwrap().contains("700")); 293 assert!(std::fs::read_to_string(dir.join("speak.json")).unwrap().contains("\"spent\":4")); 294 let again = allowance_at(&dir); 295 let limit = TokenLimit::new(Some(1000), 5).unwrap(); 296 assert_eq!(run_ready(again.thinking_usage(at(20), &limit)).unwrap().used.get(), 700); 297 assert_eq!(run_ready(again.voice_spend(at(20))).unwrap().spent_today, 4); 298} 299 300#[test] 301fn a_damaged_ledger_starts_empty_rather_than_stopping_the_service() { 302 let dir = scratch("damaged"); 303 std::fs::write(dir.join("tokens.json"), "{ nope").unwrap(); 304 let a = allowance_at(&dir); 305 let limit = TokenLimit::new(Some(1000), 5).unwrap(); 306 assert_eq!(run_ready(a.thinking_usage(at(20), &limit)).unwrap().used.get(), 0); 307}