diff options
| -rw-r--r-- | generalresearch/models/thl/product.py | 396 | ||||
| -rw-r--r-- | pyproject.toml | 2 | ||||
| -rw-r--r-- | tests/models/thl/test_product.py | 141 |
3 files changed, 275 insertions, 264 deletions
diff --git a/generalresearch/models/thl/product.py b/generalresearch/models/thl/product.py index a9024e3..0330207 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 @@ -85,6 +85,7 @@ if TYPE_CHECKING: PRODUCT_BALANCES_METRICS_CACHE_KEY = "metrics:product_balances" +PRODUCT_PRIVATE_BALANCES_METRICS_CACHE_KEY = "metrics:private_product_balances" PRODUCT_USER_WALLET_BALANCES_METRICS_CACHE_KEY = "metrics:product_user_wallet_balances" @@ -1087,14 +1088,96 @@ class Product(BaseModel, validate_assignment=True): self.bp_account = account # --- 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, + ) - def prebuild_balance( + 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_individual( self, thl_lm: ThlLedgerManager, client: Client, ds: GRLDatasets | None = None, pop_ledger: PopLedgerMerge | None = None, ) -> 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) + + pop_ledger_df = self.get_pop_ledger_df( + product_ids=[self.uuid], thl_lm=thl_lm, client=client, pop_ledger=pop_ledger + ) + self.prebuild_balance(thl_lm=thl_lm, pop_ledger_df=pop_ledger_df) + + def prebuild_balance( + self, thl_lm: ThlLedgerManager, pop_ledger_df: pd.DataFrame + ) -> None: """ This returns the Product's Balances that are calculated across all time. They are inclusive of every transaction that has ever @@ -1118,110 +1201,51 @@ class Product(BaseModel, validate_assignment=True): 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 + 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"] @@ -1234,30 +1258,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", @@ -1265,13 +1275,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" @@ -1287,9 +1293,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 @@ -1310,25 +1314,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 = [] @@ -1348,7 +1337,6 @@ class Product(BaseModel, validate_assignment=True): def prebuild_payouts( self, - thl_lm: ThlLedgerManager, bp_pem: BrokerageProductPayoutEventManager, ) -> None: LOG.debug(f"Product.prebuild_payouts({self.uuid=})") @@ -1367,133 +1355,66 @@ class Product(BaseModel, validate_assignment=True): self.payouts_total = USDCent(sum([po.amount for po in self.payouts])) self.payouts_total_str = self.payouts_total.to_usd_str() - # def prebuild_pop(self): - # account = LM.get_account(qualified_name=f"{LM.currency.value}:bp_wallet:{product.id}") - # - # from main import data - # - # gv: GlobalVar = data["gv"] - # - # ddf = gv.pop_ledger.ddf( - # force_rr_latest=False, - # include_partial=True, - # columns=numerical_col_names + ["time_idx"], - # filters=[ - # ("account_id", "==", account.uuid), - # ("time_idx", ">=", rr.start), - # ], - # ) - # - # df = gv.dask_client.compute(collections=ddf, sync=True) - # df = df.set_index("time_idx").resample(rr.freq).sum() - # - # res = [] - # for index, row in df.iterrows(): - # index: pd.Timestamp - # row: pd.DataFrame - # - # dt = index.to_pydatetime().replace(tzinfo=None) - # instance = ProductBalances.from_pandas(row) - # - # res.append( - # { - # "time": dt, - # "payout": instance.payout / 100, - # "adjustment": instance.adjustment / 100, - # "expense": instance.expense / 100, - # "net": (instance.payout + instance.adjustment + instance.expense) / 100, - # } - # ) - # - # df = pd.DataFrame.from_records(res) - - # def financial( - # product: Product = Depends(product_from_path), - # rr: ReportRequest = Depends(rr_from_query), - # ) -> Any: - # account = LM.get_account(qualified_name=f"{LM.currency.value}:bp_wallet:{product.id}") - # - # from main import data - # - # gv: GlobalVar = data["gv"] - # - # ddf = gv.pop_ledger.ddf( - # force_rr_latest=False, - # include_partial=True, - # columns=numerical_col_names + ["time_idx", "account_id"], - # filters=[("account_id", "==", account.uuid), ("time_idx", ">=", rr.start)], - # ) - # - # df = gv.dask_client.compute(collections=ddf, sync=True) - # - # # We only do it this way so it's consistent with the Business.financial view - # df = df.groupby([pd.Grouper(key="time_idx", freq=rr.interval), "account_id"]).sum() - # return POPFinancial.list_from_pandas(df, accounts=[account]) - - # def payments(self): - # """Payments are the amount of money that General Research has sent - # the owner of this Product. - # - # These are typically ACH or Wire payments to company bank accounts. - # These are not respondent payments for Products where - # - # This is Provided in a standard list without any POP Grouping to show - # the exact time and amount of any Issued Payments. - # """ - # - # account = LM.get_account(qualified_name=f"{LM.currency.value}:bp_wallet:{product.id}") - # - # from main import data - # - # gv: GlobalVar = data["gv"] - # ddf = gv.pop_ledger.ddf( - # force_rr_latest=False, - # include_partial=True, - # columns=numerical_col_names + ["time_idx", "account_id"], - # filters=[("account_id", "==", account.uuid)], - # ) - # - # df = gv.dask_client.compute(collections=ddf, sync=True) - - # --- Methods --- def set_cache( self, thl_lm: ThlLedgerManager, - ds: GRLDatasets, client: Client, bp_pem: BrokerageProductPayoutEventManager, redis_config: RedisConfig, + ds: GRLDatasets | None = None, pop_ledger: PopLedgerMerge | None = None, - ) -> None: + ): LOG.debug(f"Product.set_cache({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) - self.prefetch_bp_account(thl_lm=thl_lm) + pop_ledger_df = Product.get_pop_ledger_df( + product_ids=[self.uuid], + client=client, + thl_lm=thl_lm, + pop_ledger=pop_ledger, + ) + self.set_cache_from_pop_ledger_df( + thl_lm=thl_lm, + bp_pem=bp_pem, + redis_config=redis_config, + pop_ledger_df=pop_ledger_df, + ) + + def set_cache_from_pop_ledger_df( + self, + thl_lm: ThlLedgerManager, + bp_pem: BrokerageProductPayoutEventManager, + redis_config: RedisConfig, + pop_ledger_df: pd.DataFrame, + ) -> None: + LOG.debug(f"Product.set_cache({self.uuid=})") + 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, + 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, + pop_ledger_df=pop_ledger_df, ) - if self.user_wallet_enabled: - self.prebuild_user_wallet_balances(client=client, pop_ledger=pop_ledger) - self.prebuild_payouts(thl_lm=thl_lm, bp_pem=bp_pem) + if self.user_wallet_enabled: + self.prebuild_user_wallet_balances( + pop_ledger_df=pop_ledger_df, + ) + self.prebuild_payouts(bp_pem=bp_pem) self.prebuild_pop_financial( - thl_lm=thl_lm, ds=ds, client=client, pop_ledger=pop_ledger + thl_lm=thl_lm, + pop_ledger_df=pop_ledger_df, ) - # Validation steps. Don't save into redis until we confirm against - # the ledger. This allows parquet + db ledger balance checks - # The balance check needs to stop when the last parquet file was - # built, otherwise they'll appear unequal when it's really just - # a delay in the incite merge file not being built yet. - # bal = thl_lm.get_account_balance_timerange(time_end=) - - assert self.balance is not None rc = redis_config.create_redis_client() with rc.pipeline() as pipe: pipe.set( @@ -1501,11 +1422,18 @@ class Product(BaseModel, validate_assignment=True): value=self.model_dump_json(), ex=timedelta(days=3), ) - pipe.hset( - name=PRODUCT_BALANCES_METRICS_CACHE_KEY, - key=self.uuid, - value=self.balance.model_dump_json(), - ) + if self.balance is not None: + pipe.hset( + name=PRODUCT_BALANCES_METRICS_CACHE_KEY, + key=self.uuid, + value=self.balance.model_dump_json(), + ) + if self.private_balance is not None: + pipe.hset( + name=PRODUCT_PRIVATE_BALANCES_METRICS_CACHE_KEY, + key=self.uuid, + value=self.private_balance.model_dump_json(), + ) if self.user_wallet_balance is not None: pipe.hset( name=PRODUCT_USER_WALLET_BALANCES_METRICS_CACHE_KEY, diff --git a/pyproject.toml b/pyproject.toml index 53e4201..3424e20 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "generalresearch" -version = "3.6.9" +version = "3.7.0" description = "Python Utilities for General Research" readme = "README.md" requires-python = ">=3.14" diff --git a/tests/models/thl/test_product.py b/tests/models/thl/test_product.py index e45bd42..ce3242e 100644 --- a/tests/models/thl/test_product.py +++ b/tests/models/thl/test_product.py @@ -673,17 +673,17 @@ class TestProductFinancials: ) with pytest.raises(expected_exception=AssertionError) as cm: - p1.prebuild_balance( + p1.prebuild_balance_individual( thl_lm=thl_ledger_manager, ds=mnt_filepath, client=client_no_amm, ) - assert "Cannot build Product Balance" in str(cm.value) + assert "Cannot load Product POP ledger" 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( + p1.prebuild_balance_individual( thl_lm=thl_ledger_manager, ds=mnt_filepath, client=client_no_amm, @@ -700,7 +700,6 @@ class TestProductFinancials: body, content_type = p1.balance.to_prometheus() p1.prebuild_payouts( - thl_lm=thl_ledger_manager, bp_pem=brokerage_product_payout_event_manager, ) assert p1.payouts is not None @@ -735,7 +734,7 @@ class TestProductFinancials: ledger_collection.initial_load(client=None, sync=True) pop_ledger_merge.build(client=client_no_amm, ledger_coll=ledger_collection) - p1.prebuild_balance( + p1.prebuild_balance_individual( thl_lm=thl_ledger_manager, ds=mnt_filepath, client=client_no_amm, @@ -750,7 +749,6 @@ class TestProductFinancials: assert p1.balance.available_balance == 70 p1.prebuild_payouts( - thl_lm=thl_ledger_manager, bp_pem=brokerage_product_payout_event_manager, ) assert p1.payouts is not None @@ -783,7 +781,7 @@ class TestProductFinancials: ledger_collection.initial_load(client=None, sync=True) pop_ledger_merge.build(client=client_no_amm, ledger_coll=ledger_collection) - p1.prebuild_balance( + p1.prebuild_balance_individual( thl_lm=thl_ledger_manager, ds=mnt_filepath, client=client_no_amm, @@ -798,7 +796,6 @@ class TestProductFinancials: assert p1.balance.available_balance == 66 p1.prebuild_payouts( - thl_lm=thl_ledger_manager, bp_pem=brokerage_product_payout_event_manager, ) assert p1.payouts is not None @@ -856,22 +853,17 @@ class TestProductFinancials: 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( + df = p1.get_pop_ledger_df( + product_ids=[p1.uuid], thl_lm=thl_ledger_manager, - ds=mnt_filepath, + pop_ledger=pop_ledger_merge, client=client_no_amm, ) + + p1.prebuild_balance(thl_lm=thl_ledger_manager, pop_ledger_df=df) assert isinstance(p1.balance, ProductBalances) assert p1.balance.payout == 114 assert p1.balance.adjustment == 0 @@ -881,19 +873,105 @@ class TestProductFinancials: body, content_type = p1.balance.to_prometheus() - p1.prebuild_private_balance( - thl_lm=thl_ledger_manager, - ds=mnt_filepath, + p1.prebuild_private_balance(thl_lm=thl_ledger_manager, pop_ledger_df=df) + assert p1.private_balance.commission == 5 * 2 + + p1.prebuild_user_wallet_balances(pop_ledger_df=df) + assert p1.user_wallet_balance.outstanding_liability == 38 * 2 + body, content_type = p1.user_wallet_balance.to_prometheus() + + def test_balance_both_kinds_of_products( + 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, + thl_redis_config: RedisConfig, + brokerage_product_payout_event_manager: BrokerageProductPayoutEventManager, + ): + delete_ledger_db() + create_main_accounts() + delete_df_collection(coll=ledger_collection) + + p1: Product = product_factory( + business=gr_business, + ) + u1: User = user_factory(product=p1) + bp_wallet1 = thl_ledger_manager.get_account_or_create_bp_wallet(product=p1) + commission_wallet1 = thl_ledger_manager.get_account_or_create_bp_commission(p1) + p2: Product = product_factory( + business=gr_business, + user_wallet_config=UserWalletConfig(enabled=True), + payout_config=payout_config, + ) + u2: User = user_factory(product=p2) + bp_wallet2 = thl_ledger_manager.get_account_or_create_bp_wallet(product=p2) + user_wallet2 = thl_ledger_manager.get_account_or_create_user_wallet(user=u2) + commission_wallet2 = thl_ledger_manager.get_account_or_create_bp_commission(p2) + + session_with_tx_factory( + user=u1, + wall_req_cpi=Decimal("1.00"), + started=start, + ) + session_with_tx_factory( + user=u2, + wall_req_cpi=Decimal("1.00"), + started=start, + ) + + ledger_collection.initial_load(client=None, sync=True) + pop_ledger_merge.build(client=client_no_amm, ledger_coll=ledger_collection) + + df = Product.get_pop_ledger_df( + product_ids=[p1.uuid, p2.uuid], client=client_no_amm, + thl_lm=thl_ledger_manager, + pop_ledger=pop_ledger_merge, ) - assert p1.private_balance.commission == 5 * 2 + expected_accounts = [ + bp_wallet1.uuid, + bp_wallet2.uuid, + commission_wallet1.uuid, + commission_wallet2.uuid, + user_wallet2.uuid, + ] + assert set(df["account_id"]) == set(expected_accounts) + # The product_id column in pop ledger is set only for user txs + assert set(df["product_id"].dropna()) == {p2.uuid} - p1.prebuild_user_wallet_balances( - ds=mnt_filepath, + p1.prebuild_balance(thl_lm=thl_ledger_manager, pop_ledger_df=df) + assert p1.balance.payout == 95 + + p1.prebuild_private_balance(thl_lm=thl_ledger_manager, pop_ledger_df=df) + assert p1.private_balance.commission == 5 + + p2.prebuild_balance(thl_lm=thl_ledger_manager, pop_ledger_df=df) + assert p2.balance.payout == 57 + + p2.prebuild_private_balance(thl_lm=thl_ledger_manager, pop_ledger_df=df) + assert p2.private_balance.commission == 5 + + p2.prebuild_user_wallet_balances(pop_ledger_df=df) + assert p2.user_wallet_balance.outstanding_liability == 38 + + p1.set_cache( + thl_lm=thl_ledger_manager, client=client_no_amm, + bp_pem=brokerage_product_payout_event_manager, + redis_config=thl_redis_config, + pop_ledger=pop_ledger_merge, ) - assert p1.user_wallet_balance.outstanding_liability == 38 * 2 - body, content_type = p1.user_wallet_balance.to_prometheus() class TestProductBalance: @@ -1070,12 +1148,13 @@ class TestProductPOPFinancial: # --- test --- assert product.pop_financial is None - product.prebuild_pop_financial( + df = product.get_pop_ledger_df( + product_ids=[product.uuid], thl_lm=thl_ledger_manager, - ds=mnt_filepath, - client=client_no_amm, pop_ledger=pop_ledger_merge, + client=client_no_amm, ) + product.prebuild_pop_financial(thl_lm=thl_ledger_manager, pop_ledger_df=df) from generalresearch.models.thl.finance import POPFinancial @@ -1125,6 +1204,10 @@ class TestProductCache: create_main_accounts() delete_df_collection(coll=ledger_collection) + # In gr-api when we create a product, we create the bp wallet. + # maybe that should be standardized... + thl_ledger_manager.get_account_or_create_bp_wallet(product=product) + # Confirm the default / null behavior rc = thl_redis_config.create_redis_client() res: str | None = rc.get(product.cache_key) |
