aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--generalresearch/models/thl/finance.py322
-rw-r--r--generalresearch/models/thl/product.py5
-rw-r--r--tests/models/thl/test_product.py92
3 files changed, 418 insertions, 1 deletions
diff --git a/generalresearch/models/thl/finance.py b/generalresearch/models/thl/finance.py
index c571dc1..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
@@ -670,6 +671,167 @@ class ProductUserWalletBalances(BaseModel):
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)
@@ -992,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,
diff --git a/generalresearch/models/thl/product.py b/generalresearch/models/thl/product.py
index 63f3b39..8f8cc4b 100644
--- a/generalresearch/models/thl/product.py
+++ b/generalresearch/models/thl/product.py
@@ -1119,7 +1119,9 @@ class Product(BaseModel, validate_assignment=True):
)
from generalresearch.models.thl.finance import ProductBalances
- assert self.bp_account is not None, "Call self.prefetch_bp_account()"
+ 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
@@ -1279,6 +1281,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,
diff --git a/tests/models/thl/test_product.py b/tests/models/thl/test_product.py
index 97abf0c..5a3e488 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