flagon-io/g1t

public

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

g1t/services/billing/src/lib.rs

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