aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorGreg Stupp2026-10-05 18:46:47 +0000
committerGreg Stupp2026-10-05 18:46:47 +0000
commitc1ba00b7eb82540453a84fc4951a88ab9472b060 (patch)
tree975bc03ae615d99359f2e8b9b457200ad67804c7
parent219b31843c0f3b12ad90e76f8cc0f84cbbb7a268 (diff)
parent9c79a30fd67c151ba0d8a376ab30495fe3715ada (diff)
downloadgeneralresearch-3.6.8.tar.gz
generalresearch-3.6.8.zip
Merges pull request #6 v3.6.8
UserWalletBalances and ProductUserWalletBalances + balance histograms
-rw-r--r--generalresearch/models/thl/finance.py834
-rw-r--r--generalresearch/models/thl/product.py175
-rw-r--r--pyproject.toml3
-rw-r--r--tests/models/thl/test_product.py104
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,