diff options
| -rw-r--r-- | generalresearch/managers/base.py | 6 | ||||
| -rw-r--r-- | generalresearch/managers/thl/ledger_manager/thl_ledger.py | 2 | ||||
| -rw-r--r-- | generalresearch/managers/thl/payout.py | 997 | ||||
| -rw-r--r-- | generalresearch/models/gr/business.py | 14 | ||||
| -rw-r--r-- | generalresearch/models/thl/payout.py | 238 | ||||
| -rw-r--r-- | generalresearch/models/thl/product.py | 1 | ||||
| -rw-r--r-- | generalresearch/thl_django/event/models.py | 13 | ||||
| -rw-r--r-- | test_utils/models/conftest.py | 4 | ||||
| -rw-r--r-- | tests/managers/thl/test_payout.py | 472 | ||||
| -rw-r--r-- | tests/models/thl/test_payout.py | 113 |
10 files changed, 930 insertions, 930 deletions
diff --git a/generalresearch/managers/base.py b/generalresearch/managers/base.py index b935413..ba42e9e 100644 --- a/generalresearch/managers/base.py +++ b/generalresearch/managers/base.py @@ -1,6 +1,7 @@ from __future__ import annotations from collections.abc import Collection +from contextlib import nullcontext from enum import Enum from generalresearch.pg_helper import PostgresConfig @@ -45,6 +46,11 @@ class PostgresManager(Manager): self.pg_config = pg_config self.permissions = set(permissions) if permissions else set() + def connection(self, conn=None): + if conn is not None: + return nullcontext(conn) + return self.pg_config.make_connection() + class RedisManager(Manager): CACHE_PREFIX = None diff --git a/generalresearch/managers/thl/ledger_manager/thl_ledger.py b/generalresearch/managers/thl/ledger_manager/thl_ledger.py index 2fd0a2a..3119a46 100644 --- a/generalresearch/managers/thl/ledger_manager/thl_ledger.py +++ b/generalresearch/managers/thl/ledger_manager/thl_ledger.py @@ -837,7 +837,7 @@ class ThlLedgerManager(LedgerManager): tmc.TX_TYPE: TransactionType.BP_PAYOUT, tmc.EVENT: payoutevent_uuid, } - # This tag might will uniquely identify this tx + # This tag will uniquely identify this tx tag = f"{self.currency.value}:bp_payout:{payoutevent_uuid}" cash_account = self.get_account_cash() bp_wallet_account = self.get_account_or_create_bp_wallet(product) diff --git a/generalresearch/managers/thl/payout.py b/generalresearch/managers/thl/payout.py index 1f3742e..0efbf99 100644 --- a/generalresearch/managers/thl/payout.py +++ b/generalresearch/managers/thl/payout.py @@ -2,21 +2,27 @@ from __future__ import annotations from collections import defaultdict from collections.abc import Collection -from datetime import datetime, timedelta, timezone -from time import sleep +from datetime import datetime, timezone +from random import choice as rand_choice +from random import randint from typing import Any -from uuid import UUID, uuid4 +from uuid import uuid4 import numpy as np import pandas as pd +import psycopg from psycopg import sql from pydantic import AwareDatetime, NonNegativeInt, PositiveInt from generalresearch.currency import USDCent from generalresearch.decorators import LOG from generalresearch.managers.base import ( + Permission, PostgresManagerWithRedis, ) +from generalresearch.managers.thl.ledger_manager.exceptions import ( + LedgerTransactionConditionFailedError, +) from generalresearch.managers.thl.ledger_manager.thl_ledger import ( ThlLedgerManager, ) @@ -41,6 +47,8 @@ from generalresearch.models.thl.wallet.cashout_method import ( CashMailOrderData, CashoutRequestInfo, ) +from generalresearch.pg_helper import PostgresConfig +from generalresearch.redis_helper import RedisConfig class PayoutEventManager(PostgresManagerWithRedis): @@ -51,25 +59,6 @@ class PayoutEventManager(PostgresManagerWithRedis): """ - def set_account_lookup_table(self, thl_lm: ThlLedgerManager) -> None: - """This needs to run from grl-flow or from somewhere that has thl-redis - access - """ - - res = self.pg_config.execute_sql_query(query=f""" - SELECT uuid, reference_uuid - FROM ledger_account - WHERE qualified_name LIKE '{thl_lm.currency.value}:bp_wallet:%' - """) - account_to_product = {i["uuid"]: i["reference_uuid"] for i in res} - product_to_account = {i["reference_uuid"]: i["uuid"] for i in res} - - rc = self.redis_client - rc.hset(name="pem:account_to_product", mapping=account_to_product) - rc.hset(name="pem:product_to_account", mapping=product_to_account) - - return None - def get_by_uuid(self, pe_uuid: UUIDStr) -> PayoutEvent: res = self.pg_config.execute_sql_query( query=""" @@ -100,7 +89,7 @@ class PayoutEventManager(PostgresManagerWithRedis): order_data = order_data if order_data is not None else payout_event.order_data payout_event.update(status=status, ext_ref_id=ext_ref_id, order_data=order_data) - d = payout_event.model_dump_mysql() + d = payout_event.model_dump_postgres() query = sql.SQL(""" UPDATE event_payout SET status = %(status)s, @@ -312,7 +301,7 @@ class UserPayoutEventManager(PayoutEventManager): request_data=request_data or {}, order_data=order_data, ) - d = payout_event.model_dump_mysql() + d = payout_event.model_dump_postgres() with self.pg_config.make_connection() as conn: with conn.cursor() as c: @@ -345,43 +334,27 @@ class BrokerageProductPayoutEventManager(PayoutEventManager): def get_by_uuid( self, pe_uuid: UUIDStr, - # --- Support resources --- - account_product_mapping: dict[UUIDStr, UUIDStr] | None = None, ) -> BrokerageProductPayoutEvent: res = self.pg_config.execute_sql_query( query=""" - SELECT ep.uuid, - ep.debit_account_uuid, - ep.cashout_method_uuid, + SELECT ep.uuid, ep.debit_account_uuid, ep.cashout_method_uuid, ep.created, ep.amount, ep.status, ep.ext_ref_id, ep.payout_type, ep.request_data::jsonb, - ep.order_data::jsonb + ep.order_data::jsonb, + la.reference_uuid as product_id FROM event_payout AS ep + JOIN ledger_account la on la.uuid = debit_account_uuid WHERE ep.uuid = %s """, params=[pe_uuid], ) assert len(res) == 1, f"{pe_uuid} expected 1 result, got {len(res)}" - - d = res[0] - - # This isn't really need for creation... but we're doing it so that - # it can return back a full BrokerageProductPayoutEvent instance - if account_product_mapping is None: - rc = self.redis_client - account_product_mapping: dict = rc.hgetall(name="pem:account_to_product") - assert isinstance(account_product_mapping, dict) - - d["product_id"] = account_product_mapping[d["debit_account_uuid"]] - - return BrokerageProductPayoutEvent.model_validate(d) + return BrokerageProductPayoutEvent.model_validate(res[0]) @staticmethod def check_for_ledger_tx( thl_ledger_manager: ThlLedgerManager, - product_id: UUIDStr, - amount: USDCent, payout_event: BrokerageProductPayoutEvent, ) -> bool: """ @@ -394,6 +367,9 @@ class BrokerageProductPayoutEventManager(PayoutEventManager): are found, and raises a ValueError if something is inconsistent. """ tag = f"{thl_ledger_manager.currency.value}:bp_payout:{payout_event.uuid}" + amount = USDCent(payout_event.amount) + product_id = payout_event.product_id + txs = thl_ledger_manager.get_tx_by_tag(tag) if not txs: @@ -423,172 +399,80 @@ class BrokerageProductPayoutEventManager(PayoutEventManager): return True - def create( - self, - uuid: UUIDStr | None = None, - debit_account_uuid: UUIDStr | None = None, - created: AwareDatetimeISO = None, - amount: PositiveInt = None, - status: PayoutStatus | None = None, - ext_ref_id: str | None = None, - payout_type: PayoutType = None, - request_data: dict[str, Any] | None = None, - order_data: dict[str, Any] | CashMailOrderData | None = None, - # --- Support resources --- - account_product_mapping: dict[UUIDStr, UUIDStr] | None = None, - ) -> BrokerageProductPayoutEvent: - - if request_data is None: - request_data = dict() - - # This isn't really need for creation... but we're doing it so that - # it can return back a full BrokerageProductPayoutEvent instance - if account_product_mapping is None: - rc = self.redis_client - account_product_mapping: dict = rc.hgetall(name="pem:account_to_product") - assert isinstance(account_product_mapping, dict) - product_id = account_product_mapping[debit_account_uuid] - - bp_payout_event = BrokerageProductPayoutEvent( - uuid=uuid or uuid4().hex, - debit_account_uuid=debit_account_uuid, - cashout_method_uuid=self.CASHOUT_METHOD_UUID, - created=created or datetime.now(tz=timezone.utc), - amount=amount, - status=status, - ext_ref_id=ext_ref_id, - payout_type=payout_type, - request_data=request_data, - order_data=order_data, - product_id=product_id, - ) - d = bp_payout_event.model_dump_mysql() - - self.pg_config.execute_write( - query=""" - INSERT INTO event_payout ( - uuid, debit_account_uuid, created, cashout_method_uuid, amount, - status, ext_ref_id, payout_type, order_data, request_data - ) VALUES ( - %(uuid)s, %(debit_account_uuid)s, %(created)s, - %(cashout_method_uuid)s, %(amount)s, %(status)s, - %(ext_ref_id)s, %(payout_type)s, %(order_data)s, - %(request_data)s - ); - """, - params=d, - ) - - return bp_payout_event - def filter_by( self, - reference_uuid: str | None = None, ext_ref_id: str | None = None, debit_account_uuids: Collection[UUIDStr] | None = None, amount: int | None = None, created: datetime | None = None, created_after: datetime | None = None, - product_ids: str | None = None, - bp_user_ids: Collection[str] | None = None, + product_ids: Collection[str] | None = None, cashout_types: Collection[PayoutType] | None = None, statuses: Collection[PayoutStatus] | None = None, ) -> list[BrokerageProductPayoutEvent]: - """Try to retrieve payout events by the product_id/user_uuid, amount, - and optionally timestamp. + """Try to retrieve BP payout events. WARNING: This is only on the "payout events" table and nothing to - do with the Ledger itself. Therefore, the product_ids query - doesn't return Brokerage Product Payouts (the ACH or Wire events - to Suppliers) as part of the query. + do with the Ledger itself - *** IT IS ONLY FOR USER PAYOUTS *** - - Note: what used to be in thl-grpcs "ListCashoutRequests" calling - "list_cashout_requests" was merged into this. + *** IT IS ONLY FOR Brokerage Product PAYOUTS *** """ - args = [] + params = dict() filters = [] - if reference_uuid: - # This could be a product_id or a user_uuid - filters.append("la.reference_uuid = %s") - args.append(reference_uuid) if ext_ref_id: - # This is transaction id for tracking ACH/Wires with a banking - # institution - filters.append("ep.ext_ref_id = %s") - args.append(ext_ref_id) + # This is transaction id for tracking ACH/Wires with a banking institution + filters.append("ep.ext_ref_id = %(ext_ref_id)s") + params["ext_ref_id"] = ext_ref_id if debit_account_uuids: - # Or we could use the bp_wallet or user_wallet's account uuid - # instead of looking up by the product/user - filters.append("ep.debit_account_uuid = ANY(%s)") - args.append(debit_account_uuids) + # Or we could use the bp_wallet's account uuid + # instead of looking up by the product + filters.append("ep.debit_account_uuid = ANY(%(debit_account_uuids)s)") + params["debit_account_uuids"] = debit_account_uuids if amount: - filters.append("ep.amount = %s") - args.append(amount) + filters.append("ep.amount = %(amount)s") + params["amount"] = amount if created: - filters.append("ep.created = %s") - args.append(created.replace(tzinfo=None)) + filters.append("ep.created = %(created)s") + params["created"] = created if created_after: - filters.append("ep.created >= %s") - args.append(created_after.replace(tzinfo=None)) - if product_ids: - filters.append("product_id = ANY(%s)") - args.append(product_ids) - if bp_user_ids: - filters.append("product_user_id = ANY(%s)") - args.append(bp_user_ids) - if cashout_types: - filters.append("payout_type = ANY(%s)") - args.append([x.value for x in cashout_types]) - if statuses: - filters.append("status = ANY(%s)") - args.append([x.value for x in statuses]) + filters.append("ep.created >= %(created_after)s") + params["created_after"] = created_after + if product_ids is not None: + filters.append("la.reference_uuid = ANY(%(product_ids)s)") + params["product_ids"] = product_ids + if cashout_types is not None: + filters.append("payout_type = ANY(%(cashout_types)s)") + params["cashout_types"] = [x.value for x in cashout_types] + if statuses is not None: + filters.append("status = ANY(%(statuses)s)") + params["statuses"] = [x.value for x in statuses] assert len(filters) > 0, "must pass at least 1 filter" filter_str = " AND ".join(filters) + params["cashout_method_uuid"] = self.CASHOUT_METHOD_UUID res = self.pg_config.execute_sql_query( query=f""" - SELECT ep.uuid, - ep.debit_account_uuid, - ep.cashout_method_uuid, - ep.created, - ep.amount, ep.status, ep.ext_ref_id, ep.payout_type, + SELECT ep.uuid, ep.debit_account_uuid, ep.cashout_method_uuid, + ep.created, ep.amount, ep.status, ep.ext_ref_id, + ep.payout_type, ep.supplier_payout_id, ep.request_data::jsonb, ep.order_data::jsonb, ac.name as description, - la.reference_type as account_reference_type, - la.reference_uuid as account_reference_uuid + la.reference_uuid as product_id FROM event_payout AS ep LEFT JOIN accounting_cashoutmethod AS ac ON ep.cashout_method_uuid = ac.id LEFT JOIN ledger_account AS la ON la.uuid = ep.debit_account_uuid - LEFT JOIN thl_user u - ON la.reference_uuid = u.uuid - WHERE cashout_method_uuid = '{self.CASHOUT_METHOD_UUID}' + WHERE cashout_method_uuid = %(cashout_method_uuid)s + AND la.reference_type = 'bp' AND {filter_str} """, - params=args, + params=params, ) - - rc = self.redis_client - account_product_mapping = rc.hgetall(name="pem:account_to_product") - pes = [] - for d in res: - for k in [ - "uuid", - "debit_account_uuid", - "account_reference_uuid", - "cashout_method_uuid", - ]: - if d[k] is not None: - d[k] = UUID(d[k]).hex - - d["product_id"] = account_product_mapping[d["debit_account_uuid"]] - pes.append(BrokerageProductPayoutEvent.model_validate(d)) - + for row in res: + pes.append(BrokerageProductPayoutEvent.model_validate(row)) return pes def get_bp_payout_events_for_accounts( @@ -601,159 +485,64 @@ class BrokerageProductPayoutEventManager(PayoutEventManager): def get_bp_bp_payout_events_for_products( self, - thl_ledger_manager: ThlLedgerManager, product_uuids: Collection[UUIDStr], order_by: OrderBy | None = OrderBy.ASC, ) -> list[BrokerageProductPayoutEvent]: """This is a terrible name, but it returns the BPPayoutEvent model type rather than a list of PayoutEvents. - We do this for the Supplier centric APIs where they don't know, + We do this for the Supplier-centric APIs where they don't know or care about the underlying ledger account structure. """ assert len(product_uuids) > 0, "Must provide product_uuids" - accounts = thl_ledger_manager.get_accounts_bp_wallet_for_products( - product_uuids=product_uuids - ) - - assert len(accounts) == len(product_uuids), "Unequal Product & Account lists" + order_by = order_by or OrderBy.ASC - rc = self.redis_client - account_product_mapping = rc.hgetall(name="pem:account_to_product") - - payout_events: list[BrokerageProductPayoutEvent] = ( - self.get_bp_payout_events_for_accounts( - accounts=accounts, - ) + payout_events = self.filter_by( + product_ids=product_uuids, + cashout_types=[PayoutType.ACH], ) - - return BrokerageProductPayoutEvent.from_payout_events( - payout_events=payout_events, - account_product_mapping=account_product_mapping, - order_by=order_by, + payout_events = sorted( + payout_events, key=lambda x: x.created, reverse=order_by == OrderBy.DESC ) + return payout_events def retry_create_bp_payout_event_tx( self, thl_ledger_manager: ThlLedgerManager, product: Product, - payout_event_uuid: UUIDStr, - skip_wallet_balance_check: bool = False, - skip_one_per_day_check: bool = False, + bp_pe: BrokerageProductPayoutEvent, ) -> BrokerageProductPayoutEvent: """ If a create_bp_payout_event call fails, this can be called with the associated payoutevent. """ - bp_pe: BrokerageProductPayoutEvent = self.get_by_uuid(payout_event_uuid) - assert bp_pe.status == PayoutStatus.FAILED, "Only use this on failed payouts" - created = bp_pe.created - - assert not self.check_for_ledger_tx( - thl_ledger_manager=thl_ledger_manager, - payout_event=bp_pe, - product_id=bp_pe.product_id, - amount=bp_pe.amount_usd, - ), "Transaction exists! You should mark the payout event status as complete" + assert bp_pe.status in { + PayoutStatus.FAILED, + PayoutStatus.PENDING, + }, "Only use this on pending or failed payouts" - return self._create_tx_bp_payout_from_payout_event( - thl_ledger_manager=thl_ledger_manager, - bp_pe=bp_pe, - product=product, - amount=bp_pe.amount_usd, - created=created, - skip_one_per_day_check=skip_one_per_day_check, - skip_wallet_balance_check=skip_wallet_balance_check, - ) - - def create_bp_payout_event( - self, - thl_ledger_manager: ThlLedgerManager, - product: Product, - amount: USDCent, - payout_type: PayoutType = PayoutType.ACH, - ext_ref_id: str | None = None, - created: AwareDatetime | None = None, - skip_wallet_balance_check: bool = False, - skip_one_per_day_check: bool = False, - ) -> BrokerageProductPayoutEvent: - """This should be called when a BP is paid out money from their - wallet. Typically, this is an ACH payment. This function creates - the PayoutEvent and the Ledger entries. - - :param thl_ledger_manager: - :param product: The BP being paid. Assuming we're paying them out - of the balance of their USD wallet account. - :param amount: We're assuming everything is in USD, and we're - paying out a USD currency account. We could theoretically also - pay, for e.g. a Bitcoin account with a bitcoin transfer, but - this is not supported for now. - :param payout_type: PayoutType. default ACH - :param cashout_method_uuid: The entry in the - accounting_cashoutmethod table that records payment method - details. By default, the generic ACH cashout method (that has - no actual banking details). - - :param ext_ref_id: This is a unique ID for the Supplier Payment. - Typically it'll be from JP Morgan Chase, but may also just be - random if we can retrieve anything - - :param created: - - :param skip_wallet_balance_check: By default, this will fail unless - the BP's wallet actually has the amount requested. - - :param skip_one_per_day_check: Safety mechanism, checks if there - has already been a payout to this wallet in the past 24 hours. - - :return: - """ - - assert isinstance(amount, USDCent), "Must provide a USDCent" - - if created: - # Try to do a quick dupe check first before we create the payout event - pes = self.filter_by( - reference_uuid=product.id, amount=amount, created=created + if self.check_for_ledger_tx( + thl_ledger_manager=thl_ledger_manager, payout_event=bp_pe + ): + LOG.warning( + f"Transaction for {bp_pe.uuid=} {bp_pe.product_id=} already exists! " + f"Marking the payout event status as complete." ) - if len(pes) > 0: - raise ValueError(f"Payout event already exists!: {pes}") - - if created is None: - created = datetime.now(tz=timezone.utc) + self.update(payout_event=bp_pe, status=PayoutStatus.COMPLETE) + return bp_pe - # TODO: Explain why we're doing this. Why is it important to have - # Payout Events when the ledger has everything that should be - # needed. - bp_wallet = thl_ledger_manager.get_account_or_create_bp_wallet(product=product) - - bp_pe: BrokerageProductPayoutEvent = self.create( - debit_account_uuid=bp_wallet.uuid, - payout_type=payout_type, - amount=amount, - ext_ref_id=ext_ref_id, - created=created, - status=PayoutStatus.PENDING, - ) - return self._create_tx_bp_payout_from_payout_event( + return self.create_tx_bp_payout_from_payout_event( thl_ledger_manager=thl_ledger_manager, bp_pe=bp_pe, product=product, - amount=amount, - created=created, - skip_one_per_day_check=skip_one_per_day_check, - skip_wallet_balance_check=skip_wallet_balance_check, ) - def _create_tx_bp_payout_from_payout_event( + def create_tx_bp_payout_from_payout_event( self, thl_ledger_manager: ThlLedgerManager, bp_pe: BrokerageProductPayoutEvent, product: Product, - amount: USDCent, created: AwareDatetime | None = None, - skip_wallet_balance_check: bool = False, - skip_one_per_day_check: bool = False, ) -> BrokerageProductPayoutEvent: """ This should not be called directly. @@ -761,31 +550,30 @@ class BrokerageProductPayoutEventManager(PayoutEventManager): Handles exceptions: Check if the ledger tx actually exists or not, and set the payout event status accordingly. """ + created = created if created else bp_pe.created try: thl_ledger_manager.create_tx_bp_payout( product=product, - amount=amount, + amount=USDCent(bp_pe.amount), payoutevent_uuid=bp_pe.uuid, created=created, - skip_wallet_balance_check=skip_wallet_balance_check, - skip_one_per_day_check=skip_one_per_day_check, + skip_wallet_balance_check=True, + skip_one_per_day_check=True, ) - + except LedgerTransactionConditionFailedError as e: + if e.args[0] == "duplicate tag": + raise ValueError(f"""Payout event already exists! {e} + You are trying to create a tx that already exists. We can't know + if this is a new payout event with the same ref id, or you're + trying to run the same one twice ... So not setting the existing + payout event to FAILED, b/c the existing one is not failed! + Doing nothing ... + """) from e + self.update(payout_event=bp_pe, status=PayoutStatus.FAILED) + raise except Exception as e: - e.pe_uuid = bp_pe.uuid - if self.check_for_ledger_tx( - thl_ledger_manager=thl_ledger_manager, - product_id=product.uuid, - amount=amount, - payout_event=bp_pe, - ): - LOG.warning(f"Got exception {e} but ledger tx exists! Continuing ... ") - self.update(payout_event=bp_pe, status=PayoutStatus.COMPLETE) - return bp_pe - else: - LOG.warning(f"Got exception {e}. No ledger tx was created.") - self.update(payout_event=bp_pe, status=PayoutStatus.FAILED) - raise e + self.update(payout_event=bp_pe, status=PayoutStatus.FAILED) + raise self.update(payout_event=bp_pe, status=PayoutStatus.COMPLETE) return bp_pe @@ -814,156 +602,177 @@ class BrokerageProductPayoutEventManager(PayoutEventManager): return self.get_bp_payout_events_for_accounts(accounts=accounts) -class BusinessPayoutEventManager(BrokerageProductPayoutEventManager): +class BusinessPayoutEventManager(PostgresManagerWithRedis): - def update_ext_reference_ids( - self, - new_value: str, - current_value: str | None = None, - ) -> None: - """ - There are scenarios where an ACH/Wire payout event was saved with - a generic or anonymized reference identifier. We may want to be - able to go back and update all of those transaction IDs. + def __init__(self, *arg, **kwargs): + super().__init__(*arg, **kwargs) + self.bp_pe_manager = BrokerageProductPayoutEventManager(*arg, **kwargs) - """ + def get_by_ext_ref_id(self, ext_ref_id: str) -> BusinessPayoutEvent: + res = self.pg_config.execute_sql_query( + """ + SELECT + sp.*, + ep.bp_payouts + FROM supplier_payout sp + JOIN ( + SELECT + ep_inner.supplier_payout_id, + jsonb_agg( + to_jsonb(ep_inner) + || jsonb_build_object('product_id', la.reference_uuid) + ORDER BY ep_inner.created + ) AS bp_payouts + FROM event_payout ep_inner + JOIN ledger_account la + ON ep_inner.debit_account_uuid = la.uuid + GROUP BY ep_inner.supplier_payout_id + ) ep ON sp.id = ep.supplier_payout_id + WHERE sp.ext_ref_id = %(ext_ref_id)s + """, + {"ext_ref_id": ext_ref_id}, + ) + assert len(res) == 1, f"No Business Payout found with ext ref: {ext_ref_id}" + d = res[0] + for bp_payout in d["bp_payouts"]: + bp_payout["created"] = datetime.fromisoformat(bp_payout["created"]) + bpe = BusinessPayoutEvent.model_validate(d) + assert ( + bpe.bp_payouts is not None and len(bpe.bp_payouts) > 0 + ), "No BP payouts found for this Business Payout Event. This shouldn't happen!" + return bpe - if current_value is None: - raise ValueError("Dangerous to do ambiguous updates") + def filter_by( + self, + business_uuids: Collection[UUIDStr] | None = None, + ) -> list[BusinessPayoutEvent]: - # SELECT first to check that records exist - res = self.filter_by(ext_ref_id=current_value) - if len(res) == 0: - raise Warning("No event_payouts found to UPDATE") + params = dict() + filters = [] + if business_uuids is not None: + filters.append("business_id = ANY(%(business_uuids)s)") + params["business_uuids"] = business_uuids - # As of 2025, no single Business has more than 10,000 Products, - # leave the limit in as an additional safeguard. - query = """ - UPDATE event_payout - SET ext_ref_id = %s - WHERE ext_ref_id = %s - """ - with self.pg_config.make_connection() as conn: - with conn.cursor() as c: - c.execute(query=query, params=[new_value, current_value]) - assert c.rowcount < 10000 - conn.commit() + assert len(filters) > 0, "must pass at least 1 filter" + filter_str = " AND ".join(filters) - def delete_failed_business_payout(self, ext_ref_id: str, thl_lm: ThlLedgerManager): + res = self.pg_config.execute_sql_query( + f""" + SELECT + sp.*, + ep.bp_payouts + FROM supplier_payout sp + JOIN ( + SELECT + ep_inner.supplier_payout_id, + jsonb_agg( + to_jsonb(ep_inner) + || jsonb_build_object('product_id', la.reference_uuid) + ORDER BY ep_inner.created + ) AS bp_payouts + FROM event_payout ep_inner + JOIN ledger_account la + ON ep_inner.debit_account_uuid = la.uuid + GROUP BY ep_inner.supplier_payout_id + ) ep ON sp.id = ep.supplier_payout_id + WHERE {filter_str} + """, + params, + ) + bpes = [] + for row in res: + for bp_payout in row["bp_payouts"]: + bp_payout["created"] = datetime.fromisoformat(bp_payout["created"]) + bpe = BusinessPayoutEvent.model_validate(row) + assert ( + bpe.bp_payouts is not None and len(bpe.bp_payouts) > 0 + ), "No BP payouts found for this Business Payout Event. This shouldn't happen!" + bpes.append(bpe) + return bpes + + def validate_business_payout_in_ledger( + self, ext_ref_id: str, thl_lm: ThlLedgerManager + ): """ - Sometimes ACH/Wire payouts fail due to multiple reasons (timeouts, - Business Product having insufficient funds, etc). This is a utility - method that finds all event_payouts, and deletes them with all the - associated: - (1) Transactions - (2) Transaction Metadata - (3) Transaction Entries - - and then proceeds to delete them all in reverse order (so there is - no orphan / FK constraint issues). + Check that there exist ledger TXs for the Brokerage Product payouts + for this Business Payout Event. """ + bpe = self.get_by_ext_ref_id(ext_ref_id=ext_ref_id) + tags = [ + f"{thl_lm.currency.value}:bp_payout:{bp_pe.uuid}" + for bp_pe in bpe.bp_payouts + ] + txs = thl_lm.get_tx_ids_by_tags(tags=tags) + assert len(txs) == len( + bpe.bp_payouts + ), f"Expected {len(bpe.bp_payouts)} BP payouts but found {len(txs)}!" + return True - # (1) Find all by payout_event - event_payouts = self.filter_by(ext_ref_id=ext_ref_id) - if len(event_payouts) == 0: - raise Warning("No event_payouts found to DELETE") - - # sum([i["amount"] for i in event_payouts])/100 - event_payout_uuids = [i.uuid for i in event_payouts] - - # (2) Find all ledger_transactions - tags = [f"{thl_lm.currency.value}:bp_payout:{x}" for x in event_payout_uuids] - transactions = thl_lm.get_txs_by_tags(tags=tags) - transaction_ids = [tx.id for tx in transactions] - print("XXX1", transaction_ids) - # assert len(tags) == len(transactions) - - # (3) Find all ledger_transactionmetadata: assert two rows per tx - tx_metadata_ids = thl_lm.get_tx_metadata_ids_by_txs(transactions=transactions) - # assert len(tx_metadata) == len(transaction_ids)*2 - - # (4) Find all ledger_entry: assert two rows per tx - tx_entries = thl_lm.get_tx_entries_by_txs(transactions=transactions) - tx_entry_ids = [tx_entry.id for tx_entry in tx_entries] - # assert len(tx_entry) == len(transaction_ids)*2 - - # (5) Delete records - - # DELETE: tx_entry - self.pg_config.execute_write( - query=""" - DELETE - FROM ledger_entry - WHERE transaction_id = ANY(%s) - AND id = ANY(%s) - """, - params=[transaction_ids, tx_entry_ids], - ) + def resume_failed_business_payout( + self, ext_ref_id: str, thl_lm: ThlLedgerManager, pm: ProductManager + ): + """ + Sometimes a business payout's BP payouts fail due to multiple reasons + (timeouts, BP having insufficient funds, etc). Grab the PENDING + BP payout events and retry them. + """ + bpe = self.get_by_ext_ref_id(ext_ref_id=ext_ref_id) + assert bpe.id + assert bpe.bp_payouts - # DELETE: tx_metadata - self.pg_config.execute_write( - query=""" - DELETE - FROM ledger_transactionmetadata - WHERE transaction_id = ANY(%s) - AND id = ANY(%s) - """, - params=[transaction_ids, list(tx_metadata_ids)], - ) + if all(bp_pe.status == PayoutStatus.COMPLETE for bp_pe in bpe.bp_payouts): + try: + self.validate_business_payout_in_ledger( + ext_ref_id=ext_ref_id, thl_lm=thl_lm + ) + except AssertionError as e: + raise AssertionError( + f"Business Payout Event {ext_ref_id} is COMPLETE but BP payouts are not in the ledger! {e} " + f"This typically shouldn't happen, as if the ledger TX fails, the event_payout " + f"status won't be COMPLETE. If it does, set all the bp statuses to FAILED, and " + f"then try again. Any that do exist in the ledger will be found and marked COMPLETE." + ) from e + if bpe.status != PayoutStatus.COMPLETE: + self.update_business_payout_event( + pk=bpe.id, status=PayoutStatus.COMPLETE + ) + LOG.warning( + "All BP payouts complete, setting Business Payout Event status to COMPLETE." + ) + else: + LOG.warning( + "Nothing to do! Business Payout is COMPLETE and all Brokerage Product payouts are also COMPLETE!" + ) + return None - # DELETE: transactions - self.pg_config.execute_write( - query=""" - DELETE - FROM ledger_transaction - WHERE id = ANY(%s) - """, - params=[transaction_ids], - ) + for bp_pe in bpe.bp_payouts: + if bp_pe.status in {PayoutStatus.PENDING, PayoutStatus.FAILED}: + LOG.warning( + f"Found a {bp_pe.status} BP payout event: {bp_pe.uuid} - retrying ... " + ) + product = pm.get_by_uuid(bp_pe.product_id) + self.retry_create_bp_payout_event_tx( + thl_ledger_manager=thl_lm, bp_pe=bp_pe, product=product + ) + if bp_pe.status != PayoutStatus.COMPLETE: + raise ValueError(f"{bp_pe.uuid} has {bp_pe.status=}. Please check me.") - # DELETE: event_payouts - self.pg_config.execute_write( - query=""" - DELETE - FROM event_payout - WHERE ext_ref_id = %s - AND uuid = ANY(%s) - """, - params=[ext_ref_id, event_payout_uuids], - ) + self.validate_business_payout_in_ledger(ext_ref_id=ext_ref_id, thl_lm=thl_lm) + self.update_business_payout_event(pk=bpe.id, status=PayoutStatus.COMPLETE) return None - def get_business_payout_events_for_products( + def get_business_payout_events_for_business( self, - thl_ledger_manager: ThlLedgerManager, - product_uuids: Collection[UUIDStr], + business_uuid: UUIDStr, order_by: OrderBy | None = OrderBy.ASC, ) -> list[BusinessPayoutEvent]: - res = self.get_bp_bp_payout_events_for_products( - thl_ledger_manager=thl_ledger_manager, - product_uuids=product_uuids, - order_by=order_by, + order_by = order_by or OrderBy.ASC + bpes = self.filter_by( + business_uuids=[business_uuid], ) - - return self.from_bp_payout_events(bp_payout_events=res) - - @staticmethod - def from_bp_payout_events( - bp_payout_events: Collection[BrokerageProductPayoutEvent], - ) -> list[BusinessPayoutEvent]: - if len(bp_payout_events) == 0: - return [] - - grouped = defaultdict(list) - for bp_pe in bp_payout_events: - grouped[bp_pe.ext_ref_id].append(bp_pe) - - res = [] - for _, members in grouped.items(): - res.append(BusinessPayoutEvent.model_validate({"bp_payouts": members})) - - return res + bpes = sorted(bpes, key=lambda x: x.created, reverse=order_by == OrderBy.DESC) + return bpes @staticmethod def recoup_proportional( @@ -1091,122 +900,302 @@ class BusinessPayoutEventManager(BrokerageProductPayoutEventManager): return allocation + def create_business_payout_event( + self, + bpe: BusinessPayoutEvent, + ): + assert bpe.bp_payouts, "Must provide at least one BP Payout" + assert {bp_pe.status for bp_pe in bpe.bp_payouts} == { + PayoutStatus.PENDING + }, "All BP Payouts must be PENDING" + assert bpe.id is None, "Cannot create a BusinessPayoutEvent with an existing ID" + INSERT_SUPPLIER_PAYOUT = """ + INSERT INTO supplier_payout ( + business_id, created, amount, + status, ext_ref_id, payout_type, + request_data, order_data + ) VALUES ( + %(business_id)s, %(created)s, %(amount)s, + %(status)s, %(ext_ref_id)s, %(payout_type)s, + %(request_data)s, %(order_data)s + ) RETURNING id; + """ + INSERT_BP_PAYOUT = """ + INSERT INTO event_payout ( + uuid, debit_account_uuid, created, cashout_method_uuid, + amount, status, ext_ref_id, payout_type, order_data, + request_data, supplier_payout_id + ) VALUES ( + %(uuid)s, %(debit_account_uuid)s, %(created)s, %(cashout_method_uuid)s, + %(amount)s, %(status)s, %(ext_ref_id)s, %(payout_type)s, %(order_data)s, + %(request_data)s, %(supplier_payout_id)s + ); + """ + + with self.pg_config.make_connection() as conn: + with conn.cursor() as c: + # ext_ref_id (transaction_id) has a unique constraint + try: + c.execute(INSERT_SUPPLIER_PAYOUT, bpe.model_dump_postgres()) + except psycopg.errors.UniqueViolation as e: + if e.diag.constraint_name == "supplier_payout_ext_ref_id_key": + raise ValueError( + f"Cannot create a BusinessPayoutEvent with an existing " + f"transaction_id. {e.diag.message_detail}" + ) + raise + supplier_payout_pk = c.fetchone()["id"] + bpe.id = supplier_payout_pk + for bp_pe in bpe.bp_payouts: + c.execute( + INSERT_BP_PAYOUT, + bp_pe.model_dump_postgres() + | {"supplier_payout_id": supplier_payout_pk}, + ) + conn.commit() + def create_from_ach_or_wire( self, business: Business, amount: USDCent, + transaction_id: str, pm: ProductManager, thl_lm: ThlLedgerManager, created: datetime | None = None, - transaction_id: str | None = None, ) -> BusinessPayoutEvent | None: """This records a single banking transfer to a supplier. Takes a specific Business that was paid out and how much. It then determines how to distribute the amount to each Brokerage Product in the Business. - - :param business - :param amount - :param pm - :param thl_lm: this must have rw permissions to add transactions to - the ledger - :param created - :param transaction_id - - :return: """ assert business.balance is not None, ( "Must provide a full version of a Business in order to calculate" "the required Brokerage Product amounts." ) - assert amount > 100_00, "Must issue Supplier Payouts at least $100 minimum." - LOG.warning("Paying out ") + assert amount >= 100_00, "Must issue Supplier Payouts at least $100 minimum." + LOG.warning(f"Paying out {business.name} {amount.to_usd_str()}") if created: LOG.warning("Payouts in the past, require the parquet files to be rebuilt.") - assert created < datetime.now(tz=timezone.utc) - + assert created.tzinfo == timezone.utc, "created must be UTC" + assert created < datetime.now( + tz=timezone.utc + ), "created must be in the past" else: created = datetime.now(tz=timezone.utc) # Gather the total amount available balance from each and put into # a simple DF. We're using the available balance because we need it - # to always be positive.. and we never want to get into a negative + # to always be positive. We never want to get into a negative # situation again, so it's best to be extra conservative. - res = { + balances = { pb.product_id: pb.available_balance for pb in business.balance.product_balances } - df = pd.DataFrame.from_dict(res, orient="index").reset_index() + df = pd.DataFrame.from_dict(balances, orient="index").reset_index() df.columns = ["product_id", "available_balance"] - res = BusinessPayoutEventManager.recoup_proportional( + df = BusinessPayoutEventManager.recoup_proportional( df=df, target_amount=business.balance.recoup ) # Can't pay any Products that don't have a remaining balance - res = res[res["remaining_balance"] > 0] + df = df[df["remaining_balance"] > 0].copy() assert ( - res.deduction.sum() == business.balance.recoup + df.deduction.sum() == business.balance.recoup ), "recoup_proportional failure" - res["issue_amount"] = BusinessPayoutEventManager.distribute_amount( - df=res, amount=amount + df["issue_amount"] = BusinessPayoutEventManager.distribute_amount( + df=df, amount=amount ) - assert res.issue_amount.sum() == amount, "issue_amount failure" + assert df.issue_amount.sum() == amount, "issue_amount failure" # Can't pay any Products that don't have an issue amount - res = res[res["issue_amount"] > 0] - - recouped_amounts: list[dict[str, int]] = res[ - ["product_id", "remaining_balance", "issue_amount"] - ].to_dict(orient="records") - - # Get all of the products at once so we're not doing it for every interation - products = pm.get_by_uuids( - product_uuids=[i["product_id"] for i in recouped_amounts] + df = df[df["issue_amount"] > 0].copy() + + amounts: dict[str, dict[str, int]] = df.set_index("product_id")[ + ["remaining_balance", "issue_amount"] + ].to_dict(orient="index") + + products = pm.get_by_uuids(product_uuids=list(amounts.keys())) + product_lookup = {p.uuid: p for p in products} + + # Bulk version of this ---v + # bp_wallet = thl_lm.get_account_or_create_bp_wallet(product=product) + qualified_names = [ + f"{thl_lm.currency.value}:bp_wallet:{bp.id}" for bp in products + ] + bp_wallets = thl_lm.get_accounts(qualified_names) + wallet_lookup = {bpw.reference_uuid: bpw.uuid for bpw in bp_wallets} + + bpe = BusinessPayoutEvent( + id=None, + business_id=business.uuid, + payout_type=PayoutType.ACH, + amount=amount, + created=created, + ext_ref_id=transaction_id, + # The ACH payment was sent! We haven't yet recorded it + # in the ledger, but it was sent by the bank. This + # is kind of ambiguous the meaning, we'll say it + # is not yet COMPLETE b/c the bp payouts + # haven't all been created yet. + status=PayoutStatus.APPROVED, ) bp_payouts: list[BrokerageProductPayoutEvent] = [] - for idx, item in enumerate(recouped_amounts): - product = next((p for p in products if p.uuid == item["product_id"]), None) - assert product is not None - - try: - bp_pe: BrokerageProductPayoutEvent = self.create_bp_payout_event( - thl_ledger_manager=thl_lm, - product=product, + for product_id, item in amounts.items(): + product = product_lookup[product_id] + bp_payouts.append( + BrokerageProductPayoutEvent( + created=created, + payout_type=PayoutType.ACH, + status=PayoutStatus.PENDING, + uuid=uuid4().hex, amount=USDCent(item["issue_amount"]), - created=created + timedelta(milliseconds=idx + 1), ext_ref_id=transaction_id, - skip_wallet_balance_check=True, + product_id=product.uuid, + cashout_method_uuid=self.bp_pe_manager.CASHOUT_METHOD_UUID, + debit_account_uuid=wallet_lookup[product_id], ) + ) + bpe.bp_payouts = bp_payouts + # The supplier_payout db row and all event_payout (BP rows) are all + # created in the same DB transaction. + self.create_business_payout_event(bpe=bpe) + assert bpe.id is not None, "Something failed creating BusinessPayoutEvent" + + # Now, go through each and create ledger txs. This is resumable + # from the BrokerageProductPayoutEvents + for bp_pe in bpe.bp_payouts: + product = product_lookup[bp_pe.product_id] + self.bp_pe_manager.create_tx_bp_payout_from_payout_event( + thl_ledger_manager=thl_lm, + bp_pe=bp_pe, + product=product, + ) - assert bp_pe.status == PayoutStatus.COMPLETE - bp_payouts.append(bp_pe) + self.update_business_payout_event(pk=bpe.id, status=PayoutStatus.COMPLETE) - except (Exception,) as e: - # Cleanup bp_payouts - print("Exception", e) - return None + return bpe - if bp_pe.status == PayoutStatus.FAILED: - sleep(1) + def update_business_payout_event(self, pk: int, status: PayoutStatus): + with self.connection() as conn: + with conn.cursor() as c: + c.execute( + """ + UPDATE supplier_payout + SET status = %(status)s + WHERE id = %(pk)s""", + {"pk": pk, "status": status}, + ) + assert c.rowcount == 1, f"{id=} not found" + conn.commit() + return None - try: - bp_pe = self.retry_create_bp_payout_event_tx( - thl_ledger_manager=thl_lm, - product=product, - payout_event_uuid=bp_pe.uuid, - ) - assert bp_pe.status == PayoutStatus.COMPLETE - bp_payouts.append(bp_pe) + def create_bp_payout_event( + self, + thl_ledger_manager: ThlLedgerManager, + product: Product, + amount: USDCent, + ext_ref_id: str, + created: datetime | None = None, + ): + """ + This should NOT be called directly normally. It is just a shortcut + for tests. However, instead of just making a naked BP payout, + it created the business payout also, but with just one BP Payout + """ + created = created or datetime.now(tz=timezone.utc) + account = thl_ledger_manager.get_account( + f"{thl_ledger_manager.currency.value}:bp_wallet:{product.uuid}" + ) + bpe = BusinessPayoutEvent( + id=None, + business_id=product.business_uuid, + payout_type=PayoutType.ACH, + amount=amount, + created=created, + ext_ref_id=ext_ref_id, + status=PayoutStatus.APPROVED, + bp_payouts=[ + BrokerageProductPayoutEvent( + created=created, + payout_type=PayoutType.ACH, + status=PayoutStatus.PENDING, + uuid=uuid4().hex, + amount=amount, + ext_ref_id=ext_ref_id, + product_id=product.uuid, + cashout_method_uuid=self.bp_pe_manager.CASHOUT_METHOD_UUID, + debit_account_uuid=account.uuid, + ) + ], + ) + self.create_business_payout_event(bpe=bpe) + self.bp_pe_manager.create_tx_bp_payout_from_payout_event( + thl_ledger_manager=thl_ledger_manager, + bp_pe=bpe.bp_payouts[0], + product=product, + ) + self.update_business_payout_event(pk=bpe.id, status=PayoutStatus.COMPLETE) + return bpe + + def update_ext_reference_ids( + self, + new_value: str, + current_value: str, + ) -> None: + """ + There are scenarios where an ACH/Wire payout event was saved with + a generic or anonymized reference identifier. We may want to be + able to go back and update all of those transaction IDs. + + """ + assert new_value and current_value + + # Will raise if doesn't exist + self.get_by_ext_ref_id(ext_ref_id=current_value) + + query1 = """ + UPDATE supplier_payout + SET ext_ref_id = %(new_value)s + WHERE ext_ref_id = %(old_value)s + """ + query2 = """ + UPDATE event_payout + SET ext_ref_id = %(new_value)s + WHERE ext_ref_id = %(old_value)s + """ + params = {"new_value": new_value, "old_value": current_value} + with self.pg_config.make_connection() as conn: + with conn.cursor() as c: + c.execute(query1, params) + assert c.rowcount == 1 + c.execute(query2, params) + # As of 2025, no single Business has more than 10,000 Products, + # leave the limit in as an additional safeguard. + assert c.rowcount < 10000 + conn.commit() - except (Exception,) as e: - # Cleanup bp_payouts - return None - return BusinessPayoutEvent.model_validate({"bp_payouts": bp_payouts}) +# import duckdb +# conn = duckdb.connect() +# conn.execute(""" +# select * from read_parquet('/mnt/thl-incite/raw/df-collections/ledger/*/*.parquet') +# where event_payout is not null +# and direction =1 +# and reference_uuid in ? +# """, [b.product_uuids]) +# df = conn.fetch_df() +# df['ext_description'].value_counts() +# +# tx_ids = [35554404, 37210650] +# conn.execute(""" +# select * from read_parquet('/mnt/thl-incite/raw/df-collections/ledger/*/*.parquet') +# where tx_id in ? +# """, [tx_ids]) +# df = conn.fetch_df() diff --git a/generalresearch/models/gr/business.py b/generalresearch/models/gr/business.py index 70aafc6..51317da 100644 --- a/generalresearch/models/gr/business.py +++ b/generalresearch/models/gr/business.py @@ -434,28 +434,22 @@ class Business(BaseModel): def prebuild_payouts( self, - thl_pg_config: PostgresConfig, - thl_lm: ThlLedgerManager, bpem: BusinessPayoutEventManager, ) -> None: LOG.debug(f"Business.prebuild_payouts({self.uuid=})") - self.prefetch_products(thl_pg_config=thl_pg_config) - - self.payouts = bpem.get_business_payout_events_for_products( - thl_ledger_manager=thl_lm, - product_uuids=self.product_uuids, + self.payouts = bpem.get_business_payout_events_for_business( + business_uuid=self.uuid, order_by=OrderBy.DESC, ) - self.prebuild_payouts_total() + return None def prebuild_payouts_total(self): assert self.payouts is not None self.payouts_total = USDCent(sum([po.amount for po in self.payouts])) self.payouts_total_str = self.payouts_total.to_usd_str() - - return + return None def prebuild_pop_financial( self, diff --git a/generalresearch/models/thl/payout.py b/generalresearch/models/thl/payout.py index 1a9d534..835c3b9 100644 --- a/generalresearch/models/thl/payout.py +++ b/generalresearch/models/thl/payout.py @@ -11,11 +11,14 @@ from pydantic import ( PositiveInt, computed_field, field_validator, + model_validator, + ConfigDict, ) +from pydantic.json_schema import SkipJsonSchema from typing_extensions import Self from generalresearch.currency import USDCent -from generalresearch.models.custom_types import AwareDatetimeISO, UUIDStr +from generalresearch.models.custom_types import AwareDatetimeISO, UUIDStr, UUIDStrCoerce from generalresearch.models.thl.definitions import PayoutStatus from generalresearch.models.thl.ledger import OrderBy from generalresearch.models.thl.wallet import PayoutType @@ -39,20 +42,20 @@ class PayoutEvent(BaseModel): multiple BrokerageProductPayoutEvents. """ - uuid: UUIDStr = Field( + uuid: UUIDStrCoerce = Field( title="Payout Event Unique Identifier", default_factory=lambda: uuid4().hex, examples=["9453cd076713426cb68d05591c7145aa"], ) - debit_account_uuid: UUIDStr = Field( + debit_account_uuid: UUIDStrCoerce | None = Field( description="The LedgerAccount.uuid that money is being requested from. " "Thie User or Brokerage Product is retrievable through the " "LedgerAccount.reference_uuid", examples=["18298cb1583846fbb06e4747b5310693"], ) - cashout_method_uuid: UUIDStr = Field( + cashout_method_uuid: UUIDStrCoerce | None = Field( description="References a row in the account_cashoutmethod table. This " "is the cashout method that was used to request this " "payout. (A cashout is the same thing as a payout)", @@ -84,8 +87,8 @@ class PayoutEvent(BaseModel): description=PayoutType.as_openapi(), examples=[PayoutType.ACH] ) - request_data: dict = Field( - default_factory=dict, + request_data: dict | None = Field( + default=None, description="Stores payout-type-specific information that is used to " "request this payout from the external provider.", ) @@ -144,16 +147,15 @@ class PayoutEvent(BaseModel): # --- ORM --- - def model_dump_mysql(self, *args, **kwargs) -> dict: - d = self.model_dump(mode="json", *args, **kwargs) - - if "created" in d: - d["created"] = self.created.replace(tzinfo=None) - - if d.get("request_data") is not None: - d["request_data"] = json.dumps(self.request_data) + def model_dump_postgres(self) -> dict: + d = self.model_dump(mode="json", exclude={"request_data", "order_data"}) - if d.get("order_data") is not None: + d["request_data"] = ( + json.dumps(self.request_data) if self.request_data is not None else None + ) + if self.order_data is None: + d["order_data"] = None + else: if isinstance(self.order_data, dict): d["order_data"] = json.dumps(self.order_data) else: @@ -196,7 +198,7 @@ class BrokerageProductPayoutEvent(PayoutEvent): - created: When the Brokerage Product was paid out """ - product_id: UUIDStr = Field( + product_id: UUIDStrCoerce = Field( description="The Brokerage Product that was paid out", examples=["1108d053e4fa47c5b0dbdcd03a7981e7"], ) @@ -220,136 +222,140 @@ class BrokerageProductPayoutEvent(PayoutEvent): def amount_usd_str(self) -> str: return self.amount_usd.to_usd_str() - # --- ORM --- - @classmethod - def from_payout_event( - cls, - pe: PayoutEvent, - account_product_mapping: dict[UUIDStr, UUIDStr] | None = None, - redis_config: RedisConfig | None = None, - ) -> Self: - # TODO!: prevent re-assignment, rework this... - - if account_product_mapping is None: - rc = redis_config.create_redis_client() - account_product_mapping: dict = rc.hgetall(name="pem:account_to_product") - assert isinstance(account_product_mapping, dict) - assert pe.uuid in account_product_mapping.keys() - - d = pe.model_dump() - d["product_id"] = account_product_mapping[pe.debit_account_uuid] - return cls.model_validate(d) - - @classmethod - def from_payout_events( - cls, - payout_events: Collection[PayoutEvent], - order_by=OrderBy, - account_product_mapping: dict[UUIDStr, UUIDStr] | None = None, - redis_config: RedisConfig | None = None, - ) -> list[Self]: - # TODO!: prevent re-assignment, rework this... - - if account_product_mapping is None: - rc = redis_config.create_redis_client() - account_product_mapping: dict = rc.hgetall(name="pem:account_to_product") - assert isinstance(account_product_mapping, dict) - - res = [] - for pe in payout_events: - res.append( - cls.from_payout_event( - pe=pe, account_product_mapping=account_product_mapping - ) - ) +class BusinessPayoutEvent(BaseModel): + """A single payout event to a supplier Business.""" - match order_by: - case OrderBy.ASC: - sorted_list = sorted(res, key=lambda x: x.created, reverse=False) - case OrderBy.DESC: - sorted_list = sorted(res, key=lambda x: x.created, reverse=True) - case _: - raise ValueError("Invalid order provided..") + model_config = ConfigDict(validate_assignment=True) - return sorted_list + id: SkipJsonSchema[PositiveInt | None] = Field(exclude=True, default=None) + # Used for holding a *unique*, external, payout-type-specific identifier. + ext_ref_id: str = Field(title="Unique external reference ID") -class BusinessPayoutEvent(BaseModel): - """A single ACH or Wire event to a Business Bank Account""" + business_id: UUIDStr = Field( + description="The Business receiving this supplier payout.", + examples=[uuid4().hex], + ) - bp_payouts: list[BrokerageProductPayoutEvent] = Field( - description="Here is the list of Brokerage Product Payouts that" - "this Business Payout includes.", - min_length=1, + created: AwareDatetimeISO = Field( + default_factory=lambda: datetime.now(tz=timezone.utc) ) - @computed_field( + # In the smallest unit of the currency being transacted. For USD, this + # is cents. + amount: PositiveInt = Field( + lt=2**63 - 1, + strict=True, title="Amount", - description="The amount issued to the Bank Account", - examples=[19_823_43], - return_type=USDCent, + description="The amount issued to the supplier.", + examples=[1_982_343], ) - @property - def amount(self) -> USDCent: - return USDCent(sum([p.amount for p in self.bp_payouts])) - @computed_field( - title="Amount USD Str", - description="The amount issued to the Bank Account as a USD string", - examples=["$19,823.43"], - return_type=str, + status: PayoutStatus = Field( + default=PayoutStatus.PENDING, + description=PayoutStatus.as_openapi(), + examples=[PayoutStatus.COMPLETE], ) - @property - def amount_usd_str(self) -> str: - return self.amount.to_usd_str() - @computed_field( - title="Created", - description="This is equal to the created time of the first" - "Brokerage Product Payout Event.", - return_type=AwareDatetimeISO, + payout_type: PayoutType = Field( + description=PayoutType.as_openapi(), examples=[PayoutType.ACH] ) - @property - def created(self) -> AwareDatetimeISO: - return self.bp_payouts[0].created - @computed_field( - title="Line Items", - description="The number of sub-payments", - return_type=PositiveInt, + request_data: dict | None = Field( + default=None, + description="Stores payout-type-specific information that is used to " + "request this payout from the external provider.", + ) + + order_data: dict | None = Field( + default=None, + description="Stores payout-type-specific order information that is " + "returned from the external payout provider.", + ) + + bp_payouts: list[BrokerageProductPayoutEvent] | None = Field( + default=None, + description="The list of Brokerage Product Payouts that this Business Payout includes", + min_length=1, ) - @property - def line_items(self): - return len(self.bp_payouts) @computed_field( - title="External Reference ID", - description="ACH Transaction ID", - return_type=str | None, + title="Amount USD Str", + description="The amount issued to the supplier as a USD string", + examples=["$19,823.43"], + return_type=str, ) @property - def ext_ref_id(self): - return self.bp_payouts[0].ext_ref_id + def amount_usd_str(self) -> str: + return USDCent(self.amount).to_usd_str() # --- Validators --- + @field_validator("payout_type", mode="before") + @classmethod + def normalize_payout_type(cls, v): + if isinstance(v, str): + try: + return PayoutType[v.upper()] + except KeyError: + raise ValueError(f"Invalid payout_type: {v}") + return v + @field_validator("bp_payouts", mode="before") @classmethod - def normalize_enum(cls, v): + def validate_bp_payouts_type(cls, v): """This can be a list of Instances or Python Dictionaries depending on how it's initialized. """ + if v is None: + return v + assert isinstance(v, list) + return v - def get_field(obj, field): - if isinstance(obj, dict): - return obj.get(field) - return getattr(obj, field, None) + @model_validator(mode="after") + def validate_bp_payouts(self) -> Self: + if not self.bp_payouts: + return self - assert all( - get_field(i, "ext_ref_id") == get_field(v[0], "ext_ref_id") for i in v - ), "Not all group values are the same" + bp_payout_amount = sum([p.amount for p in self.bp_payouts]) + if bp_payout_amount != self.amount: + raise ValueError( + "BusinessPayoutEvent.amount must equal the sum of " + f"bp_payouts amounts ({self.amount=} {bp_payout_amount=})" + ) - return v + invalid_payout_types = [ + p.payout_type for p in self.bp_payouts if p.payout_type != self.payout_type + ] + if invalid_payout_types: + raise ValueError( + "All BrokerageProductPayoutEvent.payout_type values must equal " + f"BusinessPayoutEvent.payout_type ({self.payout_type})" + ) + + invalid_ext_ids = [ + p.ext_ref_id for p in self.bp_payouts if p.ext_ref_id != self.ext_ref_id + ] + if invalid_ext_ids: + raise ValueError( + "All BrokerageProductPayoutEvent.ext_ref_id values must equal " + f"BusinessPayoutEvent.ext_ref_id ({self.ext_ref_id})" + ) + + return self + + def model_dump_postgres(self): + d = self.model_dump( + mode="json", + exclude={"bp_payouts"}, + ) + d["request_data"] = ( + json.dumps(self.request_data) if self.request_data is not None else None + ) + d["order_data"] = ( + json.dumps(self.order_data) if self.order_data is not None else None + ) + return d diff --git a/generalresearch/models/thl/product.py b/generalresearch/models/thl/product.py index 3f21d92..9b7d66a 100644 --- a/generalresearch/models/thl/product.py +++ b/generalresearch/models/thl/product.py @@ -1207,7 +1207,6 @@ class Product(BaseModel, validate_assignment=True): from generalresearch.models.thl.ledger import OrderBy self.payouts = bp_pem.get_bp_bp_payout_events_for_products( - thl_ledger_manager=thl_lm, product_uuids=[self.uuid], order_by=OrderBy.DESC, ) diff --git a/generalresearch/thl_django/event/models.py b/generalresearch/thl_django/event/models.py index 51a8e2f..6d92b20 100644 --- a/generalresearch/thl_django/event/models.py +++ b/generalresearch/thl_django/event/models.py @@ -51,7 +51,11 @@ class SupplierPayout(models.Model): to the Business. """ - uuid = models.UUIDField(default=uuid.uuid4, primary_key=True) + id = models.BigAutoField(primary_key=True) + + # Used for holding a unique, external, payouttype-specific identifier. + # For ACH, this is the ACH transaction id. + ext_ref_id = models.CharField(max_length=64, unique=True) # The Business receiving this payout business_id = models.UUIDField(null=True) @@ -64,10 +68,6 @@ class SupplierPayout(models.Model): # generalresearch/models/thl/payout.py:PayoutStatus status = models.CharField(max_length=20, null=True) - # Used for holding an external, payouttype-specific identifier. - # For ACH, this is the ACH transaction id. - ext_ref_id = models.CharField(max_length=64, null=True) - # The allowed values for `payout_type` are defined in generalresearch: # generalresearch/models/thl/payout.py:PayoutType payout_type = models.CharField(max_length=14) @@ -86,7 +86,6 @@ class SupplierPayout(models.Model): indexes = [ models.Index(fields=["created"]), models.Index(fields=["business_id"]), - models.Index(fields=["ext_ref_id"]), ] @@ -133,7 +132,6 @@ class Payout(models.Model): # payout transaction that this product-level split belongs to. supplier_payout = models.ForeignKey( SupplierPayout, - db_column="supplier_payout_uuid", null=True, on_delete=models.DO_NOTHING, ) @@ -145,5 +143,4 @@ class Payout(models.Model): models.Index(fields=["created"]), models.Index(fields=["debit_account_uuid"]), models.Index(fields=["ext_ref_id"]), - models.Index(fields=["supplier_payout"]), ] diff --git a/test_utils/models/conftest.py b/test_utils/models/conftest.py index 9925a9e..89f6f32 100644 --- a/test_utils/models/conftest.py +++ b/test_utils/models/conftest.py @@ -402,10 +402,8 @@ def bp_payout_factory( thl_ledger_manager=thl_lm, product=product, amount=amount, - ext_ref_id=ext_ref_id, + ext_ref_id=ext_ref_id or uuid4().hex, created=created, - skip_wallet_balance_check=skip_wallet_balance_check, - skip_one_per_day_check=skip_one_per_day_check, ) return _inner diff --git a/tests/managers/thl/test_payout.py b/tests/managers/thl/test_payout.py index 31087b8..b79a209 100644 --- a/tests/managers/thl/test_payout.py +++ b/tests/managers/thl/test_payout.py @@ -5,10 +5,15 @@ from decimal import Decimal from random import choice as rand_choice, randint from typing import Optional from uuid import uuid4 +import io import pandas as pd import pytest +from generalresearch import pg_helper +from generalresearch.managers.thl.ledger_manager.thl_ledger import ThlLedgerManager +from generalresearch.models.thl.product import Product +from generalresearch.models.thl.user import User from generalresearch.currency import USDCent from generalresearch.managers.thl.ledger_manager.exceptions import ( LedgerTransactionConditionFailedError, @@ -16,7 +21,10 @@ from generalresearch.managers.thl.ledger_manager.exceptions import ( from generalresearch.managers.thl.payout import UserPayoutEventManager from generalresearch.models.thl.definitions import PayoutStatus from generalresearch.models.thl.ledger import LedgerEntry, Direction -from generalresearch.models.thl.payout import BusinessPayoutEvent +from generalresearch.models.thl.payout import ( + BusinessPayoutEvent, + BrokerageProductPayoutEvent, +) from generalresearch.models.thl.payout import UserPayoutEvent from generalresearch.models.thl.wallet import PayoutType from generalresearch.models.thl.ledger import LedgerAccount @@ -99,82 +107,43 @@ class TestPayout: delete_ledger_db, create_main_accounts, ): - delete_ledger_db() - create_main_accounts() - from generalresearch.models.thl.ledger import LedgerAccount - - thl_lm.get_account_or_create_bp_wallet(product=product) - brokerage_product_payout_event_manager.set_account_lookup_table(thl_lm=thl_lm) - - with pytest.raises(expected_exception=LedgerTransactionConditionFailedError): - # wallet balance failure - brokerage_product_payout_event_manager.create_bp_payout_event( - thl_ledger_manager=thl_lm, - product=product, - amount=USDCent(100), - skip_wallet_balance_check=False, - skip_one_per_day_check=False, - ) - - # (we don't have a special method for this) Put money in the BP's account - amount_cents = 100 - cash_account: LedgerAccount = thl_lm.get_account_cash() - bp_wallet: LedgerAccount = thl_lm.get_account_or_create_bp_wallet( - product=product - ) - - entries = [ - LedgerEntry( - direction=Direction.DEBIT, - account_uuid=cash_account.uuid, - amount=amount_cents, - ), - LedgerEntry( - direction=Direction.CREDIT, - account_uuid=bp_wallet.uuid, - amount=amount_cents, - ), - ] - - lm.create_tx(entries=entries) - assert 100 == lm.get_account_balance(account=bp_wallet) - - # Then run it again for $1.00 - brokerage_product_payout_event_manager.create_bp_payout_event( - thl_ledger_manager=thl_lm, - product=product, - amount=USDCent(100), - skip_wallet_balance_check=False, - skip_one_per_day_check=False, - ) - assert 0 == lm.get_account_balance(account=bp_wallet) - - # Run again should without balance check, should still fail due to day check - with pytest.raises(LedgerTransactionConditionFailedError): - brokerage_product_payout_event_manager.create_bp_payout_event( - thl_ledger_manager=thl_lm, - product=product, - amount=USDCent(100), - skip_wallet_balance_check=True, - skip_one_per_day_check=False, - ) + # create_bp_payout_event does not get called directly. We have tests + # for the ledger methods already + pass - # And then we can run again skip both checks - pe = brokerage_product_payout_event_manager.create_bp_payout_event( - thl_ledger_manager=thl_lm, - product=product, + @pytest.fixture + def pending_bp_pe( + self, + thl_web_rw, + product, + thl_lm: ThlLedgerManager, + brokerage_product_payout_event_manager, + utc_now + ) -> BrokerageProductPayoutEvent: + account = thl_lm.get_account_or_create_bp_wallet(product=product) + bp_pe = BrokerageProductPayoutEvent( + product_id=product.uuid, amount=USDCent(100), - skip_wallet_balance_check=True, - skip_one_per_day_check=True, - ) - assert -100 == lm.get_account_balance(account=bp_wallet) - - pe = brokerage_product_payout_event_manager.get_by_uuid(pe.uuid) - txs = lm.get_tx_filtered_by_metadata( - metadata_key="event_payout", metadata_value=pe.uuid - ) - - assert 1 == len(txs) + payout_type=PayoutType.ACH, + debit_account_uuid=account.uuid, + cashout_method_uuid=brokerage_product_payout_event_manager.CASHOUT_METHOD_UUID, + created=utc_now + ) + params = bp_pe.model_dump_postgres() + # This shouldn't exist. For testing only, so no supplier_payout + params['supplier_payout_id'] = None + thl_web_rw.execute_write(""" + INSERT INTO event_payout ( + uuid, debit_account_uuid, created, cashout_method_uuid, + amount, status, ext_ref_id, payout_type, order_data, + request_data, supplier_payout_id + ) VALUES ( + %(uuid)s, %(debit_account_uuid)s, %(created)s, %(cashout_method_uuid)s, + %(amount)s, %(status)s, %(ext_ref_id)s, %(payout_type)s, %(order_data)s, + %(request_data)s, %(supplier_payout_id)s + ); + """, params) + return bp_pe def test_create_bp_payout_quick_dupe( self, @@ -186,26 +155,22 @@ class TestPayout: lm, utc_now, create_main_accounts, + pending_bp_pe, ): thl_lm.get_account_or_create_bp_wallet(product=product) - brokerage_product_payout_event_manager.set_account_lookup_table(thl_lm=thl_lm) - brokerage_product_payout_event_manager.create_bp_payout_event( + brokerage_product_payout_event_manager.create_tx_bp_payout_from_payout_event( thl_ledger_manager=thl_lm, + bp_pe=pending_bp_pe, product=product, - amount=USDCent(100), - skip_wallet_balance_check=True, - skip_one_per_day_check=True, created=utc_now, ) with pytest.raises(ValueError) as cm: - brokerage_product_payout_event_manager.create_bp_payout_event( + brokerage_product_payout_event_manager.create_tx_bp_payout_from_payout_event( thl_ledger_manager=thl_lm, product=product, - amount=USDCent(100), - skip_wallet_balance_check=True, - skip_one_per_day_check=True, + bp_pe=pending_bp_pe, created=utc_now, ) assert "Payout event already exists!" in str(cm.value) @@ -290,44 +255,7 @@ class TestPayout: class TestPayoutEventManager: - def test_set_account_lookup_table( - self, payout_event_manager, thl_redis_config, thl_lm, delete_ledger_db - ): - delete_ledger_db() - rc = thl_redis_config.create_redis_client() - rc.delete("pem:account_to_product") - rc.delete("pem:product_to_account") - N = 5 - - for idx in range(N): - thl_lm.get_account_or_create_bp_wallet_by_uuid(product_uuid=uuid4().hex) - - res = rc.hgetall(name="pem:account_to_product") - assert len(res.items()) == 0 - - res = rc.hgetall(name="pem:product_to_account") - assert len(res.items()) == 0 - - payout_event_manager.set_account_lookup_table( - thl_lm=thl_lm, - ) - - res = rc.hgetall(name="pem:account_to_product") - assert len(res.items()) == N - - res = rc.hgetall(name="pem:product_to_account") - assert len(res.items()) == N - - thl_lm.get_account_or_create_bp_wallet_by_uuid(product_uuid=uuid4().hex) - payout_event_manager.set_account_lookup_table( - thl_lm=thl_lm, - ) - - res = rc.hgetall(name="pem:account_to_product") - assert len(res.items()) == N + 1 - - res = rc.hgetall(name="pem:product_to_account") - assert len(res.items()) == N + 1 + pass class TestBusinessPayoutEventManager: @@ -359,50 +287,22 @@ class TestBusinessPayoutEventManager: delete_ledger_db() create_main_accounts() - from generalresearch.models.thl.product import Product - p1: Product = product_factory(business=business) thl_lm.get_account_or_create_bp_wallet(product=p1) - business_payout_event_manager.set_account_lookup_table(thl_lm=thl_lm) ach_id1 = uuid4().hex ach_id2 = uuid4().hex - bp_payout_factory( - product=p1, - amount=USDCent(1), - ext_ref_id=None, - skip_wallet_balance_check=True, - skip_one_per_day_check=True, - ) + # ext_ref_id is required now + bp_payout_factory(product=p1,amount=USDCent(1),ext_ref_id="none") - bp_payout_factory( - product=p1, - amount=USDCent(1), - ext_ref_id=ach_id1, - skip_wallet_balance_check=True, - skip_one_per_day_check=True, - ) + bp_payout_factory(product=p1,amount=USDCent(1),ext_ref_id=ach_id1) + with pytest.raises(expected_exception=ValueError, match="Cannot create a BusinessPayoutEvent with an existing transaction_id"): + bp_payout_factory(product=p1,amount=USDCent(25),ext_ref_id=ach_id1) - bp_payout_factory( - product=p1, - amount=USDCent(25), - ext_ref_id=ach_id1, - skip_wallet_balance_check=True, - skip_one_per_day_check=True, - ) - - bp_payout_factory( - product=p1, - amount=USDCent(50), - ext_ref_id=ach_id2, - skip_wallet_balance_check=True, - skip_one_per_day_check=True, - ) + bp_payout_factory(product=p1,amount=USDCent(50),ext_ref_id=ach_id2) business.prebuild_payouts( - thl_pg_config=thl_web_rr, - thl_lm=thl_lm, bpem=business_payout_event_manager, ) @@ -410,11 +310,14 @@ class TestBusinessPayoutEventManager: assert business.payouts_total == sum([pe.amount for pe in business.payouts]) assert business.payouts[0].created > business.payouts[1].created assert len(business.payouts[0].bp_payouts) == 1 - assert len(business.payouts[1].bp_payouts) == 2 + + # Cannot pay out the same product twice in the same business payout + # assert len(business.payouts[1].bp_payouts) == 2 + assert len(business.payouts[1].bp_payouts) == 1 assert business.payouts[0].ext_ref_id == ach_id2 assert business.payouts[1].ext_ref_id == ach_id1 - assert business.payouts[2].ext_ref_id is None + assert business.payouts[2].ext_ref_id == "none" def test_update_ext_reference_ids( self, @@ -442,13 +345,9 @@ class TestBusinessPayoutEventManager: create_main_accounts() delete_df_collection(coll=ledger_collection) - from generalresearch.models.thl.product import Product - from generalresearch.models.thl.user import User - p1: Product = product_factory(business=business) u1: User = user_factory(product=p1) thl_lm.get_account_or_create_bp_wallet(product=p1) - business_payout_event_manager.set_account_lookup_table(thl_lm=thl_lm) # $250.00 to work with for idx in range(1, 10): @@ -461,12 +360,11 @@ class TestBusinessPayoutEventManager: ach_id1 = uuid4().hex ach_id2 = uuid4().hex - with pytest.raises(expected_exception=Warning) as cm: + with pytest.raises(expected_exception=AssertionError, match="No Business Payout found"): business_payout_event_manager.update_ext_reference_ids( new_value=ach_id2, current_value=ach_id1, ) - assert "No event_payouts found to UPDATE" in str(cm) # We must build the balance to issue ACH/Wire ledger_collection.initial_load(client=None, sync=True) @@ -487,6 +385,7 @@ class TestBusinessPayoutEventManager: transaction_id=ach_id1, ) assert isinstance(res, BusinessPayoutEvent) + assert business_payout_event_manager.get_by_ext_ref_id(ext_ref_id=ach_id1) # Okay, now that there is a payout_event, let's try to update the # ext_reference_id @@ -495,108 +394,10 @@ class TestBusinessPayoutEventManager: current_value=ach_id1, ) - res = business_payout_event_manager.filter_by(ext_ref_id=ach_id1) - assert len(res) == 0 - - res = business_payout_event_manager.filter_by(ext_ref_id=ach_id2) - assert len(res) == 1 - - def test_delete_failed_business_payout( - self, - brokerage_product_payout_event_manager, - business_payout_event_manager, - delete_ledger_db, - create_main_accounts, - thl_lm, - thl_web_rr, - product_factory, - bp_payout_factory, - currency, - delete_df_collection, - user_factory, - ledger_collection, - session_with_tx_factory, - pop_ledger_merge, - client_no_amm, - mnt_filepath, - lm, - product_manager, - start, - business, - ): - delete_ledger_db() - create_main_accounts() - delete_df_collection(coll=ledger_collection) - - from generalresearch.models.thl.product import Product - from generalresearch.models.thl.user import User - - p1: Product = product_factory(business=business) - u1: User = user_factory(product=p1) - thl_lm.get_account_or_create_bp_wallet(product=p1) - business_payout_event_manager.set_account_lookup_table(thl_lm=thl_lm) - - # $250.00 to work with - for idx in range(1, 10): - session_with_tx_factory( - user=u1, - wall_req_cpi=Decimal("25.00"), - started=start + timedelta(days=1, minutes=idx), - ) - - # We must build the balance to issue ACH/Wire - ledger_collection.initial_load(client=None, sync=True) - pop_ledger_merge.build(client=client_no_amm, ledger_coll=ledger_collection) - business.prebuild_balance( - thl_pg_config=thl_web_rr, - lm=lm, - ds=mnt_filepath, - client=client_no_amm, - pop_ledger=pop_ledger_merge, - ) - - ach_id1 = uuid4().hex - - res = business_payout_event_manager.create_from_ach_or_wire( - business=business, - amount=USDCent(100_01), - pm=product_manager, - thl_lm=thl_lm, - transaction_id=ach_id1, - ) - assert isinstance(res, BusinessPayoutEvent) - - # (1) Confirm the initial Event Payout, Tx, TxMeta, TxEntry all exist - event_payouts = business_payout_event_manager.filter_by(ext_ref_id=ach_id1) - event_payout_uuids = [i.uuid for i in event_payouts] - assert len(event_payout_uuids) == 1 - tags = [f"{currency.value}:bp_payout:{x}" for x in event_payout_uuids] - transactions = thl_lm.get_txs_by_tags(tags=tags) - assert len(transactions) == 1 - tx_metadata_ids = thl_lm.get_tx_metadata_ids_by_txs(transactions=transactions) - assert len(tx_metadata_ids) == 2 - tx_entries = thl_lm.get_tx_entries_by_txs(transactions=transactions) - assert len(tx_entries) == 2 - - # (2) Delete! - business_payout_event_manager.delete_failed_business_payout( - ext_ref_id=ach_id1, thl_lm=thl_lm - ) - - # (3) Confirm the initial Event Payout, Tx, TxMeta, TxEntry have - # all been deleted - res = business_payout_event_manager.filter_by(ext_ref_id=ach_id1) - assert len(res) == 0 + with pytest.raises(expected_exception=AssertionError, match="No Business Payout found"): + business_payout_event_manager.get_by_ext_ref_id(ext_ref_id=ach_id1) - # Note: b/c the event_payout shouldn't exist anymore, we are taking - # the tag strings and transactions from when they did.. - res = thl_lm.get_txs_by_tags(tags=tags) - assert len(res) == 0 - - tx_metadata_ids = thl_lm.get_tx_metadata_ids_by_txs(transactions=transactions) - assert len(tx_metadata_ids) == 0 - tx_entries = thl_lm.get_tx_entries_by_txs(transactions=transactions) - assert len(tx_entries) == 0 + assert business_payout_event_manager.get_by_ext_ref_id(ext_ref_id=ach_id2) def test_recoup_empty(self, business_payout_event_manager): res = {uuid4().hex: USDCent(0) for i in range(100)} @@ -708,7 +509,6 @@ class TestBusinessPayoutEventManager: assert int(res.deduction.sum()) == 0 def test_distribute_amount(self, business_payout_event_manager): - import io df = pd.read_csv( io.StringIO( @@ -748,7 +548,7 @@ class TestBusinessPayoutEventManager: lm, product_manager, ): - """Test having a Business with three products.. one that lost money + """Test having a Business with three products. One that lost money and two that gained money. Ensure that the Business balance reflects that to compensate for the Product in the negative and only assigns Brokerage Product payments from the 2 accounts that have @@ -759,9 +559,6 @@ class TestBusinessPayoutEventManager: create_main_accounts() delete_df_collection(coll=ledger_collection) - from generalresearch.models.thl.product import Product - from generalresearch.models.thl.user import User - p1: Product = product_factory(business=business) u1: User = user_factory(product=p1) thl_lm.get_account_or_create_bp_wallet(product=p1) @@ -776,7 +573,6 @@ class TestBusinessPayoutEventManager: wall_req_cpi=Decimal("5.00"), started=start + timedelta(days=6), ) - payout_event_manager.set_account_lookup_table(thl_lm=thl_lm) bp_payout_factory( product=u1.product, amount=USDCent(475), # 95% of $5.00 @@ -799,9 +595,131 @@ class TestBusinessPayoutEventManager: amount=USDCent(500), pm=product_manager, thl_lm=thl_lm, + transaction_id=uuid4().hex, ) assert "Must issue Supplier Payouts at least $100 minimum." in str(cm) + def test_create_from_ach_or_wire( + self, + product, + mnt_filepath, + thl_lm, + client_no_amm, + thl_redis_config, + payout_event_manager, + brokerage_product_payout_event_manager, + business_payout_event_manager, + delete_ledger_db, + create_main_accounts, + delete_df_collection, + ledger_collection, + business, + user_factory, + product_factory, + session_with_tx_factory, + pop_ledger_merge, + start, + bp_payout_factory, + adj_to_fail_with_tx_factory, + thl_web_rr, + lm, + product_manager, + rm_ledger_collection, + rm_pop_ledger_merge, + caplog, + ): + """Test having a Business with three products""" + # Now let's load it up and actually test some things + delete_ledger_db() + create_main_accounts() + delete_df_collection(coll=ledger_collection) + + p1: Product = product_factory(business=business) + p2: Product = product_factory(business=business) + p3: Product = product_factory(business=business) + u1: User = user_factory(product=p1) + u2: User = user_factory(product=p2) + u3: User = user_factory(product=p3) + thl_lm.get_account_or_create_bp_wallet(product=p1) + thl_lm.get_account_or_create_bp_wallet(product=p2) + thl_lm.get_account_or_create_bp_wallet(product=p3) + + ach_id1 = uuid4().hex + ach_id2 = uuid4().hex + + # Product 1: Complete $10 x 20 + for idx in range(20): + session_with_tx_factory( + user=u2, + wall_req_cpi=Decimal("10.00"), + started=start + timedelta(days=1, hours=2, minutes=1 + idx), + ) + + # Product 2: Complete $10 x 30 + for idx in range(30): + session_with_tx_factory( + user=u3, + wall_req_cpi=Decimal("10.00"), + started=start + timedelta(days=1, hours=3, minutes=1 + idx), + ) + + ledger_collection.initial_load(client=None, sync=True) + pop_ledger_merge.build(client=client_no_amm, ledger_coll=ledger_collection) + business.prebuild_balance( + thl_pg_config=thl_web_rr, + lm=lm, + ds=mnt_filepath, + client=client_no_amm, + pop_ledger=pop_ledger_merge, + ) + + bb = business.balance + assert bb.payout == 475_00 # $500 * .95% = $475 + assert bb.net == 475_00 + + bp1 = business_payout_event_manager.create_from_ach_or_wire( + business=business, + amount=USDCent(100_00), + pm=product_manager, + thl_lm=thl_lm, + created=start + timedelta(days=1, hours=5), + transaction_id=ach_id1, + ) + print(f"{bp1=}") + assert isinstance(bp1, BusinessPayoutEvent) + assert len(bp1.bp_payouts) == 2 + + bp2 = business_payout_event_manager.create_from_ach_or_wire( + business=business, + amount=USDCent(bb.available_balance), + pm=product_manager, + thl_lm=thl_lm, + created=start + timedelta(days=2, hours=5), + transaction_id=ach_id2, + ) + print(f"{bp2=}") + assert isinstance(bp2, BusinessPayoutEvent) + assert len(bp2.bp_payouts) == 2 + + with caplog.at_level(logging.WARNING): + business_payout_event_manager.resume_failed_business_payout( + ext_ref_id=ach_id1, thl_lm=thl_lm, pm=product_manager + ) + assert "Nothing to do!" in caplog.text + + # bpe = business_payout_event_manager.get_by_ext_ref_id(ext_ref_id=ach_id1) + # bp_pe = bpe.bp_payouts[0] + # thl_web_rr.execute_write( + # """ + # UPDATE event_payout + # SET status = %(status)s + # WHERE uuid = %(uuid)s""", + # {"uuid": bp_pe.uuid, "status": PayoutStatus.FAILED}, + # ) + + assert 1 == 0 + return None + def test_ach_payment( self, product, @@ -841,9 +759,6 @@ class TestBusinessPayoutEventManager: create_main_accounts() delete_df_collection(coll=ledger_collection) - from generalresearch.models.thl.product import Product - from generalresearch.models.thl.user import User - p1: Product = product_factory(business=business) p2: Product = product_factory(business=business) p3: Product = product_factory(business=business) @@ -863,7 +778,6 @@ class TestBusinessPayoutEventManager: wall_req_cpi=Decimal("5.00"), started=start + timedelta(days=1), ) - payout_event_manager.set_account_lookup_table(thl_lm=thl_lm) bp_payout_factory( product=u1.product, amount=USDCent(475), # 95% of $5.00 @@ -1047,9 +961,6 @@ class TestBusinessPayoutEventManager: create_main_accounts() delete_df_collection(coll=ledger_collection) - from generalresearch.models.thl.product import Product - from generalresearch.models.thl.user import User - p1: Product = product_factory(business=business) p2: Product = product_factory(business=business) p3: Product = product_factory(business=business) @@ -1081,8 +992,6 @@ class TestBusinessPayoutEventManager: pop_ledger=pop_ledger_merge, ) business.prebuild_payouts( - thl_pg_config=thl_web_rr, - thl_lm=thl_lm, bpem=business_payout_event_manager, ) @@ -1179,9 +1088,6 @@ class TestBusinessPayoutEventManager: create_main_accounts() delete_df_collection(coll=ledger_collection) - from generalresearch.models.thl.product import Product - from generalresearch.models.thl.user import User - p1: Product = product_factory(business=business) p2: Product = product_factory(business=business) p3: Product = product_factory(business=business) diff --git a/tests/models/thl/test_payout.py b/tests/models/thl/test_payout.py index 3a51328..7068a41 100644 --- a/tests/models/thl/test_payout.py +++ b/tests/models/thl/test_payout.py @@ -1,10 +1,115 @@ +from uuid import uuid4 + +import pytest +from pydantic import ValidationError + +from generalresearch.currency import USDCent +from generalresearch.models.gr import Team +from generalresearch.models.thl.payout import ( + BusinessPayoutEvent, + BrokerageProductPayoutEvent, +) +from generalresearch.models.thl.wallet import PayoutType + +from generalresearch.models.gr.business import Business, BusinessAddress, BusinessType + + class TestBusinessPayoutEvent: def test_validate(self): - from generalresearch.models.gr.business import Business - instance = Business.model_validate_json( - json_data='{"id":123,"uuid":"947f6ba5250d442b9a66cde9ee33605a","name":"Example » Demo","kind":"c","tax_number":null,"contact":null,"addresses":[],"teams":[{"id":53,"uuid":"8e4197dcaefe4f1f831a02b212e6b44a","name":"Example » Demo","memberships":null,"gr_users":null,"businesses":null,"products":null}],"products":[{"id":"fc23e741b5004581b30e6478363525df","id_int":1234,"name":"Example","enabled":true,"payments_enabled":true,"created":"2025-04-14T13:25:37.279403Z","team_id":"9e4197dcaefe4f1f831a02b212e6b44a","business_id":"857f6ba6160d442b9a66cde9ee33605a","tags":[],"commission_pct":"0.050000","redirect_url":"https://pam-api-us.reppublika.com/v2/public/4970ef00-0ef7-11f0-9962-05cb6323c84c/grl/status","harmonizer_domain":"https://talk.generalresearch.com/","sources_config":{"user_defined":[{"name":"w","active":false,"banned_countries":[],"allow_mobile_ip":true,"supplier_id":null,"allow_pii_only_buyers":false,"allow_unhashed_buyers":false,"withhold_profiling":false,"pass_unconditional_eligible_unknowns":true,"address":null,"allow_vpn":null,"distribute_harmonizer_active":null}]},"session_config":{"max_session_len":600,"max_session_hard_retry":5,"min_payout":"0.14"},"payout_config":{"payout_format":null,"payout_transformation":null},"user_wallet_config":{"enabled":false,"amt":false,"supported_payout_types":["CASH_IN_MAIL","PAYPAL","TANGO"],"min_cashout":null},"user_create_config":{"min_hourly_create_limit":0,"max_hourly_create_limit":null},"offerwall_config":{},"profiling_config":{"enabled":true,"grs_enabled":true,"n_questions":null,"max_questions":10,"avg_question_count":5.0,"task_injection_freq_mult":1.0,"non_us_mult":2.0,"hidden_questions_expiration_hours":168},"user_health_config":{"banned_countries":[],"allow_ban_iphist":true},"yield_man_config":{},"balance":null,"payouts_total_str":null,"payouts_total":null,"payouts":null,"user_wallet":{"enabled":false,"amt":false,"supported_payout_types":["CASH_IN_MAIL","PAYPAL","TANGO"],"min_cashout":null}}],"bank_accounts":[],"balance":{"product_balances":[{"product_id":"fc14e741b5004581b30e6478363414df","last_event":null,"bp_payment_credit":780251,"adjustment_credit":4678,"adjustment_debit":26446,"supplier_credit":0,"supplier_debit":451513,"user_bonus_credit":0,"user_bonus_debit":0,"issued_payment":0,"payout":780251,"payout_usd_str":"$7,802.51","adjustment":-21768,"expense":0,"net":758483,"payment":451513,"payment_usd_str":"$4,515.13","balance":306970,"retainer":76742,"retainer_usd_str":"$767.42","available_balance":230228,"available_balance_usd_str":"$2,302.28","recoup":0,"recoup_usd_str":"$0.00","adjustment_percent":0.027898714644390074}],"payout":780251,"payout_usd_str":"$7,802.51","adjustment":-21768,"expense":0,"net":758483,"net_usd_str":"$7,584.83","payment":451513,"payment_usd_str":"$4,515.13","balance":306970,"balance_usd_str":"$3,069.70","retainer":76742,"retainer_usd_str":"$767.42","available_balance":230228,"available_balance_usd_str":"$2,302.28","adjustment_percent":0.027898714644390074,"recoup":0,"recoup_usd_str":"$0.00"},"payouts_total_str":"$4,515.13","payouts_total":451513,"payouts":[{"bp_payouts":[{"uuid":"40cf2c3c341e4f9d985be4bca43e6116","debit_account_uuid":"3a058056da85493f9b7cdfe375aad0e0","cashout_method_uuid":"602113e330cf43ae85c07d94b5100291","created":"2025-08-02T09:18:20.433329Z","amount":345735,"status":"COMPLETE","ext_ref_id":null,"payout_type":"ACH","request_data":{},"order_data":null,"product_id":"fc14e741b5004581b30e6478363414df","method":"ACH","amount_usd":345735,"amount_usd_str":"$3,457.35"}],"amount":345735,"amount_usd_str":"$3,457.35","created":"2025-08-02T09:18:20.433329Z","line_items":1,"ext_ref_id":null},{"bp_payouts":[{"uuid":"63ce1787087248978919015c8fcd5ab9","debit_account_uuid":"3a058056da85493f9b7cdfe375aad0e0","cashout_method_uuid":"602113e330cf43ae85c07d94b5100291","created":"2025-06-10T22:16:18.765668Z","amount":105778,"status":"COMPLETE","ext_ref_id":"11175997868","payout_type":"ACH","request_data":{},"order_data":null,"product_id":"fc14e741b5004581b30e6478363414df","method":"ACH","amount_usd":105778,"amount_usd_str":"$1,057.78"}],"amount":105778,"amount_usd_str":"$1,057.78","created":"2025-06-10T22:16:18.765668Z","line_items":1,"ext_ref_id":"11175997868"}]}' + # Doesn't validate anymore + # instance = Business.model_validate_json( + # json_data='{"id":123,"uuid":"947f6ba5250d442b9a66cde9ee33605a","name":"Example » Demo","kind":"c","tax_number":null,"contact":null,"addresses":[],"teams":[{"id":53,"uuid":"8e4197dcaefe4f1f831a02b212e6b44a","name":"Example » Demo","memberships":null,"gr_users":null,"businesses":null,"products":null}],"products":[{"id":"fc23e741b5004581b30e6478363525df","id_int":1234,"name":"Example","enabled":true,"payments_enabled":true,"created":"2025-04-14T13:25:37.279403Z","team_id":"9e4197dcaefe4f1f831a02b212e6b44a","business_id":"857f6ba6160d442b9a66cde9ee33605a","tags":[],"commission_pct":"0.050000","redirect_url":"https://pam-api-us.reppublika.com/v2/public/4970ef00-0ef7-11f0-9962-05cb6323c84c/grl/status","harmonizer_domain":"https://talk.generalresearch.com/","sources_config":{"user_defined":[{"name":"w","active":false,"banned_countries":[],"allow_mobile_ip":true,"supplier_id":null,"allow_pii_only_buyers":false,"allow_unhashed_buyers":false,"withhold_profiling":false,"pass_unconditional_eligible_unknowns":true,"address":null,"allow_vpn":null,"distribute_harmonizer_active":null}]},"session_config":{"max_session_len":600,"max_session_hard_retry":5,"min_payout":"0.14"},"payout_config":{"payout_format":null,"payout_transformation":null},"user_wallet_config":{"enabled":false,"amt":false,"supported_payout_types":["CASH_IN_MAIL","PAYPAL","TANGO"],"min_cashout":null},"user_create_config":{"min_hourly_create_limit":0,"max_hourly_create_limit":null},"offerwall_config":{},"profiling_config":{"enabled":true,"grs_enabled":true,"n_questions":null,"max_questions":10,"avg_question_count":5.0,"task_injection_freq_mult":1.0,"non_us_mult":2.0,"hidden_questions_expiration_hours":168},"user_health_config":{"banned_countries":[],"allow_ban_iphist":true},"yield_man_config":{},"balance":null,"payouts_total_str":null,"payouts_total":null,"payouts":null,"user_wallet":{"enabled":false,"amt":false,"supported_payout_types":["CASH_IN_MAIL","PAYPAL","TANGO"],"min_cashout":null}}],"bank_accounts":[],"balance":{"product_balances":[{"product_id":"fc14e741b5004581b30e6478363414df","last_event":null,"bp_payment_credit":780251,"adjustment_credit":4678,"adjustment_debit":26446,"supplier_credit":0,"supplier_debit":451513,"user_bonus_credit":0,"user_bonus_debit":0,"issued_payment":0,"payout":780251,"payout_usd_str":"$7,802.51","adjustment":-21768,"expense":0,"net":758483,"payment":451513,"payment_usd_str":"$4,515.13","balance":306970,"retainer":76742,"retainer_usd_str":"$767.42","available_balance":230228,"available_balance_usd_str":"$2,302.28","recoup":0,"recoup_usd_str":"$0.00","adjustment_percent":0.027898714644390074}],"payout":780251,"payout_usd_str":"$7,802.51","adjustment":-21768,"expense":0,"net":758483,"net_usd_str":"$7,584.83","payment":451513,"payment_usd_str":"$4,515.13","balance":306970,"balance_usd_str":"$3,069.70","retainer":76742,"retainer_usd_str":"$767.42","available_balance":230228,"available_balance_usd_str":"$2,302.28","adjustment_percent":0.027898714644390074,"recoup":0,"recoup_usd_str":"$0.00"},"payouts_total_str":"$4,515.13","payouts_total":451513,"payouts":[{"bp_payouts":[{"uuid":"40cf2c3c341e4f9d985be4bca43e6116","debit_account_uuid":"3a058056da85493f9b7cdfe375aad0e0","cashout_method_uuid":"602113e330cf43ae85c07d94b5100291","created":"2025-08-02T09:18:20.433329Z","amount":345735,"status":"COMPLETE","ext_ref_id":null,"payout_type":"ACH","request_data":{},"order_data":null,"product_id":"fc14e741b5004581b30e6478363414df","method":"ACH","amount_usd":345735,"amount_usd_str":"$3,457.35"}],"amount":345735,"amount_usd_str":"$3,457.35","created":"2025-08-02T09:18:20.433329Z","line_items":1,"ext_ref_id":null},{"bp_payouts":[{"uuid":"63ce1787087248978919015c8fcd5ab9","debit_account_uuid":"3a058056da85493f9b7cdfe375aad0e0","cashout_method_uuid":"602113e330cf43ae85c07d94b5100291","created":"2025-06-10T22:16:18.765668Z","amount":105778,"status":"COMPLETE","ext_ref_id":"11175997868","payout_type":"ACH","request_data":{},"order_data":null,"product_id":"fc14e741b5004581b30e6478363414df","method":"ACH","amount_usd":105778,"amount_usd_str":"$1,057.78"}],"amount":105778,"amount_usd_str":"$1,057.78","created":"2025-06-10T22:16:18.765668Z","line_items":1,"ext_ref_id":"11175997868"}]}' + # ) + # assert isinstance(instance, Business) + + # Make manually + b = Business( + id=123, + uuid=uuid4().hex, + name="Example", + addresses=[ + BusinessAddress( + uuid=uuid4().hex, + city="xxx", + line_1="xxx", + state="fl", + business_id=123, + ) + ], + kind=BusinessType.COMPANY, + teams=[Team(uuid=uuid4().hex, name="Example » Demo")], + products=[], + bank_accounts=[], + ) + ext_ref_id = uuid4().hex + bpe = BusinessPayoutEvent( + business_id=uuid4().hex, + amount=USDCent(100_00), + payout_type=PayoutType.ACH, + ext_ref_id=ext_ref_id, ) + bpe.bp_payouts = [ + BrokerageProductPayoutEvent( + product_id=uuid4().hex, + payout_type=PayoutType.ACH, + amount=USDCent(47_00), + cashout_method_uuid=uuid4().hex, + debit_account_uuid=uuid4().hex, + ext_ref_id=ext_ref_id, + ), + BrokerageProductPayoutEvent( + product_id=uuid4().hex, + payout_type=PayoutType.ACH, + amount=USDCent(53_00), + cashout_method_uuid=uuid4().hex, + debit_account_uuid=uuid4().hex, + ext_ref_id=ext_ref_id, + ), + ] + + # Test validations (amount sum) + with pytest.raises( + ValidationError, + match="BusinessPayoutEvent.amount must equal the sum of bp_payouts amounts", + ): + bpe.bp_payouts = [ + BrokerageProductPayoutEvent( + product_id=uuid4().hex, + payout_type=PayoutType.ACH, + amount=USDCent(47_00), + cashout_method_uuid=uuid4().hex, + debit_account_uuid=uuid4().hex, + ext_ref_id=ext_ref_id, + ) + ] + + with pytest.raises( + ValidationError, + match="All BrokerageProductPayoutEvent.ext_ref_id values must equal", + ): + bpe.bp_payouts = [ + BrokerageProductPayoutEvent( + product_id=uuid4().hex, + payout_type=PayoutType.ACH, + amount=USDCent(100_00), + cashout_method_uuid=uuid4().hex, + debit_account_uuid=uuid4().hex, + ext_ref_id="a different value", + ) + ] - assert isinstance(instance, Business) + with pytest.raises( + ValidationError, match="All BrokerageProductPayoutEvent.payout_type values" + ): + bpe.bp_payouts = [ + BrokerageProductPayoutEvent( + product_id=uuid4().hex, + payout_type=PayoutType.PAYPAL, + amount=USDCent(100_00), + cashout_method_uuid=uuid4().hex, + debit_account_uuid=uuid4().hex, + ext_ref_id=ext_ref_id, + ) + ] |
