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

943 lines35,648 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Agents as a team: lifecycle, merge queue, billing and a new shell1//! 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//!
Paid features: a workspace turns on Deployments with a monthly plan10//! Paid features (deployments) are bought separately, as monthly plans;
11//! see `features`. They are never free.
12//!
Agents as a team: lifecycle, merge queue, billing and a new shell13//! 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
Paid features: a workspace turns on Deployments with a monthly plan19mod features;
Prices keep themselves current with what g1t pays20mod keeper;
Usage limits: unpaid usage can only go so far21mod limits;
Agents as a team: lifecycle, merge queue, billing and a new shell22mod stripe;
23
24use g1t_contracts::billing::*;
25use g1t_contracts::time::rfc3339;
26use g1t_contracts::{FailureCode, Outcome, Role, new_id};
27use g1t_kit::{args, now_ms, reply, rpc_method};
28use serde::Deserialize;
29use sha2::{Digest, Sha256};
30use worker::wasm_bindgen::JsValue;
Prices keep themselves current with what g1t pays31use worker::{Context, D1Database, Env, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
Agents as a team: lifecycle, merge queue, billing and a new shell32
33use stripe::Stripe;
34
35const MIN_TOP_UP_CENTS: u32 = 500;
36const MAX_TOP_UP_CENTS: u32 = 50_000;
37const LEDGER_PAGE: u32 = 100;
38/// A run's reported cost is believed up to this much. A sandbox cannot
39/// spend more in the time it has, so anything above is a fault.
40const MAX_RUN_COST_USD: f64 = 100.0;
41
42/// What a run is charged: its cost plus the margin, rounded up to a whole
43/// millionth of a dollar.
44pub fn charge_micros(cost_usd: f64, margin_percent: u32) -> i64 {
45 let cost_micros = (cost_usd.clamp(0.0, MAX_RUN_COST_USD) * MICROS_PER_DOLLAR as f64).ceil();
46 (cost_micros * f64::from(100 + margin_percent) / 100.0).ceil() as i64
47}
48
49fn hash(token: &str) -> String {
50 hex::encode(Sha256::digest(token.as_bytes()))
51}
52
53fn optional(value: Option<&str>) -> JsValue {
54 value.map_or(JsValue::NULL, JsValue::from)
55}
56
57#[derive(Deserialize)]
58struct AccountRow {
59 balance_micros: i64,
60 customer_id: Option<String>,
61}
62
63#[derive(Deserialize)]
64struct LedgerRow {
65 id: String,
66 kind: EntryKind,
67 amount_micros: i64,
68 description: String,
69 repo: Option<String>,
70 number: Option<u32>,
71 task: Option<String>,
72 model: Option<String>,
73 created_by: Option<String>,
74 created_at: String,
Integrations: your own model provider, alerts that open issues, tickets agents read75 billed_to: Option<String>,
Agents as a team: lifecycle, merge queue, billing and a new shell76}
77
78impl From<LedgerRow> for LedgerEntry {
79 fn from(row: LedgerRow) -> Self {
80 LedgerEntry {
81 id: row.id,
82 kind: row.kind,
83 amount_micros: row.amount_micros,
84 description: row.description,
85 repo: row.repo,
86 number: row.number,
87 task: row.task,
88 model: row.model,
Integrations: your own model provider, alerts that open issues, tickets agents read89 billed_to: row.billed_to.unwrap_or_else(|| "g1t".to_owned()),
Agents as a team: lifecycle, merge queue, billing and a new shell90 created_by: row.created_by,
91 created_at: row.created_at,
92 }
93 }
94}
95
96#[derive(Deserialize)]
97struct RunRow {
98 workspace: String,
99 repo: String,
100 number: u32,
101 task: String,
102 model: String,
103 token_hash: String,
Integrations: your own model provider, alerts that open issues, tickets agents read104 billed_to: Option<String>,
105}
106
107impl RunRow {
108 fn own_provider(&self) -> bool {
109 self.billed_to.as_deref() == Some("workspace")
110 }
Agents as a team: lifecycle, merge queue, billing and a new shell111}
112
113#[derive(Deserialize)]
114struct CheckoutRow {
115 workspace: String,
116 created_by: String,
117}
118
119/// A row an `UPDATE … RETURNING` touched.
120#[derive(Deserialize)]
121struct Touched {
122 #[allow(dead_code)]
123 id: String,
124}
125
126struct Billing {
127 db: D1Database,
128 /// Absent when no card processor is configured.
129 stripe: Option<Stripe>,
130 margin_percent: u32,
Integrations: your own model provider, alerts that open issues, tickets agents read131 /// Charged for a run on the workspace's own model provider.
132 orchestration_fee_micros: i64,
Free while g1t is being built out; agents can check out their own forks133 /// While g1t is being built out, nothing is charged (`FREE_WHILE_BUILDING`).
134 free: bool,
A free allowance on g1t's models, so anyone can try its agents135 /// The free allowance on g1t's hosted models, when there is one.
136 trial: Option<TrialConfig>,
Paid features: a workspace turns on Deployments with a monthly plan137 /// The Deployments plan's monthly price (`DEPLOYMENTS_MONTHLY_CENTS`).
138 deployments_monthly_cents: u32,
Usage limits: unpaid usage can only go so far139 /// How far unpaid usage may go; see `limits`.
140 ceilings: limits::Ceilings,
A free allowance on g1t's models, so anyone can try its agents141}
142
143/// `TRIAL_WORKSPACE_MICROS`, `TRIAL_TOTAL_MICROS` and `TRIAL_UNTIL`.
144struct TrialConfig {
145 per_workspace_micros: i64,
146 total_micros: i64,
147 /// RFC 3339, in UTC.
148 until: String,
149}
150
151#[derive(serde::Deserialize)]
152struct Sum {
153 micros: Option<i64>,
Agents as a team: lifecycle, merge queue, billing and a new shell154}
155
156impl Billing {
157 fn status(&self) -> Status {
158 Status {
159 enabled: self.stripe.is_some(),
160 live: self.stripe.as_ref().is_some_and(Stripe::live),
Free while g1t is being built out; agents can check out their own forks161 free: self.free,
Agents as a team: lifecycle, merge queue, billing and a new shell162 }
163 }
164
165 async fn row(&self, workspace: &str) -> Result<Option<AccountRow>> {
166 self.db
167 .prepare("SELECT balance_micros, customer_id FROM accounts WHERE workspace = ?")
168 .bind(&[workspace.into()])?
169 .first::<AccountRow>(None)
170 .await
171 }
172
173 async fn standing(&self, workspace: &str) -> Result<Account> {
174 Ok(Account {
175 workspace: workspace.to_owned(),
176 balance_micros: self
177 .row(workspace)
178 .await?
179 .map_or(0, |row| row.balance_micros),
180 status: self.status(),
181 margin_percent: self.margin_percent,
Integrations: your own model provider, alerts that open issues, tickets agents read182 orchestration_fee_micros: self.orchestration_fee_micros,
Agents as a team: lifecycle, merge queue, billing and a new shell183 })
184 }
185
186 /// Adds a ledger entry and moves the balance by the same amount, as
187 /// one write.
188 #[allow(clippy::too_many_arguments)]
189 async fn enter(
190 &self,
191 workspace: &str,
192 kind: EntryKind,
193 amount_micros: i64,
194 description: &str,
195 reference: &str,
196 run: Option<&RunRow>,
197 cost_micros: Option<i64>,
198 created_by: Option<&str>,
199 customer: Option<&str>,
200 ) -> Result<()> {
201 let now = now_ms();
202 let timestamp = rfc3339(now);
203 let kind = match kind {
204 EntryKind::TopUp => "top_up",
205 EntryKind::Usage => "usage",
206 };
207 self.db
208 .batch(vec![
209 self.db
210 .prepare(
211 "INSERT INTO ledger
212 (id, workspace, kind, amount_micros, description, repo, number, task,
Integrations: your own model provider, alerts that open issues, tickets agents read213 model, cost_micros, reference, created_by, created_at, billed_to)
214 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Agents as a team: lifecycle, merge queue, billing and a new shell215 )
216 .bind(&[
217 new_id("led", now).into(),
218 workspace.into(),
219 kind.into(),
220 // D1 takes numbers as doubles, which hold every
221 // amount this service will see exactly.
222 (amount_micros as f64).into(),
223 description.into(),
224 optional(run.map(|run| run.repo.as_str())),
225 run.map_or(JsValue::NULL, |run| run.number.into()),
226 optional(run.map(|run| run.task.as_str())),
227 optional(run.map(|run| run.model.as_str())),
228 cost_micros.map_or(JsValue::NULL, |cost| (cost as f64).into()),
229 reference.into(),
230 optional(created_by),
231 timestamp.as_str().into(),
Integrations: your own model provider, alerts that open issues, tickets agents read232 run.map_or("g1t", |run| if run.own_provider() { "workspace" } else { "g1t" }).into(),
Agents as a team: lifecycle, merge queue, billing and a new shell233 ])?,
234 self.db
235 .prepare(
236 "INSERT INTO accounts (workspace, balance_micros, customer_id, created_at)
237 VALUES (?1, ?2, ?3, ?4)
238 ON CONFLICT (workspace) DO UPDATE SET
239 balance_micros = balance_micros + ?2,
240 customer_id = COALESCE(?3, customer_id)",
241 )
242 .bind(&[
243 workspace.into(),
244 (amount_micros as f64).into(),
245 optional(customer),
246 timestamp.as_str().into(),
247 ])?,
248 ])
249 .await?;
250 Ok(())
251 }
252
253 async fn account(&self, a: AccountArgs) -> Result<Outcome<Account>> {
254 let workspace = a.workspace.to_lowercase();
255 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
256 return Ok(members_only());
257 }
258 Ok(Outcome::Ok(self.standing(&workspace).await?))
259 }
260
261 async fn ledger(&self, a: AccountArgs) -> Result<Outcome<Vec<LedgerEntry>>> {
262 let workspace = a.workspace.to_lowercase();
263 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
264 return Ok(members_only());
265 }
266 let rows = self
267 .db
268 .prepare("SELECT * FROM ledger WHERE workspace = ? ORDER BY id DESC LIMIT ?")
269 .bind(&[workspace.into(), LEDGER_PAGE.into()])?
270 .all()
271 .await?
272 .results::<LedgerRow>()?;
273 Ok(Outcome::Ok(
274 rows.into_iter().map(LedgerEntry::from).collect(),
275 ))
276 }
277
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request278 async fn usage(&self, a: UsageArgs) -> Result<Outcome<Usage>> {
279 let workspace = a.workspace.to_lowercase();
280 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
281 return Ok(members_only());
282 }
283 #[derive(serde::Deserialize)]
284 struct SliceRow {
285 key: Option<String>,
286 micros: Option<i64>,
287 runs: Option<u32>,
288 }
Usage while free is shown at cost; agents get rustfmt and clippy289 // While nothing is charged, what was used is what there is to show.
290 let measure = if self.free { "COALESCE(cost_micros, 0)" } else { "-amount_micros" };
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request291 let slices = |key: &str, limit: u32| {
292 format!(
Usage while free is shown at cost; agents get rustfmt and clippy293 "SELECT {key} AS key, SUM({measure}) AS micros, COUNT(*) AS runs FROM ledger
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request294 WHERE workspace = ?1 AND kind = 'usage' AND created_at >= ?2
295 GROUP BY 1 ORDER BY micros DESC LIMIT {limit}"
296 )
297 };
298 let query = |sql: String| {
299 let db = &self.db;
300 let workspace = workspace.clone();
301 let since = a.since.clone();
302 async move {
303 let rows = db
304 .prepare(sql)
305 .bind(&[workspace.into(), since.into()])?
306 .all()
307 .await?
308 .results::<SliceRow>()?;
309 Ok::<Vec<UsageSlice>, worker::Error>(
310 rows.into_iter()
311 .map(|row| UsageSlice {
312 key: row.key.unwrap_or_else(|| "other".to_owned()),
313 micros: row.micros.unwrap_or_default(),
314 runs: row.runs.unwrap_or_default(),
315 })
316 .collect(),
317 )
318 }
319 };
320 #[derive(serde::Deserialize)]
321 struct Totals {
322 spent: Option<i64>,
323 cost: Option<i64>,
Integrations: your own model provider, alerts that open issues, tickets agents read324 provider: Option<i64>,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request325 runs: Option<u32>,
326 added: Option<i64>,
327 }
328 let totals = self
329 .db
330 .prepare(
331 "SELECT
332 -SUM(CASE WHEN kind = 'usage' THEN amount_micros END) AS spent,
Integrations: your own model provider, alerts that open issues, tickets agents read333 SUM(CASE WHEN kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' THEN cost_micros END) AS cost,
334 SUM(CASE WHEN kind = 'usage' AND billed_to = 'workspace' THEN cost_micros END) AS provider,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request335 SUM(CASE WHEN kind = 'usage' THEN 1 ELSE 0 END) AS runs,
336 SUM(CASE WHEN kind = 'top_up' THEN amount_micros END) AS added
337 FROM ledger WHERE workspace = ?1 AND created_at >= ?2",
338 )
339 .bind(&[workspace.as_str().into(), a.since.as_str().into()])?
340 .first::<Totals>(None)
341 .await?;
342 let totals = totals.unwrap_or(Totals {
343 spent: None,
344 cost: None,
Integrations: your own model provider, alerts that open issues, tickets agents read345 provider: None,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request346 runs: None,
347 added: None,
348 });
349 Ok(Outcome::Ok(Usage {
350 spent_micros: totals.spent.unwrap_or_default(),
351 cost_micros: totals.cost.unwrap_or_default(),
Integrations: your own model provider, alerts that open issues, tickets agents read352 provider_micros: totals.provider.unwrap_or_default(),
Usage while free is shown at cost; agents get rustfmt and clippy353 used_micros: totals.cost.unwrap_or_default() + totals.provider.unwrap_or_default(),
354 free: self.free,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request355 runs: totals.runs.unwrap_or_default(),
356 added_micros: totals.added.unwrap_or_default(),
357 by_day: query(slices("substr(created_at, 1, 10) || '/' || COALESCE(task, 'other')", 400)).await?,
358 by_task: query(slices("task", 20)).await?,
359 by_repo: query(slices("repo", 20)).await?,
360 by_pull: query(slices("repo || '#' || number", 10)).await?,
361 by_model: query(slices("model", 10)).await?,
362 since: a.since,
363 }))
364 }
365
Agents as a team: lifecycle, merge queue, billing and a new shell366 async fn checkout(&self, a: CheckoutArgs) -> Result<Outcome<Checkout>> {
367 let workspace = a.workspace.to_lowercase();
368 if a.actor.role_in(&workspace) != Some(Role::Owner) {
369 return Ok(Outcome::fail(
370 FailureCode::Forbidden,
371 "Only an owner can add credit to a workspace.",
372 ));
373 }
374 let Some(stripe) = &self.stripe else {
375 return Ok(Outcome::fail(
376 FailureCode::Conflict,
377 "Payments are not set up on this g1t yet.",
378 ));
379 };
380 if !(MIN_TOP_UP_CENTS..=MAX_TOP_UP_CENTS).contains(&a.amount_cents) {
381 return Ok(Outcome::fail(
382 FailureCode::Invalid,
383 format!(
384 "Add between ${} and ${} at a time.",
385 MIN_TOP_UP_CENTS / 100,
386 MAX_TOP_UP_CENTS / 100
387 ),
388 ));
389 }
390 let customer = self.row(&workspace).await?.and_then(|row| row.customer_id);
Project dependencies: addresses, preview stacks, Affects, and agents who know391 let session = match stripe
392 .start_checkout(&workspace, a.amount_cents, customer.as_deref(), &a.return_url)
393 .await
394 {
395 Ok(session) => session,
396 // A customer saved under another Stripe account: start afresh.
397 Err(error) if customer.is_some() && stripe::is_missing(&error) => {
398 self.forget_customer(&workspace).await?;
399 stripe.start_checkout(&workspace, a.amount_cents, None, &a.return_url).await?
400 }
401 Err(error) => return Err(error),
402 };
Agents as a team: lifecycle, merge queue, billing and a new shell403 let Some(url) = session.url else {
404 return Err(worker::Error::RustError(
405 "the card processor returned no payment page".into(),
406 ));
407 };
408 self.db
409 .prepare(
410 "INSERT INTO checkouts (id, workspace, amount_cents, created_by, created_at)
411 VALUES (?, ?, ?, ?, ?)",
412 )
413 .bind(&[
414 session.id.into(),
415 workspace.into(),
416 a.amount_cents.into(),
417 a.actor.username.into(),
418 rfc3339(now_ms()).into(),
419 ])?
420 .run()
421 .await?;
422 Ok(Outcome::Ok(Checkout { url }))
423 }
424
425 /// Credits a payment if the processor says it was made and it has not
426 /// been credited before. The amount credited is what the processor
427 /// says was paid, not what anyone here remembers asking for.
428 async fn confirm(&self, a: ConfirmArgs) -> Result<Outcome<Account>> {
429 let workspace = a.workspace.to_lowercase();
430 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
431 return Ok(members_only());
432 }
433 let (Some(stripe), Some(checkout)) = (
434 &self.stripe,
435 self.db
436 .prepare(
437 "SELECT workspace, created_by FROM checkouts
438 WHERE id = ? AND workspace = ? AND status = 'open'",
439 )
440 .bind(&[a.session.as_str().into(), workspace.as_str().into()])?
441 .first::<CheckoutRow>(None)
442 .await?,
443 ) else {
444 // Unknown, someone else's, or already credited: nothing to do.
445 return Ok(Outcome::Ok(self.standing(&workspace).await?));
446 };
447 let session = stripe.session(&a.session).await?;
448 let paid = session
449 .amount_total
450 .filter(|_| session.payment_status == "paid");
451 if let Some(cents) = paid {
452 // Only whoever flips it from open to paid enters the credit.
453 let claimed = self
454 .db
455 .prepare(
456 "UPDATE checkouts SET status = 'paid' WHERE id = ? AND status = 'open'
457 RETURNING id",
458 )
459 .bind(&[a.session.as_str().into()])?
460 .first::<Touched>(None)
461 .await?;
462 if claimed.is_some() {
463 self.enter(
464 &checkout.workspace,
465 EntryKind::TopUp,
466 i64::from(cents) * MICROS_PER_DOLLAR / 100,
467 "Credit added by card",
468 &session.id,
469 None,
470 None,
471 Some(&checkout.created_by),
472 session.customer.as_deref(),
473 )
474 .await?;
475 }
476 }
477 Ok(Outcome::Ok(self.standing(&workspace).await?))
478 }
479
Project dependencies: addresses, preview stacks, Affects, and agents who know480 /// Drops a saved customer the card processor no longer knows.
481 pub(crate) async fn forget_customer(&self, workspace: &str) -> Result<()> {
482 self.db
483 .prepare("UPDATE accounts SET customer_id = NULL WHERE workspace = ?")
484 .bind(&[workspace.into()])?
485 .run()
486 .await?;
487 Ok(())
488 }
489
Agents as a team: lifecycle, merge queue, billing and a new shell490 /// A refusal if the workspace has no credit to start an agent with.
491 async fn out_of_credit<T>(&self, workspace: &str) -> Result<Option<Outcome<T>>> {
Free while g1t is being built out; agents can check out their own forks492 // While g1t is being built out, no one needs credit.
493 if self.free {
494 return Ok(None);
495 }
Agents as a team: lifecycle, merge queue, billing and a new shell496 let balance = self
497 .row(workspace)
498 .await?
499 .map_or(0, |row| row.balance_micros);
500 Ok((balance <= 0).then(|| {
501 Outcome::fail(
502 FailureCode::PaymentRequired,
503 format!(
504 "The {workspace} workspace has no agent credit. An owner can add some under Billing on the workspace's page."
505 ),
506 )
507 }))
508 }
509
A free allowance on g1t's models, so anyone can try its agents510 /// A workspace's free allowance on g1t's hosted models: what its runs
511 /// there have cost against its share, and the pool everyone draws on.
512 async fn trial(&self, a: TrialArgs) -> Result<Trial> {
513 let workspace = a.workspace.to_lowercase();
514 let Some(config) = &self.trial else {
515 return Ok(Trial {
516 open: false,
517 used_micros: 0,
518 limit_micros: 0,
519 ends_at: None,
520 reason: Some("off".to_owned()),
521 });
522 };
523 let used = self
524 .db
525 .prepare(
526 "SELECT SUM(cost_micros) AS micros FROM ledger
527 WHERE kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' AND workspace = ?",
528 )
529 .bind(&[workspace.as_str().into()])?
530 .first::<Sum>(None)
531 .await?
532 .and_then(|sum| sum.micros)
533 .unwrap_or_default();
534 // Everyone's, but for the workspaces open to hosted models anyway.
535 let exempt: Vec<String> = a.exempt.iter().map(|name| name.trim().to_lowercase()).collect();
536 let marks = vec!["?"; exempt.len().max(1)].join(", ");
537 let mut values: Vec<JsValue> = exempt.iter().map(|name| JsValue::from(name.as_str())).collect();
538 if values.is_empty() {
539 values.push(JsValue::from(""));
540 }
541 let pooled = self
542 .db
543 .prepare(format!(
544 "SELECT SUM(cost_micros) AS micros FROM ledger
545 WHERE kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' AND workspace NOT IN ({marks})"
546 ))
547 .bind(&values)?
548 .first::<Sum>(None)
549 .await?
550 .and_then(|sum| sum.micros)
551 .unwrap_or_default();
552 let reason = if rfc3339(now_ms()) >= config.until {
553 Some("ended")
554 } else if used >= config.per_workspace_micros {
555 Some("used")
556 } else if pooled >= config.total_micros {
557 Some("pool")
558 } else {
559 None
560 };
561 Ok(Trial {
562 open: reason.is_none(),
563 used_micros: used,
564 limit_micros: config.per_workspace_micros,
565 ends_at: Some(config.until.clone()),
566 reason: reason.map(str::to_owned),
567 })
568 }
569
Agents as a team: lifecycle, merge queue, billing and a new shell570 async fn can_start(&self, a: CanStartArgs) -> Result<Outcome<bool>> {
571 if self.stripe.is_none() {
572 return Ok(Outcome::Ok(true));
573 }
Usage limits: unpaid usage can only go so far574 if let Some(stopped) = self.stopped(&a.workspace).await? {
575 return Ok(stopped);
576 }
Agents as a team: lifecycle, merge queue, billing and a new shell577 Ok(self
578 .out_of_credit(&a.workspace.to_lowercase())
579 .await?
580 .unwrap_or(Outcome::Ok(true)))
581 }
582
583 async fn start_run(&self, a: StartRunArgs) -> Result<Outcome<Option<RunTicket>>> {
584 if self.stripe.is_none() {
585 return Ok(Outcome::Ok(None));
586 }
587 let workspace = a.workspace.to_lowercase();
Usage limits: unpaid usage can only go so far588 if let Some(stopped) = self.stopped(&workspace).await? {
589 return Ok(stopped);
590 }
Agents as a team: lifecycle, merge queue, billing and a new shell591 if let Some(refused) = self.out_of_credit(&workspace).await? {
592 return Ok(refused);
593 }
594 let now = now_ms();
595 let run_id = new_id("run", now);
596 let mut bytes = [0u8; 32];
597 getrandom::getrandom(&mut bytes).expect("no source of randomness");
598 let token = hex::encode(bytes);
599 self.db
600 .prepare(
Prices keep themselves current with what g1t pays601 "INSERT INTO runs (id, workspace, repo, number, task, model, token_hash, created_at, billed_to, session_id)
602 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Agents as a team: lifecycle, merge queue, billing and a new shell603 )
604 .bind(&[
605 run_id.as_str().into(),
606 workspace.into(),
607 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
608 a.number.into(),
609 a.task.into(),
610 a.model.into(),
611 hash(&token).into(),
612 rfc3339(now).into(),
Integrations: your own model provider, alerts that open issues, tickets agents read613 if a.billed_to == "workspace" { "workspace" } else { "g1t" }.into(),
Prices keep themselves current with what g1t pays614 optional(a.session.as_deref().filter(|_| a.billed_to != "workspace")),
Agents as a team: lifecycle, merge queue, billing and a new shell615 ])?
616 .run()
617 .await?;
618 Ok(Outcome::Ok(Some(RunTicket { run_id, token })))
619 }
620
621 async fn finish_run(&self, a: FinishRunArgs) -> Result<Outcome<bool>> {
622 let run = self
623 .db
624 .prepare(
Integrations: your own model provider, alerts that open issues, tickets agents read625 "SELECT workspace, repo, number, task, model, token_hash, billed_to FROM runs
Agents as a team: lifecycle, merge queue, billing and a new shell626 WHERE id = ? AND finished_at IS NULL",
627 )
628 .bind(&[a.run_id.as_str().into()])?
629 .first::<RunRow>(None)
630 .await?;
631 let Some(run) = run.filter(|run| run.token_hash == hash(&a.token)) else {
632 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
633 };
634 if !a.cost_usd.is_finite() || a.cost_usd < 0.0 {
635 return Ok(Outcome::fail(FailureCode::Invalid, "That is not a cost."));
636 }
637 // Only whoever closes the run charges for it.
638 let claimed = self
639 .db
640 .prepare(
641 "UPDATE runs SET finished_at = ? WHERE id = ? AND finished_at IS NULL RETURNING id",
642 )
643 .bind(&[rfc3339(now_ms()).into(), a.run_id.as_str().into()])?
644 .first::<Touched>(None)
645 .await?;
646 if claimed.is_none() {
647 return Ok(Outcome::Ok(false));
648 }
Integrations: your own model provider, alerts that open issues, tickets agents read649 // On the workspace's own provider, the model was paid for there:
650 // g1t charges its fee, and keeps the provider's cost to show.
Free while g1t is being built out; agents can check out their own forks651 let charge = if self.free {
652 // Recorded, with what it cost, but not charged.
653 0
654 } else if run.own_provider() {
Integrations: your own model provider, alerts that open issues, tickets agents read655 self.orchestration_fee_micros
656 } else {
657 charge_micros(a.cost_usd, self.margin_percent)
658 };
659 let mut description = match run.task.as_str() {
Agents as a team: lifecycle, merge queue, billing and a new shell660 "plan" => format!("Planning for {}", run.repo),
661 "review" => format!("Review of {}#{}", run.repo, run.number),
662 "update" => format!("Catching up {}#{}", run.repo, run.number),
663 _ => format!("Work on {}#{}", run.repo, run.number),
664 };
Integrations: your own model provider, alerts that open issues, tickets agents read665 if run.own_provider() {
666 description.push_str(", on your own model provider");
667 }
Free while g1t is being built out; agents can check out their own forks668 if self.free {
669 description.push_str(" (free while g1t is being built out)");
670 }
Agents as a team: lifecycle, merge queue, billing and a new shell671 self.enter(
672 &run.workspace,
673 EntryKind::Usage,
674 -charge,
675 &description,
676 &a.run_id,
677 Some(&run),
678 Some(charge_micros(a.cost_usd, 0)),
679 None,
680 None,
681 )
682 .await?;
683 Ok(Outcome::Ok(true))
684 }
685}
686
Every sandbox is metered by the second687impl Billing {
688 /// Records how long a sandbox ran: its cost always, and a charge for
689 /// the seconds past the month's free minutes.
690 async fn record_sandbox(&self, a: RecordSandboxArgs) -> Result<Outcome<bool>> {
691 if self.stripe.is_none() || a.seconds == 0 {
692 return Ok(Outcome::Ok(false));
693 }
694 let workspace = a.workspace.to_lowercase();
695 let seen = self
696 .db
697 .prepare("SELECT id FROM ledger WHERE reference = ?")
698 .bind(&[a.reference.as_str().into()])?
699 .first::<Touched>(None)
700 .await?;
701 if seen.is_some() {
702 return Ok(Outcome::Ok(false));
703 }
704 let now = now_ms();
705 let timestamp = rfc3339(now);
706 let month = &timestamp[..7];
707 #[derive(Deserialize)]
708 struct Used {
709 seconds: i64,
710 }
711 let seconds = i64::from(a.seconds);
712 let after = self
713 .db
714 .prepare(
715 "INSERT INTO sandbox_months (workspace, month, seconds) VALUES (?1, ?2, ?3)
716 ON CONFLICT (workspace, month) DO UPDATE SET seconds = seconds + ?3
717 RETURNING seconds",
718 )
719 .bind(&[workspace.as_str().into(), month.into(), (seconds as f64).into()])?
720 .first::<Used>(None)
721 .await?
722 .map_or(seconds, |used| used.seconds);
723 let billable = sandbox_billable(after - seconds, seconds);
Prices keep themselves current with what g1t pays724 // From the price book, which follows what Cloudflare bills g1t.
725 let (cost_per_second, price_per_second) = self.price("sandbox_second").await?.unwrap_or((
726 sandbox_allowance::COST_MICROS_PER_SECOND as f64,
727 sandbox_allowance::MICROS_PER_SECOND as f64,
728 ));
729 let charge = if self.free { 0 } else { (billable as f64 * price_per_second).ceil() as i64 };
Every sandbox is metered by the second730 let mut description = format!("{}: {} of sandbox time", a.description, duration(seconds));
731 if billable < seconds {
732 description.push_str(if billable == 0 {
733 ", within the month's free minutes"
734 } else {
735 ", partly within the month's free minutes"
736 });
737 }
738 if self.free && billable > 0 {
739 description.push_str(" (free while g1t is being built out)");
740 }
741 self.db
742 .batch(vec![
743 self.db
744 .prepare(
745 "INSERT INTO ledger
746 (id, workspace, kind, amount_micros, description, repo, task,
747 cost_micros, reference, created_at, billed_to)
748 VALUES (?, ?, 'usage', ?, ?, ?, 'sandbox', ?, ?, ?, 'g1t')",
749 )
750 .bind(&[
751 new_id("led", now).into(),
752 workspace.as_str().into(),
753 (-(charge as f64)).into(),
754 description.as_str().into(),
755 optional(a.repo.as_deref()),
Prices keep themselves current with what g1t pays756 (seconds as f64 * cost_per_second).ceil().into(),
Every sandbox is metered by the second757 a.reference.as_str().into(),
758 timestamp.as_str().into(),
759 ])?,
760 self.db
761 .prepare(
762 "INSERT INTO accounts (workspace, balance_micros, created_at)
763 VALUES (?1, ?2, ?3)
764 ON CONFLICT (workspace) DO UPDATE SET balance_micros = balance_micros + ?2",
765 )
766 .bind(&[
767 workspace.as_str().into(),
768 (-(charge as f64)).into(),
769 timestamp.as_str().into(),
770 ])?,
771 ])
772 .await?;
773 Ok(Outcome::Ok(true))
774 }
775}
776
777/// Of `seconds` used after `before` this month, how many are past the
778/// free minutes.
779fn sandbox_billable(before: i64, seconds: i64) -> i64 {
780 let free_left = (sandbox_allowance::FREE_SECONDS - before).max(0);
781 (seconds - free_left).max(0)
782}
783
784/// `1h 2m`, `3m 12s` or `40s`.
785fn duration(seconds: i64) -> String {
786 let (h, m, s) = (seconds / 3600, seconds % 3600 / 60, seconds % 60);
787 if h > 0 {
788 format!("{h}h {m}m")
789 } else if m > 0 {
790 format!("{m}m {s}s")
791 } else {
792 format!("{s}s")
793 }
794}
795
Agents as a team: lifecycle, merge queue, billing and a new shell796fn members_only<T>() -> Outcome<T> {
797 Outcome::fail(
798 FailureCode::Forbidden,
799 "Only members can see a workspace's billing.",
800 )
801}
802
Prices keep themselves current with what g1t pays803impl Billing {
804 fn from_env(env: &Env) -> Result<Self> {
805 Ok(Billing {
806 db: env.d1("DB")?,
807 stripe: env
808 .secret("STRIPE_SECRET_KEY")
809 .ok()
810 .map(|key| key.to_string())
811 .filter(|key| !key.is_empty())
812 .map(Stripe::new),
813 margin_percent: env
814 .var("MARGIN_PERCENT")
815 .ok()
816 .and_then(|percent| percent.to_string().parse().ok())
817 .unwrap_or(20),
818 orchestration_fee_micros: env
819 .var("ORCHESTRATION_FEE_MICROS")
820 .ok()
821 .and_then(|fee| fee.to_string().parse().ok())
822 .unwrap_or(100_000),
823 free: env.var("FREE_WHILE_BUILDING").is_ok_and(|v| v.to_string() == "true"),
824 ceilings: limits::Ceilings::from_env(&env),
825 deployments_monthly_cents: env
826 .var("DEPLOYMENTS_MONTHLY_CENTS")
827 .ok()
828 .and_then(|cents| cents.to_string().parse().ok())
829 .unwrap_or(500),
830 trial: {
831 let number = |name: &str| env.var(name).ok().and_then(|v| v.to_string().parse::<i64>().ok());
832 match (
833 number("TRIAL_WORKSPACE_MICROS"),
834 number("TRIAL_TOTAL_MICROS"),
835 env.var("TRIAL_UNTIL").ok().map(|v| v.to_string()),
836 ) {
837 (Some(per_workspace_micros), Some(total_micros), Some(until))
838 if per_workspace_micros > 0 && !until.is_empty() =>
839 {
840 Some(TrialConfig {
841 per_workspace_micros,
842 total_micros,
843 until,
844 })
845 }
846 _ => None,
847 }
848 },
849 })
850 }
851}
852
853#[event(scheduled)]
854async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
855 let Ok(billing) = Billing::from_env(&env) else {
856 return;
857 };
858 let keeper = keeper::Keeper::from_env(&env);
859 if let Err(error) = billing.settle_runs(&keeper).await {
860 worker::console_error!("settling runs failed: {error}");
861 }
The keeper reads Cloudflare as it really answers862 // Once a day, and at once if the costs were never checked: check every
863 // cost against what Cloudflare billed.
864 if event.cron() == keeper::DAILY || billing.never_checked().await.unwrap_or(false) {
Prices keep themselves current with what g1t pays865 if let Err(error) = billing.reconcile(&keeper).await {
866 worker::console_error!("checking costs against Cloudflare failed: {error}");
867 }
868 }
869}
870
Agents as a team: lifecycle, merge queue, billing and a new shell871#[event(fetch)]
872async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
873 let Some(method) = rpc_method(&request) else {
874 return Response::error("Not found", 404);
875 };
876 let body: serde_json::Value = request.json().await?;
Prices keep themselves current with what g1t pays877 let billing = Billing::from_env(&env)?;
Agents as a team: lifecycle, merge queue, billing and a new shell878 match method.as_str() {
879 "status" => reply(&billing.status()),
880 "account" => reply(&billing.account(args(body)?).await?),
881 "ledger" => reply(&billing.ledger(args(body)?).await?),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request882 "usage" => reply(&billing.usage(args(body)?).await?),
Agents as a team: lifecycle, merge queue, billing and a new shell883 "checkout" => reply(&billing.checkout(args(body)?).await?),
884 "confirm" => reply(&billing.confirm(args(body)?).await?),
885 "can_start" => reply(&billing.can_start(args(body)?).await?),
A free allowance on g1t's models, so anyone can try its agents886 "trial" => reply(&billing.trial(args(body)?).await?),
Agents as a team: lifecycle, merge queue, billing and a new shell887 "start_run" => reply(&billing.start_run(args(body)?).await?),
888 "finish_run" => reply(&billing.finish_run(args(body)?).await?),
Paid features: a workspace turns on Deployments with a monthly plan889 "features" => reply(&billing.features(args(body)?).await?),
890 "subscribe" => reply(&billing.subscribe(args(body)?).await?),
891 "confirm_subscription" => reply(&billing.confirm_subscription(args(body)?).await?),
892 "cancel_subscription" => reply(&billing.cancel_subscription(args(body)?).await?),
893 "has_feature" => reply(&billing.has_feature(args(body)?).await?),
894 "charge_feature" => reply(&billing.charge_feature(args(body)?).await?),
Every sandbox is metered by the second895 "record_sandbox" => reply(&billing.record_sandbox(args(body)?).await?),
Usage limits: unpaid usage can only go so far896 "limit" => reply(&billing.limit(args(body)?).await?),
897 "check_limit" => reply(&billing.check_limit(args(body)?).await?),
898 "set_spend_limit" => reply(&billing.set_spend_limit(args(body)?).await?),
Prices keep themselves current with what g1t pays899 "prices" => reply(&billing.prices().await?),
900 "note_pending" => reply(&billing.note_pending(args(body)?).await?),
Agents as a team: lifecycle, merge queue, billing and a new shell901 _ => Response::error("Unknown method", 404),
902 }
903}
904
905#[cfg(test)]
906mod tests {
907 use super::*;
908
909 #[test]
910 fn a_run_is_charged_its_cost_plus_the_margin() {
911 // $0.05 at 20% is six cents.
912 assert_eq!(charge_micros(0.05, 20), 60_000);
913 assert_eq!(charge_micros(1.0, 20), 1_200_000);
914 assert_eq!(charge_micros(0.05, 0), 50_000);
915 }
916
917 #[test]
918 fn fractions_of_a_millionth_round_up_and_nothing_costs_less_than_nothing() {
919 assert_eq!(charge_micros(0.000_000_4, 20), 2);
920 assert_eq!(charge_micros(0.0, 20), 0);
921 assert_eq!(charge_micros(-3.0, 20), 0);
922 }
923
924 #[test]
Every sandbox is metered by the second925 fn sandbox_seconds_are_charged_only_past_the_free_minutes() {
926 let free = sandbox_allowance::FREE_SECONDS;
927 assert_eq!(sandbox_billable(0, 600), 0);
928 assert_eq!(sandbox_billable(free - 100, 600), 500);
929 assert_eq!(sandbox_billable(free + 5, 600), 600);
930 }
931
932 #[test]
933 fn durations_read_plainly() {
934 assert_eq!(duration(40), "40s");
935 assert_eq!(duration(192), "3m 12s");
936 assert_eq!(duration(3720), "1h 2m");
937 }
938
939 #[test]
Agents as a team: lifecycle, merge queue, billing and a new shell940 fn an_absurd_cost_is_capped() {
941 assert_eq!(charge_micros(1e9, 20), 120 * MICROS_PER_DOLLAR);
942 }
943}