aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorstuppie2026-10-02 12:27:46 -0600
committerstuppie2026-10-02 12:27:46 -0600
commitdce439f35ba07efd20466d888060f8f9941548de (patch)
treef494252c23159625abfb41379d696a63caad871f
parent569bf733f2288c2892aac1b09224275d52ebac43 (diff)
downloadgeneralresearch-dce439f35ba07efd20466d888060f8f9941548de.tar.gz
generalresearch-dce439f35ba07efd20466d888060f8f9941548de.zip
methods to prebuild ProductUserWalletBalances and PrivateProductBalances from pop ledger
-rw-r--r--generalresearch/models/thl/finance.py188
-rw-r--r--generalresearch/models/thl/product.py121
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,