1//! The day's budgets, kept where every isolate sees the same numbers: one 2//! Durable Object. The arithmetic is `budget`'s; this is its clock, its 3//! storage, its view of the account, and the Worker's calls to it. 4//! 5//! The neuron budget is the account's, not this Worker's: the eval, a local 6//! dev server and anything else on the account spend the same daily 7//! allocation. So before a neuron hold, and for the page, the object reads 8//! the account's own figure from Cloudflare's analytics API and raises its 9//! meter to it. 10 11use std::cell::RefCell; 12 13use budget::{Ask, Counted, Exhausted, Hold, Meter, Status, Told, Visits, Which}; 14use serde::{Deserialize, Serialize}; 15use worker::{DurableObject, Env, Fetch, Headers, Method, Request, RequestInit, Response, State, durable_object}; 16 17/// The binding's name in `wrangler.toml`. 18const BINDING: &str = "BUDGET"; 19/// One object for the whole Worker: the budgets are shared by every visitor. 20const NAME: &str = "everyone"; 21 22/// A token that may read the account's analytics and nothing else 23/// (nixos-config, `cloudflare_account_token.lmjtfy-analytics`), and the 24/// account it reads. Without both, the meter counts this Worker alone. 25const ANALYTICS_TOKEN: &str = "CLOUDFLARE_ANALYTICS_TOKEN"; 26const ACCOUNT_ID: &str = "CLOUDFLARE_ACCOUNT_ID"; 27/// How long a reading of the account is trusted. The figure lags by about as 28/// much, and a reading is one subrequest. 29const READING_MS: f64 = 30_000.0; 30const GRAPHQL: &str = "https://api.cloudflare.com/client/v4/graphql"; 31const USAGE: &str = "query($a:String!,$d:Date!){viewer{accounts(filter:{accountTag:$a})\ 32 {aiInferenceAdaptiveGroups(limit:1,filter:{date:$d}){sum{totalNeurons}}}}}"; 33 34/// The account's neurons for a day, as last read. 35#[derive(Clone, Copy, Serialize, Deserialize)] 36struct Reading { 37 day: u64, 38 at_ms: f64, 39 neurons: f64, 40} 41 42#[durable_object] 43pub struct Budget { 44 state: State, 45 env: Env, 46 /// Each visitor's requests in their current minute. In memory, never in 47 /// storage: see `budget::Visits`. 48 visits: RefCell<Visits>, 49} 50 51fn failed(error: impl std::fmt::Display) -> worker::Error { 52 worker::Error::from(error.to_string()) 53} 54 55impl Budget { 56 // Stored as JSON text: one value, readable in the dashboard, and no 57 // serde-wasm-bindgen in between. 58 async fn load<T: for<'de> Deserialize<'de>>(&self, key: &str) -> worker::Result<Option<T>> { 59 let stored: Option<String> = self.state.storage().get(key).await?; 60 Ok(stored.and_then(|text| serde_json::from_str(&text).ok())) 61 } 62 63 async fn save<T: Serialize>(&self, key: &str, value: &T) -> worker::Result<()> { 64 self.state.storage().put(key, serde_json::to_string(value).map_err(failed)?).await 65 } 66 67 fn setting(&self, name: &str) -> Option<String> { 68 let value = self.env.secret(name).ok()?.to_string(); 69 (!value.trim().is_empty()).then(|| value.trim().to_owned()) 70 } 71 72 /// Asks Cloudflare what the account has used today. 73 async fn read_account(&self, today: u64) -> Result<f64, String> { 74 let (token, account) = self 75 .setting(ANALYTICS_TOKEN) 76 .zip(self.setting(ACCOUNT_ID)) 77 .ok_or("no analytics token or account id")?; 78 let body = serde_json::json!({ "query": USAGE, "variables": { "a": account, "d": budget::date(today) } }); 79 let headers = Headers::new(); 80 headers.set("authorization", &format!("Bearer {token}")).map_err(|e| e.to_string())?; 81 headers.set("content-type", "application/json").map_err(|e| e.to_string())?; 82 let mut init = RequestInit::new(); 83 init.with_method(Method::Post).with_headers(headers).with_body(Some(body.to_string().into())); 84 let request = Request::new_with_init(GRAPHQL, &init).map_err(|e| e.to_string())?; 85 let mut response = Fetch::Request(request).send().await.map_err(|e| e.to_string())?; 86 let text = response.text().await.map_err(|e| e.to_string())?; 87 let reply: serde_json::Value = serde_json::from_str(&text).map_err(|e| e.to_string())?; 88 let groups = reply 89 .pointer("/data/viewer/accounts/0/aiInferenceAdaptiveGroups") 90 .and_then(serde_json::Value::as_array) 91 .ok_or_else(|| format!("analytics: {}", reply.get("errors").unwrap_or(&reply)))?; 92 // No group at all means nothing has run today. 93 Ok(groups.first().and_then(|group| group.pointer("/sum/totalNeurons")).and_then(|n| n.as_f64()).unwrap_or(0.0)) 94 } 95 96 /// The account's figure for today: a recent reading, a new one, or 97 /// today's last one if Cloudflare cannot be asked. `None` when there has 98 /// never been one today. 99 async fn account(&self, today: u64, now_ms: f64) -> worker::Result<Option<f64>> { 100 let last: Option<Reading> = self.load("account").await?; 101 let today_s = last.filter(|reading| reading.day == today); 102 if let Some(reading) = today_s.filter(|reading| now_ms - reading.at_ms < READING_MS) { 103 return Ok(Some(reading.neurons)); 104 } 105 match self.read_account(today).await { 106 Ok(neurons) => { 107 self.save("account", &Reading { day: today, at_ms: now_ms, neurons }).await?; 108 Ok(Some(neurons)) 109 } 110 Err(_) => Ok(today_s.map(|reading| reading.neurons)), 111 } 112 } 113 114 /// The neuron meter, raised to the account's figure. Says whose count it is. 115 async fn neurons(&self, today: u64, now_ms: f64) -> worker::Result<(Meter, Counted)> { 116 let mut meter: Meter = self.load(Which::Neurons.key()).await?.unwrap_or_default(); 117 let counted = match self.account(today, now_ms).await? { 118 Some(neurons) => { 119 meter.observe(today, neurons); 120 Counted::Account 121 } 122 None => Counted::Own, 123 }; 124 Ok((meter, counted)) 125 } 126} 127 128impl DurableObject for Budget { 129 fn new(state: State, env: Env) -> Self { 130 Budget { state, env, visits: RefCell::new(Visits::default()) } 131 } 132 133 /// A Durable Object handles one request at a time, so reading a meter, 134 /// changing it and writing it back cannot interleave with another hold. 135 async fn fetch(&self, mut request: Request) -> worker::Result<Response> { 136 let ask: Ask = serde_json::from_str(&request.text().await?).map_err(failed)?; 137 let now_ms = js_sys::Date::now(); 138 let today = budget::day(now_ms); 139 let told = match ask { 140 Ask::Status => { 141 let (neurons, counted) = self.neurons(today, now_ms).await?; 142 self.save(Which::Neurons.key(), &neurons).await?; 143 let jev: Meter = self.load(Which::JevDollars.key()).await?.unwrap_or_default(); 144 Told::Status(Status { neurons: neurons.used(today), counted, jev_dollars: jev.used(today) }) 145 } 146 Ask::Hold { which, amount, per_day } => { 147 let mut meter = match which { 148 Which::Neurons => self.neurons(today, now_ms).await?.0, 149 Which::JevDollars => self.load(which.key()).await?.unwrap_or_default(), 150 }; 151 let told = match meter.hold(today, amount, per_day) { 152 Ok(hold) => Told::Held(hold), 153 Err(exhausted) => Told::Exhausted(exhausted), 154 }; 155 self.save(which.key(), &meter).await?; 156 told 157 } 158 Ask::Visit { what, who, per_minute } => { 159 Told::Visit(self.visits.borrow_mut().admit(&what, &who, now_ms, per_minute)) 160 } 161 Ask::Settle { which, hold, actual } => { 162 let mut meter: Meter = self.load(which.key()).await?.unwrap_or_default(); 163 meter.settle(today, hold, actual); 164 self.save(which.key(), &meter).await?; 165 Told::Settled 166 } 167 }; 168 Response::ok(serde_json::to_string(&told).map_err(failed)?) 169 } 170} 171 172async fn tell(env: &Env, ask: Ask) -> Result<Told, String> { 173 crate::object::tell(env, BINDING, NAME, "the budget", &ask).await 174} 175 176/// Why a call may not be made. 177pub enum Refused { 178 /// Today's budget cannot cover it. 179 Spent(Exhausted), 180 /// The budget could not be asked. Nothing is spent without it. 181 Unreachable(String), 182} 183 184/// Sets aside the most a call can cost. 185pub async fn hold(env: &Env, which: Which, amount: f64, per_day: f64) -> Result<Hold, Refused> { 186 match tell(env, Ask::Hold { which, amount, per_day }).await { 187 Ok(Told::Held(hold)) => Ok(hold), 188 Ok(Told::Exhausted(exhausted)) => Err(Refused::Spent(exhausted)), 189 Ok(_) => Err(Refused::Unreachable("the budget answered a hold with something else".into())), 190 Err(error) => Err(Refused::Unreachable(error)), 191 } 192} 193 194/// Replaces a hold with what the call cost. If this fails the hold stands, 195/// which overcounts: the safe direction. 196pub async fn settle(env: &Env, which: Which, hold: Hold, actual: f64) { 197 let _ = tell(env, Ask::Settle { which, hold, actual }).await; 198} 199 200/// Counts one request of kind `what` by the visitor `who`, and says whether 201/// it is within `per_minute`. If the object cannot be asked the request is 202/// let through: this limit is for fairness between visitors, and the day's 203/// budgets, which fail closed, are what stop the spending. 204pub async fn admit(env: &Env, what: &str, who: String, per_minute: u32) -> bool { 205 match tell(env, Ask::Visit { what: what.to_owned(), who, per_minute }).await { 206 Ok(Told::Visit(within)) => within, 207 Ok(_) => true, 208 Err(error) => { 209 worker::console_error!("the visitor limit could not be checked, so the request is let through: {error}"); 210 true 211 } 212 } 213} 214 215/// Today's use of both budgets, for the page. `None` if the object cannot 216/// be asked: the page then shows no figures instead of invented ones. 217pub async fn status(env: &Env) -> Option<Status> { 218 match tell(env, Ask::Status).await { 219 Ok(Told::Status(status)) => Some(status), 220 _ => None, 221 } 222}