Skip to content

g1t/services/billing/src/lib.rs

1,552 lines70,426 bytesCodeBlame
1//! The billing service: what g1t's compute costs, charged to the workspace
2//! it ran for at what g1t pays plus 20%.
3//!
4//! One paid plan, "g1t" (see `features`): $20 a month per workspace, never
5//! per person, with $10 of usage included, deployments, and more private
6//! storage. The forge is free for everyone; compute needs the plan or a
7//! card check (see `compute` and `cards`). Before anything that costs money
8//! starts, the service that starts it reserves its estimate here; when it
9//! is done, what it cost goes on the ledger plus the margin, drawn first
10//! from what the plan includes, a trial or g1t's pools (see `credits`).
11//! Every change is a ledger entry, and a balance is always the sum of its
12//! ledger.
13//!
14//! Without a card processor configured the service says so and charges
15//! nothing, so that g1t still runs where billing has not been set up.
16//!
17//! Reached only through service bindings; see `g1t_contracts::billing` for
18//! the methods and their arguments.
19
20mod accounts;
21mod ai;
22mod budget;
23mod cards;
24mod closing;
25mod compute;
26mod costs;
27mod margin;
28mod pricing;
29mod report;
30mod details;
31mod credits;
32mod grants;
33mod overages;
34mod requests;
35mod storage;
36mod invoices;
37mod sales;
38mod statement;
39mod webhooks;
40mod features;
41mod keeper;
42mod limits;
43mod rename;
44mod reset;
45mod retention;
46mod stripe;
47mod stripe_sync;
48mod subscriptions;
49mod tokens;
50
51use g1t_contracts::billing::*;
52use g1t_contracts::time::rfc3339;
53use g1t_contracts::{FailureCode, Outcome, Role, new_id};
54use g1t_kit::{args, now_ms, reply, rpc_method};
55use serde::Deserialize;
56use sha2::{Digest, Sha256};
57use worker::wasm_bindgen::JsValue;
58use worker::{Context, D1Database, Env, MessageBatch, MessageExt, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
59use futures_util::future::{try_join, try_join5};
60
61use stripe::Stripe;
62
63/// Prepaying: $25 at the least; by card up to $10,000 at a time, and by
64/// bank transfer from $1,000 to $100,000.
65const MIN_TOP_UP_CENTS: u32 = 2_500;
66const MAX_TOP_UP_CENTS: u32 = 1_000_000;
67const MIN_BANK_TRANSFER_CENTS: u32 = 100_000;
68const MAX_BANK_TRANSFER_CENTS: u32 = 10_000_000;
69const LEDGER_PAGE: u32 = 100;
70/// A run's reported cost is believed up to this much. A sandbox cannot
71/// spend more in the time it has, so anything above is a fault.
72const MAX_RUN_COST_USD: f64 = 100.0;
73
74/// What a run is charged: its cost plus the margin, rounded up to a whole
75/// millionth of a dollar.
76pub fn charge_micros(cost_usd: f64, margin_percent: u32) -> i64 {
77 let cost_micros = (cost_usd.clamp(0.0, MAX_RUN_COST_USD) * MICROS_PER_DOLLAR as f64).ceil();
78 (cost_micros * f64::from(100 + margin_percent) / 100.0).ceil() as i64
79}
80
81/// What a cost g1t trusts (the price book's, or what AI Gateway priced a
82/// run at) is charged at: plus the margin, rounded up to a whole millionth.
83/// Unlike `charge_micros`, never capped: only a sandbox's own report is
84/// held to `MAX_RUN_COST_USD`, so a run that really cost more is charged
85/// for all of it once it is settled.
86pub fn margin_on(cost_micros: i64, margin_percent: u32) -> i64 {
87 let cost = i128::from(cost_micros.max(0));
88 let charge = (cost * i128::from(100 + margin_percent) + 99) / 100;
89 i64::try_from(charge).unwrap_or(i64::MAX)
90}
91
92fn hash(token: &str) -> String {
93 hex::encode(Sha256::digest(token.as_bytes()))
94}
95
96fn optional(value: Option<&str>) -> JsValue {
97 value.map_or(JsValue::NULL, JsValue::from)
98}
99
100#[derive(Deserialize)]
101struct AccountRow {
102 balance_micros: i64,
103 customer_id: Option<String>,
104}
105
106#[derive(Deserialize)]
107struct LedgerRow {
108 id: String,
109 kind: EntryKind,
110 amount_micros: i64,
111 description: String,
112 repo: Option<String>,
113 number: Option<u32>,
114 task: Option<String>,
115 model: Option<String>,
116 created_by: Option<String>,
117 created_at: String,
118 billed_to: Option<String>,
119 #[serde(default)]
120 workspace: Option<String>,
121 #[serde(default)]
122 credit_micros: Option<i64>,
123 #[serde(default)]
124 trial_micros: Option<i64>,
125 #[serde(default)]
126 oss_micros: Option<i64>,
127 #[serde(default)]
128 given_micros: Option<i64>,
129 #[serde(default)]
130 credit_kind: Option<String>,
131 #[serde(default)]
132 discount_micros: Option<i64>,
133}
134
135impl From<LedgerRow> for LedgerEntry {
136 fn from(row: LedgerRow) -> Self {
137 LedgerEntry {
138 id: row.id,
139 kind: row.kind,
140 amount_micros: row.amount_micros,
141 description: row.description,
142 repo: row.repo,
143 number: row.number,
144 task: row.task,
145 model: row.model,
146 billed_to: row.billed_to.unwrap_or_else(|| "g1t".to_owned()),
147 created_by: row.created_by,
148 created_at: row.created_at,
149 workspace: row.workspace,
150 credit_micros: row.credit_micros.unwrap_or(0),
151 trial_micros: row.trial_micros.unwrap_or(0),
152 oss_micros: row.oss_micros.unwrap_or(0),
153 given_micros: row.given_micros.unwrap_or(0),
154 credit_kind: row.credit_kind.as_deref().and_then(CreditKind::parse),
155 discount_micros: row.discount_micros.unwrap_or(0),
156 }
157 }
158}
159
160#[derive(Deserialize)]
161struct RunRow {
162 workspace: String,
163 repo: String,
164 number: u32,
165 task: String,
166 model: String,
167 token_hash: String,
168 billed_to: Option<String>,
169}
170
171impl RunRow {
172 fn own_provider(&self) -> bool {
173 self.billed_to.as_deref() == Some("workspace")
174 }
175}
176
177#[derive(Deserialize)]
178struct CheckoutRow {
179 workspace: String,
180 created_by: String,
181}
182
183/// A payment page started, as `checkouts` keeps it.
184pub(crate) struct NewCheckout<'a> {
185 pub id: &'a str,
186 pub workspace: &'a str,
187 /// What it pays for, in cents: the prepayment, the plan's price, the AI
188 /// credit; 0 for a card check.
189 pub amount_cents: u32,
190 /// The card fee on top, for AI credit.
191 pub fee_cents: u32,
192 pub created_by: &'a str,
193 /// None for a prepayment; `plan`, `security`, `card_check` or
194 /// `ai_credit` otherwise.
195 pub feature: Option<&'a str>,
196}
197
198/// The one insert every payment page goes through, every column named, so
199/// a column added to `checkouts` with no default is caught by
200/// `tests::every_checkout_insert_fills_the_table` rather than by a 500.
201pub(crate) const CHECKOUT_INSERT: &str = "INSERT INTO checkouts (id, workspace, amount_cents, fee_cents, created_by, created_at, feature, status)
202 VALUES (?, ?, ?, ?, ?, ?, ?, 'open')";
203
204impl Billing {
205 pub(crate) async fn record_checkout(&self, c: &NewCheckout<'_>) -> Result<()> {
206 self.db
207 .prepare(CHECKOUT_INSERT)
208 .bind(&[
209 c.id.into(),
210 c.workspace.into(),
211 c.amount_cents.into(),
212 c.fee_cents.into(),
213 c.created_by.into(),
214 rfc3339(now_ms()).into(),
215 optional(c.feature),
216 ])?
217 .run()
218 .await?;
219 Ok(())
220 }
221}
222
223/// A row an `UPDATE … RETURNING` touched.
224#[derive(Deserialize)]
225struct Touched {
226 #[allow(dead_code)]
227 id: String,
228}
229
230struct Billing {
231 db: D1Database,
232 /// Absent when no card processor is configured.
233 stripe: Option<Stripe>,
234 /// The destination's signing secret from Stripe (`STRIPE_WEBHOOK_SECRET`);
235 /// without it no event is believed.
236 webhook_secret: Option<String>,
237 margin_percent: u32,
238 /// While g1t is being built out, nothing is charged (`FREE_WHILE_BUILDING`).
239 free: bool,
240 /// Whether new workspaces get trial credit (`TRIAL_WORKSPACE_MICROS`
241 /// and `TRIAL_MONTHLY_POOL_MICROS` both above zero).
242 trials_on: bool,
243 /// The plans' and pools' numbers; see `credits`.
244 plans: credits::Config,
245 /// The repos service: which repositories are public, what private ones
246 /// hold, and their git operations. Absent where it is not bound.
247 repos: Option<worker::Fetcher>,
248 /// The packages service: what each workspace's packages hold.
249 packages: Option<worker::Fetcher>,
250 /// The identity service, which emails owners. Absent where it is not
251 /// bound.
252 identity: Option<worker::Fetcher>,
253 /// How far unpaid usage may go; see `limits`.
254 ceilings: limits::Ceilings,
255 /// `PREPAID_ONLY`: the old rule, that agents need credit first.
256 prepaid_only: bool,
257 /// The caps on what g1t pays for itself; see `budget`.
258 caps: budget::Caps,
259 /// The worker's bindings, for emailing staff (`EMAIL`).
260 env: Env,
261}
262
263impl Billing {
264 fn status(&self) -> Status {
265 Status {
266 enabled: self.stripe.is_some(),
267 live: self.stripe.as_ref().is_some_and(Stripe::live),
268 free: self.free,
269 }
270 }
271
272 async fn row(&self, workspace: &str) -> Result<Option<AccountRow>> {
273 self.db
274 .prepare("SELECT balance_micros, customer_id FROM accounts WHERE workspace = ?")
275 .bind(&[workspace.into()])?
276 .first::<AccountRow>(None)
277 .await
278 }
279
280 async fn standing(&self, workspace: &str) -> Result<Account> {
281 // The card as last synced (stripe_sync.rs), not asked of Stripe.
282 let (row, card) = try_join(self.row(workspace), self.saved_card(workspace)).await?;
283 Ok(Account {
284 workspace: workspace.to_owned(),
285 balance_micros: row.map_or(0, |row| row.balance_micros),
286 status: self.status(),
287 margin_percent: self.margin_percent,
288 card,
289 })
290 }
291
292 /// The workspace's customer at Stripe, made the first time one is needed.
293 pub(crate) async fn customer_for(&self, workspace: &str) -> Result<String> {
294 let Some(stripe) = &self.stripe else {
295 return Err(worker::Error::RustError("payments are not set up".into()));
296 };
297 if let Some(customer) = self.row(workspace).await?.and_then(|row| row.customer_id) {
298 return Ok(customer);
299 }
300 let customer = stripe.create_customer(workspace).await?;
301 self.db
302 .prepare(
303 "INSERT INTO accounts (workspace, balance_micros, customer_id, created_at) VALUES (?1, 0, ?2, ?3)
304 ON CONFLICT (workspace) DO UPDATE SET customer_id = ?2",
305 )
306 .bind(&[workspace.into(), customer.as_str().into(), rfc3339(now_ms()).into()])?
307 .run()
308 .await?;
309 Ok(customer)
310 }
311
312 /// Stripe's hosted billing page for the workspace. Owners only.
313 async fn billing_portal(&self, a: BillingPortalArgs) -> Result<Outcome<Checkout>> {
314 let workspace = a.workspace.to_lowercase();
315 if a.actor.role_in(&workspace) != Some(Role::Owner) {
316 return Ok(Outcome::fail(FailureCode::Forbidden, "Only an owner can manage the workspace's billing."));
317 }
318 let Some(stripe) = &self.stripe else {
319 return Ok(Outcome::fail(FailureCode::Conflict, "Payments are not set up on this g1t."));
320 };
321 let customer = match self.customer_for(&workspace).await {
322 Ok(customer) => customer,
323 Err(error) => return Ok(Outcome::fail(FailureCode::Conflict, crate::stripe::friendly(&error))),
324 };
325 match stripe.portal_session(&customer, &a.return_url).await {
326 Ok(url) => Ok(Outcome::Ok(Checkout { url })),
327 Err(error) if stripe::is_missing(&error) => {
328 // The customer was removed at Stripe: a new one next time.
329 self.forget_customer(&workspace).await?;
330 Ok(Outcome::fail(FailureCode::Conflict, "Stripe no longer had this workspace's customer. Try again."))
331 }
332 Err(error) => Ok(Outcome::fail(FailureCode::Conflict, crate::stripe::friendly(&error))),
333 }
334 }
335
336 /// Adds a ledger entry and moves the balance by the same amount, as
337 /// one write.
338 #[allow(clippy::too_many_arguments)]
339 async fn enter(
340 &self,
341 workspace: &str,
342 kind: EntryKind,
343 amount_micros: i64,
344 description: &str,
345 reference: &str,
346 run: Option<&RunRow>,
347 cost_micros: Option<i64>,
348 created_by: Option<&str>,
349 customer: Option<&str>,
350 ) -> Result<()> {
351 let now = now_ms();
352 let timestamp = rfc3339(now);
353 let kind = match kind {
354 EntryKind::TopUp => "top_up",
355 EntryKind::Usage => "usage",
356 };
357 self.db
358 .batch(vec![
359 self.db
360 .prepare(
361 "INSERT INTO ledger
362 (id, workspace, kind, amount_micros, description, repo, number, task,
363 model, cost_micros, reference, created_by, created_at, billed_to)
364 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
365 )
366 .bind(&[
367 new_id("led", now).into(),
368 workspace.into(),
369 kind.into(),
370 // D1 takes numbers as doubles, which hold every
371 // amount this service will see exactly.
372 (amount_micros as f64).into(),
373 description.into(),
374 optional(run.map(|run| run.repo.as_str())),
375 run.map_or(JsValue::NULL, |run| run.number.into()),
376 optional(run.map(|run| run.task.as_str())),
377 optional(run.map(|run| run.model.as_str())),
378 cost_micros.map_or(JsValue::NULL, |cost| (cost as f64).into()),
379 reference.into(),
380 optional(created_by),
381 timestamp.as_str().into(),
382 run.map_or("g1t", |run| if run.own_provider() { "workspace" } else { "g1t" }).into(),
383 ])?,
384 self.db
385 .prepare(
386 "INSERT INTO accounts (workspace, balance_micros, customer_id, created_at)
387 VALUES (?1, ?2, ?3, ?4)
388 ON CONFLICT (workspace) DO UPDATE SET
389 balance_micros = balance_micros + ?2,
390 customer_id = COALESCE(?3, customer_id)",
391 )
392 .bind(&[
393 workspace.into(),
394 (amount_micros as f64).into(),
395 optional(customer),
396 timestamp.as_str().into(),
397 ])?,
398 ])
399 .await?;
400 // Money in clears a card declined at the limit. A refund or a lost
401 // dispute is a top-up of less than nothing: money out, which never
402 // clears it.
403 if kind == "top_up" && amount_micros > 0 {
404 self.db
405 .prepare("UPDATE limits SET autopay_failed_at = NULL, autopay_error = NULL WHERE workspace = ?")
406 .bind(&[workspace.into()])?
407 .run()
408 .await?;
409 }
410 Ok(())
411 }
412
413 async fn account(&self, a: AccountArgs) -> Result<Outcome<Account>> {
414 let workspace = a.workspace.to_lowercase();
415 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
416 return Ok(members_only());
417 }
418 Ok(Outcome::Ok(self.standing(&workspace).await?))
419 }
420
421 async fn ledger(&self, a: AccountArgs) -> Result<Outcome<Vec<LedgerEntry>>> {
422 let workspace = a.workspace.to_lowercase();
423 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
424 return Ok(members_only());
425 }
426 let rows = self
427 .db
428 .prepare("SELECT * FROM ledger WHERE workspace = ? ORDER BY id DESC LIMIT ?")
429 .bind(&[workspace.into(), LEDGER_PAGE.into()])?
430 .all()
431 .await?
432 .results::<LedgerRow>()?;
433 Ok(Outcome::Ok(
434 rows.into_iter().map(LedgerEntry::from).collect(),
435 ))
436 }
437
438 async fn usage(&self, a: UsageArgs) -> Result<Outcome<Usage>> {
439 let workspace = a.workspace.to_lowercase();
440 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
441 return Ok(members_only());
442 }
443 #[derive(serde::Deserialize)]
444 struct SliceRow {
445 key: Option<String>,
446 micros: Option<i64>,
447 runs: Option<u32>,
448 }
449 // While g1t is free, what was used at cost is what there is to show.
450 // With a discount, usage at its price, so a 100% discount still
451 // shows what the workspace would pay.
452 // Every slice measures usage at price (statement::PRICE_SQL), the one
453 // figure every page shows as usage, whatever paid for it.
454 let discount_percent = self.terms_of(&workspace).await?.percent_off();
455 let measure = if self.free { "COALESCE(cost_micros, 0)" } else { crate::statement::PRICE_SQL };
456 let slices = |key: &str, limit: u32| {
457 format!(
458 "SELECT {key} AS key, SUM({measure}) AS micros, COUNT(*) AS runs FROM ledger
459 WHERE workspace = ?1 AND kind = 'usage' AND created_at >= ?2
460 GROUP BY 1 ORDER BY micros DESC LIMIT {limit}"
461 )
462 };
463 let query = |sql: String| {
464 let db = &self.db;
465 let workspace = workspace.clone();
466 let since = a.since.clone();
467 async move {
468 let rows = db
469 .prepare(sql)
470 .bind(&[workspace.into(), since.into()])?
471 .all()
472 .await?
473 .results::<SliceRow>()?;
474 Ok::<Vec<UsageSlice>, worker::Error>(
475 rows.into_iter()
476 .map(|row| UsageSlice {
477 key: row.key.unwrap_or_else(|| "other".to_owned()),
478 micros: row.micros.unwrap_or_default(),
479 runs: row.runs.unwrap_or_default(),
480 })
481 .collect(),
482 )
483 }
484 };
485 #[derive(serde::Deserialize)]
486 struct Totals {
487 spent: Option<i64>,
488 cost: Option<i64>,
489 provider: Option<i64>,
490 runs: Option<u32>,
491 added: Option<i64>,
492 discount: Option<i64>,
493 covered: Option<i64>,
494 price: Option<i64>,
495 }
496 let totals = async {
497 self.db
498 .prepare(format!(
499 "SELECT
500 -SUM(CASE WHEN kind = 'usage' THEN amount_micros END) AS spent,
501 SUM(CASE WHEN kind = 'usage' THEN COALESCE(credit_micros, 0) + COALESCE(trial_micros, 0)
502 + COALESCE(oss_micros, 0) + COALESCE(given_micros, 0) END) AS covered,
503 SUM(CASE WHEN kind = 'usage' THEN {price} END) AS price,",
504 price = crate::statement::PRICE_SQL
505 ) + "
506 SUM(CASE WHEN kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' THEN cost_micros END) AS cost,
507 SUM(CASE WHEN kind = 'usage' AND billed_to = 'workspace' THEN cost_micros END) AS provider,
508 SUM(CASE WHEN kind = 'usage' AND " + crate::statement::RUN_SQL + " THEN 1 ELSE 0 END) AS runs,
509 SUM(CASE WHEN kind = 'top_up' THEN amount_micros END) AS added,
510 SUM(CASE WHEN kind = 'usage' THEN discount_micros END) AS discount
511 FROM ledger WHERE workspace = ?1 AND created_at >= ?2",
512 )
513 .bind(&[workspace.as_str().into(), a.since.as_str().into()])?
514 .first::<Totals>(None)
515 .await
516 };
517 // The totals and the five slices read the same rows independently,
518 // so they go to D1 at once: one round trip of waiting, not six.
519 let (totals, (by_day, by_task, by_repo, by_pull, by_model)) = try_join(
520 totals,
521 try_join5(
522 query(slices("substr(created_at, 1, 10) || '/' || COALESCE(task, 'other')", 400)),
523 query(slices("task", 20)),
524 query(slices("repo", 20)),
525 query(slices("repo || '#' || number", 10)),
526 query(slices("model", 10)),
527 ),
528 )
529 .await?;
530 let totals = totals.unwrap_or(Totals {
531 spent: None,
532 cost: None,
533 provider: None,
534 runs: None,
535 added: None,
536 discount: None,
537 covered: None,
538 price: None,
539 });
540 Ok(Outcome::Ok(Usage {
541 spent_micros: totals.spent.unwrap_or_default(),
542 // What the included usage, the trial, a pool or g1t paid, as the
543 // entries record it.
544 covered_micros: totals.covered.unwrap_or_default(),
545 price_micros: if self.free { totals.cost.unwrap_or_default() } else { totals.price.unwrap_or_default() },
546 cost_micros: totals.cost.unwrap_or_default(),
547 provider_micros: totals.provider.unwrap_or_default(),
548 used_micros: totals.cost.unwrap_or_default() + totals.provider.unwrap_or_default(),
549 free: self.free,
550 discount_micros: totals.discount.unwrap_or_default(),
551 discount_percent: (discount_percent > 0).then_some(discount_percent),
552 runs: totals.runs.unwrap_or_default(),
553 added_micros: totals.added.unwrap_or_default(),
554 by_day,
555 by_task,
556 by_repo,
557 by_pull,
558 by_model,
559 since: a.since,
560 }))
561 }
562
563 /// `checkout`: prepays usage, by card or (from $1,000) bank transfer.
564 async fn checkout(&self, a: CheckoutArgs) -> Result<Outcome<Checkout>> {
565 let workspace = a.workspace.to_lowercase();
566 if a.actor.role_in(&workspace) != Some(Role::Owner) {
567 return Ok(Outcome::fail(
568 FailureCode::Forbidden,
569 "Only an owner can prepay for a workspace.",
570 ));
571 }
572 let Some(stripe) = &self.stripe else {
573 return Ok(Outcome::fail(
574 FailureCode::Conflict,
575 "Payments are not set up on this g1t yet.",
576 ));
577 };
578 let bank_transfer = a.method.as_deref() == Some("bank_transfer");
579 if let Err(why) = prepay_amount(a.amount_cents, bank_transfer) {
580 return Ok(Outcome::fail(FailureCode::Invalid, why));
581 }
582 // A bank transfer needs a customer for its account details.
583 let customer = if bank_transfer {
584 match self.customer_for(&workspace).await {
585 Ok(customer) => Some(customer),
586 Err(error) => return Ok(Outcome::fail(FailureCode::Conflict, stripe::friendly(&error))),
587 }
588 } else {
589 self.row(&workspace).await?.and_then(|row| row.customer_id)
590 };
591 let started = match stripe.start_checkout(&workspace, a.amount_cents, customer.as_deref(), &a.return_url, bank_transfer).await {
592 // A customer saved under another Stripe account: start afresh.
593 Err(error) if customer.is_some() && stripe::is_missing(&error) => {
594 self.forget_customer(&workspace).await?;
595 let customer = if bank_transfer { self.customer_for(&workspace).await.ok() } else { None };
596 stripe.start_checkout(&workspace, a.amount_cents, customer.as_deref(), &a.return_url, bank_transfer).await
597 }
598 other => other,
599 };
600 self.page_opened(started, &workspace, a.amount_cents, 0, &a.actor.username, None).await
601 }
602
603 /// A payment page Stripe started (or refused), recorded in `checkouts`
604 /// and handed back; any failure as a sentence for the page, never a
605 /// raw error.
606 pub(crate) async fn page_opened(
607 &self,
608 started: Result<stripe::Session>,
609 workspace: &str,
610 amount_cents: u32,
611 fee_cents: u32,
612 created_by: &str,
613 feature: Option<&str>,
614 ) -> Result<Outcome<Checkout>> {
615 let session = match started {
616 Ok(session) => session,
617 Err(error) => {
618 worker::console_error!("{workspace}: Stripe refused a payment page ({}): {error}", feature.unwrap_or("prepay"));
619 return Ok(Outcome::fail(FailureCode::Conflict, stripe::friendly(&error)));
620 }
621 };
622 let Some(url) = session.url.clone() else {
623 return Ok(Outcome::fail(FailureCode::Conflict, "Stripe returned no payment page. Try again in a minute."));
624 };
625 let recorded = self
626 .record_checkout(&NewCheckout { id: &session.id, workspace, amount_cents, fee_cents, created_by, feature })
627 .await;
628 if let Err(error) = recorded {
629 worker::console_error!("{workspace}: a payment page could not be recorded: {error}");
630 return Ok(Outcome::fail(FailureCode::Conflict, "g1t could not keep track of the payment page. Nothing was charged; try again."));
631 }
632 Ok(Outcome::Ok(Checkout { url }))
633 }
634
635 /// Credits a payment if the processor says it was made and it has not
636 /// been credited before. The amount credited is what the processor
637 /// says was paid, not what anyone here remembers asking for.
638 async fn confirm(&self, a: ConfirmArgs) -> Result<Outcome<Account>> {
639 let workspace = a.workspace.to_lowercase();
640 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
641 return Ok(members_only());
642 }
643 let (Some(stripe), Some(checkout)) = (
644 &self.stripe,
645 self.db
646 .prepare(
647 // A prepayment only: a plan's or a card check's page is
648 // settled where it was started, never credited as money
649 // paid in advance.
650 "SELECT workspace, created_by FROM checkouts
651 WHERE id = ? AND workspace = ? AND status = 'open' AND feature IS NULL",
652 )
653 .bind(&[a.session.as_str().into(), workspace.as_str().into()])?
654 .first::<CheckoutRow>(None)
655 .await?,
656 ) else {
657 // Unknown, someone else's, or already credited: nothing to do.
658 return Ok(Outcome::Ok(self.standing(&workspace).await?));
659 };
660 let session = match stripe.session(&a.session).await {
661 Ok(session) => session,
662 Err(error) => return Ok(Outcome::fail(FailureCode::Conflict, crate::stripe::friendly(&error))),
663 };
664 let paid = session
665 .amount_total
666 .filter(|_| session.payment_status == "paid");
667 if let Some(cents) = paid {
668 // Only whoever flips it from open to paid enters the credit.
669 let claimed = self
670 .db
671 .prepare(
672 "UPDATE checkouts SET status = 'paid' WHERE id = ? AND status = 'open'
673 RETURNING id",
674 )
675 .bind(&[a.session.as_str().into()])?
676 .first::<Touched>(None)
677 .await?;
678 if claimed.is_some() {
679 self.enter(
680 &checkout.workspace,
681 EntryKind::TopUp,
682 i64::from(cents) * MICROS_PER_DOLLAR / 100,
683 "Paid in advance",
684 &session.id,
685 None,
686 None,
687 Some(&checkout.created_by),
688 session.customer.as_deref(),
689 )
690 .await?;
691 }
692 }
693 Ok(Outcome::Ok(self.standing(&workspace).await?))
694 }
695
696 /// Drops a saved customer the card processor no longer knows.
697 pub(crate) async fn forget_customer(&self, workspace: &str) -> Result<()> {
698 self.db
699 .prepare(
700 "UPDATE accounts SET customer_id = NULL, card_brand = NULL, card_last4 = NULL, card_exp_month = NULL,
701 card_exp_year = NULL, card_synced_at = NULL WHERE workspace = ?",
702 )
703 .bind(&[workspace.into()])?
704 .run()
705 .await?;
706 Ok(())
707 }
708
709 /// A refusal if the workspace has no credit to start an agent with.
710 async fn out_of_credit<T>(&self, workspace: &str) -> Result<Option<Outcome<T>>> {
711 // Billing is postpaid: usage limits decide whether work starts
712 // (see `limits`), and credit is a prepayment that lowers what is
713 // owed. A balance no longer has to be positive to start.
714 if self.free || !self.prepaid_only {
715 return Ok(None);
716 }
717 let balance = self
718 .row(workspace)
719 .await?
720 .map_or(0, |row| row.balance_micros);
721 Ok((balance <= 0).then(|| {
722 Outcome::fail(
723 FailureCode::PaymentRequired,
724 format!(
725 "The {workspace} workspace has no agent credit. An owner can add some under Billing on the workspace's page."
726 ),
727 )
728 }))
729 }
730
731 async fn can_start(&self, a: CanStartArgs) -> Result<Outcome<bool>> {
732 if self.stripe.is_none() {
733 return Ok(Outcome::Ok(true));
734 }
735 if let Some(stopped) = self.stopped(&a.workspace).await? {
736 return Ok(stopped);
737 }
738 Ok(self
739 .out_of_credit(&a.workspace.to_lowercase())
740 .await?
741 .unwrap_or(Outcome::Ok(true)))
742 }
743
744 async fn start_run(&self, a: StartRunArgs) -> Result<Outcome<Option<RunTicket>>> {
745 if self.stripe.is_none() {
746 return Ok(Outcome::Ok(None));
747 }
748 let workspace = a.workspace.to_lowercase();
749 if let Some(stopped) = self.stopped(&workspace).await? {
750 return Ok(stopped);
751 }
752 if let Some(refused) = self.out_of_credit(&workspace).await? {
753 return Ok(refused);
754 }
755 // On g1t's models, a workspace on the plan needs AI credit (or
756 // included usage) first: g1t never fronts a model's cost (ai.rs).
757 if a.billed_to != "workspace"
758 && let Some(why) = self.ai_refusal(&workspace).await?
759 {
760 return Ok(Outcome::fail(FailureCode::PaymentRequired, why));
761 }
762 let now = now_ms();
763 let run_id = new_id("run", now);
764 let mut bytes = [0u8; 32];
765 getrandom::getrandom(&mut bytes).expect("no source of randomness");
766 let token = hex::encode(bytes);
767 self.db
768 .prepare(
769 "INSERT INTO runs (id, workspace, repo, number, task, model, token_hash, created_at, billed_to, session_id, tier)
770 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
771 )
772 .bind(&[
773 run_id.as_str().into(),
774 workspace.into(),
775 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
776 a.number.into(),
777 a.task.into(),
778 a.model.into(),
779 hash(&token).into(),
780 rfc3339(now).into(),
781 if a.billed_to == "workspace" { "workspace" } else { "g1t" }.into(),
782 optional(a.session.as_deref().filter(|_| a.billed_to != "workspace")),
783 optional(
784 a.tier
785 .as_deref()
786 .filter(|tier| a.billed_to != "workspace" && matches!(*tier, "small" | "large")),
787 ),
788 ])?
789 .run()
790 .await?;
791 Ok(Outcome::Ok(Some(RunTicket { run_id, token })))
792 }
793
794 async fn finish_run(&self, a: FinishRunArgs) -> Result<Outcome<bool>> {
795 let run = self
796 .db
797 .prepare(
798 "SELECT workspace, repo, number, task, model, token_hash, billed_to FROM runs
799 WHERE id = ? AND finished_at IS NULL",
800 )
801 .bind(&[a.run_id.as_str().into()])?
802 .first::<RunRow>(None)
803 .await?;
804 let Some(run) = run.filter(|run| run.token_hash == hash(&a.token)) else {
805 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
806 };
807 if !a.cost_usd.is_finite() || a.cost_usd < 0.0 {
808 return Ok(Outcome::fail(FailureCode::Invalid, "That is not a cost."));
809 }
810 // Only whoever closes the run charges for it.
811 let claimed = self
812 .db
813 .prepare(
814 "UPDATE runs SET finished_at = ? WHERE id = ? AND finished_at IS NULL RETURNING id",
815 )
816 .bind(&[rfc3339(now_ms()).into(), a.run_id.as_str().into()])?
817 .first::<Touched>(None)
818 .await?;
819 if claimed.is_none() {
820 return Ok(Outcome::Ok(false));
821 }
822 // On the workspace's own provider, the model was paid for there,
823 // and the run's sandbox time is recorded on its own: nothing more
824 // to charge.
825 if run.own_provider() {
826 return Ok(Outcome::Ok(true));
827 }
828 // Its cost plus the margin, on the account's terms; then the plan's
829 // included usage and the trial credit pay what they can, and g1t
830 // covers a free workspace's overrun (see `credits`). Agents are
831 // never the open-source pool's.
832 // The model at the provider's price and the price book's markup on it
833 // (`agent_models`: none from 2026-10-08); g1t's own part is the agent
834 // rate, charged on a line of its own below.
835 let base = charge_micros(a.cost_usd, self.model_markup().await?);
836 let (charge, terms_note, discount) = self.charged(&run.workspace, base).await?;
837 let month = credits::month_of(&rfc3339(now_ms()));
838 let eligible = credits::eligible_for(Some(ComputeKind::Agent), None);
839 let drawn = self.draw(&run.workspace, charge, &month, &eligible).await?;
840 let mut description = match run.task.as_str() {
841 "plan" => format!("Planning for {}", run.repo),
842 "review" => format!("Review of {}#{}", run.repo, run.number),
843 "update" => format!("Catching up {}#{}", run.repo, run.number),
844 _ => format!("Work on {}#{}", run.repo, run.number),
845 };
846 description.push_str(&terms_note);
847 description.push_str(&drawn.note());
848 self.enter(
849 &run.workspace,
850 EntryKind::Usage,
851 -(charge - drawn.total()),
852 &description,
853 &a.run_id,
854 Some(&run),
855 Some(charge_micros(a.cost_usd, 0)),
856 None,
857 None,
858 )
859 .await?;
860 self.record_drawn(&a.run_id, &drawn).await?;
861 self.record_discount(&a.run_id, discount).await?;
862 self.count_spend(&run.workspace, charge_micros(a.cost_usd, 0), charge - drawn.total(), &drawn).await;
863 self.charge_agent_rate(&a.run_id, &run).await?;
864 Ok(Outcome::Ok(true))
865 }
866}
867
868impl Billing {
869 /// Records how long a sandbox ran, with its cost and its charge: every
870 /// second, from the first, at the price book's price; on its own CPU
871 /// when it reports it. Settles its reservation, if it names one.
872 async fn record_sandbox(&self, a: RecordSandboxArgs) -> Result<Outcome<bool>> {
873 // Self-hosted time is recorded wherever g1t runs, for its minutes;
874 // anything else only where there is a bill to put it on.
875 if (self.stripe.is_none() && !a.self_hosted) || a.seconds == 0 {
876 return Ok(Outcome::Ok(false));
877 }
878 let workspace = a.workspace.to_lowercase();
879 let seen = self
880 .db
881 .prepare("SELECT id FROM ledger WHERE reference = ?")
882 .bind(&[a.reference.as_str().into()])?
883 .first::<Touched>(None)
884 .await?;
885 if seen.is_some() {
886 return Ok(Outcome::Ok(false));
887 }
888 let now = now_ms();
889 let timestamp = rfc3339(now);
890 let seconds = i64::from(a.seconds);
891 // On the workspace's own machine: its minutes go on usage, at $0,
892 // and whatever was reserved for it is given back.
893 if a.self_hosted {
894 let description = format!("{}: {} of self-hosted runner time, $0", a.description, duration(seconds));
895 self.db
896 .prepare(
897 "INSERT INTO ledger
898 (id, workspace, kind, amount_micros, description, repo, task,
899 cost_micros, reference, created_at, billed_to, credit_micros, trial_micros, oss_micros, given_micros, quantity)
900 VALUES (?, ?, 'usage', 0, ?, ?, 'self_hosted', 0, ?, ?, 'workspace', 0, 0, 0, 0, ?)",
901 )
902 .bind(&[
903 new_id("led", now).into(),
904 workspace.as_str().into(),
905 description.as_str().into(),
906 optional(a.repo.as_deref()),
907 a.reference.as_str().into(),
908 timestamp.as_str().into(),
909 (seconds as f64).into(),
910 ])?
911 .run()
912 .await?;
913 if let Some(reservation) = &a.reservation_id {
914 self.settle_reservation(SettleArgs { reservation_id: reservation.clone(), actual_micros: 0 }).await?;
915 }
916 return Ok(Outcome::Ok(true));
917 }
918 // From the price book, which follows what Cloudflare bills g1t. A
919 // sandbox is the same container as a build, so without a row it is
920 // a build second's cost plus the margin.
921 let (cost_per_second, _) = self.price("sandbox_second").await?.unwrap_or_else(|| {
922 let cost = deployment_costs::MICROS_PER_BUILD_SECOND as f64;
923 (cost, Price::price_for(cost, self.margin_percent))
924 });
925 // A larger machine (`runs-on: g1t-4core`): its memory and disk
926 // cost more each second, and without its own CPU it is priced at
927 // its vCPUs as busy as the standard machine's.
928 let instance = a.instance.as_deref().and_then(g1t_contracts::actions::instance_named);
929 let base_scale = instance.map_or(1.0, |i| keeper::base_scale(i.memory_gib, i.disk_gb));
930 let price_scale = instance.map_or(1.0, |i| i.price_scale);
931 let vcpus = instance.map_or(4.0, |i| i.vcpu.max(4.0));
932 // Its own CPU when the sandbox reports it; otherwise the average.
933 let parts = match a.cpu_seconds.filter(|cpu| cpu.is_finite() && *cpu >= 0.0) {
934 Some(cpu) => match (self.price("sandbox_base_second").await?, self.price("sandbox_cpu_second").await?) {
935 (Some((base, _)), Some((vcpu, _))) => {
936 Some((keeper::run_cost(seconds, cpu.min(seconds as f64 * vcpus), base * base_scale, vcpu), cpu))
937 }
938 _ => None,
939 },
940 None => None,
941 };
942 let cost = parts.map_or(seconds as f64 * cost_per_second * price_scale, |(cost, _)| cost).ceil() as i64;
943 let (charge, terms_note, discount) = self.charged(&workspace, credits::with_margin(cost, self.margin_percent)).await?;
944 let eligible = credits::eligible_for(a.kind, a.repo.as_deref());
945 let drawn = self.draw(&workspace, charge, &credits::month_of(&timestamp), &eligible).await?;
946 let charge = charge - drawn.total();
947 let cpu_note = parts.map_or(String::new(), |(_, cpu)| format!(", {cpu:.0} vCPU-seconds"));
948 let description = format!("{}: {} of sandbox time{cpu_note}{terms_note}{}", a.description, duration(seconds), drawn.note());
949 // The price versions it was charged at (pricing.rs).
950 let meters: &[&str] = if parts.is_some() { &["sandbox_base_second", "sandbox_cpu_second"] } else { &["sandbox_second"] };
951 let mut versions = Vec::new();
952 for meter in meters {
953 versions.extend(self.version_now(meter).await?);
954 }
955 let price_version = versions.join(",");
956 self.db
957 .batch(vec![
958 self.db
959 .prepare(
960 "INSERT INTO ledger
961 (id, workspace, kind, amount_micros, description, repo, task,
962 cost_micros, reference, created_at, billed_to, credit_micros, trial_micros, oss_micros, given_micros, price_version, quantity, compute)
963 VALUES (?, ?, 'usage', ?, ?, ?, 'sandbox', ?, ?, ?, 'g1t', ?, ?, ?, ?, ?, ?, ?)",
964 )
965 .bind(&[
966 new_id("led", now).into(),
967 workspace.as_str().into(),
968 (-(charge as f64)).into(),
969 description.as_str().into(),
970 optional(a.repo.as_deref()),
971 (cost as f64).into(),
972 a.reference.as_str().into(),
973 timestamp.as_str().into(),
974 (drawn.credit as f64).into(),
975 (drawn.trial as f64).into(),
976 (drawn.oss as f64).into(),
977 (drawn.given as f64).into(),
978 optional(Some(price_version.as_str()).filter(|v| !v.is_empty())),
979 (seconds as f64).into(),
980 optional(a.kind.map(ComputeKind::as_str)),
981 ])?,
982 self.db
983 .prepare(
984 "INSERT INTO accounts (workspace, balance_micros, created_at)
985 VALUES (?1, ?2, ?3)
986 ON CONFLICT (workspace) DO UPDATE SET balance_micros = balance_micros + ?2",
987 )
988 .bind(&[
989 workspace.as_str().into(),
990 (-(charge as f64)).into(),
991 timestamp.as_str().into(),
992 ])?,
993 ])
994 .await?;
995 self.record_discount(&a.reference, discount).await?;
996 self.count_spend(&workspace, cost, charge, &drawn).await;
997 if let Some(reservation) = &a.reservation_id {
998 self.settle_reservation(SettleArgs { reservation_id: reservation.clone(), actual_micros: cost }).await?;
999 }
1000 Ok(Outcome::Ok(true))
1001 }
1002}
1003
1004/// Whether an amount may be prepaid: $25 at the least by card, and from
1005/// $1,000 by bank transfer.
1006pub(crate) fn prepay_amount(cents: u32, bank_transfer: bool) -> std::result::Result<(), String> {
1007 let (min, max) = if bank_transfer { (MIN_BANK_TRANSFER_CENTS, MAX_BANK_TRANSFER_CENTS) } else { (MIN_TOP_UP_CENTS, MAX_TOP_UP_CENTS) };
1008 if (min..=max).contains(&cents) {
1009 return Ok(());
1010 }
1011 Err(if bank_transfer {
1012 format!("Prepay between ${} and ${} by bank transfer.", min / 100, group(max / 100))
1013 } else {
1014 format!("Prepay between ${} and ${} by card; from $1,000, a bank transfer works too.", min / 100, group(max / 100))
1015 })
1016}
1017
1018/// `10,000` for 10000.
1019fn group(n: u32) -> String {
1020 let digits = n.to_string();
1021 let mut out = String::new();
1022 for (i, c) in digits.chars().enumerate() {
1023 if i > 0 && (digits.len() - i).is_multiple_of(3) {
1024 out.push(',');
1025 }
1026 out.push(c);
1027 }
1028 out
1029}
1030
1031#[cfg(test)]
1032/// What `seconds` of sandbox time are charged at `price_per_second`, in
1033/// millionths of a dollar: every second, rounded up to the next millionth.
1034fn sandbox_charge(seconds: i64, price_per_second: f64) -> i64 {
1035 (seconds as f64 * price_per_second).ceil() as i64
1036}
1037
1038/// `1h 2m`, `3m 12s` or `40s`.
1039fn duration(seconds: i64) -> String {
1040 let (h, m, s) = (seconds / 3600, seconds % 3600 / 60, seconds % 60);
1041 if h > 0 {
1042 format!("{h}h {m}m")
1043 } else if m > 0 {
1044 format!("{m}m {s}s")
1045 } else {
1046 format!("{s}s")
1047 }
1048}
1049
1050impl Billing {
1051 /// What a workspace is charged for something that would be `base`:
1052 /// nothing while g1t is free, or as its account's terms say. With a
1053 /// note for the statement when it differs.
1054 /// A charge at cost plus the margin (`base`) on the account's terms:
1055 /// what is charged, the note for the statement, and what a discount
1056 /// gave away below `base`. That last is written on the entry
1057 /// (`record_discount`) so the reconciliation counts it as given, never
1058 /// as margin lost: a sold charge is worth at least its cost plus the
1059 /// margin.
1060 pub(crate) async fn charged(&self, workspace: &str, base: i64) -> Result<(i64, String, i64)> {
1061 if self.free {
1062 return Ok((0, " (free while g1t is being built out)".to_owned(), 0));
1063 }
1064 let terms = self.terms_of(workspace).await?;
1065 let (charge, discount) = terms.discounted(base);
1066 let note = match terms.discount_label() {
1067 Some(label) if base > 0 => format!(" ({label})"),
1068 _ => String::new(),
1069 };
1070 Ok((charge, note, discount))
1071 }
1072
1073 /// What a discount gave away on an entry, below cost plus the margin
1074 /// (less than nothing on a correction down).
1075 pub(crate) async fn record_discount(&self, reference: &str, micros: i64) -> Result<()> {
1076 if micros == 0 {
1077 return Ok(());
1078 }
1079 self.db
1080 .prepare("UPDATE ledger SET discount_micros = ? WHERE reference = ?")
1081 .bind(&[(micros as f64).into(), reference.into()])?
1082 .run()
1083 .await?;
1084 Ok(())
1085 }
1086}
1087
1088fn members_only<T>() -> Outcome<T> {
1089 Outcome::fail(
1090 FailureCode::Forbidden,
1091 "Only members can see a workspace's billing.",
1092 )
1093}
1094
1095impl Billing {
1096 fn from_env(env: &Env) -> Result<Self> {
1097 Ok(Billing {
1098 db: env.d1("DB")?,
1099 stripe: env
1100 .secret("STRIPE_SECRET_KEY")
1101 .ok()
1102 .map(|key| key.to_string())
1103 .filter(|key| !key.is_empty())
1104 .map(Stripe::new),
1105 webhook_secret: env
1106 .secret("STRIPE_WEBHOOK_SECRET")
1107 .ok()
1108 .map(|secret| secret.to_string().trim().to_owned())
1109 .filter(|secret| !secret.is_empty()),
1110 margin_percent: env
1111 .var("MARGIN_PERCENT")
1112 .ok()
1113 .and_then(|percent| percent.to_string().parse().ok())
1114 .unwrap_or(20),
1115 free: env.var("FREE_WHILE_BUILDING").is_ok_and(|v| v.to_string() == "true"),
1116 ceilings: limits::Ceilings::from_env(env),
1117 prepaid_only: env.var("PREPAID_ONLY").is_ok_and(|v| v.to_string() == "true"),
1118 trials_on: {
1119 let plans = credits::Config::from_env(env);
1120 plans.trial_workspace_micros > 0 && plans.trial_monthly_pool_micros > 0
1121 },
1122 plans: credits::Config::from_env(env),
1123 repos: env.service("REPOS").ok(),
1124 packages: env.service("PACKAGES").ok(),
1125 identity: env.service("IDENTITY").ok(),
1126 caps: budget::Caps::from_env(env),
1127 env: env.clone(),
1128 })
1129 }
1130}
1131
1132#[event(scheduled)]
1133async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
1134 let Ok(billing) = Billing::from_env(&env) else {
1135 return;
1136 };
1137 let keeper = keeper::Keeper::from_env(&env);
1138 if let Err(error) = billing.settle_runs(&keeper).await {
1139 worker::console_error!("settling runs failed: {error}");
1140 }
1141 // Stripe events billing never received, handled now.
1142 match billing.replay_events().await {
1143 Ok(done) => worker::console_log!("stripe replay: {done}"),
1144 Err(error) => worker::console_error!("replaying Stripe events failed: {error}"),
1145 }
1146 if let Err(error) = billing.autopay().await {
1147 worker::console_error!("paying at the limit failed: {error}");
1148 }
1149 // AI credit below a workspace's auto-reload threshold, reloaded (ai.rs).
1150 match billing.reload_ai_credit().await {
1151 Ok(0) => {}
1152 Ok(done) => worker::console_log!("auto-reloaded AI credit for {done} workspaces"),
1153 Err(error) => worker::console_error!("auto-reloading AI credit failed: {error}"),
1154 }
1155 // Last month's metered usage (scans, embeddings, storage) goes on the
1156 // ledger before the month is closed and invoiced.
1157 if let Err(error) = billing.charge_pending().await {
1158 worker::console_error!("charging last month's metered usage failed: {error}");
1159 }
1160 if let Err(error) = billing.close_months().await {
1161 worker::console_error!("closing the month failed: {error}");
1162 }
1163 if let Err(error) = billing.invoice_enterprises().await {
1164 worker::console_error!("invoicing enterprises failed: {error}");
1165 }
1166 // Comped budgets' alerts, and a tripped breaker staff were not told of.
1167 if let Err(error) = billing.watch_spend().await {
1168 worker::console_error!("watching g1t's own spend failed: {error}");
1169 }
1170 if let Ok(identity) = env.service("IDENTITY")
1171 && let Err(error) = billing.warn_limits(&identity).await {
1172 worker::console_error!("warning owners failed: {error}");
1173 }
1174 // Once a day: Stripe's endpoint kept listening to billing's events and
1175 // enabled, and saved cards and plans not read in a while read again.
1176 if event.cron() == keeper::DAILY {
1177 match billing.keep_endpoint("billing").await {
1178 Ok(done) => worker::console_log!("stripe endpoint: {done}"),
1179 Err(error) => worker::console_error!("keeping Stripe's endpoint failed: {error}"),
1180 }
1181 match billing.refresh_from_stripe().await {
1182 Ok(done) => worker::console_log!("stripe refresh: {done}"),
1183 Err(error) => worker::console_error!("refreshing from Stripe failed: {error}"),
1184 }
1185 }
1186 // Once a day, and at once if the costs were never checked: check every
1187 // cost against what Cloudflare billed.
1188 if (event.cron() == keeper::DAILY || billing.never_checked().await.unwrap_or(false))
1189 && let Err(error) = billing.reconcile(&keeper).await {
1190 worker::console_error!("checking costs against Cloudflare failed: {error}");
1191 }
1192 // Once a day: credit from g1t past its expiry stops counting
1193 // (grants.rs), before the day is reconciled.
1194 if event.cron() == keeper::DAILY {
1195 match billing.expire_credits().await {
1196 Ok(closed) => worker::console_log!("credits expired: {closed}"),
1197 Err(error) => worker::console_error!("expiring credits failed: {error}"),
1198 }
1199 }
1200 // Once a day: what Cloudflare charged, reconciled against what g1t
1201 // counted and charged; prices whose day has come; margin alerts
1202 // (margin.rs). After the keeper, so its proposals are in.
1203 if event.cron() == keeper::DAILY {
1204 match billing.costs_daily(&env, &keeper).await {
1205 Ok(run) => worker::console_log!("costs: {} lines, {} days, {} proposals, {} alerts", run.lines, run.days, run.proposals, run.alerts),
1206 Err(error) => worker::console_error!("reconciling costs failed: {error}"),
1207 }
1208 }
1209 // Once a day: what each workspace's private repositories hold, its git
1210 // operations, Deployments plans from before the g1t plan set to end,
1211 // and old reservations cleared.
1212 if event.cron() == keeper::DAILY {
1213 if let Err(error) = billing.measure_packages().await {
1214 worker::console_error!("measuring package storage failed: {error}");
1215 }
1216 if let Err(error) = billing.measure_storage().await {
1217 worker::console_error!("measuring storage failed: {error}");
1218 }
1219 if let Err(error) = billing.measure_git().await {
1220 worker::console_error!("measuring git operations failed: {error}");
1221 }
1222 if let Err(error) = billing.retire_deployments_plans().await {
1223 worker::console_error!("ending Deployments plans failed: {error}");
1224 }
1225 if let Err(error) = billing.sweep_reservations().await {
1226 worker::console_error!("clearing reservations failed: {error}");
1227 }
1228 }
1229}
1230
1231/// Events from the bus, on billing's own queue: only a workspace's rename
1232/// matters here (see `rename`).
1233#[event(queue)]
1234async fn queue(batch: MessageBatch<g1t_contracts::events::Event>, env: Env, _ctx: Context) -> Result<()> {
1235 let billing = Billing::from_env(&env)?;
1236 let identity = env.service("IDENTITY").ok();
1237 for message in batch.messages()? {
1238 // A repository transferred: its share of the open-source pool this
1239 // month follows it. What it was charged stays with the workspace
1240 // it was charged to; usage from now on is charged to the new one.
1241 if g1t_kit::transfer::on_event(&env, &billing.db, message.body(), closing::TRANSFERRED).await? {
1242 message.ack();
1243 continue;
1244 }
1245 billing.on_event(identity.as_ref(), message.body()).await?;
1246 message.ack();
1247 }
1248 Ok(())
1249}
1250
1251#[event(fetch)]
1252async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
1253 let Some(method) = rpc_method(&request) else {
1254 return Response::error("Not found", 404);
1255 };
1256 // A replica near the caller when it asks for one (crates/kit/src/d1.rs);
1257 // what other services call (can_start, start_run) asks for none.
1258 let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?;
1259 let body: serde_json::Value = request.json().await?;
1260 let mut billing = Billing::from_env(&env)?;
1261 billing.db = db;
1262 let answered = match method.as_str() {
1263 "status" => reply(&billing.status()),
1264 "account" => reply(&billing.account(args(body)?).await?),
1265 "ledger" => reply(&billing.ledger(args(body)?).await?),
1266 "credits" => reply(&billing.credits(args(body)?).await?),
1267 "usage" => reply(&billing.usage(args(body)?).await?),
1268 "usage_report" => reply(&billing.usage_report(args(body)?).await?),
1269 "ai_credit" => reply(&billing.ai_credit(args(body)?).await?),
1270 "buy_ai_credit" => reply(&billing.buy_ai_credit(args(body)?).await?),
1271 "confirm_ai_credit" => reply(&billing.confirm_ai_credit(args(body)?).await?),
1272 "set_ai_reload" => reply(&billing.set_ai_reload(args(body)?).await?),
1273 "set_budget" => reply(&billing.set_budget(args(body)?).await?),
1274 "billing_details" => reply(&billing.billing_details(args(body)?).await?),
1275 "set_billing_details" => reply(&billing.set_billing_details(args(body)?).await?),
1276 "record_tokens" => reply(&billing.record_tokens(args(body)?).await?),
1277 "token_usage" => reply(&billing.token_usage(args(body)?).await?),
1278 "checkout" => reply(&billing.checkout(args(body)?).await?),
1279 "confirm" => reply(&billing.confirm(args(body)?).await?),
1280 "can_start" => reply(&billing.can_start(args(body)?).await?),
1281 "trial" => reply(&billing.trial(args(body)?).await?),
1282 "start_run" => reply(&billing.start_run(args(body)?).await?),
1283 "finish_run" => reply(&billing.finish_run(args(body)?).await?),
1284 "features" => reply(&billing.features(args(body)?).await?),
1285 "subscribe" => reply(&billing.subscribe(args(body)?).await?),
1286 "confirm_subscription" => reply(&billing.confirm_subscription(args(body)?).await?),
1287 "cancel_subscription" => reply(&billing.cancel_subscription(args(body)?).await?),
1288 "close_workspace" => reply(&billing.close_workspace(args(body)?).await?),
1289 "has_feature" => reply(&billing.has_feature(args(body)?).await?),
1290 "charge_feature" => reply(&billing.charge_feature(args(body)?).await?),
1291 "record_sandbox" => reply(&billing.record_sandbox(args(body)?).await?),
1292 "limit" => reply(&billing.limit(args(body)?).await?),
1293 "check_limit" => reply(&billing.check_limit(args(body)?).await?),
1294 "set_spend_limit" => reply(&billing.set_spend_limit(args(body)?).await?),
1295 "prices" => reply(&billing.prices().await?),
1296 "billing_portal" => reply(&billing.billing_portal(args(body)?).await?),
1297 "admin_billing_link" => reply(&billing.admin_billing_link(args(body)?).await?),
1298 "admin_stripe" => reply(&billing.admin_stripe(args(body)?).await?),
1299 "admin_enterprise_billing" => reply(&billing.admin_enterprise_billing(args(body)?).await?),
1300 "admin_invoice_enterprise" => reply(&billing.admin_invoice_enterprise(args(body)?).await?),
1301 "stripe_webhook" => reply(&billing.stripe_webhook(args(body)?).await?),
1302 "invoices" => reply(&billing.invoices(args(body)?).await?),
1303 "statement" => reply(&billing.statement(args(body)?).await?),
1304 "statement_entries" => reply(&billing.statement_entries(args(body)?).await?),
1305 "usage_meters" => reply(&billing.usage_meters(args(body)?).await?),
1306 "admin_workspace_invoices" => {
1307 let a: AdminWorkspaceInvoicesArgs = args(body)?;
1308 reply(&billing.workspace_invoices(&a.workspace.to_lowercase()).await?)
1309 }
1310 "admin_signals" => reply(&billing.admin_signals(args(body)?).await?),
1311 "admin_overview" => reply(&billing.admin_overview(args(body)?).await?),
1312 "admin_sales" => reply(&billing.admin_sales(args(body)?).await?),
1313 "admin_set_sales" => reply(&billing.admin_set_sales(args(body)?).await?),
1314 "admin_add_note" => reply(&billing.admin_add_note(args(body)?).await?),
1315 "admin_invoices" => reply(&billing.admin_invoices(args(body)?).await?),
1316 "admin_audit" => reply(&billing.admin_audit(args(body)?).await?),
1317 // A staff change made in another service, for sudo's audit log:
1318 // identity's restores and purges of deleted workspaces.
1319 "admin_log" => {
1320 let a: AdminLogArgs = args(body)?;
1321 billing
1322 .audit(&accounts::own_account(&a.workspace), &a.action, &a.detail, &a.by)
1323 .await?;
1324 reply(&true)
1325 }
1326 "note_pending" => reply(&billing.note_pending(args(body)?).await?),
1327 "admin_accounts" => reply(&billing.admin_accounts(args(body)?).await?),
1328 "admin_account" => reply(&billing.admin_account(args(body)?).await?),
1329 "admin_set_terms" => reply(&billing.admin_set_terms(args(body)?).await?),
1330 "admin_create_enterprise" => reply(&billing.admin_create_enterprise(args(body)?).await?),
1331 "admin_attach" => reply(&billing.admin_attach(args(body)?).await?),
1332 "admin_credit" => reply(&billing.admin_credit(args(body)?).await?),
1333 "admin_credits" => reply(&billing.admin_credits(args(body)?).await?),
1334 "admin_revoke_credit" => reply(&billing.admin_revoke_credit(args(body)?).await?),
1335 "admin_reset_billing" => reply(&billing.admin_reset_billing(&env, args(body)?).await?),
1336 "admin_set_allowances" => reply(&billing.admin_set_allowances(args(body)?).await?),
1337 "entitlements" => reply(&billing.entitlements(args(body)?).await?),
1338 "audit_retention" => reply(&billing.audit_retention(args(body)?).await?),
1339 "reserve" => reply(&billing.reserve(args(body)?).await?),
1340 "settle" => reply(&billing.settle_reservation(args(body)?).await?),
1341 "card_check" => reply(&billing.card_check(args(body)?).await?),
1342 "confirm_card_check" => reply(&billing.confirm_card_check(args(body)?).await?),
1343 "request_limit" => reply(&billing.request_limit(args(body)?).await?),
1344 "limit_requests" => reply(&billing.limit_requests(args(body)?).await?),
1345 "confirm_spike" => reply(&billing.confirm_spike(args(body)?).await?),
1346 "set_caps" => reply(&billing.set_caps(args(body)?).await?),
1347 "admin_limit_requests" => reply(&billing.admin_limit_requests(args(body)?).await?),
1348 "admin_decide_limit_request" => reply(&billing.admin_decide_limit_request(args(body)?).await?),
1349 "admin_overages" => reply(&billing.admin_overages(args(body)?).await?),
1350 "admin_goodwill" => reply(&billing.admin_goodwill(args(body)?).await?),
1351 "admin_velocity" => reply(&billing.admin_velocity(args(body)?).await?),
1352 "admin_record_payment" => reply(&billing.admin_record_payment(args(body)?).await?),
1353 "admin_costs" => reply(&billing.admin_costs(args(body)?, keeper::Keeper::from_env(&env).can_read_bill()).await?),
1354 "admin_cost_alerts" => reply(&billing.admin_cost_alerts(args(body)?).await?),
1355 "admin_spend_caps" => reply(&billing.spend_caps().await?),
1356 "admin_lift_breaker" => reply(&billing.admin_lift_breaker(args(body)?).await?),
1357 "admin_decide_proposal" => reply(&billing.admin_decide_proposal(args(body)?).await?),
1358 "admin_set_cost_settings" => reply(&billing.admin_set_cost_settings(args(body)?).await?),
1359 "admin_set_cost_mapping" => reply(&billing.admin_set_cost_mapping(args(body)?).await?),
1360 "admin_run_costs" => reply(&billing.admin_run_costs(&env, args(body)?).await?),
1361 _ => Response::error("Unknown method", 404),
1362 };
1363 served.finish(answered)
1364}
1365
1366#[cfg(test)]
1367mod tests {
1368 use super::*;
1369
1370 #[test]
1371 fn a_run_is_charged_its_cost_plus_the_margin() {
1372 // $0.05 at 20% is six cents.
1373 assert_eq!(charge_micros(0.05, 20), 60_000);
1374 assert_eq!(charge_micros(1.0, 20), 1_200_000);
1375 assert_eq!(charge_micros(0.05, 0), 50_000);
1376 }
1377
1378 #[test]
1379 fn fractions_of_a_millionth_round_up_and_nothing_costs_less_than_nothing() {
1380 assert_eq!(charge_micros(0.000_000_4, 20), 2);
1381 assert_eq!(charge_micros(0.0, 20), 0);
1382 assert_eq!(charge_micros(-3.0, 20), 0);
1383 }
1384
1385 #[test]
1386 fn sandbox_time_is_charged_from_the_first_second() {
1387 // 21 millionths a second at cost, plus 20%.
1388 let price = Price::price_for(21.0, 20);
1389 assert_eq!(sandbox_charge(1, price), 26);
1390 assert_eq!(sandbox_charge(60, price), 1_512);
1391 assert_eq!(sandbox_charge(0, price), 0);
1392 }
1393
1394 #[test]
1395 fn prepaying_starts_at_twenty_five_dollars_and_bank_transfers_at_a_thousand() {
1396 assert!(prepay_amount(2_500, false).is_ok());
1397 assert!(prepay_amount(2_499, false).is_err());
1398 assert!(prepay_amount(10_000, false).is_ok());
1399 assert!(prepay_amount(1_000_000, false).is_ok());
1400 assert_eq!(prepay_amount(1_000_001, false).unwrap_err(), "Prepay between $25 and $10,000 by card; from $1,000, a bank transfer works too.");
1401 assert!(prepay_amount(99_999, true).is_err());
1402 assert!(prepay_amount(100_000, true).is_ok());
1403 }
1404
1405 #[test]
1406 fn durations_read_plainly() {
1407 assert_eq!(duration(40), "40s");
1408 assert_eq!(duration(192), "3m 12s");
1409 assert_eq!(duration(3720), "1h 2m");
1410 }
1411
1412 #[test]
1413 fn an_absurd_cost_is_capped() {
1414 assert_eq!(charge_micros(1e9, 20), 120 * MICROS_PER_DOLLAR);
1415 }
1416
1417 /// Every migration, in order, as D1 applies them.
1418 const MIGRATIONS: &[&str] = &[
1419 include_str!("../migrations/0001_init.sql"),
1420 include_str!("../migrations/0002_own_provider.sql"),
1421 include_str!("../migrations/0003_subscriptions.sql"),
1422 include_str!("../migrations/0004_sandbox_time.sql"),
1423 include_str!("../migrations/0005_limits.sql"),
1424 include_str!("../migrations/0006_prices.sql"),
1425 include_str!("../migrations/0007_pending_usage.sql"),
1426 include_str!("../migrations/0008_accounts.sql"),
1427 include_str!("../migrations/0009_autopay.sql"),
1428 include_str!("../migrations/0010_month_close.sql"),
1429 include_str!("../migrations/0011_limit_warnings.sql"),
1430 include_str!("../migrations/0012_stripe_webhooks.sql"),
1431 include_str!("../migrations/0013_invoices_trust_sales.sql"),
1432 include_str!("../migrations/0014_custom_domain_price.sql"),
1433 include_str!("../migrations/0015_sandbox_at_cost.sql"),
1434 include_str!("../migrations/0016_plans_and_pools.sql"),
1435 include_str!("../migrations/0017_flagon_comped.sql"),
1436 include_str!("../migrations/0018_one_plan.sql"),
1437 include_str!("../migrations/0019_closed_workspaces.sql"),
1438 include_str!("../migrations/0020_syntaqx_standard.sql"),
1439 include_str!("../migrations/0021_no_quotas.sql"),
1440 include_str!("../migrations/0022_costs_and_margin.sql"),
1441 include_str!("../migrations/0023_one_operation_mapping.sql"),
1442 include_str!("../migrations/0024_spend_caps.sql"),
1443 include_str!("../migrations/0025_run_tier.sql"),
1444 include_str!("../migrations/0026_ledger_by_workspace_time.sql"),
1445 include_str!("../migrations/0027_stripe_cache.sql"),
1446 include_str!("../migrations/0028_drop_stripe_webhooks.sql"),
1447 include_str!("../migrations/0029_audit_retention.sql"),
1448 include_str!("../migrations/0030_token_usage.sql"),
1449 include_str!("../migrations/0031_package_storage.sql"),
1450 include_str!("../migrations/0032_workspace_value.sql"),
1451 include_str!("../migrations/0033_given_away.sql"),
1452 include_str!("../migrations/0034_given_by_why.sql"),
1453 include_str!("../migrations/0035_cf_subscriptions.sql"),
1454 include_str!("../migrations/0036_model_costs_in_full.sql"),
1455 include_str!("../migrations/0037_security_activation.sql"),
1456 include_str!("../migrations/0038_staff_credits.sql"),
1457 include_str!("../migrations/0039_discounts_not_comped.sql"),
1458 include_str!("../migrations/0040_ai_credit.sql"),
1459 ];
1460
1461 /// The columns of `table` after the migrations: each with whether an
1462 /// insert must give it (NOT NULL, no default).
1463 fn columns(table: &str) -> Vec<(String, bool)> {
1464 let mut out = vec![];
1465 for sql in MIGRATIONS {
1466 let sql: String = sql.lines().map(|l| l.split("--").next().unwrap_or("")).collect::<Vec<_>>().join("\n");
1467 if let Some(start) = sql.find(&format!("CREATE TABLE {table} (")) {
1468 let body = &sql[start + format!("CREATE TABLE {table} (").len()..];
1469 let body = &body[..body.find(");").unwrap()];
1470 for column in body.split(",\n") {
1471 let column = column.trim();
1472 let name = column.split_whitespace().next().unwrap_or("").to_owned();
1473 if name.is_empty() || name == "PRIMARY" {
1474 continue;
1475 }
1476 let required = column.contains("NOT NULL") && !column.contains("DEFAULT") && !column.contains("PRIMARY KEY");
1477 out.push((name, required || column.contains("PRIMARY KEY")));
1478 }
1479 }
1480 for statement in sql.split(';') {
1481 let s = statement.trim();
1482 if let Some(rest) = s.strip_prefix(&format!("ALTER TABLE {table} ADD COLUMN ")) {
1483 let name = rest.split_whitespace().next().unwrap().to_owned();
1484 out.push((name, rest.contains("NOT NULL") && !rest.contains("DEFAULT")));
1485 }
1486 }
1487 }
1488 out
1489 }
1490
1491 /// The columns an `INSERT INTO table (…)` names.
1492 fn inserted(sql: &str, table: &str) -> Vec<String> {
1493 let start = sql.find(&format!("INSERT INTO {table} (")).unwrap() + format!("INSERT INTO {table} (").len();
1494 sql[start..start + sql[start..].find(')').unwrap()].split(',').map(|c| c.trim().to_owned()).collect()
1495 }
1496
1497
1498 #[test]
1499 fn ai_prices_are_price_book_data_with_dated_versions() {
1500 let sql = include_str!("../migrations/0040_ai_credit.sql");
1501 let row = |needle: &str| sql.lines().find(|l| l.contains(needle)).unwrap_or_else(|| panic!("no {needle}")).to_owned();
1502 // AI Gateway: the provider's price, no markup while in beta.
1503 assert!(row("('gateway_models', 'AI Gateway models'").contains("1000000, 0, 'list'"));
1504 // Models lose their markup from 2026-10-08, a fall that applies at once…
1505 assert!(row("('pv_agent_models_2'").contains("1000000, 0, '2026-10-08"));
1506 // …and the agent rate is $0.25 a million tokens after 14 days' notice.
1507 assert!(row("('pv_agent_tokens_2'").contains("250000, 0, '2026-10-22"));
1508 // The card fee is 2.9% + 30¢, switchable.
1509 assert!(row("('card_fee_percent'").contains("29000"));
1510 assert!(row("('card_fee_fixed'").contains("300000"));
1511 assert!(row("('card_fee', 'on'").contains("'on'"));
1512 }
1513 #[test]
1514 fn every_checkout_insert_fills_the_table() {
1515 // The plan's, the activation's, a card check's, a prepayment's and
1516 // AI credit's pages all go through one insert, and it names every
1517 // column the table needs and none it lacks: a missing NOT NULL
1518 // column, or one from a migration not yet applied, was a 500.
1519 let table = columns("checkouts");
1520 let names: Vec<&str> = table.iter().map(|(n, _)| n.as_str()).collect();
1521 assert_eq!(names, ["id", "workspace", "amount_cents", "created_by", "status", "created_at", "feature", "fee_cents"]);
1522 let given = inserted(CHECKOUT_INSERT, "checkouts");
1523 for column in &given {
1524 assert!(names.contains(&column.as_str()), "the insert names {column}, which checkouts does not have");
1525 }
1526 for (column, required) in &table {
1527 assert!(!required || given.contains(column), "checkouts needs {column}, which the insert leaves out");
1528 }
1529 // As many values as columns.
1530 let values = CHECKOUT_INSERT.split("VALUES").nth(1).unwrap();
1531 assert_eq!(values.matches('?').count() + values.matches("'open'").count(), given.len());
1532 // And no page writes a row of its own any more.
1533 for (file, source) in [
1534 ("lib.rs", include_str!("lib.rs")),
1535 ("features.rs", include_str!("features.rs")),
1536 ("cards.rs", include_str!("cards.rs")),
1537 ("ai.rs", include_str!("ai.rs")),
1538 ] {
1539 let inserts = source.matches(concat!("INSERT INTO ", "checkouts")).count();
1540 assert!(inserts <= usize::from(file == "lib.rs"), "{file} inserts into checkouts on its own");
1541 }
1542 }
1543
1544 #[test]
1545 fn a_new_ledger_column_has_a_default() {
1546 // The ledger is written from a dozen places; a column without a
1547 // default would break every one that does not name it.
1548 for (column, required) in columns("ledger") {
1549 assert!(!required || ["id", "workspace", "kind", "amount_micros", "description", "reference", "created_at"].contains(&column.as_str()), "{column}");
1550 }
1551 }
1552}