flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/billing/src/keeper.rs

673 lines27,116 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Prices keep themselves current with what g1t pays1//! Keeps every price current with what g1t actually pays.
2//!
3//! g1t passes its own costs through, so a price is only right while the
4//! cost under it is. Two jobs keep them right, on the billing service's
5//! cron:
6//!
7//! - **Settling runs** (every 15 minutes). A model run is charged when it
8//! finishes at what the sandbox reported. Each of g1t's hosted runs goes
9//! through its AI Gateway, which prices every request at the provider's
10//! current rates and logs it with the run's session. Settling sums those
11//! logs and corrects the charge to the gateway's figure, with a
12//! correction on the statement. A run whose sandbox died before
13//! reporting is charged here instead of never.
14//! - **Checking costs** (daily). What Cloudflare billed g1t's account, from
15//! its usage API, is measured against how much was used: Containers
16//! against the seconds containers ran, Workers for Platforms per request
17//! and per CPU millisecond. When a measured cost moves, the price book
18//! moves with it, since each price is its cost plus a set markup, and
19//! the change is recorded where anyone can see it. A measurement far off
20//! the current cost is not adopted, only logged, so one odd day of data
21//! cannot reprice everything.
22
23use g1t_contracts::billing::{EntryKind, MICROS_PER_DOLLAR, Price, PriceBook, PriceChange};
24use g1t_contracts::new_id;
25use g1t_contracts::time::rfc3339;
26use g1t_kit::now_ms;
27use serde::Deserialize;
28use serde_json::{Value, json};
29use worker::{Env, Fetch, Headers, Method, Request, RequestInit, Result};
30
31use crate::{Billing, RunRow, charge_micros};
32
33/// The cron that also checks costs against Cloudflare's bill.
34pub(crate) const DAILY: &str = "17 4 * * *";
35
36/// A run is settled once its logs have had time to land.
37const SETTLE_AFTER_MS: u64 = 5 * 60 * 1000;
38/// A run with no gateway logs after this is left as reported.
39const GIVE_UP_AFTER_MS: u64 = 3 * 60 * 60 * 1000;
40/// A run never finished after this died without reporting.
41const ABANDONED_AFTER_MS: u64 = 3 * 60 * 60 * 1000;
42/// Smaller moves are noise.
43const MIN_CHANGE: f64 = 0.02;
44/// A measurement outside this factor of the current cost is suspect.
45const MAX_FACTOR: f64 = 4.0;
46
47/// Where the keeper reads what g1t pays.
48pub(crate) struct Keeper {
49 /// `CLOUDFLARE_USAGE_TOKEN`: Billing, Account Analytics and AI Gateway,
50 /// read only.
51 token: Option<String>,
52 account: String,
53 gateway: String,
54}
55
56impl Keeper {
57 pub(crate) fn from_env(env: &Env) -> Self {
58 let var = |name: &str| env.var(name).map(|v| v.to_string()).unwrap_or_default();
59 Keeper {
60 token: env.secret("CLOUDFLARE_USAGE_TOKEN").ok().map(|v| v.to_string()).filter(|v| !v.is_empty()),
61 account: var("CLOUDFLARE_ACCOUNT_ID"),
62 gateway: var("AI_GATEWAY_ID"),
63 }
64 }
65
66 async fn send(&self, method: Method, url: &str, body: Option<Value>) -> Result<Value> {
67 let Some(token) = &self.token else {
68 return Err(worker::Error::RustError("no CLOUDFLARE_USAGE_TOKEN".into()));
69 };
70 let headers = Headers::new();
71 headers.set("authorization", &format!("Bearer {token}"))?;
72 headers.set("content-type", "application/json")?;
73 let mut init = RequestInit::new();
74 init.with_method(method).with_headers(headers);
75 if let Some(body) = body {
76 init.with_body(Some(body.to_string().into()));
77 }
78 let mut response = Fetch::Request(Request::new_with_init(url, &init)?).send().await?;
79 let status = response.status_code();
80 let value: Value = response.json().await.unwrap_or(Value::Null);
81 if status != 200 {
82 return Err(worker::Error::RustError(format!("Cloudflare answered {status}: {value}")));
83 }
84 Ok(value)
85 }
86
87 fn api(&self, path: &str) -> String {
88 format!("https://api.cloudflare.com/client/v4/accounts/{}{path}", self.account)
89 }
90
91 /// What AI Gateway priced a session's requests at, in dollars, and how
92 /// many there were.
93 async fn session_cost(&self, session: &str) -> Result<(f64, u32)> {
94 let mut cost = 0.0;
95 let mut count = 0;
96 for page in 1..=40 {
The keeper reads Cloudflare as it really answers97 // The filter goes as URL-encoded JSON; the bracket form is
98 // ignored, and would sum every log there is. Session ids are
99 // [a-z0-9_], which need no escaping inside it.
100 let filter = format!(
101 "%5B%7B%22key%22%3A%22metadata.value%22%2C%22operator%22%3A%22eq%22%2C%22value%22%3A%5B%22{session}%22%5D%7D%5D"
102 );
Prices keep themselves current with what g1t pays103 let url = self.api(&format!(
The keeper reads Cloudflare as it really answers104 "/ai-gateway/gateways/{}/logs?per_page=50&page={page}&filters={filter}",
Prices keep themselves current with what g1t pays105 self.gateway
106 ));
107 let body = self.send(Method::Get, &url, None).await?;
108 let logs = body["result"].as_array().cloned().unwrap_or_default();
109 for log in &logs {
110 cost += log["cost"].as_f64().unwrap_or(0.0);
111 count += 1;
112 }
113 if logs.len() < 50 {
114 break;
115 }
116 }
117 Ok((cost, count))
118 }
119
The keeper reads Cloudflare as it really answers120 /// The account's billable usage, one row per service per day, as
121 /// Cloudflare reports it.
Prices keep themselves current with what g1t pays122 async fn billable_usage(&self, from: &str, to: &str) -> Result<Vec<UsageRow>> {
123 let body = self
The keeper reads Cloudflare as it really answers124 .send(Method::Get, &self.api(&format!("/billable-usage?from={from}&to={to}")), None)
Prices keep themselves current with what g1t pays125 .await?;
126 let rows = body["result"].as_array().cloned().unwrap_or_default();
127 Ok(rows.iter().filter_map(UsageRow::from_value).collect())
128 }
129
The keeper reads Cloudflare as it really answers130 /// What g1t's containers used from `since` to `until` (dates), as
131 /// Cloudflare bills it: memory in byte-seconds, and CPU seconds.
132 async fn container_usage(&self, since: &str, until: &str) -> Result<ContainerUsage> {
133 let query = "query ($account: String!, $since: Date!, $until: Date!) {
Prices keep themselves current with what g1t pays134 viewer { accounts(filter: { accountTag: $account }) {
The keeper reads Cloudflare as it really answers135 containersUsageAdaptiveGroups(limit: 1000, filter: { date_geq: $since, date_leq: $until }) {
136 sum { cpuTimeSec allocatedMemory }
Prices keep themselves current with what g1t pays137 }
138 } }
139 }";
140 let body = self
141 .send(
142 Method::Post,
143 "https://api.cloudflare.com/client/v4/graphql",
144 Some(json!({ "query": query, "variables": { "account": self.account, "since": since, "until": until } })),
145 )
146 .await?;
The keeper reads Cloudflare as it really answers147 let groups = body["data"]["viewer"]["accounts"][0]["containersUsageAdaptiveGroups"]
Prices keep themselves current with what g1t pays148 .as_array()
149 .cloned()
150 .unwrap_or_default();
The keeper reads Cloudflare as it really answers151 Ok(groups.iter().fold(ContainerUsage::default(), |total, g| ContainerUsage {
152 cpu_seconds: total.cpu_seconds + g["sum"]["cpuTimeSec"].as_f64().unwrap_or(0.0),
153 memory_byte_seconds: total.memory_byte_seconds + g["sum"]["allocatedMemory"].as_f64().unwrap_or(0.0),
154 }))
Prices keep themselves current with what g1t pays155 }
156}
157
The keeper reads Cloudflare as it really answers158#[derive(Debug, Default, Clone, Copy)]
159pub(crate) struct ContainerUsage {
160 cpu_seconds: f64,
161 memory_byte_seconds: f64,
162}
163
164/// g1t's sandboxes: Containers' standard-1, half a vCPU, 4 GiB, 8 GB disk.
165const SANDBOX_GIB: f64 = 4.0;
166const SANDBOX_DISK_GB: f64 = 8.0;
167const GIB: f64 = 1024.0 * 1024.0 * 1024.0;
168
169/// Cloudflare's published Containers rates, in dollars, used for any rate
170/// the bill does not show yet (while usage is inside the included amount).
171const LIST_MEMORY_GIB_SECOND: f64 = 0.000_002_5;
172const LIST_DISK_GB_SECOND: f64 = 0.000_000_07;
173const LIST_VCPU_SECOND: f64 = 0.000_02;
174
175/// What one second of a sandbox costs, in millionths of a dollar: its
176/// memory and disk for the whole second, and the CPU sandboxes actually
177/// use per second of running, which is billed only while busy.
178pub(crate) fn sandbox_second_micros(usage: ContainerUsage, memory: f64, disk: f64, vcpu: f64) -> Option<f64> {
179 let instance_seconds = usage.memory_byte_seconds / (SANDBOX_GIB * GIB);
180 if instance_seconds < 3600.0 {
181 return None;
182 }
183 let cpu_share = usage.cpu_seconds / instance_seconds;
184 Some((SANDBOX_GIB * memory + SANDBOX_DISK_GB * disk + cpu_share * vcpu) * MICROS_PER_DOLLAR as f64)
185}
186
187/// A unit's marginal rate from the bill: the median, over the days that
188/// were charged, of cost over quantity. None while nothing was charged.
189pub(crate) fn billed_rate(rows: &[&UsageRow]) -> Option<f64> {
190 let mut rates: Vec<f64> = rows
191 .iter()
192 .filter(|r| r.cost > 0.0 && r.quantity > 0.0)
193 .map(|r| r.cost / r.quantity)
194 .collect();
195 if rates.is_empty() {
196 return None;
197 }
198 rates.sort_by(f64::total_cmp);
199 Some(rates[rates.len() / 2])
200}
201
Prices keep themselves current with what g1t pays202/// One line of Cloudflare's billable usage.
203#[derive(Debug, Clone)]
204pub(crate) struct UsageRow {
205 period_start: String,
206 period_end: String,
207 service: String,
208 unit: String,
209 quantity: f64,
210 cost: f64,
211}
212
213impl UsageRow {
214 /// Read leniently: the API is new, and its field names are FOCUS's.
215 fn from_value(row: &Value) -> Option<Self> {
216 let text = |keys: &[&str]| keys.iter().find_map(|k| row[*k].as_str()).unwrap_or_default().to_owned();
217 let number = |keys: &[&str]| {
218 keys.iter()
219 .find_map(|k| row[*k].as_f64().or_else(|| row[*k].as_str().and_then(|s| s.parse().ok())))
220 .unwrap_or(0.0)
221 };
222 let service = text(&["ServiceName", "service_name", "service"]);
223 if service.is_empty() {
224 return None;
225 }
226 let family = text(&["ServiceFamilyName", "service_family_name"]);
227 Some(UsageRow {
228 period_start: text(&["ChargePeriodStart", "charge_period_start"]),
229 period_end: text(&["ChargePeriodEnd", "charge_period_end"]),
230 service: if family.is_empty() { service } else { format!("{family} / {service}") },
The keeper reads Cloudflare as it really answers231 unit: text(&["PricingUnit", "ConsumedUnit", "consumed_unit"]),
Prices keep themselves current with what g1t pays232 quantity: number(&["PricingQuantity", "ConsumedQuantity", "pricing_quantity"]),
The keeper reads Cloudflare as it really answers233 // What g1t pays; list price if nothing was contracted.
234 cost: Some(number(&["ContractedCost", "BilledCost", "contracted_cost"]))
235 .filter(|cost| *cost > 0.0)
236 .unwrap_or_else(|| number(&["ListCost", "list_cost"])),
Prices keep themselves current with what g1t pays237 })
238 }
239}
240
241/// What a cost should become from a measurement, or why not.
242pub(crate) fn adopt(current: f64, measured: f64) -> std::result::Result<Option<f64>, String> {
243 if !measured.is_finite() || measured <= 0.0 {
244 return Err("nothing to measure".into());
245 }
246 let ratio = measured / current;
247 if !(1.0 / MAX_FACTOR..=MAX_FACTOR).contains(&ratio) {
248 return Err(format!("measured {measured:.4} against {current:.4}, too far off to adopt"));
249 }
250 Ok(((ratio - 1.0).abs() >= MIN_CHANGE).then_some(measured))
251}
252
253#[derive(Deserialize)]
254struct PriceRow {
255 meter: String,
256 title: String,
257 unit: String,
258 cost_micros: f64,
259 markup_percent: u32,
260 source: String,
261 checked_at: Option<String>,
262 updated_at: String,
263}
264
265#[derive(Deserialize)]
266struct ChangeRow {
267 meter: String,
268 old_cost_micros: f64,
269 new_cost_micros: f64,
270 markup_percent: u32,
271 reason: String,
272 created_at: String,
273}
274
275#[derive(Deserialize)]
276struct Unsettled {
277 id: String,
278 workspace: String,
279 repo: String,
280 number: u32,
281 task: String,
282 model: String,
283 token_hash: String,
284 billed_to: Option<String>,
285 session_id: String,
286 created_at: String,
287 finished_at: Option<String>,
288}
289
290#[derive(Deserialize)]
291struct Charged {
292 cost_micros: Option<i64>,
293 description: String,
294}
295
296fn ms(timestamp: &str) -> u64 {
297 // RFC 3339 in UTC, as g1t writes them.
298 worker::js_sys::Date::parse(timestamp) as u64
299}
300
301impl Billing {
302 pub(crate) async fn prices(&self) -> Result<PriceBook> {
303 let prices = self
304 .db
305 .prepare("SELECT * FROM prices ORDER BY rowid")
306 .all()
307 .await?
308 .results::<PriceRow>()?;
309 let changes = self
310 .db
311 .prepare("SELECT * FROM price_changes ORDER BY created_at DESC LIMIT 20")
312 .all()
313 .await?
314 .results::<ChangeRow>()?;
315 Ok(PriceBook {
316 prices: prices
317 .into_iter()
318 .map(|row| Price {
319 price_micros: Price::price_for(row.cost_micros, row.markup_percent),
320 meter: row.meter,
321 title: row.title,
322 unit: row.unit,
323 cost_micros: row.cost_micros,
324 markup_percent: row.markup_percent,
325 source: row.source,
326 checked_at: row.checked_at,
327 updated_at: row.updated_at,
328 })
329 .collect(),
330 changes: changes
331 .into_iter()
332 .map(|row| PriceChange {
333 meter: row.meter,
334 old_cost_micros: row.old_cost_micros,
335 new_cost_micros: row.new_cost_micros,
336 markup_percent: row.markup_percent,
337 reason: row.reason,
338 created_at: row.created_at,
339 })
340 .collect(),
341 model_margin_percent: self.margin_percent,
342 })
343 }
344
The keeper reads Cloudflare as it really answers345 /// Whether the costs have never been checked against Cloudflare's bill.
346 pub(crate) async fn never_checked(&self) -> Result<bool> {
347 Ok(self
348 .db
349 .prepare("SELECT meter FROM prices WHERE checked_at IS NOT NULL LIMIT 1")
350 .first::<Value>(None)
351 .await?
352 .is_none())
353 }
354
Prices keep themselves current with what g1t pays355 /// A meter's cost and price per unit, from the book.
356 pub(crate) async fn price(&self, meter: &str) -> Result<Option<(f64, f64)>> {
357 #[derive(Deserialize)]
358 struct Row {
359 cost_micros: f64,
360 markup_percent: u32,
361 }
362 Ok(self
363 .db
364 .prepare("SELECT cost_micros, markup_percent FROM prices WHERE meter = ?")
365 .bind(&[meter.into()])?
366 .first::<Row>(None)
367 .await?
368 .map(|row| (row.cost_micros, Price::price_for(row.cost_micros, row.markup_percent))))
369 }
370
371 /// Corrects finished runs to what AI Gateway priced them at, and
372 /// charges runs whose sandbox died before reporting.
373 pub(crate) async fn settle_runs(&self, keeper: &Keeper) -> Result<()> {
374 if keeper.token.is_none() || keeper.gateway.is_empty() {
375 return Ok(());
376 }
377 let now = now_ms();
378 let runs = self
379 .db
380 .prepare(
381 "SELECT id, workspace, repo, number, task, model, token_hash, billed_to, session_id, created_at, finished_at
382 FROM runs
383 WHERE session_id IS NOT NULL AND settled_at IS NULL
384 AND ((finished_at IS NOT NULL AND finished_at < ?1) OR created_at < ?2)
385 ORDER BY created_at LIMIT 10",
386 )
387 .bind(&[rfc3339(now - SETTLE_AFTER_MS).into(), rfc3339(now - ABANDONED_AFTER_MS).into()])?
388 .all()
389 .await?
390 .results::<Unsettled>()?;
391 for run in runs {
392 let (cost_usd, requests) = match keeper.session_cost(&run.session_id).await {
393 Ok(found) => found,
394 Err(error) => {
395 worker::console_error!("could not read gateway logs for {}: {error}", run.id);
396 continue;
397 }
398 };
399 let since = ms(run.finished_at.as_deref().unwrap_or(&run.created_at));
400 if requests == 0 && now.saturating_sub(since) < GIVE_UP_AFTER_MS {
401 continue;
402 }
403 self.settle(&run, cost_usd, requests).await?;
404 }
405 Ok(())
406 }
407
408 async fn settle(&self, run: &Unsettled, cost_usd: f64, requests: u32) -> Result<()> {
409 let row = RunRow {
410 workspace: run.workspace.clone(),
411 repo: run.repo.clone(),
412 number: run.number,
413 task: run.task.clone(),
414 model: run.model.clone(),
415 token_hash: run.token_hash.clone(),
416 billed_to: run.billed_to.clone(),
417 };
418 let gateway_micros = charge_micros(cost_usd, 0);
419 let charged = self
420 .db
421 .prepare("SELECT cost_micros, description FROM ledger WHERE reference = ?")
422 .bind(&[run.id.as_str().into()])?
423 .first::<Charged>(None)
424 .await?;
425 let charge_for = |micros: i64| {
426 if self.free {
427 0
428 } else {
429 charge_micros(micros as f64 / MICROS_PER_DOLLAR as f64, self.margin_percent)
430 }
431 };
432 let settled_at = rfc3339(now_ms());
433 // Claim it, so two crons never settle it twice.
434 let claimed = self
435 .db
436 .prepare("UPDATE runs SET settled_at = ?, gateway_cost_micros = ?, finished_at = COALESCE(finished_at, ?) WHERE id = ? AND settled_at IS NULL RETURNING id")
437 .bind(&[
438 settled_at.as_str().into(),
439 (gateway_micros as f64).into(),
440 settled_at.as_str().into(),
441 run.id.as_str().into(),
442 ])?
443 .first::<Value>(None)
444 .await?;
445 if claimed.is_none() || requests == 0 {
446 return Ok(());
447 }
448 let free_note = if self.free { " (free while g1t is being built out)" } else { "" };
449 match charged {
450 // Never reported: charged now, from the gateway's figure.
451 None => {
452 let description = format!(
453 "Work on {}#{}, settled from AI Gateway after the sandbox stopped without reporting{free_note}",
454 run.repo, run.number
455 );
456 self.enter(&run.workspace, EntryKind::Usage, -charge_for(gateway_micros), &description, &run.id, Some(&row), Some(gateway_micros), None, None)
457 .await?;
458 }
459 Some(charged) => {
460 let reported = charged.cost_micros.unwrap_or(0);
461 let delta = gateway_micros - reported;
462 if delta == 0 {
463 return Ok(());
464 }
465 let amount = -(charge_for(gateway_micros) - charge_for(reported));
466 let description = format!(
467 "Correction to “{}”: AI Gateway priced its {requests} model requests at {}, not {}",
468 charged.description,
469 crate::features::dollars(gateway_micros),
470 crate::features::dollars(reported),
471 );
472 self.enter(
473 &run.workspace,
474 EntryKind::Usage,
475 amount,
476 &description,
477 &format!("{}/settled", run.id),
478 Some(&row),
479 Some(delta),
480 None,
481 None,
482 )
483 .await?;
484 }
485 }
486 Ok(())
487 }
488
489 /// Checks each cost against what Cloudflare billed this month, and
490 /// moves the ones that changed.
491 pub(crate) async fn reconcile(&self, keeper: &Keeper) -> Result<()> {
492 if keeper.token.is_none() {
493 return Ok(());
494 }
495 let now = rfc3339(now_ms());
496 let today = &now[..10];
The keeper reads Cloudflare as it really answers497 let since = rfc3339(now_ms() - 30 * 24 * 60 * 60 * 1000);
498 let rows = keeper.billable_usage(&since[..10], today).await?;
Prices keep themselves current with what g1t pays499 for row in &rows {
500 self.db
501 .prepare(
502 "INSERT INTO cloudflare_usage (period_start, period_end, service, unit, quantity, cost_usd, fetched_at)
503 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
504 ON CONFLICT (period_start, service, unit) DO UPDATE SET
505 period_end = ?2, quantity = ?5, cost_usd = ?6, fetched_at = ?7",
506 )
507 .bind(&[
508 row.period_start.as_str().into(),
509 row.period_end.as_str().into(),
510 row.service.as_str().into(),
511 row.unit.as_str().into(),
512 row.quantity.into(),
513 row.cost.into(),
514 now.as_str().into(),
515 ])?
516 .run()
517 .await?;
518 }
519
The keeper reads Cloudflare as it really answers520 let named = |words: &[&str]| -> Vec<&UsageRow> {
Prices keep themselves current with what g1t pays521 rows.iter()
The keeper reads Cloudflare as it really answers522 .filter(|r| {
523 let service = r.service.to_lowercase();
524 words.iter().all(|word| service.contains(word))
525 })
526 .collect()
Prices keep themselves current with what g1t pays527 };
528
The keeper reads Cloudflare as it really answers529 // Containers: each resource at what the bill shows it costs, or
530 // the published rate while the included amount still covers it,
531 // over how much CPU g1t's sandboxes really use per second.
532 let memory = billed_rate(&named(&["container memory"]));
533 let disk = billed_rate(&named(&["container disk"]));
534 let vcpu = billed_rate(&named(&["container vcpu"]));
535 let usage = keeper.container_usage(&since[..10], today).await?;
536 if let Some(per_second) = sandbox_second_micros(
537 usage,
538 memory.unwrap_or(LIST_MEMORY_GIB_SECOND),
539 disk.unwrap_or(LIST_DISK_GB_SECOND),
540 vcpu.unwrap_or(LIST_VCPU_SECOND),
541 ) {
542 let instance_seconds = usage.memory_byte_seconds / (SANDBOX_GIB * GIB);
543 let billed = [("memory", memory), ("disk", disk), ("vCPU", vcpu)]
544 .iter()
545 .filter(|(_, rate)| rate.is_some())
546 .map(|(name, _)| *name)
547 .collect::<Vec<_>>();
548 let reason = format!(
549 "Sandboxes used {:.2} vCPU per second over {:.0} hours of Cloudflare Containers in the last 30 days; {}",
550 usage.cpu_seconds / instance_seconds,
551 instance_seconds / 3600.0,
552 if billed.is_empty() {
553 "rates are Cloudflare's published ones".to_owned()
554 } else {
555 format!("{} at what Cloudflare billed", billed.join(", "))
556 },
557 );
558 for meter in ["sandbox_second", "build_second"] {
559 self.measure(meter, per_second, &reason).await?;
Prices keep themselves current with what g1t pays560 }
561 }
The keeper reads Cloudflare as it really answers562 // Apps run as Workers: per million requests and CPU milliseconds,
563 // once the bill shows them charged.
564 let app_meters: [(&str, &[&str], &str); 2] = [
565 ("app_requests", &["workers", "requests"], "requests"),
566 ("app_cpu", &["workers cpu"], "CPU ms"),
567 ];
568 for (meter, words, unit) in app_meters {
569 if let Some(rate) = billed_rate(&named(words)) {
570 let reason = format!("Cloudflare billed Workers {unit} at ${:.2} per million", rate * 1e6);
571 self.measure(meter, rate * 1e6 * MICROS_PER_DOLLAR as f64, &reason).await?;
Prices keep themselves current with what g1t pays572 }
573 }
574 self.db
575 .prepare("UPDATE prices SET checked_at = ?")
576 .bind(&[now.as_str().into()])?
577 .run()
578 .await?;
579 Ok(())
580 }
581
582 /// Moves a meter's cost to a measurement, if it is sound and different.
583 async fn measure(&self, meter: &str, measured: f64, reason: &str) -> Result<()> {
584 let Some((current, _)) = self.price(meter).await? else {
585 return Ok(());
586 };
587 match adopt(current, measured) {
588 Err(why) => worker::console_log!("{meter}: {why}"),
589 Ok(None) => {}
590 Ok(Some(cost)) => {
591 let now = now_ms();
592 self.db
593 .batch(vec![
594 self.db
595 .prepare("UPDATE prices SET cost_micros = ?, source = 'cloudflare', updated_at = ? WHERE meter = ?")
596 .bind(&[cost.into(), rfc3339(now).into(), meter.into()])?,
597 self.db
598 .prepare(
599 "INSERT INTO price_changes (id, meter, old_cost_micros, new_cost_micros, markup_percent, reason, created_at)
600 SELECT ?, meter, ?, ?, markup_percent, ?, ? FROM prices WHERE meter = ?",
601 )
602 .bind(&[
603 new_id("prc", now).into(),
604 current.into(),
605 cost.into(),
606 reason.into(),
607 rfc3339(now).into(),
608 meter.into(),
609 ])?,
610 ])
611 .await?;
612 }
613 }
614 Ok(())
615 }
616}
617
618#[cfg(test)]
619mod tests {
620 use super::*;
621
622 #[test]
623 fn small_moves_are_noise_and_wild_ones_are_not_believed() {
624 assert_eq!(adopt(21.0, 21.2), Ok(None));
625 assert_eq!(adopt(21.0, 25.0), Ok(Some(25.0)));
626 assert_eq!(adopt(21.0, 15.0), Ok(Some(15.0)));
627 assert!(adopt(21.0, 200.0).is_err());
628 assert!(adopt(21.0, 0.0).is_err());
629 }
630
631 #[test]
The keeper reads Cloudflare as it really answers632 fn a_sandbox_second_is_its_memory_and_disk_and_the_cpu_it_uses() {
633 // An hour of sandboxes that kept a fifth of a vCPU busy.
634 let usage = ContainerUsage { cpu_seconds: 720.0, memory_byte_seconds: 3600.0 * 4.0 * GIB };
635 let micros = sandbox_second_micros(usage, LIST_MEMORY_GIB_SECOND, LIST_DISK_GB_SECOND, LIST_VCPU_SECOND).unwrap();
636 // 4 x 2.5 + 8 x 0.07 + 0.2 x 20 = 14.56
637 assert!((micros - 14.56).abs() < 1e-9, "{micros}");
638 // Too little use to say anything.
639 assert!(sandbox_second_micros(ContainerUsage { cpu_seconds: 1.0, memory_byte_seconds: GIB }, 1.0, 1.0, 1.0).is_none());
640 }
641
642 #[test]
643 fn a_billed_rate_is_the_median_of_the_charged_days() {
644 let row = |quantity: f64, cost: f64| UsageRow {
645 period_start: String::new(),
646 period_end: String::new(),
647 service: "Containers / Container Memory".into(),
648 unit: "Count".into(),
649 quantity,
650 cost,
651 };
652 let rows = [row(100.0, 0.0), row(100.0, 0.0002), row(100.0, 0.00025), row(100.0, 0.00025)];
653 assert_eq!(billed_rate(&rows.iter().collect::<Vec<_>>()), Some(0.000_002_5));
654 assert_eq!(billed_rate(&[&row(5.0, 0.0)]), None);
655 }
656
657 #[test]
Prices keep themselves current with what g1t pays658 fn usage_rows_are_read_by_their_focus_names() {
659 let row = UsageRow::from_value(&json!({
660 "ServiceFamilyName": "Containers",
661 "ServiceName": "Memory",
The keeper reads Cloudflare as it really answers662 "PricingUnit": "GiB-seconds",
Prices keep themselves current with what g1t pays663 "PricingQuantity": "1200.5",
664 "ContractedCost": 0.003,
665 "ChargePeriodStart": "2026-10-01",
666 }))
667 .unwrap();
668 assert_eq!(row.service, "Containers / Memory");
669 assert_eq!(row.quantity, 1200.5);
670 assert_eq!(row.cost, 0.003);
671 assert!(UsageRow::from_value(&json!({ "nothing": 1 })).is_none());
672 }
673}