diff --git a/__init__.py b/__init__.py index 521176b..a60bb29 100644 --- a/__init__.py +++ b/__init__.py @@ -23,7 +23,7 @@ def payroll_stop(): for task in scheduled_tasks: try: task.cancel() - except Exception as ex: + except Exception as ex: # shutdown must not raise logger.warning(ex) diff --git a/crud.py b/crud.py index 71c98a6..2c76717 100644 --- a/crud.py +++ b/crud.py @@ -12,7 +12,7 @@ from datetime import datetime, timezone from lnbits.db import Database from lnbits.helpers import urlsafe_short_hash -from .models import Contract, ContractStatus, CreateContract +from .models import Contract, ContractStatus, CreateContract, Payout, PayoutStatus db = Database("ext_payroll") @@ -87,3 +87,54 @@ async def delete_contract(contract_id: str) -> None: await db.execute( "DELETE FROM payroll.contracts WHERE id = :id", {"id": contract_id} ) + + +# --------------------------------------------------------------------------- +# Payout ledger +# --------------------------------------------------------------------------- + + +async def create_payout(payout: Payout) -> Payout: + await db.insert("payroll.payouts", payout) + return payout + + +async def get_payouts( + contract_id: str | None = None, + status: PayoutStatus | None = None, + limit: int = 200, +) -> list[Payout]: + where = [] + values: dict = {} + if contract_id: + where.append("contract_id = :cid") + values["cid"] = contract_id + if status: + where.append("status = :status") + values["status"] = status.value + clause = f"WHERE {' AND '.join(where)}" if where else "" + return await db.fetchall( + f"SELECT * FROM payroll.payouts {clause} " + "ORDER BY created_at DESC LIMIT :limit", + {**values, "limit": limit}, + Payout, + ) + + +async def count_period_failures(contract_id: str, period_index: int) -> int: + """How many times this exact period has already failed. + + Drives the bounded-retry cap. Counted from the ledger rather than from a + counter on the contract so the number survives a restart and stays + auditable — the rows that produced it are right there. + """ + row = await db.fetchone( + "SELECT COUNT(*) AS n FROM payroll.payouts " + "WHERE contract_id = :cid AND period_index = :idx AND status = :status", + { + "cid": contract_id, + "idx": period_index, + "status": PayoutStatus.failed.value, + }, + ) + return int(row["n"]) if row else 0 diff --git a/docs/operations.md b/docs/operations.md index 93e7c74..a8e06da 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -56,6 +56,32 @@ payroll: contract a1b2c3d4 period 3 failed: insufficient balance in source wallet: 12000 sat available, 80000 sat required ``` +Retries are bounded. After `MAX_PERIOD_ATTEMPTS` (5) failures on the *same* +period, the contract is **paused** and the operator has to act. Pausing +rather than abandoning the period is the point: a payday that cannot be +funded is a fact somebody needs to see, and silently dropping it is the one +outcome payroll must never produce. Because pausing does not advance the +position, resuming after topping up the source wallet retries that same +payday. + +## The payout ledger + +Every attempt — paid, skipped *and* failed — is written to +`payroll.payouts`. `GET /payroll/api/v1/payouts` returns them newest first, +filterable by `contract_id` and `status`. "Why did nobody get paid on the +1st" is the question the ledger exists to answer, which is why failures are +in it rather than only in the log. + +Rows are self-describing: each one copies the terms in force at the time +(amount, currency, both wallet ids) rather than pointing at the contract, +so a payout still reads correctly after its contract is edited or deleted. +There is deliberately no foreign key to `contracts` — deleting a contract +must not take its history with it. + +The attempt counter is derived by counting a period's prior failures in the +ledger, not from a column on the contract, so it survives a restart and the +rows that produced the number are right there to audit. + ## Back-dated start dates Creating a contract with a start date in the past is normally an diff --git a/migrations.py b/migrations.py index 2967f47..282e1fc 100644 --- a/migrations.py +++ b/migrations.py @@ -47,3 +47,44 @@ async def m001_initial(db): await db.execute( "CREATE INDEX payroll.idx_contracts_employee ON contracts (employee_id);" ) + + +async def m002_payouts(db): + """The payout ledger — one row per attempt at one period. + + Rows are self-describing (they carry the contract terms in force at the + time) so a payout still reads correctly after its contract is edited or + deleted. There is deliberately no foreign key to `contracts` for the + same reason: deleting a contract must not take its history with it. + """ + + await db.execute(f""" + CREATE TABLE payroll.payouts ( + id TEXT PRIMARY KEY, + contract_id TEXT NOT NULL, + period_index INTEGER NOT NULL, + payday TEXT NOT NULL, + status TEXT NOT NULL, + attempt INTEGER NOT NULL DEFAULT 1, + amount_msat {db.big_int}, + amount REAL NOT NULL DEFAULT 0, + currency TEXT NOT NULL DEFAULT 'sat', + employee_id TEXT NOT NULL DEFAULT '', + employee_wallet TEXT NOT NULL DEFAULT '', + source_wallet TEXT NOT NULL DEFAULT '', + payment_hash TEXT, + detail TEXT NOT NULL DEFAULT '', + created_at TIMESTAMP NOT NULL DEFAULT {db.timestamp_now} + ); + """) + + # Counting a period's prior failures runs on every failed attempt. + await db.execute( + "CREATE INDEX payroll.idx_payouts_period " + "ON payouts (contract_id, period_index);" + ) + # The employee-facing payslip view filters on the destination wallet. + await db.execute( + "CREATE INDEX payroll.idx_payouts_employee_wallet " + "ON payouts (employee_wallet);" + ) diff --git a/models.py b/models.py index c877200..00edf83 100644 --- a/models.py +++ b/models.py @@ -191,3 +191,56 @@ class DirectoryUser(BaseModel): @property def display_name(self) -> str: return self.username or self.email or self.id + + +# --------------------------------------------------------------------------- +# Payout ledger +# --------------------------------------------------------------------------- + + +class PayoutStatus(str, Enum): + paid = "paid" # money moved + skipped = "skipped" # period deliberately not paid (back-dated, paused-over) + failed = "failed" # attempted and could not settle; will be retried + + +class Payout(BaseModel): + """One attempt at one period — the audit trail. + + Deliberately self-describing rather than a thin join key. A ledger row + has to still make sense after its contract is edited or deleted, so the + terms in force at the time (amount, currency, both wallets) are copied + onto it. This is the one place in the extension where duplicating a + contract field is right: the contract holds what is true *now*, a payout + holds what was true *then*. + + `amount_msat` is the canonical settled figure, taken from the invoice + LNbits priced. It is null on a skip, and on a failure that never got as + far as raising an invoice. + """ + + id: str + contract_id: str + period_index: int + payday: str # YYYY-MM-DD + status: PayoutStatus + # 1-based, counted per (contract, period). Only failures accumulate + # attempts — a paid or skipped period is never retried. + attempt: int = 1 + + amount_msat: int | None = None + amount: float = 0 # the instruction in force at the time + currency: str = "sat" + + employee_id: str = "" + employee_wallet: str = "" + source_wallet: str = "" + + payment_hash: str | None = None + detail: str = "" + + created_at: datetime = Field(default_factory=_now) + + @property + def amount_sat(self) -> int | None: + return None if self.amount_msat is None else self.amount_msat // 1000 diff --git a/services.py b/services.py index d54df83..3d4ede0 100644 --- a/services.py +++ b/services.py @@ -22,10 +22,18 @@ from datetime import date, datetime, timedelta, timezone from lnbits.core.crud import get_wallet from lnbits.core.services import create_invoice, pay_invoice +from lnbits.helpers import urlsafe_short_hash from loguru import logger from . import crud -from .models import TERMINAL_STATUSES, Contract, ContractStatus, Frequency +from .models import ( + TERMINAL_STATUSES, + Contract, + ContractStatus, + Frequency, + Payout, + PayoutStatus, +) # Frequencies that are an exact number of days: no calendar involved, so no # clamping is possible or needed. @@ -46,6 +54,12 @@ _MONTH_STEP = { # Protects against a mistyped start date turning into a hundred transfers. MAX_CATCH_UP_PERIODS = 12 +# How many times one period may fail before the contract is paused for the +# operator to look at. Retrying an underfunded wallet forever is not +# resilience, it is a log the operator learns to ignore — and a missed +# payday deserves to be visible in the UI, not buried in journalctl. +MAX_PERIOD_ATTEMPTS = 5 + # --------------------------------------------------------------------------- # Schedule math (pure) @@ -141,7 +155,7 @@ class PeriodOutcome: index: int payday: date - status: str # "paid" | "skipped" | "failed" + status: PayoutStatus detail: str = "" amount_msat: int | None = None payment_hash: str | None = None @@ -154,7 +168,7 @@ class PeriodOutcome: what makes the next tick retry the same payday rather than dropping it. """ - return self.status in ("paid", "skipped") + return self.status in (PayoutStatus.paid, PayoutStatus.skipped) # One lock per contract id. The scheduler is a single task, but an operator @@ -207,11 +221,13 @@ async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutc internal=True, extra=extra, ) + # Broad: pricing, wallet limits and funding-source errors all surface + # here, and every one of them is a failed period rather than a crash. except Exception as exc: return PeriodOutcome( index=index, payday=payday, - status="failed", + status=PayoutStatus.failed, detail=f"could not raise invoice: {exc}", ) @@ -223,7 +239,7 @@ async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutc return PeriodOutcome( index=index, payday=payday, - status="failed", + status=PayoutStatus.failed, detail="source wallet not found", amount_msat=amount_msat, ) @@ -231,7 +247,7 @@ async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutc return PeriodOutcome( index=index, payday=payday, - status="failed", + status=PayoutStatus.failed, detail=( f"insufficient balance in source wallet: " f"{source.balance_msat // 1000} sat available, " @@ -248,11 +264,11 @@ async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutc tag="payroll", extra=extra, ) - except Exception as exc: + except Exception as exc: # a failed payment is an expected outcome return PeriodOutcome( index=index, payday=payday, - status="failed", + status=PayoutStatus.failed, detail=f"payment failed: {exc}", amount_msat=amount_msat, payment_hash=invoice.payment_hash, @@ -261,7 +277,7 @@ async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutc return PeriodOutcome( index=index, payday=payday, - status="paid", + status=PayoutStatus.paid, amount_msat=amount_msat, payment_hash=invoice.payment_hash, ) @@ -269,9 +285,12 @@ async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutc async def run_due_periods( contract: Contract, today: date | None = None -) -> list[PeriodOutcome]: +) -> list[Payout]: """Settle every period this contract owes as of `today`. + Every attempt — paid, skipped or failed — is written to the ledger, and + the ledger rows are what this returns. + Stops at the first failure so a backlog cannot pay periods out of order: if period 3 could not be funded, period 4 waits for it. Both will be retried on the next tick, because a failed period does not advance the @@ -287,22 +306,79 @@ async def run_due_periods( return [] contract = fresh - outcomes: list[PeriodOutcome] = [] + payouts: list[Payout] = [] for index in due_period_indices(contract, today): payday = occurrence_on( parse_start_date(contract), contract.frequency, index ) outcome = await _settle(contract, index, payday) - outcomes.append(outcome) + payout = await _record(contract, outcome) + payouts.append(payout) + if not outcome.consumed: + _maybe_pause_after_repeated_failure(contract, payout) break contract.periods_done = index + 1 - if outcomes: + if payouts: _maybe_complete(contract) await crud.update_contract(contract) - return outcomes + return payouts + + +async def _record(contract: Contract, outcome: PeriodOutcome) -> Payout: + """Write one attempt to the ledger. + + The attempt number is derived by counting this period's prior failures + in the ledger itself rather than from a counter on the contract, so it + survives a restart and the rows that produced it are right there to + audit. A paid or skipped period is never retried, so it is always + attempt 1 as far as that outcome is concerned. + """ + attempt = 1 + if outcome.status == PayoutStatus.failed: + attempt = await crud.count_period_failures(contract.id, outcome.index) + 1 + + return await crud.create_payout( + Payout( + id=urlsafe_short_hash()[:12], + contract_id=contract.id, + period_index=outcome.index, + payday=outcome.payday.isoformat(), + status=PayoutStatus(outcome.status), + attempt=attempt, + amount_msat=outcome.amount_msat, + # Terms in force at the time — a ledger row has to still read + # correctly after the contract is edited or deleted. + amount=contract.amount, + currency=contract.currency, + employee_id=contract.employee_id, + employee_wallet=contract.employee_wallet, + source_wallet=contract.source_wallet, + payment_hash=outcome.payment_hash, + detail=outcome.detail, + ) + ) + + +def _maybe_pause_after_repeated_failure(contract: Contract, payout: Payout) -> None: + """Give up retrying a period and hand it to the operator. + + The contract is paused rather than the period abandoned: a payday that + cannot be funded is a fact somebody needs to act on, and silently + dropping it is the one outcome payroll must never produce. Resuming + after fixing the funding retries the same period, because pausing did + not advance the position. + """ + if payout.attempt < MAX_PERIOD_ATTEMPTS: + return + contract.status = ContractStatus.paused + logger.error( + f"payroll: contract {contract.id} paused after {payout.attempt} failed " + f"attempts at period {payout.period_index} ({payout.payday}): " + f"{payout.detail}" + ) async def _settle(contract: Contract, index: int, payday: date) -> PeriodOutcome: @@ -322,12 +398,12 @@ async def _settle(contract: Contract, index: int, payday: date) -> PeriodOutcome return PeriodOutcome( index=index, payday=payday, - status="skipped", + status=PayoutStatus.skipped, detail="payday predates the contract and backfill is off", ) outcome = await pay_period(contract, index, payday) - if outcome.status == "paid": + if outcome.status == PayoutStatus.paid: logger.success( f"payroll: contract {contract.id} period {index} paid " f"{(outcome.amount_msat or 0) // 1000} sat to " @@ -336,12 +412,16 @@ async def _settle(contract: Contract, index: int, payday: date) -> PeriodOutcome else: logger.warning( f"payroll: contract {contract.id} period {index} " - f"{outcome.status}: {outcome.detail}" + f"{outcome.status.value}: {outcome.detail}" ) return outcome def _maybe_complete(contract: Contract) -> None: + # Only an active contract completes: a period that just exhausted the + # retry budget has paused this one, and that has to stick. + if contract.status != ContractStatus.active: + return if ( contract.total_periods is not None and contract.periods_done >= contract.total_periods @@ -359,7 +439,7 @@ async def tick(today: date | None = None) -> None: for contract in contracts: try: await run_due_periods(contract, today) - except Exception as exc: + except Exception as exc: # one bad row must not stop the rest of payroll logger.error(f"payroll: contract {contract.id} tick failed: {exc}") diff --git a/tasks.py b/tasks.py index ceb0d48..f64910f 100644 --- a/tasks.py +++ b/tasks.py @@ -31,6 +31,6 @@ async def scheduler_loop(): while True: try: await services.tick() - except Exception as exc: + except Exception as exc: # the loop must outlive any single tick logger.error(f"payroll: scheduler tick failed: {exc}") await asyncio.sleep(TICK_SECONDS) diff --git a/tests/test_payout.py b/tests/test_payout.py index 8a1399f..9ceaea1 100644 --- a/tests/test_payout.py +++ b/tests/test_payout.py @@ -12,8 +12,8 @@ from datetime import date import pytest -from ..models import ContractStatus -from ..services import PeriodOutcome, run_due_periods +from ..models import ContractStatus, PayoutStatus +from ..services import MAX_PERIOD_ATTEMPTS, PeriodOutcome, run_due_periods from .conftest import make_contract @@ -29,7 +29,11 @@ def payroll_stub(monkeypatch): def __init__(self): self.contract = None self.attempted: list[int] = [] - self.outcome_status = "paid" + self.outcome_status = PayoutStatus.paid + self.ledger: list = [] + # Failures already on record for (contract, period), as the real + # ledger would report them. + self.prior_failures = 0 stub = Stub() @@ -40,13 +44,21 @@ def payroll_stub(monkeypatch): stub.contract = contract return contract + async def fake_create_payout(payout): + stub.ledger.append(payout) + return payout + + async def fake_count_period_failures(_contract_id, _period_index): + return stub.prior_failures + async def fake_pay_period(contract, index, payday): stub.attempted.append(index) + paid = stub.outcome_status == PayoutStatus.paid return PeriodOutcome( index=index, payday=payday, status=stub.outcome_status, - detail="" if stub.outcome_status == "paid" else "stubbed failure", + detail="" if paid else "stubbed failure", amount_msat=int(contract.amount) * 1000, ) @@ -54,6 +66,10 @@ def payroll_stub(monkeypatch): monkeypatch.setattr(services.crud, "get_contract", fake_get_contract) monkeypatch.setattr(services.crud, "update_contract", fake_update_contract) + monkeypatch.setattr(services.crud, "create_payout", fake_create_payout) + monkeypatch.setattr( + services.crud, "count_period_failures", fake_count_period_failures + ) monkeypatch.setattr(services, "pay_period", fake_pay_period) # Locks are created inside whichever loop is running; each test drives its # own asyncio.run, so a lock left over from a previous test would raise. @@ -74,7 +90,9 @@ def test_backdated_periods_are_skipped_when_backfill_is_off(payroll_stub): outcomes = asyncio.run(run_due_periods(payroll_stub.contract, date(2026, 5, 20))) - assert [o.status for o in outcomes] == ["skipped"] * 4 + ["paid"] + assert [p.status for p in outcomes] == [PayoutStatus.skipped] * 4 + [ + PayoutStatus.paid + ] assert payroll_stub.attempted == [4] # only the first in-life period paid # Skipped periods still consume the schedule, so the next payday is June. assert payroll_stub.contract.periods_done == 5 @@ -87,7 +105,7 @@ def test_backfill_pays_the_backlog(payroll_stub): outcomes = asyncio.run(run_due_periods(payroll_stub.contract, date(2026, 3, 20))) - assert [o.status for o in outcomes] == ["paid", "paid", "paid"] + assert [p.status for p in outcomes] == [PayoutStatus.paid] * 3 assert payroll_stub.attempted == [0, 1, 2] assert payroll_stub.contract.periods_done == 3 @@ -98,11 +116,11 @@ def test_a_failure_halts_the_backlog_and_holds_the_position(payroll_stub): payroll_stub.contract = make_contract( start_date="2026-01-15", created_at="2026-01-01", backfill=True ) - payroll_stub.outcome_status = "failed" + payroll_stub.outcome_status = PayoutStatus.failed outcomes = asyncio.run(run_due_periods(payroll_stub.contract, date(2026, 4, 20))) - assert [o.status for o in outcomes] == ["failed"] + assert [p.status for p in outcomes] == [PayoutStatus.failed] assert payroll_stub.attempted == [0] # period 1 did not jump ahead assert payroll_stub.contract.periods_done == 0 # retried next tick @@ -130,3 +148,65 @@ def test_nothing_runs_before_the_first_payday(payroll_stub): assert outcomes == [] assert payroll_stub.attempted == [] + + +# --- ledger ---------------------------------------------------------------- + + +def test_every_attempt_is_recorded_including_skips(payroll_stub): + payroll_stub.contract = make_contract( + start_date="2026-01-15", created_at="2026-03-01", backfill=False + ) + + asyncio.run(run_due_periods(payroll_stub.contract, date(2026, 3, 20))) + + assert [(p.period_index, p.status) for p in payroll_stub.ledger] == [ + (0, PayoutStatus.skipped), + (1, PayoutStatus.skipped), + (2, PayoutStatus.paid), + ] + + +def test_a_ledger_row_carries_the_terms_in_force(payroll_stub): + """Rows have to still read correctly after the contract is edited or + deleted, so they copy the terms rather than pointing at them.""" + payroll_stub.contract = make_contract( + start_date="2026-01-15", created_at="2026-01-01", amount=800, currency="EUR" + ) + + asyncio.run(run_due_periods(payroll_stub.contract, date(2026, 1, 15))) + + row = payroll_stub.ledger[0] + assert (row.amount, row.currency) == (800, "EUR") + assert row.employee_wallet == "wallet-employee" + assert row.source_wallet == "wallet-treasury" + assert row.amount_sat == 800 # from the stubbed invoice, not amount x rate + + +def test_repeated_failure_pauses_the_contract(payroll_stub): + """A payday that cannot be funded is a fact somebody has to act on. + Retrying it silently forever is how a missed salary goes unnoticed.""" + payroll_stub.contract = make_contract( + start_date="2026-01-15", created_at="2026-01-01" + ) + payroll_stub.outcome_status = PayoutStatus.failed + payroll_stub.prior_failures = MAX_PERIOD_ATTEMPTS - 1 + + asyncio.run(run_due_periods(payroll_stub.contract, date(2026, 1, 20))) + + assert payroll_stub.contract.status == ContractStatus.paused + # The position did not move, so resuming retries the same payday. + assert payroll_stub.contract.periods_done == 0 + + +def test_a_failure_under_the_cap_leaves_the_contract_active(payroll_stub): + payroll_stub.contract = make_contract( + start_date="2026-01-15", created_at="2026-01-01" + ) + payroll_stub.outcome_status = PayoutStatus.failed + payroll_stub.prior_failures = MAX_PERIOD_ATTEMPTS - 2 + + asyncio.run(run_due_periods(payroll_stub.contract, date(2026, 1, 20))) + + assert payroll_stub.contract.status == ContractStatus.active + assert payroll_stub.ledger[-1].attempt == MAX_PERIOD_ATTEMPTS - 1 diff --git a/views_api.py b/views_api.py index 9f6a83c..5d977f0 100644 --- a/views_api.py +++ b/views_api.py @@ -21,7 +21,14 @@ from lnbits.utils.exchange_rates import allowed_currencies from . import crud, services from .accounts import list_directory_users, owns_wallet -from .models import Contract, CreateContract, DirectoryUser, UpdateContract +from .models import ( + Contract, + CreateContract, + DirectoryUser, + Payout, + PayoutStatus, + UpdateContract, +) payroll_api_router = APIRouter(dependencies=[Depends(check_super_user)]) @@ -215,3 +222,24 @@ async def api_cancel_contract(contract_id: str) -> Contract: mistake" escape hatch and discards the record entirely. """ return await _transition(contract_id, services.cancel) + + +# --------------------------------------------------------------------------- +# Payout ledger +# --------------------------------------------------------------------------- + + +@payroll_api_router.get("/api/v1/payouts") +async def api_list_payouts( + contract_id: str | None = None, + status: PayoutStatus | None = None, + limit: int = 200, +) -> list[Payout]: + """Every payout attempt, newest first. + + Includes failures and skips, not just successes — "why did nobody get + paid on the 1st" is the question this endpoint exists to answer. + """ + return await crud.get_payouts( + contract_id=contract_id, status=status, limit=min(limit, 1000) + )