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

255 lines13,569 bytesCodeBlame
1//! A workspace renamed: everything billing keeps under its slug moves to
2//! the new one.
3//!
4//! Identity publishes `workspace.renamed` with the old and new slugs. The
5//! handler asks identity what the workspace is called now, so that a late
6//! or repeated delivery still converges on the current slug, and moves the
7//! rows of each stale slug in one D1 batch (one transaction).
8//!
9//! Rows may already exist under the current slug: usage recorded in the
10//! seconds between the rename and this event. Those are merged, never
11//! dropped where money is concerned: balances are summed, sandbox seconds
12//! are summed, and owners' limits are kept field by field.
13
14use g1t_contracts::events::{Event, WorkspaceRenamed};
15use g1t_contracts::identity::UsernamesArgs;
16use serde::Deserialize;
17use serde_json::Value;
18use worker::wasm_bindgen::JsValue;
19use worker::{Fetcher, Result};
20
21use crate::Billing;
22use crate::accounts::own_account;
23
24/// The event this module handles.
25pub(crate) const RENAMED: &str = "workspace.renamed";
26
27/// The statements that move one stale slug's rows to the current slug, in
28/// order. Parameters: `?1` the current slug, `?2` the stale one, `?3` and
29/// `?4` their own billing accounts (`ws_<slug>`).
30///
31/// Every statement is a no-op once the stale slug has no rows, so running
32/// the batch again changes nothing.
33pub(crate) const STATEMENTS: &[&str] = &[
34 // Plain rows: many per workspace, keyed by their own id.
35 "UPDATE ledger SET workspace = ?1 WHERE workspace = ?2",
36 "UPDATE runs SET workspace = ?1 WHERE workspace = ?2",
37 "UPDATE checkouts SET workspace = ?1 WHERE workspace = ?2",
38 "UPDATE workspace_invoices SET workspace = ?1 WHERE workspace = ?2",
39 "UPDATE sales_notes SET workspace = ?1 WHERE workspace = ?2",
40 // Repositories are named `<slug>/<name>`. Compared exactly rather than
41 // with LIKE, so nothing in a slug is read as a wildcard.
42 "UPDATE ledger SET repo = ?1 || substr(repo, length(?2) + 1) WHERE substr(repo, 1, length(?2) + 1) = ?2 || '/'",
43 "UPDATE runs SET repo = ?1 || substr(repo, length(?2) + 1) WHERE substr(repo, 1, length(?2) + 1) = ?2 || '/'",
44 // The balance is the sum of the ledger, and both ledgers are now the
45 // current slug's: the balances add. The older customer (with the card
46 // and the payment history) is kept when both have one.
47 "INSERT INTO accounts (workspace, balance_micros, customer_id, created_at)
48 SELECT ?1, balance_micros, customer_id, created_at FROM accounts WHERE workspace = ?2
49 ON CONFLICT (workspace) DO UPDATE SET
50 balance_micros = accounts.balance_micros + excluded.balance_micros,
51 customer_id = COALESCE(excluded.customer_id, accounts.customer_id),
52 created_at = MIN(accounts.created_at, excluded.created_at)",
53 "DELETE FROM accounts WHERE workspace = ?2",
54 // Sandbox time counts toward the month's free minutes: the seconds add.
55 "INSERT INTO sandbox_months (workspace, month, seconds)
56 SELECT ?1, month, seconds FROM sandbox_months WHERE workspace = ?2
57 ON CONFLICT (workspace, month) DO UPDATE SET seconds = sandbox_months.seconds + excluded.seconds",
58 "DELETE FROM sandbox_months WHERE workspace = ?2",
59 // Replaced on each report with the month's whole figure: the newer
60 // report wins.
61 "INSERT INTO pending_usage (workspace, source, month, charge_micros, updated_at)
62 SELECT ?1, source, month, charge_micros, updated_at FROM pending_usage WHERE workspace = ?2
63 ON CONFLICT (workspace, source, month) DO UPDATE SET
64 charge_micros = CASE WHEN excluded.updated_at > pending_usage.updated_at
65 THEN excluded.charge_micros ELSE pending_usage.charge_micros END,
66 updated_at = MAX(pending_usage.updated_at, excluded.updated_at)",
67 "DELETE FROM pending_usage WHERE workspace = ?2",
68 // Limits, field by field: a ceiling or an owner's spend limit set under
69 // either slug is kept (the current slug's if both), a stop for a
70 // declined card stays, and the highest warning this month is kept.
71 "INSERT INTO limits (workspace, ceiling_micros, spend_limit_micros, spend_limit_full, autopay_failed_at,
72 autopay_error, warned_month, warned_level, declined_told_at, updated_at)
73 SELECT ?1, ceiling_micros, spend_limit_micros, spend_limit_full, autopay_failed_at,
74 autopay_error, warned_month, warned_level, declined_told_at, updated_at
75 FROM limits WHERE workspace = ?2
76 ON CONFLICT (workspace) DO UPDATE SET
77 ceiling_micros = COALESCE(limits.ceiling_micros, excluded.ceiling_micros),
78 spend_limit_micros = CASE WHEN limits.spend_limit_micros IS NOT NULL OR limits.spend_limit_full = 1
79 THEN limits.spend_limit_micros ELSE excluded.spend_limit_micros END,
80 spend_limit_full = CASE WHEN limits.spend_limit_micros IS NOT NULL OR limits.spend_limit_full = 1
81 THEN limits.spend_limit_full ELSE excluded.spend_limit_full END,
82 autopay_error = CASE WHEN limits.autopay_failed_at IS NOT NULL
83 THEN limits.autopay_error ELSE excluded.autopay_error END,
84 autopay_failed_at = COALESCE(limits.autopay_failed_at, excluded.autopay_failed_at),
85 warned_level = CASE WHEN excluded.warned_month > COALESCE(limits.warned_month, '') THEN excluded.warned_level
86 WHEN excluded.warned_month = limits.warned_month THEN MAX(limits.warned_level, excluded.warned_level)
87 ELSE limits.warned_level END,
88 warned_month = COALESCE(MAX(limits.warned_month, excluded.warned_month), limits.warned_month, excluded.warned_month),
89 declined_told_at = COALESCE(MAX(limits.declined_told_at, excluded.declined_told_at), limits.declined_told_at, excluded.declined_told_at),
90 updated_at = MAX(limits.updated_at, excluded.updated_at)",
91 "DELETE FROM limits WHERE workspace = ?2",
92 // Plans: one per feature. A live plan under the stale slug replaces a
93 // canceled one under the current slug; otherwise the current one stays.
94 "DELETE FROM subscriptions WHERE workspace = ?1 AND status = 'canceled'
95 AND feature IN (SELECT feature FROM subscriptions WHERE workspace = ?2 AND status <> 'canceled')",
96 "UPDATE OR IGNORE subscriptions SET workspace = ?1 WHERE workspace = ?2",
97 "DELETE FROM subscriptions WHERE workspace = ?2",
98 // Records of one per workspace (per month): the current slug's stays
99 // if it has one.
100 "UPDATE OR IGNORE month_closes SET workspace = ?1 WHERE workspace = ?2",
101 "DELETE FROM month_closes WHERE workspace = ?2",
102 "UPDATE OR IGNORE account_members SET workspace = ?1 WHERE workspace = ?2",
103 "DELETE FROM account_members WHERE workspace = ?2",
104 "UPDATE OR IGNORE sales_records SET workspace = ?1 WHERE workspace = ?2",
105 "DELETE FROM sales_records WHERE workspace = ?2",
106 // An enterprise invoice's line for the workspace: amounts add.
107 "INSERT INTO enterprise_invoice_lines (invoice_id, workspace, amount_micros)
108 SELECT invoice_id, ?1, amount_micros FROM enterprise_invoice_lines WHERE workspace = ?2
109 ON CONFLICT (invoice_id, workspace) DO UPDATE SET
110 amount_micros = enterprise_invoice_lines.amount_micros + excluded.amount_micros",
111 "DELETE FROM enterprise_invoice_lines WHERE workspace = ?2",
112 // The workspace's own billing account, `ws_<slug>`. Terms set by staff
113 // are kept: a standard row under the current slug gives way to the
114 // stale one; two with terms keep the current one.
115 "DELETE FROM billing_accounts WHERE id = ?3 AND terms_kind = 'standard'
116 AND EXISTS (SELECT 1 FROM billing_accounts WHERE id = ?4)",
117 "UPDATE OR IGNORE billing_accounts SET id = ?3, name = CASE WHEN name = ?2 THEN ?1 ELSE name END WHERE id = ?4",
118 "DELETE FROM billing_accounts WHERE id = ?4",
119 "UPDATE admin_actions SET account = ?3 WHERE account = ?4",
120 "UPDATE account_members SET account_id = ?3 WHERE account_id = ?4",
121 "UPDATE enterprise_invoices SET account_id = ?3 WHERE account_id = ?4",
122];
123
124/// The highest `?N` a statement names: how many values it is bound with.
125pub(crate) fn parameters(sql: &str) -> usize {
126 let bytes = sql.as_bytes();
127 let mut highest = 0;
128 for (i, byte) in bytes.iter().enumerate() {
129 if *byte == b'?' {
130 let digits: String = bytes[i + 1..].iter().take_while(|b| b.is_ascii_digit()).map(|b| *b as char).collect();
131 highest = highest.max(digits.parse().unwrap_or(0));
132 }
133 }
134 highest
135}
136
137/// The values `STATEMENTS` are bound with, in parameter order.
138pub(crate) fn values(stale: &str, current: &str) -> [String; 4] {
139 [current.to_owned(), stale.to_owned(), own_account(current), own_account(stale)]
140}
141
142impl Billing {
143 /// Handles one event from the bus; every type but a rename is ignored.
144 pub(crate) async fn on_event(&self, identity: Option<&Fetcher>, event: &Event) -> Result<()> {
145 if event.kind != RENAMED {
146 return Ok(());
147 }
148 let renamed: WorkspaceRenamed = serde_json::from_value(event.data.clone())?;
149 let current = match identity {
150 Some(identity) => {
151 let names: std::collections::HashMap<String, String> = g1t_kit::call(
152 identity,
153 "usernames",
154 &UsernamesArgs { ids: vec![renamed.workspace_id.clone()] },
155 )
156 .await?;
157 names.get(&renamed.workspace_id).cloned().unwrap_or_else(|| renamed.to.clone())
158 }
159 None => renamed.to.clone(),
160 };
161 self.rename_workspace(&renamed, &current.to_lowercase()).await
162 }
163
164 /// Moves every row of the rename's stale slugs to `current`.
165 pub(crate) async fn rename_workspace(&self, renamed: &WorkspaceRenamed, current: &str) -> Result<()> {
166 for stale in renamed.stale_slugs(current) {
167 let stale = stale.to_lowercase();
168 if stale.is_empty() || stale == current {
169 continue;
170 }
171 let values = values(&stale, current);
172 let mut batch = Vec::with_capacity(STATEMENTS.len());
173 for sql in STATEMENTS {
174 let binds: Vec<JsValue> = values[..parameters(sql)].iter().map(|v| JsValue::from(v.as_str())).collect();
175 batch.push(self.db.prepare(*sql).bind(&binds)?);
176 }
177 self.db.batch(batch).await?;
178 self.rename_customer(&stale, current).await;
179 }
180 Ok(())
181 }
182
183 /// The Stripe customer carries the slug as its name and metadata; both
184 /// follow the rename. A name someone changed at Stripe is left alone.
185 /// Best effort: a failure here moves no money.
186 async fn rename_customer(&self, stale: &str, current: &str) {
187 let Some(stripe) = &self.stripe else { return };
188 let Ok(Some(customer)) = self.row(current).await.map(|row| row.and_then(|row| row.customer_id)) else {
189 return;
190 };
191 #[derive(Deserialize)]
192 struct Customer {
193 name: Option<String>,
194 metadata: Option<std::collections::HashMap<String, String>>,
195 }
196 let path = format!("/customers/{customer}");
197 let found: Customer = match stripe.get(&path).await {
198 Ok(found) => found,
199 Err(error) => {
200 worker::console_error!("renaming {stale} at Stripe: could not read {customer}: {error}");
201 return;
202 }
203 };
204 let mut fields = vec![];
205 if found.name.as_deref() == Some(stale) {
206 fields.push(("name", current.to_owned()));
207 }
208 let tagged = found.metadata.as_ref().and_then(|m| m.get("workspace")).map(String::as_str);
209 if tagged != Some(current) {
210 fields.push(("metadata[workspace]", current.to_owned()));
211 }
212 if fields.is_empty() {
213 return;
214 }
215 if let Err(error) = stripe.post::<Value>(&path, &fields).await {
216 worker::console_error!("renaming {stale} at Stripe: could not update {customer}: {error}");
217 }
218 }
219}
220
221#[cfg(test)]
222mod tests {
223 use super::*;
224
225 #[test]
226 fn every_statement_is_bound_with_what_it_names() {
227 for sql in STATEMENTS {
228 let n = parameters(sql);
229 assert!((1..=4).contains(&n), "{sql}");
230 // Numbered only: a bare `?` would take the wrong value.
231 assert!(!sql.contains("? ") && !sql.ends_with('?'), "{sql}");
232 }
233 assert_eq!(parameters("UPDATE t SET a = ?1 WHERE b = ?2"), 2);
234 assert_eq!(parameters("UPDATE t SET a = ?3 WHERE b = ?4"), 4);
235 }
236
237 #[test]
238 fn the_values_name_both_slugs_and_their_accounts() {
239 assert_eq!(values("acme", "acme-co"), ["acme-co".to_owned(), "acme".into(), "ws_acme-co".into(), "ws_acme".into()]);
240 }
241
242 #[test]
243 fn every_table_keyed_by_a_slug_is_moved() {
244 let all = STATEMENTS.join("\n");
245 for table in [
246 "ledger", "runs", "checkouts", "workspace_invoices", "sales_notes", "accounts", "sandbox_months",
247 "pending_usage", "limits", "subscriptions", "month_closes", "account_members", "sales_records",
248 "enterprise_invoice_lines", "billing_accounts", "admin_actions", "enterprise_invoices",
249 ] {
250 assert!(all.contains(&format!("FROM {table} WHERE workspace = ?2"))
251 || all.contains(&format!("UPDATE {table} SET"))
252 || all.contains(&format!("UPDATE OR IGNORE {table} SET")), "{table}");
253 }
254 }
255}