lmjtfy.git / apps / lmjtfy / src / meter.rs
meter.rsannotatedmeter.rssource222 lines · 9.8 KB · raw
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}