mod.rsannotatedmod.rssource244 lines · 9.1 KB · raw
1//! The fixtures that hold the SQLite adapters to `whiskers-conformance`'s laws: each on a database of its own, a
2//! restart being the same file opened again. Shared by this crate's tests and by `whiskersd`'s, which includes this
3//! file by path, so the adapters are held to the laws in one place.
4
5#![allow(dead_code)]
6
7use std::path::PathBuf;
8use std::sync::Arc;
9use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
10
11use whiskers_conformance::suite::fixture::*;
12use whiskers_conformance::suite::hubs::WriteSwitch;
13use whiskers_core::{ChatError, ChatState, ChatStore, TokenLimit};
14use whiskers_ports::{
15    Allowance, Appended, DeviceCursor, DeviceName, EntryLine, Hub, Journal, LogCursor, Merged, Millis, Page, PictureShelf, Put, ReserveError, SafePictureId,
16    Settlement, StoreError,
17};
18use whiskers_store::{SqliteAllowance, SqliteChat, SqliteChatHub, SqliteHouseholdHub, SqliteJournal, SqliteMemoryHub, SqlitePictureShelf, Store};
19
20static COUNTER: AtomicUsize = AtomicUsize::new(0);
21
22pub fn scratch(what: &str) -> PathBuf {
23    let d = std::env::temp_dir().join(format!("whiskers-store-conformance-{what}-{}-{}", std::process::id(), COUNTER.fetch_add(1, Ordering::SeqCst)));
24    let _ = std::fs::remove_dir_all(&d);
25    std::fs::create_dir_all(&d).unwrap();
26    d.join("w.db")
27}
28
29/// An adapter and the database it keeps its state in, so a restart can open the same file again.
30pub struct On<T> {
31    pub port: T,
32    pub store: Store,
33}
34
35impl<H: Hub> Hub for On<H> {
36    type Document = H::Document;
37    async fn merge(&self, theirs: H::Document) -> Result<Merged<H::Document>, StoreError> {
38        self.port.merge(theirs).await
39    }
40    async fn current(&self) -> Result<H::Document, StoreError> {
41        self.port.current().await
42    }
43}
44
45impl<J: Journal> Journal for On<J> {
46    async fn held(&self, d: &DeviceName) -> Result<DeviceCursor, StoreError> {
47        self.port.held(d).await
48    }
49    async fn append(&self, d: &DeviceName, at: DeviceCursor, lines: &[EntryLine]) -> Result<Appended, StoreError> {
50        self.port.append(d, at, lines).await
51    }
52    async fn pull(&self, since: LogCursor) -> Result<Page, StoreError> {
53        self.port.pull(since).await
54    }
55}
56
57impl<S: PictureShelf> PictureShelf for On<S> {
58    async fn has(&self, id: &SafePictureId) -> Result<bool, StoreError> {
59        self.port.has(id).await
60    }
61    async fn put(&self, id: &SafePictureId, bytes: &[u8]) -> Result<Put, StoreError> {
62        self.port.put(id, bytes).await
63    }
64    async fn get(&self, id: &SafePictureId) -> Result<Option<Vec<u8>>, StoreError> {
65        self.port.get(id).await
66    }
67}
68
69impl<A: Allowance> Allowance for On<A> {
70    async fn thinking_room(&self, now: Millis, l: &TokenLimit) -> Result<whiskers_ports::ThinkingRoom, StoreError> {
71        self.port.thinking_room(now, l).await
72    }
73    async fn charge_thinking(&self, now: Millis, c: whiskers_ports::Tokens) -> Result<(), StoreError> {
74        self.port.charge_thinking(now, c).await
75    }
76    async fn thinking_usage(&self, now: Millis, l: &TokenLimit) -> Result<whiskers_ports::ThinkingUsage, StoreError> {
77        self.port.thinking_usage(now, l).await
78    }
79    async fn reserve_voice(&self, now: Millis, chars: u32, cap: u32) -> Result<whiskers_ports::VoiceReservation, ReserveError> {
80        self.port.reserve_voice(now, chars, cap).await
81    }
82    async fn settle_voice(&self, r: whiskers_ports::VoiceReservation, how: Settlement) -> Result<(), StoreError> {
83        self.port.settle_voice(r, how).await
84    }
85    async fn voice_spend(&self, now: Millis) -> Result<whiskers_ports::VoiceSpend, StoreError> {
86        self.port.voice_spend(now).await
87    }
88    async fn record_voice_hit(&self, now: Millis, chars: u32) -> Result<(), StoreError> {
89        self.port.record_voice_hit(now, chars).await
90    }
91}
92
93/// A speech cache, its database and its cap, so a restart reopens the same file over the same cap.
94pub struct CacheOn {
95    pub port: whiskers_store::SqliteSpeechCache,
96    pub store: Store,
97    pub cap: u64,
98}
99
100impl whiskers_ports::SpeechCache for CacheOn {
101    async fn find(&self, v: &whiskers_ports::VoiceId, l: &whiskers_ports::SpeechLine, now: Millis) -> Result<Option<whiskers_ports::TimedAudio>, StoreError> {
102        self.port.find(v, l, now).await
103    }
104    async fn keep(&self, v: &whiskers_ports::VoiceId, l: &whiskers_ports::SpeechLine, s: &whiskers_ports::TimedAudio, now: Millis) -> Result<(), StoreError> {
105        self.port.keep(v, l, s, now).await
106    }
107    async fn size(&self) -> Result<whiskers_ports::CacheSize, StoreError> {
108        self.port.size().await
109    }
110}
111
112impl SpeechCacheFixture for SqliteFixtures {
113    type Port = CacheOn;
114    fn fresh(&self, cap: u64) -> CacheOn {
115        let store = Store::open(&scratch("speech")).unwrap();
116        CacheOn { port: whiskers_store::SqliteSpeechCache::new(store.clone(), cap), store, cap }
117    }
118    fn restart(&self, p: CacheOn) -> CacheOn {
119        let store = reopened(&p.store);
120        let cap = p.cap;
121        drop(p);
122        CacheOn { port: whiskers_store::SqliteSpeechCache::new(store.clone(), cap), store, cap }
123    }
124}
125
126pub struct SqliteFixtures;
127
128fn reopened(on: &Store) -> Store {
129    Store::open(on.path().expect("a database with a file")).unwrap()
130}
131
132impl MemoryFixture for SqliteFixtures {
133    type Port = On<SqliteMemoryHub>;
134    fn fresh(&self) -> Self::Port {
135        let store = Store::open(&scratch("memory")).unwrap();
136        On { port: SqliteMemoryHub::open(store.clone()).unwrap(), store }
137    }
138    fn restart(&self, p: Self::Port) -> Self::Port {
139        let store = reopened(&p.store);
140        drop(p);
141        On { port: SqliteMemoryHub::open(store.clone()).unwrap(), store }
142    }
143}
144
145impl HouseholdFixture for SqliteFixtures {
146    type Port = On<SqliteHouseholdHub>;
147    fn fresh(&self) -> Self::Port {
148        let store = Store::open(&scratch("household")).unwrap();
149        On { port: SqliteHouseholdHub::open(store.clone()), store }
150    }
151    fn restart(&self, p: Self::Port) -> Self::Port {
152        let store = reopened(&p.store);
153        drop(p);
154        On { port: SqliteHouseholdHub::open(store.clone()), store }
155    }
156}
157
158impl ChatFixture for SqliteFixtures {
159    type Port = On<SqliteChatHub>;
160    fn fresh(&self) -> Self::Port {
161        let store = Store::open(&scratch("chat")).unwrap();
162        On { port: SqliteChatHub::open(store.clone()), store }
163    }
164    fn restart(&self, p: Self::Port) -> Self::Port {
165        let store = reopened(&p.store);
166        drop(p);
167        On { port: SqliteChatHub::open(store.clone()), store }
168    }
169}
170
171/// A chat store that is the real database until told to refuse writes, as a full disk would.
172pub struct RefusingStore {
173    inner: SqliteChat,
174    refusing: Arc<AtomicBool>,
175}
176
177impl ChatStore for RefusingStore {
178    fn load(&self) -> Result<ChatState, ChatError> {
179        self.inner.load()
180    }
181    fn save(&mut self, state: &ChatState) -> Result<(), ChatError> {
182        if self.refusing.load(Ordering::SeqCst) {
183            return Err(ChatError("no space left on device".into()));
184        }
185        self.inner.save(state)
186    }
187}
188
189impl ChatFaultFixture for SqliteFixtures {
190    type Port = On<SqliteChatHub<RefusingStore>>;
191    fn fresh_refusing(&self) -> (Self::Port, WriteSwitch) {
192        let store = Store::open(&scratch("chat-broken")).unwrap();
193        let refusing = Arc::new(AtomicBool::new(false));
194        let hub = SqliteChatHub::with_store(ChatState::default(), RefusingStore { inner: SqliteChat::new(store.clone()), refusing: refusing.clone() });
195        (On { port: hub, store }, Box::new(move |refuse| refusing.store(refuse, Ordering::SeqCst)))
196    }
197    fn restart_refusing(&self, p: Self::Port) -> Self::Port {
198        let store = reopened(&p.store);
199        drop(p);
200        let chat = SqliteChat::new(store.clone());
201        let state = chat.load().expect("what the hub kept can be read");
202        On { port: SqliteChatHub::with_store(state, RefusingStore { inner: chat, refusing: Arc::new(AtomicBool::new(false)) }), store }
203    }
204}
205
206impl JournalFixture for SqliteFixtures {
207    type Port = On<SqliteJournal>;
208    fn fresh(&self) -> Self::Port {
209        let store = Store::open(&scratch("journal")).unwrap();
210        On { port: SqliteJournal::new(store.clone()), store }
211    }
212    fn restart(&self, p: Self::Port) -> Self::Port {
213        let store = reopened(&p.store);
214        drop(p);
215        On { port: SqliteJournal::new(store.clone()), store }
216    }
217}
218
219impl PicturesFixture for SqliteFixtures {
220    type Port = On<SqlitePictureShelf>;
221    fn fresh(&self) -> Self::Port {
222        let store = Store::open(&scratch("pictures")).unwrap();
223        On { port: SqlitePictureShelf::new(store.clone()), store }
224    }
225    fn restart(&self, p: Self::Port) -> Self::Port {
226        let store = reopened(&p.store);
227        drop(p);
228        On { port: SqlitePictureShelf::new(store.clone()), store }
229    }
230}
231
232impl AllowanceFixture for SqliteFixtures {
233    type Port = On<SqliteAllowance>;
234    fn fresh(&self) -> Self::Port {
235        let store = Store::open(&scratch("allowance")).unwrap();
236        On { port: SqliteAllowance::open(store.clone()), store }
237    }
238    fn restart(&self, p: Self::Port) -> Self::Port {
239        let store = reopened(&p.store);
240        drop(p);
241        On { port: SqliteAllowance::open(store.clone()), store }
242    }
243}
244