g1t/services/billing/src/stripe_sync.rs

364 lines14,910 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.
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
16use crate::Billing;
17use crate::stripe::form;
18use crate::webhooks::EVENTS;
19use g1t_contracts::billing::Card;
20use g1t_contracts::time::rfc3339;
21use g1t_kit::now_ms;
22use serde::Deserialize;
23use serde_json::Value;
24use 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.
28const 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.
31const SETTLE_SECONDS: i64 = 5 * 60;
32/// Where the first replay starts.
33const FIRST_LOOK_SECONDS: i64 = 3 * 24 * 60 * 60;
34/// Pages of 100 events one replay reads at most.
35const MAX_PAGES: usize = 10;
36/// A claim this old that never finished: its handler died, so the event
37/// may be handled again.
38const STALE_CLAIM_MS: u64 = 10 * 60 * 1000;
39/// How long a saved card is believed before the weekly pass asks again.
40const CARD_FRESH_DAYS: u64 = 7;
41/// Accounts or plans the daily pass refreshes at most.
42const REFRESH_BATCH: u32 = 25;
43
44#[derive(Deserialize)]
45struct 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.
57pub(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
65impl 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
324fn 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)]
334mod 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}