g1t/services/billing/src/lib.rs
//! The billing service: what agents cost, charged to the workspace they
//! worked for.
//!
//! A workspace buys credit with a card. Before the runner starts an agent
//! it asks here, and is refused if the workspace has none. When the
//! agent's sandbox finishes it reports what the model cost, and that plus
//! g1t's margin comes off the balance. Every change is a ledger entry, and
//! a balance is always the sum of its ledger.
//!
//! Paid features (deployments) are bought separately, as monthly plans;
//! see `features`. They are never free.
//!
//! Without a card processor configured the service says so and charges
//! nothing, so that g1t still runs where billing has not been set up.
//!
//! Reached only through service bindings; see `g1t_contracts::billing` for
//! the methods and their arguments.
mod accounts;
mod invoices;
mod sales;
mod webhooks;
mod features;
mod keeper;
mod limits;
mod stripe;
use g1t_contracts::billing::*;
use g1t_contracts::time::rfc3339;
use g1t_contracts::{FailureCode, Outcome, Role, new_id};
use g1t_contracts::billing::TermsKind;
use g1t_kit::{args, now_ms, reply, rpc_method};
use serde::Deserialize;
use sha2::{Digest, Sha256};
use worker::wasm_bindgen::JsValue;
use worker::{Context, D1Database, Env, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
use stripe::Stripe;
const MIN_TOP_UP_CENTS: u32 = 500;
const MAX_TOP_UP_CENTS: u32 = 50_000;
const LEDGER_PAGE: u32 = 100;
/// A run's reported cost is believed up to this much. A sandbox cannot
/// spend more in the time it has, so anything above is a fault.
const MAX_RUN_COST_USD: f64 = 100.0;
/// What a run is charged: its cost plus the margin, rounded up to a whole
/// millionth of a dollar.
pub fn charge_micros(cost_usd: f64, margin_percent: u32) -> i64 {
let cost_micros = (cost_usd.clamp(0.0, MAX_RUN_COST_USD) * MICROS_PER_DOLLAR as f64).ceil();
(cost_micros * f64::from(100 + margin_percent) / 100.0).ceil() as i64
}
fn hash(token: &str) -> String {
hex::encode(Sha256::digest(token.as_bytes()))
}
fn optional(value: Option<&str>) -> JsValue {
value.map_or(JsValue::NULL, JsValue::from)
}
#[derive(Deserialize)]
struct AccountRow {
balance_micros: i64,
customer_id: Option<String>,
}
#[derive(Deserialize)]
struct LedgerRow {
id: String,
kind: EntryKind,
amount_micros: i64,
description: String,
repo: Option<String>,
number: Option<u32>,
task: Option<String>,
model: Option<String>,
created_by: Option<String>,
created_at: String,
billed_to: Option<String>,
#[serde(default)]
workspace: Option<String>,
}
impl From<LedgerRow> for LedgerEntry {
fn from(row: LedgerRow) -> Self {
LedgerEntry {
id: row.id,
kind: row.kind,
amount_micros: row.amount_micros,
description: row.description,
repo: row.repo,
number: row.number,
task: row.task,
model: row.model,
billed_to: row.billed_to.unwrap_or_else(|| "g1t".to_owned()),
created_by: row.created_by,
created_at: row.created_at,
workspace: row.workspace,
}
}
}
#[derive(Deserialize)]
struct RunRow {
workspace: String,
repo: String,
number: u32,
task: String,
model: String,
token_hash: String,
billed_to: Option<String>,
}
impl RunRow {
fn own_provider(&self) -> bool {
self.billed_to.as_deref() == Some("workspace")
}
}
#[derive(Deserialize)]
struct CheckoutRow {
workspace: String,
created_by: String,
}
/// A row an `UPDATE … RETURNING` touched.
#[derive(Deserialize)]
struct Touched {
#[allow(dead_code)]
id: String,
}
struct Billing {
db: D1Database,
/// Absent when no card processor is configured.
stripe: Option<Stripe>,
margin_percent: u32,
/// Charged for a run on the workspace's own model provider.
orchestration_fee_micros: i64,
/// While g1t is being built out, nothing is charged (`FREE_WHILE_BUILDING`).
free: bool,
/// The free allowance on g1t's hosted models, when there is one.
trial: Option<TrialConfig>,
/// The Deployments plan's monthly price (`DEPLOYMENTS_MONTHLY_CENTS`).
deployments_monthly_cents: u32,
/// How far unpaid usage may go; see `limits`.
ceilings: limits::Ceilings,
/// `PREPAID_ONLY`: the old rule, that agents need credit first.
prepaid_only: bool,
}
/// `TRIAL_WORKSPACE_MICROS`, `TRIAL_TOTAL_MICROS` and `TRIAL_UNTIL`.
struct TrialConfig {
per_workspace_micros: i64,
total_micros: i64,
/// RFC 3339, in UTC.
until: String,
}
#[derive(serde::Deserialize)]
struct Sum {
micros: Option<i64>,
}
impl Billing {
fn status(&self) -> Status {
Status {
enabled: self.stripe.is_some(),
live: self.stripe.as_ref().is_some_and(Stripe::live),
free: self.free,
}
}
async fn row(&self, workspace: &str) -> Result<Option<AccountRow>> {
self.db
.prepare("SELECT balance_micros, customer_id FROM accounts WHERE workspace = ?")
.bind(&[workspace.into()])?
.first::<AccountRow>(None)
.await
}
async fn standing(&self, workspace: &str) -> Result<Account> {
let row = self.row(workspace).await?;
let card = match (&self.stripe, row.as_ref().and_then(|row| row.customer_id.as_deref())) {
(Some(stripe), Some(customer)) => stripe.card(customer).await.ok().flatten().map(|card| Card {
brand: card.brand,
last4: card.last4,
exp_month: card.exp_month,
exp_year: card.exp_year,
}),
_ => None,
};
Ok(Account {
workspace: workspace.to_owned(),
balance_micros: row.map_or(0, |row| row.balance_micros),
status: self.status(),
margin_percent: self.margin_percent,
orchestration_fee_micros: self.orchestration_fee_micros,
card,
})
}
/// The workspace's customer at Stripe, made the first time one is needed.
pub(crate) async fn customer_for(&self, workspace: &str) -> Result<String> {
let Some(stripe) = &self.stripe else {
return Err(worker::Error::RustError("payments are not set up".into()));
};
if let Some(customer) = self.row(workspace).await?.and_then(|row| row.customer_id) {
return Ok(customer);
}
let customer = stripe.create_customer(workspace).await?;
self.db
.prepare(
"INSERT INTO accounts (workspace, balance_micros, customer_id, created_at) VALUES (?1, 0, ?2, ?3)
ON CONFLICT (workspace) DO UPDATE SET customer_id = ?2",
)
.bind(&[workspace.into(), customer.as_str().into(), rfc3339(now_ms()).into()])?
.run()
.await?;
Ok(customer)
}
/// Stripe's hosted billing page for the workspace. Owners only.
async fn billing_portal(&self, a: BillingPortalArgs) -> Result<Outcome<Checkout>> {
let workspace = a.workspace.to_lowercase();
if a.actor.role_in(&workspace) != Some(Role::Owner) {
return Ok(Outcome::fail(FailureCode::Forbidden, "Only an owner can manage the workspace's billing."));
}
let Some(stripe) = &self.stripe else {
return Ok(Outcome::fail(FailureCode::Conflict, "Payments are not set up on this g1t."));
};
let customer = match self.customer_for(&workspace).await {
Ok(customer) => customer,
Err(error) => return Ok(Outcome::fail(FailureCode::Conflict, format!("Stripe could not be reached: {error}"))),
};
match stripe.portal_session(&customer, &a.return_url).await {
Ok(url) => Ok(Outcome::Ok(Checkout { url })),
Err(error) if stripe::is_missing(&error) => {
// The customer was removed at Stripe: a new one next time.
self.forget_customer(&workspace).await?;
Ok(Outcome::fail(FailureCode::Conflict, "Stripe no longer had this workspace's customer. Try again."))
}
Err(error) => Ok(Outcome::fail(FailureCode::Conflict, format!("Stripe's billing page could not be opened: {error}"))),
}
}
/// Adds a ledger entry and moves the balance by the same amount, as
/// one write.
#[allow(clippy::too_many_arguments)]
async fn enter(
&self,
workspace: &str,
kind: EntryKind,
amount_micros: i64,
description: &str,
reference: &str,
run: Option<&RunRow>,
cost_micros: Option<i64>,
created_by: Option<&str>,
customer: Option<&str>,
) -> Result<()> {
let now = now_ms();
let timestamp = rfc3339(now);
let kind = match kind {
EntryKind::TopUp => "top_up",
EntryKind::Usage => "usage",
};
self.db
.batch(vec![
self.db
.prepare(
"INSERT INTO ledger
(id, workspace, kind, amount_micros, description, repo, number, task,
model, cost_micros, reference, created_by, created_at, billed_to)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
)
.bind(&[
new_id("led", now).into(),
workspace.into(),
kind.into(),
// D1 takes numbers as doubles, which hold every
// amount this service will see exactly.
(amount_micros as f64).into(),
description.into(),
optional(run.map(|run| run.repo.as_str())),
run.map_or(JsValue::NULL, |run| run.number.into()),
optional(run.map(|run| run.task.as_str())),
optional(run.map(|run| run.model.as_str())),
cost_micros.map_or(JsValue::NULL, |cost| (cost as f64).into()),
reference.into(),
optional(created_by),
timestamp.as_str().into(),
run.map_or("g1t", |run| if run.own_provider() { "workspace" } else { "g1t" }).into(),
])?,
self.db
.prepare(
"INSERT INTO accounts (workspace, balance_micros, customer_id, created_at)
VALUES (?1, ?2, ?3, ?4)
ON CONFLICT (workspace) DO UPDATE SET
balance_micros = balance_micros + ?2,
customer_id = COALESCE(?3, customer_id)",
)
.bind(&[
workspace.into(),
(amount_micros as f64).into(),
optional(customer),
timestamp.as_str().into(),
])?,
])
.await?;
// Money in clears a card declined at the limit.
if kind == "top_up" {
self.db
.prepare("UPDATE limits SET autopay_failed_at = NULL, autopay_error = NULL WHERE workspace = ?")
.bind(&[workspace.into()])?
.run()
.await?;
}
Ok(())
}
async fn account(&self, a: AccountArgs) -> Result<Outcome<Account>> {
let workspace = a.workspace.to_lowercase();
if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
return Ok(members_only());
}
Ok(Outcome::Ok(self.standing(&workspace).await?))
}
async fn ledger(&self, a: AccountArgs) -> Result<Outcome<Vec<LedgerEntry>>> {
let workspace = a.workspace.to_lowercase();
if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
return Ok(members_only());
}
let rows = self
.db
.prepare("SELECT * FROM ledger WHERE workspace = ? ORDER BY id DESC LIMIT ?")
.bind(&[workspace.into(), LEDGER_PAGE.into()])?
.all()
.await?
.results::<LedgerRow>()?;
Ok(Outcome::Ok(
rows.into_iter().map(LedgerEntry::from).collect(),
))
}
async fn usage(&self, a: UsageArgs) -> Result<Outcome<Usage>> {
let workspace = a.workspace.to_lowercase();
if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
return Ok(members_only());
}
#[derive(serde::Deserialize)]
struct SliceRow {
key: Option<String>,
micros: Option<i64>,
runs: Option<u32>,
}
// While nothing is charged, what was used is what there is to show.
let measure = if self.free { "COALESCE(cost_micros, 0)" } else { "-amount_micros" };
let slices = |key: &str, limit: u32| {
format!(
"SELECT {key} AS key, SUM({measure}) AS micros, COUNT(*) AS runs FROM ledger
WHERE workspace = ?1 AND kind = 'usage' AND created_at >= ?2
GROUP BY 1 ORDER BY micros DESC LIMIT {limit}"
)
};
let query = |sql: String| {
let db = &self.db;
let workspace = workspace.clone();
let since = a.since.clone();
async move {
let rows = db
.prepare(sql)
.bind(&[workspace.into(), since.into()])?
.all()
.await?
.results::<SliceRow>()?;
Ok::<Vec<UsageSlice>, worker::Error>(
rows.into_iter()
.map(|row| UsageSlice {
key: row.key.unwrap_or_else(|| "other".to_owned()),
micros: row.micros.unwrap_or_default(),
runs: row.runs.unwrap_or_default(),
})
.collect(),
)
}
};
#[derive(serde::Deserialize)]
struct Totals {
spent: Option<i64>,
cost: Option<i64>,
provider: Option<i64>,
runs: Option<u32>,
added: Option<i64>,
}
let totals = self
.db
.prepare(
"SELECT
-SUM(CASE WHEN kind = 'usage' THEN amount_micros END) AS spent,
SUM(CASE WHEN kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' THEN cost_micros END) AS cost,
SUM(CASE WHEN kind = 'usage' AND billed_to = 'workspace' THEN cost_micros END) AS provider,
SUM(CASE WHEN kind = 'usage' THEN 1 ELSE 0 END) AS runs,
SUM(CASE WHEN kind = 'top_up' THEN amount_micros END) AS added
FROM ledger WHERE workspace = ?1 AND created_at >= ?2",
)
.bind(&[workspace.as_str().into(), a.since.as_str().into()])?
.first::<Totals>(None)
.await?;
let totals = totals.unwrap_or(Totals {
spent: None,
cost: None,
provider: None,
runs: None,
added: None,
});
Ok(Outcome::Ok(Usage {
spent_micros: totals.spent.unwrap_or_default(),
cost_micros: totals.cost.unwrap_or_default(),
provider_micros: totals.provider.unwrap_or_default(),
used_micros: totals.cost.unwrap_or_default() + totals.provider.unwrap_or_default(),
free: self.free,
runs: totals.runs.unwrap_or_default(),
added_micros: totals.added.unwrap_or_default(),
by_day: query(slices("substr(created_at, 1, 10) || '/' || COALESCE(task, 'other')", 400)).await?,
by_task: query(slices("task", 20)).await?,
by_repo: query(slices("repo", 20)).await?,
by_pull: query(slices("repo || '#' || number", 10)).await?,
by_model: query(slices("model", 10)).await?,
since: a.since,
}))
}
async fn checkout(&self, a: CheckoutArgs) -> Result<Outcome<Checkout>> {
let workspace = a.workspace.to_lowercase();
if a.actor.role_in(&workspace) != Some(Role::Owner) {
return Ok(Outcome::fail(
FailureCode::Forbidden,
"Only an owner can add credit to a workspace.",
));
}
let Some(stripe) = &self.stripe else {
return Ok(Outcome::fail(
FailureCode::Conflict,
"Payments are not set up on this g1t yet.",
));
};
if !(MIN_TOP_UP_CENTS..=MAX_TOP_UP_CENTS).contains(&a.amount_cents) {
return Ok(Outcome::fail(
FailureCode::Invalid,
format!(
"Add between ${} and ${} at a time.",
MIN_TOP_UP_CENTS / 100,
MAX_TOP_UP_CENTS / 100
),
));
}
let customer = self.row(&workspace).await?.and_then(|row| row.customer_id);
let session = match stripe
.start_checkout(&workspace, a.amount_cents, customer.as_deref(), &a.return_url)
.await
{
Ok(session) => session,
// A customer saved under another Stripe account: start afresh.
Err(error) if customer.is_some() && stripe::is_missing(&error) => {
self.forget_customer(&workspace).await?;
stripe.start_checkout(&workspace, a.amount_cents, None, &a.return_url).await?
}
Err(error) => return Err(error),
};
let Some(url) = session.url else {
return Err(worker::Error::RustError(
"the card processor returned no payment page".into(),
));
};
self.db
.prepare(
"INSERT INTO checkouts (id, workspace, amount_cents, created_by, created_at)
VALUES (?, ?, ?, ?, ?)",
)
.bind(&[
session.id.into(),
workspace.into(),
a.amount_cents.into(),
a.actor.username.into(),
rfc3339(now_ms()).into(),
])?
.run()
.await?;
Ok(Outcome::Ok(Checkout { url }))
}
/// Credits a payment if the processor says it was made and it has not
/// been credited before. The amount credited is what the processor
/// says was paid, not what anyone here remembers asking for.
async fn confirm(&self, a: ConfirmArgs) -> Result<Outcome<Account>> {
let workspace = a.workspace.to_lowercase();
if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
return Ok(members_only());
}
let (Some(stripe), Some(checkout)) = (
&self.stripe,
self.db
.prepare(
"SELECT workspace, created_by FROM checkouts
WHERE id = ? AND workspace = ? AND status = 'open'",
)
.bind(&[a.session.as_str().into(), workspace.as_str().into()])?
.first::<CheckoutRow>(None)
.await?,
) else {
// Unknown, someone else's, or already credited: nothing to do.
return Ok(Outcome::Ok(self.standing(&workspace).await?));
};
let session = stripe.session(&a.session).await?;
let paid = session
.amount_total
.filter(|_| session.payment_status == "paid");
if let Some(cents) = paid {
// Only whoever flips it from open to paid enters the credit.
let claimed = self
.db
.prepare(
"UPDATE checkouts SET status = 'paid' WHERE id = ? AND status = 'open'
RETURNING id",
)
.bind(&[a.session.as_str().into()])?
.first::<Touched>(None)
.await?;
if claimed.is_some() {
self.enter(
&checkout.workspace,
EntryKind::TopUp,
i64::from(cents) * MICROS_PER_DOLLAR / 100,
"Credit added by card",
&session.id,
None,
None,
Some(&checkout.created_by),
session.customer.as_deref(),
)
.await?;
}
}
Ok(Outcome::Ok(self.standing(&workspace).await?))
}
/// Drops a saved customer the card processor no longer knows.
pub(crate) async fn forget_customer(&self, workspace: &str) -> Result<()> {
self.db
.prepare("UPDATE accounts SET customer_id = NULL WHERE workspace = ?")
.bind(&[workspace.into()])?
.run()
.await?;
Ok(())
}
/// A refusal if the workspace has no credit to start an agent with.
async fn out_of_credit<T>(&self, workspace: &str) -> Result<Option<Outcome<T>>> {
// Billing is postpaid: usage limits decide whether work starts
// (see `limits`), and credit is a prepayment that lowers what is
// owed. A balance no longer has to be positive to start.
if self.free || !self.prepaid_only {
return Ok(None);
}
let balance = self
.row(workspace)
.await?
.map_or(0, |row| row.balance_micros);
Ok((balance <= 0).then(|| {
Outcome::fail(
FailureCode::PaymentRequired,
format!(
"The {workspace} workspace has no agent credit. An owner can add some under Billing on the workspace's page."
),
)
}))
}
/// A workspace's free allowance on g1t's hosted models: what its runs
/// there have cost against its share, and the pool everyone draws on.
/// How much of a hosted model run's cost the workspace's free allowance
/// covers, if it is still open.
async fn trial_covers(&self, workspace: &str, cost_micros: i64) -> Result<i64> {
let trial = self.trial(TrialArgs { workspace: workspace.to_owned(), exempt: vec![] }).await?;
if !trial.open {
return Ok(0);
}
Ok(cost_micros.min((trial.limit_micros - trial.used_micros).max(0)))
}
async fn trial(&self, a: TrialArgs) -> Result<Trial> {
let workspace = a.workspace.to_lowercase();
let Some(config) = &self.trial else {
return Ok(Trial {
open: false,
used_micros: 0,
limit_micros: 0,
ends_at: None,
reason: Some("off".to_owned()),
});
};
let used = self
.db
.prepare(
"SELECT SUM(cost_micros) AS micros FROM ledger
WHERE kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' AND COALESCE(task, '') NOT IN ('sandbox', 'deployments') AND workspace = ?",
)
.bind(&[workspace.as_str().into()])?
.first::<Sum>(None)
.await?
.and_then(|sum| sum.micros)
.unwrap_or_default();
// Everyone's, but for the workspaces open to hosted models anyway.
let exempt: Vec<String> = a.exempt.iter().map(|name| name.trim().to_lowercase()).collect();
let marks = vec!["?"; exempt.len().max(1)].join(", ");
let mut values: Vec<JsValue> = exempt.iter().map(|name| JsValue::from(name.as_str())).collect();
if values.is_empty() {
values.push(JsValue::from(""));
}
let pooled = self
.db
.prepare(format!(
"SELECT SUM(cost_micros) AS micros FROM ledger
WHERE kind = 'usage' AND COALESCE(billed_to, 'g1t') = 'g1t' AND COALESCE(task, '') NOT IN ('sandbox', 'deployments') AND workspace NOT IN ({marks})"
))
.bind(&values)?
.first::<Sum>(None)
.await?
.and_then(|sum| sum.micros)
.unwrap_or_default();
let reason = if rfc3339(now_ms()) >= config.until {
Some("ended")
} else if used >= config.per_workspace_micros {
Some("used")
} else if pooled >= config.total_micros {
Some("pool")
} else {
None
};
Ok(Trial {
open: reason.is_none(),
used_micros: used,
limit_micros: config.per_workspace_micros,
ends_at: Some(config.until.clone()),
reason: reason.map(str::to_owned),
})
}
async fn can_start(&self, a: CanStartArgs) -> Result<Outcome<bool>> {
if self.stripe.is_none() {
return Ok(Outcome::Ok(true));
}
if let Some(stopped) = self.stopped(&a.workspace).await? {
return Ok(stopped);
}
Ok(self
.out_of_credit(&a.workspace.to_lowercase())
.await?
.unwrap_or(Outcome::Ok(true)))
}
async fn start_run(&self, a: StartRunArgs) -> Result<Outcome<Option<RunTicket>>> {
if self.stripe.is_none() {
return Ok(Outcome::Ok(None));
}
let workspace = a.workspace.to_lowercase();
if let Some(stopped) = self.stopped(&workspace).await? {
return Ok(stopped);
}
if let Some(refused) = self.out_of_credit(&workspace).await? {
return Ok(refused);
}
let now = now_ms();
let run_id = new_id("run", now);
let mut bytes = [0u8; 32];
getrandom::getrandom(&mut bytes).expect("no source of randomness");
let token = hex::encode(bytes);
self.db
.prepare(
"INSERT INTO runs (id, workspace, repo, number, task, model, token_hash, created_at, billed_to, session_id)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
)
.bind(&[
run_id.as_str().into(),
workspace.into(),
format!("{}/{}", a.repo.namespace, a.repo.name).into(),
a.number.into(),
a.task.into(),
a.model.into(),
hash(&token).into(),
rfc3339(now).into(),
if a.billed_to == "workspace" { "workspace" } else { "g1t" }.into(),
optional(a.session.as_deref().filter(|_| a.billed_to != "workspace")),
])?
.run()
.await?;
Ok(Outcome::Ok(Some(RunTicket { run_id, token })))
}
async fn finish_run(&self, a: FinishRunArgs) -> Result<Outcome<bool>> {
let run = self
.db
.prepare(
"SELECT workspace, repo, number, task, model, token_hash, billed_to FROM runs
WHERE id = ? AND finished_at IS NULL",
)
.bind(&[a.run_id.as_str().into()])?
.first::<RunRow>(None)
.await?;
let Some(run) = run.filter(|run| run.token_hash == hash(&a.token)) else {
return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
};
if !a.cost_usd.is_finite() || a.cost_usd < 0.0 {
return Ok(Outcome::fail(FailureCode::Invalid, "That is not a cost."));
}
// Only whoever closes the run charges for it.
let claimed = self
.db
.prepare(
"UPDATE runs SET finished_at = ? WHERE id = ? AND finished_at IS NULL RETURNING id",
)
.bind(&[rfc3339(now_ms()).into(), a.run_id.as_str().into()])?
.first::<Touched>(None)
.await?;
if claimed.is_none() {
return Ok(Outcome::Ok(false));
}
// On the workspace's own provider, the model was paid for there:
// g1t charges its fee, and keeps the provider's cost to show.
let base = if run.own_provider() {
self.orchestration_fee_micros
} else {
// The free allowance on g1t's models covers what it can.
let cost = charge_micros(a.cost_usd, 0);
let covered = self.trial_covers(&run.workspace, cost).await?;
charge_micros((cost - covered) as f64 / MICROS_PER_DOLLAR as f64, self.margin_percent)
};
let (charge, terms_note) = self.charged(&run.workspace, base).await?;
let mut description = match run.task.as_str() {
"plan" => format!("Planning for {}", run.repo),
"review" => format!("Review of {}#{}", run.repo, run.number),
"update" => format!("Catching up {}#{}", run.repo, run.number),
_ => format!("Work on {}#{}", run.repo, run.number),
};
if run.own_provider() {
description.push_str(", on your own model provider");
}
description.push_str(&terms_note);
self.enter(
&run.workspace,
EntryKind::Usage,
-charge,
&description,
&a.run_id,
Some(&run),
Some(charge_micros(a.cost_usd, 0)),
None,
None,
)
.await?;
Ok(Outcome::Ok(true))
}
}
impl Billing {
/// Records how long a sandbox ran: its cost always, and a charge for
/// the seconds past the month's free minutes.
async fn record_sandbox(&self, a: RecordSandboxArgs) -> Result<Outcome<bool>> {
if self.stripe.is_none() || a.seconds == 0 {
return Ok(Outcome::Ok(false));
}
let workspace = a.workspace.to_lowercase();
let seen = self
.db
.prepare("SELECT id FROM ledger WHERE reference = ?")
.bind(&[a.reference.as_str().into()])?
.first::<Touched>(None)
.await?;
if seen.is_some() {
return Ok(Outcome::Ok(false));
}
let now = now_ms();
let timestamp = rfc3339(now);
let month = ×tamp[..7];
#[derive(Deserialize)]
struct Used {
seconds: i64,
}
let seconds = i64::from(a.seconds);
let after = self
.db
.prepare(
"INSERT INTO sandbox_months (workspace, month, seconds) VALUES (?1, ?2, ?3)
ON CONFLICT (workspace, month) DO UPDATE SET seconds = seconds + ?3
RETURNING seconds",
)
.bind(&[workspace.as_str().into(), month.into(), (seconds as f64).into()])?
.first::<Used>(None)
.await?
.map_or(seconds, |used| used.seconds);
let billable = sandbox_billable(after - seconds, seconds);
// From the price book, which follows what Cloudflare bills g1t.
let (cost_per_second, price_per_second) = self.price("sandbox_second").await?.unwrap_or((
sandbox_allowance::COST_MICROS_PER_SECOND as f64,
sandbox_allowance::MICROS_PER_SECOND as f64,
));
let (charge, terms_note) = self.charged(&workspace, (billable as f64 * price_per_second).ceil() as i64).await?;
let mut description = format!("{}: {} of sandbox time", a.description, duration(seconds));
if billable < seconds {
description.push_str(if billable == 0 {
", within the month's free minutes"
} else {
", partly within the month's free minutes"
});
}
if billable > 0 {
description.push_str(&terms_note);
}
self.db
.batch(vec![
self.db
.prepare(
"INSERT INTO ledger
(id, workspace, kind, amount_micros, description, repo, task,
cost_micros, reference, created_at, billed_to)
VALUES (?, ?, 'usage', ?, ?, ?, 'sandbox', ?, ?, ?, 'g1t')",
)
.bind(&[
new_id("led", now).into(),
workspace.as_str().into(),
(-(charge as f64)).into(),
description.as_str().into(),
optional(a.repo.as_deref()),
(seconds as f64 * cost_per_second).ceil().into(),
a.reference.as_str().into(),
timestamp.as_str().into(),
])?,
self.db
.prepare(
"INSERT INTO accounts (workspace, balance_micros, created_at)
VALUES (?1, ?2, ?3)
ON CONFLICT (workspace) DO UPDATE SET balance_micros = balance_micros + ?2",
)
.bind(&[
workspace.as_str().into(),
(-(charge as f64)).into(),
timestamp.as_str().into(),
])?,
])
.await?;
Ok(Outcome::Ok(true))
}
}
/// Of `seconds` used after `before` this month, how many are past the
/// free minutes.
fn sandbox_billable(before: i64, seconds: i64) -> i64 {
let free_left = (sandbox_allowance::FREE_SECONDS - before).max(0);
(seconds - free_left).max(0)
}
/// `1h 2m`, `3m 12s` or `40s`.
fn duration(seconds: i64) -> String {
let (h, m, s) = (seconds / 3600, seconds % 3600 / 60, seconds % 60);
if h > 0 {
format!("{h}h {m}m")
} else if m > 0 {
format!("{m}m {s}s")
} else {
format!("{s}s")
}
}
impl Billing {
/// What a workspace is charged for something that would be `base`:
/// nothing while g1t is free, or as its account's terms say. With a
/// note for the statement when it differs.
pub(crate) async fn charged(&self, workspace: &str, base: i64) -> Result<(i64, String)> {
if self.free {
return Ok((0, " (free while g1t is being built out)".to_owned()));
}
let terms = self.terms_of(workspace).await?;
let charge = terms.apply(base);
let note = match terms.kind {
TermsKind::Comped => " (comped)".to_owned(),
TermsKind::Custom if terms.discount_percent > 0 && base > 0 => format!(" ({}% off)", terms.discount_percent),
_ => String::new(),
};
Ok((charge, note))
}
}
fn members_only<T>() -> Outcome<T> {
Outcome::fail(
FailureCode::Forbidden,
"Only members can see a workspace's billing.",
)
}
impl Billing {
fn from_env(env: &Env) -> Result<Self> {
Ok(Billing {
db: env.d1("DB")?,
stripe: env
.secret("STRIPE_SECRET_KEY")
.ok()
.map(|key| key.to_string())
.filter(|key| !key.is_empty())
.map(Stripe::new),
margin_percent: env
.var("MARGIN_PERCENT")
.ok()
.and_then(|percent| percent.to_string().parse().ok())
.unwrap_or(20),
orchestration_fee_micros: env
.var("ORCHESTRATION_FEE_MICROS")
.ok()
.and_then(|fee| fee.to_string().parse().ok())
.unwrap_or(100_000),
free: env.var("FREE_WHILE_BUILDING").is_ok_and(|v| v.to_string() == "true"),
ceilings: limits::Ceilings::from_env(&env),
prepaid_only: env.var("PREPAID_ONLY").is_ok_and(|v| v.to_string() == "true"),
deployments_monthly_cents: env
.var("DEPLOYMENTS_MONTHLY_CENTS")
.ok()
.and_then(|cents| cents.to_string().parse().ok())
.unwrap_or(500),
trial: {
let number = |name: &str| env.var(name).ok().and_then(|v| v.to_string().parse::<i64>().ok());
match (
number("TRIAL_WORKSPACE_MICROS"),
number("TRIAL_TOTAL_MICROS"),
env.var("TRIAL_UNTIL").ok().map(|v| v.to_string()),
) {
(Some(per_workspace_micros), Some(total_micros), Some(until))
if per_workspace_micros > 0 && !until.is_empty() =>
{
Some(TrialConfig {
per_workspace_micros,
total_micros,
until,
})
}
_ => None,
}
},
})
}
}
#[event(scheduled)]
async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
let Ok(billing) = Billing::from_env(&env) else {
return;
};
let keeper = keeper::Keeper::from_env(&env);
if let Err(error) = billing.settle_runs(&keeper).await {
worker::console_error!("settling runs failed: {error}");
}
if let Err(error) = billing.autopay().await {
worker::console_error!("paying at the limit failed: {error}");
}
if let Err(error) = billing.close_months().await {
worker::console_error!("closing the month failed: {error}");
}
if let Err(error) = billing.invoice_enterprises().await {
worker::console_error!("invoicing enterprises failed: {error}");
}
if let Ok(identity) = env.service("IDENTITY") {
if let Err(error) = billing.warn_limits(&identity).await {
worker::console_error!("warning owners failed: {error}");
}
}
// Once a day, and at once if the costs were never checked: check every
// cost against what Cloudflare billed.
if event.cron() == keeper::DAILY || billing.never_checked().await.unwrap_or(false) {
if let Err(error) = billing.reconcile(&keeper).await {
worker::console_error!("checking costs against Cloudflare failed: {error}");
}
}
}
#[event(fetch)]
async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
let Some(method) = rpc_method(&request) else {
return Response::error("Not found", 404);
};
let body: serde_json::Value = request.json().await?;
let billing = Billing::from_env(&env)?;
match method.as_str() {
"status" => reply(&billing.status()),
"account" => reply(&billing.account(args(body)?).await?),
"ledger" => reply(&billing.ledger(args(body)?).await?),
"usage" => reply(&billing.usage(args(body)?).await?),
"checkout" => reply(&billing.checkout(args(body)?).await?),
"confirm" => reply(&billing.confirm(args(body)?).await?),
"can_start" => reply(&billing.can_start(args(body)?).await?),
"trial" => reply(&billing.trial(args(body)?).await?),
"start_run" => reply(&billing.start_run(args(body)?).await?),
"finish_run" => reply(&billing.finish_run(args(body)?).await?),
"features" => reply(&billing.features(args(body)?).await?),
"subscribe" => reply(&billing.subscribe(args(body)?).await?),
"confirm_subscription" => reply(&billing.confirm_subscription(args(body)?).await?),
"cancel_subscription" => reply(&billing.cancel_subscription(args(body)?).await?),
"has_feature" => reply(&billing.has_feature(args(body)?).await?),
"charge_feature" => reply(&billing.charge_feature(args(body)?).await?),
"record_sandbox" => reply(&billing.record_sandbox(args(body)?).await?),
"limit" => reply(&billing.limit(args(body)?).await?),
"check_limit" => reply(&billing.check_limit(args(body)?).await?),
"set_spend_limit" => reply(&billing.set_spend_limit(args(body)?).await?),
"prices" => reply(&billing.prices().await?),
"billing_portal" => reply(&billing.billing_portal(args(body)?).await?),
"admin_billing_link" => reply(&billing.admin_billing_link(args(body)?).await?),
"admin_stripe" => reply(&billing.admin_stripe(args(body)?).await?),
"admin_enterprise_billing" => reply(&billing.admin_enterprise_billing(args(body)?).await?),
"admin_invoice_enterprise" => reply(&billing.admin_invoice_enterprise(args(body)?).await?),
"stripe_webhook" => reply(&billing.stripe_webhook(args(body)?).await?),
"invoices" => reply(&billing.invoices(args(body)?).await?),
"admin_workspace_invoices" => {
let a: AdminWorkspaceInvoicesArgs = args(body)?;
reply(&billing.workspace_invoices(&a.workspace.to_lowercase()).await?)
}
"admin_signals" => reply(&billing.admin_signals(args(body)?).await?),
"admin_overview" => reply(&billing.admin_overview(args(body)?).await?),
"admin_sales" => reply(&billing.admin_sales(args(body)?).await?),
"admin_set_sales" => reply(&billing.admin_set_sales(args(body)?).await?),
"admin_add_note" => reply(&billing.admin_add_note(args(body)?).await?),
"admin_invoices" => reply(&billing.admin_invoices(args(body)?).await?),
"admin_audit" => reply(&billing.admin_audit(args(body)?).await?),
"note_pending" => reply(&billing.note_pending(args(body)?).await?),
"admin_accounts" => reply(&billing.admin_accounts(args(body)?).await?),
"admin_account" => reply(&billing.admin_account(args(body)?).await?),
"admin_set_terms" => reply(&billing.admin_set_terms(args(body)?).await?),
"admin_create_enterprise" => reply(&billing.admin_create_enterprise(args(body)?).await?),
"admin_attach" => reply(&billing.admin_attach(args(body)?).await?),
"admin_credit" => reply(&billing.admin_credit(args(body)?).await?),
_ => Response::error("Unknown method", 404),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_run_is_charged_its_cost_plus_the_margin() {
// $0.05 at 20% is six cents.
assert_eq!(charge_micros(0.05, 20), 60_000);
assert_eq!(charge_micros(1.0, 20), 1_200_000);
assert_eq!(charge_micros(0.05, 0), 50_000);
}
#[test]
fn fractions_of_a_millionth_round_up_and_nothing_costs_less_than_nothing() {
assert_eq!(charge_micros(0.000_000_4, 20), 2);
assert_eq!(charge_micros(0.0, 20), 0);
assert_eq!(charge_micros(-3.0, 20), 0);
}
#[test]
fn sandbox_seconds_are_charged_only_past_the_free_minutes() {
let free = sandbox_allowance::FREE_SECONDS;
assert_eq!(sandbox_billable(0, 600), 0);
assert_eq!(sandbox_billable(free - 100, 600), 500);
assert_eq!(sandbox_billable(free + 5, 600), 600);
}
#[test]
fn durations_read_plainly() {
assert_eq!(duration(40), "40s");
assert_eq!(duration(192), "3m 12s");
assert_eq!(duration(3720), "1h 2m");
}
#[test]
fn an_absurd_cost_is_capped() {
assert_eq!(charge_micros(1e9, 20), 120 * MICROS_PER_DOLLAR);
}
}