From 049603bf8e0cf08faf7583241f3be99874607654 Mon Sep 17 00:00:00 2001 From: Joseph Shearer Date: Wed, 1 Apr 2026 16:32:59 -0400 Subject: [PATCH 1/6] billing: add `--dry-run` mode The invoice generator's trial-run workflow ran against a sandbox Stripe account, but produced inaccurate results because the sandbox lacked livemode invoice state (manual bills with existing open/paid invoices appeared as phantom creates). `--dry-run` runs against livemode Stripe in read-only mode, showing what would happen without creating invoices, customers, or modifying anything. Structural changes: * Split `upsert_invoice` into `classify` (read-only: validation, Stripe searches) and `execute` (writes: customer/invoice creation, line items, verification). `--dry-run` stops after classify. * Decompose `get_or_create_customer_for_tenant` into `find_customer` (read-only search) and `ensure_customer_for_invoicing` (find-or-create + email backfill). * Reorder classify checks so cheap local validations (FreeTier, FutureTrialStart, LessThanMinimum) run before any Stripe API calls. * Multi-month manual bills that already have an `open`, `paid`, `void`, or `uncollectible` invoice in Stripe are now classified as `AlreadyProcessed` instead of erroring. These are expected when date-range-overlapping manual bills were invoiced in a previous billing run. * Per-tenant summary output annotates manual bills with their date range (`[manual: 2026-01-01 - 2026-06-30]`). * Dry-run with `--clean-up` previews which stale draft invoices would be deleted. `--recreate-finalized` logs which invoices would be deleted and recreated. --- crates/billing-integrations/src/publish.rs | 720 +++++++++++++-------- 1 file changed, 464 insertions(+), 256 deletions(-) diff --git a/crates/billing-integrations/src/publish.rs b/crates/billing-integrations/src/publish.rs index 9e46a32af36..abb5f5b9e24 100644 --- a/crates/billing-integrations/src/publish.rs +++ b/crates/billing-integrations/src/publish.rs @@ -58,6 +58,10 @@ pub struct PublishInvoice { /// Clean up dangling invoices that are not in the database #[clap(long, default_value_t = false)] pub clean_up: bool, + /// Run in read-only mode: classify all invoices and report what would + /// happen, without creating or modifying anything in Stripe. + #[clap(long, default_value_t = false)] + pub dry_run: bool, } fn parse_date(arg: &str) -> Result { @@ -90,20 +94,32 @@ enum InvoiceResult { FutureTrialStart, NoDataMoved, NoFullPipeline, + AlreadyProcessed, Error, } impl InvoiceResult { - pub fn message(&self) -> String { + pub fn message(&self, dry_run: bool) -> String { match self { InvoiceResult::Created(provider) => { + let verb = if dry_run { + "Would publish" + } else { + "Published" + }; if provider == &PaymentProvider::Stripe { - "Published new invoice".to_string() + format!("{verb} new invoice") } else { - format!("Published new invoice for tenant using {provider:?} provider") + format!("{verb} new invoice for tenant using {provider:?} provider") + } + } + InvoiceResult::Updated => { + if dry_run { + "Would update existing invoice".to_string() + } else { + "Updated existing invoice".to_string() } } - InvoiceResult::Updated => "Updated existing invoice".to_string(), InvoiceResult::LessThanMinimum => { "Skipping invoice for less than the minimum chargable amount ($0.50)".to_string() } @@ -117,11 +133,31 @@ impl InvoiceResult { InvoiceResult::NoFullPipeline => { "Skipping invoice for tenant without an active pipeline".to_string() } + InvoiceResult::AlreadyProcessed => { + "Skipping invoice already processed in a previous billing run".to_string() + } InvoiceResult::Error => "Error publishing invoices".to_string(), } } } +/// The outcome of the classify phase: what action should be taken for this invoice. +enum InvoiceAction { + /// Invoice should not be created. Carries the skip reason and the + /// customer (if found) for potential clean-up of stale drafts. + Skip { + result: InvoiceResult, + customer: Option, + }, + /// Create a new invoice. `replace` is set when --recreate-finalized + /// requires deleting an existing invoice first. + Create { replace: Option }, + /// Update an existing draft invoice's line items. + Update { + existing_invoice_id: stripe::InvoiceId, + }, +} + #[derive(Serialize, Deserialize, Debug, Clone, sqlx::FromRow)] struct Invoice { subtotal: i64, @@ -167,93 +203,215 @@ impl Invoice { Ok(invoice_search.into_iter().next()) } - #[tracing::instrument(skip(self, client, db_client), fields(tenant=self.billed_prefix, invoice_type=format!("{:?}",self.invoice_type), subtotal=format!("${:.2}", self.subtotal as f64 / 100.0)))] - async fn upsert_invoice( + /// Read-only classification: determines what action should be taken for this + /// invoice without making any writes to Stripe. + #[tracing::instrument(skip(self, client), fields(tenant=self.billed_prefix, invoice_type=format!("{:?}",self.invoice_type), subtotal=format!("${:.2}", self.subtotal as f64 / 100.0)))] + async fn classify( &self, client: &stripe::Client, - db_client: &Pool, recreate_finalized: bool, - mode: ChargeType, - ) -> anyhow::Result { + ) -> anyhow::Result { + // --- Phase 1: Cheap local checks (no Stripe calls) --- + match (&self.invoice_type, &self.extra) { (InvoiceType::Preview, _) => { bail!("Should not create Stripe invoices for preview invoices") } - (InvoiceType::Final, Some(extra)) => { - // If we have a payment method, don't skip the invoice - // If `has_payment_method` is Some, then there is a stripe customer to check - let validated_has_payment_method = - if let Some(has_payment_method) = self.has_payment_method { - // The Stripe capture in the database has been known to be unreliable. - // Let's double-check with Stripe to make sure it agrees that we really - // do not have a payment method set. - let real_default_payment_method = get_or_create_customer_for_tenant( - client, - db_client, - self.billed_prefix.to_owned(), - false, // If there's no customer, there's no way there can be a payment method - ) - .await? - .and_then(|customer| customer.invoice_settings) - .and_then(|i| i.default_payment_method); - - if has_payment_method != real_default_payment_method.is_some() { - tracing::warn!( - ?has_payment_method, - stripe_payment_method = real_default_payment_method.is_some(), - "Inconsistent payment method state" - ); - } - - real_default_payment_method.is_some() - } else { - false - }; - - let unwrapped_extra = extra.clone().0.expect( - "This is just a sqlx quirk, if the outer Option is Some then this will be Some", - ); - - if !validated_has_payment_method { - if unwrapped_extra.processed_data_gb.unwrap_or_default() == 0.0 - && !matches!(&self.invoice_type, InvoiceType::Manual) - { - return Ok(InvoiceResult::NoDataMoved); - } - - if !self.has_full_pipeline && !matches!(&self.invoice_type, InvoiceType::Manual) - { - return Ok(InvoiceResult::NoFullPipeline); - } - } - } (InvoiceType::Final, None) => { bail!("Invoice should have extra") } _ => {} }; - // An invoice should be generated in Stripe if the tenant is on a paid plan, which means: - // * The tenant has a free trial start date - // * The tenant's free trial start date is before the invoice period's end date if let InvoiceType::Final = self.invoice_type { match self.tenant_trial_start { Some(trial_start) if self.date_end < trial_start => { - return Ok(InvoiceResult::FutureTrialStart); + return Ok(InvoiceAction::Skip { + result: InvoiceResult::FutureTrialStart, + customer: None, + }); } None => { - return Ok(InvoiceResult::FreeTier); + return Ok(InvoiceAction::Skip { + result: InvoiceResult::FreeTier, + customer: None, + }); } _ => {} } } - // The minimum chargable amount of USD in Stripe is $0.50. - // https://stripe.com/docs/currencies#minimum-and-maximum-charge-amounts if self.subtotal < 50 { - return Ok(InvoiceResult::LessThanMinimum); + return Ok(InvoiceAction::Skip { + result: InvoiceResult::LessThanMinimum, + customer: None, + }); } + // --- Phase 2: Stripe calls (only for invoices that survived Phase 1) --- + + // For Final invoices, verify the payment method state with Stripe. + // The DB capture has been known to be unreliable, so Stripe is the + // source of truth. If the tenant has no payment method, skip on + // NoDataMoved / NoFullPipeline. + let mut found_customer: Option> = None; + + if let (InvoiceType::Final, Some(extra)) = (&self.invoice_type, &self.extra) { + let validated_has_payment_method = + if let Some(has_payment_method) = self.has_payment_method { + let customer = find_customer(client, &self.billed_prefix).await?; + let real_has_pm = customer + .as_ref() + .and_then(|c| c.invoice_settings.as_ref()) + .and_then(|i| i.default_payment_method.as_ref()) + .is_some(); + + if has_payment_method != real_has_pm { + tracing::warn!( + ?has_payment_method, + stripe_payment_method = real_has_pm, + "Inconsistent payment method state" + ); + } + + found_customer = Some(customer); + real_has_pm + } else { + false + }; + + if !validated_has_payment_method { + let unwrapped_extra = extra.clone().0.expect( + "This is just a sqlx quirk, if the outer Option is Some then this will be Some", + ); + + if unwrapped_extra.processed_data_gb.unwrap_or_default() == 0.0 { + return Ok(InvoiceAction::Skip { + result: InvoiceResult::NoDataMoved, + customer: found_customer.flatten(), + }); + } + + if !self.has_full_pipeline { + return Ok(InvoiceAction::Skip { + result: InvoiceResult::NoFullPipeline, + customer: found_customer.flatten(), + }); + } + } + } + + // Look up customer (reuse if already fetched during payment method validation) + let customer = match found_customer { + Some(c) => c, + None => find_customer(client, &self.billed_prefix).await?, + }; + + let customer = match customer { + Some(c) => c, + // No customer in Stripe means no existing invoice is possible + None => return Ok(InvoiceAction::Create { replace: None }), + }; + + let customer_id = customer.id.to_string(); + + // Search for an existing invoice in Stripe + if let Some(invoice) = self + .get_stripe_invoice(client, customer_id.as_str()) + .await? + { + match invoice.status { + Some(stripe::InvoiceStatus::Open | stripe::InvoiceStatus::Draft) + if recreate_finalized => + { + Ok(InvoiceAction::Create { + replace: Some(invoice.id), + }) + } + Some(stripe::InvoiceStatus::Draft) => { + tracing::debug!( + "Found existing draft invoice {id}", + id = invoice.id.to_string() + ); + Ok(InvoiceAction::Update { + existing_invoice_id: invoice.id, + }) + } + Some(stripe::InvoiceStatus::Open) + if matches!(self.invoice_type, InvoiceType::Manual) => + { + tracing::debug!( + "Manual invoice {id} already open, skipping", + id = invoice.id.to_string() + ); + Ok(InvoiceAction::Skip { + result: InvoiceResult::AlreadyProcessed, + customer: Some(customer), + }) + } + Some(stripe::InvoiceStatus::Open) => { + bail!( + "Found open invoice {id}. Pass --recreate-finalized to delete and recreate this invoice.", + id = invoice.id.to_string() + ) + } + Some( + status @ (stripe::InvoiceStatus::Paid + | stripe::InvoiceStatus::Void + | stripe::InvoiceStatus::Uncollectible), + ) if matches!(self.invoice_type, InvoiceType::Manual) => { + tracing::debug!( + "Manual invoice {id} already in state {status}, skipping", + id = invoice.id.to_string(), + status = status + ); + Ok(InvoiceAction::Skip { + result: InvoiceResult::AlreadyProcessed, + customer: Some(customer), + }) + } + Some(status) => { + bail!( + "Found invoice {id} in unsupported state {status}, skipping.", + id = invoice.id.to_string(), + status = status + ); + } + None => { + bail!( + "Unexpected missing status from invoice {id}", + id = invoice.id.to_string() + ); + } + } + } else { + Ok(InvoiceAction::Create { replace: None }) + } + } + + /// Execute the classified action: performs all Stripe writes (customer creation, + /// invoice creation/update, line item management, verification). + #[tracing::instrument(skip(self, client, db_client, action), fields(tenant=self.billed_prefix, invoice_type=format!("{:?}",self.invoice_type), subtotal=format!("${:.2}", self.subtotal as f64 / 100.0)))] + async fn execute( + &self, + client: &stripe::Client, + db_client: &Pool, + action: InvoiceAction, + mode: ChargeType, + ) -> anyhow::Result { + let (is_update, replace, existing_invoice_id) = match action { + InvoiceAction::Skip { result, .. } => return Ok(result), + InvoiceAction::Create { replace, .. } => (false, replace, None), + InvoiceAction::Update { + existing_invoice_id, + .. + } => (true, None, Some(existing_invoice_id)), + }; + + // Ensure customer exists and has an email (required for invoicing) + let customer = + ensure_customer_for_invoicing(client, db_client, &self.billed_prefix).await?; + // Anything before 12:00:00 renders as the previous day in Stripe let date_start_secs = self .date_start @@ -278,114 +436,89 @@ impl Invoice { let date_start_repr = self.date_start.format("%F").to_string(); let date_end_repr = self.date_end.format("%F").to_string(); - let customer = get_or_create_customer_for_tenant( - client, - db_client, - self.billed_prefix.to_owned(), - true, - ) - .await? - .expect("Should never return None"); - let customer_id = customer.id.to_string(); - - let maybe_invoice = if let Some(invoice) = self - .get_stripe_invoice(&client, customer_id.as_str()) - .await? - { - match invoice.status { - Some(state @ (stripe::InvoiceStatus::Open | stripe::InvoiceStatus::Draft)) - if recreate_finalized => - { + // Delete existing invoice if --recreate-finalized was used + if let Some(ref replace_id) = replace { + // Re-verify the invoice status before deleting (guard against race conditions) + let existing = stripe::Invoice::retrieve(client, replace_id, &[]).await?; + match existing.status { + Some(state @ (stripe::InvoiceStatus::Open | stripe::InvoiceStatus::Draft)) => { tracing::warn!( - "Found invoice {id} in state {state} deleting and recreating", - id = invoice.id.to_string(), + "Found invoice {id} in state {state}, deleting and recreating", + id = replace_id.to_string(), state = state ); - stripe::Invoice::delete(client, &invoice.id).await?; - None - } - Some(stripe::InvoiceStatus::Draft) => { - tracing::debug!( - "Updating existing invoice {id}", - id = invoice.id.to_string() - ); - Some(invoice) - } - Some(stripe::InvoiceStatus::Open) => { - bail!( - "Found open invoice {id}. Pass --recreate-finalized to delete and recreate this invoice.", - id = invoice.id.to_string() - ) + stripe::Invoice::delete(client, replace_id).await?; } Some(status) => { bail!( - "Found invoice {id} in unsupported state {status}, skipping.", - id = invoice.id.to_string(), + "Invoice {id} changed to state {status} since classification, cannot delete.", + id = replace_id.to_string(), status = status ); } None => { bail!( "Unexpected missing status from invoice {id}", - id = invoice.id.to_string() + id = replace_id.to_string() ); } } - } else { - None - }; + } - let invoice = match maybe_invoice.clone() { - Some(inv) => inv, - None => { - let invoice = stripe::Invoice::create( - client, - stripe::CreateInvoice { - customer: Some(customer.id.to_owned()), - // Stripe timestamps are measured in _seconds_ since epoch - // Due date must be in the future. Bill net-30, so 30 days from today - due_date: match mode { - ChargeType::SendInvoice => Some((Utc::now() + Duration::days(30)).timestamp()), - ChargeType::AutoCharge => None - }, - description: Some( - format!( - "Your Flow bill for the billing period between {date_start_human} - {date_end_human}. Tenant: {tenant}", - tenant=self.billed_prefix.to_owned() - ) - .as_str(), - ), - collection_method: Some(match mode { - ChargeType::AutoCharge => stripe::CollectionMethod::ChargeAutomatically, - ChargeType::SendInvoice => stripe::CollectionMethod::SendInvoice, - }), - auto_advance: Some(false), - custom_fields: Some(vec![ - stripe::CreateInvoiceCustomFields { - name: "Billing Period Start".to_string(), - value: date_start_human.to_owned(), - }, - stripe::CreateInvoiceCustomFields { - name: "Billing Period End".to_string(), - value: date_end_human.to_owned(), - }, - ]), - metadata: Some( - InvoiceMetadata { - tenant: self.billed_prefix.to_owned(), - invoice_type: self.invoice_type, - period_start: date_start_repr, - period_end: date_end_repr, - } - .to_metadata_map(), - ), - ..Default::default() + // Create or reuse the invoice + let invoice = if let Some(existing_id) = existing_invoice_id { + tracing::debug!( + "Updating existing invoice {id}", + id = existing_id.to_string() + ); + stripe::Invoice::retrieve(client, &existing_id, &[]).await? + } else { + let description_text = format!( + "Your Flow bill for the billing period between {date_start_human} - {date_end_human}. Tenant: {tenant}", + tenant = self.billed_prefix + ); + let invoice = stripe::Invoice::create( + client, + stripe::CreateInvoice { + customer: Some(customer.id.to_owned()), + due_date: match mode { + ChargeType::SendInvoice => { + Some((Utc::now() + Duration::days(30)).timestamp()) + } + ChargeType::AutoCharge => None, }, - ) - .await.context("Creating a new invoice")?; - tracing::debug!("Created a new invoice {id}", id = invoice.id); - invoice - } + description: Some(description_text.as_str()), + collection_method: Some(match mode { + ChargeType::AutoCharge => stripe::CollectionMethod::ChargeAutomatically, + ChargeType::SendInvoice => stripe::CollectionMethod::SendInvoice, + }), + auto_advance: Some(false), + custom_fields: Some(vec![ + stripe::CreateInvoiceCustomFields { + name: "Billing Period Start".to_string(), + value: date_start_human.to_owned(), + }, + stripe::CreateInvoiceCustomFields { + name: "Billing Period End".to_string(), + value: date_end_human.to_owned(), + }, + ]), + metadata: Some( + InvoiceMetadata { + tenant: self.billed_prefix.to_owned(), + invoice_type: self.invoice_type, + period_start: date_start_repr, + period_end: date_end_repr, + } + .to_metadata_map(), + ), + ..Default::default() + }, + ) + .await + .context("Creating a new invoice")?; + tracing::debug!("Created a new invoice {id}", id = invoice.id); + invoice }; // Clear out line items from invoice, if there are any @@ -438,11 +571,10 @@ impl Invoice { ); } - // Let's double-check that the invoice total matches the desired total + // Re-fetch invoice and customer for fresh data (balance may have changed) let check_invoice = stripe::Invoice::retrieve(client, &invoice.id, &[]).await?; - - // Customers can have an invoice credit balance, so let's make sure we take that into account. - let credit_balance = customer.balance.unwrap_or(0); + let fresh_customer = stripe::Customer::retrieve(client, &customer.id, &[]).await?; + let credit_balance = fresh_customer.balance.unwrap_or(0); let expected = (self.subtotal + (diff.ceil() as i64) + credit_balance).max(0); @@ -454,10 +586,10 @@ impl Invoice { ) } - if maybe_invoice.is_some() { - return Ok(InvoiceResult::Updated); + if is_update { + Ok(InvoiceResult::Updated) } else { - return Ok(InvoiceResult::Created(self.payment_provider)); + Ok(InvoiceResult::Created(self.payment_provider)) } } } @@ -592,76 +724,85 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { .or_default() += 1; }); - tracing::info!( - "Processing {usage} usage-based invoices, and {manual} manually-entered invoices.", - usage = invoice_type_counter - .remove(&InvoiceType::Final) - .unwrap_or_default(), - manual = invoice_type_counter - .remove(&InvoiceType::Manual) - .unwrap_or_default(), - ); + if cmd.dry_run { + tracing::info!( + "[DRY RUN] Classifying {usage} usage-based invoices and {manual} manually-entered invoices without making any changes to Stripe.", + usage = invoice_type_counter + .remove(&InvoiceType::Final) + .unwrap_or_default(), + manual = invoice_type_counter + .remove(&InvoiceType::Manual) + .unwrap_or_default(), + ); + } else { + tracing::info!( + "Processing {usage} usage-based invoices, and {manual} manually-entered invoices.", + usage = invoice_type_counter + .remove(&InvoiceType::Final) + .unwrap_or_default(), + manual = invoice_type_counter + .remove(&InvoiceType::Manual) + .unwrap_or_default(), + ); + } let invoice_futures: Vec<_> = invoices .iter() .map(|response| { let client = stripe_client.clone(); let db_pool = db_pool.clone(); + + let annotation = match response.invoice_type { + InvoiceType::Manual => Some(format!( + "[manual: {} - {}]", + response.date_start.format("%Y-%m-%d"), + response.date_end.format("%Y-%m-%d") + )), + _ => None, + }; + async move { - let res = response - .upsert_invoice( - &client, - &db_pool, - cmd.recreate_finalized, - cmd.charge_type, - ) + let action = response + .classify(&client, cmd.recreate_finalized) .await; - match res { + + match action { Err(err) => { let formatted = format!( - "Error publishing {invoice_type:?} invoice for {tenant}", + "Error classifying {invoice_type:?} invoice for {tenant}", tenant = response.billed_prefix, invoice_type = response.invoice_type ); - Err(anyhow::anyhow!(format!( - "{}: {err:#}", - formatted, - err = err - ))) + Err(anyhow::anyhow!("{formatted}: {err:#}")) } - Ok(res) => { + Ok(InvoiceAction::Skip { result, customer }) => { tracing::debug!( tenant = response.billed_prefix, invoice_type = format!("{:?}", response.invoice_type), subtotal = format!("${:.2}", response.subtotal as f64 / 100.0), "{}", - res.message() + result.message(cmd.dry_run) ); - match res { - InvoiceResult::Created(_) - | InvoiceResult::Updated - | InvoiceResult::Error => {} - // Remove any incorrectly created invoices that are now skipped for whatever reason - _ if cmd.clean_up => { - let task_res: Result<(), anyhow::Error> = async move { - let customer = match get_or_create_customer_for_tenant( - &client, - &db_pool, - response.billed_prefix.to_owned(), - false, - ) - .await? - { - Some(c) => c, - None => return Ok(()), - }; - - let customer_id = customer.id.to_string(); - - if let Some(invoice) = - response.get_stripe_invoice(&client, &customer_id).await? - { - if let Some(InvoiceStatus::Draft) = invoice.status { + + if cmd.clean_up { + let task_res: Result<(), anyhow::Error> = async { + let customer = match customer { + Some(c) => c, + None => return Ok(()), + }; + let customer_id = customer.id.to_string(); + + if let Some(invoice) = + response.get_stripe_invoice(&client, &customer_id).await? + { + if let Some(InvoiceStatus::Draft) = invoice.status { + if cmd.dry_run { + tracing::warn!( + tenant = response.billed_prefix.to_string(), + "[dry-run] Would delete stale draft invoice {}", + invoice.id + ); + } else { tracing::warn!( tenant = response.billed_prefix.to_string(), "Deleting draft invoice!" @@ -669,18 +810,67 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { stripe::Invoice::delete(&client, &invoice.id).await?; } } - - Ok(()) } - .await; + Ok(()) + } + .await; - if let Err(e) = task_res { - tracing::warn!("Failed to check for or clear potential leaked draft invoices for {}, this is probably not a problem: {e:#}", response.billed_prefix.to_owned()); - } - }, - _ => {} + if let Err(e) = task_res { + tracing::warn!("Failed to check for or clear potential leaked draft invoices for {}, this is probably not a problem: {e:#}", response.billed_prefix.to_owned()); + } + } + + Ok((result, response.subtotal, response.billed_prefix.to_owned(), annotation)) + } + Ok(action) if cmd.dry_run => { + let result = match &action { + InvoiceAction::Create { replace: Some(id), .. } => { + tracing::info!( + tenant = response.billed_prefix, + "[dry-run] Would delete existing invoice {} and recreate", + id + ); + InvoiceResult::Created(response.payment_provider) + } + InvoiceAction::Create { .. } => { + InvoiceResult::Created(response.payment_provider) + } + InvoiceAction::Update { .. } => InvoiceResult::Updated, + InvoiceAction::Skip { .. } => unreachable!(), + }; + tracing::debug!( + tenant = response.billed_prefix, + invoice_type = format!("{:?}", response.invoice_type), + subtotal = format!("${:.2}", response.subtotal as f64 / 100.0), + "[dry-run] {}", + result.message(cmd.dry_run) + ); + Ok((result, response.subtotal, response.billed_prefix.to_owned(), annotation)) + } + Ok(action) => { + let res = response + .execute(&client, &db_pool, action, cmd.charge_type) + .await; + match res { + Err(err) => { + let formatted = format!( + "Error publishing {invoice_type:?} invoice for {tenant}", + tenant = response.billed_prefix, + invoice_type = response.invoice_type + ); + Err(anyhow::anyhow!("{formatted}: {err:#}")) + } + Ok(res) => { + tracing::debug!( + tenant = response.billed_prefix, + invoice_type = format!("{:?}", response.invoice_type), + subtotal = format!("${:.2}", response.subtotal as f64 / 100.0), + "{}", + res.message(cmd.dry_run) + ); + Ok((res, response.subtotal, response.billed_prefix.to_owned(), annotation)) + } } - Ok((res, response.subtotal, response.billed_prefix.to_owned())) } } } @@ -691,22 +881,22 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { let total = invoice_futures.len(); - let collected: HashMap)> = + let collected: HashMap)>)> = futures::stream::iter(invoice_futures) .buffer_unordered(cmd.concurrency) .or_else(|(err, invoice)| async move { if !cmd.fail_fast { tracing::error!("[{}]: {err:#}", invoice.billed_prefix); - Ok((InvoiceResult::Error, 0, invoice.billed_prefix)) + Ok((InvoiceResult::Error, 0, invoice.billed_prefix, None)) } else { Err(err) } }) .try_fold( HashMap::new(), - |mut map, (res, subtotal, tenant)| async move { + |mut map, (res, subtotal, tenant, annotation)| async move { let overall_count = map.values().map(|(_, count, _)| *count).sum::() + 1; - let msg = res.message(); + let msg = res.message(cmd.dry_run); let (subtotal_sum, count_for_result_type, tenants) = map.entry(res).or_insert((0, 0, vec![])); @@ -714,7 +904,7 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { *count_for_result_type += 1; tracing::info!("[{overall_count}/{total}, {tenant}]: {msg}"); - tenants.push((tenant, subtotal)); + tenants.push((tenant, subtotal, annotation)); Ok(map) }, ) @@ -724,7 +914,7 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { tracing::info!( "[{:4} invoices]: {:70}${:.2}", count, - status.message(), + status.message(cmd.dry_run), *subtotal_agg as f64 / 100.0 ); let limit = match status { @@ -732,18 +922,24 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { InvoiceResult::NoDataMoved | InvoiceResult::NoFullPipeline | InvoiceResult::LessThanMinimum - | InvoiceResult::FreeTier => 0, + | InvoiceResult::FreeTier + | InvoiceResult::AlreadyProcessed => 0, _ => 10, }; let sorted_tenants = tenants .iter() - .sorted_by(|(_, a), (_, b)| b.cmp(a)) + .sorted_by(|(_, a, _), (_, b, _)| b.cmp(a)) .collect_vec(); let (displayed_tenants, remainder_tenants) = sorted_tenants.split_at(limit.min(tenants.len())); - for (tenant, subtotal) in displayed_tenants { - tracing::info!(" - {:} ${:.2}", tenant, *subtotal as f64 / 100.0); + for (tenant, subtotal, annotation) in displayed_tenants { + match annotation { + Some(note) => { + tracing::info!(" - {} ${:.2} {}", tenant, *subtotal as f64 / 100.0, note) + } + None => tracing::info!(" - {} ${:.2}", tenant, *subtotal as f64 / 100.0), + } } if limit > 0 && remainder_tenants.len() > 0 { tracing::info!(" - ... {} Others", remainder_tenants.len(),); @@ -753,35 +949,50 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { Ok(()) } -#[tracing::instrument(skip(client, db_client))] -async fn get_or_create_customer_for_tenant( +/// Read-only: search Stripe for an existing customer by tenant metadata. +#[tracing::instrument(skip(client))] +async fn find_customer( client: &stripe::Client, - db_client: &Pool, - tenant: String, - create: bool, + tenant: &str, ) -> anyhow::Result> { - let billing_row = sqlx::query!( - r#"SELECT billing_email, billing_address FROM tenants WHERE tenant = $1"#, - tenant, - ) - .fetch_optional(db_client) - .await?; - let customers: Vec = stripe_search( client, "customers", SearchParams { - query: customer_search_query(&tenant), + query: customer_search_query(tenant), ..Default::default() }, ) .await .context(format!("Searching for tenant {tenant}"))?; - let customer = if let Some(customer) = customers.into_iter().next() { + if let Some(customer) = customers.into_iter().next() { tracing::debug!("Found existing customer {id}", id = customer.id.to_string()); + Ok(Some(customer)) + } else { + Ok(None) + } +} + +/// Ensures a Stripe customer exists for this tenant and is ready for invoicing. +/// Finds an existing customer or creates a new one, then ensures the customer +/// has an email set (looking up the earliest admin on the tenant if needed). +#[tracing::instrument(skip(client, db_client))] +async fn ensure_customer_for_invoicing( + client: &stripe::Client, + db_client: &Pool, + tenant: &str, +) -> anyhow::Result { + let billing_row = sqlx::query!( + r#"SELECT billing_email, billing_address FROM tenants WHERE tenant = $1"#, + tenant, + ) + .fetch_optional(db_client) + .await?; + + let customer = if let Some(customer) = find_customer(client, tenant).await? { customer - } else if create { + } else { tracing::debug!("Creating new customer"); // Match the deterministic Idempotency-Key used by the GraphQL path so // a setup-intent flow racing against billing automation can't produce @@ -789,7 +1000,7 @@ async fn get_or_create_customer_for_tenant( let create_client = client .clone() .with_strategy(stripe::RequestStrategy::Idempotent( - customer_create_idempotency_key(&tenant), + customer_create_idempotency_key(tenant), )); let billing_email = billing_row @@ -802,17 +1013,16 @@ async fn get_or_create_customer_for_tenant( .transpose() .context("deserializing billing_address")?; + let description = format!("Represents the billing entity for Flow tenant '{tenant}'"); let new_customer = stripe::Customer::create( &create_client, stripe::CreateCustomer { - name: Some(tenant.as_str()), + name: Some(tenant), email: billing_email, address: billing_address, - description: Some( - format!("Represents the billing entity for Flow tenant '{tenant}'").as_str(), - ), + description: Some(description.as_str()), metadata: Some({ - let mut metadata = tenant_metadata(tenant.as_str()); + let mut metadata = tenant_metadata(tenant); metadata.insert( CREATED_BY_BILLING_AUTOMATION.to_string(), "true".to_string(), @@ -830,8 +1040,6 @@ async fn get_or_create_customer_for_tenant( // Waking it during an invoicing run would add a Stripe lookup per new // customer for no contact change. new_customer - } else { - return Ok(None); }; if customer.email.is_none() { @@ -893,5 +1101,5 @@ async fn get_or_create_customer_for_tenant( } } } - Ok(Some(customer)) + Ok(customer) } From f542dd5f384dca98a4fc915456a95a5495967b68 Mon Sep 17 00:00:00 2001 From: Joseph Shearer Date: Fri, 10 Apr 2026 17:53:41 -0400 Subject: [PATCH 2/6] billing: always send manual invoices rather than auto-charging Customers' stored payment methods are for monthly usage overages. Manual bills (contracts, one-off charges, etc.) should be sent as invoices so the customer can decide how to pay, rather than being automatically charged to their payment method. * Override `charge_type` to `SendInvoice` for manual invoices during creation in `publish` * Switch manual invoices from `charge_automatically` to `send_invoice` during the send phase, even if the customer has a payment method on file --- crates/billing-integrations/src/publish.rs | 8 ++++++++ crates/billing-integrations/src/send.rs | 6 ++++-- crates/billing-integrations/src/stripe_utils.rs | 10 +++++++++- 3 files changed, 21 insertions(+), 3 deletions(-) diff --git a/crates/billing-integrations/src/publish.rs b/crates/billing-integrations/src/publish.rs index abb5f5b9e24..ecc4f70e08e 100644 --- a/crates/billing-integrations/src/publish.rs +++ b/crates/billing-integrations/src/publish.rs @@ -466,6 +466,14 @@ impl Invoice { } // Create or reuse the invoice + // Manual invoices should always be sent as invoices rather than + // charged to the customer's payment method. + let mode = if self.invoice_type == InvoiceType::Manual { + ChargeType::SendInvoice + } else { + mode + }; + let invoice = if let Some(existing_id) = existing_invoice_id { tracing::debug!( "Updating existing invoice {id}", diff --git a/crates/billing-integrations/src/send.rs b/crates/billing-integrations/src/send.rs index bd626870726..7fb1f58221c 100644 --- a/crates/billing-integrations/src/send.rs +++ b/crates/billing-integrations/src/send.rs @@ -183,13 +183,15 @@ async fn update_draft_collection_methods( stripe_client: &Client, mut to_update: Vec, ) -> anyhow::Result> { - // Identify invoices that are `charge_automatically` but don't have a default payment method + // Identify invoices that need to be switched to `send_invoice`: + // - Manual invoices should always be sent as invoices, never auto-charged + // - Auto-charge invoices without a payment method on file must be sent as invoices let needs_update: HashSet = to_update .iter() .filter(|inv| { inv.collection_method().map_or(false, |cm| { cm == stripe::CollectionMethod::ChargeAutomatically - }) && !inv.has_cc() + }) && (inv.is_manual() || !inv.has_cc()) }) .map(|inv| inv.id().clone()) .collect::>(); diff --git a/crates/billing-integrations/src/stripe_utils.rs b/crates/billing-integrations/src/stripe_utils.rs index fa35f97f562..f93b50758b7 100644 --- a/crates/billing-integrations/src/stripe_utils.rs +++ b/crates/billing-integrations/src/stripe_utils.rs @@ -1,4 +1,4 @@ -use billing_types::{InvoiceMetadata, SearchParams, stripe_search}; +use billing_types::{InvoiceMetadata, InvoiceType, SearchParams, stripe_search}; use num_format::{Locale, ToFormattedString}; use std::ops::{Deref, DerefMut}; @@ -84,6 +84,14 @@ impl Invoice { self.0.status.clone() } + pub fn is_manual(&self) -> bool { + self.0 + .metadata + .as_ref() + .and_then(InvoiceMetadata::from_metadata_map) + .map_or(false, |m| m.invoice_type == InvoiceType::Manual) + } + pub fn period_start(&self) -> Option { self.0 .metadata From 2e634d1fda4c386441c6ed2d2d88529309650f39 Mon Sep 17 00:00:00 2001 From: Joseph Shearer Date: Wed, 29 Jul 2026 14:39:14 -0400 Subject: [PATCH 3/6] billing: harden invoice publication workflows Preserve cleanup behavior and dry-run fidelity across the classification and execution split, including billing-email validation, fail-fast propagation, customer reuse, and consistent reporting. Replacements are reported as a distinct `Replaced` result, so run summaries show how many invoices were (or would be) voided or deleted and reissued. Reconcile collection methods toward `send_invoice` only, so refreshing a draft never reverts the send workflow's correction for tenants without a payment method, and re-check invoice state before update or replacement. Void open invoices during replacement and recover cleanly from interrupted runs. Determine payment-method state directly from Stripe rather than the DB's `stripe.customers` capture, which has been unreliable. --- ...0583e9e6de51321af2bf6effac7488684a326.json | 144 ---- ...52cba7a2db38f044b2c229c90cd6908e9c750.json | 87 ++ ...edda921cc72a19fb8b5301b27c31053766883.json | 143 ---- ...5c4507359f243a2307e0ff08c5a83b9a01a05.json | 86 ++ crates/billing-integrations/src/publish.rs | 767 +++++++++++------- crates/billing-integrations/src/send.rs | 5 +- .../billing-integrations/src/stripe_utils.rs | 2 +- 7 files changed, 658 insertions(+), 576 deletions(-) delete mode 100644 .sqlx/query-36025f03b3f4ef7e88f17f50e210583e9e6de51321af2bf6effac7488684a326.json create mode 100644 .sqlx/query-402d978379e462304be791a274e52cba7a2db38f044b2c229c90cd6908e9c750.json delete mode 100644 .sqlx/query-42edae75728926abb9b438ea400edda921cc72a19fb8b5301b27c31053766883.json create mode 100644 .sqlx/query-c586f9fd13f5c26682d5516ab985c4507359f243a2307e0ff08c5a83b9a01a05.json diff --git a/.sqlx/query-36025f03b3f4ef7e88f17f50e210583e9e6de51321af2bf6effac7488684a326.json b/.sqlx/query-36025f03b3f4ef7e88f17f50e210583e9e6de51321af2bf6effac7488684a326.json deleted file mode 100644 index 33718217d6e..00000000000 --- a/.sqlx/query-36025f03b3f4ef7e88f17f50e210583e9e6de51321af2bf6effac7488684a326.json +++ /dev/null @@ -1,144 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n select\n date_start as \"date_start!\",\n date_end as \"date_end!\",\n billed_prefix as \"billed_prefix!\",\n invoice_type as \"invoice_type!: InvoiceType\",\n line_items as \"line_items!: sqlx::types::Json>\",\n subtotal::bigint as \"subtotal!\",\n extra as \"extra: sqlx::types::Json>\",\n customer.has_payment_method as has_payment_method,\n coalesce(dataflow.has_full_pipeline, false) as \"has_full_pipeline!\",\n tenants.payment_provider as \"payment_provider!: PaymentProvider\",\n tenants.trial_start as tenant_trial_start\n from invoices_ext\n left join tenants on tenants.tenant = billed_prefix\n left join lateral(\n \tselect bool_or(\"invoice_settings/default_payment_method\" is not null) as has_payment_method\n \tfrom stripe.customers\n \twhere customers.metadata->>'estuary.dev/tenant_name' = billed_prefix\n \tgroup by billed_prefix\n ) as customer on true\n left join lateral(\n \tselect\n \t\tsum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0\n \t\tand sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'materialization') > 0\n \t\tas has_full_pipeline\n from catalog_stats_monthly\n join live_specs on live_specs.catalog_name ^@ catalog_stats_monthly.catalog_name\n where\n \tcatalog_stats_monthly.catalog_name = billed_prefix\n \tand tstzrange(date_trunc('day', $1::date), date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day') @> catalog_stats_monthly.ts\n ) as dataflow on true\n where ((\n date_start >= date_trunc('day', $1::date)\n and date_end <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and invoice_type = 'final'\n ) or (\n date_start <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and date_end >= date_trunc('day', $1::date)\n and invoice_type = 'manual'\n ))\n and billed_prefix = any($2)\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "date_start!", - "type_info": "Date", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "date_start" - } - } - }, - { - "ordinal": 1, - "name": "date_end!", - "type_info": "Date", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "date_end" - } - } - }, - { - "ordinal": 2, - "name": "billed_prefix!", - "type_info": "Text", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "billed_prefix" - } - } - }, - { - "ordinal": 3, - "name": "invoice_type!: InvoiceType", - "type_info": "Text", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "invoice_type" - } - } - }, - { - "ordinal": 4, - "name": "line_items!: sqlx::types::Json>", - "type_info": "Jsonb", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "line_items" - } - } - }, - { - "ordinal": 5, - "name": "subtotal!", - "type_info": "Int8", - "origin": "Expression" - }, - { - "ordinal": 6, - "name": "extra: sqlx::types::Json>", - "type_info": "Jsonb", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "extra" - } - } - }, - { - "ordinal": 7, - "name": "has_payment_method", - "type_info": "Bool", - "origin": "Expression" - }, - { - "ordinal": 8, - "name": "has_full_pipeline!", - "type_info": "Bool", - "origin": "Expression" - }, - { - "ordinal": 9, - "name": "payment_provider!: PaymentProvider", - "type_info": { - "Custom": { - "name": "payment_provider_type", - "kind": { - "Enum": [ - "stripe", - "external" - ] - } - } - }, - "origin": { - "Table": { - "table": "tenants", - "name": "payment_provider" - } - } - }, - { - "ordinal": 10, - "name": "tenant_trial_start", - "type_info": "Date", - "origin": { - "Table": { - "table": "tenants", - "name": "trial_start" - } - } - } - ], - "parameters": { - "Left": [ - "Date", - "TextArray" - ] - }, - "nullable": [ - true, - true, - true, - true, - true, - null, - true, - null, - null, - true, - true - ] - }, - "hash": "36025f03b3f4ef7e88f17f50e210583e9e6de51321af2bf6effac7488684a326" -} diff --git a/.sqlx/query-402d978379e462304be791a274e52cba7a2db38f044b2c229c90cd6908e9c750.json b/.sqlx/query-402d978379e462304be791a274e52cba7a2db38f044b2c229c90cd6908e9c750.json new file mode 100644 index 00000000000..033a7905d4a --- /dev/null +++ b/.sqlx/query-402d978379e462304be791a274e52cba7a2db38f044b2c229c90cd6908e9c750.json @@ -0,0 +1,87 @@ +{ + "db_name": "PostgreSQL", + "query": "\n select\n date_start as \"date_start!\",\n date_end as \"date_end!\",\n billed_prefix as \"billed_prefix!\",\n invoice_type as \"invoice_type!: InvoiceType\",\n line_items as \"line_items!: sqlx::types::Json>\",\n subtotal::bigint as \"subtotal!\",\n extra as \"extra: sqlx::types::Json>\",\n coalesce(dataflow.has_full_pipeline, false) as \"has_full_pipeline!\",\n tenants.payment_provider as \"payment_provider!: PaymentProvider\",\n tenants.trial_start as tenant_trial_start\n from invoices_ext\n left join tenants on tenants.tenant = billed_prefix\n left join lateral(\n \tselect\n \t\tsum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0\n \t\tand sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'materialization') > 0\n \t\tas has_full_pipeline\n from catalog_stats_monthly\n join live_specs on live_specs.catalog_name ^@ catalog_stats_monthly.catalog_name\n where\n \tcatalog_stats_monthly.catalog_name = billed_prefix\n \tand tstzrange(date_trunc('day', $1::date), date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day') @> catalog_stats_monthly.ts\n ) as dataflow on true\n where ((\n date_start >= date_trunc('day', $1::date)\n and date_end <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and invoice_type = 'final'\n ) or (\n date_start <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and date_end >= date_trunc('day', $1::date)\n and invoice_type = 'manual'\n ))\n and billed_prefix = any($2)\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "date_start!", + "type_info": "Date" + }, + { + "ordinal": 1, + "name": "date_end!", + "type_info": "Date" + }, + { + "ordinal": 2, + "name": "billed_prefix!", + "type_info": "Text" + }, + { + "ordinal": 3, + "name": "invoice_type!: InvoiceType", + "type_info": "Text" + }, + { + "ordinal": 4, + "name": "line_items!: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 5, + "name": "subtotal!", + "type_info": "Int8" + }, + { + "ordinal": 6, + "name": "extra: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 7, + "name": "has_full_pipeline!", + "type_info": "Bool" + }, + { + "ordinal": 8, + "name": "payment_provider!: PaymentProvider", + "type_info": { + "Custom": { + "name": "payment_provider_type", + "kind": { + "Enum": [ + "stripe", + "external" + ] + } + } + } + }, + { + "ordinal": 9, + "name": "tenant_trial_start", + "type_info": "Date" + } + ], + "parameters": { + "Left": [ + "Date", + "TextArray" + ] + }, + "nullable": [ + true, + true, + true, + true, + true, + null, + true, + null, + true, + true + ] + }, + "hash": "402d978379e462304be791a274e52cba7a2db38f044b2c229c90cd6908e9c750" +} diff --git a/.sqlx/query-42edae75728926abb9b438ea400edda921cc72a19fb8b5301b27c31053766883.json b/.sqlx/query-42edae75728926abb9b438ea400edda921cc72a19fb8b5301b27c31053766883.json deleted file mode 100644 index a52fd329a23..00000000000 --- a/.sqlx/query-42edae75728926abb9b438ea400edda921cc72a19fb8b5301b27c31053766883.json +++ /dev/null @@ -1,143 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n select\n date_start as \"date_start!\",\n date_end as \"date_end!\",\n billed_prefix as \"billed_prefix!\",\n invoice_type as \"invoice_type!: InvoiceType\",\n line_items as \"line_items!: sqlx::types::Json>\",\n subtotal::bigint as \"subtotal!\",\n extra as \"extra: sqlx::types::Json>\",\n customer.has_payment_method as has_payment_method,\n coalesce(dataflow.has_full_pipeline, false) as \"has_full_pipeline!\",\n tenants.payment_provider as \"payment_provider!: PaymentProvider\",\n tenants.trial_start as tenant_trial_start\n from invoices_ext\n left join tenants on tenants.tenant = billed_prefix\n left join lateral(\n \tselect bool_or(\"invoice_settings/default_payment_method\" is not null) as has_payment_method\n \tfrom stripe.customers\n \twhere customers.metadata->>'estuary.dev/tenant_name' = billed_prefix\n \tgroup by billed_prefix\n ) as customer on true\n left join lateral(\n \tselect\n \t\tsum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0\n \t\tand sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'materialization') > 0\n \t\tas has_full_pipeline\n from catalog_stats_monthly\n join live_specs on live_specs.catalog_name ^@ catalog_stats_monthly.catalog_name\n where\n \tcatalog_stats_monthly.catalog_name = billed_prefix\n \tand tstzrange(date_trunc('day', $1::date), date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day') @> catalog_stats_monthly.ts\n ) as dataflow on true\n where (\n date_start >= date_trunc('day', $1::date)\n and date_end <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and invoice_type = 'final'\n ) or (\n date_start <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and date_end >= date_trunc('day', $1::date)\n and invoice_type = 'manual'\n )\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "date_start!", - "type_info": "Date", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "date_start" - } - } - }, - { - "ordinal": 1, - "name": "date_end!", - "type_info": "Date", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "date_end" - } - } - }, - { - "ordinal": 2, - "name": "billed_prefix!", - "type_info": "Text", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "billed_prefix" - } - } - }, - { - "ordinal": 3, - "name": "invoice_type!: InvoiceType", - "type_info": "Text", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "invoice_type" - } - } - }, - { - "ordinal": 4, - "name": "line_items!: sqlx::types::Json>", - "type_info": "Jsonb", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "line_items" - } - } - }, - { - "ordinal": 5, - "name": "subtotal!", - "type_info": "Int8", - "origin": "Expression" - }, - { - "ordinal": 6, - "name": "extra: sqlx::types::Json>", - "type_info": "Jsonb", - "origin": { - "Table": { - "table": "invoices_ext", - "name": "extra" - } - } - }, - { - "ordinal": 7, - "name": "has_payment_method", - "type_info": "Bool", - "origin": "Expression" - }, - { - "ordinal": 8, - "name": "has_full_pipeline!", - "type_info": "Bool", - "origin": "Expression" - }, - { - "ordinal": 9, - "name": "payment_provider!: PaymentProvider", - "type_info": { - "Custom": { - "name": "payment_provider_type", - "kind": { - "Enum": [ - "stripe", - "external" - ] - } - } - }, - "origin": { - "Table": { - "table": "tenants", - "name": "payment_provider" - } - } - }, - { - "ordinal": 10, - "name": "tenant_trial_start", - "type_info": "Date", - "origin": { - "Table": { - "table": "tenants", - "name": "trial_start" - } - } - } - ], - "parameters": { - "Left": [ - "Date" - ] - }, - "nullable": [ - true, - true, - true, - true, - true, - null, - true, - null, - null, - true, - true - ] - }, - "hash": "42edae75728926abb9b438ea400edda921cc72a19fb8b5301b27c31053766883" -} diff --git a/.sqlx/query-c586f9fd13f5c26682d5516ab985c4507359f243a2307e0ff08c5a83b9a01a05.json b/.sqlx/query-c586f9fd13f5c26682d5516ab985c4507359f243a2307e0ff08c5a83b9a01a05.json new file mode 100644 index 00000000000..c00968f1d4e --- /dev/null +++ b/.sqlx/query-c586f9fd13f5c26682d5516ab985c4507359f243a2307e0ff08c5a83b9a01a05.json @@ -0,0 +1,86 @@ +{ + "db_name": "PostgreSQL", + "query": "\n select\n date_start as \"date_start!\",\n date_end as \"date_end!\",\n billed_prefix as \"billed_prefix!\",\n invoice_type as \"invoice_type!: InvoiceType\",\n line_items as \"line_items!: sqlx::types::Json>\",\n subtotal::bigint as \"subtotal!\",\n extra as \"extra: sqlx::types::Json>\",\n coalesce(dataflow.has_full_pipeline, false) as \"has_full_pipeline!\",\n tenants.payment_provider as \"payment_provider!: PaymentProvider\",\n tenants.trial_start as tenant_trial_start\n from invoices_ext\n left join tenants on tenants.tenant = billed_prefix\n left join lateral(\n \tselect\n \t\tsum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0\n \t\tand sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'materialization') > 0\n \t\tas has_full_pipeline\n from catalog_stats_monthly\n join live_specs on live_specs.catalog_name ^@ catalog_stats_monthly.catalog_name\n where\n \tcatalog_stats_monthly.catalog_name = billed_prefix\n \tand tstzrange(date_trunc('day', $1::date), date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day') @> catalog_stats_monthly.ts\n ) as dataflow on true\n where (\n date_start >= date_trunc('day', $1::date)\n and date_end <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and invoice_type = 'final'\n ) or (\n date_start <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and date_end >= date_trunc('day', $1::date)\n and invoice_type = 'manual'\n )\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "date_start!", + "type_info": "Date" + }, + { + "ordinal": 1, + "name": "date_end!", + "type_info": "Date" + }, + { + "ordinal": 2, + "name": "billed_prefix!", + "type_info": "Text" + }, + { + "ordinal": 3, + "name": "invoice_type!: InvoiceType", + "type_info": "Text" + }, + { + "ordinal": 4, + "name": "line_items!: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 5, + "name": "subtotal!", + "type_info": "Int8" + }, + { + "ordinal": 6, + "name": "extra: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 7, + "name": "has_full_pipeline!", + "type_info": "Bool" + }, + { + "ordinal": 8, + "name": "payment_provider!: PaymentProvider", + "type_info": { + "Custom": { + "name": "payment_provider_type", + "kind": { + "Enum": [ + "stripe", + "external" + ] + } + } + } + }, + { + "ordinal": 9, + "name": "tenant_trial_start", + "type_info": "Date" + } + ], + "parameters": { + "Left": [ + "Date" + ] + }, + "nullable": [ + true, + true, + true, + true, + true, + null, + true, + null, + true, + true + ] + }, + "hash": "c586f9fd13f5c26682d5516ab985c4507359f243a2307e0ff08c5a83b9a01a05" +} diff --git a/crates/billing-integrations/src/publish.rs b/crates/billing-integrations/src/publish.rs index ecc4f70e08e..28e06692bf3 100644 --- a/crates/billing-integrations/src/publish.rs +++ b/crates/billing-integrations/src/publish.rs @@ -10,7 +10,6 @@ use serde::{Deserialize, Serialize}; use sqlx::Postgres; use sqlx::{Pool, postgres::PgPoolOptions, types::chrono::NaiveDate}; use std::collections::HashMap; -use stripe::InvoiceStatus; const CREATED_BY_BILLING_AUTOMATION: &str = "estuary.dev/created_by_automation"; @@ -40,7 +39,8 @@ pub struct PublishInvoice { /// The month to generate invoices for, in format "YYYY-MM-DD" #[clap(long, value_parser = parse_date)] month: NaiveDate, - /// Whether to delete and recreate finalized invoices + /// Whether to recreate existing invoices: drafts are deleted and open + /// invoices are voided before the replacement is created. #[clap(long)] recreate_finalized: bool, /// Stop execution after first failure @@ -60,6 +60,11 @@ pub struct PublishInvoice { pub clean_up: bool, /// Run in read-only mode: classify all invoices and report what would /// happen, without creating or modifying anything in Stripe. + /// + /// The preview checks that each invoice could be issued (a billing email is + /// resolvable), but it cannot reconcile the final invoice total against + /// Stripe, since that requires actually creating the invoice and its line + /// items. #[clap(long, default_value_t = false)] pub dry_run: bool, } @@ -89,6 +94,7 @@ struct LineItem { enum InvoiceResult { Created(PaymentProvider), Updated, + Replaced, LessThanMinimum, FreeTier, FutureTrialStart, @@ -120,6 +126,13 @@ impl InvoiceResult { "Updated existing invoice".to_string() } } + InvoiceResult::Replaced => { + if dry_run { + "Would void/delete and replace existing invoice".to_string() + } else { + "Replaced existing invoice".to_string() + } + } InvoiceResult::LessThanMinimum => { "Skipping invoice for less than the minimum chargable amount ($0.50)".to_string() } @@ -136,7 +149,13 @@ impl InvoiceResult { InvoiceResult::AlreadyProcessed => { "Skipping invoice already processed in a previous billing run".to_string() } - InvoiceResult::Error => "Error publishing invoices".to_string(), + InvoiceResult::Error => { + if dry_run { + "Would fail to publish invoices".to_string() + } else { + "Error publishing invoices".to_string() + } + } } } } @@ -150,14 +169,33 @@ enum InvoiceAction { customer: Option, }, /// Create a new invoice. `replace` is set when --recreate-finalized - /// requires deleting an existing invoice first. - Create { replace: Option }, - /// Update an existing draft invoice's line items. + /// requires deleting an existing invoice first. `customer` is the customer + /// found during classification, or None when none exists yet (execute then + /// creates one). + Create { + replace: Option, + customer: Option, + }, + /// Update an existing draft invoice's line items. `customer` is the invoice's + /// owner, already found during classification. Update { existing_invoice_id: stripe::InvoiceId, + customer: stripe::Customer, }, } +impl InvoiceAction { + /// The customer located during classification, if one exists in Stripe. + fn customer(&self) -> Option<&stripe::Customer> { + match self { + InvoiceAction::Skip { customer, .. } | InvoiceAction::Create { customer, .. } => { + customer.as_ref() + } + InvoiceAction::Update { customer, .. } => Some(customer), + } + } +} + #[derive(Serialize, Deserialize, Debug, Clone, sqlx::FromRow)] struct Invoice { subtotal: i64, @@ -167,7 +205,6 @@ struct Invoice { billed_prefix: String, invoice_type: InvoiceType, extra: Option>>, - has_payment_method: Option, has_full_pipeline: bool, payment_provider: PaymentProvider, tenant_trial_start: Option, @@ -200,7 +237,19 @@ impl Invoice { .await .context("Searching for an invoice")?; - Ok(invoice_search.into_iter().next()) + // Prefer a live invoice over a voided one: after --recreate-finalized + // voids and recreates an invoice, both match this search, and + // classification should act on the replacement. A voided invoice is + // still returned when it's the only match, so that a voided manual + // invoice classifies as AlreadyProcessed rather than being recreated. + let mut voided = None; + for invoice in invoice_search { + if invoice.status != Some(stripe::InvoiceStatus::Void) { + return Ok(Some(invoice)); + } + voided.get_or_insert(invoice); + } + Ok(voided) } /// Read-only classification: determines what action should be taken for this @@ -211,8 +260,6 @@ impl Invoice { client: &stripe::Client, recreate_finalized: bool, ) -> anyhow::Result { - // --- Phase 1: Cheap local checks (no Stripe calls) --- - match (&self.invoice_type, &self.extra) { (InvoiceType::Preview, _) => { bail!("Should not create Stripe invoices for preview invoices") @@ -248,84 +295,88 @@ impl Invoice { }); } - // --- Phase 2: Stripe calls (only for invoices that survived Phase 1) --- - - // For Final invoices, verify the payment method state with Stripe. - // The DB capture has been known to be unreliable, so Stripe is the - // source of truth. If the tenant has no payment method, skip on - // NoDataMoved / NoFullPipeline. - let mut found_customer: Option> = None; + // Reuse this lookup for payment-method checks, invoice search, and execution. + let customer = find_customer(client, &self.billed_prefix).await?; if let (InvoiceType::Final, Some(extra)) = (&self.invoice_type, &self.extra) { - let validated_has_payment_method = - if let Some(has_payment_method) = self.has_payment_method { - let customer = find_customer(client, &self.billed_prefix).await?; - let real_has_pm = customer - .as_ref() - .and_then(|c| c.invoice_settings.as_ref()) - .and_then(|i| i.default_payment_method.as_ref()) - .is_some(); - - if has_payment_method != real_has_pm { - tracing::warn!( - ?has_payment_method, - stripe_payment_method = real_has_pm, - "Inconsistent payment method state" - ); - } - - found_customer = Some(customer); - real_has_pm - } else { - false - }; - - if !validated_has_payment_method { - let unwrapped_extra = extra.clone().0.expect( - "This is just a sqlx quirk, if the outer Option is Some then this will be Some", + // If the outer Option is Some the inner should be too; treat a + // null inner as a data error rather than panicking and aborting + // the whole run. + let Some(unwrapped_extra) = extra.0.as_ref() else { + bail!( + "Final invoice for {tenant} has a null `extra` payload", + tenant = self.billed_prefix ); + }; + + let has_payment_method = customer + .as_ref() + .and_then(|c| c.invoice_settings.as_ref()) + .and_then(|i| i.default_payment_method.as_ref()) + .is_some(); + if !has_payment_method { if unwrapped_extra.processed_data_gb.unwrap_or_default() == 0.0 { return Ok(InvoiceAction::Skip { result: InvoiceResult::NoDataMoved, - customer: found_customer.flatten(), + customer, }); } if !self.has_full_pipeline { return Ok(InvoiceAction::Skip { result: InvoiceResult::NoFullPipeline, - customer: found_customer.flatten(), + customer, }); } } } - // Look up customer (reuse if already fetched during payment method validation) - let customer = match found_customer { - Some(c) => c, - None => find_customer(client, &self.billed_prefix).await?, - }; - let customer = match customer { Some(c) => c, // No customer in Stripe means no existing invoice is possible - None => return Ok(InvoiceAction::Create { replace: None }), + None => { + return Ok(InvoiceAction::Create { + replace: None, + customer: None, + }); + } }; let customer_id = customer.id.to_string(); - // Search for an existing invoice in Stripe if let Some(invoice) = self .get_stripe_invoice(client, customer_id.as_str()) .await? { match invoice.status { + // Manual invoices are excluded from --recreate-finalized: an + // already-sent (open) manual invoice must not be voided and + // reissued, and a manual draft is refreshed via the Update arm below. Some(stripe::InvoiceStatus::Open | stripe::InvoiceStatus::Draft) - if recreate_finalized => + if recreate_finalized && !matches!(self.invoice_type, InvoiceType::Manual) => { Ok(InvoiceAction::Create { replace: Some(invoice.id), + customer: Some(customer), + }) + } + // A voided invoice can be neither deleted nor updated, so treat + // it as absent and create a fresh replacement. This also + // recovers a --recreate-finalized run that voided an open + // invoice but failed before creating its replacement. Without + // the flag, the unsupported-state error below keeps a voided + // invoice loud rather than silently re-billing the tenant. + Some(stripe::InvoiceStatus::Void) + if recreate_finalized && !matches!(self.invoice_type, InvoiceType::Manual) => + { + tracing::warn!( + "Found voided invoice {id}; treating it as absent", + id = invoice.id.to_string() + ); + Ok(InvoiceAction::Create { + replace: None, + customer: Some(customer), }) } Some(stripe::InvoiceStatus::Draft) => { @@ -335,6 +386,7 @@ impl Invoice { ); Ok(InvoiceAction::Update { existing_invoice_id: invoice.id, + customer, }) } Some(stripe::InvoiceStatus::Open) @@ -351,7 +403,7 @@ impl Invoice { } Some(stripe::InvoiceStatus::Open) => { bail!( - "Found open invoice {id}. Pass --recreate-finalized to delete and recreate this invoice.", + "Found open invoice {id}. Pass --recreate-finalized to void and recreate this invoice.", id = invoice.id.to_string() ) } @@ -385,7 +437,10 @@ impl Invoice { } } } else { - Ok(InvoiceAction::Create { replace: None }) + Ok(InvoiceAction::Create { + replace: None, + customer: Some(customer), + }) } } @@ -399,18 +454,21 @@ impl Invoice { action: InvoiceAction, mode: ChargeType, ) -> anyhow::Result { - let (is_update, replace, existing_invoice_id) = match action { - InvoiceAction::Skip { result, .. } => return Ok(result), - InvoiceAction::Create { replace, .. } => (false, replace, None), + let (replace, existing_invoice_id, found_customer) = match action { + InvoiceAction::Skip { .. } => { + unreachable!("Skip actions are resolved by the caller and never executed") + } + InvoiceAction::Create { replace, customer } => (replace, None, customer), InvoiceAction::Update { existing_invoice_id, - .. - } => (true, None, Some(existing_invoice_id)), + customer, + } => (None, Some(existing_invoice_id), Some(customer)), }; + let is_update = existing_invoice_id.is_some(); - // Ensure customer exists and has an email (required for invoicing) let customer = - ensure_customer_for_invoicing(client, db_client, &self.billed_prefix).await?; + ensure_customer_for_invoicing(client, db_client, &self.billed_prefix, found_customer) + .await?; // Anything before 12:00:00 renders as the previous day in Stripe let date_start_secs = self @@ -436,22 +494,30 @@ impl Invoice { let date_start_repr = self.date_start.format("%F").to_string(); let date_end_repr = self.date_end.format("%F").to_string(); - // Delete existing invoice if --recreate-finalized was used + // Remove the existing invoice if --recreate-finalized was used if let Some(ref replace_id) = replace { - // Re-verify the invoice status before deleting (guard against race conditions) + // Re-verify the invoice status before removing it (guard against race conditions) let existing = stripe::Invoice::retrieve(client, replace_id, &[]).await?; match existing.status { - Some(state @ (stripe::InvoiceStatus::Open | stripe::InvoiceStatus::Draft)) => { + Some(stripe::InvoiceStatus::Draft) => { tracing::warn!( - "Found invoice {id} in state {state}, deleting and recreating", + "Found draft invoice {id}, deleting and recreating", id = replace_id.to_string(), - state = state ); stripe::Invoice::delete(client, replace_id).await?; } + // Stripe only permits deleting drafts; a finalized invoice + // must be voided instead. + Some(stripe::InvoiceStatus::Open) => { + tracing::warn!( + "Found open invoice {id}, voiding and recreating", + id = replace_id.to_string(), + ); + stripe::Invoice::void(client, replace_id).await?; + } Some(status) => { bail!( - "Invoice {id} changed to state {status} since classification, cannot delete.", + "Invoice {id} changed to state {status} since classification, cannot replace.", id = replace_id.to_string(), status = status ); @@ -465,7 +531,6 @@ impl Invoice { } } - // Create or reuse the invoice // Manual invoices should always be sent as invoices rather than // charged to the customer's payment method. let mode = if self.invoice_type == InvoiceType::Manual { @@ -473,13 +538,78 @@ impl Invoice { } else { mode }; + let collection_method = match mode { + ChargeType::AutoCharge => stripe::CollectionMethod::ChargeAutomatically, + ChargeType::SendInvoice => stripe::CollectionMethod::SendInvoice, + }; + // `send_invoice` requires a due date; `charge_automatically` must not carry one. + let due_date = match mode { + ChargeType::SendInvoice => Some((Utc::now() + Duration::days(30)).timestamp()), + ChargeType::AutoCharge => None, + }; let invoice = if let Some(existing_id) = existing_invoice_id { tracing::debug!( "Updating existing invoice {id}", id = existing_id.to_string() ); - stripe::Invoice::retrieve(client, &existing_id, &[]).await? + let existing = stripe::Invoice::retrieve(client, &existing_id, &[]).await?; + + // Re-verify the invoice is still an updatable draft. It was a draft at + // classification time, but a concurrent run or a human finalizing / + // paying / voiding it in Stripe could have changed that; mutating a + // finalized invoice's line items would otherwise fail with an opaque + // Stripe error. Mirrors the guard on the delete/recreate path above. + match existing.status { + Some(stripe::InvoiceStatus::Draft) => {} + Some(status) => { + bail!( + "Invoice {id} changed to state {status} since classification, cannot update.", + id = existing_id.to_string(), + status = status + ); + } + None => { + bail!( + "Unexpected missing status from invoice {id}", + id = existing_id.to_string() + ); + } + } + + // The create path sets the collection method from `mode`, but the update + // path reuses whatever the invoice was originally created with. Reconcile + // toward `send_invoice` only, so a manual invoice previously stored as + // `charge_automatically` stops silently auto-charging the card. The + // reverse direction is deliberately not applied: it would undo the send + // workflow's switch to `send_invoice` for tenants without a payment + // method on file. + if collection_method == stripe::CollectionMethod::SendInvoice + && existing.collection_method != Some(collection_method) + { + #[derive(serde::Serialize)] + struct UpdateInvoice { + collection_method: stripe::CollectionMethod, + due_date: i64, + } + tracing::debug!( + "Reconciling collection method of invoice {id} to {collection_method:?}", + id = existing_id.to_string() + ); + let updated: stripe::Invoice = client + .post_form( + &format!("/invoices/{existing_id}"), + UpdateInvoice { + collection_method, + due_date: due_date.expect("send_invoice mode always sets a due date"), + }, + ) + .await + .context("Reconciling collection method of existing invoice")?; + updated + } else { + existing + } } else { let description_text = format!( "Your Flow bill for the billing period between {date_start_human} - {date_end_human}. Tenant: {tenant}", @@ -489,17 +619,9 @@ impl Invoice { client, stripe::CreateInvoice { customer: Some(customer.id.to_owned()), - due_date: match mode { - ChargeType::SendInvoice => { - Some((Utc::now() + Duration::days(30)).timestamp()) - } - ChargeType::AutoCharge => None, - }, + due_date, description: Some(description_text.as_str()), - collection_method: Some(match mode { - ChargeType::AutoCharge => stripe::CollectionMethod::ChargeAutomatically, - ChargeType::SendInvoice => stripe::CollectionMethod::SendInvoice, - }), + collection_method: Some(collection_method), auto_advance: Some(false), custom_fields: Some(vec![ stripe::CreateInvoiceCustomFields { @@ -596,6 +718,8 @@ impl Invoice { if is_update { Ok(InvoiceResult::Updated) + } else if replace.is_some() { + Ok(InvoiceResult::Replaced) } else { Ok(InvoiceResult::Created(self.payment_provider)) } @@ -635,18 +759,11 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { line_items as "line_items!: sqlx::types::Json>", subtotal::bigint as "subtotal!", extra as "extra: sqlx::types::Json>", - customer.has_payment_method as has_payment_method, coalesce(dataflow.has_full_pipeline, false) as "has_full_pipeline!", tenants.payment_provider as "payment_provider!: PaymentProvider", tenants.trial_start as tenant_trial_start from invoices_ext left join tenants on tenants.tenant = billed_prefix - left join lateral( - select bool_or("invoice_settings/default_payment_method" is not null) as has_payment_method - from stripe.customers - where customers.metadata->>'estuary.dev/tenant_name' = billed_prefix - group by billed_prefix - ) as customer on true left join lateral( select sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0 @@ -686,18 +803,11 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { line_items as "line_items!: sqlx::types::Json>", subtotal::bigint as "subtotal!", extra as "extra: sqlx::types::Json>", - customer.has_payment_method as has_payment_method, coalesce(dataflow.has_full_pipeline, false) as "has_full_pipeline!", tenants.payment_provider as "payment_provider!: PaymentProvider", tenants.trial_start as tenant_trial_start from invoices_ext left join tenants on tenants.tenant = billed_prefix - left join lateral( - select bool_or("invoice_settings/default_payment_method" is not null) as has_payment_method - from stripe.customers - where customers.metadata->>'estuary.dev/tenant_name' = billed_prefix - group by billed_prefix - ) as customer on true left join lateral( select sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0 @@ -732,25 +842,19 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { .or_default() += 1; }); + let usage = invoice_type_counter + .remove(&InvoiceType::Final) + .unwrap_or_default(); + let manual = invoice_type_counter + .remove(&InvoiceType::Manual) + .unwrap_or_default(); if cmd.dry_run { tracing::info!( - "[DRY RUN] Classifying {usage} usage-based invoices and {manual} manually-entered invoices without making any changes to Stripe.", - usage = invoice_type_counter - .remove(&InvoiceType::Final) - .unwrap_or_default(), - manual = invoice_type_counter - .remove(&InvoiceType::Manual) - .unwrap_or_default(), + "[dry-run] Classifying {usage} usage-based invoices and {manual} manually-entered invoices without making any changes to Stripe." ); } else { tracing::info!( - "Processing {usage} usage-based invoices, and {manual} manually-entered invoices.", - usage = invoice_type_counter - .remove(&InvoiceType::Final) - .unwrap_or_default(), - manual = invoice_type_counter - .remove(&InvoiceType::Manual) - .unwrap_or_default(), + "Processing {usage} usage-based invoices, and {manual} manually-entered invoices." ); } @@ -760,127 +864,14 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { let client = stripe_client.clone(); let db_pool = db_pool.clone(); - let annotation = match response.invoice_type { - InvoiceType::Manual => Some(format!( - "[manual: {} - {}]", - response.date_start.format("%Y-%m-%d"), - response.date_end.format("%Y-%m-%d") - )), - _ => None, - }; - async move { - let action = response - .classify(&client, cmd.recreate_finalized) - .await; - - match action { - Err(err) => { - let formatted = format!( - "Error classifying {invoice_type:?} invoice for {tenant}", - tenant = response.billed_prefix, - invoice_type = response.invoice_type - ); - Err(anyhow::anyhow!("{formatted}: {err:#}")) - } - Ok(InvoiceAction::Skip { result, customer }) => { - tracing::debug!( - tenant = response.billed_prefix, - invoice_type = format!("{:?}", response.invoice_type), - subtotal = format!("${:.2}", response.subtotal as f64 / 100.0), - "{}", - result.message(cmd.dry_run) - ); - - if cmd.clean_up { - let task_res: Result<(), anyhow::Error> = async { - let customer = match customer { - Some(c) => c, - None => return Ok(()), - }; - let customer_id = customer.id.to_string(); - - if let Some(invoice) = - response.get_stripe_invoice(&client, &customer_id).await? - { - if let Some(InvoiceStatus::Draft) = invoice.status { - if cmd.dry_run { - tracing::warn!( - tenant = response.billed_prefix.to_string(), - "[dry-run] Would delete stale draft invoice {}", - invoice.id - ); - } else { - tracing::warn!( - tenant = response.billed_prefix.to_string(), - "Deleting draft invoice!" - ); - stripe::Invoice::delete(&client, &invoice.id).await?; - } - } - } - Ok(()) - } - .await; - - if let Err(e) = task_res { - tracing::warn!("Failed to check for or clear potential leaked draft invoices for {}, this is probably not a problem: {e:#}", response.billed_prefix.to_owned()); - } - } - - Ok((result, response.subtotal, response.billed_prefix.to_owned(), annotation)) - } - Ok(action) if cmd.dry_run => { - let result = match &action { - InvoiceAction::Create { replace: Some(id), .. } => { - tracing::info!( - tenant = response.billed_prefix, - "[dry-run] Would delete existing invoice {} and recreate", - id - ); - InvoiceResult::Created(response.payment_provider) - } - InvoiceAction::Create { .. } => { - InvoiceResult::Created(response.payment_provider) - } - InvoiceAction::Update { .. } => InvoiceResult::Updated, - InvoiceAction::Skip { .. } => unreachable!(), - }; - tracing::debug!( - tenant = response.billed_prefix, - invoice_type = format!("{:?}", response.invoice_type), - subtotal = format!("${:.2}", response.subtotal as f64 / 100.0), - "[dry-run] {}", - result.message(cmd.dry_run) - ); - Ok((result, response.subtotal, response.billed_prefix.to_owned(), annotation)) - } - Ok(action) => { - let res = response - .execute(&client, &db_pool, action, cmd.charge_type) - .await; - match res { - Err(err) => { - let formatted = format!( - "Error publishing {invoice_type:?} invoice for {tenant}", - tenant = response.billed_prefix, - invoice_type = response.invoice_type - ); - Err(anyhow::anyhow!("{formatted}: {err:#}")) - } - Ok(res) => { - tracing::debug!( - tenant = response.billed_prefix, - invoice_type = format!("{:?}", response.invoice_type), - subtotal = format!("${:.2}", response.subtotal as f64 / 100.0), - "{}", - res.message(cmd.dry_run) - ); - Ok((res, response.subtotal, response.billed_prefix.to_owned(), annotation)) - } - } - } - } + let result = process_invoice(cmd, &client, &db_pool, response).await?; + anyhow::Ok(( + result, + response.subtotal, + response.billed_prefix.to_owned(), + manual_annotation(response), + )) } .boxed() .map_err(|e| (e, response.clone())) @@ -895,7 +886,8 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { .or_else(|(err, invoice)| async move { if !cmd.fail_fast { tracing::error!("[{}]: {err:#}", invoice.billed_prefix); - Ok((InvoiceResult::Error, 0, invoice.billed_prefix, None)) + let annotation = manual_annotation(&invoice); + Ok((InvoiceResult::Error, 0, invoice.billed_prefix, annotation)) } else { Err(err) } @@ -926,7 +918,7 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { *subtotal_agg as f64 / 100.0 ); let limit = match status { - InvoiceResult::Created(_) | InvoiceResult::Updated => 9999, + InvoiceResult::Created(_) | InvoiceResult::Updated | InvoiceResult::Replaced => 9999, InvoiceResult::NoDataMoved | InvoiceResult::NoFullPipeline | InvoiceResult::LessThanMinimum @@ -957,13 +949,162 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { Ok(()) } +/// Process a single invoice end-to-end: classify it, then execute the action +/// (or preview it under --dry-run), returning the summary bucket for the run +/// report. Skips optionally clean up a stale draft left behind in Stripe. +async fn process_invoice( + cmd: &PublishInvoice, + client: &stripe::Client, + db_pool: &Pool, + invoice: &Invoice, +) -> anyhow::Result { + let action = invoice + .classify(client, cmd.recreate_finalized) + .await + .with_context(|| { + format!( + "Error classifying {invoice_type:?} invoice for {tenant}", + tenant = invoice.billed_prefix, + invoice_type = invoice.invoice_type + ) + })?; + + // A real run resolves a billing email inside execute(), which bails when + // none can be found. Mirror that read-only so the preview reports such + // tenants as errors instead of claiming the invoice would be published. + // (The invoice total cannot be reconciled without creating the invoice in + // Stripe, so that check is necessarily skipped.) + if cmd.dry_run && !matches!(action, InvoiceAction::Skip { .. }) { + if !billing_email_available(db_pool, action.customer(), &invoice.billed_prefix).await? { + let reason = "no customer email, tenants.billing_email, or admin user to invoice"; + tracing::warn!( + tenant = invoice.billed_prefix, + "[dry-run] Would fail: {reason}" + ); + if cmd.fail_fast { + bail!( + "[dry-run] Would fail to publish invoice for {tenant}: {reason}", + tenant = invoice.billed_prefix + ); + } + return Ok(InvoiceResult::Error); + } + } + + let result = match (action, cmd.dry_run) { + (InvoiceAction::Skip { result, customer }, _) => { + if cmd.clean_up { + clean_up_stale_draft(client, invoice, customer, cmd.dry_run).await; + } + result + } + ( + InvoiceAction::Create { + replace: Some(id), .. + }, + true, + ) => { + tracing::info!( + tenant = invoice.billed_prefix, + "[dry-run] Would replace existing invoice {}", + id + ); + InvoiceResult::Replaced + } + (InvoiceAction::Create { .. }, true) => InvoiceResult::Created(invoice.payment_provider), + (InvoiceAction::Update { .. }, true) => InvoiceResult::Updated, + (action, false) => invoice + .execute(client, db_pool, action, cmd.charge_type) + .await + .with_context(|| { + format!( + "Error publishing {invoice_type:?} invoice for {tenant}", + tenant = invoice.billed_prefix, + invoice_type = invoice.invoice_type + ) + })?, + }; + + tracing::debug!( + tenant = invoice.billed_prefix, + invoice_type = format!("{:?}", invoice.invoice_type), + subtotal = format!("${:.2}", invoice.subtotal as f64 / 100.0), + "{}", + result.message(cmd.dry_run) + ); + Ok(result) +} + +/// Best-effort removal of a stale draft invoice for a bill that was skipped: +/// checks Stripe for a matching draft and deletes it (or reports it under +/// --dry-run). Failures are logged rather than failing the run. +async fn clean_up_stale_draft( + client: &stripe::Client, + invoice: &Invoice, + customer: Option, + dry_run: bool, +) { + let task_res: anyhow::Result<()> = async { + // Locally resolved skips do not carry a Stripe customer, so find it + // here when cleaning up stale drafts. + let customer = match customer { + Some(c) => c, + None => match find_customer(client, &invoice.billed_prefix).await? { + Some(c) => c, + None => return Ok(()), + }, + }; + let customer_id = customer.id.to_string(); + + if let Some(stripe_invoice) = invoice.get_stripe_invoice(client, &customer_id).await? { + if let Some(stripe::InvoiceStatus::Draft) = stripe_invoice.status { + if dry_run { + tracing::warn!( + tenant = invoice.billed_prefix.to_string(), + "[dry-run] Would delete stale draft invoice {}", + stripe_invoice.id + ); + } else { + tracing::warn!( + tenant = invoice.billed_prefix.to_string(), + "Deleting draft invoice!" + ); + stripe::Invoice::delete(client, &stripe_invoice.id).await?; + } + } + } + Ok(()) + } + .await; + + if let Err(e) = task_res { + tracing::warn!( + "Failed to check for or clear potential leaked draft invoices for {}, this is probably not a problem: {e:#}", + invoice.billed_prefix.to_owned() + ); + } +} + +/// Display tag distinguishing a manually-entered invoice (and its period) in the +/// run summary; usage invoices get no tag. +fn manual_annotation(invoice: &Invoice) -> Option { + match invoice.invoice_type { + InvoiceType::Manual => Some(format!( + "[manual: {} - {}]", + invoice.date_start.format("%Y-%m-%d"), + invoice.date_end.format("%Y-%m-%d") + )), + _ => None, + } +} + /// Read-only: search Stripe for an existing customer by tenant metadata. #[tracing::instrument(skip(client))] async fn find_customer( client: &stripe::Client, tenant: &str, ) -> anyhow::Result> { - let customers: Vec = stripe_search( + let mut customers: Vec = stripe_search( client, "customers", SearchParams { @@ -974,7 +1115,7 @@ async fn find_customer( .await .context(format!("Searching for tenant {tenant}"))?; - if let Some(customer) = customers.into_iter().next() { + if let Some(customer) = customers.drain(..).next() { tracing::debug!("Found existing customer {id}", id = customer.id.to_string()); Ok(Some(customer)) } else { @@ -983,22 +1124,24 @@ async fn find_customer( } /// Ensures a Stripe customer exists for this tenant and is ready for invoicing. -/// Finds an existing customer or creates a new one, then ensures the customer -/// has an email set (looking up the earliest admin on the tenant if needed). -#[tracing::instrument(skip(client, db_client))] +/// Uses `found` (the customer located during classification) when present, +/// otherwise finds an existing customer or creates a new one (carrying the +/// tenant's billing email and address from the control-plane DB), then ensures +/// the customer has an email set (falling back to the earliest admin on the +/// tenant if needed). +#[tracing::instrument(skip(client, db_client, found))] async fn ensure_customer_for_invoicing( client: &stripe::Client, db_client: &Pool, tenant: &str, + found: Option, ) -> anyhow::Result { - let billing_row = sqlx::query!( - r#"SELECT billing_email, billing_address FROM tenants WHERE tenant = $1"#, - tenant, - ) - .fetch_optional(db_client) - .await?; + let existing = match found { + Some(customer) => Some(customer), + None => find_customer(client, tenant).await?, + }; - let customer = if let Some(customer) = find_customer(client, tenant).await? { + let customer = if let Some(customer) = existing { customer } else { tracing::debug!("Creating new customer"); @@ -1011,6 +1154,7 @@ async fn ensure_customer_for_invoicing( customer_create_idempotency_key(tenant), )); + let billing_row = tenant_billing_row(db_client, tenant).await?; let billing_email = billing_row .as_ref() .and_then(|r| r.billing_email.as_deref()); @@ -1051,22 +1195,71 @@ async fn ensure_customer_for_invoicing( }; if customer.email.is_none() { - let db_email = billing_row.as_ref().and_then(|r| r.billing_email.clone()); + let db_email = tenant_billing_row(db_client, tenant) + .await? + .and_then(|r| r.billing_email); - if let Some(email) = db_email { - tracing::info!("Using billing_email from tenants table: {email}"); - stripe::Customer::update( - client, - &customer.id, - stripe::UpdateCustomer { - email: Some(&email), - ..Default::default() - }, - ) - .await?; - } else { - let responses = sqlx::query!( - r#" + let email = match db_email { + Some(email) => { + tracing::info!("Using billing_email from tenants table: {email}"); + email + } + None => match earliest_admin_email(db_client, tenant).await? { + Some(email) => { + tracing::warn!( + "Stripe customer object is missing an email. Going with {email}, an admin on that tenant." + ); + email + } + None => bail!( + "Stripe customer object is missing an email. No admins found for tenant {tenant}, unable to create invoice without email. Skipping" + ), + }, + }; + + stripe::Customer::update( + client, + &customer.id, + stripe::UpdateCustomer { + email: Some(&email), + ..Default::default() + }, + ) + .await?; + } + Ok(customer) +} + +/// A tenant's billing contact info from the control-plane DB. +struct TenantBillingRow { + billing_email: Option, + billing_address: Option, +} + +/// Read-only: the tenant's billing contact row, if the tenant exists. +async fn tenant_billing_row( + db_client: &Pool, + tenant: &str, +) -> anyhow::Result> { + Ok(sqlx::query_as!( + TenantBillingRow, + r#"SELECT billing_email, billing_address FROM tenants WHERE tenant = $1"#, + tenant, + ) + .fetch_optional(db_client) + .await?) +} + +/// Read-only: the email of the earliest-created admin user on the tenant, if any. +async fn earliest_admin_email( + db_client: &Pool, + tenant: &str, +) -> anyhow::Result> { + // NOTE: the SQL text below is intentionally indented to stay byte-identical + // to the query this was extracted from, so it keeps hitting the same + // offline sqlx cache entry. + let responses = sqlx::query!( + r#" select users.email as email from user_grants join auth.users as users on user_grants.user_id = users.id @@ -1079,35 +1272,35 @@ async fn ensure_customer_for_invoicing( ) order by users.created_at asc "#, - tenant - ) - .fetch_all(db_client) - .await?; + tenant + ) + .fetch_all(db_client) + .await?; - if let Some(email) = responses - .iter() - .find_map(|response| response.email.to_owned()) - { - tracing::warn!( - "Stripe customer object is missing an email. Going with {email}, an admin on that tenant." - ); - stripe::Customer::update( - client, - &customer.id, - stripe::UpdateCustomer { - email: Some(&email), - ..Default::default() - }, - ) - .await?; - } else { - bail!( - "Stripe customer object is missing an email. No admins found for tenant {tenant}, unable to create invoice without email. Found users: {found:?} Skipping", - found = responses, - tenant = tenant - ); - } - } + Ok(responses.into_iter().find_map(|r| r.email)) +} + +/// Read-only mirror of the email requirement enforced by +/// `ensure_customer_for_invoicing`: the classified customer's email, else +/// `tenants.billing_email`, else the earliest tenant admin. Returns whether an +/// email could be resolved without performing any writes, so --dry-run can +/// report tenants that would fail to invoice rather than over-reporting success. +async fn billing_email_available( + db_client: &Pool, + customer: Option<&stripe::Customer>, + tenant: &str, +) -> anyhow::Result { + if customer.is_some_and(|c| c.email.is_some()) { + return Ok(true); } - Ok(customer) + + if tenant_billing_row(db_client, tenant) + .await? + .and_then(|r| r.billing_email) + .is_some() + { + return Ok(true); + } + + Ok(earliest_admin_email(db_client, tenant).await?.is_some()) } diff --git a/crates/billing-integrations/src/send.rs b/crates/billing-integrations/src/send.rs index 7fb1f58221c..8dcbe842fac 100644 --- a/crates/billing-integrations/src/send.rs +++ b/crates/billing-integrations/src/send.rs @@ -250,6 +250,9 @@ async fn update_collection_methods( let pb = ProgressBar::new(invoices.len() as u64); pb.set_message("updating collection method"); pb.set_style(ProgressStyle::with_template(PROGRESS_BAR_TEMPLATE).unwrap()); + // Continue past per-invoice failures: the following invoice table shows the + // true state of every draft (failed switches stay flagged as missing a + // payment method) before the finalize prompt, so the operator can abort. for inv in invoices { let res: Result = stripe_client .post_form( @@ -262,7 +265,7 @@ async fn update_collection_methods( .await; match res { Ok(_) => { - inv.collection_method = Some(method.clone()); + inv.collection_method = Some(method); } Err(e) => { pb.println(format!( diff --git a/crates/billing-integrations/src/stripe_utils.rs b/crates/billing-integrations/src/stripe_utils.rs index f93b50758b7..37bb20ea0ca 100644 --- a/crates/billing-integrations/src/stripe_utils.rs +++ b/crates/billing-integrations/src/stripe_utils.rs @@ -89,7 +89,7 @@ impl Invoice { .metadata .as_ref() .and_then(InvoiceMetadata::from_metadata_map) - .map_or(false, |m| m.invoice_type == InvoiceType::Manual) + .is_some_and(|m| m.invoice_type == InvoiceType::Manual) } pub fn period_start(&self) -> Option { From 757f2b27252954e4bf2c4a5b54477cb51293831a Mon Sep 17 00:00:00 2001 From: Joseph Shearer Date: Mon, 3 Aug 2026 15:26:47 -0400 Subject: [PATCH 4/6] billing: preserve tenant filtering for invoice reads --- ...087811da95a1ea41930b4bac33ad42a66cabec3557fd929b.json} | 4 ++-- ...bcbf5e0248595385a5a0c05d17a27672e81183d706b8afa5.json} | 4 ++-- crates/billing-integrations/src/publish.rs | 8 ++++---- 3 files changed, 8 insertions(+), 8 deletions(-) rename .sqlx/{query-402d978379e462304be791a274e52cba7a2db38f044b2c229c90cd6908e9c750.json => query-49dafee439e40ad5087811da95a1ea41930b4bac33ad42a66cabec3557fd929b.json} (59%) rename .sqlx/{query-c586f9fd13f5c26682d5516ab985c4507359f243a2307e0ff08c5a83b9a01a05.json => query-86637e47572df3efbcbf5e0248595385a5a0c05d17a27672e81183d706b8afa5.json} (60%) diff --git a/.sqlx/query-402d978379e462304be791a274e52cba7a2db38f044b2c229c90cd6908e9c750.json b/.sqlx/query-49dafee439e40ad5087811da95a1ea41930b4bac33ad42a66cabec3557fd929b.json similarity index 59% rename from .sqlx/query-402d978379e462304be791a274e52cba7a2db38f044b2c229c90cd6908e9c750.json rename to .sqlx/query-49dafee439e40ad5087811da95a1ea41930b4bac33ad42a66cabec3557fd929b.json index 033a7905d4a..f7c277f851e 100644 --- a/.sqlx/query-402d978379e462304be791a274e52cba7a2db38f044b2c229c90cd6908e9c750.json +++ b/.sqlx/query-49dafee439e40ad5087811da95a1ea41930b4bac33ad42a66cabec3557fd929b.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n select\n date_start as \"date_start!\",\n date_end as \"date_end!\",\n billed_prefix as \"billed_prefix!\",\n invoice_type as \"invoice_type!: InvoiceType\",\n line_items as \"line_items!: sqlx::types::Json>\",\n subtotal::bigint as \"subtotal!\",\n extra as \"extra: sqlx::types::Json>\",\n coalesce(dataflow.has_full_pipeline, false) as \"has_full_pipeline!\",\n tenants.payment_provider as \"payment_provider!: PaymentProvider\",\n tenants.trial_start as tenant_trial_start\n from invoices_ext\n left join tenants on tenants.tenant = billed_prefix\n left join lateral(\n \tselect\n \t\tsum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0\n \t\tand sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'materialization') > 0\n \t\tas has_full_pipeline\n from catalog_stats_monthly\n join live_specs on live_specs.catalog_name ^@ catalog_stats_monthly.catalog_name\n where\n \tcatalog_stats_monthly.catalog_name = billed_prefix\n \tand tstzrange(date_trunc('day', $1::date), date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day') @> catalog_stats_monthly.ts\n ) as dataflow on true\n where ((\n date_start >= date_trunc('day', $1::date)\n and date_end <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and invoice_type = 'final'\n ) or (\n date_start <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and date_end >= date_trunc('day', $1::date)\n and invoice_type = 'manual'\n ))\n and billed_prefix = any($2)\n ", + "query": "\n select\n date_start as \"date_start!\",\n date_end as \"date_end!\",\n billed_prefix as \"billed_prefix!\",\n invoice_type as \"invoice_type!: InvoiceType\",\n line_items as \"line_items!: sqlx::types::Json>\",\n subtotal::bigint as \"subtotal!\",\n extra as \"extra: sqlx::types::Json>\",\n coalesce(dataflow.has_full_pipeline, false) as \"has_full_pipeline!\",\n tenants.payment_provider as \"payment_provider!: PaymentProvider\",\n tenants.trial_start as tenant_trial_start\n from internal.invoices_ext\n join tenants on tenants.tenant = billed_prefix\n left join lateral(\n \tselect\n \t\tsum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0\n \t\tand sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'materialization') > 0\n \t\tas has_full_pipeline\n from catalog_stats_monthly\n join live_specs on live_specs.catalog_name ^@ catalog_stats_monthly.catalog_name\n where\n \tcatalog_stats_monthly.catalog_name = billed_prefix\n \tand tstzrange(date_trunc('day', $1::date), date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day') @> catalog_stats_monthly.ts\n ) as dataflow on true\n where ((\n date_start >= date_trunc('day', $1::date)\n and date_end <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and invoice_type = 'final'\n ) or (\n date_start <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and date_end >= date_trunc('day', $1::date)\n and invoice_type = 'manual'\n ))\n and billed_prefix = any($2)\n ", "describe": { "columns": [ { @@ -83,5 +83,5 @@ true ] }, - "hash": "402d978379e462304be791a274e52cba7a2db38f044b2c229c90cd6908e9c750" + "hash": "49dafee439e40ad5087811da95a1ea41930b4bac33ad42a66cabec3557fd929b" } diff --git a/.sqlx/query-c586f9fd13f5c26682d5516ab985c4507359f243a2307e0ff08c5a83b9a01a05.json b/.sqlx/query-86637e47572df3efbcbf5e0248595385a5a0c05d17a27672e81183d706b8afa5.json similarity index 60% rename from .sqlx/query-c586f9fd13f5c26682d5516ab985c4507359f243a2307e0ff08c5a83b9a01a05.json rename to .sqlx/query-86637e47572df3efbcbf5e0248595385a5a0c05d17a27672e81183d706b8afa5.json index c00968f1d4e..d698b34bd8b 100644 --- a/.sqlx/query-c586f9fd13f5c26682d5516ab985c4507359f243a2307e0ff08c5a83b9a01a05.json +++ b/.sqlx/query-86637e47572df3efbcbf5e0248595385a5a0c05d17a27672e81183d706b8afa5.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n select\n date_start as \"date_start!\",\n date_end as \"date_end!\",\n billed_prefix as \"billed_prefix!\",\n invoice_type as \"invoice_type!: InvoiceType\",\n line_items as \"line_items!: sqlx::types::Json>\",\n subtotal::bigint as \"subtotal!\",\n extra as \"extra: sqlx::types::Json>\",\n coalesce(dataflow.has_full_pipeline, false) as \"has_full_pipeline!\",\n tenants.payment_provider as \"payment_provider!: PaymentProvider\",\n tenants.trial_start as tenant_trial_start\n from invoices_ext\n left join tenants on tenants.tenant = billed_prefix\n left join lateral(\n \tselect\n \t\tsum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0\n \t\tand sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'materialization') > 0\n \t\tas has_full_pipeline\n from catalog_stats_monthly\n join live_specs on live_specs.catalog_name ^@ catalog_stats_monthly.catalog_name\n where\n \tcatalog_stats_monthly.catalog_name = billed_prefix\n \tand tstzrange(date_trunc('day', $1::date), date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day') @> catalog_stats_monthly.ts\n ) as dataflow on true\n where (\n date_start >= date_trunc('day', $1::date)\n and date_end <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and invoice_type = 'final'\n ) or (\n date_start <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and date_end >= date_trunc('day', $1::date)\n and invoice_type = 'manual'\n )\n ", + "query": "\n select\n date_start as \"date_start!\",\n date_end as \"date_end!\",\n billed_prefix as \"billed_prefix!\",\n invoice_type as \"invoice_type!: InvoiceType\",\n line_items as \"line_items!: sqlx::types::Json>\",\n subtotal::bigint as \"subtotal!\",\n extra as \"extra: sqlx::types::Json>\",\n coalesce(dataflow.has_full_pipeline, false) as \"has_full_pipeline!\",\n tenants.payment_provider as \"payment_provider!: PaymentProvider\",\n tenants.trial_start as tenant_trial_start\n from internal.invoices_ext\n join tenants on tenants.tenant = billed_prefix\n left join lateral(\n \tselect\n \t\tsum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0\n \t\tand sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'materialization') > 0\n \t\tas has_full_pipeline\n from catalog_stats_monthly\n join live_specs on live_specs.catalog_name ^@ catalog_stats_monthly.catalog_name\n where\n \tcatalog_stats_monthly.catalog_name = billed_prefix\n \tand tstzrange(date_trunc('day', $1::date), date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day') @> catalog_stats_monthly.ts\n ) as dataflow on true\n where (\n date_start >= date_trunc('day', $1::date)\n and date_end <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and invoice_type = 'final'\n ) or (\n date_start <= date_trunc('day', ($1::date)) + interval '1 month' - interval '1 day'\n and date_end >= date_trunc('day', $1::date)\n and invoice_type = 'manual'\n )\n ", "describe": { "columns": [ { @@ -82,5 +82,5 @@ true ] }, - "hash": "c586f9fd13f5c26682d5516ab985c4507359f243a2307e0ff08c5a83b9a01a05" + "hash": "86637e47572df3efbcbf5e0248595385a5a0c05d17a27672e81183d706b8afa5" } diff --git a/crates/billing-integrations/src/publish.rs b/crates/billing-integrations/src/publish.rs index 28e06692bf3..0a2c86e9c25 100644 --- a/crates/billing-integrations/src/publish.rs +++ b/crates/billing-integrations/src/publish.rs @@ -762,8 +762,8 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { coalesce(dataflow.has_full_pipeline, false) as "has_full_pipeline!", tenants.payment_provider as "payment_provider!: PaymentProvider", tenants.trial_start as tenant_trial_start - from invoices_ext - left join tenants on tenants.tenant = billed_prefix + from internal.invoices_ext + join tenants on tenants.tenant = billed_prefix left join lateral( select sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0 @@ -806,8 +806,8 @@ pub async fn do_publish_invoices(cmd: &PublishInvoice) -> anyhow::Result<()> { coalesce(dataflow.has_full_pipeline, false) as "has_full_pipeline!", tenants.payment_provider as "payment_provider!: PaymentProvider", tenants.trial_start as tenant_trial_start - from invoices_ext - left join tenants on tenants.tenant = billed_prefix + from internal.invoices_ext + join tenants on tenants.tenant = billed_prefix left join lateral( select sum(catalog_stats_monthly.usage_seconds) filter (where live_specs.spec_type = 'capture') > 0 From afbabbe7aa3206ea812e8a24a4a468851fdb0ab8 Mon Sep 17 00:00:00 2001 From: Joseph Shearer Date: Mon, 21 Sep 2026 16:44:39 -0400 Subject: [PATCH 5/6] billing: guard invoice finalization and clarify CLI options --- Cargo.lock | 1 + crates/billing-integrations/Cargo.toml | 3 + crates/billing-integrations/README.md | 18 ++ crates/billing-integrations/src/lib.rs | 65 +++++ crates/billing-integrations/src/publish.rs | 38 ++- crates/billing-integrations/src/send.rs | 291 ++++++++++++++++++--- 6 files changed, 353 insertions(+), 63 deletions(-) create mode 100644 crates/billing-integrations/README.md diff --git a/Cargo.lock b/Cargo.lock index 221dcc01b3b..4989d3bdfe6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1344,6 +1344,7 @@ version = "0.0.0" dependencies = [ "anyhow", "async-stripe", + "axum", "billing-types", "chrono", "clap", diff --git a/crates/billing-integrations/Cargo.toml b/crates/billing-integrations/Cargo.toml index 96bf62c4d2b..d1581b590b5 100644 --- a/crates/billing-integrations/Cargo.toml +++ b/crates/billing-integrations/Cargo.toml @@ -28,3 +28,6 @@ tracing-subscriber = { workspace = true } num-format = { workspace = true } itertools = { workspace = true } indicatif = { workspace = true } + +[dev-dependencies] +axum = { workspace = true } diff --git a/crates/billing-integrations/README.md b/crates/billing-integrations/README.md new file mode 100644 index 00000000000..8704c5815bd --- /dev/null +++ b/crates/billing-integrations/README.md @@ -0,0 +1,18 @@ +# Billing integrations + +Operator CLI for publishing control-plane bills to Stripe and approving collection. +`lib.rs` defines the commands and shared calendar-month parser. + +- `publish.rs` classifies bills and creates or refreshes draft invoices. `--dry-run` + performs read-only classification against Stripe. Manual bills always use + `send_invoice`. `--recreate-existing` (alias `--recreate-finalized`) replaces + eligible usage invoices; it does not reissue manual bills. +- `send.rs` checks current Stripe state before approval and again before + finalization. Usage drafts without a payment method can switch to Net 30 with + operator approval. Failed switches are excluded from finalization; incorrectly + configured manual drafts require republishing. A changed collection decision + after approval requires another send run. +- `stripe_utils.rs` wraps Stripe invoices for metadata access and operator tables. + +Both commands accept `--month YYYY-MM` or the legacy `YYYY-MM-01` form. +Run the focused checks with `mise exec -- cargo nextest run -p billing-integrations`. diff --git a/crates/billing-integrations/src/lib.rs b/crates/billing-integrations/src/lib.rs index a4048d74c09..a8786086571 100644 --- a/crates/billing-integrations/src/lib.rs +++ b/crates/billing-integrations/src/lib.rs @@ -28,3 +28,68 @@ impl Cli { } } } + +fn parse_month(arg: &str) -> anyhow::Result { + use chrono::Datelike; + let date = chrono::NaiveDate::parse_from_str( + &if arg.len() == 7 { + format!("{arg}-01") + } else { + arg.to_owned() + }, + "%Y-%m-%d", + )?; + anyhow::ensure!(date.day() == 1, "expected YYYY-MM or YYYY-MM-01"); + Ok(date) +} + +#[cfg(test)] +mod tests { + use clap::Parser; + + #[test] + fn billing_month_and_recreate_alias() { + assert_eq!( + super::parse_month("2026-08").unwrap(), + chrono::NaiveDate::from_ymd_opt(2026, 8, 1).unwrap() + ); + for command in ["publish-invoices", "send-invoices"] { + for month in ["2026-08", "2026-08-01", "2026-08-15", "2026-13", "invalid"] { + let mut args = vec![ + "billing-integrations", + command, + "--stripe-api-key", + "test", + "--all-tenants", + "--month", + month, + ]; + if command == "publish-invoices" { + args.extend(["--connection-string", "test"]); + } + assert_eq!( + super::Cli::try_parse_from(args).is_ok(), + matches!(month, "2026-08" | "2026-08-01"), + "{command}: {month}" + ); + } + } + for flag in ["--recreate-existing", "--recreate-finalized"] { + assert!( + super::Cli::try_parse_from([ + "billing-integrations", + "publish-invoices", + "--stripe-api-key", + "test", + "--connection-string", + "test", + "--all-tenants", + "--month", + "2026-08", + flag + ]) + .is_ok() + ); + } + } +} diff --git a/crates/billing-integrations/src/publish.rs b/crates/billing-integrations/src/publish.rs index 0a2c86e9c25..66e70b50066 100644 --- a/crates/billing-integrations/src/publish.rs +++ b/crates/billing-integrations/src/publish.rs @@ -3,7 +3,7 @@ use billing_types::{ InvoiceMetadata, InvoiceSearch, InvoiceType, PaymentProvider, SearchParams, customer_create_idempotency_key, customer_search_query, stripe_search, tenant_metadata, }; -use chrono::{Duration, ParseError, Utc}; +use chrono::{Duration, Utc}; use futures::{FutureExt, StreamExt, TryFutureExt, TryStreamExt}; use itertools::Itertools; use serde::{Deserialize, Serialize}; @@ -36,13 +36,13 @@ pub struct PublishInvoice { /// Generate invoices for all tenants that have bills in the provided month. #[clap(long, conflicts_with("tenants"))] all_tenants: bool, - /// The month to generate invoices for, in format "YYYY-MM-DD" - #[clap(long, value_parser = parse_date)] + /// The month to generate invoices for: YYYY-MM (also accepts YYYY-MM-01) + #[clap(long, value_parser = crate::parse_month)] month: NaiveDate, - /// Whether to recreate existing invoices: drafts are deleted and open - /// invoices are voided before the replacement is created. - #[clap(long)] - recreate_finalized: bool, + /// Recreate existing usage invoices: delete drafts or void open invoices. + /// Manual invoices are excluded; paid and uncollectible invoices are not replaced. + #[clap(long, alias = "recreate-finalized")] + recreate_existing: bool, /// Stop execution after first failure #[clap(long)] fail_fast: bool, @@ -69,10 +69,6 @@ pub struct PublishInvoice { pub dry_run: bool, } -fn parse_date(arg: &str) -> Result { - NaiveDate::parse_from_str(arg, "%Y-%m-%d") -} - #[derive(Serialize, Deserialize, Debug, Clone)] struct Extra { trial_start: Option, @@ -168,7 +164,7 @@ enum InvoiceAction { result: InvoiceResult, customer: Option, }, - /// Create a new invoice. `replace` is set when --recreate-finalized + /// Create a new invoice. `replace` is set when --recreate-existing /// requires deleting an existing invoice first. `customer` is the customer /// found during classification, or None when none exists yet (execute then /// creates one). @@ -237,7 +233,7 @@ impl Invoice { .await .context("Searching for an invoice")?; - // Prefer a live invoice over a voided one: after --recreate-finalized + // Prefer a live invoice over a voided one: after --recreate-existing // voids and recreates an invoice, both match this search, and // classification should act on the replacement. A voided invoice is // still returned when it's the only match, so that a voided manual @@ -258,7 +254,7 @@ impl Invoice { async fn classify( &self, client: &stripe::Client, - recreate_finalized: bool, + recreate_existing: bool, ) -> anyhow::Result { match (&self.invoice_type, &self.extra) { (InvoiceType::Preview, _) => { @@ -350,11 +346,11 @@ impl Invoice { .await? { match invoice.status { - // Manual invoices are excluded from --recreate-finalized: an + // Manual invoices are excluded from --recreate-existing: an // already-sent (open) manual invoice must not be voided and // reissued, and a manual draft is refreshed via the Update arm below. Some(stripe::InvoiceStatus::Open | stripe::InvoiceStatus::Draft) - if recreate_finalized && !matches!(self.invoice_type, InvoiceType::Manual) => + if recreate_existing && !matches!(self.invoice_type, InvoiceType::Manual) => { Ok(InvoiceAction::Create { replace: Some(invoice.id), @@ -363,12 +359,12 @@ impl Invoice { } // A voided invoice can be neither deleted nor updated, so treat // it as absent and create a fresh replacement. This also - // recovers a --recreate-finalized run that voided an open + // recovers a --recreate-existing run that voided an open // invoice but failed before creating its replacement. Without // the flag, the unsupported-state error below keeps a voided // invoice loud rather than silently re-billing the tenant. Some(stripe::InvoiceStatus::Void) - if recreate_finalized && !matches!(self.invoice_type, InvoiceType::Manual) => + if recreate_existing && !matches!(self.invoice_type, InvoiceType::Manual) => { tracing::warn!( "Found voided invoice {id}; treating it as absent", @@ -403,7 +399,7 @@ impl Invoice { } Some(stripe::InvoiceStatus::Open) => { bail!( - "Found open invoice {id}. Pass --recreate-finalized to void and recreate this invoice.", + "Found open invoice {id}. Pass --recreate-existing to void and recreate this invoice.", id = invoice.id.to_string() ) } @@ -494,7 +490,7 @@ impl Invoice { let date_start_repr = self.date_start.format("%F").to_string(); let date_end_repr = self.date_end.format("%F").to_string(); - // Remove the existing invoice if --recreate-finalized was used + // Remove the existing invoice if --recreate-existing was used if let Some(ref replace_id) = replace { // Re-verify the invoice status before removing it (guard against race conditions) let existing = stripe::Invoice::retrieve(client, replace_id, &[]).await?; @@ -959,7 +955,7 @@ async fn process_invoice( invoice: &Invoice, ) -> anyhow::Result { let action = invoice - .classify(client, cmd.recreate_finalized) + .classify(client, cmd.recreate_existing) .await .with_context(|| { format!( diff --git a/crates/billing-integrations/src/send.rs b/crates/billing-integrations/src/send.rs index 8dcbe842fac..31cc8cb6f40 100644 --- a/crates/billing-integrations/src/send.rs +++ b/crates/billing-integrations/src/send.rs @@ -1,4 +1,5 @@ use crate::stripe_utils::{Invoice, fetch_invoices}; +use anyhow::Context; use billing_types::{InvoiceSearch, InvoiceType, StatusFilter}; use chrono::{Datelike, Duration, NaiveDate, Utc}; use clap::Args; @@ -7,7 +8,7 @@ use indicatif::{ProgressBar, ProgressStyle}; use itertools::Itertools; use num_format::{Locale, ToFormattedString}; use std::collections::HashSet; -use stripe::{Client, FinalizeInvoiceParams, Invoice as StripeInvoice, InvoiceId}; +use stripe::{Client, FinalizeInvoiceParams, Invoice as StripeInvoice}; const PROGRESS_BAR_TEMPLATE: &str = "{spinner} [{elapsed_precise}] [{bar:40}] {pos}/{len} {msg}"; @@ -18,8 +19,8 @@ pub struct SendInvoices { /// Stripe API key. #[clap(long)] pub stripe_api_key: String, - /// The month to send invoices for, in format "YYYY-MM-DD" - #[clap(long)] + /// The month to send invoices for: YYYY-MM (also accepts YYYY-MM-01) + #[clap(long, value_parser = crate::parse_month)] pub month: NaiveDate, /// A list of tenants to exclude #[clap(long, value_delimiter = ',', conflicts_with = "tenants")] @@ -145,7 +146,9 @@ pub async fn do_send_invoices(cmd: &SendInvoices) -> anyhow::Result<()> { if !draft_invoices.is_empty() { // 2a. Update collection methods for any drafts that need it draft_invoices = update_draft_collection_methods(&stripe_client, draft_invoices).await?; + } + if !draft_invoices.is_empty() { print_invoice_table("Invoices to finalize", &draft_invoices); prompt_to_continue("Enter Y to finalize these invoices, or anything else to abort: ") .await?; @@ -183,18 +186,33 @@ async fn update_draft_collection_methods( stripe_client: &Client, mut to_update: Vec, ) -> anyhow::Result> { - // Identify invoices that need to be switched to `send_invoice`: - // - Manual invoices should always be sent as invoices, never auto-charged - // - Auto-charge invoices without a payment method on file must be sent as invoices - let needs_update: HashSet = to_update - .iter() - .filter(|inv| { - inv.collection_method().map_or(false, |cm| { - cm == stripe::CollectionMethod::ChargeAutomatically - }) && (inv.is_manual() || !inv.has_cc()) - }) - .map(|inv| inv.id().clone()) - .collect::>(); + // Search results and an earlier publish run may have stale payment-method state. + let mut refreshed = Vec::new(); + let mut needs_update = HashSet::new(); + for inv in to_update { + let result = async { + let current = Invoice::from( + StripeInvoice::retrieve(stripe_client, inv.id(), &["customer"]).await?, + ); + anyhow::ensure!( + current.status() == Some(stripe::InvoiceStatus::Draft), + "invoice is no longer a draft" + ); + let needs_update = collection_method_needs_update(¤t)?; + Ok::<_, anyhow::Error>((current, needs_update)) + } + .await; + match result { + Ok((current, update)) => { + if update { + needs_update.insert(current.id().clone()); + } + refreshed.push(current); + } + Err(error) => tracing::error!(invoice = %inv.id(), error = %error, "Skipping invoice"), + } + } + to_update = refreshed; // Modify the table row for those that need to be updated showing the transition let table_rows = to_update @@ -215,23 +233,18 @@ async fn update_draft_collection_methods( if !table_rows.is_empty() { let table = build_invoice_table(table_rows, None); println!( - "\nThe following draft invoices will updated to use the 'send_invoice' collection method:" + "\nThe following draft invoices will be updated to use the 'send_invoice' collection method:" ); println!("{}", table); prompt_to_continue("Enter Y to update collection methods, or anything else to abort: ") .await?; - let to_update = to_update - .iter_mut() - .filter(|inv| needs_update.contains(inv.id())) - .collect::>(); - update_collection_methods( - stripe_client, - to_update, - stripe::CollectionMethod::SendInvoice, - ) - .await?; + let (updates, mut unchanged): (Vec<_>, Vec<_>) = to_update + .into_iter() + .partition(|inv| needs_update.contains(inv.id())); + unchanged.extend(update_collection_methods(stripe_client, updates).await?); + to_update = unchanged; } Ok(to_update) @@ -239,9 +252,8 @@ async fn update_draft_collection_methods( async fn update_collection_methods( stripe_client: &Client, - invoices: Vec<&mut Invoice>, - method: stripe::CollectionMethod, -) -> anyhow::Result<()> { + invoices: Vec, +) -> anyhow::Result> { #[derive(serde::Serialize)] struct PostBody { collection_method: stripe::CollectionMethod, @@ -250,26 +262,22 @@ async fn update_collection_methods( let pb = ProgressBar::new(invoices.len() as u64); pb.set_message("updating collection method"); pb.set_style(ProgressStyle::with_template(PROGRESS_BAR_TEMPLATE).unwrap()); - // Continue past per-invoice failures: the following invoice table shows the - // true state of every draft (failed switches stay flagged as missing a - // payment method) before the finalize prompt, so the operator can abort. + let mut updated = Vec::new(); for inv in invoices { let res: Result = stripe_client .post_form( &format!("/invoices/{}", inv.id()), PostBody { - collection_method: method, + collection_method: stripe::CollectionMethod::SendInvoice, due_date: Some((Utc::now() + Duration::days(30)).timestamp()), }, ) .await; match res { - Ok(_) => { - inv.collection_method = Some(method); - } + Ok(invoice) => updated.push(Invoice::from(invoice)), Err(e) => { pb.println(format!( - "Failed to update collection method for invoice {}: {}", + "Skipping invoice {} after collection method update failed: {}", inv.id(), e )); @@ -278,7 +286,25 @@ async fn update_collection_methods( pb.inc(1); } pb.finish_with_message("Collection method update complete"); - Ok(()) + Ok(updated) +} + +// Missing payment methods are a late collection decision. Incorrect manual +// invoice configuration must be repaired by publish before operator approval. +fn collection_method_needs_update(invoice: &Invoice) -> anyhow::Result { + let method = invoice.collection_method()?; + anyhow::ensure!( + !invoice.is_manual() || method == stripe::CollectionMethod::SendInvoice, + "manual invoice must use send_invoice; rerun publish-invoices" + ); + if method == stripe::CollectionMethod::SendInvoice { + return Ok(false); + } + anyhow::ensure!( + invoice.customer().is_some(), + "missing expanded Stripe customer" + ); + Ok(!invoice.has_cc()) } // Finalizes the invoices and re-fetches them to ensure we have the correct state @@ -294,6 +320,26 @@ async fn finalize_invoices( let stripe_client = stripe_client; let pb = pb.clone(); async move { + // The operator can pause at the prompt. Re-read before enabling collection, + // and require another send run if the approved collection decision changed. + let current = Invoice::from( + StripeInvoice::retrieve(stripe_client, row.id(), &["customer"]) + .await + .with_context(|| { + format!("Refreshing invoice {} before finalization", row.id()) + })?, + ); + anyhow::ensure!( + current.status() == Some(stripe::InvoiceStatus::Draft), + "invoice {} is no longer a draft", + row.id() + ); + anyhow::ensure!( + current.collection_method()? == row.collection_method()? + && !collection_method_needs_update(¤t)?, + "invoice {} collection decision changed; rerun send-invoices", + row.id() + ); StripeInvoice::finalize( stripe_client, row.id(), @@ -407,12 +453,19 @@ async fn update_auto_advance(stripe_client: &Client, invoices: Vec) -> pb.set_style(ProgressStyle::with_template(PROGRESS_BAR_TEMPLATE).unwrap()); for inv in invoices { - let res: Result = stripe_client - .post_form( + let res: anyhow::Result = async { + let current = Invoice::from(StripeInvoice::retrieve(stripe_client, inv.id(), &["customer"]).await?); + anyhow::ensure!(current.status() == Some(stripe::InvoiceStatus::Open), "invoice is no longer open"); + anyhow::ensure!( + !collection_method_needs_update(¤t) + .context("Open invoice requires explicit correction in Stripe")?, + "open invoice needs collection-method correction in Stripe; leaving auto_advance disabled" + ); + stripe_client.post_form( &format!("/invoices/{}", inv.id()), UpdateAutoAdvance { auto_advance: true }, - ) - .await; + ).await.context("Enabling invoice collection") + }.await; match res { Ok(_) => { @@ -519,3 +572,157 @@ async fn prompt_to_continue(message: &str) -> anyhow::Result<()> { Err(anyhow::anyhow!("Aborted by user.")) } } + +#[cfg(test)] +mod tests { + use super::*; + + fn invoice(manual: bool, has_payment_method: bool) -> Invoice { + Invoice::from(stripe::Invoice { + id: "in_test".parse().unwrap(), + status: Some(stripe::InvoiceStatus::Draft), + collection_method: Some(stripe::CollectionMethod::ChargeAutomatically), + customer: Some(stripe::Expandable::Object(Box::new(stripe::Customer { + id: "cus_test".parse().unwrap(), + invoice_settings: Some(stripe::InvoiceSettingCustomerSetting { + default_payment_method: has_payment_method + .then(|| stripe::Expandable::Id("pm_test".parse().unwrap())), + ..Default::default() + }), + ..Default::default() + }))), + metadata: Some( + billing_types::InvoiceMetadata { + tenant: "acmeCo/".to_string(), + invoice_type: if manual { + InvoiceType::Manual + } else { + InvoiceType::Final + }, + period_start: "2026-08-01".to_string(), + period_end: "2026-08-31".to_string(), + } + .to_metadata_map(), + ), + ..Default::default() + }) + } + + #[test] + fn collection_policy() { + for has_payment_method in [false, true] { + let mut manual = invoice(true, has_payment_method); + assert!( + collection_method_needs_update(&manual) + .unwrap_err() + .to_string() + .contains("publish-invoices") + ); + manual.collection_method = Some(stripe::CollectionMethod::SendInvoice); + assert!(!collection_method_needs_update(&manual).unwrap()); + + let mut usage = invoice(false, has_payment_method); + assert_eq!( + collection_method_needs_update(&usage).unwrap(), + !has_payment_method + ); + usage.collection_method = Some(stripe::CollectionMethod::SendInvoice); + assert!(!collection_method_needs_update(&usage).unwrap()); + usage.collection_method = None; + assert!(collection_method_needs_update(&usage).is_err()); + } + } + + async fn stripe_stub( + responses: Vec<(axum::http::StatusCode, serde_json::Value)>, + ) -> ( + stripe::Client, + std::sync::Arc>>, + ) { + let requests = std::sync::Arc::new(std::sync::Mutex::new(Vec::new())); + let recorded = requests.clone(); + let responses = std::sync::Arc::new(std::sync::Mutex::new(responses.into_iter())); + let app = axum::Router::new().fallback(move |method: axum::http::Method, uri: axum::http::Uri| { + let recorded = recorded.clone(); + let responses = responses.clone(); + async move { + recorded.lock().unwrap().push(format!("{method} {}", uri.path())); + let (status, body) = responses.lock().unwrap().next().unwrap_or(( + axum::http::StatusCode::BAD_REQUEST, + serde_json::json!({"error": {"type": "invalid_request_error", "message": "unexpected request"}}), + )); + (status, axum::Json(body)) + } + }); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}/", listener.local_addr().unwrap()); + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + (stripe::Client::from_url(url.as_str(), "test"), requests) + } + + #[tokio::test] + async fn failed_conversion_is_excluded_from_finalization() { + let mut successful = invoice(false, false); + successful.id = "in_success".parse().unwrap(); + let mut converted = successful.clone(); + converted.collection_method = Some(stripe::CollectionMethod::SendInvoice); + let mut finalized = converted.clone(); + finalized.status = Some(stripe::InvoiceStatus::Open); + let (client, requests) = stripe_stub(vec![ + (axum::http::StatusCode::BAD_REQUEST, serde_json::json!({"error": {"type": "invalid_request_error", "message": "update failed"}})), + (axum::http::StatusCode::OK, serde_json::to_value(&*converted).unwrap()), + (axum::http::StatusCode::OK, serde_json::to_value(&*converted).unwrap()), + (axum::http::StatusCode::OK, serde_json::to_value(&*finalized).unwrap()), + (axum::http::StatusCode::OK, serde_json::to_value(&*finalized).unwrap()), + ]).await; + let updated = update_collection_methods(&client, vec![invoice(false, false), successful]) + .await + .unwrap(); + assert_eq!(updated.len(), 1); + let finalized = finalize_invoices(&client, updated).await.unwrap(); + assert_eq!(finalized.len(), 1); + assert_eq!( + *requests.lock().unwrap(), + [ + "POST /v1/invoices/in_test", + "POST /v1/invoices/in_success", + "GET /v1/invoices/in_success", + "POST /v1/invoices/in_success/finalize", + "GET /v1/invoices/in_success", + ] + ); + } + + #[tokio::test] + async fn payment_method_removed_after_confirmation_prevents_finalization() { + let (client, requests) = stripe_stub(vec![( + axum::http::StatusCode::OK, + serde_json::to_value(&*invoice(false, false)).unwrap(), + )]) + .await; + assert!( + finalize_invoices(&client, vec![invoice(false, true)]) + .await + .unwrap() + .is_empty() + ); + assert_eq!(*requests.lock().unwrap(), ["GET /v1/invoices/in_test"]); + } + + #[tokio::test] + async fn invalid_open_invoices_cannot_enable_collection() { + for manual in [false, true] { + let mut current = invoice(manual, manual); + current.status = Some(stripe::InvoiceStatus::Open); + let (client, requests) = stripe_stub(vec![( + axum::http::StatusCode::OK, + serde_json::to_value(&*current).unwrap(), + )]) + .await; + update_auto_advance(&client, vec![current]).await.unwrap(); + assert_eq!(*requests.lock().unwrap(), ["GET /v1/invoices/in_test"]); + } + } +} From a954e01dd73169fcceefb2ff79ba9d735094f526 Mon Sep 17 00:00:00 2001 From: Joseph Shearer Date: Mon, 21 Sep 2026 17:00:23 -0400 Subject: [PATCH 6/6] billing: report incomplete send runs as failures Continue processing valid invoices, but return an error when preparation, finalization, or auto-advance updates fail. Include tenant and invoice context and distinguish draft republishing from open-invoice correction. --- ...1ea41930b4bac33ad42a66cabec3557fd929b.json | 68 +++++++-- ...95385a5a0c05d17a27672e81183d706b8afa5.json | 68 +++++++-- crates/billing-integrations/README.md | 4 +- crates/billing-integrations/src/send.rs | 143 ++++++++++++------ 4 files changed, 214 insertions(+), 69 deletions(-) diff --git a/.sqlx/query-49dafee439e40ad5087811da95a1ea41930b4bac33ad42a66cabec3557fd929b.json b/.sqlx/query-49dafee439e40ad5087811da95a1ea41930b4bac33ad42a66cabec3557fd929b.json index f7c277f851e..e27cddd00d2 100644 --- a/.sqlx/query-49dafee439e40ad5087811da95a1ea41930b4bac33ad42a66cabec3557fd929b.json +++ b/.sqlx/query-49dafee439e40ad5087811da95a1ea41930b4bac33ad42a66cabec3557fd929b.json @@ -6,42 +6,80 @@ { "ordinal": 0, "name": "date_start!", - "type_info": "Date" + "type_info": "Date", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "date_start" + } + } }, { "ordinal": 1, "name": "date_end!", - "type_info": "Date" + "type_info": "Date", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "date_end" + } + } }, { "ordinal": 2, "name": "billed_prefix!", - "type_info": "Text" + "type_info": "Text", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "billed_prefix" + } + } }, { "ordinal": 3, "name": "invoice_type!: InvoiceType", - "type_info": "Text" + "type_info": "Text", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "invoice_type" + } + } }, { "ordinal": 4, "name": "line_items!: sqlx::types::Json>", - "type_info": "Jsonb" + "type_info": "Jsonb", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "line_items" + } + } }, { "ordinal": 5, "name": "subtotal!", - "type_info": "Int8" + "type_info": "Int8", + "origin": "Expression" }, { "ordinal": 6, "name": "extra: sqlx::types::Json>", - "type_info": "Jsonb" + "type_info": "Jsonb", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "extra" + } + } }, { "ordinal": 7, "name": "has_full_pipeline!", - "type_info": "Bool" + "type_info": "Bool", + "origin": "Expression" }, { "ordinal": 8, @@ -56,12 +94,24 @@ ] } } + }, + "origin": { + "Table": { + "table": "tenants", + "name": "payment_provider" + } } }, { "ordinal": 9, "name": "tenant_trial_start", - "type_info": "Date" + "type_info": "Date", + "origin": { + "Table": { + "table": "tenants", + "name": "trial_start" + } + } } ], "parameters": { diff --git a/.sqlx/query-86637e47572df3efbcbf5e0248595385a5a0c05d17a27672e81183d706b8afa5.json b/.sqlx/query-86637e47572df3efbcbf5e0248595385a5a0c05d17a27672e81183d706b8afa5.json index d698b34bd8b..5102517f9ca 100644 --- a/.sqlx/query-86637e47572df3efbcbf5e0248595385a5a0c05d17a27672e81183d706b8afa5.json +++ b/.sqlx/query-86637e47572df3efbcbf5e0248595385a5a0c05d17a27672e81183d706b8afa5.json @@ -6,42 +6,80 @@ { "ordinal": 0, "name": "date_start!", - "type_info": "Date" + "type_info": "Date", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "date_start" + } + } }, { "ordinal": 1, "name": "date_end!", - "type_info": "Date" + "type_info": "Date", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "date_end" + } + } }, { "ordinal": 2, "name": "billed_prefix!", - "type_info": "Text" + "type_info": "Text", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "billed_prefix" + } + } }, { "ordinal": 3, "name": "invoice_type!: InvoiceType", - "type_info": "Text" + "type_info": "Text", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "invoice_type" + } + } }, { "ordinal": 4, "name": "line_items!: sqlx::types::Json>", - "type_info": "Jsonb" + "type_info": "Jsonb", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "line_items" + } + } }, { "ordinal": 5, "name": "subtotal!", - "type_info": "Int8" + "type_info": "Int8", + "origin": "Expression" }, { "ordinal": 6, "name": "extra: sqlx::types::Json>", - "type_info": "Jsonb" + "type_info": "Jsonb", + "origin": { + "Table": { + "table": "internal.invoices_ext", + "name": "extra" + } + } }, { "ordinal": 7, "name": "has_full_pipeline!", - "type_info": "Bool" + "type_info": "Bool", + "origin": "Expression" }, { "ordinal": 8, @@ -56,12 +94,24 @@ ] } } + }, + "origin": { + "Table": { + "table": "tenants", + "name": "payment_provider" + } } }, { "ordinal": 9, "name": "tenant_trial_start", - "type_info": "Date" + "type_info": "Date", + "origin": { + "Table": { + "table": "tenants", + "name": "trial_start" + } + } } ], "parameters": { diff --git a/crates/billing-integrations/README.md b/crates/billing-integrations/README.md index 8704c5815bd..0e2d2ac6c46 100644 --- a/crates/billing-integrations/README.md +++ b/crates/billing-integrations/README.md @@ -11,7 +11,9 @@ Operator CLI for publishing control-plane bills to Stripe and approving collecti finalization. Usage drafts without a payment method can switch to Net 30 with operator approval. Failed switches are excluded from finalization; incorrectly configured manual drafts require republishing. A changed collection decision - after approval requires another send run. + after approval requires another send run. Per-invoice failures do not prevent + other invoices from completing, but cause the command to exit unsuccessfully. + Errors use structured logging so invoice context is retained in redirected output. - `stripe_utils.rs` wraps Stripe invoices for metadata access and operator tables. Both commands accept `--month YYYY-MM` or the legacy `YYYY-MM-01` form. diff --git a/crates/billing-integrations/src/send.rs b/crates/billing-integrations/src/send.rs index 31cc8cb6f40..b767ac85bee 100644 --- a/crates/billing-integrations/src/send.rs +++ b/crates/billing-integrations/src/send.rs @@ -3,7 +3,10 @@ use anyhow::Context; use billing_types::{InvoiceSearch, InvoiceType, StatusFilter}; use chrono::{Datelike, Duration, NaiveDate, Utc}; use clap::Args; -use futures::stream::{self, StreamExt}; +use futures::{ + TryFutureExt, + stream::{self, StreamExt}, +}; use indicatif::{ProgressBar, ProgressStyle}; use itertools::Itertools; use num_format::{Locale, ToFormattedString}; @@ -143,10 +146,9 @@ pub async fn do_send_invoices(cmd: &SendInvoices) -> anyhow::Result<()> { draft_invoices.len() ); - if !draft_invoices.is_empty() { - // 2a. Update collection methods for any drafts that need it - draft_invoices = update_draft_collection_methods(&stripe_client, draft_invoices).await?; - } + let selected_drafts = draft_invoices.len(); + draft_invoices = update_draft_collection_methods(&stripe_client, draft_invoices).await?; + let mut failures = selected_drafts - draft_invoices.len(); if !draft_invoices.is_empty() { print_invoice_table("Invoices to finalize", &draft_invoices); @@ -154,17 +156,19 @@ pub async fn do_send_invoices(cmd: &SendInvoices) -> anyhow::Result<()> { .await?; // 2b. Move the draft invoices to the `open` state - finalized_invoices.append(&mut finalize_invoices(&stripe_client, draft_invoices).await?); + let selected = draft_invoices.len(); + let mut finalized = finalize_invoices(&stripe_client, draft_invoices).await?; + failures += selected - finalized.len(); + finalized_invoices.append(&mut finalized); } if finalized_invoices.is_empty() { tracing::info!("No invoices to send for {month_human_repr}"); - return Ok(()); } // 2c. Check for and fix auto_advance if flag is set if cmd.fix_auto_advance { - finalized_invoices = check_and_fix_auto_advance(&stripe_client, finalized_invoices).await?; + failures += check_and_fix_auto_advance(&stripe_client, &finalized_invoices).await?; } // 3. Show final status of invoices (auto-advance will handle charging automatically) @@ -175,6 +179,10 @@ pub async fn do_send_invoices(cmd: &SendInvoices) -> anyhow::Result<()> { finalized_invoices.len() ); } + anyhow::ensure!( + failures == 0, + "Failed to process {failures} invoice(s); review the errors above before retrying" + ); Ok(()) } @@ -198,7 +206,8 @@ async fn update_draft_collection_methods( current.status() == Some(stripe::InvoiceStatus::Draft), "invoice is no longer a draft" ); - let needs_update = collection_method_needs_update(¤t)?; + let needs_update = collection_method_needs_update(¤t) + .context("Invalid draft collection configuration; rerun publish-invoices")?; Ok::<_, anyhow::Error>((current, needs_update)) } .await; @@ -209,7 +218,9 @@ async fn update_draft_collection_methods( } refreshed.push(current); } - Err(error) => tracing::error!(invoice = %inv.id(), error = %error, "Skipping invoice"), + Err(error) => { + tracing::error!(invoice = %inv.id(), tenant = %inv.tenant(), error = %format!("{error:#}"), "Skipping invoice") + } } } to_update = refreshed; @@ -274,13 +285,17 @@ async fn update_collection_methods( ) .await; match res { - Ok(invoice) => updated.push(Invoice::from(invoice)), + Ok(mut invoice) => { + invoice.customer = inv.customer.clone(); + updated.push(Invoice::from(invoice)); + } Err(e) => { - pb.println(format!( - "Skipping invoice {} after collection method update failed: {}", - inv.id(), - e - )); + tracing::error!( + invoice = %inv.id(), + tenant = %inv.tenant(), + error = %format!("{e:#}"), + "Skipping invoice after collection method update failed" + ); } } pb.inc(1); @@ -295,7 +310,7 @@ fn collection_method_needs_update(invoice: &Invoice) -> anyhow::Result { let method = invoice.collection_method()?; anyhow::ensure!( !invoice.is_manual() || method == stripe::CollectionMethod::SendInvoice, - "manual invoice must use send_invoice; rerun publish-invoices" + "manual invoice must use send_invoice" ); if method == stripe::CollectionMethod::SendInvoice { return Ok(false); @@ -319,6 +334,7 @@ async fn finalize_invoices( let finalize_futs = to_finalize.into_iter().map(|row| { let stripe_client = stripe_client; let pb = pb.clone(); + let context = format!("Invoice {} (tenant: {})", row.id(), row.tenant()); async move { // The operator can pause at the prompt. Re-read before enabling collection, // and require another send run if the approved collection decision changed. @@ -336,7 +352,9 @@ async fn finalize_invoices( ); anyhow::ensure!( current.collection_method()? == row.collection_method()? - && !collection_method_needs_update(¤t)?, + && !collection_method_needs_update(¤t).context( + "Invalid draft collection configuration; rerun publish-invoices" + )?, "invoice {} collection decision changed; rerun send-invoices", row.id() ); @@ -348,17 +366,21 @@ async fn finalize_invoices( }, ) .await - .map_err(|e| { - pb.println(format!("Error finalizing invoice {}: {}", row.id(), e)); - anyhow::Error::from(e) - })?; + .context("Finalizing invoice")?; pb.inc(1); - let invoice = - StripeInvoice::retrieve(stripe_client, row.id(), vec!["customer"].as_slice()) - .await?; + let invoice = StripeInvoice::retrieve( + stripe_client, + row.id(), + vec!["customer"].as_slice(), + ) + .await + .context( + "Invoice was finalized but could not be read back; check Stripe before retrying", + )?; Ok(Invoice::from(invoice)) } + .map_err(move |error: anyhow::Error| error.context(context)) }); let finalize_results = stream::iter(finalize_futs) .buffer_unordered(10) @@ -367,7 +389,7 @@ async fn finalize_invoices( .into_iter() .filter(|res: &anyhow::Result| { if let Err(e) = res { - tracing::error!(error = ?e, "Error finalizing invoice"); + tracing::error!(error = %format!("{e:#}"), "Failed to process invoice"); return false; } true @@ -382,8 +404,8 @@ async fn finalize_invoices( async fn check_and_fix_auto_advance( stripe_client: &Client, - invoices: Vec, -) -> anyhow::Result> { + invoices: &[Invoice], +) -> anyhow::Result { // Find invoices with auto_advance turned off let needs_auto_advance_fix: Vec = invoices .iter() @@ -395,7 +417,7 @@ async fn check_and_fix_auto_advance( .collect(); if needs_auto_advance_fix.is_empty() { - return Ok(invoices); + return Ok(0); } // Show table of invoices that need auto_advance fixed @@ -436,13 +458,16 @@ async fn check_and_fix_auto_advance( .await?; // Update auto_advance to true - update_auto_advance(stripe_client, needs_auto_advance_fix).await?; + return update_auto_advance(stripe_client, needs_auto_advance_fix).await; } - Ok(invoices) + Ok(0) } -async fn update_auto_advance(stripe_client: &Client, invoices: Vec) -> anyhow::Result<()> { +async fn update_auto_advance( + stripe_client: &Client, + invoices: Vec, +) -> anyhow::Result { #[derive(serde::Serialize)] struct UpdateAutoAdvance { auto_advance: bool, @@ -452,10 +477,16 @@ async fn update_auto_advance(stripe_client: &Client, invoices: Vec) -> pb.set_message("updating auto_advance"); pb.set_style(ProgressStyle::with_template(PROGRESS_BAR_TEMPLATE).unwrap()); + let mut failures = 0; for inv in invoices { let res: anyhow::Result = async { - let current = Invoice::from(StripeInvoice::retrieve(stripe_client, inv.id(), &["customer"]).await?); - anyhow::ensure!(current.status() == Some(stripe::InvoiceStatus::Open), "invoice is no longer open"); + let current = Invoice::from( + StripeInvoice::retrieve(stripe_client, inv.id(), &["customer"]).await? + ); + anyhow::ensure!( + current.status() == Some(stripe::InvoiceStatus::Open), + "invoice is no longer open" + ); anyhow::ensure!( !collection_method_needs_update(¤t) .context("Open invoice requires explicit correction in Stripe")?, @@ -476,19 +507,20 @@ async fn update_auto_advance(stripe_client: &Client, invoices: Vec) -> )); } Err(e) => { - pb.println(format!( - "Failed to update auto_advance for invoice {} (tenant: {}): {}", - inv.id(), - inv.tenant(), - e - )); + failures += 1; + tracing::error!( + invoice = %inv.id(), + tenant = %inv.tenant(), + error = %format!("{e:#}"), + "Failed to update auto_advance for invoice" + ); } } pb.inc(1); } pb.finish_with_message("Auto-advance updates complete"); - Ok(()) + Ok(failures) } fn build_invoice_table(rows: I, subtotal: Option) -> comfy_table::Table @@ -616,7 +648,7 @@ mod tests { collection_method_needs_update(&manual) .unwrap_err() .to_string() - .contains("publish-invoices") + .contains("manual invoice") ); manual.collection_method = Some(stripe::CollectionMethod::SendInvoice); assert!(!collection_method_needs_update(&manual).unwrap()); @@ -712,17 +744,28 @@ mod tests { } #[tokio::test] - async fn invalid_open_invoices_cannot_enable_collection() { - for manual in [false, true] { - let mut current = invoice(manual, manual); + async fn only_valid_open_invoices_can_enable_collection() { + for (manual, has_payment_method, valid) in [ + (false, false, false), + (true, true, false), + (false, true, true), + ] { + let mut current = invoice(manual, has_payment_method); current.status = Some(stripe::InvoiceStatus::Open); - let (client, requests) = stripe_stub(vec![( + let response = ( axum::http::StatusCode::OK, serde_json::to_value(&*current).unwrap(), - )]) - .await; - update_auto_advance(&client, vec![current]).await.unwrap(); - assert_eq!(*requests.lock().unwrap(), ["GET /v1/invoices/in_test"]); + ); + let (client, requests) = stripe_stub(vec![response.clone(), response]).await; + // The second response is used only when collection is allowed. + let expected = if valid { + vec!["GET /v1/invoices/in_test", "POST /v1/invoices/in_test"] + } else { + vec!["GET /v1/invoices/in_test"] + }; + let failures = update_auto_advance(&client, vec![current]).await.unwrap(); + assert_eq!(failures, usize::from(!valid)); + assert_eq!(*requests.lock().unwrap(), expected); } } }