spend.rsannotatedspend.rssource186 lines · 8.3 KB · raw
1//! What the service has spent, in `thinking_spend` and `voice_spend`. The rules (the rolling window, the day
2//! turning, reserve and refund) are `whiskers-ports`'; this keeps them behind one lock, so every call is atomic
3//! with respect to every other, and writes them down.
4//!
5//! The thinking ledger is saved when a question is charged. The voice spend is saved when a line is settled,
6//! not when it is reserved: a reservation is a promise about a line that has not been spoken yet, and a crash
7//! between the two forgets it (the spend is then a line short, never a line over).
8
9use std::sync::{Mutex, MutexGuard, PoisonError};
10
11use ::log::{debug, error, info, warn};
12use rusqlite::params;
13use whiskers_core::TokenLimit;
14use whiskers_ports::{
15    Allowance, Day, Millis, ReserveError, Settlement, SpeakError, StoreError, ThinkingLedger, ThinkingRoom, ThinkingUsage, Tokens, VoiceDay,
16    VoiceReservation, VoiceSaved, VoiceSpend,
17};
18
19use crate::{DbError, Store};
20
21struct State {
22    thinking: ThinkingLedger,
23    voice: VoiceDay,
24    saved: VoiceSaved,
25    /// Why the last attempt to speak failed, until one succeeds. Not kept across a restart.
26    last_failure: Option<SpeakError>,
27}
28
29pub struct SqliteAllowance {
30    store: Store,
31    state: Mutex<State>,
32}
33
34/// The ledger and the voice day are the ports' types, whose fields are theirs; they are moved through their
35/// JSON, which is what each has always been, so that no rule is written twice.
36fn load_thinking(store: &Store) -> Result<ThinkingLedger, DbError> {
37    let conn = store.lock();
38    let mut st = conn.prepare("SELECT at_ms, tokens FROM thinking_spend ORDER BY at_ms, rowid")?;
39    let spent: Vec<(u64, u64)> = st.query_map([], |r| Ok((r.get::<_, i64>(0)? as u64, r.get::<_, i64>(1)? as u64)))?.collect::<Result<_, _>>()?;
40    serde_json::from_value(serde_json::json!({ "spent": spent })).map_err(|e| DbError::Damaged(format!("the thinking ledger: {e}")))
41}
42
43fn load_saved(store: &Store) -> Result<VoiceSaved, DbError> {
44    let conn = store.lock();
45    let latest: Option<(i64, i64, i64)> = rusqlite::OptionalExtension::optional(
46        conn.query_row("SELECT day, hits, saved_chars FROM voice_spend ORDER BY day DESC LIMIT 1", [], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?))),
47    )?;
48    let (day, hits, chars) = latest.unwrap_or((0, 0, 0));
49    let day: Day = serde_json::from_value(serde_json::json!(day)).map_err(|e| DbError::Damaged(format!("the voice day: {e}")))?;
50    Ok(VoiceSaved::restore(day, hits as u32, chars as u32))
51}
52
53fn load_voice(store: &Store) -> Result<VoiceDay, DbError> {
54    let conn = store.lock();
55    let latest: Option<(i64, i64)> =
56        rusqlite::OptionalExtension::optional(conn.query_row("SELECT day, chars FROM voice_spend ORDER BY day DESC LIMIT 1", [], |r| Ok((r.get(0)?, r.get(1)?))))?;
57    let (day, spent) = latest.unwrap_or((0, 0));
58    serde_json::from_value(serde_json::json!({ "day": day, "spent": spent })).map_err(|e| DbError::Damaged(format!("the voice spend: {e}")))
59}
60
61impl SqliteAllowance {
62    /// A ledger that cannot be read starts empty (the allowance is then fresh, and the log says so).
63    pub fn open(store: Store) -> Self {
64        let thinking = load_thinking(&store).unwrap_or_else(|e| {
65            warn!("the token ledger cannot be read ({e}); starting empty");
66            ThinkingLedger::default()
67        });
68        let voice = load_voice(&store).unwrap_or_else(|e| {
69            warn!("the voice spend cannot be read ({e}); starting empty");
70            VoiceDay::default()
71        });
72        let saved = load_saved(&store).unwrap_or_else(|e| {
73            warn!("the speech cache's savings cannot be read ({e}); starting empty");
74            VoiceSaved::default()
75        });
76        debug!("token ledger opened: {} entries", thinking.entries());
77        info!("voice: {} chars already spent on day {}", voice.spent_on(voice.day()), voice.day().number());
78        Self { store, state: Mutex::new(State { thinking, voice, saved, last_failure: None }) }
79    }
80
81    fn locked(&self) -> MutexGuard<'_, State> {
82        self.state.lock().unwrap_or_else(PoisonError::into_inner)
83    }
84
85    fn save_thinking(&self, ledger: &ThinkingLedger) -> Result<(), StoreError> {
86        let json = serde_json::to_value(ledger).map_err(|e| {
87            error!("token ledger does not serialise: {e}");
88            StoreError::Unavailable
89        })?;
90        let spent: Vec<(u64, u64)> = serde_json::from_value(json["spent"].clone()).map_err(|e| {
91            error!("token ledger has an unexpected shape: {e}");
92            StoreError::Unavailable
93        })?;
94        self.store
95            .transaction(|tx| {
96                tx.execute("DELETE FROM thinking_spend", [])?;
97                for (at, tokens) in &spent {
98                    tx.execute("INSERT INTO thinking_spend (at_ms, tokens) VALUES (?1, ?2)", params![*at as i64, *tokens as i64])?;
99                }
100                Ok(())
101            })
102            .map_err(|e| {
103                warn!("token ledger not saved ({e}); a restart may forget this spending");
104                StoreError::Unavailable
105            })
106    }
107
108    fn save_voice(&self, voice: &VoiceDay) -> Result<(), StoreError> {
109        let json = serde_json::to_value(voice).map_err(|_| StoreError::Unavailable)?;
110        let (day, spent) = (json["day"].as_i64().unwrap_or(0), json["spent"].as_i64().unwrap_or(0));
111        self.store
112            .transaction(|tx| {
113                tx.execute(
114                    "INSERT INTO voice_spend (day, chars) VALUES (?1, ?2) ON CONFLICT (day) DO UPDATE SET chars = excluded.chars",
115                    params![day, spent],
116                )?;
117                Ok(())
118            })
119            .map_err(|e| {
120                warn!("voice spend not saved ({e}); a restart may forget this spending");
121                StoreError::Unavailable
122            })
123    }
124}
125
126impl Allowance for SqliteAllowance {
127    async fn thinking_room(&self, now: Millis, limit: &TokenLimit) -> Result<ThinkingRoom, StoreError> {
128        Ok(self.locked().thinking.room(now, limit))
129    }
130
131    async fn charge_thinking(&self, now: Millis, cost: Tokens) -> Result<(), StoreError> {
132        let mut state = self.locked();
133        state.thinking.charge(now, cost);
134        self.save_thinking(&state.thinking)
135    }
136
137    async fn thinking_usage(&self, now: Millis, limit: &TokenLimit) -> Result<ThinkingUsage, StoreError> {
138        Ok(self.locked().thinking.usage(now, limit))
139    }
140
141    async fn reserve_voice(&self, now: Millis, chars: u32, cap: u32) -> Result<VoiceReservation, ReserveError> {
142        self.locked().voice.reserve(now, chars, cap)
143    }
144
145    async fn settle_voice(&self, reservation: VoiceReservation, how: Settlement) -> Result<(), StoreError> {
146        let mut state = self.locked();
147        match how {
148            Settlement::Spoken => {
149                state.last_failure = None;
150                self.save_voice(&state.voice)
151            }
152            Settlement::NotSpoken(why) => {
153                state.voice.release(reservation);
154                state.last_failure = Some(why);
155                Ok(())
156            }
157        }
158    }
159
160    async fn voice_spend(&self, now: Millis) -> Result<VoiceSpend, StoreError> {
161        let state = self.locked();
162        let today = Day::containing(now);
163        let (hits_today, saved_today) = state.saved.on(today);
164        Ok(VoiceSpend { spent_today: state.voice.spent_on(today), hits_today, saved_today, last_failure: state.last_failure.clone() })
165    }
166
167    async fn record_voice_hit(&self, now: Millis, chars: u32) -> Result<(), StoreError> {
168        let mut state = self.locked();
169        state.saved.record(now, chars);
170        let json = serde_json::to_value(state.saved).map_err(|_| StoreError::Unavailable)?;
171        let (day, hits, saved) = (json["day"].as_i64().unwrap_or(0), json["hits"].as_i64().unwrap_or(0), json["chars"].as_i64().unwrap_or(0));
172        self.store
173            .transaction(|tx| {
174                tx.execute(
175                    "INSERT INTO voice_spend (day, chars, hits, saved_chars) VALUES (?1, 0, ?2, ?3)
176                     ON CONFLICT (day) DO UPDATE SET hits = excluded.hits, saved_chars = excluded.saved_chars",
177                    params![day, hits, saved],
178                )?;
179                Ok(())
180            })
181            .map_err(|e| {
182                warn!("speech cache savings not saved ({e})");
183                StoreError::Unavailable
184            })
185    }
186}