What the service has spent, in thinking_spend and voice_spend. The rules (the rolling window, the day
turning, reserve and refund) are whiskers-ports'; this keeps them behind one lock, so every call is atomic
with respect to every other, and writes them down.
The thinking ledger is saved when a question is charged. The voice spend is saved when a line is settled, not when it is reserved: a reservation is a promise about a line that has not been spoken yet, and a crash between the two forgets it (the spend is then a line short, never a line over).
9use std::sync::{Mutex, MutexGuard, PoisonError};
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,
Why the last attempt to speak failed, until one succeeds. Not kept across a restart.
The ledger and the voice day are the ports' types, whose fields are theirs; they are moved through their 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}
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 {
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 }
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}