Skip to content

g1t/services/billing/src/stripe_sync.rs

353 lines14,836 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept1//! 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.
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard12//! - The destination, kept: once a day it is enabled again if Stripe turned
13//! it off after failures, and given any event billing handles that it
14//! does not send. Its signing secret, `STRIPE_WEBHOOK_SECRET`, is not
15//! touched.
Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept16
17use crate::Billing;
18use crate::stripe::form;
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard19use crate::webhooks::{EVENTS, WEBHOOK_URL, missing_events};
Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept20use g1t_contracts::billing::Card;
21use g1t_contracts::time::rfc3339;
22use g1t_kit::now_ms;
23use serde::Deserialize;
24use serde_json::Value;
25use worker::Result;
26
27/// How far back each replay starts before where the last one got to, for
28/// events Stripe lists a little late.
29const OVERLAP_SECONDS: i64 = 60 * 60;
30/// How recent an event may be and still be listed later: the replay's
31/// mark stays this far behind the clock.
32const SETTLE_SECONDS: i64 = 5 * 60;
33/// Where the first replay starts.
34const FIRST_LOOK_SECONDS: i64 = 3 * 24 * 60 * 60;
35/// Pages of 100 events one replay reads at most.
36const MAX_PAGES: usize = 10;
37/// A claim this old that never finished: its handler died, so the event
38/// may be handled again.
39const STALE_CLAIM_MS: u64 = 10 * 60 * 1000;
40/// How long a saved card is believed before the weekly pass asks again.
41const CARD_FRESH_DAYS: u64 = 7;
42/// Accounts or plans the daily pass refreshes at most.
43const REFRESH_BATCH: u32 = 25;
44
45#[derive(Deserialize)]
46struct CardRow {
47 customer_id: Option<String>,
48 card_brand: Option<String>,
49 card_last4: Option<String>,
50 card_exp_month: Option<u32>,
51 card_exp_year: Option<u32>,
52 card_synced_at: Option<String>,
53}
54
55/// Where the replay's mark goes after a run: the first event that failed
56/// (so it is listed again), the old mark when the list was not read to its
57/// end, else just behind the clock.
58pub(crate) fn next_through(previous: i64, now: i64, first_failed: Option<i64>, complete: bool) -> i64 {
59 match first_failed {
60 Some(created) => created.min(previous.max(created)),
61 None if !complete => previous,
62 None => previous.max(now - SETTLE_SECONDS),
63 }
64}
65
66impl Billing {
67 /// The workspace's saved card as last synced. Stripe is asked only when
68 /// it never was, or g1t changed the card since; a Stripe that does not
69 /// answer then shows no card, as before.
70 pub(crate) async fn saved_card(&self, workspace: &str) -> Result<Option<Card>> {
71 let Some(row) = self
72 .db
73 .prepare(
74 "SELECT customer_id, card_brand, card_last4, card_exp_month, card_exp_year, card_synced_at
75 FROM accounts WHERE workspace = ?",
76 )
77 .bind(&[workspace.into()])?
78 .first::<CardRow>(None)
79 .await?
80 else {
81 return Ok(None);
82 };
83 if row.card_synced_at.is_some() {
84 return Ok(card_of(&row));
85 }
86 let Some(customer) = row.customer_id.as_deref() else { return Ok(None) };
87 if self.stripe.is_none() {
88 return Ok(None);
89 }
90 Ok(match self.fetch_card(customer).await {
91 Ok(card) => card,
92 Err(error) => {
93 worker::console_error!("reading {workspace}'s card from Stripe failed: {error}");
94 None
95 }
96 })
97 }
98
99 /// Asks Stripe for the customer's card and keeps it on every account
100 /// with that customer.
101 async fn fetch_card(&self, customer: &str) -> Result<Option<Card>> {
102 let Some(stripe) = &self.stripe else { return Ok(None) };
Usage, Billing settings and prepaid AI credit; fixes from the UX audit103 // The card invoices are charged to (the default), not merely the
104 // newest: what Stripe's billing page made the default is the one.
105 let card = stripe.default_payment_method(customer).await?.filter(|m| m.kind == "card").and_then(|m| {
106 Some(Card { brand: m.brand?, last4: m.last4?, exp_month: m.exp_month?, exp_year: m.exp_year? })
Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept107 });
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
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard228 /// Keeps the destination at billing's address enabled and sending
229 /// every event billing handles. Its signing secret is not touched.
230 /// Nothing when Stripe has no destination there.
231 pub(crate) async fn keep_endpoint(&self, by: &str) -> Result<String> {
232 let (Some(stripe), Some(destination)) = (&self.stripe, self.destination().await?) else {
233 return Ok(format!("no destination at {WEBHOOK_URL}"));
Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept234 };
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard235 let path = format!("/webhook_endpoints/{}", destination.id);
Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept236 let mut fields: Vec<(String, String)> = Vec::new();
237 let mut changes = Vec::new();
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard238 if destination.status != "enabled" {
Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept239 fields.push(("disabled".to_owned(), "false".to_owned()));
240 changes.push("enabled again".to_owned());
241 }
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard242 let missing = missing_events(&destination.enabled_events);
243 if !missing.is_empty() {
244 // Stripe replaces the list: what it sends now, plus what is missing.
245 let all = destination.enabled_events.iter().cloned().chain(missing.iter().cloned());
246 fields.extend(all.enumerate().map(|(i, event)| (format!("enabled_events[{i}]"), event)));
247 changes.push(format!("added {}", missing.join(", ")));
Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept248 }
249 if changes.is_empty() {
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard250 return Ok("destination as it should be".to_owned());
Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept251 }
252 let fields: Vec<(&str, String)> = fields.iter().map(|(name, value)| (name.as_str(), value.clone())).collect();
253 let _: Value = stripe.post(&path, &fields).await?;
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard254 let done = changes.join("; ");
255 self.audit("stripe", "webhook", &format!("Destination {}: {done}", destination.id), by).await?;
Billing keeps Stripe's view itself: the saved card on the account, missed events replayed every 15 minutes, and the endpoint kept256 Ok(done)
257 }
258
259 /// Reads again the saved cards and plans not read in a while, a few a
260 /// day, in case an event never came.
261 pub(crate) async fn refresh_from_stripe(&self) -> Result<String> {
262 if self.stripe.is_none() {
263 return Ok("payments are not set up".to_owned());
264 }
265 #[derive(Deserialize)]
266 struct Customer {
267 customer_id: String,
268 }
269 #[derive(Deserialize)]
270 struct Plan {
271 subscription_id: String,
272 }
273 let card_cutoff = rfc3339(now_ms().saturating_sub(CARD_FRESH_DAYS * 24 * 60 * 60 * 1000));
274 let customers = self
275 .db
276 .prepare(
277 "SELECT DISTINCT customer_id FROM accounts
278 WHERE customer_id IS NOT NULL AND card_synced_at IS NOT NULL AND card_synced_at < ?1
279 ORDER BY card_synced_at LIMIT ?2",
280 )
281 .bind(&[card_cutoff.into(), REFRESH_BATCH.into()])?
282 .all()
283 .await?
284 .results::<Customer>()?;
285 let plan_cutoff = rfc3339(now_ms().saturating_sub(24 * 60 * 60 * 1000));
286 let plans = self
287 .db
288 .prepare(
289 "SELECT subscription_id FROM subscriptions WHERE status <> 'canceled' AND updated_at < ?1
290 ORDER BY updated_at LIMIT ?2",
291 )
292 .bind(&[plan_cutoff.into(), REFRESH_BATCH.into()])?
293 .all()
294 .await?
295 .results::<Plan>()?;
296 let mut failed = 0;
297 for customer in &customers {
298 if let Err(error) = self.fetch_card(&customer.customer_id).await {
299 worker::console_error!("refreshing {}'s card failed: {error}", customer.customer_id);
300 failed += 1;
301 }
302 }
303 for plan in &plans {
304 if let Err(error) = self.settle_subscription(&plan.subscription_id).await {
305 worker::console_error!("refreshing plan {} failed: {error}", plan.subscription_id);
306 failed += 1;
307 }
308 }
309 Ok(format!("{} cards and {} plans read again, {failed} failed", customers.len(), plans.len()))
310 }
311}
312
313fn card_of(row: &CardRow) -> Option<Card> {
314 Some(Card {
315 brand: row.card_brand.clone()?,
316 last4: row.card_last4.clone()?,
317 exp_month: row.card_exp_month?,
318 exp_year: row.card_exp_year?,
319 })
320}
321
322#[cfg(test)]
323mod tests {
324 use super::*;
325
326 #[test]
327 fn the_replay_moves_on_only_past_what_it_handled() {
328 let now = 1_791_000_000;
329 // All handled: just behind the clock.
330 assert_eq!(next_through(now - 900, now, None, true), now - SETTLE_SECONDS);
331 // Never backwards.
332 assert_eq!(next_through(now, now, None, true), now);
333 // A failure: back to it, so it is listed again.
334 assert_eq!(next_through(now - 900, now, Some(now - 600), true), now - 600);
335 assert_eq!(next_through(now - 900, now, Some(now - 1200), true), now - 1200);
336 // The list not read to its end: stay.
337 assert_eq!(next_through(now - 900, now, None, false), now - 900);
338 }
339
340 #[test]
341 fn a_card_needs_every_part() {
342 let row = |brand: Option<&str>| CardRow {
343 customer_id: Some("cus_1".into()),
344 card_brand: brand.map(str::to_owned),
345 card_last4: Some("4242".into()),
346 card_exp_month: Some(12),
347 card_exp_year: Some(2030),
348 card_synced_at: Some("2026-10-06T00:00:00Z".into()),
349 };
350 assert_eq!(card_of(&row(Some("visa"))).map(|c| c.last4), Some("4242".into()));
351 assert!(card_of(&row(None)).is_none());
352 }
353}