diff options
| author | Greg Stupp | 2026-10-05 18:46:47 +0000 |
|---|---|---|
| committer | Greg Stupp | 2026-10-05 18:46:47 +0000 |
| commit | c1ba00b7eb82540453a84fc4951a88ab9472b060 (patch) | |
| tree | 975bc03ae615d99359f2e8b9b457200ad67804c7 | |
| parent | 219b31843c0f3b12ad90e76f8cc0f84cbbb7a268 (diff) | |
| parent | 9c79a30fd67c151ba0d8a376ab30495fe3715ada (diff) | |
| download | generalresearch-c1ba00b7eb82540453a84fc4951a88ab9472b060.tar.gz generalresearch-c1ba00b7eb82540453a84fc4951a88ab9472b060.zip | |
Merges pull request #6
v3.6.8
UserWalletBalances and ProductUserWalletBalances + balance histograms
| -rw-r--r-- | generalresearch/models/thl/finance.py | 834 | ||||
| -rw-r--r-- | generalresearch/models/thl/product.py | 175 | ||||
| -rw-r--r-- | pyproject.toml | 3 | ||||
| -rw-r--r-- | tests/models/thl/test_product.py | 104 |
4 files changed, 1096 insertions, 20 deletions
diff --git a/generalresearch/models/thl/finance.py b/generalresearch/models/thl/finance.py index f012cbf..0eb4f8f 100644 --- a/generalresearch/models/thl/finance.py +++ b/generalresearch/models/thl/finance.py @@ -1,6 +1,7 @@ from __future__ import annotations import random +from collections.abc import Iterable from datetime import UTC from typing import TYPE_CHECKING from uuid import uuid4 @@ -22,12 +23,67 @@ from generalresearch.currency import USDCent from generalresearch.decorators import LOG from generalresearch.models.custom_types import AwareDatetimeISO, UUIDStr from generalresearch.models.thl.definitions import SessionAdjustedStatus +from generalresearch.models.thl.user_identifiers import BPUIDStr payout_example = random.randint(150, 750 * 100) adjustment_example = random.randint(-1_000, 50 * 100) -if TYPE_CHECKING: +# Keep these boundaries stable so Prometheus can aggregate histograms across +# Products and Grafana can compare the same ranges over time. Values are USD +# cents; ``None`` represents the required +Inf bucket. +USER_WALLET_BALANCE_BUCKETS: tuple[int | None, ...] = ( + -1000_00, + -500_00, + -250_00, + -100_00, + -50_00, + -25_00, + -10_00, + -5_00, + -1_00, + -1, + 0, + 1, + 5, + 1_00, + 2_00, + 3_00, + 4_00, + 5_00, + 6_00, + 7_00, + 8_00, + 9_00, + 10_00, + 15_00, + 20_00, + 25_00, + 30_00, + 40_00, + 50_00, + 75_00, + 100_00, + 250_00, + 500_00, + None, +) +USER_WALLET_CREDIT_COLUMNS = ( + "bp_payment.CREDIT", + "bp_adjustment.CREDIT", + "user_bonus.CREDIT", + "user_payout_cancel.CREDIT", + "close_contest.CREDIT", + "user_milestone.CREDIT", +) +USER_WALLET_DEBIT_COLUMNS = ( + "bp_adjustment.DEBIT", + "user_bonus.DEBIT", + "user_payout_request.DEBIT", + "user_enter_contest.DEBIT", +) + +if TYPE_CHECKING: from generalresearch.managers.thl.product import ProductManager from generalresearch.models.thl.ledger import LedgerAccount from generalresearch.models.thl.product import Product @@ -183,6 +239,600 @@ class POPFinancial(BaseModel): return res +class UserWalletBalances(BaseModel): + """Cumulative ledger activity and balance for one user's USD wallet. + + All monetary values are integer USD cents and are expressed from the user's + perspective: credits increase the wallet balance and debits decrease it. + A positive balance is an outstanding liability for the Brokerage Product; + a negative balance is tracked separately and must not offset liabilities to + other users when producing Product-level totals. + """ + + model_config = ConfigDict(extra="ignore", populate_by_name=True) + + product_id: UUIDStr = Field( + description="Brokerage Product that owns this user wallet.", + examples=[uuid4().hex], + ) + product_user_id: BPUIDStr = Field() + + last_event: AwareDatetimeISO | None = Field( + default=None, + description=( + "Timestamp of the most recent ledger event included in these totals, " + "or null when the wallet has no events." + ), + ) + + bp_payment_credit: NonNegativeInt = Field( + default=0, + validation_alias="bp_payment.CREDIT", + description="Total task-completion earnings credited to this wallet.", + examples=[18_837], + ) + + adjustment_credit: NonNegativeInt = Field( + default=0, + validation_alias="bp_adjustment.CREDIT", + description="Total positive task reconciliations credited to this wallet.", + examples=[2], + ) + + adjustment_debit: NonNegativeInt = Field( + default=0, + validation_alias="bp_adjustment.DEBIT", + description="Total negative task reconciliations debited from this wallet.", + examples=[753], + ) + + user_bonus_credit: NonNegativeInt = Field( + default=0, + validation_alias="user_bonus.CREDIT", + description="Total non-task bonuses credited to this wallet.", + examples=[0], + ) + + user_bonus_debit: NonNegativeInt = Field( + default=0, + validation_alias="user_bonus.DEBIT", + description="Total bonus reversals debited from this wallet.", + examples=[2_745], + ) + + user_payout_request: NonNegativeInt = Field( + default=0, + validation_alias="user_payout_request.DEBIT", + description=( + "Total payout requests debited from this wallet. These amounts are no " + "longer part of the unredeemed wallet balance." + ), + examples=[18_837], + ) + + user_payout_cancel: NonNegativeInt = Field( + default=0, + validation_alias="user_payout_cancel.CREDIT", + description="Total cancelled payout requests returned to this wallet.", + examples=[18_837], + ) + + user_enter_contest_debit: NonNegativeInt = Field( + default=0, + validation_alias="user_enter_contest.DEBIT", + description="Total cash entries transferred from this wallet to contests.", + examples=[500], + ) + + close_contest_credit: NonNegativeInt = Field( + default=0, + validation_alias="close_contest.CREDIT", + description="Total cash prizes credited when contests closed.", + examples=[2_500], + ) + + user_milestone_credit: NonNegativeInt = Field( + default=0, + validation_alias="user_milestone.CREDIT", + description="Total cash milestone awards credited to this wallet.", + examples=[1_000], + ) + + @computed_field(description="Total debits that decreased this wallet.") + @property + def debit(self) -> int: + return ( + self.adjustment_debit + + self.user_bonus_debit + + self.user_payout_request + + self.user_enter_contest_debit + ) + + @computed_field(description="Total credits that increased this wallet.") + @property + def credit(self) -> int: + return ( + self.bp_payment_credit + + self.adjustment_credit + + self.user_bonus_credit + + self.user_payout_cancel + + self.close_contest_credit + + self.user_milestone_credit + ) + + @computed_field( + title="Wallet Balance", + description=( + "Current wallet balance from the user's perspective. A positive value " + "is owed to the user; a negative value means the user owes or must earn " + "back that amount." + ), + examples=[5_341], + ) + @property + def balance(self) -> int: + return self.credit + (self.debit * -1) + + @classmethod + def from_pop_ledger( + cls, + input_data: pd.DataFrame, + product_id: UUIDStr, + product_user_id: BPUIDStr, + ) -> UserWalletBalances: + """Build one user's wallet balance from ``PopLedgerSchema`` rows.""" + amount_columns = [ + *USER_WALLET_CREDIT_COLUMNS, + *USER_WALLET_DEBIT_COLUMNS, + ] + required_columns = { + "time_idx", + "product_id", + "product_user_id", + *amount_columns, + } + missing_columns = required_columns.difference(input_data.columns) + assert not missing_columns, ( + f"PopLedgerSchema columns missing: {sorted(missing_columns)}" + ) + + wallet_rows = input_data.loc[ + input_data["product_id"].eq(product_id) + & input_data["product_user_id"].eq(product_user_id) + ] + wallet_rows = wallet_rows.loc[ + wallet_rows[amount_columns].sum(axis="columns").gt(0) + ] + + data = wallet_rows[amount_columns].sum().to_dict() + data.update( + product_id=product_id, + product_user_id=product_user_id, + last_event=( + None if wallet_rows.empty else wallet_rows["time_idx"].max() + ), + ) + return cls.model_validate(data) + + +class UserWalletBalanceHistogramBucket(BaseModel): + """One cumulative bucket in the user-wallet balance distribution. + + Monetary boundaries are integer USD cents. ``upper_bound=None`` represents + the required +Inf bucket containing every wallet in the distribution. + """ + + upper_bound: int | None = Field( + description="Inclusive upper boundary, or null for +Inf." + ) + cumulative_count: NonNegativeInt = Field( + description="Number of wallets at or below the upper boundary." + ) + + +def _empty_user_wallet_balance_histogram() -> list[UserWalletBalanceHistogramBucket]: + return [ + UserWalletBalanceHistogramBucket(upper_bound=bound, cumulative_count=0) + for bound in USER_WALLET_BALANCE_BUCKETS + ] + + +class ProductUserWalletBalances(BaseModel): + """Aggregate user-wallet activity for one Brokerage Product. + + All monetary values are integer USD cents. + + These totals should be computed directly by the data layer rather than by + materializing every ``UserWalletBalances`` instance. Products may have many + thousands of user wallets. + + Positive and negative wallet balances are aggregated separately because + a negative user balance MUST not reduce the BP's outstanding liability to + users with positive balances. + + Pending payouts are reported separately because a payout request removes funds + from the user's wallet before disbursement. + """ + + product_id: UUIDStr = Field( + description="Brokerage Product represented by this aggregation.", + examples=[uuid4().hex], + ) + + debit: NonNegativeInt = Field( + default=0, + description="Sum of ledger debits across the Product's user wallets.", + ) + # Note the credit won't equal the net user_task_payments, b/c a user wallet + # could also have credits from things such as contests or bribes. + credit: NonNegativeInt = Field( + default=0, + description="Sum of ledger credits across the Product's user wallets.", + ) + + user_task_payment: NonNegativeInt = Field( + default=0, + description="Total task-completion payments credited to user wallets.", + ) + user_task_adjustment_credit: NonNegativeInt = Field( + default=0, + description="Total positive task adjustments credited to user wallets.", + ) + user_task_adjustment_debit: NonNegativeInt = Field( + default=0, + description="Total negative task adjustments debited from user wallets.", + ) + + outstanding_liability: NonNegativeInt = Field( + default=0, + description="Sum of positive user-wallet balances owed by the Product.", + ) + negative_balance_total: NonNegativeInt = Field( + default=0, + description=( + "Absolute sum of negative user-wallet balances. This does not reduce " + "outstanding_liability." + ), + ) + + positive_wallet_count: NonNegativeInt = Field( + default=0, + description="Number of user wallets with a positive balance.", + ) + zero_wallet_count: NonNegativeInt = Field( + default=0, + description="Number of user wallets with a zero balance.", + ) + negative_wallet_count: NonNegativeInt = Field( + default=0, + description="Number of user wallets with a negative balance.", + ) + + wallet_balance_buckets: list[UserWalletBalanceHistogramBucket] = Field( + default_factory=_empty_user_wallet_balance_histogram, + description=( + "Cumulative current-balance buckets for a Prometheus gauge histogram." + ), + ) + + oldest_event: AwareDatetimeISO | None = Field( + default=None, + description=( + "Oldest wallet-event timestamp included in this aggregation, or null " + "when no events were included." + ), + ) + newest_event: AwareDatetimeISO | None = Field( + default=None, + description=( + "Newest wallet-event timestamp included in this aggregation, or null " + "when no events were included." + ), + ) + + # We can't determine the amount per payout method without querying the + # payout-event table. + pending_payout_amount: NonNegativeInt = Field( + default=0, + description=( + "Total value of requested payouts that have not completed or been " + "cancelled." + ), + ) + + # We cannot determine the pending payout count, only the balance, b/c individual + # txs are aggregated in the POP ledger. We'd have to query the payout-event table. + pending_payout_count: NonNegativeInt | None = Field( + default=None, + description=( + "Number of payout requests that have not completed or been cancelled, " + "or null when the source cannot determine pending status." + ), + ) + + @computed_field(description="Number of wallets represented by the histogram.") + @property + def wallet_balance_count(self) -> int: + return ( + self.positive_wallet_count + + self.zero_wallet_count + + self.negative_wallet_count + ) + + @computed_field(description="Sum of balances represented by the histogram.") + @property + def wallet_balance_sum(self) -> int: + return self.balance + + @model_validator(mode="after") + def validate_wallet_balance_buckets(self) -> ProductUserWalletBalances: + boundaries = tuple(bucket.upper_bound for bucket in self.wallet_balance_buckets) + assert boundaries == USER_WALLET_BALANCE_BUCKETS, ( + "wallet_balance_buckets must use USER_WALLET_BALANCE_BUCKETS" + ) + + counts = [bucket.cumulative_count for bucket in self.wallet_balance_buckets] + assert counts == sorted(counts), ( + "wallet balance bucket counts must be cumulative" + ) + assert counts[-1] == self.wallet_balance_count, ( + "+Inf bucket must equal wallet_balance_count" + ) + return self + + @classmethod + def from_pop_ledger( + cls, + input_data: pd.DataFrame, + product_id: UUIDStr, + ) -> ProductUserWalletBalances: + """Build a current wallet-liability snapshot from ``PopLedgerSchema`` rows. + + The caller can pass rows for multiple Products. This method selects the + requested Product and rows associated with a user, then collapses the + minute/account grain to one lifetime balance per ``product_user_id``. + + Pending payout metrics are intentionally left at their defaults. The POP + ledger merge drops payout identity, so it cannot distinguish individual + payout requests or reliably determine whether each request is still open. + """ + required_columns = { + "time_idx", + "product_id", + "product_user_id", + *USER_WALLET_CREDIT_COLUMNS, + *USER_WALLET_DEBIT_COLUMNS, + } + missing_columns = required_columns.difference(input_data.columns) + assert not missing_columns, ( + f"PopLedgerSchema columns missing: {sorted(missing_columns)}" + ) + + amount_columns = [ + *USER_WALLET_CREDIT_COLUMNS, + *USER_WALLET_DEBIT_COLUMNS, + ] + wallet_rows = input_data.loc[ + input_data["product_id"].eq(product_id) + & input_data["product_user_id"].notna() + ] + wallet_rows = wallet_rows.loc[ + wallet_rows[amount_columns].sum(axis="columns").gt(0) + ] + if wallet_rows.empty: + return cls(product_id=product_id) + + users = wallet_rows.groupby("product_user_id", observed=True)[ + amount_columns + ].sum() + credits = users[list(USER_WALLET_CREDIT_COLUMNS)].sum(axis="columns") + debits = users[list(USER_WALLET_DEBIT_COLUMNS)].sum(axis="columns") + balances = credits - debits + + buckets = [ + UserWalletBalanceHistogramBucket( + upper_bound=upper_bound, + cumulative_count=( + len(balances) + if upper_bound is None + else int(balances.le(upper_bound).sum()) + ), + ) + for upper_bound in USER_WALLET_BALANCE_BUCKETS + ] + + return cls( + product_id=product_id, + debit=int(debits.sum()), + credit=int(credits.sum()), + user_task_payment=int(users["bp_payment.CREDIT"].sum()), + user_task_adjustment_credit=int(users["bp_adjustment.CREDIT"].sum()), + user_task_adjustment_debit=int(users["bp_adjustment.DEBIT"].sum()), + outstanding_liability=int(balances[balances > 0].sum()), + negative_balance_total=abs(int(balances[balances < 0].sum())), + positive_wallet_count=int(balances.gt(0).sum()), + zero_wallet_count=int(balances.eq(0).sum()), + negative_wallet_count=int(balances.lt(0).sum()), + wallet_balance_buckets=buckets, + oldest_event=wallet_rows["time_idx"].min(), + newest_event=wallet_rows["time_idx"].max(), + ) + + @computed_field( + title="Wallet Balance", + description=( + "Net balance across all user wallets (credits minus debits). This is a " + "net position, not the BP's outstanding liability when any individual " + "wallet has a negative balance." + ), + examples=[5_341], + ) + @property + def balance(self) -> int: + return self.credit + (self.debit * -1) + + def to_prometheus(self) -> tuple[bytes, str]: + """Render this Product as a Prometheus response body and content type.""" + return self.many_to_prometheus([self]) + + @classmethod + def many_to_prometheus( + cls, + snapshots: Iterable[ProductUserWalletBalances], + ) -> tuple[bytes, str]: + """Render multiple Products in one Prometheus response. + + Each Product is represented by samples sharing the same metric families + and distinguished by the ``product_id`` label. + + Usage: + balances: list[ProductUserWalletBalances] = ... + body, content_type = ProductUserWalletBalances.many_to_prometheus(balances) + return HttpResponse(body, content_type=content_type) + """ + from prometheus_client import ( + CONTENT_TYPE_LATEST, + CollectorRegistry, + generate_latest, + ) + from prometheus_client.core import ( + GaugeHistogramMetricFamily, + GaugeMetricFamily, + ) + + products = tuple(snapshots) + product_ids = [str(product.product_id) for product in products] + if len(product_ids) != len(set(product_ids)): + raise ValueError("snapshots must contain unique product_id values") + + class ProductUserWalletCollector: + def collect(collector_self): + monetary_metrics = ( + ("debit_usd", "Ledger debits across user wallets.", "debit"), + ("credit_usd", "Ledger credits across user wallets.", "credit"), + ( + "user_task_payment_usd", + "Task-completion payments credited to user wallets.", + "user_task_payment", + ), + ( + "user_task_adjustment_credit_usd", + "Positive task adjustments credited to user wallets.", + "user_task_adjustment_credit", + ), + ( + "user_task_adjustment_debit_usd", + "Negative task adjustments debited from user wallets.", + "user_task_adjustment_debit", + ), + ( + "outstanding_liability_usd", + "Positive user-wallet balances owed by the Product.", + "outstanding_liability", + ), + ( + "negative_balance_total_usd", + "Absolute sum of negative user-wallet balances.", + "negative_balance_total", + ), + ( + "pending_payout_amount_usd", + "Requested payouts awaiting completion or cancellation.", + "pending_payout_amount", + ), + ( + "net_balance_usd", + "Net balance across all user wallets.", + "balance", + ), + ) + for suffix, description, attribute in monetary_metrics: + metric = GaugeMetricFamily( + f"grl_product_user_wallet_{suffix}", + description, + labels=["product_id"], + ) + for product in products: + metric.add_metric( + [str(product.product_id)], + int(getattr(product, attribute)) / 100, + ) + yield metric + + wallet_counts = GaugeMetricFamily( + "grl_product_user_wallet_count", + "Number of user wallets grouped by balance sign.", + labels=["product_id", "balance_sign"], + ) + for product in products: + product_label = str(product.product_id) + for sign, count in ( + ("positive", product.positive_wallet_count), + ("zero", product.zero_wallet_count), + ("negative", product.negative_wallet_count), + ): + wallet_counts.add_metric([product_label, sign], count) + yield wallet_counts + + if any(p.pending_payout_count is not None for p in products): + pending_count = GaugeMetricFamily( + "grl_product_user_wallet_pending_payout_count", + "Payout requests awaiting completion or cancellation.", + labels=["product_id"], + ) + for product in products: + if product.pending_payout_count is not None: + pending_count.add_metric( + [str(product.product_id)], + product.pending_payout_count, + ) + yield pending_count + + for event_name, attribute in ( + ("oldest", "oldest_event"), + ("newest", "newest_event"), + ): + event_timestamp = GaugeMetricFamily( + f"grl_product_user_wallet_{event_name}_event_timestamp_seconds", + f"Unix timestamp of the {event_name} included wallet event.", + labels=["product_id"], + ) + for product in products: + event_time = getattr(product, attribute) + if event_time is not None: + event_timestamp.add_metric( + [str(product.product_id)], + event_time.timestamp(), + ) + if event_timestamp.samples: + yield event_timestamp + + histogram = GaugeHistogramMetricFamily( + "grl_product_user_wallet_balance_usd", + "Current user-wallet balance distribution.", + labels=["product_id"], + ) + for product in products: + histogram.add_metric( + [str(product.product_id)], + [ + ( + "+Inf" + if bucket.upper_bound is None + else str(bucket.upper_bound / 100), + bucket.cumulative_count, + ) + for bucket in product.wallet_balance_buckets + ], + product.wallet_balance_sum / 100, + ) + yield histogram + + registry = CollectorRegistry() + registry.register(ProductUserWalletCollector()) + return generate_latest(registry), CONTENT_TYPE_LATEST + + class ProductBalances(BaseModel): model_config = ConfigDict(extra="ignore", populate_by_name=True) @@ -211,7 +861,6 @@ class ProductBalances(BaseModel): plug_credit: SkipJsonSchema[NonNegativeInt] = Field( default=0, exclude=True, validation_alias="plug.CREDIT" ) - plug_debit: SkipJsonSchema[NonNegativeInt] = Field( default=0, exclude=True, validation_alias="plug.DEBIT" ) @@ -272,6 +921,12 @@ class ProductBalances(BaseModel): examples=[2_745], ) + user_payout_complete_debit: NonNegativeInt = Field( + default=0, + validation_alias="user_payout_complete.DEBIT", + description="Payout processing fees charged to the Product.", + ) + # --- Hidden helper values --- issued_payment: NonNegativeInt = Field( @@ -346,7 +1001,7 @@ class ProductBalances(BaseModel): ) @property def expense(self) -> int: - return self.user_bonus_credit + (self.user_bonus_debit * -1) + return self.user_bonus_credit - self.user_bonus_debit - self.user_payout_complete_debit # --- Properties: account related --- @computed_field( @@ -499,6 +1154,166 @@ class ProductBalances(BaseModel): return abs(self.adjustment) / self.payout + def to_prometheus(self) -> tuple[bytes, str]: + """Render this Product as a Prometheus response body and content type.""" + return self.many_to_prometheus([self]) + + @staticmethod + def many_to_prometheus( + snapshots: Iterable[ProductBalances], + ) -> tuple[bytes, str]: + """Render Product balance snapshots as one Prometheus response. + + Fields documented as always zero and formatted USD strings are omitted. + Every included Product must have a unique ``product_id``. + """ + from prometheus_client import ( + CONTENT_TYPE_LATEST, + CollectorRegistry, + generate_latest, + ) + from prometheus_client.core import GaugeMetricFamily + + products = tuple(snapshots) + if any(product.product_id is None for product in products): + raise ValueError("every snapshot must have a product_id") + + product_ids = [str(product.product_id) for product in products] + if len(product_ids) != len(set(product_ids)): + raise ValueError("snapshots must contain unique product_id values") + + class ProductBalanceCollector: + def collect(collector_self): + usd_metrics = ( + ( + "task_payment_credit_usd", + "Task-completion earnings credited to the Product account.", + "bp_payment_credit", + ), + ( + "adjustment_credit_usd", + "Positive task reconciliations credited to the Product.", + "adjustment_credit", + ), + ( + "adjustment_debit_usd", + "Negative task reconciliations debited from the Product.", + "adjustment_debit", + ), + ( + "supplier_credit_usd", + "Supplier funds received to recoup a negative balance.", + "supplier_credit", + ), + ( + "supplier_debit_usd", + "Supplier payments sent by ACH or wire.", + "supplier_debit", + ), + ( + "user_bonus_debit_usd", + "Bonuses and other non-task payments sent to user wallets.", + "user_bonus_debit", + ), + ( + "user_payout_fee_usd", + "User payout processing fees charged to the Product.", + "user_payout_complete_debit", + ), + ( + "issued_payment_usd", + "Amount credited as taken from this Product for payment.", + "issued_payment", + ), + ( + "payout_usd", + "Total task payouts earned by the Product.", + "payout", + ), + ( + "adjustment_usd", + "Net value of all task adjustments.", + "adjustment", + ), + ( + "expense_usd", + "Net Product expenses, including bonuses and payout fees.", + "expense", + ), + ( + "net_earnings_usd", + "Task payouts after adjustments and Product expenses.", + "net", + ), + ( + "supplier_payment_usd", + "Net supplier payments made by ACH or wire.", + "payment", + ), + ( + "balance_usd", + "Product net earnings after supplier payments.", + "balance", + ), + ( + "retainer_usd", + "Amount held to cover possible future task adjustments.", + "retainer", + ), + ( + "available_balance_usd", + "Product balance currently available for withdrawal.", + "available_balance", + ), + ( + "recoup_usd", + "Amount required to recover a negative Product balance.", + "recoup", + ), + ) + for suffix, description, attribute in usd_metrics: + metric = GaugeMetricFamily( + f"grl_product_balance_{suffix}", + description, + labels=["product_id"], + ) + for product in products: + metric.add_metric( + [str(product.product_id)], + int(getattr(product, attribute)) / 100, + ) + yield metric + + adjustment_ratio = GaugeMetricFamily( + "grl_product_balance_adjustment_ratio", + "Absolute net adjustment divided by total task payouts.", + labels=["product_id"], + ) + for product in products: + adjustment_ratio.add_metric( + [str(product.product_id)], + product.adjustment_percent, + ) + yield adjustment_ratio + + last_event = GaugeMetricFamily( + "grl_product_balance_last_event_timestamp_seconds", + "Unix timestamp of the most recent included ledger event.", + labels=["product_id"], + ) + for product in products: + if product.last_event is not None: + last_event.add_metric( + [str(product.product_id)], + product.last_event.timestamp(), + ) + if last_event.samples: + yield last_event + + registry = CollectorRegistry() + registry.register(ProductBalanceCollector()) + return generate_latest(registry), CONTENT_TYPE_LATEST + @staticmethod def from_pandas( input_data: pd.DataFrame | pd.Series, @@ -535,6 +1350,19 @@ class ProductBalances(BaseModel): ).replace("$-", "-$") +class PrivateProductBalances(ProductBalances): + """Product balances with internal revenue visible to administrative APIs.""" + + commission: int = Field( + default=0, + description=( + "Net commission revenue earned by GRL from this Brokerage Product. " + "Positive adjustments increase this value and reversals decrease it." + ), + examples=[5_038], + ) + + class BusinessBalances(BaseModel): product_balances: list[ProductBalances] = Field(default_factory=list) diff --git a/generalresearch/models/thl/product.py b/generalresearch/models/thl/product.py index ad5b137..16288f6 100644 --- a/generalresearch/models/thl/product.py +++ b/generalresearch/models/thl/product.py @@ -7,6 +7,7 @@ import math import warnings from collections import defaultdict from collections.abc import Callable +from datetime import timedelta from decimal import Decimal from enum import StrEnum from functools import cached_property, partial @@ -35,6 +36,7 @@ from pydantic import ( ) from pydantic.json_schema import SkipJsonSchema +from generalresearch.config import JAMES_BILLINGS_BPID, JAMES_BILLINGS_TX_CUTOFF from generalresearch.currency import USDCent from generalresearch.decorators import LOG from generalresearch.models.custom_types import ( @@ -46,7 +48,9 @@ from generalresearch.models.custom_types import ( from generalresearch.models.definitions import Source from generalresearch.models.thl.finance import ( POPFinancial, + PrivateProductBalances, ProductBalances, + ProductUserWalletBalances, ) from generalresearch.models.thl.payout import ( BrokerageProductPayoutEvent, @@ -80,6 +84,12 @@ if TYPE_CHECKING: from generalresearch.models.thl.ledger import LedgerAccount +PRODUCT_BALANCES_METRICS_CACHE_KEY = "metrics:product_balances" +PRODUCT_USER_WALLET_BALANCES_METRICS_CACHE_KEY = ( + "metrics:product_user_wallet_balances" +) + + # fmt: off GRS_SKINS = [ "mmfwcl.com", "profile.generalresearch.com", @@ -961,6 +971,10 @@ class Product(BaseModel, validate_assignment=True): # Initialization is deferred until unless it's called # (see .prebuild_***()) balance: ProductBalances | None = Field(default=None, description="Product Balance") + private_balance: PrivateProductBalances | None = Field( + default=None, description="Product Balance including private keys" + ) + user_wallet_balance: ProductUserWalletBalances | None = Field(default=None) payouts_total_str: str | None = Field(default=None) payouts_total: USDCent | None = Field(default=None) @@ -1071,6 +1085,7 @@ class Product(BaseModel, validate_assignment=True): # --- Prefetch --- def prefetch_bp_account(self, thl_lm: ThlLedgerManager) -> None: account = thl_lm.get_account_or_create_bp_wallet(product=self) + assert self.id == account.reference_uuid self.bp_account = account # --- Prebuild --- @@ -1078,8 +1093,8 @@ class Product(BaseModel, validate_assignment=True): def prebuild_balance( self, thl_lm: ThlLedgerManager, - ds: GRLDatasets, client: Client, + ds: GRLDatasets | None = None, pop_ledger: PopLedgerMerge | None = None, ) -> None: """ @@ -1105,26 +1120,35 @@ class Product(BaseModel, validate_assignment=True): volume levels. """ LOG.debug(f"Product.prebuild_balance({self.uuid=})") - from generalresearch.incite.schemas.mergers.pop_ledger import ( numerical_col_names, ) + from generalresearch.models.thl.finance import ProductBalances - account: LedgerAccount = thl_lm.get_account_or_create_bp_wallet(product=self) - assert self.id == account.reference_uuid + if self.bp_account is None: + self.prefetch_bp_account(thl_lm=thl_lm) + assert self.bp_account is not None if pop_ledger is None: + assert ds is not None from generalresearch.incite.defaults import pop_ledger as plm pop_ledger = plm(ds=ds) + account: LedgerAccount = self.bp_account + assert self.id == account.reference_uuid + + filters = [ + ("account_id", "==", account.uuid), + ] + if self.uuid == JAMES_BILLINGS_BPID: + filters.append(("time_idx", ">", JAMES_BILLINGS_TX_CUTOFF)) + ddf = pop_ledger.ddf( force_rr_latest=False, include_partial=True, columns=numerical_col_names + ["time_idx"], - filters=[ - ("account_id", "==", account.uuid), - ], + filters=filters, ) if ddf is None: @@ -1143,18 +1167,112 @@ class Product(BaseModel, validate_assignment=True): ) df = df.set_index("time_idx") - from generalresearch.models.thl.finance import ProductBalances balance = ProductBalances.from_pandas(df) balance.product_id = self.uuid - bal: int = thl_lm.get_account_balance_timerange( - account=account, time_end=balance.last_event - ) - assert bal == balance.balance, "Sql and Parquet Balance inconsistent" + # This will time out ... + # bal: int = thl_lm.get_account_balance_timerange( + # account=account, time_end=balance.last_event + # ) + # assert bal == balance.balance, "Sql and Parquet Balance inconsistent" self.balance = balance + def prebuild_private_balance( + self, + thl_lm: ThlLedgerManager, + client: Client, + ds: GRLDatasets | None = None, + pop_ledger: PopLedgerMerge | None = None, + ) -> None: + from generalresearch.models.thl.finance import PrivateProductBalances + + assert self.balance is not None, "Must call self.prebuild_balance() first" + + if pop_ledger is None: + assert ds is not None + from generalresearch.incite.defaults import pop_ledger as plm + + pop_ledger = plm(ds=ds) + + commission_account = thl_lm.get_account_or_create_bp_commission_by_uuid( + self.uuid + ) + filters = [ + ("account_id", "==", commission_account.uuid), + ] + if self.uuid == JAMES_BILLINGS_BPID: + filters.append(("time_idx", ">", JAMES_BILLINGS_TX_CUTOFF)) + + ddf = pop_ledger.ddf( + columns=[ + "time_idx", + "bp_payment.CREDIT", + "bp_adjustment.CREDIT", + "bp_adjustment.DEBIT", + ], + filters=filters, + ) + df = client.compute(collections=ddf, sync=True) + + s = df.set_index("time_idx").sum() + + commission = ( + s["bp_payment.CREDIT"] + + s["bp_adjustment.CREDIT"] + - s["bp_adjustment.DEBIT"] + ) + self.private_balance = PrivateProductBalances.model_validate( + self.balance.model_dump() | {"commission": commission} + ) + + def prebuild_user_wallet_balances( + self, + client: Client, + ds: GRLDatasets | None = None, + pop_ledger: PopLedgerMerge | None = None, + ) -> None: + assert self.user_wallet_enabled + + if pop_ledger is None: + assert ds is not None + from generalresearch.incite.defaults import pop_ledger as plm + + pop_ledger = plm(ds=ds) + + from generalresearch.models.thl.finance import ( + USER_WALLET_CREDIT_COLUMNS, + USER_WALLET_DEBIT_COLUMNS, + ProductUserWalletBalances, + ) + + filters = [ + ("product_id", "==", self.uuid), + ] + if self.uuid == JAMES_BILLINGS_BPID: + filters.append(("time_idx", ">", JAMES_BILLINGS_TX_CUTOFF)) + + columns = [ + "time_idx", + "product_id", + "product_user_id", + *USER_WALLET_CREDIT_COLUMNS, + *USER_WALLET_DEBIT_COLUMNS, + ] + ddf = pop_ledger.ddf( + columns=columns, + filters=filters, + include_partial=True, + force_rr_latest=False, + ) + user_wallet_df = client.compute(ddf, sync=True) + result = ProductUserWalletBalances.from_pop_ledger( + input_data=user_wallet_df, + product_id=self.uuid, + ) + self.user_wallet_balance = result + def prebuild_pop_financial( self, thl_lm: ThlLedgerManager, @@ -1169,6 +1287,7 @@ class Product(BaseModel, validate_assignment=True): """ if self.bp_account is None: self.prefetch_bp_account(thl_lm=thl_lm) + assert self.bp_account is not None from generalresearch.incite.schemas.mergers.pop_ledger import ( numerical_col_names, @@ -1337,13 +1456,19 @@ class Product(BaseModel, validate_assignment=True): ) -> None: LOG.debug(f"Product.set_cache({self.uuid=})") - ex_secs = 60 * 60 * 24 * 3 # 3 days + if pop_ledger is None: + from generalresearch.incite.defaults import pop_ledger as plm + + pop_ledger = plm(ds=ds) self.prefetch_bp_account(thl_lm=thl_lm) - self.prebuild_balance( - thl_lm=thl_lm, ds=ds, client=client, pop_ledger=pop_ledger + self.prebuild_balance(thl_lm=thl_lm, client=client, pop_ledger=pop_ledger) + self.prebuild_private_balance( + thl_lm=thl_lm, client=client, pop_ledger=pop_ledger ) + if self.user_wallet_enabled: + self.prebuild_user_wallet_balances(client=client, pop_ledger=pop_ledger) self.prebuild_payouts(thl_lm=thl_lm, bp_pem=bp_pem) self.prebuild_pop_financial( thl_lm=thl_lm, ds=ds, client=client, pop_ledger=pop_ledger @@ -1356,8 +1481,26 @@ class Product(BaseModel, validate_assignment=True): # a delay in the incite merge file not being built yet. # bal = thl_lm.get_account_balance_timerange(time_end=) + assert self.balance is not None rc = redis_config.create_redis_client() - rc.set(name=self.cache_key, value=self.model_dump_json(), ex=ex_secs) + with rc.pipeline() as pipe: + pipe.set( + name=self.cache_key, + value=self.model_dump_json(), + ex=timedelta(days=3), + ) + pipe.hset( + name=PRODUCT_BALANCES_METRICS_CACHE_KEY, + key=self.uuid, + value=self.balance.model_dump_json(), + ) + if self.user_wallet_balance is not None: + pipe.hset( + name=PRODUCT_USER_WALLET_BALANCES_METRICS_CACHE_KEY, + key=self.uuid, + value=self.user_wallet_balance.model_dump_json(), + ) + pipe.execute() def determine_bp_payment(self, thl_net: Decimal) -> Decimal: """ diff --git a/pyproject.toml b/pyproject.toml index 5da029e..e1d8824 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "generalresearch" -version = "3.6.7" +version = "3.6.8" description = "Python Utilities for General Research" readme = "README.md" requires-python = ">=3.14" @@ -22,6 +22,7 @@ dependencies = [ "pandas", "pandera", "protobuf", + "prometheus-client", "pyarrow", "pycountry", "pydantic-extra-types[phonenumbers]", diff --git a/tests/models/thl/test_product.py b/tests/models/thl/test_product.py index 97abf0c..e45bd42 100644 --- a/tests/models/thl/test_product.py +++ b/tests/models/thl/test_product.py @@ -26,6 +26,7 @@ from generalresearch.models.thl.product import ( SourcesConfig, SupplyConfig, SupplyPolicy, + UserWalletConfig, ) if TYPE_CHECKING: @@ -696,6 +697,8 @@ class TestProductFinancials: assert p1.balance.retainer == 35 assert p1.balance.available_balance == 108 + body, content_type = p1.balance.to_prometheus() + p1.prebuild_payouts( thl_lm=thl_ledger_manager, bp_pem=brokerage_product_payout_event_manager, @@ -803,6 +806,95 @@ class TestProductFinancials: assert p1.payouts_total == 55 assert p1.payouts_total_str == "$0.55" + def test_balance_user_wallet( + self, + gr_business: Business, + product_factory: Callable[..., Product], + user_factory: Callable[..., User], + mnt_filepath: GRLDatasets, + thl_ledger_manager: ThlLedgerManager, + start: datetime, + session_with_tx_factory: Callable[..., Session], + delete_ledger_db: Callable[..., None], + create_main_accounts: Callable[..., None], + client_no_amm: DaskClient, + ledger_collection: LedgerDFCollection, + pop_ledger_merge: PopLedgerMerge, + delete_df_collection: Callable[..., None], + payout_config, + ): + delete_ledger_db() + create_main_accounts() + delete_df_collection(coll=ledger_collection) + + p1: Product = product_factory( + business=gr_business, + user_wallet_config=UserWalletConfig(enabled=True), + payout_config=payout_config, + ) + u1: User = user_factory(product=p1) + bp_wallet = thl_ledger_manager.get_account_or_create_bp_wallet(product=p1) + user_wallet = thl_ledger_manager.get_account_or_create_user_wallet(user=u1) + + session_with_tx_factory( + user=u1, + wall_req_cpi=Decimal("1.00"), + started=start + timedelta(days=1), + ) + assert ( + thl_ledger_manager.get_account_balance(account=bp_wallet) == 57 + ) # 95 * 60% + assert ( + thl_ledger_manager.get_account_balance(account=user_wallet) == 38 + ) # 95-57 + + session_with_tx_factory( + user=u1, + wall_req_cpi=Decimal("1.00"), + started=start + timedelta(days=2), + ) + txs = thl_ledger_manager.get_tx_filtered_by_account(bp_wallet.uuid) + assert len(txs) == 2 + + with pytest.raises(expected_exception=AssertionError) as cm: + p1.prebuild_balance( + thl_lm=thl_ledger_manager, + ds=mnt_filepath, + client=client_no_amm, + ) + assert "Cannot build Product Balance" in str(cm.value) + + ledger_collection.initial_load(client=None, sync=True) + pop_ledger_merge.build(client=client_no_amm, ledger_coll=ledger_collection) + + p1.prebuild_balance( + thl_lm=thl_ledger_manager, + ds=mnt_filepath, + client=client_no_amm, + ) + assert isinstance(p1.balance, ProductBalances) + assert p1.balance.payout == 114 + assert p1.balance.adjustment == 0 + assert p1.balance.expense == 0 + assert p1.balance.net == 114 + assert p1.balance.balance == 114 + + body, content_type = p1.balance.to_prometheus() + + p1.prebuild_private_balance( + thl_lm=thl_ledger_manager, + ds=mnt_filepath, + client=client_no_amm, + ) + assert p1.private_balance.commission == 5 * 2 + + p1.prebuild_user_wallet_balances( + ds=mnt_filepath, + client=client_no_amm, + ) + assert p1.user_wallet_balance.outstanding_liability == 38 * 2 + body, content_type = p1.user_wallet_balance.to_prometheus() + class TestProductBalance: @pytest.fixture @@ -1082,6 +1174,18 @@ class TestProductCache: assert p1.balance.retainer_usd_str == "$0.17" assert p1.balance.available_balance_usd_str == "$0.54" + from generalresearch.models.thl.product import ( + PRODUCT_BALANCES_METRICS_CACHE_KEY, + ) + + metrics_balance_json = rc.hget( + PRODUCT_BALANCES_METRICS_CACHE_KEY, + product.uuid, + ) + assert isinstance(metrics_balance_json, str) + metrics_balance = ProductBalances.model_validate_json(metrics_balance_json) + assert metrics_balance == p1.balance + def test_neg_balance_cache( self, product: Product, |
