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

633 lines22,900 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//! Without a card processor configured the service says so and charges
11//! nothing, so that g1t still runs where billing has not been set up.
12//!
13//! Reached only through service bindings; see `g1t_contracts::billing` for
14//! the methods and their arguments.
15
16mod stripe;
17
18use g1t_contracts::billing::*;
19use g1t_contracts::time::rfc3339;
20use g1t_contracts::{FailureCode, Outcome, Role, new_id};
21use g1t_kit::{args, now_ms, reply, rpc_method};
22use serde::Deserialize;
23use sha2::{Digest, Sha256};
24use worker::wasm_bindgen::JsValue;
25use worker::{Context, D1Database, Env, Request, Response, Result, event};
26
27use stripe::Stripe;
28
29const MIN_TOP_UP_CENTS: u32 = 500;
30const MAX_TOP_UP_CENTS: u32 = 50_000;
31const LEDGER_PAGE: u32 = 100;
32/// A run's reported cost is believed up to this much. A sandbox cannot
33/// spend more in the time it has, so anything above is a fault.
34const MAX_RUN_COST_USD: f64 = 100.0;
35
36/// What a run is charged: its cost plus the margin, rounded up to a whole
37/// millionth of a dollar.
38pub fn charge_micros(cost_usd: f64, margin_percent: u32) -> i64 {
39 let cost_micros = (cost_usd.clamp(0.0, MAX_RUN_COST_USD) * MICROS_PER_DOLLAR as f64).ceil();
40 (cost_micros * f64::from(100 + margin_percent) / 100.0).ceil() as i64
41}
42
43fn hash(token: &str) -> String {
44 hex::encode(Sha256::digest(token.as_bytes()))
45}
46
47fn optional(value: Option<&str>) -> JsValue {
48 value.map_or(JsValue::NULL, JsValue::from)
49}
50
51#[derive(Deserialize)]
52struct AccountRow {
53 balance_micros: i64,
54 customer_id: Option<String>,
55}
56
57#[derive(Deserialize)]
58struct LedgerRow {
59 id: String,
60 kind: EntryKind,
61 amount_micros: i64,
62 description: String,
63 repo: Option<String>,
64 number: Option<u32>,
65 task: Option<String>,
66 model: Option<String>,
67 created_by: Option<String>,
68 created_at: String,
69 billed_to: Option<String>,
70}
71
72impl From<LedgerRow> for LedgerEntry {
73 fn from(row: LedgerRow) -> Self {
74 LedgerEntry {
75 id: row.id,
76 kind: row.kind,
77 amount_micros: row.amount_micros,
78 description: row.description,
79 repo: row.repo,
80 number: row.number,
81 task: row.task,
82 model: row.model,
83 billed_to: row.billed_to.unwrap_or_else(|| "g1t".to_owned()),
84 created_by: row.created_by,
85 created_at: row.created_at,
86 }
87 }
88}
89
90#[derive(Deserialize)]
91struct RunRow {
92 workspace: String,
93 repo: String,
94 number: u32,
95 task: String,
96 model: String,
97 token_hash: String,
98 billed_to: Option<String>,
99}
100
101impl RunRow {
102 fn own_provider(&self) -> bool {
103 self.billed_to.as_deref() == Some("workspace")
104 }
105}
106
107#[derive(Deserialize)]
108struct CheckoutRow {
109 workspace: String,
110 created_by: String,
111}
112
113/// A row an `UPDATE … RETURNING` touched.
114#[derive(Deserialize)]
115struct Touched {
116 #[allow(dead_code)]
117 id: String,
118}
119
120struct Billing {
121 db: D1Database,
122 /// Absent when no card processor is configured.
123 stripe: Option<Stripe>,
124 margin_percent: u32,
125 /// Charged for a run on the workspace's own model provider.
126 orchestration_fee_micros: i64,
127}
128
129impl Billing {
130 fn status(&self) -> Status {
131 Status {
132 enabled: self.stripe.is_some(),
133 live: self.stripe.as_ref().is_some_and(Stripe::live),
134 }
135 }
136
137 async fn row(&self, workspace: &str) -> Result<Option<AccountRow>> {
138 self.db
139 .prepare("SELECT balance_micros, customer_id FROM accounts WHERE workspace = ?")
140 .bind(&[workspace.into()])?
141 .first::<AccountRow>(None)
142 .await
143 }
144
145 async fn standing(&self, workspace: &str) -> Result<Account> {
146 Ok(Account {
147 workspace: workspace.to_owned(),
148 balance_micros: self
149 .row(workspace)
150 .await?
151 .map_or(0, |row| row.balance_micros),
152 status: self.status(),
153 margin_percent: self.margin_percent,
154 orchestration_fee_micros: self.orchestration_fee_micros,
155 })
156 }
157
158 /// Adds a ledger entry and moves the balance by the same amount, as
159 /// one write.
160 #[allow(clippy::too_many_arguments)]
161 async fn enter(
162 &self,
163 workspace: &str,
164 kind: EntryKind,
165 amount_micros: i64,
166 description: &str,
167 reference: &str,
168 run: Option<&RunRow>,
169 cost_micros: Option<i64>,
170 created_by: Option<&str>,
171 customer: Option<&str>,
172 ) -> Result<()> {
173 let now = now_ms();
174 let timestamp = rfc3339(now);
175 let kind = match kind {
176 EntryKind::TopUp => "top_up",
177 EntryKind::Usage => "usage",
178 };
179 self.db
180 .batch(vec![
181 self.db
182 .prepare(
183 "INSERT INTO ledger
184 (id, workspace, kind, amount_micros, description, repo, number, task,
185 model, cost_micros, reference, created_by, created_at, billed_to)
186 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
187 )
188 .bind(&[
189 new_id("led", now).into(),
190 workspace.into(),
191 kind.into(),
192 // D1 takes numbers as doubles, which hold every
193 // amount this service will see exactly.
194 (amount_micros as f64).into(),
195 description.into(),
196 optional(run.map(|run| run.repo.as_str())),
197 run.map_or(JsValue::NULL, |run| run.number.into()),
198 optional(run.map(|run| run.task.as_str())),
199 optional(run.map(|run| run.model.as_str())),
200 cost_micros.map_or(JsValue::NULL, |cost| (cost as f64).into()),
201 reference.into(),
202 optional(created_by),
203 timestamp.as_str().into(),
204 run.map_or("g1t", |run| if run.own_provider() { "workspace" } else { "g1t" }).into(),
205 ])?,
206 self.db
207 .prepare(
208 "INSERT INTO accounts (workspace, balance_micros, customer_id, created_at)
209 VALUES (?1, ?2, ?3, ?4)
210 ON CONFLICT (workspace) DO UPDATE SET
211 balance_micros = balance_micros + ?2,
212 customer_id = COALESCE(?3, customer_id)",
213 )
214 .bind(&[
215 workspace.into(),
216 (amount_micros as f64).into(),
217 optional(customer),
218 timestamp.as_str().into(),
219 ])?,
220 ])
221 .await?;
222 Ok(())
223 }
224
225 async fn account(&self, a: AccountArgs) -> Result<Outcome<Account>> {
226 let workspace = a.workspace.to_lowercase();
227 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
228 return Ok(members_only());
229 }
230 Ok(Outcome::Ok(self.standing(&workspace).await?))
231 }
232
233 async fn ledger(&self, a: AccountArgs) -> Result<Outcome<Vec<LedgerEntry>>> {
234 let workspace = a.workspace.to_lowercase();
235 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
236 return Ok(members_only());
237 }
238 let rows = self
239 .db
240 .prepare("SELECT * FROM ledger WHERE workspace = ? ORDER BY id DESC LIMIT ?")
241 .bind(&[workspace.into(), LEDGER_PAGE.into()])?
242 .all()
243 .await?
244 .results::<LedgerRow>()?;
245 Ok(Outcome::Ok(
246 rows.into_iter().map(LedgerEntry::from).collect(),
247 ))
248 }
249
250 async fn usage(&self, a: UsageArgs) -> Result<Outcome<Usage>> {
251 let workspace = a.workspace.to_lowercase();
252 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
253 return Ok(members_only());
254 }
255 #[derive(serde::Deserialize)]
256 struct SliceRow {
257 key: Option<String>,
258 micros: Option<i64>,
259 runs: Option<u32>,
260 }
261 let slices = |key: &str, limit: u32| {
262 format!(
263 "SELECT {key} AS key, -SUM(amount_micros) AS micros, COUNT(*) AS runs FROM ledger
264 WHERE workspace = ?1 AND kind = 'usage' AND created_at >= ?2
265 GROUP BY 1 ORDER BY micros DESC LIMIT {limit}"
266 )
267 };
268 let query = |sql: String| {
269 let db = &self.db;
270 let workspace = workspace.clone();
271 let since = a.since.clone();
272 async move {
273 let rows = db
274 .prepare(sql)
275 .bind(&[workspace.into(), since.into()])?
276 .all()
277 .await?
278 .results::<SliceRow>()?;
279 Ok::<Vec<UsageSlice>, worker::Error>(
280 rows.into_iter()
281 .map(|row| UsageSlice {
282 key: row.key.unwrap_or_else(|| "other".to_owned()),
283 micros: row.micros.unwrap_or_default(),
284 runs: row.runs.unwrap_or_default(),
285 })
286 .collect(),
287 )
288 }
289 };
290 #[derive(serde::Deserialize)]
291 struct Totals {
292 spent: Option<i64>,
293 cost: Option<i64>,
294 provider: Option<i64>,
295 runs: Option<u32>,
296 added: Option<i64>,
297 }
298 let totals = self
299 .db
300 .prepare(
301 "SELECT
302 -SUM(CASE WHEN kind = 'usage' THEN amount_micros END) AS spent,
303 SUM(CASE WHEN kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' THEN cost_micros END) AS cost,
304 SUM(CASE WHEN kind = 'usage' AND billed_to = 'workspace' THEN cost_micros END) AS provider,
305 SUM(CASE WHEN kind = 'usage' THEN 1 ELSE 0 END) AS runs,
306 SUM(CASE WHEN kind = 'top_up' THEN amount_micros END) AS added
307 FROM ledger WHERE workspace = ?1 AND created_at >= ?2",
308 )
309 .bind(&[workspace.as_str().into(), a.since.as_str().into()])?
310 .first::<Totals>(None)
311 .await?;
312 let totals = totals.unwrap_or(Totals {
313 spent: None,
314 cost: None,
315 provider: None,
316 runs: None,
317 added: None,
318 });
319 Ok(Outcome::Ok(Usage {
320 spent_micros: totals.spent.unwrap_or_default(),
321 cost_micros: totals.cost.unwrap_or_default(),
322 provider_micros: totals.provider.unwrap_or_default(),
323 runs: totals.runs.unwrap_or_default(),
324 added_micros: totals.added.unwrap_or_default(),
325 by_day: query(slices("substr(created_at, 1, 10) || '/' || COALESCE(task, 'other')", 400)).await?,
326 by_task: query(slices("task", 20)).await?,
327 by_repo: query(slices("repo", 20)).await?,
328 by_pull: query(slices("repo || '#' || number", 10)).await?,
329 by_model: query(slices("model", 10)).await?,
330 since: a.since,
331 }))
332 }
333
334 async fn checkout(&self, a: CheckoutArgs) -> Result<Outcome<Checkout>> {
335 let workspace = a.workspace.to_lowercase();
336 if a.actor.role_in(&workspace) != Some(Role::Owner) {
337 return Ok(Outcome::fail(
338 FailureCode::Forbidden,
339 "Only an owner can add credit to a workspace.",
340 ));
341 }
342 let Some(stripe) = &self.stripe else {
343 return Ok(Outcome::fail(
344 FailureCode::Conflict,
345 "Payments are not set up on this g1t yet.",
346 ));
347 };
348 if !(MIN_TOP_UP_CENTS..=MAX_TOP_UP_CENTS).contains(&a.amount_cents) {
349 return Ok(Outcome::fail(
350 FailureCode::Invalid,
351 format!(
352 "Add between ${} and ${} at a time.",
353 MIN_TOP_UP_CENTS / 100,
354 MAX_TOP_UP_CENTS / 100
355 ),
356 ));
357 }
358 let customer = self.row(&workspace).await?.and_then(|row| row.customer_id);
359 let session = stripe
360 .start_checkout(
361 &workspace,
362 a.amount_cents,
363 customer.as_deref(),
364 &a.return_url,
365 )
366 .await?;
367 let Some(url) = session.url else {
368 return Err(worker::Error::RustError(
369 "the card processor returned no payment page".into(),
370 ));
371 };
372 self.db
373 .prepare(
374 "INSERT INTO checkouts (id, workspace, amount_cents, created_by, created_at)
375 VALUES (?, ?, ?, ?, ?)",
376 )
377 .bind(&[
378 session.id.into(),
379 workspace.into(),
380 a.amount_cents.into(),
381 a.actor.username.into(),
382 rfc3339(now_ms()).into(),
383 ])?
384 .run()
385 .await?;
386 Ok(Outcome::Ok(Checkout { url }))
387 }
388
389 /// Credits a payment if the processor says it was made and it has not
390 /// been credited before. The amount credited is what the processor
391 /// says was paid, not what anyone here remembers asking for.
392 async fn confirm(&self, a: ConfirmArgs) -> Result<Outcome<Account>> {
393 let workspace = a.workspace.to_lowercase();
394 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
395 return Ok(members_only());
396 }
397 let (Some(stripe), Some(checkout)) = (
398 &self.stripe,
399 self.db
400 .prepare(
401 "SELECT workspace, created_by FROM checkouts
402 WHERE id = ? AND workspace = ? AND status = 'open'",
403 )
404 .bind(&[a.session.as_str().into(), workspace.as_str().into()])?
405 .first::<CheckoutRow>(None)
406 .await?,
407 ) else {
408 // Unknown, someone else's, or already credited: nothing to do.
409 return Ok(Outcome::Ok(self.standing(&workspace).await?));
410 };
411 let session = stripe.session(&a.session).await?;
412 let paid = session
413 .amount_total
414 .filter(|_| session.payment_status == "paid");
415 if let Some(cents) = paid {
416 // Only whoever flips it from open to paid enters the credit.
417 let claimed = self
418 .db
419 .prepare(
420 "UPDATE checkouts SET status = 'paid' WHERE id = ? AND status = 'open'
421 RETURNING id",
422 )
423 .bind(&[a.session.as_str().into()])?
424 .first::<Touched>(None)
425 .await?;
426 if claimed.is_some() {
427 self.enter(
428 &checkout.workspace,
429 EntryKind::TopUp,
430 i64::from(cents) * MICROS_PER_DOLLAR / 100,
431 "Credit added by card",
432 &session.id,
433 None,
434 None,
435 Some(&checkout.created_by),
436 session.customer.as_deref(),
437 )
438 .await?;
439 }
440 }
441 Ok(Outcome::Ok(self.standing(&workspace).await?))
442 }
443
444 /// A refusal if the workspace has no credit to start an agent with.
445 async fn out_of_credit<T>(&self, workspace: &str) -> Result<Option<Outcome<T>>> {
446 let balance = self
447 .row(workspace)
448 .await?
449 .map_or(0, |row| row.balance_micros);
450 Ok((balance <= 0).then(|| {
451 Outcome::fail(
452 FailureCode::PaymentRequired,
453 format!(
454 "The {workspace} workspace has no agent credit. An owner can add some under Billing on the workspace's page."
455 ),
456 )
457 }))
458 }
459
460 async fn can_start(&self, a: CanStartArgs) -> Result<Outcome<bool>> {
461 if self.stripe.is_none() {
462 return Ok(Outcome::Ok(true));
463 }
464 Ok(self
465 .out_of_credit(&a.workspace.to_lowercase())
466 .await?
467 .unwrap_or(Outcome::Ok(true)))
468 }
469
470 async fn start_run(&self, a: StartRunArgs) -> Result<Outcome<Option<RunTicket>>> {
471 if self.stripe.is_none() {
472 return Ok(Outcome::Ok(None));
473 }
474 let workspace = a.workspace.to_lowercase();
475 if let Some(refused) = self.out_of_credit(&workspace).await? {
476 return Ok(refused);
477 }
478 let now = now_ms();
479 let run_id = new_id("run", now);
480 let mut bytes = [0u8; 32];
481 getrandom::getrandom(&mut bytes).expect("no source of randomness");
482 let token = hex::encode(bytes);
483 self.db
484 .prepare(
485 "INSERT INTO runs (id, workspace, repo, number, task, model, token_hash, created_at, billed_to)
486 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
487 )
488 .bind(&[
489 run_id.as_str().into(),
490 workspace.into(),
491 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
492 a.number.into(),
493 a.task.into(),
494 a.model.into(),
495 hash(&token).into(),
496 rfc3339(now).into(),
497 if a.billed_to == "workspace" { "workspace" } else { "g1t" }.into(),
498 ])?
499 .run()
500 .await?;
501 Ok(Outcome::Ok(Some(RunTicket { run_id, token })))
502 }
503
504 async fn finish_run(&self, a: FinishRunArgs) -> Result<Outcome<bool>> {
505 let run = self
506 .db
507 .prepare(
508 "SELECT workspace, repo, number, task, model, token_hash, billed_to FROM runs
509 WHERE id = ? AND finished_at IS NULL",
510 )
511 .bind(&[a.run_id.as_str().into()])?
512 .first::<RunRow>(None)
513 .await?;
514 let Some(run) = run.filter(|run| run.token_hash == hash(&a.token)) else {
515 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
516 };
517 if !a.cost_usd.is_finite() || a.cost_usd < 0.0 {
518 return Ok(Outcome::fail(FailureCode::Invalid, "That is not a cost."));
519 }
520 // Only whoever closes the run charges for it.
521 let claimed = self
522 .db
523 .prepare(
524 "UPDATE runs SET finished_at = ? WHERE id = ? AND finished_at IS NULL RETURNING id",
525 )
526 .bind(&[rfc3339(now_ms()).into(), a.run_id.as_str().into()])?
527 .first::<Touched>(None)
528 .await?;
529 if claimed.is_none() {
530 return Ok(Outcome::Ok(false));
531 }
532 // On the workspace's own provider, the model was paid for there:
533 // g1t charges its fee, and keeps the provider's cost to show.
534 let charge = if run.own_provider() {
535 self.orchestration_fee_micros
536 } else {
537 charge_micros(a.cost_usd, self.margin_percent)
538 };
539 let mut description = match run.task.as_str() {
540 "plan" => format!("Planning for {}", run.repo),
541 "review" => format!("Review of {}#{}", run.repo, run.number),
542 "update" => format!("Catching up {}#{}", run.repo, run.number),
543 _ => format!("Work on {}#{}", run.repo, run.number),
544 };
545 if run.own_provider() {
546 description.push_str(", on your own model provider");
547 }
548 self.enter(
549 &run.workspace,
550 EntryKind::Usage,
551 -charge,
552 &description,
553 &a.run_id,
554 Some(&run),
555 Some(charge_micros(a.cost_usd, 0)),
556 None,
557 None,
558 )
559 .await?;
560 Ok(Outcome::Ok(true))
561 }
562}
563
564fn members_only<T>() -> Outcome<T> {
565 Outcome::fail(
566 FailureCode::Forbidden,
567 "Only members can see a workspace's billing.",
568 )
569}
570
571#[event(fetch)]
572async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
573 let Some(method) = rpc_method(&request) else {
574 return Response::error("Not found", 404);
575 };
576 let body: serde_json::Value = request.json().await?;
577 let billing = Billing {
578 db: env.d1("DB")?,
579 stripe: env
580 .secret("STRIPE_SECRET_KEY")
581 .ok()
582 .map(|key| key.to_string())
583 .filter(|key| !key.is_empty())
584 .map(Stripe::new),
585 margin_percent: env
586 .var("MARGIN_PERCENT")
587 .ok()
588 .and_then(|percent| percent.to_string().parse().ok())
589 .unwrap_or(20),
590 orchestration_fee_micros: env
591 .var("ORCHESTRATION_FEE_MICROS")
592 .ok()
593 .and_then(|fee| fee.to_string().parse().ok())
594 .unwrap_or(100_000),
595 };
596 match method.as_str() {
597 "status" => reply(&billing.status()),
598 "account" => reply(&billing.account(args(body)?).await?),
599 "ledger" => reply(&billing.ledger(args(body)?).await?),
600 "usage" => reply(&billing.usage(args(body)?).await?),
601 "checkout" => reply(&billing.checkout(args(body)?).await?),
602 "confirm" => reply(&billing.confirm(args(body)?).await?),
603 "can_start" => reply(&billing.can_start(args(body)?).await?),
604 "start_run" => reply(&billing.start_run(args(body)?).await?),
605 "finish_run" => reply(&billing.finish_run(args(body)?).await?),
606 _ => Response::error("Unknown method", 404),
607 }
608}
609
610#[cfg(test)]
611mod tests {
612 use super::*;
613
614 #[test]
615 fn a_run_is_charged_its_cost_plus_the_margin() {
616 // $0.05 at 20% is six cents.
617 assert_eq!(charge_micros(0.05, 20), 60_000);
618 assert_eq!(charge_micros(1.0, 20), 1_200_000);
619 assert_eq!(charge_micros(0.05, 0), 50_000);
620 }
621
622 #[test]
623 fn fractions_of_a_millionth_round_up_and_nothing_costs_less_than_nothing() {
624 assert_eq!(charge_micros(0.000_000_4, 20), 2);
625 assert_eq!(charge_micros(0.0, 20), 0);
626 assert_eq!(charge_micros(-3.0, 20), 0);
627 }
628
629 #[test]
630 fn an_absurd_cost_is_capped() {
631 assert_eq!(charge_micros(1e9, 20), 120 * MICROS_PER_DOLLAR);
632 }
633}