whiskers.git / crates / whiskersd / tests / conformance.rs
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}