pr_01m47d24b0e6n91zwymwxg0vpx/services/billing/src/lib.rs

1,082 lines43,119 bytesCodeBlame
1//! The billing service: what agents cost, charged to the workspace they
2//! worked for.
3//!
4//! A workspace buys credit with a card. Before the runner starts an agent
5//! it asks here, and is refused if the workspace has none. When the
6//! agent's sandbox finishes it reports what the model cost, and that plus
7//! g1t's margin comes off the balance. Every change is a ledger entry, and
8//! a balance is always the sum of its ledger.
9//!
10//! Paid features (deployments) are bought separately, as monthly plans;
11//! see `features`. They are never free.
12//!
13//! Without a card processor configured the service says so and charges
14//! nothing, so that g1t still runs where billing has not been set up.
15//!
16//! Reached only through service bindings; see `g1t_contracts::billing` for
17//! the methods and their arguments.
18
19mod accounts;
20mod invoices;
21mod sales;
22mod webhooks;
23mod features;
24mod keeper;
25mod limits;
26mod stripe;
27
28use g1t_contracts::billing::*;
29use g1t_contracts::time::rfc3339;
30use g1t_contracts::{FailureCode, Outcome, Role, new_id};
31use g1t_contracts::billing::TermsKind;
32use g1t_kit::{args, now_ms, reply, rpc_method};
33use serde::Deserialize;
34use sha2::{Digest, Sha256};
35use worker::wasm_bindgen::JsValue;
36use worker::{Context, D1Database, Env, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
37
38use stripe::Stripe;
39
40const MIN_TOP_UP_CENTS: u32 = 500;
41const MAX_TOP_UP_CENTS: u32 = 50_000;
42const LEDGER_PAGE: u32 = 100;
43/// A run's reported cost is believed up to this much. A sandbox cannot
44/// spend more in the time it has, so anything above is a fault.
45const MAX_RUN_COST_USD: f64 = 100.0;
46
47/// What a run is charged: its cost plus the margin, rounded up to a whole
48/// millionth of a dollar.
49pub fn charge_micros(cost_usd: f64, margin_percent: u32) -> i64 {
50 let cost_micros = (cost_usd.clamp(0.0, MAX_RUN_COST_USD) * MICROS_PER_DOLLAR as f64).ceil();
51 (cost_micros * f64::from(100 + margin_percent) / 100.0).ceil() as i64
52}
53
54fn hash(token: &str) -> String {
55 hex::encode(Sha256::digest(token.as_bytes()))
56}
57
58fn optional(value: Option<&str>) -> JsValue {
59 value.map_or(JsValue::NULL, JsValue::from)
60}
61
62#[derive(Deserialize)]
63struct AccountRow {
64 balance_micros: i64,
65 customer_id: Option<String>,
66}
67
68#[derive(Deserialize)]
69struct LedgerRow {
70 id: String,
71 kind: EntryKind,
72 amount_micros: i64,
73 description: String,
74 repo: Option<String>,
75 number: Option<u32>,
76 task: Option<String>,
77 model: Option<String>,
78 created_by: Option<String>,
79 created_at: String,
80 billed_to: Option<String>,
81 #[serde(default)]
82 workspace: Option<String>,
83}
84
85impl From<LedgerRow> for LedgerEntry {
86 fn from(row: LedgerRow) -> Self {
87 LedgerEntry {
88 id: row.id,
89 kind: row.kind,
90 amount_micros: row.amount_micros,
91 description: row.description,
92 repo: row.repo,
93 number: row.number,
94 task: row.task,
95 model: row.model,
96 billed_to: row.billed_to.unwrap_or_else(|| "g1t".to_owned()),
97 created_by: row.created_by,
98 created_at: row.created_at,
99 workspace: row.workspace,
100 }
101 }
102}
103
104#[derive(Deserialize)]
105struct RunRow {
106 workspace: String,
107 repo: String,
108 number: u32,
109 task: String,
110 model: String,
111 token_hash: String,
112 billed_to: Option<String>,
113}
114
115impl RunRow {
116 fn own_provider(&self) -> bool {
117 self.billed_to.as_deref() == Some("workspace")
118 }
119}
120
121#[derive(Deserialize)]
122struct CheckoutRow {
123 workspace: String,
124 created_by: String,
125}
126
127/// A row an `UPDATE … RETURNING` touched.
128#[derive(Deserialize)]
129struct Touched {
130 #[allow(dead_code)]
131 id: String,
132}
133
134struct Billing {
135 db: D1Database,
136 /// Absent when no card processor is configured.
137 stripe: Option<Stripe>,
138 margin_percent: u32,
139 /// Charged for a run on the workspace's own model provider.
140 orchestration_fee_micros: i64,
141 /// While g1t is being built out, nothing is charged (`FREE_WHILE_BUILDING`).
142 free: bool,
143 /// The free allowance on g1t's hosted models, when there is one.
144 trial: Option<TrialConfig>,
145 /// The Deployments plan's monthly price (`DEPLOYMENTS_MONTHLY_CENTS`).
146 deployments_monthly_cents: u32,
147 /// How far unpaid usage may go; see `limits`.
148 ceilings: limits::Ceilings,
149 /// `PREPAID_ONLY`: the old rule, that agents need credit first.
150 prepaid_only: bool,
151}
152
153/// `TRIAL_WORKSPACE_MICROS`, `TRIAL_TOTAL_MICROS` and `TRIAL_UNTIL`.
154struct TrialConfig {
155 per_workspace_micros: i64,
156 total_micros: i64,
157 /// RFC 3339, in UTC.
158 until: String,
159}
160
161#[derive(serde::Deserialize)]
162struct Sum {
163 micros: Option<i64>,
164}
165
166impl Billing {
167 fn status(&self) -> Status {
168 Status {
169 enabled: self.stripe.is_some(),
170 live: self.stripe.as_ref().is_some_and(Stripe::live),
171 free: self.free,
172 }
173 }
174
175 async fn row(&self, workspace: &str) -> Result<Option<AccountRow>> {
176 self.db
177 .prepare("SELECT balance_micros, customer_id FROM accounts WHERE workspace = ?")
178 .bind(&[workspace.into()])?
179 .first::<AccountRow>(None)
180 .await
181 }
182
183 async fn standing(&self, workspace: &str) -> Result<Account> {
184 let row = self.row(workspace).await?;
185 let card = match (&self.stripe, row.as_ref().and_then(|row| row.customer_id.as_deref())) {
186 (Some(stripe), Some(customer)) => stripe.card(customer).await.ok().flatten().map(|card| Card {
187 brand: card.brand,
188 last4: card.last4,
189 exp_month: card.exp_month,
190 exp_year: card.exp_year,
191 }),
192 _ => None,
193 };
194 Ok(Account {
195 workspace: workspace.to_owned(),
196 balance_micros: row.map_or(0, |row| row.balance_micros),
197 status: self.status(),
198 margin_percent: self.margin_percent,
199 orchestration_fee_micros: self.orchestration_fee_micros,
200 card,
201 })
202 }
203
204 /// The workspace's customer at Stripe, made the first time one is needed.
205 pub(crate) async fn customer_for(&self, workspace: &str) -> Result<String> {
206 let Some(stripe) = &self.stripe else {
207 return Err(worker::Error::RustError("payments are not set up".into()));
208 };
209 if let Some(customer) = self.row(workspace).await?.and_then(|row| row.customer_id) {
210 return Ok(customer);
211 }
212 let customer = stripe.create_customer(workspace).await?;
213 self.db
214 .prepare(
215 "INSERT INTO accounts (workspace, balance_micros, customer_id, created_at) VALUES (?1, 0, ?2, ?3)
216 ON CONFLICT (workspace) DO UPDATE SET customer_id = ?2",
217 )
218 .bind(&[workspace.into(), customer.as_str().into(), rfc3339(now_ms()).into()])?
219 .run()
220 .await?;
221 Ok(customer)
222 }
223
224 /// Stripe's hosted billing page for the workspace. Owners only.
225 async fn billing_portal(&self, a: BillingPortalArgs) -> Result<Outcome<Checkout>> {
226 let workspace = a.workspace.to_lowercase();
227 if a.actor.role_in(&workspace) != Some(Role::Owner) {
228 return Ok(Outcome::fail(FailureCode::Forbidden, "Only an owner can manage the workspace's billing."));
229 }
230 let Some(stripe) = &self.stripe else {
231 return Ok(Outcome::fail(FailureCode::Conflict, "Payments are not set up on this g1t."));
232 };
233 let customer = match self.customer_for(&workspace).await {
234 Ok(customer) => customer,
235 Err(error) => return Ok(Outcome::fail(FailureCode::Conflict, format!("Stripe could not be reached: {error}"))),
236 };
237 match stripe.portal_session(&customer, &a.return_url).await {
238 Ok(url) => Ok(Outcome::Ok(Checkout { url })),
239 Err(error) if stripe::is_missing(&error) => {
240 // The customer was removed at Stripe: a new one next time.
241 self.forget_customer(&workspace).await?;
242 Ok(Outcome::fail(FailureCode::Conflict, "Stripe no longer had this workspace's customer. Try again."))
243 }
244 Err(error) => Ok(Outcome::fail(FailureCode::Conflict, format!("Stripe's billing page could not be opened: {error}"))),
245 }
246 }
247
248 /// Adds a ledger entry and moves the balance by the same amount, as
249 /// one write.
250 #[allow(clippy::too_many_arguments)]
251 async fn enter(
252 &self,
253 workspace: &str,
254 kind: EntryKind,
255 amount_micros: i64,
256 description: &str,
257 reference: &str,
258 run: Option<&RunRow>,
259 cost_micros: Option<i64>,
260 created_by: Option<&str>,
261 customer: Option<&str>,
262 ) -> Result<()> {
263 let now = now_ms();
264 let timestamp = rfc3339(now);
265 let kind = match kind {
266 EntryKind::TopUp => "top_up",
267 EntryKind::Usage => "usage",
268 };
269 self.db
270 .batch(vec![
271 self.db
272 .prepare(
273 "INSERT INTO ledger
274 (id, workspace, kind, amount_micros, description, repo, number, task,
275 model, cost_micros, reference, created_by, created_at, billed_to)
276 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
277 )
278 .bind(&[
279 new_id("led", now).into(),
280 workspace.into(),
281 kind.into(),
282 // D1 takes numbers as doubles, which hold every
283 // amount this service will see exactly.
284 (amount_micros as f64).into(),
285 description.into(),
286 optional(run.map(|run| run.repo.as_str())),
287 run.map_or(JsValue::NULL, |run| run.number.into()),
288 optional(run.map(|run| run.task.as_str())),
289 optional(run.map(|run| run.model.as_str())),
290 cost_micros.map_or(JsValue::NULL, |cost| (cost as f64).into()),
291 reference.into(),
292 optional(created_by),
293 timestamp.as_str().into(),
294 run.map_or("g1t", |run| if run.own_provider() { "workspace" } else { "g1t" }).into(),
295 ])?,
296 self.db
297 .prepare(
298 "INSERT INTO accounts (workspace, balance_micros, customer_id, created_at)
299 VALUES (?1, ?2, ?3, ?4)
300 ON CONFLICT (workspace) DO UPDATE SET
301 balance_micros = balance_micros + ?2,
302 customer_id = COALESCE(?3, customer_id)",
303 )
304 .bind(&[
305 workspace.into(),
306 (amount_micros as f64).into(),
307 optional(customer),
308 timestamp.as_str().into(),
309 ])?,
310 ])
311 .await?;
312 // Money in clears a card declined at the limit.
313 if kind == "top_up" {
314 self.db
315 .prepare("UPDATE limits SET autopay_failed_at = NULL, autopay_error = NULL WHERE workspace = ?")
316 .bind(&[workspace.into()])?
317 .run()
318 .await?;
319 }
320 Ok(())
321 }
322
323 async fn account(&self, a: AccountArgs) -> Result<Outcome<Account>> {
324 let workspace = a.workspace.to_lowercase();
325 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
326 return Ok(members_only());
327 }
328 Ok(Outcome::Ok(self.standing(&workspace).await?))
329 }
330
331 async fn ledger(&self, a: AccountArgs) -> Result<Outcome<Vec<LedgerEntry>>> {
332 let workspace = a.workspace.to_lowercase();
333 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
334 return Ok(members_only());
335 }
336 let rows = self
337 .db
338 .prepare("SELECT * FROM ledger WHERE workspace = ? ORDER BY id DESC LIMIT ?")
339 .bind(&[workspace.into(), LEDGER_PAGE.into()])?
340 .all()
341 .await?
342 .results::<LedgerRow>()?;
343 Ok(Outcome::Ok(
344 rows.into_iter().map(LedgerEntry::from).collect(),
345 ))
346 }
347
348 async fn usage(&self, a: UsageArgs) -> Result<Outcome<Usage>> {
349 let workspace = a.workspace.to_lowercase();
350 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
351 return Ok(members_only());
352 }
353 #[derive(serde::Deserialize)]
354 struct SliceRow {
355 key: Option<String>,
356 micros: Option<i64>,
357 runs: Option<u32>,
358 }
359 // While nothing is charged, what was used is what there is to show.
360 let measure = if self.free { "COALESCE(cost_micros, 0)" } else { "-amount_micros" };
361 let slices = |key: &str, limit: u32| {
362 format!(
363 "SELECT {key} AS key, SUM({measure}) AS micros, COUNT(*) AS runs FROM ledger
364 WHERE workspace = ?1 AND kind = 'usage' AND created_at >= ?2
365 GROUP BY 1 ORDER BY micros DESC LIMIT {limit}"
366 )
367 };
368 let query = |sql: String| {
369 let db = &self.db;
370 let workspace = workspace.clone();
371 let since = a.since.clone();
372 async move {
373 let rows = db
374 .prepare(sql)
375 .bind(&[workspace.into(), since.into()])?
376 .all()
377 .await?
378 .results::<SliceRow>()?;
379 Ok::<Vec<UsageSlice>, worker::Error>(
380 rows.into_iter()
381 .map(|row| UsageSlice {
382 key: row.key.unwrap_or_else(|| "other".to_owned()),
383 micros: row.micros.unwrap_or_default(),
384 runs: row.runs.unwrap_or_default(),
385 })
386 .collect(),
387 )
388 }
389 };
390 #[derive(serde::Deserialize)]
391 struct Totals {
392 spent: Option<i64>,
393 cost: Option<i64>,
394 provider: Option<i64>,
395 runs: Option<u32>,
396 added: Option<i64>,
397 }
398 let totals = self
399 .db
400 .prepare(
401 "SELECT
402 -SUM(CASE WHEN kind = 'usage' THEN amount_micros END) AS spent,
403 SUM(CASE WHEN kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' THEN cost_micros END) AS cost,
404 SUM(CASE WHEN kind = 'usage' AND billed_to = 'workspace' THEN cost_micros END) AS provider,
405 SUM(CASE WHEN kind = 'usage' THEN 1 ELSE 0 END) AS runs,
406 SUM(CASE WHEN kind = 'top_up' THEN amount_micros END) AS added
407 FROM ledger WHERE workspace = ?1 AND created_at >= ?2",
408 )
409 .bind(&[workspace.as_str().into(), a.since.as_str().into()])?
410 .first::<Totals>(None)
411 .await?;
412 let totals = totals.unwrap_or(Totals {
413 spent: None,
414 cost: None,
415 provider: None,
416 runs: None,
417 added: None,
418 });
419 Ok(Outcome::Ok(Usage {
420 spent_micros: totals.spent.unwrap_or_default(),
421 cost_micros: totals.cost.unwrap_or_default(),
422 provider_micros: totals.provider.unwrap_or_default(),
423 used_micros: totals.cost.unwrap_or_default() + totals.provider.unwrap_or_default(),
424 free: self.free,
425 runs: totals.runs.unwrap_or_default(),
426 added_micros: totals.added.unwrap_or_default(),
427 by_day: query(slices("substr(created_at, 1, 10) || '/' || COALESCE(task, 'other')", 400)).await?,
428 by_task: query(slices("task", 20)).await?,
429 by_repo: query(slices("repo", 20)).await?,
430 by_pull: query(slices("repo || '#' || number", 10)).await?,
431 by_model: query(slices("model", 10)).await?,
432 since: a.since,
433 }))
434 }
435
436 async fn checkout(&self, a: CheckoutArgs) -> Result<Outcome<Checkout>> {
437 let workspace = a.workspace.to_lowercase();
438 if a.actor.role_in(&workspace) != Some(Role::Owner) {
439 return Ok(Outcome::fail(
440 FailureCode::Forbidden,
441 "Only an owner can add credit to a workspace.",
442 ));
443 }
444 let Some(stripe) = &self.stripe else {
445 return Ok(Outcome::fail(
446 FailureCode::Conflict,
447 "Payments are not set up on this g1t yet.",
448 ));
449 };
450 if !(MIN_TOP_UP_CENTS..=MAX_TOP_UP_CENTS).contains(&a.amount_cents) {
451 return Ok(Outcome::fail(
452 FailureCode::Invalid,
453 format!(
454 "Add between ${} and ${} at a time.",
455 MIN_TOP_UP_CENTS / 100,
456 MAX_TOP_UP_CENTS / 100
457 ),
458 ));
459 }
460 let customer = self.row(&workspace).await?.and_then(|row| row.customer_id);
461 let session = match stripe
462 .start_checkout(&workspace, a.amount_cents, customer.as_deref(), &a.return_url)
463 .await
464 {
465 Ok(session) => session,
466 // A customer saved under another Stripe account: start afresh.
467 Err(error) if customer.is_some() && stripe::is_missing(&error) => {
468 self.forget_customer(&workspace).await?;
469 stripe.start_checkout(&workspace, a.amount_cents, None, &a.return_url).await?
470 }
471 Err(error) => return Err(error),
472 };
473 let Some(url) = session.url else {
474 return Err(worker::Error::RustError(
475 "the card processor returned no payment page".into(),
476 ));
477 };
478 self.db
479 .prepare(
480 "INSERT INTO checkouts (id, workspace, amount_cents, created_by, created_at)
481 VALUES (?, ?, ?, ?, ?)",
482 )
483 .bind(&[
484 session.id.into(),
485 workspace.into(),
486 a.amount_cents.into(),
487 a.actor.username.into(),
488 rfc3339(now_ms()).into(),
489 ])?
490 .run()
491 .await?;
492 Ok(Outcome::Ok(Checkout { url }))
493 }
494
495 /// Credits a payment if the processor says it was made and it has not
496 /// been credited before. The amount credited is what the processor
497 /// says was paid, not what anyone here remembers asking for.
498 async fn confirm(&self, a: ConfirmArgs) -> Result<Outcome<Account>> {
499 let workspace = a.workspace.to_lowercase();
500 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
501 return Ok(members_only());
502 }
503 let (Some(stripe), Some(checkout)) = (
504 &self.stripe,
505 self.db
506 .prepare(
507 "SELECT workspace, created_by FROM checkouts
508 WHERE id = ? AND workspace = ? AND status = 'open'",
509 )
510 .bind(&[a.session.as_str().into(), workspace.as_str().into()])?
511 .first::<CheckoutRow>(None)
512 .await?,
513 ) else {
514 // Unknown, someone else's, or already credited: nothing to do.
515 return Ok(Outcome::Ok(self.standing(&workspace).await?));
516 };
517 let session = stripe.session(&a.session).await?;
518 let paid = session
519 .amount_total
520 .filter(|_| session.payment_status == "paid");
521 if let Some(cents) = paid {
522 // Only whoever flips it from open to paid enters the credit.
523 let claimed = self
524 .db
525 .prepare(
526 "UPDATE checkouts SET status = 'paid' WHERE id = ? AND status = 'open'
527 RETURNING id",
528 )
529 .bind(&[a.session.as_str().into()])?
530 .first::<Touched>(None)
531 .await?;
532 if claimed.is_some() {
533 self.enter(
534 &checkout.workspace,
535 EntryKind::TopUp,
536 i64::from(cents) * MICROS_PER_DOLLAR / 100,
537 "Credit added by card",
538 &session.id,
539 None,
540 None,
541 Some(&checkout.created_by),
542 session.customer.as_deref(),
543 )
544 .await?;
545 }
546 }
547 Ok(Outcome::Ok(self.standing(&workspace).await?))
548 }
549
550 /// Drops a saved customer the card processor no longer knows.
551 pub(crate) async fn forget_customer(&self, workspace: &str) -> Result<()> {
552 self.db
553 .prepare("UPDATE accounts SET customer_id = NULL WHERE workspace = ?")
554 .bind(&[workspace.into()])?
555 .run()
556 .await?;
557 Ok(())
558 }
559
560 /// A refusal if the workspace has no credit to start an agent with.
561 async fn out_of_credit<T>(&self, workspace: &str) -> Result<Option<Outcome<T>>> {
562 // Billing is postpaid: usage limits decide whether work starts
563 // (see `limits`), and credit is a prepayment that lowers what is
564 // owed. A balance no longer has to be positive to start.
565 if self.free || !self.prepaid_only {
566 return Ok(None);
567 }
568 let balance = self
569 .row(workspace)
570 .await?
571 .map_or(0, |row| row.balance_micros);
572 Ok((balance <= 0).then(|| {
573 Outcome::fail(
574 FailureCode::PaymentRequired,
575 format!(
576 "The {workspace} workspace has no agent credit. An owner can add some under Billing on the workspace's page."
577 ),
578 )
579 }))
580 }
581
582 /// A workspace's free allowance on g1t's hosted models: what its runs
583 /// there have cost against its share, and the pool everyone draws on.
584 /// How much of a hosted model run's cost the workspace's free allowance
585 /// covers, if it is still open.
586 async fn trial_covers(&self, workspace: &str, cost_micros: i64) -> Result<i64> {
587 let trial = self.trial(TrialArgs { workspace: workspace.to_owned(), exempt: vec![] }).await?;
588 if !trial.open {
589 return Ok(0);
590 }
591 Ok(cost_micros.min((trial.limit_micros - trial.used_micros).max(0)))
592 }
593
594 async fn trial(&self, a: TrialArgs) -> Result<Trial> {
595 let workspace = a.workspace.to_lowercase();
596 let Some(config) = &self.trial else {
597 return Ok(Trial {
598 open: false,
599 used_micros: 0,
600 limit_micros: 0,
601 ends_at: None,
602 reason: Some("off".to_owned()),
603 });
604 };
605 let used = self
606 .db
607 .prepare(
608 "SELECT SUM(cost_micros) AS micros FROM ledger
609 WHERE kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' AND COALESCE(task, '') NOT IN ('sandbox', 'deployments') AND workspace = ?",
610 )
611 .bind(&[workspace.as_str().into()])?
612 .first::<Sum>(None)
613 .await?
614 .and_then(|sum| sum.micros)
615 .unwrap_or_default();
616 // Everyone's, but for the workspaces open to hosted models anyway.
617 let exempt: Vec<String> = a.exempt.iter().map(|name| name.trim().to_lowercase()).collect();
618 let marks = vec!["?"; exempt.len().max(1)].join(", ");
619 let mut values: Vec<JsValue> = exempt.iter().map(|name| JsValue::from(name.as_str())).collect();
620 if values.is_empty() {
621 values.push(JsValue::from(""));
622 }
623 let pooled = self
624 .db
625 .prepare(format!(
626 "SELECT SUM(cost_micros) AS micros FROM ledger
627 WHERE kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' AND COALESCE(task, '') NOT IN ('sandbox', 'deployments') AND workspace NOT IN ({marks})"
628 ))
629 .bind(&values)?
630 .first::<Sum>(None)
631 .await?
632 .and_then(|sum| sum.micros)
633 .unwrap_or_default();
634 let reason = if rfc3339(now_ms()) >= config.until {
635 Some("ended")
636 } else if used >= config.per_workspace_micros {
637 Some("used")
638 } else if pooled >= config.total_micros {
639 Some("pool")
640 } else {
641 None
642 };
643 Ok(Trial {
644 open: reason.is_none(),
645 used_micros: used,
646 limit_micros: config.per_workspace_micros,
647 ends_at: Some(config.until.clone()),
648 reason: reason.map(str::to_owned),
649 })
650 }
651
652 async fn can_start(&self, a: CanStartArgs) -> Result<Outcome<bool>> {
653 if self.stripe.is_none() {
654 return Ok(Outcome::Ok(true));
655 }
656 if let Some(stopped) = self.stopped(&a.workspace).await? {
657 return Ok(stopped);
658 }
659 Ok(self
660 .out_of_credit(&a.workspace.to_lowercase())
661 .await?
662 .unwrap_or(Outcome::Ok(true)))
663 }
664
665 async fn start_run(&self, a: StartRunArgs) -> Result<Outcome<Option<RunTicket>>> {
666 if self.stripe.is_none() {
667 return Ok(Outcome::Ok(None));
668 }
669 let workspace = a.workspace.to_lowercase();
670 if let Some(stopped) = self.stopped(&workspace).await? {
671 return Ok(stopped);
672 }
673 if let Some(refused) = self.out_of_credit(&workspace).await? {
674 return Ok(refused);
675 }
676 let now = now_ms();
677 let run_id = new_id("run", now);
678 let mut bytes = [0u8; 32];
679 getrandom::getrandom(&mut bytes).expect("no source of randomness");
680 let token = hex::encode(bytes);
681 self.db
682 .prepare(
683 "INSERT INTO runs (id, workspace, repo, number, task, model, token_hash, created_at, billed_to, session_id)
684 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
685 )
686 .bind(&[
687 run_id.as_str().into(),
688 workspace.into(),
689 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
690 a.number.into(),
691 a.task.into(),
692 a.model.into(),
693 hash(&token).into(),
694 rfc3339(now).into(),
695 if a.billed_to == "workspace" { "workspace" } else { "g1t" }.into(),
696 optional(a.session.as_deref().filter(|_| a.billed_to != "workspace")),
697 ])?
698 .run()
699 .await?;
700 Ok(Outcome::Ok(Some(RunTicket { run_id, token })))
701 }
702
703 async fn finish_run(&self, a: FinishRunArgs) -> Result<Outcome<bool>> {
704 let run = self
705 .db
706 .prepare(
707 "SELECT workspace, repo, number, task, model, token_hash, billed_to FROM runs
708 WHERE id = ? AND finished_at IS NULL",
709 )
710 .bind(&[a.run_id.as_str().into()])?
711 .first::<RunRow>(None)
712 .await?;
713 let Some(run) = run.filter(|run| run.token_hash == hash(&a.token)) else {
714 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
715 };
716 if !a.cost_usd.is_finite() || a.cost_usd < 0.0 {
717 return Ok(Outcome::fail(FailureCode::Invalid, "That is not a cost."));
718 }
719 // Only whoever closes the run charges for it.
720 let claimed = self
721 .db
722 .prepare(
723 "UPDATE runs SET finished_at = ? WHERE id = ? AND finished_at IS NULL RETURNING id",
724 )
725 .bind(&[rfc3339(now_ms()).into(), a.run_id.as_str().into()])?
726 .first::<Touched>(None)
727 .await?;
728 if claimed.is_none() {
729 return Ok(Outcome::Ok(false));
730 }
731 // On the workspace's own provider, the model was paid for there:
732 // g1t charges its fee, and keeps the provider's cost to show.
733 let base = if run.own_provider() {
734 self.orchestration_fee_micros
735 } else {
736 // The free allowance on g1t's models covers what it can.
737 let cost = charge_micros(a.cost_usd, 0);
738 let covered = self.trial_covers(&run.workspace, cost).await?;
739 charge_micros((cost - covered) as f64 / MICROS_PER_DOLLAR as f64, self.margin_percent)
740 };
741 let (charge, terms_note) = self.charged(&run.workspace, base).await?;
742 let mut description = match run.task.as_str() {
743 "plan" => format!("Planning for {}", run.repo),
744 "review" => format!("Review of {}#{}", run.repo, run.number),
745 "update" => format!("Catching up {}#{}", run.repo, run.number),
746 _ => format!("Work on {}#{}", run.repo, run.number),
747 };
748 if run.own_provider() {
749 description.push_str(", on your own model provider");
750 }
751 description.push_str(&terms_note);
752 self.enter(
753 &run.workspace,
754 EntryKind::Usage,
755 -charge,
756 &description,
757 &a.run_id,
758 Some(&run),
759 Some(charge_micros(a.cost_usd, 0)),
760 None,
761 None,
762 )
763 .await?;
764 Ok(Outcome::Ok(true))
765 }
766}
767
768impl Billing {
769 /// Records how long a sandbox ran: its cost always, and a charge for
770 /// the seconds past the month's free minutes.
771 async fn record_sandbox(&self, a: RecordSandboxArgs) -> Result<Outcome<bool>> {
772 if self.stripe.is_none() || a.seconds == 0 {
773 return Ok(Outcome::Ok(false));
774 }
775 let workspace = a.workspace.to_lowercase();
776 let seen = self
777 .db
778 .prepare("SELECT id FROM ledger WHERE reference = ?")
779 .bind(&[a.reference.as_str().into()])?
780 .first::<Touched>(None)
781 .await?;
782 if seen.is_some() {
783 return Ok(Outcome::Ok(false));
784 }
785 let now = now_ms();
786 let timestamp = rfc3339(now);
787 let month = &timestamp[..7];
788 #[derive(Deserialize)]
789 struct Used {
790 seconds: i64,
791 }
792 let seconds = i64::from(a.seconds);
793 let after = self
794 .db
795 .prepare(
796 "INSERT INTO sandbox_months (workspace, month, seconds) VALUES (?1, ?2, ?3)
797 ON CONFLICT (workspace, month) DO UPDATE SET seconds = seconds + ?3
798 RETURNING seconds",
799 )
800 .bind(&[workspace.as_str().into(), month.into(), (seconds as f64).into()])?
801 .first::<Used>(None)
802 .await?
803 .map_or(seconds, |used| used.seconds);
804 let billable = sandbox_billable(after - seconds, seconds);
805 // From the price book, which follows what Cloudflare bills g1t.
806 let (cost_per_second, price_per_second) = self.price("sandbox_second").await?.unwrap_or((
807 sandbox_allowance::COST_MICROS_PER_SECOND as f64,
808 sandbox_allowance::MICROS_PER_SECOND as f64,
809 ));
810 let (charge, terms_note) = self.charged(&workspace, (billable as f64 * price_per_second).ceil() as i64).await?;
811 let mut description = format!("{}: {} of sandbox time", a.description, duration(seconds));
812 if billable < seconds {
813 description.push_str(if billable == 0 {
814 ", within the month's free minutes"
815 } else {
816 ", partly within the month's free minutes"
817 });
818 }
819 if billable > 0 {
820 description.push_str(&terms_note);
821 }
822 self.db
823 .batch(vec![
824 self.db
825 .prepare(
826 "INSERT INTO ledger
827 (id, workspace, kind, amount_micros, description, repo, task,
828 cost_micros, reference, created_at, billed_to)
829 VALUES (?, ?, 'usage', ?, ?, ?, 'sandbox', ?, ?, ?, 'g1t')",
830 )
831 .bind(&[
832 new_id("led", now).into(),
833 workspace.as_str().into(),
834 (-(charge as f64)).into(),
835 description.as_str().into(),
836 optional(a.repo.as_deref()),
837 (seconds as f64 * cost_per_second).ceil().into(),
838 a.reference.as_str().into(),
839 timestamp.as_str().into(),
840 ])?,
841 self.db
842 .prepare(
843 "INSERT INTO accounts (workspace, balance_micros, created_at)
844 VALUES (?1, ?2, ?3)
845 ON CONFLICT (workspace) DO UPDATE SET balance_micros = balance_micros + ?2",
846 )
847 .bind(&[
848 workspace.as_str().into(),
849 (-(charge as f64)).into(),
850 timestamp.as_str().into(),
851 ])?,
852 ])
853 .await?;
854 Ok(Outcome::Ok(true))
855 }
856}
857
858/// Of `seconds` used after `before` this month, how many are past the
859/// free minutes.
860fn sandbox_billable(before: i64, seconds: i64) -> i64 {
861 let free_left = (sandbox_allowance::FREE_SECONDS - before).max(0);
862 (seconds - free_left).max(0)
863}
864
865/// `1h 2m`, `3m 12s` or `40s`.
866fn duration(seconds: i64) -> String {
867 let (h, m, s) = (seconds / 3600, seconds % 3600 / 60, seconds % 60);
868 if h > 0 {
869 format!("{h}h {m}m")
870 } else if m > 0 {
871 format!("{m}m {s}s")
872 } else {
873 format!("{s}s")
874 }
875}
876
877impl Billing {
878 /// What a workspace is charged for something that would be `base`:
879 /// nothing while g1t is free, or as its account's terms say. With a
880 /// note for the statement when it differs.
881 pub(crate) async fn charged(&self, workspace: &str, base: i64) -> Result<(i64, String)> {
882 if self.free {
883 return Ok((0, " (free while g1t is being built out)".to_owned()));
884 }
885 let terms = self.terms_of(workspace).await?;
886 let charge = terms.apply(base);
887 let note = match terms.kind {
888 TermsKind::Comped => " (comped)".to_owned(),
889 TermsKind::Custom if terms.discount_percent > 0 && base > 0 => format!(" ({}% off)", terms.discount_percent),
890 _ => String::new(),
891 };
892 Ok((charge, note))
893 }
894}
895
896fn members_only<T>() -> Outcome<T> {
897 Outcome::fail(
898 FailureCode::Forbidden,
899 "Only members can see a workspace's billing.",
900 )
901}
902
903impl Billing {
904 fn from_env(env: &Env) -> Result<Self> {
905 Ok(Billing {
906 db: env.d1("DB")?,
907 stripe: env
908 .secret("STRIPE_SECRET_KEY")
909 .ok()
910 .map(|key| key.to_string())
911 .filter(|key| !key.is_empty())
912 .map(Stripe::new),
913 margin_percent: env
914 .var("MARGIN_PERCENT")
915 .ok()
916 .and_then(|percent| percent.to_string().parse().ok())
917 .unwrap_or(20),
918 orchestration_fee_micros: env
919 .var("ORCHESTRATION_FEE_MICROS")
920 .ok()
921 .and_then(|fee| fee.to_string().parse().ok())
922 .unwrap_or(100_000),
923 free: env.var("FREE_WHILE_BUILDING").is_ok_and(|v| v.to_string() == "true"),
924 ceilings: limits::Ceilings::from_env(&env),
925 prepaid_only: env.var("PREPAID_ONLY").is_ok_and(|v| v.to_string() == "true"),
926 deployments_monthly_cents: env
927 .var("DEPLOYMENTS_MONTHLY_CENTS")
928 .ok()
929 .and_then(|cents| cents.to_string().parse().ok())
930 .unwrap_or(500),
931 trial: {
932 let number = |name: &str| env.var(name).ok().and_then(|v| v.to_string().parse::<i64>().ok());
933 match (
934 number("TRIAL_WORKSPACE_MICROS"),
935 number("TRIAL_TOTAL_MICROS"),
936 env.var("TRIAL_UNTIL").ok().map(|v| v.to_string()),
937 ) {
938 (Some(per_workspace_micros), Some(total_micros), Some(until))
939 if per_workspace_micros > 0 && !until.is_empty() =>
940 {
941 Some(TrialConfig {
942 per_workspace_micros,
943 total_micros,
944 until,
945 })
946 }
947 _ => None,
948 }
949 },
950 })
951 }
952}
953
954#[event(scheduled)]
955async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
956 let Ok(billing) = Billing::from_env(&env) else {
957 return;
958 };
959 let keeper = keeper::Keeper::from_env(&env);
960 if let Err(error) = billing.settle_runs(&keeper).await {
961 worker::console_error!("settling runs failed: {error}");
962 }
963 if let Err(error) = billing.autopay().await {
964 worker::console_error!("paying at the limit failed: {error}");
965 }
966 if let Err(error) = billing.close_months().await {
967 worker::console_error!("closing the month failed: {error}");
968 }
969 if let Err(error) = billing.invoice_enterprises().await {
970 worker::console_error!("invoicing enterprises failed: {error}");
971 }
972 if let Ok(identity) = env.service("IDENTITY") {
973 if let Err(error) = billing.warn_limits(&identity).await {
974 worker::console_error!("warning owners failed: {error}");
975 }
976 }
977 // Once a day, and at once if the costs were never checked: check every
978 // cost against what Cloudflare billed.
979 if event.cron() == keeper::DAILY || billing.never_checked().await.unwrap_or(false) {
980 if let Err(error) = billing.reconcile(&keeper).await {
981 worker::console_error!("checking costs against Cloudflare failed: {error}");
982 }
983 }
984}
985
986#[event(fetch)]
987async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
988 let Some(method) = rpc_method(&request) else {
989 return Response::error("Not found", 404);
990 };
991 let body: serde_json::Value = request.json().await?;
992 let billing = Billing::from_env(&env)?;
993 match method.as_str() {
994 "status" => reply(&billing.status()),
995 "account" => reply(&billing.account(args(body)?).await?),
996 "ledger" => reply(&billing.ledger(args(body)?).await?),
997 "usage" => reply(&billing.usage(args(body)?).await?),
998 "checkout" => reply(&billing.checkout(args(body)?).await?),
999 "confirm" => reply(&billing.confirm(args(body)?).await?),
1000 "can_start" => reply(&billing.can_start(args(body)?).await?),
1001 "trial" => reply(&billing.trial(args(body)?).await?),
1002 "start_run" => reply(&billing.start_run(args(body)?).await?),
1003 "finish_run" => reply(&billing.finish_run(args(body)?).await?),
1004 "features" => reply(&billing.features(args(body)?).await?),
1005 "subscribe" => reply(&billing.subscribe(args(body)?).await?),
1006 "confirm_subscription" => reply(&billing.confirm_subscription(args(body)?).await?),
1007 "cancel_subscription" => reply(&billing.cancel_subscription(args(body)?).await?),
1008 "has_feature" => reply(&billing.has_feature(args(body)?).await?),
1009 "charge_feature" => reply(&billing.charge_feature(args(body)?).await?),
1010 "record_sandbox" => reply(&billing.record_sandbox(args(body)?).await?),
1011 "limit" => reply(&billing.limit(args(body)?).await?),
1012 "check_limit" => reply(&billing.check_limit(args(body)?).await?),
1013 "set_spend_limit" => reply(&billing.set_spend_limit(args(body)?).await?),
1014 "prices" => reply(&billing.prices().await?),
1015 "billing_portal" => reply(&billing.billing_portal(args(body)?).await?),
1016 "admin_billing_link" => reply(&billing.admin_billing_link(args(body)?).await?),
1017 "admin_stripe" => reply(&billing.admin_stripe(args(body)?).await?),
1018 "admin_enterprise_billing" => reply(&billing.admin_enterprise_billing(args(body)?).await?),
1019 "admin_invoice_enterprise" => reply(&billing.admin_invoice_enterprise(args(body)?).await?),
1020 "stripe_webhook" => reply(&billing.stripe_webhook(args(body)?).await?),
1021 "invoices" => reply(&billing.invoices(args(body)?).await?),
1022 "admin_workspace_invoices" => {
1023 let a: AdminWorkspaceInvoicesArgs = args(body)?;
1024 reply(&billing.workspace_invoices(&a.workspace.to_lowercase()).await?)
1025 }
1026 "admin_signals" => reply(&billing.admin_signals(args(body)?).await?),
1027 "admin_overview" => reply(&billing.admin_overview(args(body)?).await?),
1028 "admin_sales" => reply(&billing.admin_sales(args(body)?).await?),
1029 "admin_set_sales" => reply(&billing.admin_set_sales(args(body)?).await?),
1030 "admin_add_note" => reply(&billing.admin_add_note(args(body)?).await?),
1031 "admin_invoices" => reply(&billing.admin_invoices(args(body)?).await?),
1032 "admin_audit" => reply(&billing.admin_audit(args(body)?).await?),
1033 "note_pending" => reply(&billing.note_pending(args(body)?).await?),
1034 "admin_accounts" => reply(&billing.admin_accounts(args(body)?).await?),
1035 "admin_account" => reply(&billing.admin_account(args(body)?).await?),
1036 "admin_set_terms" => reply(&billing.admin_set_terms(args(body)?).await?),
1037 "admin_create_enterprise" => reply(&billing.admin_create_enterprise(args(body)?).await?),
1038 "admin_attach" => reply(&billing.admin_attach(args(body)?).await?),
1039 "admin_credit" => reply(&billing.admin_credit(args(body)?).await?),
1040 _ => Response::error("Unknown method", 404),
1041 }
1042}
1043
1044#[cfg(test)]
1045mod tests {
1046 use super::*;
1047
1048 #[test]
1049 fn a_run_is_charged_its_cost_plus_the_margin() {
1050 // $0.05 at 20% is six cents.
1051 assert_eq!(charge_micros(0.05, 20), 60_000);
1052 assert_eq!(charge_micros(1.0, 20), 1_200_000);
1053 assert_eq!(charge_micros(0.05, 0), 50_000);
1054 }
1055
1056 #[test]
1057 fn fractions_of_a_millionth_round_up_and_nothing_costs_less_than_nothing() {
1058 assert_eq!(charge_micros(0.000_000_4, 20), 2);
1059 assert_eq!(charge_micros(0.0, 20), 0);
1060 assert_eq!(charge_micros(-3.0, 20), 0);
1061 }
1062
1063 #[test]
1064 fn sandbox_seconds_are_charged_only_past_the_free_minutes() {
1065 let free = sandbox_allowance::FREE_SECONDS;
1066 assert_eq!(sandbox_billable(0, 600), 0);
1067 assert_eq!(sandbox_billable(free - 100, 600), 500);
1068 assert_eq!(sandbox_billable(free + 5, 600), 600);
1069 }
1070
1071 #[test]
1072 fn durations_read_plainly() {
1073 assert_eq!(duration(40), "40s");
1074 assert_eq!(duration(192), "3m 12s");
1075 assert_eq!(duration(3720), "1h 2m");
1076 }
1077
1078 #[test]
1079 fn an_absurd_cost_is_capped() {
1080 assert_eq!(charge_micros(1e9, 20), 120 * MICROS_PER_DOLLAR);
1081 }
1082}