aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--generalresearch/models/thl/product.py396
-rw-r--r--pyproject.toml2
-rw-r--r--tests/models/thl/test_product.py141
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)