The day's budgets, kept where every isolate sees the same numbers: one
Durable Object. The arithmetic is budget's; this is its clock, its
storage, its view of the account, and the Worker's calls to it.
The neuron budget is the account's, not this Worker's: the eval, a local dev server and anything else on the account spend the same daily allocation. So before a neuron hold, and for the page, the object reads the account's own figure from Cloudflare's analytics API and raises its meter to it.
11use std::cell::RefCell;
The binding's name in wrangler.toml.
18const BINDING: &str = "BUDGET";
One object for the whole Worker: the budgets are shared by every visitor.
20const NAME: &str = "everyone";
A token that may read the account's analytics and nothing else
(nixos-config, cloudflare_account_token.lmjtfy-analytics), and the
account it reads. Without both, the meter counts this Worker alone.
How long a reading of the account is trusted. The figure lags by about as much, and a reading is one subrequest.
The account's neurons for a day, as last read.
Each visitor's requests in their current minute. In memory, never in
storage: see budget::Visits.
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 }
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 }
The account's figure for today: a recent reading, a new one, or
today's last one if Cloudflare cannot be asked. None when there has
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 }
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}
A Durable Object handles one request at a time, so reading a meter, 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}
Why a call may not be made.
177pub enum Refused {
Today's budget cannot cover it.
179 Spent(Exhausted),
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}
Replaces a hold with what the call cost. If this fails the hold stands, which overcounts: the safe direction.
Counts one request of kind what by the visitor who, and says whether
it is within per_minute. If the object cannot be asked the request is
let through: this limit is for fairness between visitors, and the day's
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}
Today's use of both budgets, for the page. None if the object cannot
be asked: the page then shows no figures instead of invented ones.