//! Fan+ consumer subscription queries. use chrono::{DateTime, Utc}; use sqlx::PgPool; use super::enums::SubscriptionStatus; use super::id_types::*; use super::models::DbFanPlusSubscription; use crate::error::Result; /// Create or reactivate a Fan+ subscription record. /// /// Uses ON CONFLICT DO UPDATE on the user_id unique constraint to handle /// both duplicate webhooks and re-subscription after cancellation. #[tracing::instrument(skip_all)] pub async fn create_fan_plus_subscription<'e>( executor: impl sqlx::PgExecutor<'e>, user_id: UserId, stripe_subscription_id: &str, stripe_customer_id: &str, ) -> Result> { let sub = sqlx::query_as::<_, DbFanPlusSubscription>( r#" INSERT INTO fan_plus_subscriptions (user_id, stripe_subscription_id, stripe_customer_id) VALUES ($1, $2, $3) ON CONFLICT (user_id) DO UPDATE SET stripe_subscription_id = EXCLUDED.stripe_subscription_id, stripe_customer_id = EXCLUDED.stripe_customer_id, status = 'active', canceled_at = NULL RETURNING * "#, ) .bind(user_id) .bind(stripe_subscription_id) .bind(stripe_customer_id) .fetch_optional(executor) .await?; Ok(sub) } /// Look up a Fan+ subscription by its Stripe subscription ID. #[tracing::instrument(skip_all)] pub async fn get_fan_plus_by_stripe_id( pool: &PgPool, stripe_subscription_id: &str, ) -> Result> { let sub = sqlx::query_as::<_, DbFanPlusSubscription>( "SELECT * FROM fan_plus_subscriptions WHERE stripe_subscription_id = $1", ) .bind(stripe_subscription_id) .fetch_optional(pool) .await?; Ok(sub) } /// Update the status of a Fan+ subscription. /// Sets canceled_at when transitioning to canceled, preserving existing value. #[tracing::instrument(skip_all)] pub async fn update_fan_plus_status<'e>( executor: impl sqlx::PgExecutor<'e>, stripe_subscription_id: &str, status: SubscriptionStatus, ) -> Result> { let sub = sqlx::query_as::<_, DbFanPlusSubscription>( r#" UPDATE fan_plus_subscriptions SET status = $2, canceled_at = CASE WHEN $2 = 'canceled' THEN COALESCE(canceled_at, NOW()) ELSE canceled_at END WHERE stripe_subscription_id = $1 RETURNING * "#, ) .bind(stripe_subscription_id) .bind(status) .fetch_optional(executor) .await?; Ok(sub) } /// Update the billing period of a Fan+ subscription. #[tracing::instrument(skip_all)] pub async fn update_fan_plus_period<'e>( executor: impl sqlx::PgExecutor<'e>, stripe_subscription_id: &str, start: DateTime, end: DateTime, ) -> Result<()> { sqlx::query( r#" UPDATE fan_plus_subscriptions SET current_period_start = $2, current_period_end = $3 WHERE stripe_subscription_id = $1 "#, ) .bind(stripe_subscription_id) .bind(start) .bind(end) .execute(executor) .await?; Ok(()) } /// Cancel a Fan+ subscription (set status + canceled_at). #[tracing::instrument(skip_all)] pub async fn cancel_fan_plus( pool: &PgPool, stripe_subscription_id: &str, ) -> Result> { let sub = sqlx::query_as::<_, DbFanPlusSubscription>( r#" UPDATE fan_plus_subscriptions SET status = 'canceled', canceled_at = COALESCE(canceled_at, NOW()) WHERE stripe_subscription_id = $1 RETURNING * "#, ) .bind(stripe_subscription_id) .fetch_optional(pool) .await?; Ok(sub) } /// Check whether a user has an active Fan+ subscription. #[tracing::instrument(skip_all)] pub async fn is_fan_plus_active(pool: &PgPool, user_id: UserId) -> Result { let exists = sqlx::query_scalar::<_, bool>( "SELECT EXISTS(SELECT 1 FROM fan_plus_subscriptions WHERE user_id = $1 AND status = 'active')", ) .bind(user_id) .fetch_one(pool) .await?; Ok(exists) } /// Mark a Fan+ subscription as scheduled to cancel at period end (or undo). /// /// Sets the local flag; Stripe is the source of truth and re-asserts it via /// the `customer.subscription.updated` webhook. Called from the dashboard /// Cancel/Resume buttons and from the webhook handler. #[tracing::instrument(skip_all)] pub async fn set_cancel_at_period_end( pool: &PgPool, stripe_subscription_id: &str, cancel: bool, ) -> Result> { let sub = sqlx::query_as::<_, DbFanPlusSubscription>( "UPDATE fan_plus_subscriptions SET cancel_at_period_end = $2 WHERE stripe_subscription_id = $1 RETURNING *", ) .bind(stripe_subscription_id) .bind(cancel) .fetch_optional(pool) .await?; Ok(sub) } /// Get a user's Fan+ subscription (any status). #[tracing::instrument(skip_all)] pub async fn get_fan_plus_by_user( pool: &PgPool, user_id: UserId, ) -> Result> { let sub = sqlx::query_as::<_, DbFanPlusSubscription>( "SELECT * FROM fan_plus_subscriptions WHERE user_id = $1", ) .bind(user_id) .fetch_optional(pool) .await?; Ok(sub) }