schedule payout runs

This commit is contained in:
aecsocket
2026-08-19 17:37:25 +09:00
parent 311d603168
commit 2f9956f507
12 changed files with 482 additions and 29 deletions
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "\n\t\t\tSELECT\n\t\t\t\tpayout_periods.period,\n\t\t\t\tpayout_periods.raw_actual_aditude_revenue_usd,\n\t\t\t\tpayout_periods.adjustments,\n\t\t\t\tEXISTS (\n\t\t\t\t\tSELECT 1\n\t\t\t\t\tFROM payout_runs\n\t\t\t\t\tWHERE payout_runs.period = payout_periods.period\n\t\t\t\t\t\tAND payout_runs.status IN ('scheduled', 'running')\n\t\t\t\t) AS \"has_active_run!\",\n\t\t\t\tEXISTS (\n\t\t\t\t\tSELECT 1\n\t\t\t\t\tFROM payout_runs\n\t\t\t\t\tWHERE payout_runs.period = payout_periods.period\n\t\t\t\t\t\tAND payout_runs.status = 'succeeded'\n\t\t\t\t) AS \"has_succeeded_run!\"\n\t\t\tFROM payout_periods\n\t\t\tWHERE payout_periods.period = ANY($1)\n\t\t\t",
"query": "\n\t\t\tSELECT\n\t\t\t\tpayout_periods.period,\n\t\t\t\tpayout_periods.raw_actual_aditude_revenue_usd,\n\t\t\t\tpayout_periods.adjustments,\n\t\t\t\tEXISTS (\n\t\t\t\t\tSELECT 1\n\t\t\t\t\tFROM payout_runs\n\t\t\t\t\tWHERE payout_runs.period = payout_periods.period\n\t\t\t\t\t\tAND payout_runs.status = 'scheduled'\n\t\t\t\t) AS \"has_scheduled_run!\",\n\t\t\t\tEXISTS (\n\t\t\t\t\tSELECT 1\n\t\t\t\t\tFROM payout_runs\n\t\t\t\t\tWHERE payout_runs.period = payout_periods.period\n\t\t\t\t\t\tAND payout_runs.status = 'running'\n\t\t\t\t) AS \"has_running_run!\",\n\t\t\t\tEXISTS (\n\t\t\t\t\tSELECT 1\n\t\t\t\t\tFROM payout_runs\n\t\t\t\t\tWHERE payout_runs.period = payout_periods.period\n\t\t\t\t\t\tAND payout_runs.status = 'succeeded'\n\t\t\t\t) AS \"has_succeeded_run!\"\n\t\t\tFROM payout_periods\n\t\t\tWHERE payout_periods.period = ANY($1)\n\t\t\t",
"describe": {
"columns": [
{
@@ -20,11 +20,16 @@
},
{
"ordinal": 3,
"name": "has_active_run!",
"name": "has_scheduled_run!",
"type_info": "Bool"
},
{
"ordinal": 4,
"name": "has_running_run!",
"type_info": "Bool"
},
{
"ordinal": 5,
"name": "has_succeeded_run!",
"type_info": "Bool"
}
@@ -39,8 +44,9 @@
false,
false,
null,
null,
null
]
},
"hash": "90237a2f470766b030df7bb9642f2b1876f262e5c9e165f9f1256aaa58850a0c"
"hash": "066e14c739f3467fbeef939879074a4816912b069894c4b06972300f19ebd191"
}
@@ -0,0 +1,14 @@
{
"db_name": "PostgreSQL",
"query": "\n DELETE FROM payout_period_days\n WHERE period = $1\n ",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Date"
]
},
"nullable": []
},
"hash": "364343e1f73721e33c2eaf3480a471705b60baaa598f5878bcc5a5eb048d0abc"
}
@@ -0,0 +1,28 @@
{
"db_name": "PostgreSQL",
"query": "\n SELECT\n EXISTS (\n SELECT 1\n FROM payout_runs\n WHERE period = $1\n AND status IN ('scheduled', 'running')\n ) AS \"has_active_run!\",\n EXISTS (\n SELECT 1\n FROM payout_runs\n WHERE period = $1\n AND status = 'succeeded'\n ) AS \"has_succeeded_run!\"\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "has_active_run!",
"type_info": "Bool"
},
{
"ordinal": 1,
"name": "has_succeeded_run!",
"type_info": "Bool"
}
],
"parameters": {
"Left": [
"Date"
]
},
"nullable": [
null,
null
]
},
"hash": "63e02c2f770542021a83afbf919d20d7cc5458d00d9067ae8af932434d1c1ad6"
}
@@ -0,0 +1,22 @@
{
"db_name": "PostgreSQL",
"query": "\n SELECT totp_secret\n FROM users\n WHERE id = $1\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "totp_secret",
"type_info": "Varchar"
}
],
"parameters": {
"Left": [
"Int8"
]
},
"nullable": [
true
]
},
"hash": "85d855e515bba733ac8c2a7369da703b762a2622b37af05fe8afb74c5052d59e"
}
@@ -0,0 +1,17 @@
{
"db_name": "PostgreSQL",
"query": "\n INSERT INTO payout_period_days (\n period,\n date,\n raw_estimated_aditude_revenue_usd,\n aditude_impressions\n )\n VALUES ($1, $2, $3, $4)\n ",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Date",
"Date",
"Numeric",
"Int8"
]
},
"nullable": []
},
"hash": "9a40aa709d370075823de2f37fefddb654be5145c7054c8d85602fdabac737e2"
}
@@ -0,0 +1,26 @@
{
"db_name": "PostgreSQL",
"query": "\n SELECT\n NOW() AS \"started_at!\",\n NOW() + INTERVAL '2 minutes' AS \"execute_at!\"\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "started_at!",
"type_info": "Timestamptz"
},
{
"ordinal": 1,
"name": "execute_at!",
"type_info": "Timestamptz"
}
],
"parameters": {
"Left": []
},
"nullable": [
null,
null
]
},
"hash": "b22226128dbe12c7c4144928409480eb891bd6b793bf3b21691746f8eb5de113"
}
@@ -0,0 +1,16 @@
{
"db_name": "PostgreSQL",
"query": "\n INSERT INTO payout_periods (\n period,\n raw_actual_aditude_revenue_usd,\n adjustments\n )\n VALUES ($1, $2, $3)\n ON CONFLICT (period) DO UPDATE SET\n raw_actual_aditude_revenue_usd =\n EXCLUDED.raw_actual_aditude_revenue_usd,\n adjustments = EXCLUDED.adjustments\n ",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Date",
"Numeric",
"Jsonb"
]
},
"nullable": []
},
"hash": "df7d47467f7344cb4538ffddaf3805bc9666c9651152f6ed648434b6e2b7ee61"
}
@@ -0,0 +1,14 @@
{
"db_name": "PostgreSQL",
"query": "\n UPDATE payout_runs\n SET\n status = 'cancelled',\n cancelled_at = NOW(),\n cancelled_by = $1\n WHERE status = 'scheduled'\n ",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Int8"
]
},
"nullable": []
},
"hash": "fc9f0a0ba7c6d28c51d72127e215a3a90918001d9184e403ffa0df66414f8378"
}
@@ -11,7 +11,8 @@ pub struct DBPayoutPeriod {
pub raw_actual_aditude_revenue_usd: Decimal,
pub adjustments: serde_json::Value,
pub days: Vec<DBPayoutPeriodDay>,
pub has_active_run: bool,
pub has_scheduled_run: bool,
pub has_running_run: bool,
pub has_succeeded_run: bool,
}
@@ -41,8 +42,14 @@ impl DBPayoutPeriod {
SELECT 1
FROM payout_runs
WHERE payout_runs.period = payout_periods.period
AND payout_runs.status IN ('scheduled', 'running')
) AS "has_active_run!",
AND payout_runs.status = 'scheduled'
) AS "has_scheduled_run!",
EXISTS (
SELECT 1
FROM payout_runs
WHERE payout_runs.period = payout_periods.period
AND payout_runs.status = 'running'
) AS "has_running_run!",
EXISTS (
SELECT 1
FROM payout_runs
@@ -68,7 +75,8 @@ impl DBPayoutPeriod {
.raw_actual_aditude_revenue_usd,
adjustments: row.adjustments,
days: Vec::new(),
has_active_run: row.has_active_run,
has_scheduled_run: row.has_scheduled_run,
has_running_run: row.has_running_run,
has_succeeded_run: row.has_succeeded_run,
},
)
@@ -0,0 +1,290 @@
use actix_web::{HttpRequest, post, web};
use chrono::{DateTime, Utc};
use rust_decimal::Decimal;
use serde::{Deserialize, Serialize};
use xredis::RedisPool;
use super::Adjustment;
use crate::{
auth::{
AuthenticationError, get_user_from_headers, two_factor::verify_2fa_code,
},
database::{
PgPool,
models::{
DBUserId, generate_payout_run_id,
payout_run_item::{DBPayoutRun, PayoutRunStatus},
},
},
models::{ids::PayoutRunId, pats::Scopes},
queue::{payout_run::estimate, session::AuthQueue},
routes::ApiError,
util::{
error::Context,
time::{YearMonth, net_60_payout_available_at},
},
};
#[derive(Debug, Deserialize)]
pub struct StartPayoutRun {
pub period: YearMonth,
pub two_factor_code: String,
#[serde(with = "rust_decimal::serde::float")]
pub raw_actual_revenue_usd: Decimal,
pub adjustments: Vec<Adjustment>,
}
#[derive(Debug, Serialize)]
pub struct StartPayoutRunResponse {
pub id: PayoutRunId,
pub execute_at: DateTime<Utc>,
}
#[post("/start")]
pub async fn start_run(
req: HttpRequest,
pool: web::Data<PgPool>,
redis: web::Data<RedisPool>,
aditude: web::Data<aditude::Client>,
session_queue: web::Data<AuthQueue>,
web::Json(body): web::Json<StartPayoutRun>,
) -> Result<web::Json<StartPayoutRunResponse>, ApiError> {
let user = get_user_from_headers(
&req,
&**pool,
&redis,
&session_queue,
Scopes::SESSION_ACCESS,
)
.await
.wrap_auth_err("authenticating API request")?
.1;
if !user.role.is_admin() {
return Err(ApiError::Auth(eyre::eyre!(
AuthenticationError::InvalidCredentials,
)));
}
let user_id = DBUserId::from(user.id);
let secret = sqlx::query_scalar!(
r#"
SELECT totp_secret
FROM users
WHERE id = $1
"#,
user_id.0,
)
.fetch_one(&**pool)
.await
.wrap_internal_err("fetching user two-factor secret")?
.wrap_auth_err_with(|| AuthenticationError::InvalidCredentials)?;
let valid_totp =
verify_2fa_code(&body.two_factor_code, &secret, user_id, &redis)
.await
.wrap_auth_err("verifying two-factor code")?;
if !valid_totp {
return Err(ApiError::Auth(eyre::eyre!(
AuthenticationError::InvalidCredentials,
)));
}
if body.raw_actual_revenue_usd.is_sign_negative() {
return Err(ApiError::Request(eyre::eyre!(
"`raw_actual_revenue_usd` cannot be negative",
)));
}
let available_at = net_60_payout_available_at(body.period)
.wrap_request_err("calculating payout period availability")?;
if Utc::now() < available_at {
return Err(ApiError::Conflict(eyre::eyre!(
"payout period is still open",
)));
}
let mut estimates =
estimate(aditude.get_ref(), redis.get_ref(), &[body.period])
.await
.wrap_internal_err("fetching payout estimate")?;
let estimate = estimates
.pop()
.wrap_internal_err("missing requested payout estimate")?;
let adjustments = serde_json::to_value(&body.adjustments)
.wrap_internal_err("serializing payout adjustments")?;
let mut transaction = pool
.begin()
.await
.wrap_internal_err("starting database transaction")?;
sqlx::query!(
r#"
INSERT INTO payout_periods (
period,
raw_actual_aditude_revenue_usd,
adjustments
)
VALUES ($1, $2, $3)
ON CONFLICT (period) DO UPDATE SET
raw_actual_aditude_revenue_usd =
EXCLUDED.raw_actual_aditude_revenue_usd,
adjustments = EXCLUDED.adjustments
"#,
body.period.date(),
body.raw_actual_revenue_usd,
&adjustments,
)
.execute(&mut transaction)
.await
.wrap_internal_err("storing payout period")?;
let run_state = sqlx::query!(
r#"
SELECT
EXISTS (
SELECT 1
FROM payout_runs
WHERE period = $1
AND status IN ('scheduled', 'running')
) AS "has_active_run!",
EXISTS (
SELECT 1
FROM payout_runs
WHERE period = $1
AND status = 'succeeded'
) AS "has_succeeded_run!"
"#,
body.period.date(),
)
.fetch_one(&mut transaction)
.await
.wrap_internal_err("checking payout period run status")?;
if run_state.has_active_run || run_state.has_succeeded_run {
return Err(ApiError::Conflict(eyre::eyre!(
"payout period already has an active or succeeded run",
)));
}
sqlx::query!(
r#"
DELETE FROM payout_period_days
WHERE period = $1
"#,
body.period.date(),
)
.execute(&mut transaction)
.await
.wrap_internal_err("clearing stored payout period days")?;
for day in estimate.days {
let impressions = i64::try_from(day.impressions)
.wrap_internal_err("converting Aditude impressions")?;
sqlx::query!(
r#"
INSERT INTO payout_period_days (
period,
date,
raw_estimated_aditude_revenue_usd,
aditude_impressions
)
VALUES ($1, $2, $3, $4)
"#,
body.period.date(),
day.date,
day.raw_estimated_revenue_usd,
impressions,
)
.execute(&mut transaction)
.await
.wrap_internal_err("storing payout period day")?;
}
let timing = sqlx::query!(
r#"
SELECT
NOW() AS "started_at!",
NOW() + INTERVAL '2 minutes' AS "execute_at!"
"#,
)
.fetch_one(&mut transaction)
.await
.wrap_internal_err("calculating payout run schedule")?;
let id = generate_payout_run_id(&mut transaction)
.await
.wrap_internal_err("generating payout run ID")?;
let run = DBPayoutRun {
id,
period: body.period.date(),
status: PayoutRunStatus::Scheduled,
started_at: timing.started_at,
started_by: user_id,
execute_at: timing.execute_at,
processing_started_at: None,
finished_at: None,
cancelled_at: None,
cancelled_by: None,
error: None,
};
run.upsert(&mut transaction)
.await
.wrap_internal_err("creating scheduled payout run")?;
transaction
.commit()
.await
.wrap_internal_err("committing payout run")?;
Ok(web::Json(StartPayoutRunResponse {
id: id.into(),
execute_at: run.execute_at,
}))
}
#[derive(Debug, Serialize)]
pub struct CancelPayoutRunsResponse {
pub cancelled: u64,
}
#[post("/cancel")]
pub async fn cancel_runs(
req: HttpRequest,
pool: web::Data<PgPool>,
redis: web::Data<RedisPool>,
session_queue: web::Data<AuthQueue>,
) -> Result<web::Json<CancelPayoutRunsResponse>, ApiError> {
let user = get_user_from_headers(
&req,
&**pool,
&redis,
&session_queue,
Scopes::SESSION_ACCESS,
)
.await
.wrap_auth_err("authenticating API request")?
.1;
if !user.role.is_admin() {
return Err(ApiError::Auth(eyre::eyre!(
AuthenticationError::InvalidCredentials,
)));
}
let user_id = DBUserId::from(user.id);
let result = sqlx::query!(
r#"
UPDATE payout_runs
SET
status = 'cancelled',
cancelled_at = NOW(),
cancelled_by = $1
WHERE status = 'scheduled'
"#,
user_id.0,
)
.execute(&**pool)
.await
.wrap_internal_err("cancelling scheduled payout runs")?;
Ok(web::Json(CancelPayoutRunsResponse {
cancelled: result.rows_affected(),
}))
}
@@ -2,10 +2,10 @@ use std::collections::HashMap;
use actix_web::{HttpRequest, get, web};
use chrono::{Months, NaiveDate, Utc};
use rust_decimal::Decimal;
use serde::{Deserialize, Serialize};
use xredis::RedisPool;
use super::Adjustment;
use crate::{
auth::get_user_from_headers,
database::{
@@ -30,10 +30,6 @@ use crate::{
},
};
pub fn config(cfg: &mut web::ServiceConfig) {
cfg.service(get_runs);
}
#[derive(Debug, Serialize, Deserialize)]
pub struct PayoutRuns {
pub periods: Vec<PayoutRunPeriod>,
@@ -44,7 +40,7 @@ pub struct PayoutRunPeriod {
pub period: YearMonth,
pub status: PayoutPeriodStatus,
pub days: Vec<PayoutRunDay>,
pub adjustments: Vec<PayoutRunAdjustment>,
pub adjustments: Vec<Adjustment>,
}
/// Has revenue been distributed for a specific payout period month yet?
@@ -57,6 +53,8 @@ pub enum PayoutPeriodStatus {
/// Revenue should have been received for the platform by now; waiting for
/// an admin to manually execute the payout run.
InReview,
/// A payout run is waiting for its cancellation window to expire.
Scheduled,
/// Payout run is currently executing.
Running,
/// Payout run has been paid out to creators.
@@ -70,19 +68,6 @@ pub struct PayoutRunDay {
pub actual: Option<DayDistribution>,
}
/// Manual admin-input adjustment to a [`PayoutRunPeriod`].
#[derive(Debug, Serialize, Deserialize)]
pub struct PayoutRunAdjustment {
/// Total value of the adjustment.
#[serde(with = "rust_decimal::serde::float")]
pub amount_usd: Decimal,
/// Why this adjustment was applied.
///
/// Only visible to admins.
#[serde(skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
}
#[get("")]
pub async fn get_runs(
req: HttpRequest,
@@ -165,7 +150,7 @@ pub async fn get_runs(
.map(|estimate| -> Result<_, ApiError> {
if let Some(period) = stored_periods.get(&estimate.period.date()) {
let mut adjustments =
serde_json::from_value::<Vec<PayoutRunAdjustment>>(
serde_json::from_value::<Vec<Adjustment>>(
period.adjustments.clone(),
)
.wrap_internal_err("deserializing payout adjustments")?;
@@ -178,7 +163,7 @@ pub async fn get_runs(
.days
.iter()
.map(|day| day.raw_estimated_aditude_revenue_usd)
.sum::<Decimal>();
.sum();
let actual_flow = compute_actual_distribution_flow(
total_estimated_revenue_usd,
period.raw_actual_aditude_revenue_usd,
@@ -210,8 +195,10 @@ pub async fn get_runs(
.collect::<Result<Vec<_>, _>>()?;
let status = if period.has_succeeded_run {
PayoutPeriodStatus::Paid
} else if period.has_active_run {
} else if period.has_running_run {
PayoutPeriodStatus::Running
} else if period.has_scheduled_run {
PayoutPeriodStatus::Scheduled
} else {
PayoutPeriodStatus::InReview
};
@@ -0,0 +1,25 @@
use actix_web::web;
use rust_decimal::Decimal;
use serde::{Deserialize, Serialize};
mod admin;
mod fetch;
pub fn config(cfg: &mut web::ServiceConfig) {
cfg.service(fetch::get_runs)
.service(admin::start_run)
.service(admin::cancel_runs);
}
/// Manual admin-input adjustment to a payout period.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Adjustment {
/// Total value of the adjustment.
#[serde(with = "rust_decimal::serde::float")]
pub amount_usd: Decimal,
/// Why this adjustment was applied.
///
/// Only visible to admins.
#[serde(skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
}