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/webhooks.rs

842 lines35,841 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 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 "customer.subscription.updated",
47 "customer.subscription.deleted",
48 "invoice.paid",
49 "invoice.payment_failed",
50 "invoice.overdue",
51 "invoice.voided",
52 "charge.refunded",
53 "charge.dispute.created",
54 "charge.dispute.closed",
55];
56
57/// How old a signed event may be, so a captured one cannot be replayed.
58const TOLERANCE_SECONDS: i64 = 5 * 60;
59
60/// Whether `header` (`t=…,v1=…`) signs `payload` with `secret`, within the
61/// tolerance of `now_seconds`.
62pub(crate) fn verify(payload: &str, header: &str, secret: &str, now_seconds: i64) -> bool {
63 let mut timestamp = None;
64 let mut signatures = vec![];
65 for part in header.split(',') {
66 match part.trim().split_once('=') {
67 Some(("t", value)) => timestamp = value.parse::<i64>().ok(),
68 Some(("v1", value)) => signatures.push(value.to_owned()),
69 _ => {}
70 }
71 }
72 let Some(timestamp) = timestamp else { return false };
73 if (now_seconds - timestamp).abs() > TOLERANCE_SECONDS {
74 return false;
75 }
76 let Ok(mut mac) = Hmac::<Sha256>::new_from_slice(secret.as_bytes()) else { return false };
77 mac.update(format!("{timestamp}.{payload}").as_bytes());
78 let expected = mac.finalize().into_bytes();
79 signatures.iter().any(|signature| {
80 hex::decode(signature).is_ok_and(|given| {
81 // Constant time: compare every byte whatever the first difference.
82 given.len() == expected.len() && given.iter().zip(expected.iter()).fold(0u8, |acc, (a, b)| acc | (a ^ b)) == 0
83 })
84 })
85}
86
87#[derive(Deserialize)]
88struct WebhookRow {
89 endpoint_id: String,
90 secret: String,
91 url: String,
92 events: String,
93 created_by: String,
94 created_at: String,
95}
96
97#[derive(Deserialize)]
98struct EventRow {
99 id: String,
100 r#type: String,
101 outcome: String,
102 received_at: String,
103}
104
105#[derive(Deserialize)]
106struct InvoiceRow {
107 invoice_id: String,
108 period: String,
109 amount_micros: i64,
110 status: String,
111 hosted_url: Option<String>,
112 created_at: String,
113}
114
115#[derive(Deserialize)]
116struct LineRow {
117 workspace: String,
118 amount_micros: i64,
119}
120
121impl Billing {
122 fn mode(&self) -> &'static str {
123 match &self.stripe {
124 None => "off",
125 Some(stripe) if stripe.live() => "live",
126 Some(_) => "test",
127 }
128 }
129
130 async fn webhook_row(&self) -> Result<Option<WebhookRow>> {
131 self.db
132 .prepare("SELECT * FROM stripe_webhooks WHERE mode = ?")
133 .bind(&[self.mode().into()])?
134 .first::<WebhookRow>(None)
135 .await
136 }
137
138 // --- Staff ------------------------------------------------------------
139
140 pub(crate) async fn admin_stripe(&self, a: AdminStripeArgs) -> Result<StripeStatus> {
141 let mut error = None;
142 if a.setup {
143 if let Err(e) = self.register_webhook(a.by.as_deref().unwrap_or("sudo")).await {
144 error = Some(e.to_string());
145 }
146 }
147 let webhook = self.webhook_row().await?.map(|row| StripeWebhook {
148 url: row.url,
149 endpoint_id: row.endpoint_id,
150 events: row.events.split(',').map(str::to_owned).collect(),
151 created_by: row.created_by,
152 created_at: row.created_at,
153 });
154 let recent_events = self
155 .db
156 .prepare("SELECT * FROM stripe_events ORDER BY received_at DESC LIMIT 25")
157 .all()
158 .await?
159 .results::<EventRow>()?
160 .into_iter()
161 .map(|row| StripeEventSummary { id: row.id, kind: row.r#type, outcome: row.outcome, received_at: row.received_at })
162 .collect();
163 Ok(StripeStatus { mode: self.mode().to_owned(), webhook, recent_events, error })
164 }
165
166 /// Registers billing's endpoint at Stripe for the current mode,
167 /// replacing any it made before, and keeps the new signing secret.
168 async fn register_webhook(&self, by: &str) -> Result<()> {
169 let Some(stripe) = &self.stripe else {
170 return Err(worker::Error::RustError("payments are not set up".into()));
171 };
172 #[derive(Deserialize)]
173 struct Endpoint {
174 id: String,
175 url: String,
176 #[serde(default)]
177 secret: Option<String>,
178 }
179 #[derive(Deserialize)]
180 struct List {
181 data: Vec<Endpoint>,
182 }
183 // Ours from before, whose secret cannot be read again: replaced.
184 let existing: List = stripe.get("/webhook_endpoints?limit=100").await?;
185 for endpoint in existing.data.iter().filter(|e| e.url == WEBHOOK_URL) {
186 let _: Value = stripe.delete(&format!("/webhook_endpoints/{}", endpoint.id)).await?;
187 }
188 let mut fields = vec![
189 ("url", WEBHOOK_URL.to_owned()),
190 ("description", "g1t billing".to_owned()),
191 ("metadata[g1t]", "billing".to_owned()),
192 ];
193 let names: Vec<String> = (0..EVENTS.len()).map(|i| format!("enabled_events[{i}]")).collect();
194 for (name, event) in names.iter().zip(EVENTS) {
195 fields.push((name.as_str(), (*event).to_owned()));
196 }
197 let created: Endpoint = stripe.post("/webhook_endpoints", &fields).await?;
198 let Some(secret) = created.secret else {
199 return Err(worker::Error::RustError("Stripe returned no signing secret".into()));
200 };
201 self.db
202 .prepare(
203 "INSERT INTO stripe_webhooks (mode, endpoint_id, secret, url, events, created_by, created_at)
204 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
205 ON CONFLICT (mode) DO UPDATE SET endpoint_id = ?2, secret = ?3, url = ?4, events = ?5,
206 created_by = ?6, created_at = ?7",
207 )
208 .bind(&[
209 self.mode().into(),
210 created.id.as_str().into(),
211 secret.as_str().into(),
212 created.url.as_str().into(),
213 EVENTS.join(",").into(),
214 by.into(),
215 rfc3339(now_ms()).into(),
216 ])?
217 .run()
218 .await?;
219 self.audit("stripe", "webhook", &format!("Registered {WEBHOOK_URL} ({} mode)", self.mode()), by).await?;
220 Ok(())
221 }
222
223 // --- Events -----------------------------------------------------------
224
225 pub(crate) async fn stripe_webhook(&self, a: StripeWebhookArgs) -> Result<Outcome<bool>> {
226 let Some(webhook) = self.webhook_row().await? else {
227 return Ok(Outcome::fail(FailureCode::Conflict, "No webhook is registered for this mode."));
228 };
229 let now_seconds = (now_ms() / 1000) as i64;
230 if !verify(&a.payload, &a.signature, &webhook.secret, now_seconds) {
231 return Ok(Outcome::fail(FailureCode::Forbidden, "The signature does not match."));
232 }
233 let event: Value = serde_json::from_str(&a.payload).map_err(|e| worker::Error::RustError(e.to_string()))?;
234 let id = event["id"].as_str().unwrap_or_default().to_owned();
235 let kind = event["type"].as_str().unwrap_or_default().to_owned();
236 if id.is_empty() {
237 return Ok(Outcome::fail(FailureCode::Invalid, "Not an event."));
238 }
239 // Once each: the first to record it handles it.
240 let claimed = self
241 .db
242 .prepare("INSERT OR IGNORE INTO stripe_events (id, type, outcome, received_at) VALUES (?, ?, 'handling', ?) RETURNING id")
243 .bind(&[id.as_str().into(), kind.as_str().into(), rfc3339(now_ms()).into()])?
244 .first::<Value>(None)
245 .await?;
246 if claimed.is_none() {
247 return Ok(Outcome::Ok(false));
248 }
249 let object = &event["data"]["object"];
250 let outcome = match self.handle(&kind, object).await {
251 Ok(outcome) => outcome,
252 Err(error) => {
253 // Let Stripe send it again: forget it was seen.
254 self.db.prepare("DELETE FROM stripe_events WHERE id = ?").bind(&[id.as_str().into()])?.run().await?;
255 return Err(error);
256 }
257 };
258 self.db
259 .prepare("UPDATE stripe_events SET outcome = ? WHERE id = ?")
260 .bind(&[outcome.as_str().into(), id.as_str().into()])?
261 .run()
262 .await?;
263 Ok(Outcome::Ok(true))
264 }
265
266 async fn handle(&self, kind: &str, object: &Value) -> Result<String> {
267 let text = |key: &str| object[key].as_str().unwrap_or_default().to_owned();
268 Ok(match kind {
269 "checkout.session.completed" => self.settle_checkout(&text("id")).await?,
270 "customer.subscription.updated" | "customer.subscription.deleted" => {
271 self.settle_subscription(&text("id")).await?
272 }
273 "invoice.paid" => {
274 if let Some(subscription) = object["subscription"].as_str() {
275 self.settle_subscription(subscription).await?
276 } else if let Some(done) = self.workspace_invoice_paid(&text("id")).await? {
277 done
278 } else {
279 self.enterprise_invoice_paid(&text("id")).await?
280 }
281 }
282 "invoice.payment_failed" => match (object["subscription"].as_str(), object["metadata"]["g1t_workspace"].as_str()) {
283 (Some(subscription), _) => self.settle_subscription(subscription).await?,
284 (None, Some(tagged)) => {
285 // The invoice's own row names the workspace as it is
286 // now; the metadata keeps the slug it was sent under.
287 let workspace = self.workspace_of_invoice(&text("id")).await?.unwrap_or_else(|| tagged.to_owned());
288 let workspace = workspace.as_str();
289 self.mark_declined(workspace, "the card was declined for an invoice").await?;
290 format!("{workspace}: invoice payment failed; work stopped")
291 }
292 _ => "ignored: not g1t's".to_owned(),
293 },
294 "invoice.overdue" => self.enterprise_invoice_status(&text("id"), "overdue").await?,
295 "invoice.voided" => self.enterprise_invoice_status(&text("id"), "void").await?,
296 "charge.refunded" => self.refunded(object).await?,
297 "charge.dispute.created" => self.disputed(object, true).await?,
298 "charge.dispute.closed" => self.disputed(object, object["status"].as_str() == Some("lost")).await?,
299 _ => "ignored".to_owned(),
300 })
301 }
302
303 /// A payment page done, whether or not the person came back to g1t.
304 async fn settle_checkout(&self, session_id: &str) -> Result<String> {
305 #[derive(Deserialize)]
306 struct Open {
307 workspace: String,
308 created_by: String,
309 feature: Option<String>,
310 }
311 let Some(open) = self
312 .db
313 .prepare("SELECT workspace, created_by, feature FROM checkouts WHERE id = ? AND status = 'open'")
314 .bind(&[session_id.into()])?
315 .first::<Open>(None)
316 .await?
317 else {
318 return Ok("ignored: already settled or not g1t's".to_owned());
319 };
320 let Some(stripe) = &self.stripe else { return Ok("ignored: payments off".to_owned()) };
321 let session = stripe.session(session_id).await?;
322 if session.payment_status != "paid" && open.feature.is_none() {
323 return Ok("ignored: not paid".to_owned());
324 }
325 let claimed = self
326 .db
327 .prepare("UPDATE checkouts SET status = 'paid' WHERE id = ? AND status = 'open' RETURNING id")
328 .bind(&[session_id.into()])?
329 .first::<Value>(None)
330 .await?;
331 if claimed.is_none() {
332 return Ok("ignored: settled meanwhile".to_owned());
333 }
334 match open.feature.as_deref() {
335 None => {
336 let cents = i64::from(session.amount_total.unwrap_or(0));
337 self.enter(
338 &open.workspace,
339 EntryKind::TopUp,
340 cents * 10_000,
341 "Credit added by card",
342 &session.id,
343 None,
344 None,
345 Some(&open.created_by),
346 session.customer.as_deref(),
347 )
348 .await?;
349 Ok(format!("credited {} to {}", crate::features::dollars(cents * 10_000), open.workspace))
350 }
351 Some(feature) => {
352 let Some(feature) = g1t_contracts::billing::Feature::parse(feature) else {
353 return Ok("ignored: unknown feature".to_owned());
354 };
355 if let Some(subscription_id) = &session.subscription {
356 let subscription = stripe.subscription(subscription_id).await?;
357 self.record(&open.workspace, feature, &subscription, &open.created_by).await?;
358 }
359 self.db
360 .prepare(
361 "INSERT INTO accounts (workspace, balance_micros, customer_id, created_at) VALUES (?1, 0, ?2, ?3)
362 ON CONFLICT (workspace) DO UPDATE SET customer_id = COALESCE(customer_id, ?2)",
363 )
364 .bind(&[open.workspace.as_str().into(), crate::optional(session.customer.as_deref()), rfc3339(now_ms()).into()])?
365 .run()
366 .await?;
367 Ok(format!("{} plan started for {}", feature.title(), open.workspace))
368 }
369 }
370 }
371
372 /// A plan that changed at Stripe: renewed, failed, canceled.
373 async fn settle_subscription(&self, subscription_id: &str) -> Result<String> {
374 #[derive(Deserialize)]
375 struct Plan {
376 workspace: String,
377 feature: String,
378 started_by: String,
379 }
380 let Some(plan) = self
381 .db
382 .prepare("SELECT workspace, feature, started_by FROM subscriptions WHERE subscription_id = ?")
383 .bind(&[subscription_id.into()])?
384 .first::<Plan>(None)
385 .await?
386 else {
387 return Ok("ignored: not a g1t plan".to_owned());
388 };
389 let (Some(stripe), Some(feature)) = (&self.stripe, g1t_contracts::billing::Feature::parse(&plan.feature)) else {
390 return Ok("ignored".to_owned());
391 };
392 let subscription = stripe.subscription(subscription_id).await?;
393 self.record(&plan.workspace, feature, &subscription, &plan.started_by).await?;
394 Ok(format!("{} plan for {} is {}", feature.title(), plan.workspace, subscription.status))
395 }
396
397 /// The workspace a Stripe customer belongs to.
398 async fn workspace_of_customer(&self, customer: &str) -> Result<Option<String>> {
399 #[derive(Deserialize)]
400 struct Row {
401 workspace: String,
402 }
403 Ok(self
404 .db
405 .prepare("SELECT workspace FROM accounts WHERE customer_id = ?")
406 .bind(&[customer.into()])?
407 .first::<Row>(None)
408 .await?
409 .map(|row| row.workspace))
410 }
411
412 /// The workspace a workspace invoice was sent to, under its slug now.
413 async fn workspace_of_invoice(&self, invoice_id: &str) -> Result<Option<String>> {
414 #[derive(Deserialize)]
415 struct Row {
416 workspace: String,
417 }
418 Ok(self
419 .db
420 .prepare("SELECT workspace FROM workspace_invoices WHERE invoice_id = ?")
421 .bind(&[invoice_id.into()])?
422 .first::<Row>(None)
423 .await?
424 .map(|row| row.workspace))
425 }
426
427 /// Money given back: what was paid is less, by the refund.
428 async fn refunded(&self, charge: &Value) -> Result<String> {
429 let Some(customer) = charge["customer"].as_str() else { return Ok("ignored: no customer".to_owned()) };
430 let Some(workspace) = self.workspace_of_customer(customer).await? else {
431 return Ok("ignored: not a workspace's customer".to_owned());
432 };
433 let refunded = charge["amount_refunded"].as_i64().unwrap_or(0);
434 if refunded <= 0 {
435 return Ok("ignored: nothing refunded".to_owned());
436 }
437 // Each refund total once, so partial refunds add up correctly.
438 let charge_id = charge["id"].as_str().unwrap_or_default();
439 #[derive(Deserialize)]
440 struct Sum {
441 micros: Option<i64>,
442 }
443 let already = self
444 .db
445 .prepare("SELECT -SUM(amount_micros) AS micros FROM ledger WHERE reference LIKE ?")
446 .bind(&[format!("refund/{charge_id}/%").into()])?
447 .first::<Sum>(None)
448 .await?
449 .and_then(|s| s.micros)
450 .unwrap_or(0);
451 let new = refunded * 10_000 - already;
452 if new <= 0 {
453 return Ok("ignored: refund already recorded".to_owned());
454 }
455 self.enter(
456 &workspace,
457 EntryKind::TopUp,
458 -new,
459 "Refunded to the card",
460 &format!("refund/{charge_id}/{refunded}"),
461 None,
462 None,
463 None,
464 None,
465 )
466 .await?;
467 Ok(format!("refund of {} recorded for {workspace}", crate::features::dollars(new)))
468 }
469
470 /// A disputed payment stops the workspace's work until it is resolved;
471 /// one closed in the workspace's favour lets it go on.
472 async fn disputed(&self, dispute: &Value, stop: bool) -> Result<String> {
473 let Some(stripe) = &self.stripe else { return Ok("ignored".to_owned()) };
474 let Some(charge_id) = dispute["charge"].as_str() else { return Ok("ignored: no charge".to_owned()) };
475 let charge: Value = stripe.get(&format!("/charges/{charge_id}")).await?;
476 let Some(workspace) = (match charge["customer"].as_str() {
477 Some(customer) => self.workspace_of_customer(customer).await?,
478 None => None,
479 }) else {
480 return Ok("ignored: not a workspace's customer".to_owned());
481 };
482 let now = rfc3339(now_ms());
483 // The disputed payment never counts toward trust again.
484 for reference in [charge["payment_intent"].as_str(), charge["invoice"].as_str()].into_iter().flatten() {
485 self.db
486 .prepare("UPDATE ledger SET disputed = ? WHERE workspace = ? AND reference = ?")
487 .bind(&[(if stop { 1 } else { 0 }).into(), workspace.as_str().into(), reference.into()])?
488 .run()
489 .await?;
490 }
491 if stop {
492 self.db
493 .prepare(
494 "INSERT INTO limits (workspace, autopay_failed_at, autopay_error, updated_at) VALUES (?1, ?2, ?3, ?2)
495 ON CONFLICT (workspace) DO UPDATE SET autopay_failed_at = ?2, autopay_error = ?3, updated_at = ?2",
496 )
497 .bind(&[workspace.as_str().into(), now.as_str().into(), "a payment was disputed with the card's bank".into()])?
498 .run()
499 .await?;
500 } else {
501 self.db
502 .prepare("UPDATE limits SET autopay_failed_at = NULL, autopay_error = NULL WHERE workspace = ?")
503 .bind(&[workspace.as_str().into()])?
504 .run()
505 .await?;
506 }
507 let account = self.account_of(&workspace).await?;
508 let what = if stop { "dispute: work stopped" } else { "dispute closed in the workspace's favour" };
509 self.audit(&account.id, "dispute", &format!("{workspace}: {what}"), "stripe").await?;
510 Ok(format!("{workspace}: {what}"))
511 }
512
513 // --- Enterprise invoices ------------------------------------------------
514
515 pub(crate) async fn admin_enterprise_billing(&self, a: AdminEnterpriseBillingArgs) -> Result<Outcome<BillingAccount>> {
516 let email = a.email.trim().to_lowercase();
517 if !email.contains('@') || email.contains(char::is_whitespace) || a.by.trim().is_empty() {
518 return Ok(Outcome::fail(FailureCode::Invalid, "Give the email the enterprise's invoices go to."));
519 }
520 let Some(stripe) = &self.stripe else {
521 return Ok(Outcome::fail(FailureCode::Conflict, "Payments are not set up on this g1t."));
522 };
523 #[derive(Deserialize)]
524 struct Row {
525 name: String,
526 kind: String,
527 customer_id: Option<String>,
528 }
529 let Some(row) = self
530 .db
531 .prepare("SELECT name, kind, customer_id FROM billing_accounts WHERE id = ?")
532 .bind(&[a.id.as_str().into()])?
533 .first::<Row>(None)
534 .await?
535 .filter(|row| row.kind == "enterprise")
536 else {
537 return Ok(Outcome::fail(FailureCode::NotFound, "No such enterprise."));
538 };
539 #[derive(Deserialize)]
540 struct Customer {
541 id: String,
542 }
543 let fields = [
544 ("name", row.name.clone()),
545 ("email", email.clone()),
546 ("metadata[g1t_enterprise]", a.id.clone()),
547 ];
548 let customer: Customer = match &row.customer_id {
549 Some(id) => stripe.post(&format!("/customers/{id}"), &fields).await?,
550 None => stripe.post("/customers", &fields).await?,
551 };
552 self.db
553 .prepare("UPDATE billing_accounts SET billing_email = ?, customer_id = ? WHERE id = ?")
554 .bind(&[email.as_str().into(), customer.id.as_str().into(), a.id.as_str().into()])?
555 .run()
556 .await?;
557 self.audit(&a.id, "billing_email", &format!("Invoices go to {email}"), &a.by).await?;
558 Ok(match self.enterprise(&a.id).await? {
559 Some(account) => Outcome::Ok(account),
560 None => Outcome::fail(FailureCode::NotFound, "No such enterprise."),
561 })
562 }
563
564 pub(crate) async fn admin_invoice_enterprise(&self, a: AdminInvoiceEnterpriseArgs) -> Result<Outcome<EnterpriseInvoice>> {
565 if a.by.trim().is_empty() {
566 return Ok(Outcome::fail(FailureCode::Invalid, "Say who is sending it."));
567 }
568 match self.invoice_enterprise(&a.id, "now", &a.by).await? {
569 Ok(invoice) => Ok(Outcome::Ok(invoice)),
570 Err(why) => Ok(Outcome::fail(FailureCode::Conflict, why)),
571 }
572 }
573
574 /// Invoices each enterprise for the month that closed. Live payments
575 /// only; sudo can send one sooner in test mode.
576 pub(crate) async fn invoice_enterprises(&self) -> Result<()> {
577 if !self.stripe.as_ref().is_some_and(crate::stripe::Stripe::live) {
578 return Ok(());
579 }
580 let closing = crate::limits::previous_month(&rfc3339(now_ms())[..7]);
581 #[derive(Deserialize)]
582 struct Id {
583 id: String,
584 }
585 let due = self
586 .db
587 .prepare(
588 "SELECT id FROM billing_accounts WHERE kind = 'enterprise' AND customer_id IS NOT NULL
589 AND terms_kind <> 'comped'
590 AND NOT EXISTS (SELECT 1 FROM enterprise_invoices i WHERE i.account_id = billing_accounts.id AND i.period = ?)",
591 )
592 .bind(&[closing.as_str().into()])?
593 .all()
594 .await?
595 .results::<Id>()?;
596 for Id { id } in due {
597 if let Err(why) = self.invoice_enterprise(&id, &closing, "month close").await? {
598 worker::console_log!("enterprise {id} not invoiced for {closing}: {why}");
599 }
600 }
601 Ok(())
602 }
603
604 /// One invoice for what each of the enterprise's workspaces owes now.
605 async fn invoice_enterprise(&self, id: &str, period: &str, by: &str) -> Result<std::result::Result<EnterpriseInvoice, String>> {
606 let Some(stripe) = &self.stripe else { return Ok(Err("Payments are not set up.".into())) };
607 let Some(account) = self.enterprise(id).await? else { return Ok(Err("No such enterprise.".into())) };
608 #[derive(Deserialize)]
609 struct Customer {
610 customer_id: Option<String>,
611 }
612 let Some(customer) = self
613 .db
614 .prepare("SELECT customer_id FROM billing_accounts WHERE id = ?")
615 .bind(&[id.into()])?
616 .first::<Customer>(None)
617 .await?
618 .and_then(|row| row.customer_id)
619 else {
620 return Ok(Err("Set where the enterprise's invoices go first.".into()));
621 };
622 // What each workspace owes: its charges less what it has paid, and
623 // less what is on invoices still open.
624 let mut lines = vec![];
625 for workspace in &account.workspaces {
626 let balance = self.row(workspace).await?.map_or(0, |row| row.balance_micros);
627 #[derive(Deserialize)]
628 struct Sum {
629 micros: Option<i64>,
630 }
631 let invoiced = self
632 .db
633 .prepare(
634 "SELECT SUM(l.amount_micros) AS micros FROM enterprise_invoice_lines l
635 JOIN enterprise_invoices i ON i.invoice_id = l.invoice_id
636 WHERE l.workspace = ? AND i.status IN ('open', 'overdue')",
637 )
638 .bind(&[workspace.as_str().into()])?
639 .first::<Sum>(None)
640 .await?
641 .and_then(|s| s.micros)
642 .unwrap_or(0);
643 let owed = (-balance).max(0) - invoiced;
644 if owed >= 10_000 {
645 lines.push(InvoiceLine { workspace: workspace.clone(), amount_micros: owed });
646 }
647 }
648 if lines.is_empty() {
649 return Ok(Err("Its workspaces owe nothing to invoice.".into()));
650 }
651 for line in &lines {
652 let cents = (line.amount_micros + 9_999) / 10_000;
653 let fields = [
654 ("customer", customer.clone()),
655 ("amount", cents.to_string()),
656 ("currency", "usd".to_owned()),
657 ("description", format!("{}: g1t usage", line.workspace)),
658 ("metadata[workspace]", line.workspace.clone()),
659 ];
660 let _: Value = stripe.post("/invoiceitems", &fields).await?;
661 }
662 let fields = [
663 ("customer", customer.clone()),
664 ("collection_method", "send_invoice".to_owned()),
665 ("days_until_due", "30".to_owned()),
666 ("pending_invoice_items_behavior", "include".to_owned()),
667 ("description", format!("g1t usage for the {} enterprise", account.name)),
668 ("metadata[g1t_enterprise]", id.to_owned()),
669 ("metadata[period]", period.to_owned()),
670 ];
671 #[derive(Deserialize)]
672 struct Invoice {
673 id: String,
674 #[serde(default)]
675 hosted_invoice_url: Option<String>,
676 #[serde(default)]
677 amount_due: i64,
678 }
679 let draft: Invoice = stripe.post("/invoices", &fields).await?;
680 let _: Value = stripe.post(&format!("/invoices/{}/finalize", draft.id), &[]).await?;
681 let sent: Invoice = stripe.post(&format!("/invoices/{}/send", draft.id), &[]).await?;
682 let now = rfc3339(now_ms());
683 let total = lines.iter().map(|l| l.amount_micros).sum::<i64>().max(sent.amount_due * 10_000);
684 let mut writes = vec![self
685 .db
686 .prepare(
687 "INSERT INTO enterprise_invoices (invoice_id, account_id, period, amount_micros, status, hosted_url, created_by, created_at)
688 VALUES (?, ?, ?, ?, 'open', ?, ?, ?)",
689 )
690 .bind(&[
691 sent.id.as_str().into(),
692 id.into(),
693 period.into(),
694 (total as f64).into(),
695 crate::optional(sent.hosted_invoice_url.as_deref()),
696 by.into(),
697 now.as_str().into(),
698 ])?];
699 for line in &lines {
700 writes.push(
701 self.db
702 .prepare("INSERT INTO enterprise_invoice_lines (invoice_id, workspace, amount_micros) VALUES (?, ?, ?)")
703 .bind(&[sent.id.as_str().into(), line.workspace.as_str().into(), (line.amount_micros as f64).into()])?,
704 );
705 }
706 self.db.batch(writes).await?;
707 self.audit(id, "invoice", &format!("Invoice {} for {} sent ({period})", sent.id, crate::features::dollars(total)), by)
708 .await?;
709 Ok(Ok(EnterpriseInvoice {
710 invoice_id: sent.id,
711 hosted_url: sent.hosted_invoice_url,
712 amount_micros: total,
713 status: "open".to_owned(),
714 period: period.to_owned(),
715 lines,
716 created_at: now,
717 }))
718 }
719
720 /// An enterprise invoice paid: each workspace is credited its line, and
721 /// any stop for the invoice is lifted.
722 async fn enterprise_invoice_paid(&self, invoice_id: &str) -> Result<String> {
723 let claimed = self
724 .db
725 .prepare(
726 "UPDATE enterprise_invoices SET status = 'paid', paid_at = ? WHERE invoice_id = ? AND status <> 'paid'
727 RETURNING invoice_id",
728 )
729 .bind(&[rfc3339(now_ms()).into(), invoice_id.into()])?
730 .first::<Value>(None)
731 .await?;
732 if claimed.is_none() {
733 return Ok("ignored: not an open enterprise invoice".to_owned());
734 }
735 let lines = self.invoice_lines(invoice_id).await?;
736 for line in &lines {
737 self.enter(
738 &line.workspace,
739 EntryKind::TopUp,
740 line.amount_micros,
741 &format!("Paid on the enterprise's invoice {invoice_id}"),
742 &format!("inv/{invoice_id}/{}", line.workspace),
743 None,
744 None,
745 None,
746 None,
747 )
748 .await?;
749 }
750 Ok(format!("invoice {invoice_id} paid; {} workspaces credited", lines.len()))
751 }
752
753 /// An enterprise invoice that went overdue stops its workspaces' work;
754 /// one voided is simply closed.
755 async fn enterprise_invoice_status(&self, invoice_id: &str, status: &str) -> Result<String> {
756 let updated = self
757 .db
758 .prepare("UPDATE enterprise_invoices SET status = ? WHERE invoice_id = ? AND status <> 'paid' RETURNING invoice_id")
759 .bind(&[status.into(), invoice_id.into()])?
760 .first::<Value>(None)
761 .await?;
762 if updated.is_none() {
763 return Ok("ignored: not an open enterprise invoice".to_owned());
764 }
765 if status == "overdue" {
766 let now = rfc3339(now_ms());
767 for line in self.invoice_lines(invoice_id).await? {
768 self.db
769 .prepare(
770 "INSERT INTO limits (workspace, autopay_failed_at, autopay_error, updated_at) VALUES (?1, ?2, ?3, ?2)
771 ON CONFLICT (workspace) DO UPDATE SET autopay_failed_at = ?2, autopay_error = ?3, updated_at = ?2",
772 )
773 .bind(&[line.workspace.as_str().into(), now.as_str().into(), format!("the enterprise's invoice {invoice_id} is overdue").into()])?
774 .run()
775 .await?;
776 }
777 }
778 Ok(format!("invoice {invoice_id} is {status}"))
779 }
780
781 async fn invoice_lines(&self, invoice_id: &str) -> Result<Vec<InvoiceLine>> {
782 Ok(self
783 .db
784 .prepare("SELECT workspace, amount_micros FROM enterprise_invoice_lines WHERE invoice_id = ?")
785 .bind(&[invoice_id.into()])?
786 .all()
787 .await?
788 .results::<LineRow>()?
789 .into_iter()
790 .map(|row| InvoiceLine { workspace: row.workspace, amount_micros: row.amount_micros })
791 .collect())
792 }
793
794 /// An enterprise's invoices, newest first.
795 pub(crate) async fn enterprise_invoices(&self, id: &str) -> Result<Vec<EnterpriseInvoice>> {
796 let rows = self
797 .db
798 .prepare("SELECT * FROM enterprise_invoices WHERE account_id = ? ORDER BY created_at DESC LIMIT 24")
799 .bind(&[JsValue::from(id)])?
800 .all()
801 .await?
802 .results::<InvoiceRow>()?;
803 let mut invoices = vec![];
804 for row in rows {
805 invoices.push(EnterpriseInvoice {
806 lines: self.invoice_lines(&row.invoice_id).await?,
807 invoice_id: row.invoice_id,
808 hosted_url: row.hosted_url,
809 amount_micros: row.amount_micros,
810 status: row.status,
811 period: row.period,
812 created_at: row.created_at,
813 });
814 }
815 Ok(invoices)
816 }
817}
818
819#[cfg(test)]
820mod tests {
821 use super::*;
822
823 fn sign(payload: &str, secret: &str, t: i64) -> String {
824 let mut mac = Hmac::<Sha256>::new_from_slice(secret.as_bytes()).unwrap();
825 mac.update(format!("{t}.{payload}").as_bytes());
826 format!("t={t},v1={}", hex::encode(mac.finalize().into_bytes()))
827 }
828
829 #[test]
830 fn a_signed_event_is_believed_only_as_signed_and_only_fresh() {
831 let payload = r#"{"id":"evt_1","type":"invoice.paid"}"#;
832 let header = sign(payload, "whsec_test", 1_000_000);
833 assert!(verify(payload, &header, "whsec_test", 1_000_010));
834 assert!(!verify(payload, &header, "whsec_other", 1_000_010));
835 assert!(!verify(&payload.replace("paid", "voided"), &header, "whsec_test", 1_000_010));
836 assert!(!verify(payload, &header, "whsec_test", 1_000_000 + 301));
837 assert!(!verify(payload, "v1=abc", "whsec_test", 1_000_000));
838 // Stripe may sign with more than one secret while one is rolled.
839 let both = format!("{},v1=00ff", sign(payload, "whsec_test", 1_000_000));
840 assert!(verify(payload, &both, "whsec_test", 1_000_000));
841 }
842}