diff options
| -rw-r--r-- | generalresearch/models/thl/finance.py | 322 | ||||
| -rw-r--r-- | generalresearch/models/thl/product.py | 5 | ||||
| -rw-r--r-- | tests/models/thl/test_product.py | 92 |
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 |
