aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorstuppie2026-10-05 14:22:17 -0600
committerstuppie2026-10-05 14:22:17 -0600
commit3e26170be6bb0bc690e15d0cd7b7ace9e5d3827c (patch)
tree4585b64f7c59b535920008cd7bff0f1ddd22e52c
parent63831d97cc1fc4791cd1d35ba19f9e2395202ff7 (diff)
downloadgeneralresearch-3e26170be6bb0bc690e15d0cd7b7ace9e5d3827c.tar.gz
generalresearch-3e26170be6bb0bc690e15d0cd7b7ace9e5d3827c.zip
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
-rw-r--r--generalresearch/models/thl/product.py251
1 files 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()