Skip to content

g1t/services/billing/src/webhooks.rs

1,079 lines49,115 bytesCodeBlame
1//! Stripe telling billing what happened, and enterprise invoices.
2//!
3//! Most of billing asks Stripe when it needs to know: a payment page is
4//! confirmed when the person comes back, a plan is checked when its period
5//! ends. That misses whatever happens while no one is looking: a page paid
6//! for and closed, a renewal that failed, a refund, a dispute, an invoice
7//! paid by bank transfer a week later. Stripe sends each as an event to
8//! `https://api.g1t.sh/stripe/webhook`; the API passes the raw body and its
9//! signature here, untouched.
10//!
11//! - The endpoint is registered by billing itself, from sudo, once per
12//! mode, and its signing secret is kept in billing's database. It is never
13//! shown, and nothing can be posted here without it.
14//! - Each event is handled once, by id, and recorded with what was done.
15//! - Every handler is safe alongside the paths that ask Stripe directly:
16//! both claim the same rows.
17//!
18//! Enterprises are invoiced: one Stripe invoice per month (or sooner, from
19//! sudo), with a line per workspace for what it owes, sent to the
20//! enterprise's billing email and paid on Stripe's hosted invoice page.
21//! When it is paid, each workspace is credited its line; when it goes
22//! overdue, their work stops until it is paid.
23
24use g1t_contracts::billing::{
25 AdminEnterpriseAddressArgs, AdminEnterpriseBillingArgs, AdminInvoiceEnterpriseArgs, AdminStripeArgs, BillingAccount, EnterpriseInvoice,
26 EntryKind, InvoiceLine, StripeEventSummary, StripeStatus, StripeWebhook, StripeWebhookArgs,
27};
28use g1t_contracts::time::rfc3339;
29use g1t_contracts::{FailureCode, Outcome};
30use g1t_kit::now_ms;
31use hmac::{Hmac, Mac};
32use serde::Deserialize;
33use serde_json::Value;
34use sha2::Sha256;
35use worker::Result;
36use worker::wasm_bindgen::JsValue;
37
38use crate::Billing;
39
40/// Where Stripe sends events.
41pub(crate) const WEBHOOK_URL: &str = "https://api.g1t.sh/stripe/webhook";
42
43/// The events billing acts on.
44pub(crate) const EVENTS: &[&str] = &[
45 "checkout.session.completed",
46 // A prepayment by bank transfer: the page completes when the transfer
47 // is set up, and this comes when the money arrives.
48 "checkout.session.async_payment_succeeded",
49 "customer.subscription.updated",
50 "customer.subscription.deleted",
51 "invoice.paid",
52 "invoice.payment_failed",
53 "invoice.overdue",
54 "invoice.voided",
55 "charge.refunded",
56 "charge.dispute.created",
57 "charge.dispute.closed",
58 // The saved card, kept on the account row (stripe_sync.rs).
59 "customer.updated",
60 "payment_method.attached",
61 "payment_method.detached",
62 "payment_method.updated",
63 "payment_method.automatically_updated",
64 "setup_intent.succeeded",
65];
66
67/// The subscription an invoice is for: `subscription` at the API version
68/// billing asks in, `parent.subscription_details.subscription` in events
69/// sent at a newer one (a destination keeps the version it was made at).
70pub(crate) fn invoice_subscription(invoice: &Value) -> Option<&str> {
71 invoice["subscription"].as_str().or_else(|| invoice["parent"]["subscription_details"]["subscription"].as_str())
72}
73
74/// Events billing handles that `has` does not include, in billing's order.
75pub(crate) fn missing_events(has: &[String]) -> Vec<String> {
76 // `*` is every event.
77 if has.iter().any(|event| event == "*") {
78 return Vec::new();
79 }
80 EVENTS.iter().filter(|event| !has.iter().any(|h| h == *event)).map(|event| (*event).to_owned()).collect()
81}
82
83/// How old a signed event may be, so a captured one cannot be replayed.
84const TOLERANCE_SECONDS: i64 = 5 * 60;
85
86/// Whether `header` (`t=…,v1=…`) signs `payload` with `secret`, within the
87/// tolerance of `now_seconds`.
88pub(crate) fn verify(payload: &str, header: &str, secret: &str, now_seconds: i64) -> bool {
89 let mut timestamp = None;
90 let mut signatures = vec![];
91 for part in header.split(',') {
92 match part.trim().split_once('=') {
93 Some(("t", value)) => timestamp = value.parse::<i64>().ok(),
94 Some(("v1", value)) => signatures.push(value.to_owned()),
95 _ => {}
96 }
97 }
98 let Some(timestamp) = timestamp else { return false };
99 if (now_seconds - timestamp).abs() > TOLERANCE_SECONDS {
100 return false;
101 }
102 let Ok(mut mac) = Hmac::<Sha256>::new_from_slice(secret.as_bytes()) else { return false };
103 mac.update(format!("{timestamp}.{payload}").as_bytes());
104 let expected = mac.finalize().into_bytes();
105 signatures.iter().any(|signature| {
106 hex::decode(signature).is_ok_and(|given| {
107 // Constant time: compare every byte whatever the first difference.
108 given.len() == expected.len() && given.iter().zip(expected.iter()).fold(0u8, |acc, (a, b)| acc | (a ^ b)) == 0
109 })
110 })
111}
112
113/// The destination at billing's address, as Stripe lists it.
114#[derive(Deserialize)]
115pub(crate) struct Destination {
116 pub(crate) id: String,
117 url: String,
118 /// `enabled` or `disabled`.
119 pub(crate) status: String,
120 pub(crate) enabled_events: Vec<String>,
121 created: i64,
122}
123
124#[derive(Deserialize)]
125struct EventRow {
126 id: String,
127 r#type: String,
128 outcome: String,
129 received_at: String,
130}
131
132#[derive(Deserialize)]
133struct InvoiceRow {
134 invoice_id: String,
135 period: String,
136 amount_micros: i64,
137 status: String,
138 hosted_url: Option<String>,
139 created_at: String,
140}
141
142#[derive(Deserialize)]
143struct LineRow {
144 workspace: String,
145 amount_micros: i64,
146}
147
148/// The idempotency key of an enterprise's invoice: the same enterprise,
149/// period and amounts are the same invoice however often it is attempted.
150pub(crate) fn enterprise_invoice_key(id: &str, period: &str, lines: &[(&str, i64)]) -> String {
151 // FNV-1a over every line, so the key stays short however many there are.
152 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
153 for (workspace, cents) in lines {
154 for byte in format!("{workspace}={cents};").bytes() {
155 hash = (hash ^ u64::from(byte)).wrapping_mul(0x0100_0000_01b3);
156 }
157 }
158 format!("ent-invoice/{id}/{period}/{hash:016x}")
159}
160
161impl Billing {
162 pub(crate) fn mode(&self) -> &'static str {
163 match &self.stripe {
164 None => "off",
165 Some(stripe) if stripe.live() => "live",
166 Some(_) => "test",
167 }
168 }
169
170 /// The destination at billing's address in Stripe, for this key's
171 /// mode: an enabled one first. None when there is none.
172 pub(crate) async fn destination(&self) -> Result<Option<Destination>> {
173 let Some(stripe) = &self.stripe else { return Ok(None) };
174 #[derive(Deserialize)]
175 struct List {
176 data: Vec<Destination>,
177 }
178 let mut ours: Vec<Destination> =
179 stripe.get::<List>("/webhook_endpoints?limit=100").await?.data.into_iter().filter(|d| d.url == WEBHOOK_URL).collect();
180 ours.sort_by_key(|d| d.status != "enabled");
181 Ok(ours.into_iter().next())
182 }
183
184 // --- Staff ------------------------------------------------------------
185
186 pub(crate) async fn admin_stripe(&self, a: AdminStripeArgs) -> Result<StripeStatus> {
187 let mut error = None;
188 if a.fix
189 && let Err(e) = self.keep_endpoint(a.by.as_deref().unwrap_or("sudo")).await
190 {
191 error = Some(e.to_string());
192 }
193 let (webhook, missing_events) = match self.destination().await {
194 Ok(Some(d)) => {
195 let missing = missing_events(&d.enabled_events);
196 let webhook = StripeWebhook {
197 url: d.url,
198 endpoint_id: d.id,
199 status: d.status,
200 events: d.enabled_events,
201 created_at: rfc3339(d.created.max(0) as u64 * 1000),
202 };
203 (Some(webhook), missing)
204 }
205 Ok(None) => (None, Vec::new()),
206 Err(e) => {
207 error = error.or(Some(format!("Stripe could not be read: {e}")));
208 (None, Vec::new())
209 }
210 };
211 let recent_events = self
212 .db
213 .prepare("SELECT * FROM stripe_events ORDER BY received_at DESC LIMIT 25")
214 .all()
215 .await?
216 .results::<EventRow>()?
217 .into_iter()
218 .map(|row| StripeEventSummary { id: row.id, kind: row.r#type, outcome: row.outcome, received_at: row.received_at })
219 .collect();
220 Ok(StripeStatus {
221 mode: self.mode().to_owned(),
222 secret_set: self.webhook_secret.is_some(),
223 webhook,
224 missing_events,
225 recent_events,
226 error,
227 })
228 }
229
230 // --- Events -----------------------------------------------------------
231
232 pub(crate) async fn stripe_webhook(&self, a: StripeWebhookArgs) -> Result<Outcome<bool>> {
233 let Some(secret) = &self.webhook_secret else {
234 return Ok(Outcome::fail(
235 FailureCode::Conflict,
236 "STRIPE_WEBHOOK_SECRET is not set, so billing cannot check that events come from Stripe.",
237 ));
238 };
239 let now_seconds = (now_ms() / 1000) as i64;
240 if !verify(&a.payload, &a.signature, secret, now_seconds) {
241 return Ok(Outcome::fail(FailureCode::Forbidden, "The signature does not match."));
242 }
243 let event: Value = serde_json::from_str(&a.payload).map_err(|e| worker::Error::RustError(e.to_string()))?;
244 if event["id"].as_str().unwrap_or_default().is_empty() {
245 return Ok(Outcome::fail(FailureCode::Invalid, "Not an event."));
246 }
247 Ok(Outcome::Ok(self.process_event(&event).await?))
248 }
249
250 /// Handles a Stripe event once, however often it arrives: by webhook,
251 /// or again from the event list (stripe_sync.rs). False when it was
252 /// seen before.
253 pub(crate) async fn process_event(&self, event: &Value) -> Result<bool> {
254 let id = event["id"].as_str().unwrap_or_default().to_owned();
255 let kind = event["type"].as_str().unwrap_or_default().to_owned();
256 // Once each: the first to record it handles it.
257 let claimed = self
258 .db
259 .prepare("INSERT OR IGNORE INTO stripe_events (id, type, outcome, received_at) VALUES (?, ?, 'handling', ?) RETURNING id")
260 .bind(&[id.as_str().into(), kind.as_str().into(), rfc3339(now_ms()).into()])?
261 .first::<Value>(None)
262 .await?;
263 if claimed.is_none() {
264 return Ok(false);
265 }
266 let outcome = match self.handle(&kind, event).await {
267 Ok(outcome) => outcome,
268 Err(error) => {
269 // Let Stripe send it again: forget it was seen.
270 self.db.prepare("DELETE FROM stripe_events WHERE id = ?").bind(&[id.as_str().into()])?.run().await?;
271 return Err(error);
272 }
273 };
274 self.db
275 .prepare("UPDATE stripe_events SET outcome = ? WHERE id = ?")
276 .bind(&[outcome.as_str().into(), id.as_str().into()])?
277 .run()
278 .await?;
279 Ok(true)
280 }
281
282 async fn handle(&self, kind: &str, event: &Value) -> Result<String> {
283 let object = &event["data"]["object"];
284 let text = |key: &str| object[key].as_str().unwrap_or_default().to_owned();
285 // The customer whose saved card may have changed. A detached card
286 // has no customer any more; the event says whose it was.
287 let card_owner = match kind {
288 "customer.updated" => object["id"].as_str(),
289 "payment_method.detached" => event["data"]["previous_attributes"]["customer"].as_str(),
290 _ => object["customer"].as_str(),
291 };
292 Ok(match kind {
293 "checkout.session.completed" | "checkout.session.async_payment_succeeded" => self.settle_checkout(&text("id")).await?,
294 "customer.subscription.updated" | "customer.subscription.deleted" => {
295 self.settle_subscription(&text("id")).await?
296 }
297 "invoice.paid" => {
298 if let Some(subscription) = invoice_subscription(object) {
299 self.plan_paid(subscription, object).await?;
300 self.settle_subscription(subscription).await?
301 } else if let Some(done) = self.workspace_invoice_paid(&text("id")).await? {
302 done
303 } else {
304 self.enterprise_invoice_paid(object).await?
305 }
306 }
307 "invoice.payment_failed" => match (invoice_subscription(object), object["metadata"]["g1t_workspace"].as_str()) {
308 (Some(subscription), _) => self.settle_subscription(subscription).await?,
309 (None, Some(tagged)) => {
310 // The invoice's own row names the workspace as it is
311 // now; the metadata keeps the slug it was sent under.
312 let workspace = self.workspace_of_invoice(&text("id")).await?.unwrap_or_else(|| tagged.to_owned());
313 let workspace = workspace.as_str();
314 self.mark_declined(workspace, "the card was declined for an invoice").await?;
315 format!("{workspace}: invoice payment failed; work stopped")
316 }
317 _ => "ignored: not g1t's".to_owned(),
318 },
319 "invoice.overdue" => self.enterprise_invoice_status(&text("id"), "overdue").await?,
320 "invoice.voided" => self.enterprise_invoice_status(&text("id"), "void").await?,
321 "charge.refunded" => self.refunded(object).await?,
322 "charge.dispute.created" => self.disputed(object, true).await?,
323 "charge.dispute.closed" => self.disputed(object, object["status"].as_str() == Some("lost")).await?,
324 "customer.updated"
325 | "payment_method.attached"
326 | "payment_method.detached"
327 | "payment_method.updated"
328 | "payment_method.automatically_updated"
329 | "setup_intent.succeeded" => match card_owner {
330 Some(customer) => self.sync_card_of(customer).await?,
331 None => "ignored: no customer".to_owned(),
332 },
333 _ => "ignored".to_owned(),
334 })
335 }
336
337 /// A payment page done, whether or not the person came back to g1t.
338 async fn settle_checkout(&self, session_id: &str) -> Result<String> {
339 #[derive(Deserialize)]
340 struct Open {
341 workspace: String,
342 created_by: String,
343 feature: Option<String>,
344 #[serde(default)]
345 fee_cents: Option<u32>,
346 }
347 let Some(open) = self
348 .db
349 .prepare("SELECT workspace, created_by, feature, fee_cents FROM checkouts WHERE id = ? AND status = 'open'")
350 .bind(&[session_id.into()])?
351 .first::<Open>(None)
352 .await?
353 else {
354 return Ok("ignored: already settled or not g1t's".to_owned());
355 };
356 // A card check: saved and verified, never charged.
357 if open.feature.as_deref() == Some(crate::cards::CARD_CHECK) {
358 return Ok(match self.settle_card_check(session_id).await? {
359 Ok(done) => done,
360 Err(why) => format!("card check not passed: {why}"),
361 });
362 }
363 // AI credit: credited once, by whichever of this and the person
364 // coming back claims the page first (ai.rs).
365 if open.feature.as_deref() == Some(crate::ai::AI_CREDIT) {
366 return Ok(match self.settle_ai_credit(session_id).await? {
367 Ok(done) => done,
368 Err(why) => format!("AI credit not credited: {why}"),
369 });
370 }
371 let Some(stripe) = &self.stripe else { return Ok("ignored: payments off".to_owned()) };
372 let session = stripe.session(session_id).await?;
373 if session.payment_status != "paid" && open.feature.is_none() {
374 return Ok("ignored: not paid".to_owned());
375 }
376 let claimed = self
377 .db
378 .prepare("UPDATE checkouts SET status = 'paid' WHERE id = ? AND status = 'open' RETURNING id")
379 .bind(&[session_id.into()])?
380 .first::<Value>(None)
381 .await?;
382 if claimed.is_none() {
383 return Ok("ignored: settled meanwhile".to_owned());
384 }
385 match open.feature.as_deref() {
386 None => {
387 let cents = self.credit_prepayment(&open.workspace, &session, open.fee_cents.unwrap_or(0), &open.created_by).await?;
388 Ok(format!("credited {} to {}", crate::features::dollars(cents * 10_000), open.workspace))
389 }
390 Some(feature) => {
391 let Some(feature) = g1t_contracts::billing::Feature::parse(feature) else {
392 return Ok("ignored: unknown feature".to_owned());
393 };
394 if let Some(subscription_id) = &session.subscription {
395 let subscription = stripe.subscription(subscription_id).await?;
396 self.record(&open.workspace, feature, &subscription, &open.created_by).await?;
397 }
398 self.db
399 .prepare(
400 "INSERT INTO accounts (workspace, balance_micros, customer_id, created_at) VALUES (?1, 0, ?2, ?3)
401 ON CONFLICT (workspace) DO UPDATE SET customer_id = COALESCE(customer_id, ?2)",
402 )
403 .bind(&[open.workspace.as_str().into(), crate::optional(session.customer.as_deref()), rfc3339(now_ms()).into()])?
404 .run()
405 .await?;
406 Ok(format!("{} plan started for {}", feature.title(), open.workspace))
407 }
408 }
409 }
410
411 /// The plan's monthly price, paid: revenue that never goes through the
412 /// ledger, recorded once per invoice for sudo's figures and for trust.
413 /// Only the plan's own price is revenue: the invoice's tax and card fee
414 /// are kept apart (`tax_and_fees`).
415 async fn plan_paid(&self, subscription_id: &str, invoice: &Value) -> Result<()> {
416 let invoice_id = invoice["id"].as_str().unwrap_or_default();
417 let split = crate::stripe::invoice_split(invoice, invoice["amount_paid"].as_i64().unwrap_or(0));
418 let amount_cents = split.net_cents;
419 if amount_cents <= 0 || invoice_id.is_empty() {
420 return Ok(());
421 }
422 #[derive(Deserialize)]
423 struct Owner {
424 workspace: String,
425 }
426 if let Some(owner) = self
427 .db
428 .prepare("SELECT workspace FROM subscriptions WHERE subscription_id = ?")
429 .bind(&[subscription_id.into()])?
430 .first::<Owner>(None)
431 .await?
432 {
433 let extras = crate::tax::Extras { tax_cents: split.tax_cents, fee_cents: split.fee_cents };
434 self.record_extras(&owner.workspace, invoice_id, invoice["payment_intent"].as_str(), extras, None).await?;
435 }
436 self.db
437 .prepare(
438 "INSERT INTO plan_payments (invoice_id, workspace, amount_micros, paid_at)
439 SELECT ?1, workspace, ?2, ?3 FROM subscriptions WHERE subscription_id = ?4
440 ON CONFLICT (invoice_id) DO NOTHING",
441 )
442 .bind(&[
443 invoice_id.into(),
444 ((amount_cents * 10_000) as f64).into(),
445 rfc3339(now_ms()).into(),
446 subscription_id.into(),
447 ])?
448 .run()
449 .await?;
450 Ok(())
451 }
452
453 /// A plan that changed at Stripe: renewed, failed, canceled.
454 pub(crate) async fn settle_subscription(&self, subscription_id: &str) -> Result<String> {
455 #[derive(Deserialize)]
456 struct Plan {
457 workspace: String,
458 feature: String,
459 started_by: String,
460 }
461 let Some(plan) = self
462 .db
463 .prepare("SELECT workspace, feature, started_by FROM subscriptions WHERE subscription_id = ?")
464 .bind(&[subscription_id.into()])?
465 .first::<Plan>(None)
466 .await?
467 else {
468 return Ok("ignored: not a g1t plan".to_owned());
469 };
470 let (Some(stripe), Some(feature)) = (&self.stripe, g1t_contracts::billing::Feature::parse(&plan.feature)) else {
471 return Ok("ignored".to_owned());
472 };
473 let subscription = stripe.subscription(subscription_id).await?;
474 self.record(&plan.workspace, feature, &subscription, &plan.started_by).await?;
475 Ok(format!("{} plan for {} is {}", feature.title(), plan.workspace, subscription.status))
476 }
477
478 /// The workspace a Stripe customer belongs to.
479 pub(crate) async fn workspace_of_customer(&self, customer: &str) -> Result<Option<String>> {
480 #[derive(Deserialize)]
481 struct Row {
482 workspace: String,
483 }
484 Ok(self
485 .db
486 .prepare("SELECT workspace FROM accounts WHERE customer_id = ?")
487 .bind(&[customer.into()])?
488 .first::<Row>(None)
489 .await?
490 .map(|row| row.workspace))
491 }
492
493 /// The workspace a workspace invoice was sent to, under its slug now.
494 async fn workspace_of_invoice(&self, invoice_id: &str) -> Result<Option<String>> {
495 #[derive(Deserialize)]
496 struct Row {
497 workspace: String,
498 }
499 Ok(self
500 .db
501 .prepare("SELECT workspace FROM workspace_invoices WHERE invoice_id = ?")
502 .bind(&[invoice_id.into()])?
503 .first::<Row>(None)
504 .await?
505 .map(|row| row.workspace))
506 }
507
508 /// Money given back: what was paid is less, by the refund.
509 async fn refunded(&self, charge: &Value) -> Result<String> {
510 let Some(customer) = charge["customer"].as_str() else { return Ok("ignored: no customer".to_owned()) };
511 let Some(workspace) = self.workspace_of_customer(customer).await? else {
512 return Ok("ignored: not a workspace's customer".to_owned());
513 };
514 let refunded = charge["amount_refunded"].as_i64().unwrap_or(0);
515 if refunded <= 0 {
516 return Ok("ignored: nothing refunded".to_owned());
517 }
518 // Each refund total once, so partial refunds add up correctly.
519 let charge_id = charge["id"].as_str().unwrap_or_default();
520 #[derive(Deserialize)]
521 struct Sum {
522 micros: Option<i64>,
523 }
524 let already = self
525 .db
526 .prepare("SELECT -SUM(amount_micros) AS micros FROM ledger WHERE reference LIKE ?")
527 .bind(&[format!("refund/{charge_id}/%").into()])?
528 .first::<Sum>(None)
529 .await?
530 .and_then(|s| s.micros)
531 .unwrap_or(0);
532 // A refund gives back the payment's tax and card fee in proportion:
533 // only the rest comes off the balance, which never held them.
534 let (extras, transaction) = self.extras_of(charge["payment_intent"].as_str(), charge["invoice"].as_str()).await?;
535 let paid = charge["amount"].as_i64().unwrap_or(refunded);
536 let (balance_cents, tax_cents, fee_cents) = crate::tax::refund_split(refunded, paid, extras.tax_cents, extras.fee_cents);
537 let (given_tax, given_fee) = self.refunded_extras(charge_id).await?;
538 let new = balance_cents * 10_000 - already;
539 let new_extras = crate::tax::Extras { tax_cents: -(tax_cents - given_tax), fee_cents: -(fee_cents - given_fee) };
540 if new <= 0 && new_extras == crate::tax::Extras::default() {
541 return Ok("ignored: refund already recorded".to_owned());
542 }
543 let reference = format!("refund/{charge_id}/{refunded}");
544 if new > 0 {
545 self.enter(&workspace, EntryKind::TopUp, -new, "Refunded to the card", &reference, None, None, None, None).await?;
546 }
547 self.record_extras(&workspace, &reference, charge["payment_intent"].as_str(), new_extras, None).await?;
548 // An off-session charge's tax was recorded by g1t (auto-reload):
549 // reverse what this refund gave back. Checkout's and invoices' tax
550 // Stripe Tax keeps itself.
551 if let (Some(stripe), Some(transaction)) = (&self.stripe, transaction.as_deref()) {
552 let given_back = new.max(0) / 10_000 - new_extras.tax_cents - new_extras.fee_cents;
553 if given_back > 0
554 && let Err(error) = stripe.reverse_tax(transaction, &reference, given_back).await
555 {
556 worker::console_error!("{workspace}: the tax on refund {reference} was not reversed: {error}");
557 }
558 }
559 Ok(format!("refund of {} recorded for {workspace}", crate::features::dollars(new.max(0))))
560 }
561
562 /// What earlier refunds of a charge gave back of its tax and card fee,
563 /// in cents, as positive amounts.
564 async fn refunded_extras(&self, charge_id: &str) -> Result<(i64, i64)> {
565 #[derive(Deserialize)]
566 struct Row {
567 kind: String,
568 micros: Option<i64>,
569 }
570 let rows = self
571 .db
572 .prepare("SELECT kind, SUM(amount_micros) AS micros FROM tax_and_fees WHERE reference LIKE ? GROUP BY kind")
573 .bind(&[format!("refund/{charge_id}/%").into()])?
574 .all()
575 .await?
576 .results::<Row>()?;
577 let of = |kind: &str| rows.iter().filter(|r| r.kind == kind).map(|r| -r.micros.unwrap_or(0) / 10_000).sum::<i64>();
578 Ok((of("tax"), of("card_fee")))
579 }
580
581 /// A disputed payment stops the workspace's work until it is resolved;
582 /// one closed in the workspace's favour lets it go on.
583 async fn disputed(&self, dispute: &Value, stop: bool) -> Result<String> {
584 let Some(stripe) = &self.stripe else { return Ok("ignored".to_owned()) };
585 let Some(charge_id) = dispute["charge"].as_str() else { return Ok("ignored: no charge".to_owned()) };
586 let charge: Value = stripe.get(&format!("/charges/{charge_id}")).await?;
587 let Some(workspace) = (match charge["customer"].as_str() {
588 Some(customer) => self.workspace_of_customer(customer).await?,
589 None => None,
590 }) else {
591 return Ok("ignored: not a workspace's customer".to_owned());
592 };
593 let now = rfc3339(now_ms());
594 // The disputed payment never counts toward trust again.
595 for reference in [charge["payment_intent"].as_str(), charge["invoice"].as_str()].into_iter().flatten() {
596 self.db
597 .prepare("UPDATE ledger SET disputed = ? WHERE workspace = ? AND reference = ?")
598 .bind(&[(if stop { 1 } else { 0 }).into(), workspace.as_str().into(), reference.into()])?
599 .run()
600 .await?;
601 }
602 if stop {
603 self.db
604 .prepare(
605 "INSERT INTO limits (workspace, autopay_failed_at, autopay_error, updated_at) VALUES (?1, ?2, ?3, ?2)
606 ON CONFLICT (workspace) DO UPDATE SET autopay_failed_at = ?2, autopay_error = ?3, updated_at = ?2",
607 )
608 .bind(&[workspace.as_str().into(), now.as_str().into(), "a payment was disputed with the card's bank".into()])?
609 .run()
610 .await?;
611 } else {
612 self.db
613 .prepare("UPDATE limits SET autopay_failed_at = NULL, autopay_error = NULL WHERE workspace = ?")
614 .bind(&[workspace.as_str().into()])?
615 .run()
616 .await?;
617 }
618 let account = self.account_of(&workspace).await?;
619 let what = if stop { "dispute: work stopped" } else { "dispute closed in the workspace's favour" };
620 self.audit(&account.id, "dispute", &format!("{workspace}: {what}"), "stripe").await?;
621 Ok(format!("{workspace}: {what}"))
622 }
623
624 // --- Enterprise invoices ------------------------------------------------
625
626 pub(crate) async fn admin_enterprise_billing(&self, a: AdminEnterpriseBillingArgs) -> Result<Outcome<BillingAccount>> {
627 let email = a.email.trim().to_lowercase();
628 if !email.contains('@') || email.contains(char::is_whitespace) || a.by.trim().is_empty() {
629 return Ok(Outcome::fail(FailureCode::Invalid, "Give the email the enterprise's invoices go to."));
630 }
631 let Some(stripe) = &self.stripe else {
632 return Ok(Outcome::fail(FailureCode::Conflict, "Payments are not set up on this g1t."));
633 };
634 #[derive(Deserialize)]
635 struct Row {
636 name: String,
637 kind: String,
638 customer_id: Option<String>,
639 }
640 let Some(row) = self
641 .db
642 .prepare("SELECT name, kind, customer_id FROM billing_accounts WHERE id = ?")
643 .bind(&[a.id.as_str().into()])?
644 .first::<Row>(None)
645 .await?
646 .filter(|row| row.kind == "enterprise")
647 else {
648 return Ok(Outcome::fail(FailureCode::NotFound, "No such enterprise."));
649 };
650 #[derive(Deserialize)]
651 struct Customer {
652 id: String,
653 }
654 let fields = [
655 ("name", row.name.clone()),
656 ("email", email.clone()),
657 ("metadata[g1t_enterprise]", a.id.clone()),
658 ];
659 let customer: Customer = match &row.customer_id {
660 Some(id) => stripe.post(&format!("/customers/{id}"), &fields).await?,
661 None => stripe.post("/customers", &fields).await?,
662 };
663 self.db
664 .prepare("UPDATE billing_accounts SET billing_email = ?, customer_id = ? WHERE id = ?")
665 .bind(&[email.as_str().into(), customer.id.as_str().into(), a.id.as_str().into()])?
666 .run()
667 .await?;
668 self.audit(&a.id, "billing_email", &format!("Invoices go to {email}"), &a.by).await?;
669 Ok(match self.enterprise(&a.id).await? {
670 Some(account) => Outcome::Ok(account),
671 None => Outcome::fail(FailureCode::NotFound, "No such enterprise."),
672 })
673 }
674
675 /// `admin_enterprise_address`: the address Stripe Tax works the
676 /// enterprise's invoices out from, and its tax ID, on its customer.
677 pub(crate) async fn admin_enterprise_address(&self, a: AdminEnterpriseAddressArgs) -> Result<Outcome<bool>> {
678 if a.by.trim().is_empty() {
679 return Ok(Outcome::fail(FailureCode::Invalid, "Say who is changing it."));
680 }
681 if let Some(why) = crate::details::address_invalid(&a.address).or_else(|| crate::details::tax_id_invalid(a.tax_id_type.as_deref(), a.tax_id.as_deref())) {
682 return Ok(Outcome::fail(FailureCode::Invalid, why));
683 }
684 let place = serde_json::json!({ "country": a.address.country.trim(), "postal_code": a.address.postal_code.trim(), "state": a.address.state.trim() });
685 if !crate::stripe::address_places_customer(&place) {
686 return Ok(Outcome::fail(FailureCode::Invalid, "Give at least the country, and in the US the ZIP code: Stripe Tax needs them."));
687 }
688 let Some(stripe) = &self.stripe else {
689 return Ok(Outcome::fail(FailureCode::Conflict, "Payments are not set up on this g1t."));
690 };
691 #[derive(Deserialize)]
692 struct Row {
693 kind: String,
694 customer_id: Option<String>,
695 }
696 let Some(row) = self
697 .db
698 .prepare("SELECT kind, customer_id FROM billing_accounts WHERE id = ?")
699 .bind(&[a.id.as_str().into()])?
700 .first::<Row>(None)
701 .await?
702 .filter(|row| row.kind == "enterprise")
703 else {
704 return Ok(Outcome::fail(FailureCode::NotFound, "No such enterprise."));
705 };
706 let Some(customer) = row.customer_id else {
707 return Ok(Outcome::fail(FailureCode::Conflict, "Set where its invoices go first: that makes its Stripe customer."));
708 };
709 let fields: Vec<(&str, String)> = vec![
710 ("address[line1]", a.address.line1.trim().to_owned()),
711 ("address[line2]", a.address.line2.trim().to_owned()),
712 ("address[city]", a.address.city.trim().to_owned()),
713 ("address[state]", a.address.state.trim().to_owned()),
714 ("address[postal_code]", a.address.postal_code.trim().to_owned()),
715 ("address[country]", a.address.country.trim().to_uppercase()),
716 ];
717 if let Err(error) = stripe.post::<Value>(&format!("/customers/{customer}"), &fields).await {
718 return Ok(Outcome::fail(FailureCode::Conflict, crate::stripe::friendly(&error)));
719 }
720 if let (Some(kind), Some(value)) = (a.tax_id_type.as_deref(), a.tax_id.as_deref())
721 && let Err(why) = crate::details::replace_tax_id(stripe, &customer, &a.id, kind, value).await
722 {
723 return Ok(Outcome::fail(FailureCode::Invalid, why));
724 }
725 self.audit(&a.id, "billing_address", &format!("Billing address set ({})", a.address.country.trim().to_uppercase()), &a.by).await?;
726 Ok(Outcome::Ok(true))
727 }
728
729 pub(crate) async fn admin_invoice_enterprise(&self, a: AdminInvoiceEnterpriseArgs) -> Result<Outcome<EnterpriseInvoice>> {
730 if a.by.trim().is_empty() {
731 return Ok(Outcome::fail(FailureCode::Invalid, "Say who is sending it."));
732 }
733 match self.invoice_enterprise(&a.id, "now", &a.by).await? {
734 Ok(invoice) => Ok(Outcome::Ok(invoice)),
735 Err(why) => Ok(Outcome::fail(FailureCode::Conflict, why)),
736 }
737 }
738
739 /// Invoices each enterprise for the month that closed. Live payments
740 /// only; sudo can send one sooner in test mode.
741 pub(crate) async fn invoice_enterprises(&self) -> Result<()> {
742 if !self.stripe.as_ref().is_some_and(crate::stripe::Stripe::live) {
743 return Ok(());
744 }
745 let closing = crate::limits::previous_month(&rfc3339(now_ms())[..7]);
746 #[derive(Deserialize)]
747 struct Id {
748 id: String,
749 }
750 let due = self
751 .db
752 .prepare(format!(
753 "SELECT id FROM billing_accounts WHERE kind = 'enterprise' AND customer_id IS NOT NULL
754 AND NOT {full}
755 AND NOT EXISTS (SELECT 1 FROM enterprise_invoices i WHERE i.account_id = billing_accounts.id AND i.period = ?)",
756 full = crate::sales::FULL_DISCOUNT_SQL
757 ))
758 .bind(&[closing.as_str().into()])?
759 .all()
760 .await?
761 .results::<Id>()?;
762 for Id { id } in due {
763 if let Err(why) = self.invoice_enterprise(&id, &closing, "month close").await? {
764 worker::console_log!("enterprise {id} not invoiced for {closing}: {why}");
765 }
766 }
767 Ok(())
768 }
769
770 /// One invoice for what each of the enterprise's workspaces owes now.
771 async fn invoice_enterprise(&self, id: &str, period: &str, by: &str) -> Result<std::result::Result<EnterpriseInvoice, String>> {
772 let Some(stripe) = &self.stripe else { return Ok(Err("Payments are not set up.".into())) };
773 let Some(account) = self.enterprise(id).await? else { return Ok(Err("No such enterprise.".into())) };
774 #[derive(Deserialize)]
775 struct Customer {
776 customer_id: Option<String>,
777 }
778 let Some(customer) = self
779 .db
780 .prepare("SELECT customer_id FROM billing_accounts WHERE id = ?")
781 .bind(&[id.into()])?
782 .first::<Customer>(None)
783 .await?
784 .and_then(|row| row.customer_id)
785 else {
786 return Ok(Err("Set where the enterprise's invoices go first.".into()));
787 };
788 // What each workspace owes: its charges less what it has paid, and
789 // less what is on invoices still open.
790 let mut lines = vec![];
791 for workspace in &account.workspaces {
792 let balance = self.row(workspace).await?.map_or(0, |row| row.balance_micros);
793 #[derive(Deserialize)]
794 struct Sum {
795 micros: Option<i64>,
796 }
797 let invoiced = self
798 .db
799 .prepare(
800 "SELECT SUM(l.amount_micros) AS micros FROM enterprise_invoice_lines l
801 JOIN enterprise_invoices i ON i.invoice_id = l.invoice_id
802 WHERE l.workspace = ? AND i.status IN ('open', 'overdue')",
803 )
804 .bind(&[workspace.as_str().into()])?
805 .first::<Sum>(None)
806 .await?
807 .and_then(|s| s.micros)
808 .unwrap_or(0);
809 let owed = (-balance).max(0) - invoiced;
810 if owed >= 10_000 {
811 lines.push(InvoiceLine { workspace: workspace.clone(), amount_micros: owed });
812 }
813 }
814 if lines.is_empty() {
815 return Ok(Err("Its workspaces owe nothing to invoice.".into()));
816 }
817 // Stripe Tax places the customer by its address; without one the
818 // invoice could not be finalized, so it is not started.
819 match stripe.customer_placed(&customer).await {
820 Ok(true) => {}
821 Ok(false) => {
822 return Ok(Err(format!(
823 "Add {}'s billing address first (Invoices → Billing address): Stripe needs it to work out tax.",
824 account.name
825 )));
826 }
827 Err(error) => return Ok(Err(crate::stripe::friendly(&error))),
828 }
829 // Invoiced, never by card: no card fee. Tax on top, by Stripe Tax.
830 let mut fields = vec![
831 ("customer", customer.clone()),
832 ("collection_method", "send_invoice".to_owned()),
833 ("days_until_due", "30".to_owned()),
834 ("pending_invoice_items_behavior", "exclude".to_owned()),
835 ("description", format!("g1t usage for the {} enterprise", account.name)),
836 ("metadata[g1t_enterprise]", id.to_owned()),
837 ("metadata[period]", period.to_owned()),
838 ];
839 fields.extend(crate::stripe::invoice_tax_fields());
840 #[derive(Deserialize)]
841 struct Invoice {
842 id: String,
843 #[serde(default)]
844 hosted_invoice_url: Option<String>,
845 #[serde(default)]
846 amount_due: i64,
847 /// The tax Stripe added, in cents (`tax`, in this API version).
848 #[serde(default)]
849 tax: Option<i64>,
850 }
851 // The draft first, keyed on what it bills, and its lines put on it:
852 // an attempt that failed half way is found again, never billed again
853 // by an invoice sweeping its pending lines in.
854 let cents: Vec<i64> = lines.iter().map(|line| (line.amount_micros + 9_999) / 10_000).collect();
855 let keyed: Vec<(&str, i64)> = lines.iter().zip(&cents).map(|(line, cents)| (line.workspace.as_str(), *cents)).collect();
856 let key = enterprise_invoice_key(id, period, &keyed);
857 let draft: Invoice = stripe.post_idempotent("/invoices", &fields, &key).await?;
858 for (line, cents) in lines.iter().zip(&cents) {
859 let mut fields = vec![
860 ("customer", customer.clone()),
861 ("invoice", draft.id.clone()),
862 ("amount", cents.to_string()),
863 ("currency", "usd".to_owned()),
864 ("description", format!("{}: g1t usage", line.workspace)),
865 ("metadata[workspace]", line.workspace.clone()),
866 ];
867 fields.extend(crate::stripe::item_tax_fields());
868 let _: Value = stripe.post_idempotent("/invoiceitems", &fields, &format!("{key}/item/{}", line.workspace)).await?;
869 }
870 // A retry finds it finalized already; that is fine.
871 let _ = stripe.post::<Value>(&format!("/invoices/{}/finalize", draft.id), &[]).await;
872 let sent: Invoice = match stripe.post(&format!("/invoices/{}/send", draft.id), &[]).await {
873 Ok(sent) => sent,
874 Err(error) if crate::stripe::is_tax_location_error(&error) => {
875 return Ok(Err(format!("Stripe could not work out tax for {}: add its billing address, then send it again.", account.name)));
876 }
877 Err(error) => return Err(error),
878 };
879 let now = rfc3339(now_ms());
880 // What the workspaces owe, before tax: the tax is the enterprise's
881 // to pay on top, never their usage, and is kept apart.
882 let tax_cents = sent.tax.unwrap_or(0).max(0);
883 let total = lines.iter().map(|l| l.amount_micros).sum::<i64>().max((sent.amount_due - tax_cents) * 10_000);
884 let mut writes = vec![self
885 .db
886 .prepare(
887 "INSERT INTO enterprise_invoices (invoice_id, account_id, period, amount_micros, status, hosted_url, created_by, created_at)
888 VALUES (?, ?, ?, ?, 'open', ?, ?, ?)",
889 )
890 .bind(&[
891 sent.id.as_str().into(),
892 id.into(),
893 period.into(),
894 (total as f64).into(),
895 crate::optional(sent.hosted_invoice_url.as_deref()),
896 by.into(),
897 now.as_str().into(),
898 ])?];
899 for line in &lines {
900 writes.push(
901 self.db
902 .prepare("INSERT INTO enterprise_invoice_lines (invoice_id, workspace, amount_micros) VALUES (?, ?, ?)")
903 .bind(&[sent.id.as_str().into(), line.workspace.as_str().into(), (line.amount_micros as f64).into()])?,
904 );
905 }
906 self.db.batch(writes).await?;
907 self.audit(id, "invoice", &format!("Invoice {} for {} sent ({period})", sent.id, crate::features::dollars(total)), by)
908 .await?;
909 Ok(Ok(EnterpriseInvoice {
910 invoice_id: sent.id,
911 hosted_url: sent.hosted_invoice_url,
912 amount_micros: total,
913 status: "open".to_owned(),
914 period: period.to_owned(),
915 lines,
916 created_at: now,
917 }))
918 }
919
920 /// An enterprise invoice paid: each workspace is credited its line, and
921 /// any stop for the invoice is lifted.
922 async fn enterprise_invoice_paid(&self, invoice: &Value) -> Result<String> {
923 let invoice_id = invoice["id"].as_str().unwrap_or_default();
924 let claimed = self
925 .db
926 .prepare(
927 "UPDATE enterprise_invoices SET status = 'paid', paid_at = ? WHERE invoice_id = ? AND status <> 'paid'
928 RETURNING invoice_id",
929 )
930 .bind(&[rfc3339(now_ms()).into(), invoice_id.into()])?
931 .first::<Value>(None)
932 .await?;
933 if claimed.is_none() {
934 return Ok("ignored: not an open enterprise invoice".to_owned());
935 }
936 let lines = self.invoice_lines(invoice_id).await?;
937 for line in &lines {
938 self.enter(
939 &line.workspace,
940 EntryKind::TopUp,
941 line.amount_micros,
942 &format!("Paid on the enterprise's invoice {invoice_id}"),
943 &format!("inv/{invoice_id}/{}", line.workspace),
944 None,
945 None,
946 None,
947 None,
948 )
949 .await?;
950 }
951 // Its tax is the enterprise's, paid on top of what its workspaces
952 // owed, and never credited to them.
953 #[derive(Deserialize)]
954 struct Account {
955 account_id: String,
956 }
957 if let Some(account) = self
958 .db
959 .prepare("SELECT account_id FROM enterprise_invoices WHERE invoice_id = ?")
960 .bind(&[invoice_id.into()])?
961 .first::<Account>(None)
962 .await?
963 {
964 let extras = crate::tax::Extras { tax_cents: invoice["tax"].as_i64().unwrap_or(0).max(0), fee_cents: 0 };
965 self.record_extras(&account.account_id, invoice_id, invoice["payment_intent"].as_str(), extras, None).await?;
966 }
967 Ok(format!("invoice {invoice_id} paid; {} workspaces credited", lines.len()))
968 }
969
970 /// An enterprise invoice that went overdue stops its workspaces' work;
971 /// one voided is simply closed.
972 async fn enterprise_invoice_status(&self, invoice_id: &str, status: &str) -> Result<String> {
973 let updated = self
974 .db
975 .prepare("UPDATE enterprise_invoices SET status = ? WHERE invoice_id = ? AND status <> 'paid' RETURNING invoice_id")
976 .bind(&[status.into(), invoice_id.into()])?
977 .first::<Value>(None)
978 .await?;
979 if updated.is_none() {
980 return Ok("ignored: not an open enterprise invoice".to_owned());
981 }
982 if status == "overdue" {
983 let now = rfc3339(now_ms());
984 for line in self.invoice_lines(invoice_id).await? {
985 self.db
986 .prepare(
987 "INSERT INTO limits (workspace, autopay_failed_at, autopay_error, updated_at) VALUES (?1, ?2, ?3, ?2)
988 ON CONFLICT (workspace) DO UPDATE SET autopay_failed_at = ?2, autopay_error = ?3, updated_at = ?2",
989 )
990 .bind(&[line.workspace.as_str().into(), now.as_str().into(), format!("the enterprise's invoice {invoice_id} is overdue").into()])?
991 .run()
992 .await?;
993 }
994 }
995 Ok(format!("invoice {invoice_id} is {status}"))
996 }
997
998 async fn invoice_lines(&self, invoice_id: &str) -> Result<Vec<InvoiceLine>> {
999 Ok(self
1000 .db
1001 .prepare("SELECT workspace, amount_micros FROM enterprise_invoice_lines WHERE invoice_id = ?")
1002 .bind(&[invoice_id.into()])?
1003 .all()
1004 .await?
1005 .results::<LineRow>()?
1006 .into_iter()
1007 .map(|row| InvoiceLine { workspace: row.workspace, amount_micros: row.amount_micros })
1008 .collect())
1009 }
1010
1011 /// An enterprise's invoices, newest first.
1012 pub(crate) async fn enterprise_invoices(&self, id: &str) -> Result<Vec<EnterpriseInvoice>> {
1013 let rows = self
1014 .db
1015 .prepare("SELECT * FROM enterprise_invoices WHERE account_id = ? ORDER BY created_at DESC LIMIT 24")
1016 .bind(&[JsValue::from(id)])?
1017 .all()
1018 .await?
1019 .results::<InvoiceRow>()?;
1020 let mut invoices = vec![];
1021 for row in rows {
1022 invoices.push(EnterpriseInvoice {
1023 lines: self.invoice_lines(&row.invoice_id).await?,
1024 invoice_id: row.invoice_id,
1025 hosted_url: row.hosted_url,
1026 amount_micros: row.amount_micros,
1027 status: row.status,
1028 period: row.period,
1029 created_at: row.created_at,
1030 });
1031 }
1032 Ok(invoices)
1033 }
1034}
1035
1036#[cfg(test)]
1037mod tests {
1038 use super::*;
1039
1040 #[test]
1041 fn an_enterprise_invoice_attempted_again_is_the_same_invoice() {
1042 let lines = [("acme", 1_250), ("acme-labs", 99)];
1043 let key = enterprise_invoice_key("ent_1", "2026-09", &lines);
1044 assert_eq!(key, enterprise_invoice_key("ent_1", "2026-09", &lines));
1045 assert!(key.len() < 255 && key.starts_with("ent-invoice/ent_1/2026-09/"));
1046 // Different amounts, even with the same total, are a different invoice.
1047 assert_ne!(key, enterprise_invoice_key("ent_1", "2026-09", &[("acme", 1_300), ("acme-labs", 49)]));
1048 assert_ne!(key, enterprise_invoice_key("ent_1", "2026-10", &lines));
1049 }
1050
1051 fn sign(payload: &str, secret: &str, t: i64) -> String {
1052 let mut mac = Hmac::<Sha256>::new_from_slice(secret.as_bytes()).unwrap();
1053 mac.update(format!("{t}.{payload}").as_bytes());
1054 format!("t={t},v1={}", hex::encode(mac.finalize().into_bytes()))
1055 }
1056
1057 #[test]
1058 fn a_destination_misses_the_events_billing_handles_that_it_does_not_send() {
1059 let has: Vec<String> = EVENTS.iter().take(11).map(|e| (*e).to_owned()).collect();
1060 assert_eq!(missing_events(&has), EVENTS[11..].iter().map(|e| (*e).to_owned()).collect::<Vec<_>>());
1061 let all: Vec<String> = EVENTS.iter().map(|e| (*e).to_owned()).collect();
1062 assert!(missing_events(&all).is_empty());
1063 assert!(missing_events(&["*".to_owned()]).is_empty());
1064 }
1065
1066 #[test]
1067 fn a_signed_event_is_believed_only_as_signed_and_only_fresh() {
1068 let payload = r#"{"id":"evt_1","type":"invoice.paid"}"#;
1069 let header = sign(payload, "whsec_test", 1_000_000);
1070 assert!(verify(payload, &header, "whsec_test", 1_000_010));
1071 assert!(!verify(payload, &header, "whsec_other", 1_000_010));
1072 assert!(!verify(&payload.replace("paid", "voided"), &header, "whsec_test", 1_000_010));
1073 assert!(!verify(payload, &header, "whsec_test", 1_000_000 + 301));
1074 assert!(!verify(payload, "v1=abc", "whsec_test", 1_000_000));
1075 // Stripe may sign with more than one secret while one is rolled.
1076 let both = format!("{},v1=00ff", sign(payload, "whsec_test", 1_000_000));
1077 assert!(verify(payload, &both, "whsec_test", 1_000_000));
1078 }
1079}