from __future__ import annotations import logging from collections.abc import Callable from datetime import UTC, datetime, timedelta from decimal import Decimal from random import randint from typing import TYPE_CHECKING from uuid import uuid4 import pytest import redis from pydantic import RedisDsn from redis.lock import Lock from generalresearch.currency import USDCent from generalresearch.managers.base import Permission from generalresearch.managers.thl.ledger_manager.exceptions import ( LedgerTransactionConditionFailedError, LedgerTransactionCreateError, LedgerTransactionFlagAlreadyExistsError, LedgerTransactionReleaseLockError, ) from generalresearch.managers.thl.ledger_manager.ledger import LedgerTransaction from generalresearch.managers.thl.ledger_manager.thl_ledger import ThlLedgerManager from generalresearch.models.definitions import Source from generalresearch.models.thl.definitions import PayoutStatus from generalresearch.models.thl.ledger import Direction, TransactionType from generalresearch.models.thl.session import ( Session, Status, StatusCode1, Wall, ) from generalresearch.redis_helper import RedisConfig if TYPE_CHECKING: from generalresearch.currency import LedgerCurrency from generalresearch.managers.thl.payout import ( BrokerageProductPayoutEventManager, BusinessPayoutEventManager, ) from generalresearch.models.thl.product import Product from generalresearch.models.thl.user import User from generalresearch.pg_helper import PostgresConfig def broken_acquire(self, *args, **kwargs): raise redis.exceptions.TimeoutError("Simulated timeout during acquire") def broken_release(self, *args, **kwargs): raise redis.exceptions.TimeoutError("Simulated timeout during release") class TestThlLedgerManagerBPPayout: @pytest.fixture(autouse=True) def setup(self, create_main_accounts): create_main_accounts() def test_create_tx_with_bp_payment( self, user_factory: Callable[..., User], product_user_wallet_no: Product, caplog, thl_ledger_manager: ThlLedgerManager, ): now = datetime.now(UTC) - timedelta(hours=1) user: User = user_factory(product=product_user_wallet_no) wall1 = Wall( user_id=user.user_id, source=Source.DYNATA, req_survey_id="xxx", req_cpi=Decimal("6.00"), session_id=1, status=Status.COMPLETE, status_code_1=StatusCode1.COMPLETE, started=now, finished=now + timedelta(seconds=1), ) tx = thl_ledger_manager.create_tx_task_complete( wall=wall1, user=user, created=wall1.started ) assert isinstance(tx, LedgerTransaction) session = Session(started=wall1.started, user=user, wall_events=[wall1]) status, status_code_1 = session.determine_session_status() _, _, bp_pay, user_pay = session.determine_payments() session.update( status=status, status_code_1=status_code_1, finished=now + timedelta(minutes=10), payout=bp_pay, user_payout=user_pay, ) thl_ledger_manager.create_tx_bp_payment(session=session, created=wall1.started) lock_key = f"test:bp_payout:{user.product.id}" flag_name = f"{thl_ledger_manager.cache_prefix}:transaction_flag:{lock_key}" thl_ledger_manager.redis_client.delete(flag_name) payoutevent_uuid = uuid4().hex thl_ledger_manager.create_tx_bp_payout( product=user.product, amount=USDCent(200), created=now, payoutevent_uuid=payoutevent_uuid, ) payoutevent_uuid = uuid4().hex thl_ledger_manager.create_tx_bp_payout( product=user.product, amount=USDCent(200), created=now + timedelta(minutes=2), skip_one_per_day_check=True, payoutevent_uuid=payoutevent_uuid, ) cash = thl_ledger_manager.get_account_cash() bp_wallet_account = thl_ledger_manager.get_account_or_create_bp_wallet( user.product ) assert 170 == thl_ledger_manager.get_account_balance(bp_wallet_account) assert 200 == thl_ledger_manager.get_account_balance(cash) with pytest.raises(expected_exception=LedgerTransactionFlagAlreadyExistsError): thl_ledger_manager.create_tx_bp_payout( user.product, amount=USDCent(200), created=now + timedelta(minutes=2), skip_one_per_day_check=False, skip_wallet_balance_check=False, payoutevent_uuid=payoutevent_uuid, ) payoutevent_uuid = uuid4().hex with ( caplog.at_level(logging.INFO), pytest.raises(LedgerTransactionConditionFailedError), ): thl_ledger_manager.create_tx_bp_payout( user.product, amount=USDCent(10_000), created=now + timedelta(minutes=2), skip_one_per_day_check=True, skip_wallet_balance_check=False, payoutevent_uuid=payoutevent_uuid, ) assert "failed condition check balance:" in caplog.text thl_ledger_manager.create_tx_bp_payout( product=user.product, amount=USDCent(10_00), created=now + timedelta(minutes=2), skip_one_per_day_check=True, skip_wallet_balance_check=True, payoutevent_uuid=payoutevent_uuid, ) assert 170 - 1000 == thl_ledger_manager.get_account_balance(bp_wallet_account) def test_create_tx( self, product: Product, caplog, thl_ledger_manager: ThlLedgerManager, currency: LedgerCurrency, ): rand_amount: USDCent = USDCent(randint(100, 1_000)) payoutevent_uuid = uuid4().hex # Create a BP Payout for a Product without any activity. By issuing, # the skip_* checks, we should be able to force it to work, and will # then ultimately result in a negative balance tx = thl_ledger_manager.create_tx_bp_payout( product=product, amount=rand_amount, payoutevent_uuid=payoutevent_uuid, created=datetime.now(tz=UTC), skip_wallet_balance_check=True, skip_one_per_day_check=True, skip_flag_check=True, ) # Check the basic attributes assert isinstance(tx, LedgerTransaction) assert tx.ext_description == "BP Payout" assert ( tx.tag == f"{currency.value}:{TransactionType.BP_PAYOUT.value}:{payoutevent_uuid}" ) assert tx.entries[0].amount == rand_amount assert tx.entries[1].amount == rand_amount # Check the Product's balance, it should be negative the amount that was # paid out. That's because the Product earned nothing.. and then was # sent something. balance = thl_ledger_manager.get_account_balance( account=thl_ledger_manager.get_account_or_create_bp_wallet(product=product) ) assert balance == int(rand_amount) * -1 # Test some basic assertions with ( caplog.at_level(logging.INFO), pytest.raises(expected_exception=LedgerTransactionConditionFailedError), ): thl_ledger_manager.create_tx_bp_payout( product=product, amount=rand_amount, payoutevent_uuid=uuid4().hex, created=datetime.now(tz=UTC), skip_wallet_balance_check=False, skip_one_per_day_check=False, skip_flag_check=False, ) assert "failed condition check >1 tx per day" in caplog.text def test_create_tx_redis_failure( self, product: Product, thl_web_rw: PostgresConfig, thl_ledger_manager: ThlLedgerManager, ): rand_amount: USDCent = USDCent(randint(100, 1_000)) payoutevent_uuid = uuid4().hex now = datetime.now(tz=UTC) thl_ledger_manager.create_tx_plug_bp_wallet( product, rand_amount, now, direction=Direction.CREDIT ) # Non routable IP address. Redis will fail thl_lm_redis_0 = ThlLedgerManager( pg_config=thl_web_rw, permissions=[ Permission.CREATE, Permission.READ, Permission.UPDATE, Permission.DELETE, ], testing=True, redis_config=RedisConfig( dsn=RedisDsn("redis://10.255.255.1:6379"), socket_connect_timeout=0.1, ), ) with pytest.raises(expected_exception=Exception) as e: thl_lm_redis_0.create_tx_bp_payout( product=product, amount=rand_amount, payoutevent_uuid=payoutevent_uuid, created=datetime.now(tz=UTC), ) assert e.type is redis.exceptions.TimeoutError # No txs were created bp_wallet_account = thl_ledger_manager.get_account_or_create_bp_wallet( product=product ) txs = thl_ledger_manager.get_tx_filtered_by_account( account_uuid=bp_wallet_account.uuid ) txs = [tx for tx in txs if tx.metadata["tx_type"] != "plug"] assert len(txs) == 0 def test_create_tx_multiple_per_day( self, product: Product, thl_ledger_manager: ThlLedgerManager ): rand_amount: USDCent = USDCent(randint(100, 1_000)) payoutevent_uuid = uuid4().hex now = datetime.now(tz=UTC) thl_ledger_manager.create_tx_plug_bp_wallet( product, rand_amount * USDCent(2), now, direction=Direction.CREDIT ) thl_ledger_manager.create_tx_bp_payout( product=product, amount=rand_amount, payoutevent_uuid=payoutevent_uuid, created=datetime.now(tz=UTC), ) # Try to create another # Will fail b/c it has the same payout event uuid with pytest.raises(expected_exception=Exception) as e: thl_ledger_manager.create_tx_bp_payout( product=product, amount=rand_amount, payoutevent_uuid=payoutevent_uuid, created=datetime.now(tz=UTC), ) assert e.type is LedgerTransactionFlagAlreadyExistsError # Try to create another # Will fail due to multiple per day payoutevent_uuid2 = uuid4().hex with pytest.raises(expected_exception=Exception) as e: thl_ledger_manager.create_tx_bp_payout( product=product, amount=rand_amount, payoutevent_uuid=payoutevent_uuid2, created=datetime.now(tz=UTC), ) assert e.type is LedgerTransactionConditionFailedError assert str(e.value) == ">1 tx per day" # Make it run by skipping one per day check thl_ledger_manager.create_tx_bp_payout( product=product, amount=rand_amount, payoutevent_uuid=payoutevent_uuid2, created=datetime.now(tz=UTC), skip_one_per_day_check=True, ) def test_create_tx_redis_lock_release_error( self, product: Product, thl_ledger_manager: ThlLedgerManager, monkeypatch: pytest.MonkeyPatch, ): rand_amount: USDCent = USDCent(randint(100, 1_000)) payoutevent_uuid = uuid4().hex now = datetime.now(tz=UTC) bp_wallet_account = thl_ledger_manager.get_account_or_create_bp_wallet( product=product ) thl_ledger_manager.create_tx_plug_bp_wallet( product, rand_amount * USDCent(2), now, direction=Direction.CREDIT ) # Create TX will fail on lock enter, no tx will actually get created with monkeypatch.context() as m: m.setattr(Lock, "acquire", broken_acquire) with pytest.raises(expected_exception=Exception) as e: thl_ledger_manager.create_tx_bp_payout( product=product, amount=rand_amount, payoutevent_uuid=payoutevent_uuid, created=datetime.now(tz=UTC), ) assert e.type is LedgerTransactionCreateError assert str(e.value) == "Redis error: Simulated timeout during acquire" txs = thl_ledger_manager.get_tx_filtered_by_account( account_uuid=bp_wallet_account.uuid ) txs = [tx for tx in txs if tx.metadata["tx_type"] != "plug"] assert len(txs) == 0 # Create TX will fail on lock exit, after the tx was created! with monkeypatch.context() as m: m.setattr(Lock, "release", broken_release) with pytest.raises(LedgerTransactionReleaseLockError) as e: thl_ledger_manager.create_tx_bp_payout( product=product, amount=rand_amount, payoutevent_uuid=payoutevent_uuid, created=datetime.now(tz=UTC), ) assert str(e.value) == "Redis error: Simulated timeout during release" # Transaction was still created! txs = thl_ledger_manager.get_tx_filtered_by_account( account_uuid=bp_wallet_account.uuid ) txs = [tx for tx in txs if tx.metadata["tx_type"] != "plug"] assert len(txs) == 1 class TestPayoutEventManagerBPPayout: @pytest.fixture(autouse=True) def setup(self, create_main_accounts): create_main_accounts() def test_create( self, product: Product, thl_ledger_manager: ThlLedgerManager, brokerage_product_payout_event_manager: BrokerageProductPayoutEventManager, business_payout_event_manager: BusinessPayoutEventManager, ): rand_amount: USDCent = USDCent(randint(100, 1_000)) now = datetime.now(tz=UTC) bp_wallet_account = thl_ledger_manager.get_account_or_create_bp_wallet( product=product ) assert thl_ledger_manager.get_account_balance(bp_wallet_account) == 0 thl_ledger_manager.create_tx_plug_bp_wallet( product, rand_amount, now, direction=Direction.CREDIT ) assert thl_ledger_manager.get_account_balance(bp_wallet_account) == rand_amount bpe = business_payout_event_manager.create_bp_payout_event( thl_ledger_manager=thl_ledger_manager, product=product, created=now, amount=rand_amount, ext_ref_id=uuid4().hex, ) bp_pe = bpe.bp_payouts[0] assert brokerage_product_payout_event_manager.check_for_ledger_tx( thl_ledger_manager=thl_ledger_manager, payout_event=bp_pe, ) assert thl_ledger_manager.get_account_balance(bp_wallet_account) == 0 def test_create_with_redis_error( self, product: Product, caplog, thl_ledger_manager: ThlLedgerManager, brokerage_product_payout_event_manager: BrokerageProductPayoutEventManager, business_payout_event_manager: BusinessPayoutEventManager, monkeypatch: pytest.MonkeyPatch, ): caplog.set_level("WARNING") ext_ref_id = uuid4().hex rand_amount: USDCent = USDCent(randint(100, 1_000)) now = datetime.now(tz=UTC) bp_wallet_account = thl_ledger_manager.get_account_or_create_bp_wallet( product=product ) assert thl_ledger_manager.get_account_balance(bp_wallet_account) == 0 thl_ledger_manager.create_tx_plug_bp_wallet( product=product, amount=rand_amount, created=now, direction=Direction.CREDIT ) assert thl_ledger_manager.get_account_balance(bp_wallet_account) == rand_amount # Will fail on lock enter, no tx will actually get created with monkeypatch.context() as m: m.setattr(Lock, "acquire", broken_acquire) with pytest.raises(LedgerTransactionCreateError) as e: business_payout_event_manager.create_bp_payout_event( thl_ledger_manager=thl_ledger_manager, product=product, created=now, amount=rand_amount, ext_ref_id=ext_ref_id, ) assert str(e.value) == "Redis error: Simulated timeout during acquire" txs = thl_ledger_manager.get_tx_filtered_by_account( account_uuid=bp_wallet_account.uuid ) txs = [tx for tx in txs if tx.metadata["tx_type"] != "plug"] # One payout event is created, status is failed, and no ledger txs exist assert len(txs) == 0 pes = ( brokerage_product_payout_event_manager.get_bp_bp_payout_events_for_products( product_uuids=[product.id] ) ) assert len(pes) == 1 assert pes[0].status == PayoutStatus.FAILED pe = pes[0] # Try to fix the failed payout, by trying ledger tx again brokerage_product_payout_event_manager.retry_create_bp_payout_event_tx( product=product, thl_ledger_manager=thl_ledger_manager, bp_pe=pe, ) txs = thl_ledger_manager.get_tx_filtered_by_account( account_uuid=bp_wallet_account.uuid ) txs = [tx for tx in txs if tx.metadata["tx_type"] != "plug"] assert len(txs) == 1 assert thl_ledger_manager.get_account_balance(bp_wallet_account) == 0 # And then try to run it again, it'll fail because a payout event with the same info exists with pytest.raises(expected_exception=ValueError) as e: pe = business_payout_event_manager.create_bp_payout_event( thl_ledger_manager=thl_ledger_manager, product=product, created=now, amount=rand_amount, ext_ref_id=ext_ref_id, ) assert ( "Cannot create a BusinessPayoutEvent with an existing transaction_id" in str(e.value) ) # We wouldn't do this in practice, because this is paying out the BP again, but # we can if want to. # Change the ext_ref_id so it'll create a new payout event pe = business_payout_event_manager.create_bp_payout_event( thl_ledger_manager=thl_ledger_manager, product=product, created=now, amount=rand_amount, ext_ref_id=uuid4().hex, ) txs = thl_ledger_manager.get_tx_filtered_by_account( account_uuid=bp_wallet_account.uuid ) txs = [tx for tx in txs if tx.metadata["tx_type"] != "plug"] assert len(txs) == 2 # since they were paid twice assert ( thl_ledger_manager.get_account_balance(bp_wallet_account) == 0 - rand_amount ) def test_create_with_redis_error_release( self, product: Product, thl_ledger_manager: ThlLedgerManager, business_payout_event_manager: BusinessPayoutEventManager, brokerage_product_payout_event_manager: BrokerageProductPayoutEventManager, monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture, ): caplog.set_level("WARNING") rand_amount: USDCent = USDCent(randint(100, 1_000)) now = datetime.now(tz=UTC) bp_wallet_account = thl_ledger_manager.get_account_or_create_bp_wallet( product=product ) assert thl_ledger_manager.get_account_balance(bp_wallet_account) == 0 thl_ledger_manager.create_tx_plug_bp_wallet( product, rand_amount, now, direction=Direction.CREDIT ) assert thl_ledger_manager.get_account_balance(bp_wallet_account) == rand_amount # Will fail on lock exit, after the tx was created! # But it'll see that the tx was created and so everything will be fine caplog.clear() with monkeypatch.context() as m, caplog.at_level("WARNING"): m.setattr(Lock, "release", broken_release) business_payout_event_manager.create_bp_payout_event( thl_ledger_manager=thl_ledger_manager, product=product, created=now, amount=rand_amount, ext_ref_id=uuid4().hex, ) assert "Redis error: Simulated timeout during release" in caplog.messages txs = thl_ledger_manager.get_tx_filtered_by_account( account_uuid=bp_wallet_account.uuid ) txs = [tx for tx in txs if tx.metadata["tx_type"] != "plug"] assert len(txs) == 1 pes = ( brokerage_product_payout_event_manager.get_bp_bp_payout_events_for_products( product_uuids=[product.uuid] ) ) assert len(pes) == 1 assert pes[0].status == PayoutStatus.COMPLETE