mod.rsannotatedmod.rssource244 lines · 9.1 KB · raw

The fixtures that hold the SQLite adapters to whiskers-conformance's laws: each on a database of its own, a restart being the same file opened again. Shared by this crate's tests and by whiskersd's, which includes this file by path, so the adapters are held to the laws in one place.

5#![allow(dead_code)]
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}

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}
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}

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}
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}

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}
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}