g1t/services/billing/src/rename.rs
| 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 | |
| 14 | use g1t_contracts::events::{Event, WorkspaceRenamed}; |
| 15 | use g1t_contracts::identity::UsernamesArgs; |
| 16 | use serde::Deserialize; |
| 17 | use serde_json::Value; |
| 18 | use worker::wasm_bindgen::JsValue; |
| 19 | use worker::{Fetcher, Result}; |
| 20 | |
| 21 | use crate::Billing; |
| 22 | use crate::accounts::own_account; |
| 23 | |
| 24 | /// The event this module handles. |
| 25 | pub(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. |
| 33 | pub(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. |
| 125 | pub(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. |
| 138 | pub(crate) fn values(stale: &str, current: &str) -> [String; 4] { |
| 139 | [current.to_owned(), stale.to_owned(), own_account(current), own_account(stale)] |
| 140 | } |
| 141 | |
| 142 | impl 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, ¤t.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)] |
| 222 | mod 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 | } |