Commit

Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept

- The saved card is kept on the account row (0027) and read from there; Stripe is asked only the first time, or after g1t set a default card. The billing page no longer calls Stripe on every view. - New webhook events keep it current: customer.updated, payment_method.*, setup_intent.succeeded. - Every cron run replays events Stripe sent that billing never saw, from its event list, through the same once-only claim as the webhook (stripe_sync, process_event). A lost delivery costs at most 15 minutes. - Daily: the endpoint gets billing's event list in place, keeping its secret, and is enabled again if Stripe disabled it; stale cards and plans are read again, 25 of each. - Clippy is quiet in billing too. - docs/BILLING_OPERATIONS.md: a Stripe section.

syntaqxcommitted Parentd5df0deBrowse files
13 files+526−700/13 viewed
+39−0
281281 `breaker_lifted`). It resets by itself at 00:00 UTC.
282282 - **Change a default**: edit the variable in `services/billing/wrangler.jsonc`
283283 and deploy g1t-billing.
284+
285+## Stripe
286+
287+Billing keeps what it needs from Stripe so reads never wait on it, and
288+hears of changes three ways (`webhooks.rs`, `stripe_sync.rs`).
289+
290+**The webhook.** sudo → Stripe → **Register webhook** creates the endpoint
291+`https://api.g1t.sh/stripe/webhook` through Stripe's API. Stripe returns the
292+signing secret only in that answer; billing stores it in `stripe_webhooks`,
293+one row per key mode (test, live). Never add the endpoint in Stripe's
294+dashboard: its secret would not reach billing, and every event would fail
295+with "The signature does not match". Register again in each mode (at
296+launch, after the key changes to live). Events are claimed once each in
297+`stripe_events`; a handler that fails forgets its claim, and Stripe retries.
298+
299+**What is kept, and how it stays current**
300+
301+| Kept | Where | Refreshed by |
302+| --- | --- | --- |
303+| Saved card (brand, last 4, expiry) | `accounts.card_*`, `card_synced_at` | `customer.updated`, `payment_method.*`, `setup_intent.succeeded`; forgotten when g1t sets a default card, so the next read asks once; the daily pass for cards older than 7 days |
304+| Plans | `subscriptions` | `customer.subscription.*`, `invoice.*`; a read asks Stripe at most hourly once a period is over; the daily pass for rows older than a day |
305+| Payments, refunds, disputes | ledger, `checkouts`, invoices | their events |
306+
307+**Every cron run (every 15 minutes)** replays missed events: Stripe's event
308+list from an hour before `stripe_sync.through`, oldest first, through the
309+same once-only claim. `through` moves to 5 minutes before now when all were
310+handled, back to the first failure otherwise, and stays when more than
311+1,000 events were listed. Claims stuck at `handling` for 10 minutes are
312+dropped so the replay retries them. The first run reads 3 days back.
313+
314+**Daily** (`keeper::DAILY`): the endpoint is given billing's event list in
315+place (its secret stays) and enabled again if Stripe disabled it, both
316+audited as `stripe`/`webhook`; then up to 25 stale cards and 25 stale plans
317+are read again.
318+
319+**What still calls Stripe on a request**: starting a payment page, a plan
320+or a card check; opening the billing portal; settling a page the person
321+came back from; renaming a workspace (the customer's name). Nothing a page
322+view reads.
+22−0
1+-- What billing keeps of Stripe so reads never wait on it (stripe_sync.rs).
2+
3+-- The workspace's saved card as Stripe last said, refreshed by card and
4+-- customer events and a weekly pass. card_synced_at NULL: never read, or
5+-- forgotten after g1t itself changed the card; the next read asks Stripe.
6+ALTER TABLE accounts ADD COLUMN card_brand TEXT;
7+ALTER TABLE accounts ADD COLUMN card_last4 TEXT;
8+ALTER TABLE accounts ADD COLUMN card_exp_month INTEGER;
9+ALTER TABLE accounts ADD COLUMN card_exp_year INTEGER;
10+ALTER TABLE accounts ADD COLUMN card_synced_at TEXT;
11+
12+-- Card and customer events name the customer, not the workspace.
13+CREATE INDEX IF NOT EXISTS accounts_by_customer ON accounts (customer_id);
14+
15+-- How far missed events have been replayed from Stripe's event list, per
16+-- key mode: every event created before `through` (Unix seconds) has been
17+-- seen, by its webhook or by the replay.
18+CREATE TABLE stripe_sync (
19+ mode TEXT PRIMARY KEY,
20+ through INTEGER NOT NULL,
21+ checked_at TEXT NOT NULL
22+);
+5−7
197197 .bind(&[workspace.as_str().into()])?
198198 .first::<Link>(None)
199199 .await?;
200− if let Some(link) = linked {
201− if let Some(row) = self.account_row(&link.account_id).await? {
200+ if let Some(link) = linked
201+ && let Some(row) = self.account_row(&link.account_id).await? {
202202 let members = self.members(&row.id).await?;
203203 return Ok(self.to_account(&row, members));
204204 }
205− }
206205 let id = own_account(&workspace);
207206 Ok(match self.account_row(&id).await? {
208207 Some(row) => self.to_account(&row, vec![workspace]),
398397 if seen.contains(&id) {
399398 continue;
400399 }
401− if let Some(account) = self.find_account(&id).await? {
402− if a.query.as_deref().is_none_or(|q| account.name.to_lowercase().contains(&q.to_lowercase())) {
400+ if let Some(account) = self.find_account(&id).await?
401+ && a.query.as_deref().is_none_or(|q| account.name.to_lowercase().contains(&q.to_lowercase())) {
403402 seen.insert(id);
404403 summaries.push(self.summary(account).await?);
405404 }
406− }
407405 }
408− summaries.sort_by(|x, y| y.limit.exposure_micros.cmp(&x.limit.exposure_micros));
406+ summaries.sort_by_key(|x| std::cmp::Reverse(x.limit.exposure_micros));
409407 Ok(summaries)
410408 }
411409
+4−3
8989 .bind(&[a.session.as_str().into(), workspace.as_str().into(), CARD_CHECK.into()])?
9090 .first::<serde_json::Value>(None)
9191 .await?;
92− if mine.is_some() {
93− if let Err(why) = self.settle_card_check(&a.session).await? {
92+ if mine.is_some()
93+ && let Err(why) = self.settle_card_check(&a.session).await? {
9494 return Ok(Outcome::fail(FailureCode::Conflict, why));
9595 }
96− }
9796 Ok(Outcome::Ok(self.entitlements(EntitlementsArgs { workspace }).await?))
9897 }
9998
179178 if let Err(error) = stripe.set_default_card(customer, &card.payment_method).await {
180179 worker::console_error!("{workspace}: the checked card was not made the default: {error}");
181180 }
181+ // The page this lands on shows the new card, not the old.
182+ self.forget_card(&workspace).await?;
182183 self.db
183184 .prepare("UPDATE accounts SET customer_id = COALESCE(customer_id, ?2) WHERE workspace = ?1")
184185 .bind(&[workspace.as_str().into(), customer.into()])?
+5−7
245245 /// own workspaces, which are watched in sudo but never paused.
246246 async fn spike_pause(&self, workspace: &str, plan: PlanKind) -> Result<Option<Spike>> {
247247 let latest = self.latest_spike(workspace).await?;
248− if let Some(spike) = &latest {
249− if spike.status == "open" || spike.status == "stopped" {
248+ if let Some(spike) = &latest
249+ && (spike.status == "open" || spike.status == "stopped") {
250250 return Ok(latest);
251251 }
252− }
253252 if plan == PlanKind::Internal || self.stripe.is_none() {
254253 return Ok(None);
255254 }
256255 let pace = self.pace(workspace).await?;
257256 let now = rfc3339(now_ms());
258− if let Some(spike) = &latest {
259− if spike.status == "continued" && still_continued(spike.until.as_deref(), &now, spike.hour_micros, pace.last_hour) {
257+ if let Some(spike) = &latest
258+ && spike.status == "continued" && still_continued(spike.until.as_deref(), &now, spike.hour_micros, pace.last_hour) {
260259 return Ok(None);
261260 }
262− }
263261 if !is_spike(pace.last_hour, pace.usual_hour, self.plans.spike_factor, self.plans.spike_floor_micros) {
264262 return Ok(None);
265263 }
583581 let held = self.held(&account.workspaces).await?;
584582 let (paid_by, hold) = match place(&room, held, estimate, has_plan) {
585583 Ok(placed) => placed,
586− Err(short) => return Ok(self.short(&workspace, plan, &a, verified, &room, short).await?),
584+ Err(short) => return self.short(&workspace, plan, &a, verified, &room, short).await,
587585 };
588586 let id = new_id("rsv", now);
589587 let mut values: Vec<JsValue> = vec![
+3−5
238238 self.enter(workspace, EntryKind::TopUp, amount, &format!("Paid invoice {}", invoice.id), &invoice.id, None, None, None, None)
239239 .await?;
240240 // Prepaid cards pay, but never raise the limit.
241− if let (Some(stripe), Some(charge)) = (&self.stripe, &invoice.charge) {
242− if let Ok(charge) = stripe.get::<Value>(&format!("/charges/{charge}")).await {
243− if let Some(funding) = charge["payment_method_details"]["card"]["funding"].as_str() {
241+ if let (Some(stripe), Some(charge)) = (&self.stripe, &invoice.charge)
242+ && let Ok(charge) = stripe.get::<Value>(&format!("/charges/{charge}")).await
243+ && let Some(funding) = charge["payment_method_details"]["card"]["funding"].as_str() {
244244 self.db
245245 .prepare("UPDATE ledger SET funding = ? WHERE reference = ?")
246246 .bind(&[funding.into(), invoice.id.as_str().into()])?
247247 .run()
248248 .await?;
249249 }
250− }
251− }
252250 Ok(true)
253251 }
254252
+29−18
3838 mod limits;
3939 mod rename;
4040 mod stripe;
41+mod stripe_sync;
4142
4243 use g1t_contracts::billing::*;
4344 use g1t_contracts::time::rfc3339;
208209 }
209210
210211 async fn standing(&self, workspace: &str) -> Result<Account> {
211− let row = self.row(workspace).await?;
212− let card = match (&self.stripe, row.as_ref().and_then(|row| row.customer_id.as_deref())) {
213− (Some(stripe), Some(customer)) => stripe.card(customer).await.ok().flatten().map(|card| Card {
214− brand: card.brand,
215− last4: card.last4,
216− exp_month: card.exp_month,
217− exp_year: card.exp_year,
218− }),
219− _ => None,
220− };
212+ // The card as last synced (stripe_sync.rs), not asked of Stripe.
213+ let (row, card) = try_join(self.row(workspace), self.saved_card(workspace)).await?;
221214 Ok(Account {
222215 workspace: workspace.to_owned(),
223216 balance_micros: row.map_or(0, |row| row.balance_micros),
594587 /// Drops a saved customer the card processor no longer knows.
595588 pub(crate) async fn forget_customer(&self, workspace: &str) -> Result<()> {
596589 self.db
597− .prepare("UPDATE accounts SET customer_id = NULL WHERE workspace = ?")
590+ .prepare(
591+ "UPDATE accounts SET customer_id = NULL, card_brand = NULL, card_last4 = NULL, card_exp_month = NULL,
592+ card_exp_year = NULL, card_synced_at = NULL WHERE workspace = ?",
593+ )
598594 .bind(&[workspace.into()])?
599595 .run()
600596 .await?;
899895 let digits = n.to_string();
900896 let mut out = String::new();
901897 for (i, c) in digits.chars().enumerate() {
902− if i > 0 && (digits.len() - i) % 3 == 0 {
898+ if i > 0 && (digits.len() - i).is_multiple_of(3) {
903899 out.push(',');
904900 }
905901 out.push(c);
968964 .and_then(|percent| percent.to_string().parse().ok())
969965 .unwrap_or(20),
970966 free: env.var("FREE_WHILE_BUILDING").is_ok_and(|v| v.to_string() == "true"),
971− ceilings: limits::Ceilings::from_env(&env),
967+ ceilings: limits::Ceilings::from_env(env),
972968 prepaid_only: env.var("PREPAID_ONLY").is_ok_and(|v| v.to_string() == "true"),
973969 trials_on: {
974970 let plans = credits::Config::from_env(env);
992988 if let Err(error) = billing.settle_runs(&keeper).await {
993989 worker::console_error!("settling runs failed: {error}");
994990 }
991+ // Stripe events billing never received, handled now.
992+ match billing.replay_events().await {
993+ Ok(done) => worker::console_log!("stripe replay: {done}"),
994+ Err(error) => worker::console_error!("replaying Stripe events failed: {error}"),
995+ }
995996 if let Err(error) = billing.autopay().await {
996997 worker::console_error!("paying at the limit failed: {error}");
997998 }
10101011 if let Err(error) = billing.watch_spend().await {
10111012 worker::console_error!("watching g1t's own spend failed: {error}");
10121013 }
1013− if let Ok(identity) = env.service("IDENTITY") {
1014− if let Err(error) = billing.warn_limits(&identity).await {
1014+ if let Ok(identity) = env.service("IDENTITY")
1015+ && let Err(error) = billing.warn_limits(&identity).await {
10151016 worker::console_error!("warning owners failed: {error}");
10161017 }
1018+ // Once a day: Stripe's endpoint kept listening to billing's events and
1019+ // enabled, and saved cards and plans not read in a while read again.
1020+ if event.cron() == keeper::DAILY {
1021+ match billing.keep_endpoint().await {
1022+ Ok(done) => worker::console_log!("stripe endpoint: {done}"),
1023+ Err(error) => worker::console_error!("keeping Stripe's endpoint failed: {error}"),
1024+ }
1025+ match billing.refresh_from_stripe().await {
1026+ Ok(done) => worker::console_log!("stripe refresh: {done}"),
1027+ Err(error) => worker::console_error!("refreshing from Stripe failed: {error}"),
1028+ }
10171029 }
10181030 // Once a day, and at once if the costs were never checked: check every
10191031 // cost against what Cloudflare billed.
1020− if event.cron() == keeper::DAILY || billing.never_checked().await.unwrap_or(false) {
1021− if let Err(error) = billing.reconcile(&keeper).await {
1032+ if (event.cron() == keeper::DAILY || billing.never_checked().await.unwrap_or(false))
1033+ && let Err(error) = billing.reconcile(&keeper).await {
10221034 worker::console_error!("checking costs against Cloudflare failed: {error}");
10231035 }
1024− }
10251036 // Once a day: what Cloudflare charged, reconciled against what g1t
10261037 // counted and charged; prices whose day has come; margin alerts
10271038 // (margin.rs). After the keeper, so its proposals are in.
+3−6
423423 let risk = state(exposure, ceiling);
424424 let budget = state(spent, spend_limit);
425425 let over_budget = budget == LimitState::Stopped;
426− let state = if declined.is_some() && exposure > 0 {
427− LimitState::Stopped
428− } else if risk == LimitState::Stopped || over_budget {
426+ let state = if (declined.is_some() && exposure > 0) || risk == LimitState::Stopped || over_budget {
429427 LimitState::Stopped
430428 } else if risk == LimitState::Warning || budget == LimitState::Warning {
431429 LimitState::Warning
839837 let billing = format!("https://g1t.sh/{workspace}/-/billing");
840838
841839 // A declined card, once per decline.
842− if let Some(Told { autopay_failed_at: Some(failed), declined_told_at }) = &told {
843− if declined_told_at.as_deref().is_none_or(|at| at < failed.as_str()) {
840+ if let Some(Told { autopay_failed_at: Some(failed), declined_told_at }) = &told
841+ && declined_told_at.as_deref().is_none_or(|at| at < failed.as_str()) {
844842 let limit = self.limit_of(&workspace).await?;
845843 let sent = notify(
846844 identity,
859857 .await?;
860858 }
861859 }
862− }
863860
864861 // 50, 75, 90 and 100%, once each a month and meter: only the
865862 // highest new level is emailed.
+2−2
4444 }
4545 sorted.sort_unstable();
4646 let middle = sorted.len() / 2;
47− if sorted.len() % 2 == 0 { (sorted[middle - 1] + sorted[middle]) / 2 } else { sorted[middle] }
47+ if sorted.len().is_multiple_of(2) { (sorted[middle - 1] + sorted[middle]) / 2 } else { sorted[middle] }
4848 }
4949
5050 /// Whether this month is well past the typical one.
214214 queue.push(overage);
215215 }
216216 }
217− queue.sort_by(|a, b| b.goodwill.overage_micros.cmp(&a.goodwill.overage_micros));
217+ queue.sort_by_key(|a| std::cmp::Reverse(a.goodwill.overage_micros));
218218 Ok(queue)
219219 }
220220
+2−3
357357 push(SignalKind::Declined, limit.message.clone().unwrap_or_default(), limit.exposure_micros);
358358 } else if limit.state == LimitState::Stopped {
359359 push(SignalKind::AtLimit, limit.message.clone().unwrap_or_default(), limit.exposure_micros.max(limit.spent_micros));
360− } else if let Some(available) = limit.available_micros.filter(|a| *a > 0) {
361− if limit.exposure_micros * 5 >= available * 4 {
360+ } else if let Some(available) = limit.available_micros.filter(|a| *a > 0)
361+ && limit.exposure_micros * 5 >= available * 4 {
362362 push(
363363 SignalKind::NearCeiling,
364364 format!(
370370 limit.exposure_micros,
371371 );
372372 }
373− }
374373 if last_month >= HIGH_SPEND_MICROS {
375374 let terms = self.terms_of(&row.workspace).await?;
376375 if terms.kind == TermsKind::Standard && !limit.account.starts_with("ent_") {
+1−1
214214 group.charged_micros = group.lines.iter().filter(|l| l.charged_micros > 0).map(|l| l.charged_micros).sum();
215215 }
216216 if by_project {
217− groups.sort_by(|a, b| b.charged_micros.cmp(&a.charged_micros));
217+ groups.sort_by_key(|a| std::cmp::Reverse(a.charged_micros));
218218 } else {
219219 groups.sort_by(|a, b| b.key.cmp(&a.key));
220220 }
+364−0
1+//! What billing keeps of Stripe, so that reads never wait on it and nothing
2+//! Stripe says is lost:
3+//!
4+//! - The saved card, on the account row: read from there, asked of Stripe
5+//! only the first time, and refreshed by card and customer events
6+//! (webhooks.rs) and a weekly pass.
7+//! - Missed events, replayed: every run of the cron reads Stripe's event
8+//! list from where the last run got to and handles any event not seen,
9+//! through the same once-only claim the webhook uses. A delivery that
10+//! failed for good, or an endpoint not registered yet, costs at most one
11+//! cron interval.
12+//! - The endpoint, kept: once a day its events are made billing's list in
13+//! place (its signing secret stays), and it is enabled again if Stripe
14+//! turned it off after failures.
15+
16+use crate::Billing;
17+use crate::stripe::form;
18+use crate::webhooks::EVENTS;
19+use g1t_contracts::billing::Card;
20+use g1t_contracts::time::rfc3339;
21+use g1t_kit::now_ms;
22+use serde::Deserialize;
23+use serde_json::Value;
24+use worker::Result;
25+
26+/// How far back each replay starts before where the last one got to, for
27+/// events Stripe lists a little late.
28+const OVERLAP_SECONDS: i64 = 60 * 60;
29+/// How recent an event may be and still be listed later: the replay's
30+/// mark stays this far behind the clock.
31+const SETTLE_SECONDS: i64 = 5 * 60;
32+/// Where the first replay starts.
33+const FIRST_LOOK_SECONDS: i64 = 3 * 24 * 60 * 60;
34+/// Pages of 100 events one replay reads at most.
35+const MAX_PAGES: usize = 10;
36+/// A claim this old that never finished: its handler died, so the event
37+/// may be handled again.
38+const STALE_CLAIM_MS: u64 = 10 * 60 * 1000;
39+/// How long a saved card is believed before the weekly pass asks again.
40+const CARD_FRESH_DAYS: u64 = 7;
41+/// Accounts or plans the daily pass refreshes at most.
42+const REFRESH_BATCH: u32 = 25;
43+
44+#[derive(Deserialize)]
45+struct CardRow {
46+ customer_id: Option<String>,
47+ card_brand: Option<String>,
48+ card_last4: Option<String>,
49+ card_exp_month: Option<u32>,
50+ card_exp_year: Option<u32>,
51+ card_synced_at: Option<String>,
52+}
53+
54+/// Where the replay's mark goes after a run: the first event that failed
55+/// (so it is listed again), the old mark when the list was not read to its
56+/// end, else just behind the clock.
57+pub(crate) fn next_through(previous: i64, now: i64, first_failed: Option<i64>, complete: bool) -> i64 {
58+ match first_failed {
59+ Some(created) => created.min(previous.max(created)),
60+ None if !complete => previous,
61+ None => previous.max(now - SETTLE_SECONDS),
62+ }
63+}
64+
65+impl Billing {
66+ /// The workspace's saved card as last synced. Stripe is asked only when
67+ /// it never was, or g1t changed the card since; a Stripe that does not
68+ /// answer then shows no card, as before.
69+ pub(crate) async fn saved_card(&self, workspace: &str) -> Result<Option<Card>> {
70+ let Some(row) = self
71+ .db
72+ .prepare(
73+ "SELECT customer_id, card_brand, card_last4, card_exp_month, card_exp_year, card_synced_at
74+ FROM accounts WHERE workspace = ?",
75+ )
76+ .bind(&[workspace.into()])?
77+ .first::<CardRow>(None)
78+ .await?
79+ else {
80+ return Ok(None);
81+ };
82+ if row.card_synced_at.is_some() {
83+ return Ok(card_of(&row));
84+ }
85+ let Some(customer) = row.customer_id.as_deref() else { return Ok(None) };
86+ if self.stripe.is_none() {
87+ return Ok(None);
88+ }
89+ Ok(match self.fetch_card(customer).await {
90+ Ok(card) => card,
91+ Err(error) => {
92+ worker::console_error!("reading {workspace}'s card from Stripe failed: {error}");
93+ None
94+ }
95+ })
96+ }
97+
98+ /// Asks Stripe for the customer's card and keeps it on every account
99+ /// with that customer.
100+ async fn fetch_card(&self, customer: &str) -> Result<Option<Card>> {
101+ let Some(stripe) = &self.stripe else { return Ok(None) };
102+ let card = stripe.card(customer).await?.map(|card| Card {
103+ brand: card.brand,
104+ last4: card.last4,
105+ exp_month: card.exp_month,
106+ exp_year: card.exp_year,
107+ });
108+ self.db
109+ .prepare(
110+ "UPDATE accounts SET card_brand = ?2, card_last4 = ?3, card_exp_month = ?4, card_exp_year = ?5,
111+ card_synced_at = ?6 WHERE customer_id = ?1",
112+ )
113+ .bind(&[
114+ customer.into(),
115+ card.as_ref().map(|c| c.brand.clone()).into(),
116+ card.as_ref().map(|c| c.last4.clone()).into(),
117+ card.as_ref().map(|c| c.exp_month as f64).into(),
118+ card.as_ref().map(|c| c.exp_year as f64).into(),
119+ rfc3339(now_ms()).into(),
120+ ])?
121+ .run()
122+ .await?;
123+ Ok(card)
124+ }
125+
126+ /// A card or customer event: the saved card read again. The event's
127+ /// outcome, for `stripe_events`.
128+ pub(crate) async fn sync_card_of(&self, customer: &str) -> Result<String> {
129+ if self.workspace_of_customer(customer).await?.is_none() {
130+ return Ok("ignored: not a workspace's customer".to_owned());
131+ }
132+ Ok(match self.fetch_card(customer).await? {
133+ Some(card) => format!("card saved: {} ending {}", card.brand, card.last4),
134+ None => "card saved: none".to_owned(),
135+ })
136+ }
137+
138+ /// Forgets the saved card after g1t itself changed it, so the next
139+ /// read asks Stripe instead of waiting for the event.
140+ pub(crate) async fn forget_card(&self, workspace: &str) -> Result<()> {
141+ self.db
142+ .prepare("UPDATE accounts SET card_synced_at = NULL WHERE workspace = ?")
143+ .bind(&[workspace.into()])?
144+ .run()
145+ .await?;
146+ Ok(())
147+ }
148+
149+ /// Handles events Stripe sent that billing never saw, from its event
150+ /// list. What happened, for the log.
151+ pub(crate) async fn replay_events(&self) -> Result<String> {
152+ let Some(stripe) = &self.stripe else { return Ok("payments are not set up".to_owned()) };
153+ #[derive(Deserialize)]
154+ struct Mark {
155+ through: i64,
156+ }
157+ #[derive(Deserialize)]
158+ struct Page {
159+ data: Vec<Value>,
160+ has_more: bool,
161+ }
162+ let mode = self.mode();
163+ let now = (now_ms() / 1000) as i64;
164+ // A claim whose handler died never finished: let it be handled again.
165+ self.db
166+ .prepare("DELETE FROM stripe_events WHERE outcome = 'handling' AND received_at < ?")
167+ .bind(&[rfc3339(now_ms().saturating_sub(STALE_CLAIM_MS)).into()])?
168+ .run()
169+ .await?;
170+ let previous = self
171+ .db
172+ .prepare("SELECT through FROM stripe_sync WHERE mode = ?")
173+ .bind(&[mode.into()])?
174+ .first::<Mark>(None)
175+ .await?
176+ .map_or(now - FIRST_LOOK_SECONDS, |mark| mark.through);
177+ let since = previous - OVERLAP_SECONDS;
178+
179+ // Stripe lists newest first; read back to `since`.
180+ let mut listed: Vec<Value> = Vec::new();
181+ let mut complete = false;
182+ let mut after: Option<String> = None;
183+ for _ in 0..MAX_PAGES {
184+ let mut query: Vec<(&str, String)> = vec![("limit", "100".to_owned()), ("created[gte]", since.to_string())];
185+ query.extend(EVENTS.iter().map(|event| ("types[]", (*event).to_owned())));
186+ if let Some(after) = &after {
187+ query.push(("starting_after", after.clone()));
188+ }
189+ let page: Page = stripe.get(&format!("/events?{}", form(&query))).await?;
190+ after = page.data.last().and_then(|event| event["id"].as_str()).map(str::to_owned);
191+ listed.extend(page.data);
192+ if !page.has_more || after.is_none() {
193+ complete = true;
194+ break;
195+ }
196+ }
197+ if !complete {
198+ worker::console_error!("Stripe listed more than {} events since {since}; replaying the newest", MAX_PAGES * 100);
199+ }
200+
201+ // Oldest first, as they happened.
202+ listed.sort_by_key(|event| event["created"].as_i64().unwrap_or(0));
203+ let mut handled = 0;
204+ let mut first_failed = None;
205+ for event in &listed {
206+ match self.process_event(event).await {
207+ Ok(true) => handled += 1,
208+ Ok(false) => {}
209+ Err(error) => {
210+ let id = event["id"].as_str().unwrap_or_default();
211+ worker::console_error!("replaying Stripe event {id} failed: {error}");
212+ first_failed = first_failed.or(event["created"].as_i64());
213+ }
214+ }
215+ }
216+ let through = next_through(previous, now, first_failed, complete);
217+ self.db
218+ .prepare(
219+ "INSERT INTO stripe_sync (mode, through, checked_at) VALUES (?1, ?2, ?3)
220+ ON CONFLICT (mode) DO UPDATE SET through = ?2, checked_at = ?3",
221+ )
222+ .bind(&[mode.into(), (through as f64).into(), rfc3339(now_ms()).into()])?
223+ .run()
224+ .await?;
225+ Ok(format!("{} listed, {handled} not seen before and handled", listed.len()))
226+ }
227+
228+ /// Keeps the registered endpoint listening to billing's events and
229+ /// enabled. Nothing when no endpoint is registered for this mode.
230+ pub(crate) async fn keep_endpoint(&self) -> Result<String> {
231+ let (Some(stripe), Some(webhook)) = (&self.stripe, self.webhook_row().await?) else {
232+ return Ok("no endpoint registered".to_owned());
233+ };
234+ #[derive(Deserialize)]
235+ struct Endpoint {
236+ status: String,
237+ enabled_events: Vec<String>,
238+ }
239+ let path = format!("/webhook_endpoints/{}", webhook.endpoint_id);
240+ let endpoint: Endpoint = stripe.get(&path).await?;
241+ let mut fields: Vec<(String, String)> = Vec::new();
242+ let mut changes = Vec::new();
243+ if endpoint.status != "enabled" {
244+ fields.push(("disabled".to_owned(), "false".to_owned()));
245+ changes.push("enabled again".to_owned());
246+ }
247+ let mut wanted: Vec<&str> = EVENTS.to_vec();
248+ let mut has: Vec<&str> = endpoint.enabled_events.iter().map(String::as_str).collect();
249+ wanted.sort_unstable();
250+ has.sort_unstable();
251+ if wanted != has {
252+ fields.extend(EVENTS.iter().enumerate().map(|(i, event)| (format!("enabled_events[{i}]"), (*event).to_owned())));
253+ changes.push(format!("events set to billing's {}", EVENTS.len()));
254+ }
255+ if changes.is_empty() {
256+ return Ok("endpoint as it should be".to_owned());
257+ }
258+ let fields: Vec<(&str, String)> = fields.iter().map(|(name, value)| (name.as_str(), value.clone())).collect();
259+ let _: Value = stripe.post(&path, &fields).await?;
260+ self.db
261+ .prepare("UPDATE stripe_webhooks SET events = ? WHERE mode = ?")
262+ .bind(&[EVENTS.join(",").into(), self.mode().into()])?
263+ .run()
264+ .await?;
265+ let done = changes.join(", ");
266+ self.audit("stripe", "webhook", &format!("Endpoint {}: {done}", webhook.endpoint_id), "billing").await?;
267+ Ok(done)
268+ }
269+
270+ /// Reads again the saved cards and plans not read in a while, a few a
271+ /// day, in case an event never came.
272+ pub(crate) async fn refresh_from_stripe(&self) -> Result<String> {
273+ if self.stripe.is_none() {
274+ return Ok("payments are not set up".to_owned());
275+ }
276+ #[derive(Deserialize)]
277+ struct Customer {
278+ customer_id: String,
279+ }
280+ #[derive(Deserialize)]
281+ struct Plan {
282+ subscription_id: String,
283+ }
284+ let card_cutoff = rfc3339(now_ms().saturating_sub(CARD_FRESH_DAYS * 24 * 60 * 60 * 1000));
285+ let customers = self
286+ .db
287+ .prepare(
288+ "SELECT DISTINCT customer_id FROM accounts
289+ WHERE customer_id IS NOT NULL AND card_synced_at IS NOT NULL AND card_synced_at < ?1
290+ ORDER BY card_synced_at LIMIT ?2",
291+ )
292+ .bind(&[card_cutoff.into(), REFRESH_BATCH.into()])?
293+ .all()
294+ .await?
295+ .results::<Customer>()?;
296+ let plan_cutoff = rfc3339(now_ms().saturating_sub(24 * 60 * 60 * 1000));
297+ let plans = self
298+ .db
299+ .prepare(
300+ "SELECT subscription_id FROM subscriptions WHERE status <> 'canceled' AND updated_at < ?1
301+ ORDER BY updated_at LIMIT ?2",
302+ )
303+ .bind(&[plan_cutoff.into(), REFRESH_BATCH.into()])?
304+ .all()
305+ .await?
306+ .results::<Plan>()?;
307+ let mut failed = 0;
308+ for customer in &customers {
309+ if let Err(error) = self.fetch_card(&customer.customer_id).await {
310+ worker::console_error!("refreshing {}'s card failed: {error}", customer.customer_id);
311+ failed += 1;
312+ }
313+ }
314+ for plan in &plans {
315+ if let Err(error) = self.settle_subscription(&plan.subscription_id).await {
316+ worker::console_error!("refreshing plan {} failed: {error}", plan.subscription_id);
317+ failed += 1;
318+ }
319+ }
320+ Ok(format!("{} cards and {} plans read again, {failed} failed", customers.len(), plans.len()))
321+ }
322+}
323+
324+fn card_of(row: &CardRow) -> Option<Card> {
325+ Some(Card {
326+ brand: row.card_brand.clone()?,
327+ last4: row.card_last4.clone()?,
328+ exp_month: row.card_exp_month?,
329+ exp_year: row.card_exp_year?,
330+ })
331+}
332+
333+#[cfg(test)]
334+mod tests {
335+ use super::*;
336+
337+ #[test]
338+ fn the_replay_moves_on_only_past_what_it_handled() {
339+ let now = 1_791_000_000;
340+ // All handled: just behind the clock.
341+ assert_eq!(next_through(now - 900, now, None, true), now - SETTLE_SECONDS);
342+ // Never backwards.
343+ assert_eq!(next_through(now, now, None, true), now);
344+ // A failure: back to it, so it is listed again.
345+ assert_eq!(next_through(now - 900, now, Some(now - 600), true), now - 600);
346+ assert_eq!(next_through(now - 900, now, Some(now - 1200), true), now - 1200);
347+ // The list not read to its end: stay.
348+ assert_eq!(next_through(now - 900, now, None, false), now - 900);
349+ }
350+
351+ #[test]
352+ fn a_card_needs_every_part() {
353+ let row = |brand: Option<&str>| CardRow {
354+ customer_id: Some("cus_1".into()),
355+ card_brand: brand.map(str::to_owned),
356+ card_last4: Some("4242".into()),
357+ card_exp_month: Some(12),
358+ card_exp_year: Some(2030),
359+ card_synced_at: Some("2026-10-06T00:00:00Z".into()),
360+ };
361+ assert_eq!(card_of(&row(Some("visa"))).map(|c| c.last4), Some("4242".into()));
362+ assert!(card_of(&row(None)).is_none());
363+ }
364+}
+47−18
5555 "charge.refunded",
5656 "charge.dispute.created",
5757 "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",
5865 ];
5966
6067 /// How old a signed event may be, so a captured one cannot be replayed.
8895 }
8996
9097 #[derive(Deserialize)]
91−struct WebhookRow {
92− endpoint_id: String,
98+pub(crate) struct WebhookRow {
99+ pub(crate) endpoint_id: String,
93100 secret: String,
94101 url: String,
95− events: String,
102+ pub(crate) events: String,
96103 created_by: String,
97104 created_at: String,
98105 }
122129 }
123130
124131 impl Billing {
125− fn mode(&self) -> &'static str {
132+ pub(crate) fn mode(&self) -> &'static str {
126133 match &self.stripe {
127134 None => "off",
128135 Some(stripe) if stripe.live() => "live",
130137 }
131138 }
132139
133− async fn webhook_row(&self) -> Result<Option<WebhookRow>> {
140+ pub(crate) async fn webhook_row(&self) -> Result<Option<WebhookRow>> {
134141 self.db
135142 .prepare("SELECT * FROM stripe_webhooks WHERE mode = ?")
136143 .bind(&[self.mode().into()])?
142149
143150 pub(crate) async fn admin_stripe(&self, a: AdminStripeArgs) -> Result<StripeStatus> {
144151 let mut error = None;
145− if a.setup {
146− if let Err(e) = self.register_webhook(a.by.as_deref().unwrap_or("sudo")).await {
152+ if a.setup
153+ && let Err(e) = self.register_webhook(a.by.as_deref().unwrap_or("sudo")).await {
147154 error = Some(e.to_string());
148155 }
149− }
150156 let webhook = self.webhook_row().await?.map(|row| StripeWebhook {
151157 url: row.url,
152158 endpoint_id: row.endpoint_id,
234240 return Ok(Outcome::fail(FailureCode::Forbidden, "The signature does not match."));
235241 }
236242 let event: Value = serde_json::from_str(&a.payload).map_err(|e| worker::Error::RustError(e.to_string()))?;
237− let id = event["id"].as_str().unwrap_or_default().to_owned();
238− let kind = event["type"].as_str().unwrap_or_default().to_owned();
239− if id.is_empty() {
243+ if event["id"].as_str().unwrap_or_default().is_empty() {
240244 return Ok(Outcome::fail(FailureCode::Invalid, "Not an event."));
241245 }
246+ Ok(Outcome::Ok(self.process_event(&event).await?))
247+ }
248+
249+ /// Handles a Stripe event once, however often it arrives: by webhook,
250+ /// or again from the event list (stripe_sync.rs). False when it was
251+ /// seen before.
252+ pub(crate) async fn process_event(&self, event: &Value) -> Result<bool> {
253+ let id = event["id"].as_str().unwrap_or_default().to_owned();
254+ let kind = event["type"].as_str().unwrap_or_default().to_owned();
242255 // Once each: the first to record it handles it.
243256 let claimed = self
244257 .db
247260 .first::<Value>(None)
248261 .await?;
249262 if claimed.is_none() {
250− return Ok(Outcome::Ok(false));
263+ return Ok(false);
251264 }
252− let object = &event["data"]["object"];
253− let outcome = match self.handle(&kind, object).await {
265+ let outcome = match self.handle(&kind, event).await {
254266 Ok(outcome) => outcome,
255267 Err(error) => {
256268 // Let Stripe send it again: forget it was seen.
263275 .bind(&[outcome.as_str().into(), id.as_str().into()])?
264276 .run()
265277 .await?;
266− Ok(Outcome::Ok(true))
278+ Ok(true)
267279 }
268280
269− async fn handle(&self, kind: &str, object: &Value) -> Result<String> {
281+ async fn handle(&self, kind: &str, event: &Value) -> Result<String> {
282+ let object = &event["data"]["object"];
270283 let text = |key: &str| object[key].as_str().unwrap_or_default().to_owned();
284+ // The customer whose saved card may have changed. A detached card
285+ // has no customer any more; the event says whose it was.
286+ let card_owner = match kind {
287+ "customer.updated" => object["id"].as_str(),
288+ "payment_method.detached" => event["data"]["previous_attributes"]["customer"].as_str(),
289+ _ => object["customer"].as_str(),
290+ };
271291 Ok(match kind {
272292 "checkout.session.completed" | "checkout.session.async_payment_succeeded" => self.settle_checkout(&text("id")).await?,
273293 "customer.subscription.updated" | "customer.subscription.deleted" => {
300320 "charge.refunded" => self.refunded(object).await?,
301321 "charge.dispute.created" => self.disputed(object, true).await?,
302322 "charge.dispute.closed" => self.disputed(object, object["status"].as_str() == Some("lost")).await?,
323+ "customer.updated"
324+ | "payment_method.attached"
325+ | "payment_method.detached"
326+ | "payment_method.updated"
327+ | "payment_method.automatically_updated"
328+ | "setup_intent.succeeded" => match card_owner {
329+ Some(customer) => self.sync_card_of(customer).await?,
330+ None => "ignored: no customer".to_owned(),
331+ },
303332 _ => "ignored".to_owned(),
304333 })
305334 }
404433 }
405434
406435 /// A plan that changed at Stripe: renewed, failed, canceled.
407− async fn settle_subscription(&self, subscription_id: &str) -> Result<String> {
436+ pub(crate) async fn settle_subscription(&self, subscription_id: &str) -> Result<String> {
408437 #[derive(Deserialize)]
409438 struct Plan {
410439 workspace: String,
429458 }
430459
431460 /// The workspace a Stripe customer belongs to.
432− async fn workspace_of_customer(&self, customer: &str) -> Result<Option<String>> {
461+ pub(crate) async fn workspace_of_customer(&self, customer: &str) -> Result<Option<String>> {
433462 #[derive(Deserialize)]
434463 struct Row {
435464 workspace: String,