From 3e26170be6bb0bc690e15d0cd7b7ace9e5d3827c Mon Sep 17 00:00:00 2001 From: stuppie Date: Mon, 5 Oct 2026 14:22:17 -0600 Subject: Product cache: working on a bulk version. so we dont have to do pop_ledger.ddf tens of thousands of times, which by itself will take days across all products --- generalresearch/models/thl/product.py | 251 ++++++++++++++++------------------ 1 file changed, 121 insertions(+), 130 deletions(-) diff --git a/generalresearch/models/thl/product.py b/generalresearch/models/thl/product.py index f1fea95..7818e52 100644 --- a/generalresearch/models/thl/product.py +++ b/generalresearch/models/thl/product.py @@ -6,7 +6,7 @@ import json import math import warnings from collections import defaultdict -from collections.abc import Callable +from collections.abc import Callable, Collection from datetime import timedelta from decimal import Decimal from enum import StrEnum @@ -1089,12 +1089,79 @@ class Product(BaseModel, validate_assignment=True): # --- Prebuild --- + @staticmethod + def get_pop_ledger_df( + product_ids: Collection[UUIDStr], + client: Client, + thl_lm: ThlLedgerManager, + pop_ledger: PopLedgerMerge, + ) -> pd.DataFrame: + """Load all POP-ledger rows needed to cache a batch of Products.""" + from generalresearch.incite.schemas.mergers.pop_ledger import ( + numerical_col_names, + ) + + product_ids = tuple(product_ids) + if not product_ids: + return pd.DataFrame() + if len(product_ids) != len(set(product_ids)): + raise ValueError("product_ids must be unique") + + accounts: list[LedgerAccount] = thl_lm.get_accounts_if_exists( + qualified_names=[ + f"{thl_lm.currency.value}:bp_wallet:{bpid}" for bpid in product_ids + ] + ) + if len(accounts) != len(product_ids): + accounts_ref = {a.reference_uuid for a in accounts} + missing = set(product_ids) - accounts_ref + raise ValueError(f"Inconsistent BP Wallet Accounts (missing {missing})") + account_ids = [account.uuid for account in accounts] + + commission_accounts: list[LedgerAccount] = thl_lm.get_accounts_if_exists( + qualified_names=[ + f"{thl_lm.currency.value}:revenue:bp_commission:{bpid}" + for bpid in product_ids + ] + ) + account_ids.extend([account.uuid for account in commission_accounts]) + + ddf = pop_ledger.ddf( + force_rr_latest=False, + include_partial=True, + columns=numerical_col_names + + ["time_idx", "account_id", "product_id", "product_user_id"], + # outer lists are OR'd, inner lists are AND'd + filters=[ + [("product_id", "in", list(product_ids))], + [("account_id", "in", list(account_ids))], + ], + ) + if ddf is None: + raise AssertionError("Cannot load Product POP ledger") + + df: pd.DataFrame = client.compute(collections=ddf, sync=True) + if JAMES_BILLINGS_BPID in product_ids and not df.empty: + jb_account = thl_lm.get_account_or_create_bp_wallet_by_uuid( + JAMES_BILLINGS_BPID + ).uuid + jb_commission = thl_lm.get_account_or_create_bp_commission_by_uuid( + JAMES_BILLINGS_BPID + ).uuid + df = df.loc[ + ( + df["product_id"].ne(JAMES_BILLINGS_BPID) + & df["account_id"].ne(jb_account) + & df["account_id"].ne(jb_commission) + ) + | df["time_idx"].gt(JAMES_BILLINGS_TX_CUTOFF) + ] + return df + def prebuild_balance( self, thl_lm: ThlLedgerManager, - client: Client, - ds: GRLDatasets | None = None, - pop_ledger: PopLedgerMerge | None = None, + pop_ledger_df: pd.DataFrame, ) -> None: """ This returns the Product's Balances that are calculated across @@ -1118,111 +1185,52 @@ class Product(BaseModel, validate_assignment=True): absolutely required, Smart Retainer accounts are supported for any 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 + LOG.debug(f"Product.prebuild_balance_from_pop_ledger_df({self.uuid=})") + self.balance = None 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), + balance_df = pop_ledger_df.loc[ + pop_ledger_df["account_id"].eq(self.bp_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=filters, - ) - - if ddf is None: - raise AssertionError("Cannot build Product Balance") - - df = client.compute(collections=ddf, sync=True) - - if df.empty: - # A Product may not have any ledger transactional events. Don't - # attempt to build a balance, leave it as None rather than - # all zeros + if balance_df.empty: LOG.warning(f"Product({self.uuid=}).prebuild_balance empty dataframe") - assert thl_lm.get_account_balance_timerange(account=account) == 0, ( - "If the df is empty, we can also assume that there should be no " - "transactions in the ledger." - ) return - df = df.set_index("time_idx") - - balance = ProductBalances.from_pandas(df) + balance_df = balance_df.drop( + columns=["account_id", "product_id", "product_user_id"] + ).set_index("time_idx") + balance = ProductBalances.from_pandas(balance_df) balance.product_id = self.uuid - - # 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, + pop_ledger_df: pd.DataFrame | 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, - include_partial=True, - ) - df = client.compute(collections=ddf, sync=True) + df = pop_ledger_df.loc[pop_ledger_df["account_id"].eq(commission_account.uuid)] if df.empty: LOG.warning( f"Product({self.uuid=}).prebuild_private_balance empty dataframe" ) return - s = df.set_index("time_idx").sum() + s = df[ + [ + "bp_payment.CREDIT", + "bp_adjustment.CREDIT", + "bp_adjustment.DEBIT", + ] + ].sum() commission = ( s["bp_payment.CREDIT"] @@ -1235,30 +1243,16 @@ class Product(BaseModel, validate_assignment=True): def prebuild_user_wallet_balances( self, - client: Client, - ds: GRLDatasets | None = None, - pop_ledger: PopLedgerMerge | None = None, + pop_ledger_df: pd.DataFrame | 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", @@ -1266,13 +1260,9 @@ class Product(BaseModel, validate_assignment=True): *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) + user_wallet_df = pop_ledger_df.loc[ + pop_ledger_df["product_id"].eq(self.uuid), columns + ] if user_wallet_df.empty: LOG.warning( f"Product({self.uuid=}).prebuild_user_wallet_balances empty dataframe" @@ -1288,9 +1278,7 @@ class Product(BaseModel, validate_assignment=True): def prebuild_pop_financial( self, thl_lm: ThlLedgerManager, - ds: GRLDatasets, - client: Client, - pop_ledger: PopLedgerMerge | None = None, + pop_ledger_df: pd.DataFrame | None = None, ) -> None: """This is very similar to the Product POP Financial endpoint; however, it returns more than one item for a single time interval. This is @@ -1311,25 +1299,10 @@ class Product(BaseModel, validate_assignment=True): rr = ReportRequest(report_type=ReportType.POP_LEDGER, interval="5min") - if pop_ledger is None: - from generalresearch.incite.defaults import pop_ledger as plm - - pop_ledger = plm(ds=ds) - - ddf = pop_ledger.ddf( - force_rr_latest=False, - include_partial=True, - columns=numerical_col_names + ["time_idx", "account_id"], - filters=[ - ("account_id", "==", self.bp_account.uuid), - ("time_idx", ">=", pop_ledger.start), - ], - ) - if ddf is None: - self.pop_financial = [] - return - - df = client.compute(collections=ddf, sync=True) + df = pop_ledger_df.loc[ + pop_ledger_df["account_id"].eq(self.bp_account.uuid), + numerical_col_names + ["time_idx", "account_id"], + ] if df.empty: self.pop_financial = [] @@ -1465,6 +1438,7 @@ class Product(BaseModel, validate_assignment=True): bp_pem: BrokerageProductPayoutEventManager, redis_config: RedisConfig, pop_ledger: PopLedgerMerge | None = None, + pop_ledger_df: pd.DataFrame | None = None, ) -> None: LOG.debug(f"Product.set_cache({self.uuid=})") @@ -1473,18 +1447,35 @@ class Product(BaseModel, validate_assignment=True): pop_ledger = plm(ds=ds) - self.prefetch_bp_account(thl_lm=thl_lm) + if self.bp_account is None: + self.prefetch_bp_account(thl_lm=thl_lm) - self.prebuild_balance(thl_lm=thl_lm, client=client, pop_ledger=pop_ledger) + self.prebuild_balance( + thl_lm=thl_lm, + client=client, + pop_ledger=pop_ledger, + pop_ledger_df=pop_ledger_df, + ) if self.balance: self.prebuild_private_balance( - thl_lm=thl_lm, client=client, pop_ledger=pop_ledger + thl_lm=thl_lm, + client=client, + pop_ledger=pop_ledger, + pop_ledger_df=pop_ledger_df, ) if self.user_wallet_enabled: - self.prebuild_user_wallet_balances(client=client, pop_ledger=pop_ledger) + self.prebuild_user_wallet_balances( + client=client, + pop_ledger=pop_ledger, + pop_ledger_df=pop_ledger_df, + ) 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 + thl_lm=thl_lm, + ds=ds, + client=client, + pop_ledger=pop_ledger, + pop_ledger_df=pop_ledger_df, ) rc = redis_config.create_redis_client() -- cgit v1.2.3