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