1//! The allowance: what Whiskers may spend, kept by the hub so one allowance covers every device and 2//! no child's tablet can raise it. 3//! 4//! Two different things are metered, and they are metered differently on purpose: 5//! 6//! - **Thinking** is a rolling window, like a session limit: so many tokens in any `window_hours`. 7//! What a question costs is only known after it is answered, so it is *asked first, charged after*: 8//! the last question may overshoot. That is the documented intent of a rolling allowance, not a race. 9//! - **The natural voice** is a hard daily cap that protects a paid plan. What a line costs is known 10//! before it is spoken (its characters), so the cap is enforced by **reserve, then settle**: the 11//! characters are taken up front, atomically, and given back if the line did not come out. Two lines 12//! asked at once can therefore never both pass a cap with room for one. 13//! 14//! The rules (what rolls off the window, when a day turns, what refunds) are pure structs here 15//! ([`ThinkingLedger`], [`VoiceDay`]) so every backend applies the same ones; an adapter only decides 16//! where the struct is kept and how access to it is serialised. 17 18use ::log::{debug, info, trace, warn}; 19use serde::{Deserialize, Serialize}; 20use whiskers_core::TokenLimit; 21 22use crate::error::StoreError; 23use crate::time::{Day, Millis}; 24use crate::voice::SpeakError; 25 26const HOUR_MS: u64 = 3_600_000; 27 28/// The longest window the parents can choose is a day, but a changed window must find the spending 29/// still there, so spending is kept a week. 30const KEEP_MS: u64 = 7 * 24 * HOUR_MS; 31 32/// A count of the model's tokens. 33#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)] 34#[serde(transparent)] 35pub struct Tokens(u64); 36 37impl Tokens { 38 pub const fn new(n: u64) -> Self { 39 Self(n) 40 } 41 42 pub const fn get(self) -> u64 { 43 self.0 44 } 45} 46 47/// What the parents' screen shows about thinking. 48#[derive(Clone, Copy, Debug, PartialEq, Eq)] 49pub struct ThinkingUsage { 50 pub used: Tokens, 51 pub limit: Option<Tokens>, 52 pub window_hours: u64, 53 /// When the oldest spending leaves the window, so room starts to return; `None` if nothing is spent. 54 pub frees_up_at: Option<Millis>, 55} 56 57/// Whether another question may be asked now. 58#[derive(Clone, Copy, Debug, PartialEq, Eq)] 59pub enum ThinkingRoom { 60 Room, 61 /// The window is spent. `frees_up_at` is when the oldest spending leaves it. (If the window is spent 62 /// something was spent, so there is always such a moment.) 63 Spent { frees_up_at: Millis }, 64} 65 66/// What was spent on thinking and when. How much *may* be spent is the grown-ups' choice, held in the 67/// household document, so it is passed in on every question and can change at any time. The JSON is 68/// the service's token ledger file, unchanged. 69#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] 70pub struct ThinkingLedger { 71 /// (when in ms since the epoch, tokens) 72 spent: Vec<(u64, u64)>, 73} 74 75impl ThinkingLedger { 76 pub fn entries(&self) -> usize { 77 self.spent.len() 78 } 79 80 fn used_in(&self, now: Millis, window_hours: u64) -> u64 { 81 let window = window_hours.saturating_mul(HOUR_MS); 82 self.spent.iter().filter(|(at, _)| at.saturating_add(window) > now.get()).map(|(_, t)| t).sum() 83 } 84 85 fn oldest_leaving(&self, now: Millis, window_hours: u64) -> Option<Millis> { 86 let window = window_hours.saturating_mul(HOUR_MS); 87 self.spent.iter().filter(|(at, _)| at.saturating_add(window) > now.get()).map(|(at, _)| Millis::new(at.saturating_add(window))).min() 88 } 89 90 /// Whether another question may be asked now under `limit`. 91 pub fn room(&self, now: Millis, limit: &TokenLimit) -> ThinkingRoom { 92 let used = self.used_in(now, limit.window_hours()); 93 if limit.per_window().is_none_or(|l| used < l) { 94 trace!("thinking allowance: room ({used} used of {:?})", limit.per_window()); 95 return ThinkingRoom::Room; 96 } 97 info!("thinking allowance: no room ({used} used of {:?} in {} h)", limit.per_window(), limit.window_hours()); 98 match self.oldest_leaving(now, limit.window_hours()) { 99 Some(frees_up_at) => ThinkingRoom::Spent { frees_up_at }, 100 // Unreachable: a spent window has spending in it. Say "now" rather than invent a wait. 101 None => ThinkingRoom::Room, 102 } 103 } 104 105 /// Records what a question cost. 106 pub fn charge(&mut self, now: Millis, cost: Tokens) { 107 let before = self.spent.len(); 108 self.spent.retain(|(at, _)| at.saturating_add(KEEP_MS) > now.get()); 109 if self.spent.len() != before { 110 trace!("thinking ledger pruned {} old entries", before - self.spent.len()); 111 } 112 self.spent.push((now.get(), cost.get())); 113 debug!("thinking allowance: {} tokens charged; {} entries kept", cost.get(), self.spent.len()); 114 } 115 116 pub fn usage(&self, now: Millis, limit: &TokenLimit) -> ThinkingUsage { 117 ThinkingUsage { 118 used: Tokens::new(self.used_in(now, limit.window_hours())), 119 limit: limit.per_window().map(Tokens::new), 120 window_hours: limit.window_hours(), 121 frees_up_at: self.oldest_leaving(now, limit.window_hours()), 122 } 123 } 124} 125 126/// The tokens a Messages API reply cost: what went in plus what came out, cache writes included. 127/// Cache reads are not counted: they are the same prompt read back cheaply, and Anthropic does not 128/// count them against its own limits either. A reply that does not parse costs nothing. 129pub fn tokens_in(reply: &str) -> Tokens { 130 let Ok(v) = serde_json::from_str::<serde_json::Value>(reply) else { 131 warn!("tokens_in: the gateway reply ({} bytes) is not JSON; charging nothing", reply.len()); 132 return Tokens::new(0); 133 }; 134 let u = &v["usage"]; 135 Tokens::new(["input_tokens", "output_tokens", "cache_creation_input_tokens"].iter().filter_map(|k| u[*k].as_u64()).sum()) 136} 137 138/// Characters of natural voice taken from today's allowance for one line, until it is settled. 139/// Only [`VoiceDay::reserve`] makes one, and an adapter must settle it ([`Allowance::settle_voice`]). 140#[derive(Debug, PartialEq, Eq)] 141#[must_use = "a reservation holds characters of today's allowance until it is settled"] 142pub struct VoiceReservation { 143 day: Day, 144 chars: u32, 145} 146 147impl VoiceReservation { 148 pub fn chars(&self) -> u32 { 149 self.chars 150 } 151 152 pub fn day(&self) -> Day { 153 self.day 154 } 155} 156 157/// Why a line may not be spoken. 158#[derive(Clone, Copy, Debug, PartialEq, Eq)] 159pub enum ReserveError { 160 /// Today's characters are spent. `frees_up_at` is when the next day begins. 161 DailyCap { spent: u32, cap: u32, frees_up_at: Millis }, 162 Store(StoreError), 163} 164 165impl From<StoreError> for ReserveError { 166 fn from(e: StoreError) -> Self { 167 ReserveError::Store(e) 168 } 169} 170 171/// How many characters of the natural voice have been spent on which day. The JSON is the service's 172/// voice spend file, unchanged. 173#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] 174pub struct VoiceDay { 175 day: Day, 176 spent: u32, 177} 178 179impl VoiceDay { 180 /// The day the stored spend belongs to. 181 pub fn day(&self) -> Day { 182 self.day 183 } 184 185 /// What has been spent on `today` (a day that has turned since is a fresh allowance). 186 pub fn spent_on(&self, today: Day) -> u32 { 187 if self.day == today { self.spent } else { 0 } 188 } 189 190 /// Takes `chars` from the day containing `now`, or refuses if they would pass `cap`. 191 pub fn reserve(&mut self, now: Millis, chars: u32, cap: u32) -> Result<VoiceReservation, ReserveError> { 192 let today = Day::containing(now); 193 if today != self.day { 194 info!("voice: a new day ({} -> {}); the spend starts from zero (was {})", self.day.number(), today.number(), self.spent); 195 *self = Self { day: today, spent: 0 }; 196 } 197 if self.spent.saturating_add(chars) > cap { 198 info!("speech refused: {chars} chars on top of {} spent would pass today's cap of {cap}", self.spent); 199 return Err(ReserveError::DailyCap { spent: self.spent, cap, frees_up_at: today.ends_at() }); 200 } 201 self.spent += chars; 202 debug!("speech reserved: {chars} chars, {} of {cap} spent today", self.spent); 203 Ok(VoiceReservation { day: today, chars }) 204 } 205 206 /// The line was not spoken: gives its characters back, if the day has not turned since. 207 pub fn release(&mut self, reservation: VoiceReservation) { 208 if self.day == reservation.day { 209 self.spent = self.spent.saturating_sub(reservation.chars); 210 debug!("speech reservation of {} chars released; {} spent today", reservation.chars, self.spent); 211 } else { 212 debug!("speech reservation of {} chars released after the day turned; nothing to give back", reservation.chars); 213 } 214 } 215} 216 217/// How a reservation ended. 218#[derive(Clone, Debug, PartialEq, Eq)] 219pub enum Settlement { 220 /// The line was spoken: the characters stay spent, and any earlier failure is forgotten. 221 Spoken, 222 /// The line was not spoken: the characters are given back, and this is remembered as the last failure. 223 NotSpoken(SpeakError), 224} 225 226/// What the parents' screen shows about the natural voice's allowance. 227#[derive(Clone, Debug, PartialEq, Eq)] 228pub struct VoiceSpend { 229 /// Characters spent on the day containing the `now` asked about. 230 pub spent_today: u32, 231 /// Why the last attempt to speak failed, until one succeeds: what shows when the natural voice 232 /// has quietly turned into the tablet's. 233 pub last_failure: Option<SpeakError>, 234} 235 236/// The allowance. 237/// 238/// Contract, for every adapter: 239/// 240/// - **Atomic.** Every method is atomic with respect to every other. In particular two `reserve_voice` 241/// calls never both succeed if together they pass the cap. 242/// - **Shared.** One allowance per household, whichever device or request asks. 243/// - **Durable.** Spending survives the hub restarting: a restart must not grant a fresh allowance. 244/// A charge is durable once `charge_thinking` returns `Ok`, a voice line once it is settled (a 245/// reservation is a promise about a line not yet spoken, and a crash may forget it: the spend is 246/// then a line short, never a line over). If the store cannot be written the call still keeps the 247/// spending in memory but returns `Err` rather than pretend it was kept. 248/// - **The rules are the ledgers'.** An adapter holds a [`ThinkingLedger`] and a [`VoiceDay`] and calls 249/// their methods; it adds nothing to the arithmetic. 250#[expect(async_fn_in_trait, reason = "a Worker's futures hold JavaScript values and cannot be Send, so no Send bound may be required here")] 251pub trait Allowance { 252 /// Whether another question may be asked now under `limit`. 253 async fn thinking_room(&self, now: Millis, limit: &TokenLimit) -> Result<ThinkingRoom, StoreError>; 254 255 /// Records what an answered question cost. 256 async fn charge_thinking(&self, now: Millis, cost: Tokens) -> Result<(), StoreError>; 257 258 async fn thinking_usage(&self, now: Millis, limit: &TokenLimit) -> Result<ThinkingUsage, StoreError>; 259 260 /// Takes `chars` from today's allowance, or refuses with [`ReserveError::DailyCap`]. 261 async fn reserve_voice(&self, now: Millis, chars: u32, cap: u32) -> Result<VoiceReservation, ReserveError>; 262 263 /// Ends a reservation: keeps the characters spent if the line was spoken, gives them back if not. 264 async fn settle_voice(&self, reservation: VoiceReservation, how: Settlement) -> Result<(), StoreError>; 265 266 async fn voice_spend(&self, now: Millis) -> Result<VoiceSpend, StoreError>; 267} 268 269#[cfg(test)] 270mod tests { 271 use super::*; 272 273 const H: u64 = HOUR_MS; 274 275 fn limit(tokens: u64, hours: u64) -> TokenLimit { 276 TokenLimit::new(Some(tokens), hours).unwrap() 277 } 278 279 fn at(ms: u64) -> Millis { 280 Millis::new(ms) 281 } 282 283 #[test] 284 fn it_has_room_until_the_window_is_spent_and_then_stops() { 285 let (mut b, l) = (ThinkingLedger::default(), limit(1000, 5)); 286 assert_eq!(b.room(at(0), &l), ThinkingRoom::Room); 287 b.charge(at(0), Tokens::new(600)); 288 assert_eq!(b.room(at(1), &l), ThinkingRoom::Room); 289 b.charge(at(1), Tokens::new(600)); 290 assert!(matches!(b.room(at(2), &l), ThinkingRoom::Spent { .. }), "1200 of 1000"); 291 } 292 293 #[test] 294 fn spending_ages_out_of_a_rolling_window_not_a_clock_day() { 295 let (mut b, l) = (ThinkingLedger::default(), limit(1000, 5)); 296 b.charge(at(0), Tokens::new(1000)); 297 assert_eq!(b.room(at(4 * H), &l), ThinkingRoom::Spent { frees_up_at: at(5 * H) }, "and it says when"); 298 assert_eq!(b.room(at(5 * H + 1), &l), ThinkingRoom::Room, "the old spending has left the window"); 299 assert_eq!(b.usage(at(5 * H + 1), &l).used, Tokens::new(0)); 300 } 301 302 #[test] 303 fn changing_the_limit_or_the_window_takes_effect_at_once() { 304 let mut b = ThinkingLedger::default(); 305 b.charge(at(0), Tokens::new(900)); 306 assert!(matches!(b.room(at(H), &limit(500, 5)), ThinkingRoom::Spent { .. }), "a lower allowance stops it now"); 307 assert_eq!(b.room(at(H), &limit(5000, 5)), ThinkingRoom::Room, "a higher one reopens it"); 308 assert_eq!(b.room(at(3 * H), &limit(500, 2)), ThinkingRoom::Room, "a shorter window forgets the spending sooner"); 309 assert_eq!(b.room(at(H), &TokenLimit::new(None, 5).unwrap()), ThinkingRoom::Room, "no limit at all"); 310 } 311 312 #[test] 313 fn it_says_when_room_starts_to_return() { 314 let (mut b, l) = (ThinkingLedger::default(), limit(1000, 5)); 315 assert_eq!(b.usage(at(0), &l).frees_up_at, None); 316 b.charge(at(2 * H), Tokens::new(300)); 317 b.charge(at(3 * H), Tokens::new(300)); 318 assert_eq!(b.usage(at(4 * H), &l).frees_up_at, Some(at(7 * H))); 319 } 320 321 #[test] 322 fn the_ledger_is_the_service_token_file_unchanged() { 323 let mut b = ThinkingLedger::default(); 324 b.charge(at(10), Tokens::new(700)); 325 let json = serde_json::to_string(&b).unwrap(); 326 assert_eq!(json, r#"{"spent":[[10,700]]}"#); 327 let again: ThinkingLedger = serde_json::from_str(&json).unwrap(); 328 assert_eq!(again.usage(at(20), &limit(1000, 5)).used, Tokens::new(700)); 329 } 330 331 #[test] 332 fn it_reads_what_a_reply_cost() { 333 let reply = r#"{"content":[],"usage":{"input_tokens":120,"output_tokens":30,"cache_read_input_tokens":50}}"#; 334 assert_eq!(tokens_in(reply), Tokens::new(150), "cache reads are free"); 335 assert_eq!(tokens_in("not json"), Tokens::new(0)); 336 assert_eq!(tokens_in(r#"{"content":[]}"#), Tokens::new(0)); 337 } 338 339 #[test] 340 fn the_voice_cap_is_reserved_up_front_and_given_back_when_the_line_does_not_come_out() { 341 let mut v = VoiceDay::default(); 342 let a = v.reserve(at(1000), 6, 10).unwrap(); 343 assert_eq!(v.spent_on(Day::containing(at(1000))), 6); 344 assert!(matches!(v.reserve(at(1001), 6, 10), Err(ReserveError::DailyCap { spent: 6, cap: 10, .. })), "the second would pass the cap"); 345 v.release(a); 346 assert_eq!(v.spent_on(Day::containing(at(1000))), 0); 347 assert!(v.reserve(at(1002), 6, 10).is_ok(), "the characters came back"); 348 } 349 350 #[test] 351 fn a_refused_line_says_when_the_day_turns_and_a_new_day_starts_from_zero() { 352 let mut v = VoiceDay::default(); 353 let _kept = v.reserve(at(5), 10, 10).unwrap(); 354 match v.reserve(at(6), 1, 10) { 355 Err(ReserveError::DailyCap { frees_up_at, .. }) => assert_eq!(frees_up_at, Day::containing(at(6)).ends_at()), 356 other => panic!("{other:?}"), 357 } 358 assert!(v.reserve(Day::containing(at(5)).ends_at(), 10, 10).is_ok(), "tomorrow is a fresh allowance"); 359 } 360 361 #[test] 362 fn a_reservation_released_after_the_day_turned_gives_nothing_back() { 363 let mut v = VoiceDay::default(); 364 let old = v.reserve(at(5), 4, 10).unwrap(); 365 let tomorrow = Day::containing(at(5)).ends_at(); 366 let _new = v.reserve(tomorrow, 3, 10).unwrap(); 367 v.release(old); 368 assert_eq!(v.spent_on(Day::containing(tomorrow)), 3, "yesterday's refund does not touch today's"); 369 } 370 371 #[test] 372 fn the_voice_file_keeps_its_old_shape() { 373 let mut v = VoiceDay::default(); 374 let _r = v.reserve(at(86_400_000 * 3 + 5), 7, 10).unwrap(); 375 assert_eq!(serde_json::to_string(&v).unwrap(), r#"{"day":3,"spent":7}"#); 376 } 377}