"""Abacus statement period logic methods.""" from datetime import datetime, timezone from http import HTTPStatus from typing import List from collaborator.constants import error from collaborator.constants.abacus_statement_period import PAYONEER_PAYMENT_DESCRIPTION from collaborator.logic.dp_payment import get_dp_payments from collaborator.models.ows import ows_payee from collaborator.models.rds.abacus_statement_period_persister import ( AbacusStatementPeriodPersister, ) from collaborator.models.rds.dp_payment import PayoneerPaymentStatus from collaborator.models.rds.dp_payment_persister import ( DpPaymentPersister, ) from collaborator.models.snowflake.abacus_statement_period_persister import ( AbacusStatementPeriodPersister as AbacusStatementPeriodPersisterSF, ) from collaborator.schemas.dp_payment import DpPaymentSchema, PayoneerPayoutSchema from collaborator.utils.error import OwsError from collaborator.utils.typing import User def abacus_statement_period_dataloader(abacus_statement_period_ids: List[int]): """Dataloader logic for Abacus statement periods.""" periods = AbacusStatementPeriodPersister.get_abacus_statement_periods( abacus_statement_period_ids ) periods_by_id = { period.abacus_statement_period_id: dict(period._mapping) for period in periods } return periods_by_id def direct_payment_balances(abacus_statement_period_id: int): """Logic for getting direct payment balances for an abacus statement period.""" results = AbacusStatementPeriodPersisterSF.get_direct_payment_balances( abacus_statement_period_id, ) total_amount: float = 0 balances = [] for result in results: total_amount += float(result.balance_after_tax) balances.append( { "agreement_type": result.agreement_type, "account_name": result.account_name, "account_id": result.account_id, "currency_code": result.currency_code, "current_balance": float(result.balance_after_tax), "current_statement_period_name": result.current_statement_period_name, "payoneer_program_id": result.payoneer_program_id, "payoneer_program_name": result.payoneer_program_name, "payoneer_client_reference_id": result.payoneer_client_reference_id, "balance_after_tax": float(result.balance_after_tax), "collaborator_name": result.collaborator_name, "collaborator_id": result.collaborator_id, } ) return { "currency_agnostic_total_amount": total_amount, "total_count": len(balances), "balances": balances, } def replace_dp_payments( abacus_statement_period_id: int, user: User ) -> tuple[list[DpPaymentSchema], float]: """Logic for replacing DP Payment data for an abacus statement period.""" payments, total_amount = get_dp_payments( abacus_statement_period_id=abacus_statement_period_id, ) if len(payments) and all(p.approved_date is not None for p in payments): return payments, total_amount results = AbacusStatementPeriodPersisterSF.get_direct_payment_balances( abacus_statement_period_id, ) total_amount = 0 dp_payments_to_create = [] for result in results: total_amount += float(result.balance_after_tax) payoneer_payment_id = ( f"split:{abacus_statement_period_id}:{result.collaborator_id}" ) dp_payments_to_create.append( { "abacus_statement_period_id": abacus_statement_period_id, "abacus_statement_period_name": result.current_statement_period_name, "collaborator_id": result.collaborator_id, "collaborator_name": result.collaborator_name, "amount": result.balance_after_tax, "currency": result.currency_code, "payoneer_program_id": result.payoneer_program_id, "payoneer_program_name": result.payoneer_program_name, "payoneer_client_reference_id": result.payoneer_client_reference_id, "created_date": datetime.now(tz=timezone.utc), "account_id": result.account_id, "account_name": result.account_name, "payee_id": result.payee_id, "agreement_type": result.agreement_type, "payoneer_payment_id": payoneer_payment_id, "created_by": user.id, } ) created_payments = DpPaymentPersister.replace_by_abacus_statement_period_id( abacus_statement_period_id=abacus_statement_period_id, payments=dp_payments_to_create, ) return [ DpPaymentSchema.parse(payment) for payment in created_payments ], total_amount def create_dp_payment_transactions( abacus_statement_period_id: int, user: User, ): """Logic for creating DP payment transactions for an abacus statement period.""" AbacusStatementPeriodPersister.create_dp_payment_transactions( abacus_statement_period_id=abacus_statement_period_id, user=user, ) def approve_dp_payments( abacus_statement_period_id: int, user: User, ) -> list[DpPaymentSchema]: """Logic for approving DP payments for an abacus statement period.""" approved_payments = DpPaymentPersister.approve_by_abacus_statement_period_id( abacus_statement_period_id=abacus_statement_period_id, user=user, ) return [DpPaymentSchema.parse(payment) for payment in approved_payments] def submit_dp_payments(abacus_statement_period_id: int, user: User) -> dict: """Submit approved DP payments to ows-payee as mass payout batches. Fetches all DP payments for the given statement period, filters to those that are approved but not yet submitted (approved_date set, payoneer_payment_status is NULL), groups them by `payoneer_program_id`, and calls `ows_payee.create_mass_payouts` once per program group. Args: abacus_statement_period_id (int): Abacus statement period ID. Returns: dict: `{'submitted_count': N}` Raises: OwsError: 422 when no submittable payments exist for the period. """ rows = DpPaymentPersister.get_by_filters( abacus_statement_period_id=abacus_statement_period_id, ) pending_submission = [ dp for dp, _ in rows if dp.approved_date is not None and dp.payoneer_payment_status is None ] if not pending_submission: raise OwsError( code=error.ERROR_CODE_DP_PAYMENT_NOT_APPROVED, message=error.ERROR_MESSAGE_DP_PAYMENT_NOT_APPROVED, status=HTTPStatus.UNPROCESSABLE_ENTITY, ) payment_group_by_program_id: dict[int, list[dict]] = {} payment_ids_by_program_id: dict[int, list[int]] = {} for dp in pending_submission: program_id = int(dp.payoneer_program_id) payout = PayoneerPayoutSchema( client_reference_id=dp.payoneer_payment_id, payee_id=dp.payoneer_client_reference_id, amount=float(dp.amount), currency=dp.currency, description=PAYONEER_PAYMENT_DESCRIPTION, ) payment_group_by_program_id.setdefault(program_id, []).append(payout.dump()) payment_ids_by_program_id.setdefault(program_id, []).append(dp.dp_payment_id) for program_id, payout_list in payment_group_by_program_id.items(): ows_payee.create_mass_payouts(program_id, payout_list) DpPaymentPersister.update_payoneer_status_bulk( payment_ids_by_program_id[program_id], PayoneerPaymentStatus.INIT, event_type=None, reason=None, user=user, ) return { "submitted_count": len(pending_submission), }