diff options
| -rw-r--r-- | generalresearch/models/thl/finance.py | 188 | ||||
| -rw-r--r-- | generalresearch/models/thl/product.py | 121 |
2 files changed, 291 insertions, 18 deletions
diff --git a/generalresearch/models/thl/finance.py b/generalresearch/models/thl/finance.py index 4586460..c571dc1 100644 --- a/generalresearch/models/thl/finance.py +++ b/generalresearch/models/thl/finance.py @@ -67,6 +67,21 @@ USER_WALLET_BALANCE_BUCKETS: tuple[int | None, ...] = ( 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 @@ -357,6 +372,47 @@ class UserWalletBalances(BaseModel): 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. @@ -397,15 +453,35 @@ class ProductUserWalletBalances(BaseModel): 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.", @@ -453,6 +529,8 @@ class ProductUserWalletBalances(BaseModel): ), ) + # We can't determine the amount per payout method without querying the + # payout-event table. pending_payout_amount: NonNegativeInt = Field( default=0, description=( @@ -460,10 +538,14 @@ class ProductUserWalletBalances(BaseModel): "cancelled." ), ) - pending_payout_count: NonNegativeInt = Field( - default=0, + + # 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." + "Number of payout requests that have not completed or been cancelled, " + "or null when the source cannot determine pending status." ), ) @@ -497,6 +579,84 @@ class ProductUserWalletBalances(BaseModel): ) 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=( @@ -539,7 +699,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" ) @@ -600,6 +759,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( @@ -674,7 +839,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( @@ -863,6 +1028,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..997b6e6 100644 --- a/generalresearch/models/thl/product.py +++ b/generalresearch/models/thl/product.py @@ -35,6 +35,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 +47,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, @@ -961,6 +964,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) @@ -1078,8 +1085,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: """ @@ -1106,25 +1113,31 @@ class Product(BaseModel, validate_assignment=True): """ LOG.debug(f"Product.prebuild_balance({self.uuid=})") + 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.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 pop_ledger is None: - from generalresearch.incite.defaults import pop_ledger as plm - - pop_ledger = plm(ds=ds) + 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 +1156,100 @@ 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" + + 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, + thl_lm: ThlLedgerManager, + client: Client, + ds: GRLDatasets | None = None, + pop_ledger: PopLedgerMerge | None = None, + ) -> None: + from generalresearch.models.thl.finance import ( + USER_WALLET_CREDIT_COLUMNS, + USER_WALLET_DEBIT_COLUMNS, + ProductUserWalletBalances, + ) + + assert self.user_wallet_enabled + 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, |
