g1t/services/billing/src/stripe_sync.rs

354 lines14,673 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) };
103 let card = stripe.card(customer).await?.map(|card| Card {
104 brand: card.brand,
105 last4: card.last4,
106 exp_month: card.exp_month,
107 exp_year: card.exp_year,
108 });
109 self.db
110 .prepare(
111 "UPDATE accounts SET card_brand = ?2, card_last4 = ?3, card_exp_month = ?4, card_exp_year = ?5,
112 card_synced_at = ?6 WHERE customer_id = ?1",
113 )
114 .bind(&[
115 customer.into(),
116 card.as_ref().map(|c| c.brand.clone()).into(),
117 card.as_ref().map(|c| c.last4.clone()).into(),
118 card.as_ref().map(|c| c.exp_month as f64).into(),
119 card.as_ref().map(|c| c.exp_year as f64).into(),
120 rfc3339(now_ms()).into(),
121 ])?
122 .run()
123 .await?;
124 Ok(card)
125 }
126
127 /// A card or customer event: the saved card read again. The event's
128 /// outcome, for `stripe_events`.
129 pub(crate) async fn sync_card_of(&self, customer: &str) -> Result<String> {
130 if self.workspace_of_customer(customer).await?.is_none() {
131 return Ok("ignored: not a workspace's customer".to_owned());
132 }
133 Ok(match self.fetch_card(customer).await? {
134 Some(card) => format!("card saved: {} ending {}", card.brand, card.last4),
135 None => "card saved: none".to_owned(),
136 })
137 }
138
139 /// Forgets the saved card after g1t itself changed it, so the next
140 /// read asks Stripe instead of waiting for the event.
141 pub(crate) async fn forget_card(&self, workspace: &str) -> Result<()> {
142 self.db
143 .prepare("UPDATE accounts SET card_synced_at = NULL WHERE workspace = ?")
144 .bind(&[workspace.into()])?
145 .run()
146 .await?;
147 Ok(())
148 }
149
150 /// Handles events Stripe sent that billing never saw, from its event
151 /// list. What happened, for the log.
152 pub(crate) async fn replay_events(&self) -> Result<String> {
153 let Some(stripe) = &self.stripe else { return Ok("payments are not set up".to_owned()) };
154 #[derive(Deserialize)]
155 struct Mark {
156 through: i64,
157 }
158 #[derive(Deserialize)]
159 struct Page {
160 data: Vec<Value>,
161 has_more: bool,
162 }
163 let mode = self.mode();
164 let now = (now_ms() / 1000) as i64;
165 // A claim whose handler died never finished: let it be handled again.
166 self.db
167 .prepare("DELETE FROM stripe_events WHERE outcome = 'handling' AND received_at < ?")
168 .bind(&[rfc3339(now_ms().saturating_sub(STALE_CLAIM_MS)).into()])?
169 .run()
170 .await?;
171 let previous = self
172 .db
173 .prepare("SELECT through FROM stripe_sync WHERE mode = ?")
174 .bind(&[mode.into()])?
175 .first::<Mark>(None)
176 .await?
177 .map_or(now - FIRST_LOOK_SECONDS, |mark| mark.through);
178 let since = previous - OVERLAP_SECONDS;
179
180 // Stripe lists newest first; read back to `since`.
181 let mut listed: Vec<Value> = Vec::new();
182 let mut complete = false;
183 let mut after: Option<String> = None;
184 for _ in 0..MAX_PAGES {
185 let mut query: Vec<(&str, String)> = vec![("limit", "100".to_owned()), ("created[gte]", since.to_string())];
186 query.extend(EVENTS.iter().map(|event| ("types[]", (*event).to_owned())));
187 if let Some(after) = &after {
188 query.push(("starting_after", after.clone()));
189 }
190 let page: Page = stripe.get(&format!("/events?{}", form(&query))).await?;
191 after = page.data.last().and_then(|event| event["id"].as_str()).map(str::to_owned);
192 listed.extend(page.data);
193 if !page.has_more || after.is_none() {
194 complete = true;
195 break;
196 }
197 }
198 if !complete {
199 worker::console_error!("Stripe listed more than {} events since {since}; replaying the newest", MAX_PAGES * 100);
200 }
201
202 // Oldest first, as they happened.
203 listed.sort_by_key(|event| event["created"].as_i64().unwrap_or(0));
204 let mut handled = 0;
205 let mut first_failed = None;
206 for event in &listed {
207 match self.process_event(event).await {
208 Ok(true) => handled += 1,
209 Ok(false) => {}
210 Err(error) => {
211 let id = event["id"].as_str().unwrap_or_default();
212 worker::console_error!("replaying Stripe event {id} failed: {error}");
213 first_failed = first_failed.or(event["created"].as_i64());
214 }
215 }
216 }
217 let through = next_through(previous, now, first_failed, complete);
218 self.db
219 .prepare(
220 "INSERT INTO stripe_sync (mode, through, checked_at) VALUES (?1, ?2, ?3)
221 ON CONFLICT (mode) DO UPDATE SET through = ?2, checked_at = ?3",
222 )
223 .bind(&[mode.into(), (through as f64).into(), rfc3339(now_ms()).into()])?
224 .run()
225 .await?;
226 Ok(format!("{} listed, {handled} not seen before and handled", listed.len()))
227 }
228
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard229 /// Keeps the destination at billing's address enabled and sending
230 /// every event billing handles. Its signing secret is not touched.
231 /// Nothing when Stripe has no destination there.
232 pub(crate) async fn keep_endpoint(&self, by: &str) -> Result<String> {
233 let (Some(stripe), Some(destination)) = (&self.stripe, self.destination().await?) else {
234 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 kept235 };
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard236 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 kept237 let mut fields: Vec<(String, String)> = Vec::new();
238 let mut changes = Vec::new();
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard239 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 kept240 fields.push(("disabled".to_owned(), "false".to_owned()));
241 changes.push("enabled again".to_owned());
242 }
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard243 let missing = missing_events(&destination.enabled_events);
244 if !missing.is_empty() {
245 // Stripe replaces the list: what it sends now, plus what is missing.
246 let all = destination.enabled_events.iter().cloned().chain(missing.iter().cloned());
247 fields.extend(all.enumerate().map(|(i, event)| (format!("enabled_events[{i}]"), event)));
248 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 kept249 }
250 if changes.is_empty() {
Stripe's webhook secret is a Worker secret, STRIPE_WEBHOOK_SECRET, from a destination made in Stripe's dashboard251 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 kept252 }
253 let fields: Vec<(&str, String)> = fields.iter().map(|(name, value)| (name.as_str(), value.clone())).collect();
254 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 dashboard255 let done = changes.join("; ");
256 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 kept257 Ok(done)
258 }
259
260 /// Reads again the saved cards and plans not read in a while, a few a
261 /// day, in case an event never came.
262 pub(crate) async fn refresh_from_stripe(&self) -> Result<String> {
263 if self.stripe.is_none() {
264 return Ok("payments are not set up".to_owned());
265 }
266 #[derive(Deserialize)]
267 struct Customer {
268 customer_id: String,
269 }
270 #[derive(Deserialize)]
271 struct Plan {
272 subscription_id: String,
273 }
274 let card_cutoff = rfc3339(now_ms().saturating_sub(CARD_FRESH_DAYS * 24 * 60 * 60 * 1000));
275 let customers = self
276 .db
277 .prepare(
278 "SELECT DISTINCT customer_id FROM accounts
279 WHERE customer_id IS NOT NULL AND card_synced_at IS NOT NULL AND card_synced_at < ?1
280 ORDER BY card_synced_at LIMIT ?2",
281 )
282 .bind(&[card_cutoff.into(), REFRESH_BATCH.into()])?
283 .all()
284 .await?
285 .results::<Customer>()?;
286 let plan_cutoff = rfc3339(now_ms().saturating_sub(24 * 60 * 60 * 1000));
287 let plans = self
288 .db
289 .prepare(
290 "SELECT subscription_id FROM subscriptions WHERE status <> 'canceled' AND updated_at < ?1
291 ORDER BY updated_at LIMIT ?2",
292 )
293 .bind(&[plan_cutoff.into(), REFRESH_BATCH.into()])?
294 .all()
295 .await?
296 .results::<Plan>()?;
297 let mut failed = 0;
298 for customer in &customers {
299 if let Err(error) = self.fetch_card(&customer.customer_id).await {
300 worker::console_error!("refreshing {}'s card failed: {error}", customer.customer_id);
301 failed += 1;
302 }
303 }
304 for plan in &plans {
305 if let Err(error) = self.settle_subscription(&plan.subscription_id).await {
306 worker::console_error!("refreshing plan {} failed: {error}", plan.subscription_id);
307 failed += 1;
308 }
309 }
310 Ok(format!("{} cards and {} plans read again, {failed} failed", customers.len(), plans.len()))
311 }
312}
313
314fn card_of(row: &CardRow) -> Option<Card> {
315 Some(Card {
316 brand: row.card_brand.clone()?,
317 last4: row.card_last4.clone()?,
318 exp_month: row.card_exp_month?,
319 exp_year: row.card_exp_year?,
320 })
321}
322
323#[cfg(test)]
324mod tests {
325 use super::*;
326
327 #[test]
328 fn the_replay_moves_on_only_past_what_it_handled() {
329 let now = 1_791_000_000;
330 // All handled: just behind the clock.
331 assert_eq!(next_through(now - 900, now, None, true), now - SETTLE_SECONDS);
332 // Never backwards.
333 assert_eq!(next_through(now, now, None, true), now);
334 // A failure: back to it, so it is listed again.
335 assert_eq!(next_through(now - 900, now, Some(now - 600), true), now - 600);
336 assert_eq!(next_through(now - 900, now, Some(now - 1200), true), now - 1200);
337 // The list not read to its end: stay.
338 assert_eq!(next_through(now - 900, now, None, false), now - 900);
339 }
340
341 #[test]
342 fn a_card_needs_every_part() {
343 let row = |brand: Option<&str>| CardRow {
344 customer_id: Some("cus_1".into()),
345 card_brand: brand.map(str::to_owned),
346 card_last4: Some("4242".into()),
347 card_exp_month: Some(12),
348 card_exp_year: Some(2030),
349 card_synced_at: Some("2026-10-06T00:00:00Z".into()),
350 };
351 assert_eq!(card_of(&row(Some("visa"))).map(|c| c.last4), Some("4242".into()));
352 assert!(card_of(&row(None)).is_none());
353 }
354}