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