From 760a1adbe42821a01611bc0919e2f63243d91157 Mon Sep 17 00:00:00 2001 From: aecsocket <43144841+aecsocket@users.noreply.github.com> Date: Mon, 17 Aug 2026 19:08:18 +0900 Subject: [PATCH] payout runs/payout periods split --- ...0d46bbe78053eda81fb2c227d95262cb4140a.json | 20 ++++ ...4f83fcfc49fea35488ad7c9d6f0487760b219.json | 26 ++++ ...a42b68c675c0f38cb760c42acd1786a20c241.json | 24 ++++ ...d9a7d0399e5658657c66d639bbfb0075e321b.json | 22 ++++ ...f3cf2778a413311eae646de8756df53c9488c.json | 15 +++ .../20260817120000_payout-period-revenue.sql | 60 ++++++++++ apps/labrinth/src/database/models/ids.rs | 11 +- apps/labrinth/src/database/models/mod.rs | 2 + .../src/database/models/payout_run_item.rs | 112 ++++++++++++++++++ .../database/models/payout_variance_item.rs | 50 ++++++++ apps/labrinth/src/models/v3/ids.rs | 1 + apps/labrinth/src/queue/payout_run/mod.rs | 4 +- 12 files changed, 342 insertions(+), 5 deletions(-) create mode 100644 apps/labrinth/.sqlx/query-154d32137068456f7c0739271790d46bbe78053eda81fb2c227d95262cb4140a.json create mode 100644 apps/labrinth/.sqlx/query-1ae95b6e5e2480d2d7d4b24384e4f83fcfc49fea35488ad7c9d6f0487760b219.json create mode 100644 apps/labrinth/.sqlx/query-3de025bb99952ac0ba5d1bbf270a42b68c675c0f38cb760c42acd1786a20c241.json create mode 100644 apps/labrinth/.sqlx/query-b7d1f7fc294346cc723b6ed999ad9a7d0399e5658657c66d639bbfb0075e321b.json create mode 100644 apps/labrinth/.sqlx/query-d66c0e72d492d142925f105013ef3cf2778a413311eae646de8756df53c9488c.json create mode 100644 apps/labrinth/migrations/20260817120000_payout-period-revenue.sql create mode 100644 apps/labrinth/src/database/models/payout_run_item.rs create mode 100644 apps/labrinth/src/database/models/payout_variance_item.rs diff --git a/apps/labrinth/.sqlx/query-154d32137068456f7c0739271790d46bbe78053eda81fb2c227d95262cb4140a.json b/apps/labrinth/.sqlx/query-154d32137068456f7c0739271790d46bbe78053eda81fb2c227d95262cb4140a.json new file mode 100644 index 0000000000..e911aa561a --- /dev/null +++ b/apps/labrinth/.sqlx/query-154d32137068456f7c0739271790d46bbe78053eda81fb2c227d95262cb4140a.json @@ -0,0 +1,20 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT MAX(created)\n FROM payouts_values\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "max", + "type_info": "Timestamptz" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + null + ] + }, + "hash": "154d32137068456f7c0739271790d46bbe78053eda81fb2c227d95262cb4140a" +} diff --git a/apps/labrinth/.sqlx/query-1ae95b6e5e2480d2d7d4b24384e4f83fcfc49fea35488ad7c9d6f0487760b219.json b/apps/labrinth/.sqlx/query-1ae95b6e5e2480d2d7d4b24384e4f83fcfc49fea35488ad7c9d6f0487760b219.json new file mode 100644 index 0000000000..f3a398c3e2 --- /dev/null +++ b/apps/labrinth/.sqlx/query-1ae95b6e5e2480d2d7d4b24384e4f83fcfc49fea35488ad7c9d6f0487760b219.json @@ -0,0 +1,26 @@ +{ + "db_name": "PostgreSQL", + "query": "\n\t\t\tSELECT applied_on, variance\n\t\t\tFROM payouts_variance\n\t\t\tORDER BY applied_on\n\t\t\t", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "applied_on", + "type_info": "Date" + }, + { + "ordinal": 1, + "name": "variance", + "type_info": "Numeric" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + false, + false + ] + }, + "hash": "1ae95b6e5e2480d2d7d4b24384e4f83fcfc49fea35488ad7c9d6f0487760b219" +} diff --git a/apps/labrinth/.sqlx/query-3de025bb99952ac0ba5d1bbf270a42b68c675c0f38cb760c42acd1786a20c241.json b/apps/labrinth/.sqlx/query-3de025bb99952ac0ba5d1bbf270a42b68c675c0f38cb760c42acd1786a20c241.json new file mode 100644 index 0000000000..33b482037e --- /dev/null +++ b/apps/labrinth/.sqlx/query-3de025bb99952ac0ba5d1bbf270a42b68c675c0f38cb760c42acd1786a20c241.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "\n\t\t\tINSERT INTO payout_runs (\n\t\t\t\tid,\n\t\t\t\tperiod,\n\t\t\t\tstatus,\n\t\t\t\tstarted_at,\n\t\t\t\tstarted_by,\n\t\t\t\texecute_at,\n\t\t\t\tprocessing_started_at,\n\t\t\t\tfinished_at,\n\t\t\t\tcancelled_at,\n\t\t\t\tcancelled_by,\n\t\t\t\terror\n\t\t\t)\n\t\t\tVALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)\n\t\t\tON CONFLICT (id) DO UPDATE SET\n\t\t\t\tperiod = EXCLUDED.period,\n\t\t\t\tstatus = EXCLUDED.status,\n\t\t\t\tstarted_at = EXCLUDED.started_at,\n\t\t\t\tstarted_by = EXCLUDED.started_by,\n\t\t\t\texecute_at = EXCLUDED.execute_at,\n\t\t\t\tprocessing_started_at = EXCLUDED.processing_started_at,\n\t\t\t\tfinished_at = EXCLUDED.finished_at,\n\t\t\t\tcancelled_at = EXCLUDED.cancelled_at,\n\t\t\t\tcancelled_by = EXCLUDED.cancelled_by,\n\t\t\t\terror = EXCLUDED.error\n\t\t\t", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Int8", + "Date", + "Text", + "Timestamptz", + "Int8", + "Timestamptz", + "Timestamptz", + "Timestamptz", + "Timestamptz", + "Int8", + "Jsonb" + ] + }, + "nullable": [] + }, + "hash": "3de025bb99952ac0ba5d1bbf270a42b68c675c0f38cb760c42acd1786a20c241" +} diff --git a/apps/labrinth/.sqlx/query-b7d1f7fc294346cc723b6ed999ad9a7d0399e5658657c66d639bbfb0075e321b.json b/apps/labrinth/.sqlx/query-b7d1f7fc294346cc723b6ed999ad9a7d0399e5658657c66d639bbfb0075e321b.json new file mode 100644 index 0000000000..94dbe7b4c1 --- /dev/null +++ b/apps/labrinth/.sqlx/query-b7d1f7fc294346cc723b6ed999ad9a7d0399e5658657c66d639bbfb0075e321b.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT EXISTS(SELECT 1 FROM payout_runs WHERE id=$1)", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "exists", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Int8" + ] + }, + "nullable": [ + null + ] + }, + "hash": "b7d1f7fc294346cc723b6ed999ad9a7d0399e5658657c66d639bbfb0075e321b" +} diff --git a/apps/labrinth/.sqlx/query-d66c0e72d492d142925f105013ef3cf2778a413311eae646de8756df53c9488c.json b/apps/labrinth/.sqlx/query-d66c0e72d492d142925f105013ef3cf2778a413311eae646de8756df53c9488c.json new file mode 100644 index 0000000000..ca81a5a998 --- /dev/null +++ b/apps/labrinth/.sqlx/query-d66c0e72d492d142925f105013ef3cf2778a413311eae646de8756df53c9488c.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "\n\t\t\tINSERT INTO payouts_variance (applied_on, variance)\n\t\t\tVALUES ($1, $2)\n\t\t\tON CONFLICT (applied_on) DO UPDATE\n\t\t\tSET variance = EXCLUDED.variance\n\t\t\t", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Date", + "Numeric" + ] + }, + "nullable": [] + }, + "hash": "d66c0e72d492d142925f105013ef3cf2778a413311eae646de8756df53c9488c" +} diff --git a/apps/labrinth/migrations/20260817120000_payout-period-revenue.sql b/apps/labrinth/migrations/20260817120000_payout-period-revenue.sql new file mode 100644 index 0000000000..aa84557d32 --- /dev/null +++ b/apps/labrinth/migrations/20260817120000_payout-period-revenue.sql @@ -0,0 +1,60 @@ +CREATE TABLE payout_periods ( + period DATE PRIMARY KEY + CHECK (EXTRACT(DAY FROM period) = 1), + raw_actual_aditude_revenue_usd NUMERIC(40, 20) NOT NULL + CHECK (raw_actual_aditude_revenue_usd >= 0), + adjustments JSONB NOT NULL +); + +CREATE TABLE payout_period_days ( + period DATE NOT NULL REFERENCES payout_periods(period), + date DATE NOT NULL, + raw_estimated_aditude_revenue_usd NUMERIC(40, 20) NOT NULL + CHECK (raw_estimated_aditude_revenue_usd >= 0), + aditude_impressions BIGINT NOT NULL + CHECK (aditude_impressions >= 0), + CHECK (date >= period AND date < period + INTERVAL '1 month'), + PRIMARY KEY (period, date) +); + +CREATE TABLE payout_runs ( + id BIGINT PRIMARY KEY, + period DATE NOT NULL REFERENCES payout_periods(period), + status TEXT NOT NULL, + started_at TIMESTAMPTZ NOT NULL, + started_by BIGINT NOT NULL REFERENCES users(id) ON DELETE CASCADE, + execute_at TIMESTAMPTZ NOT NULL, + processing_started_at TIMESTAMPTZ, + finished_at TIMESTAMPTZ, + cancelled_at TIMESTAMPTZ, + cancelled_by BIGINT REFERENCES users(id) ON DELETE CASCADE, + error JSONB, + CHECK (execute_at >= started_at), + CHECK ( + processing_started_at IS NULL + OR processing_started_at >= started_at + ), + CHECK (finished_at IS NULL OR finished_at >= started_at), + CHECK (cancelled_at IS NULL OR cancelled_at >= started_at) +); + +CREATE UNIQUE INDEX payout_runs_active_period + ON payout_runs (period) + WHERE status IN ('scheduled', 'running'); + +CREATE UNIQUE INDEX payout_runs_succeeded_period + ON payout_runs (period) + WHERE status = 'succeeded'; + +CREATE INDEX payout_runs_scheduled_execute_at + ON payout_runs (execute_at) + WHERE status = 'scheduled'; + +CREATE TABLE payouts_variance ( + applied_on DATE PRIMARY KEY, + variance NUMERIC(40, 20) NOT NULL + CHECK (variance BETWEEN 0 AND 1) +); + +INSERT INTO payouts_variance (applied_on, variance) +VALUES ('1970-01-01', 0.1); diff --git a/apps/labrinth/src/database/models/ids.rs b/apps/labrinth/src/database/models/ids.rs index 5f92ece21d..861b01799a 100644 --- a/apps/labrinth/src/database/models/ids.rs +++ b/apps/labrinth/src/database/models/ids.rs @@ -4,9 +4,10 @@ use crate::models::ids::{ AffiliateCodeId, AnalyticsEventId, AttributionGroupId, CampaignDonationId, ChargeId, CollectionId, FileId, ImageId, NotificationId, OAuthAccessTokenId, OAuthClientAuthorizationId, OAuthClientId, - OAuthRedirectUriId, OrganizationId, PasskeyId, PatId, PayoutId, ProductId, - ProductPriceId, ProjectId, ReportId, SessionId, TeamId, TeamMemberId, - ThreadId, ThreadMessageId, UserSubscriptionId, VersionId, + OAuthRedirectUriId, OrganizationId, PasskeyId, PatId, PayoutId, + PayoutRunId, ProductId, ProductPriceId, ProjectId, ReportId, SessionId, + TeamId, TeamMemberId, ThreadId, ThreadMessageId, UserSubscriptionId, + VersionId, }; use ariadne::ids::base62_impl::to_base62; use ariadne::ids::{UserId, random_base62_rng, random_base62_rng_range}; @@ -217,6 +218,10 @@ db_id_interface!( PayoutId, generator: generate_payout_id @ "payouts", ); +db_id_interface!( + PayoutRunId, + generator: generate_payout_run_id @ "payout_runs", +); db_id_interface!( ProductId, generator: generate_product_id @ "products", diff --git a/apps/labrinth/src/database/models/mod.rs b/apps/labrinth/src/database/models/mod.rs index 896791e7d3..484742c621 100644 --- a/apps/labrinth/src/database/models/mod.rs +++ b/apps/labrinth/src/database/models/mod.rs @@ -27,6 +27,8 @@ pub mod organization_item; pub mod passkey_item; pub mod pat_item; pub mod payout_item; +pub mod payout_run_item; +pub mod payout_variance_item; pub mod payouts_values_notifications; pub mod product_item; pub mod products_tax_identifier_item; diff --git a/apps/labrinth/src/database/models/payout_run_item.rs b/apps/labrinth/src/database/models/payout_run_item.rs new file mode 100644 index 0000000000..d66bd520cd --- /dev/null +++ b/apps/labrinth/src/database/models/payout_run_item.rs @@ -0,0 +1,112 @@ +use chrono::{DateTime, NaiveDate, Utc}; +use serde::{Deserialize, Serialize}; +use sqlx::types::Json; +use strum::{EnumString, IntoStaticStr}; + +use super::{DBPayoutRunId, DBUserId, DatabaseError}; + +#[derive( + Debug, + Clone, + Copy, + PartialEq, + Eq, + Serialize, + Deserialize, + EnumString, + IntoStaticStr, +)] +#[serde(rename_all = "snake_case")] +#[strum(serialize_all = "snake_case")] +pub enum PayoutRunStatus { + Scheduled, + Running, + Cancelled, + Failed, + Succeeded, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct PayoutRunError { + pub error: String, + pub description: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub details: Option, +} + +impl From> for PayoutRunError { + fn from(error: crate::models::error::ApiError<'_>) -> Self { + Self { + error: error.error.to_owned(), + description: error.description, + details: error.details, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DBPayoutRun { + pub id: DBPayoutRunId, + pub period: NaiveDate, + pub status: PayoutRunStatus, + pub started_at: DateTime, + pub started_by: DBUserId, + pub execute_at: DateTime, + pub processing_started_at: Option>, + pub finished_at: Option>, + pub cancelled_at: Option>, + pub cancelled_by: Option, + pub error: Option, +} + +impl DBPayoutRun { + pub async fn upsert( + &self, + exec: impl crate::database::Executor<'_, Database = sqlx::Postgres>, + ) -> Result<(), DatabaseError> { + sqlx::query!( + r#" + INSERT INTO payout_runs ( + id, + period, + status, + started_at, + started_by, + execute_at, + processing_started_at, + finished_at, + cancelled_at, + cancelled_by, + error + ) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11) + ON CONFLICT (id) DO UPDATE SET + period = EXCLUDED.period, + status = EXCLUDED.status, + started_at = EXCLUDED.started_at, + started_by = EXCLUDED.started_by, + execute_at = EXCLUDED.execute_at, + processing_started_at = EXCLUDED.processing_started_at, + finished_at = EXCLUDED.finished_at, + cancelled_at = EXCLUDED.cancelled_at, + cancelled_by = EXCLUDED.cancelled_by, + error = EXCLUDED.error + "#, + self.id.0, + self.period, + <&'static str>::from(self.status), + self.started_at, + self.started_by.0, + self.execute_at, + self.processing_started_at, + self.finished_at, + self.cancelled_at, + self.cancelled_by.map(|id| id.0), + self.error.as_ref().map(Json) as Option>, + ) + .execute(exec) + .await?; + + Ok(()) + } +} diff --git a/apps/labrinth/src/database/models/payout_variance_item.rs b/apps/labrinth/src/database/models/payout_variance_item.rs new file mode 100644 index 0000000000..1afef897db --- /dev/null +++ b/apps/labrinth/src/database/models/payout_variance_item.rs @@ -0,0 +1,50 @@ +use chrono::NaiveDate; +use rust_decimal::Decimal; +use serde::{Deserialize, Serialize}; + +use super::DatabaseError; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DBPayoutVariance { + pub applied_on: NaiveDate, + pub variance: Decimal, +} + +impl DBPayoutVariance { + pub async fn get_all( + exec: impl crate::database::Executor<'_, Database = sqlx::Postgres>, + ) -> Result, DatabaseError> { + let variances = sqlx::query_as!( + Self, + r#" + SELECT applied_on, variance + FROM payouts_variance + ORDER BY applied_on + "#, + ) + .fetch_all(exec) + .await?; + + Ok(variances) + } + + pub async fn upsert( + &self, + exec: impl crate::database::Executor<'_, Database = sqlx::Postgres>, + ) -> Result<(), DatabaseError> { + sqlx::query!( + r#" + INSERT INTO payouts_variance (applied_on, variance) + VALUES ($1, $2) + ON CONFLICT (applied_on) DO UPDATE + SET variance = EXCLUDED.variance + "#, + self.applied_on, + self.variance, + ) + .execute(exec) + .await?; + + Ok(()) + } +} diff --git a/apps/labrinth/src/models/v3/ids.rs b/apps/labrinth/src/models/v3/ids.rs index 292c253b7f..8ab03bcd35 100644 --- a/apps/labrinth/src/models/v3/ids.rs +++ b/apps/labrinth/src/models/v3/ids.rs @@ -14,6 +14,7 @@ base62_id!(OAuthRedirectUriId); base62_id!(OrganizationId); base62_id!(PatId); base62_id!(PayoutId); +base62_id!(PayoutRunId); base62_id!(ProductId); base62_id!(ProductPriceId); base62_id!(ProjectId); diff --git a/apps/labrinth/src/queue/payout_run/mod.rs b/apps/labrinth/src/queue/payout_run/mod.rs index 9fc0a634d4..099b87dbc8 100644 --- a/apps/labrinth/src/queue/payout_run/mod.rs +++ b/apps/labrinth/src/queue/payout_run/mod.rs @@ -58,8 +58,8 @@ //! ## Variance //! //! We store a table `payouts_variance` with columns: -//! - a timestamp from when this variance value applies (first entry at Unix -//! epoch) +//! - a date from when this variance value applies (first entry on the Unix +//! epoch date) //! - the decimal fraction of variance to apply use chrono::NaiveDate;