feat: payout ledger with bounded retry
Closes the gap the scheduler commit left open: a failed payday was retried indefinitely with nothing but a log line to show for it. Every attempt — paid, skipped and failed — is now written to payroll.payouts and exposed at GET /api/v1/payouts. Failures are in the ledger, not only in the log, because "why did nobody get paid on the 1st" is the question the ledger exists to answer. Ledger rows are self-describing: each copies the terms in force at the time (amount, currency, both wallets) instead of pointing at the contract, and there is no foreign key to contracts. A payout has to still read correctly after its contract is edited, and deleting a contract must not take its history with it. This is the one place duplicating a contract field is right — the contract holds what is true now, a payout holds what was true then. Retries are bounded at 5 attempts per period, after which the contract is paused rather than the period abandoned. A payday that cannot be funded is a fact somebody has to act on; dropping it silently is the one outcome payroll must never produce. Pausing does not advance the position, so resuming after topping up retries the same payday. The attempt count is derived from the ledger rather than a column on the contract, so it survives a restart and stays auditable. Also restores the explanatory comments on the broad `except Exception` handlers, which ruff's RUF100 stripped along with their now-unused noqa directives. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_018jy52j9GRZ6XKa1Zt21LLj
This commit is contained in:
parent
95bfa86f79
commit
8b054d27ea
9 changed files with 389 additions and 30 deletions
|
|
@ -23,7 +23,7 @@ def payroll_stop():
|
||||||
for task in scheduled_tasks:
|
for task in scheduled_tasks:
|
||||||
try:
|
try:
|
||||||
task.cancel()
|
task.cancel()
|
||||||
except Exception as ex:
|
except Exception as ex: # shutdown must not raise
|
||||||
logger.warning(ex)
|
logger.warning(ex)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
53
crud.py
53
crud.py
|
|
@ -12,7 +12,7 @@ from datetime import datetime, timezone
|
||||||
from lnbits.db import Database
|
from lnbits.db import Database
|
||||||
from lnbits.helpers import urlsafe_short_hash
|
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")
|
db = Database("ext_payroll")
|
||||||
|
|
||||||
|
|
@ -87,3 +87,54 @@ async def delete_contract(contract_id: str) -> None:
|
||||||
await db.execute(
|
await db.execute(
|
||||||
"DELETE FROM payroll.contracts WHERE id = :id", {"id": contract_id}
|
"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
|
||||||
|
|
|
||||||
|
|
@ -56,6 +56,32 @@ payroll: contract a1b2c3d4 period 3 failed: insufficient balance in source
|
||||||
wallet: 12000 sat available, 80000 sat required
|
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
|
## Back-dated start dates
|
||||||
|
|
||||||
Creating a contract with a start date in the past is normally an
|
Creating a contract with a start date in the past is normally an
|
||||||
|
|
|
||||||
|
|
@ -47,3 +47,44 @@ async def m001_initial(db):
|
||||||
await db.execute(
|
await db.execute(
|
||||||
"CREATE INDEX payroll.idx_contracts_employee ON contracts (employee_id);"
|
"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);"
|
||||||
|
)
|
||||||
|
|
|
||||||
53
models.py
53
models.py
|
|
@ -191,3 +191,56 @@ class DirectoryUser(BaseModel):
|
||||||
@property
|
@property
|
||||||
def display_name(self) -> str:
|
def display_name(self) -> str:
|
||||||
return self.username or self.email or self.id
|
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
|
||||||
|
|
|
||||||
116
services.py
116
services.py
|
|
@ -22,10 +22,18 @@ from datetime import date, datetime, timedelta, timezone
|
||||||
|
|
||||||
from lnbits.core.crud import get_wallet
|
from lnbits.core.crud import get_wallet
|
||||||
from lnbits.core.services import create_invoice, pay_invoice
|
from lnbits.core.services import create_invoice, pay_invoice
|
||||||
|
from lnbits.helpers import urlsafe_short_hash
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
|
|
||||||
from . import crud
|
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
|
# Frequencies that are an exact number of days: no calendar involved, so no
|
||||||
# clamping is possible or needed.
|
# clamping is possible or needed.
|
||||||
|
|
@ -46,6 +54,12 @@ _MONTH_STEP = {
|
||||||
# Protects against a mistyped start date turning into a hundred transfers.
|
# Protects against a mistyped start date turning into a hundred transfers.
|
||||||
MAX_CATCH_UP_PERIODS = 12
|
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)
|
# Schedule math (pure)
|
||||||
|
|
@ -141,7 +155,7 @@ class PeriodOutcome:
|
||||||
|
|
||||||
index: int
|
index: int
|
||||||
payday: date
|
payday: date
|
||||||
status: str # "paid" | "skipped" | "failed"
|
status: PayoutStatus
|
||||||
detail: str = ""
|
detail: str = ""
|
||||||
amount_msat: int | None = None
|
amount_msat: int | None = None
|
||||||
payment_hash: str | 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
|
what makes the next tick retry the same payday rather than dropping
|
||||||
it.
|
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
|
# 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,
|
internal=True,
|
||||||
extra=extra,
|
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:
|
except Exception as exc:
|
||||||
return PeriodOutcome(
|
return PeriodOutcome(
|
||||||
index=index,
|
index=index,
|
||||||
payday=payday,
|
payday=payday,
|
||||||
status="failed",
|
status=PayoutStatus.failed,
|
||||||
detail=f"could not raise invoice: {exc}",
|
detail=f"could not raise invoice: {exc}",
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -223,7 +239,7 @@ async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutc
|
||||||
return PeriodOutcome(
|
return PeriodOutcome(
|
||||||
index=index,
|
index=index,
|
||||||
payday=payday,
|
payday=payday,
|
||||||
status="failed",
|
status=PayoutStatus.failed,
|
||||||
detail="source wallet not found",
|
detail="source wallet not found",
|
||||||
amount_msat=amount_msat,
|
amount_msat=amount_msat,
|
||||||
)
|
)
|
||||||
|
|
@ -231,7 +247,7 @@ async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutc
|
||||||
return PeriodOutcome(
|
return PeriodOutcome(
|
||||||
index=index,
|
index=index,
|
||||||
payday=payday,
|
payday=payday,
|
||||||
status="failed",
|
status=PayoutStatus.failed,
|
||||||
detail=(
|
detail=(
|
||||||
f"insufficient balance in source wallet: "
|
f"insufficient balance in source wallet: "
|
||||||
f"{source.balance_msat // 1000} sat available, "
|
f"{source.balance_msat // 1000} sat available, "
|
||||||
|
|
@ -248,11 +264,11 @@ async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutc
|
||||||
tag="payroll",
|
tag="payroll",
|
||||||
extra=extra,
|
extra=extra,
|
||||||
)
|
)
|
||||||
except Exception as exc:
|
except Exception as exc: # a failed payment is an expected outcome
|
||||||
return PeriodOutcome(
|
return PeriodOutcome(
|
||||||
index=index,
|
index=index,
|
||||||
payday=payday,
|
payday=payday,
|
||||||
status="failed",
|
status=PayoutStatus.failed,
|
||||||
detail=f"payment failed: {exc}",
|
detail=f"payment failed: {exc}",
|
||||||
amount_msat=amount_msat,
|
amount_msat=amount_msat,
|
||||||
payment_hash=invoice.payment_hash,
|
payment_hash=invoice.payment_hash,
|
||||||
|
|
@ -261,7 +277,7 @@ async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutc
|
||||||
return PeriodOutcome(
|
return PeriodOutcome(
|
||||||
index=index,
|
index=index,
|
||||||
payday=payday,
|
payday=payday,
|
||||||
status="paid",
|
status=PayoutStatus.paid,
|
||||||
amount_msat=amount_msat,
|
amount_msat=amount_msat,
|
||||||
payment_hash=invoice.payment_hash,
|
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(
|
async def run_due_periods(
|
||||||
contract: Contract, today: date | None = None
|
contract: Contract, today: date | None = None
|
||||||
) -> list[PeriodOutcome]:
|
) -> list[Payout]:
|
||||||
"""Settle every period this contract owes as of `today`.
|
"""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:
|
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
|
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
|
retried on the next tick, because a failed period does not advance the
|
||||||
|
|
@ -287,22 +306,79 @@ async def run_due_periods(
|
||||||
return []
|
return []
|
||||||
contract = fresh
|
contract = fresh
|
||||||
|
|
||||||
outcomes: list[PeriodOutcome] = []
|
payouts: list[Payout] = []
|
||||||
for index in due_period_indices(contract, today):
|
for index in due_period_indices(contract, today):
|
||||||
payday = occurrence_on(
|
payday = occurrence_on(
|
||||||
parse_start_date(contract), contract.frequency, index
|
parse_start_date(contract), contract.frequency, index
|
||||||
)
|
)
|
||||||
outcome = await _settle(contract, index, payday)
|
outcome = await _settle(contract, index, payday)
|
||||||
outcomes.append(outcome)
|
payout = await _record(contract, outcome)
|
||||||
|
payouts.append(payout)
|
||||||
|
|
||||||
if not outcome.consumed:
|
if not outcome.consumed:
|
||||||
|
_maybe_pause_after_repeated_failure(contract, payout)
|
||||||
break
|
break
|
||||||
contract.periods_done = index + 1
|
contract.periods_done = index + 1
|
||||||
|
|
||||||
if outcomes:
|
if payouts:
|
||||||
_maybe_complete(contract)
|
_maybe_complete(contract)
|
||||||
await crud.update_contract(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:
|
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(
|
return PeriodOutcome(
|
||||||
index=index,
|
index=index,
|
||||||
payday=payday,
|
payday=payday,
|
||||||
status="skipped",
|
status=PayoutStatus.skipped,
|
||||||
detail="payday predates the contract and backfill is off",
|
detail="payday predates the contract and backfill is off",
|
||||||
)
|
)
|
||||||
|
|
||||||
outcome = await pay_period(contract, index, payday)
|
outcome = await pay_period(contract, index, payday)
|
||||||
if outcome.status == "paid":
|
if outcome.status == PayoutStatus.paid:
|
||||||
logger.success(
|
logger.success(
|
||||||
f"payroll: contract {contract.id} period {index} paid "
|
f"payroll: contract {contract.id} period {index} paid "
|
||||||
f"{(outcome.amount_msat or 0) // 1000} sat to "
|
f"{(outcome.amount_msat or 0) // 1000} sat to "
|
||||||
|
|
@ -336,12 +412,16 @@ async def _settle(contract: Contract, index: int, payday: date) -> PeriodOutcome
|
||||||
else:
|
else:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
f"payroll: contract {contract.id} period {index} "
|
f"payroll: contract {contract.id} period {index} "
|
||||||
f"{outcome.status}: {outcome.detail}"
|
f"{outcome.status.value}: {outcome.detail}"
|
||||||
)
|
)
|
||||||
return outcome
|
return outcome
|
||||||
|
|
||||||
|
|
||||||
def _maybe_complete(contract: Contract) -> None:
|
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 (
|
if (
|
||||||
contract.total_periods is not None
|
contract.total_periods is not None
|
||||||
and contract.periods_done >= contract.total_periods
|
and contract.periods_done >= contract.total_periods
|
||||||
|
|
@ -359,7 +439,7 @@ async def tick(today: date | None = None) -> None:
|
||||||
for contract in contracts:
|
for contract in contracts:
|
||||||
try:
|
try:
|
||||||
await run_due_periods(contract, today)
|
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}")
|
logger.error(f"payroll: contract {contract.id} tick failed: {exc}")
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
2
tasks.py
2
tasks.py
|
|
@ -31,6 +31,6 @@ async def scheduler_loop():
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
await services.tick()
|
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}")
|
logger.error(f"payroll: scheduler tick failed: {exc}")
|
||||||
await asyncio.sleep(TICK_SECONDS)
|
await asyncio.sleep(TICK_SECONDS)
|
||||||
|
|
|
||||||
|
|
@ -12,8 +12,8 @@ from datetime import date
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from ..models import ContractStatus
|
from ..models import ContractStatus, PayoutStatus
|
||||||
from ..services import PeriodOutcome, run_due_periods
|
from ..services import MAX_PERIOD_ATTEMPTS, PeriodOutcome, run_due_periods
|
||||||
from .conftest import make_contract
|
from .conftest import make_contract
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -29,7 +29,11 @@ def payroll_stub(monkeypatch):
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
self.contract = None
|
self.contract = None
|
||||||
self.attempted: list[int] = []
|
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()
|
stub = Stub()
|
||||||
|
|
||||||
|
|
@ -40,13 +44,21 @@ def payroll_stub(monkeypatch):
|
||||||
stub.contract = contract
|
stub.contract = contract
|
||||||
return 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):
|
async def fake_pay_period(contract, index, payday):
|
||||||
stub.attempted.append(index)
|
stub.attempted.append(index)
|
||||||
|
paid = stub.outcome_status == PayoutStatus.paid
|
||||||
return PeriodOutcome(
|
return PeriodOutcome(
|
||||||
index=index,
|
index=index,
|
||||||
payday=payday,
|
payday=payday,
|
||||||
status=stub.outcome_status,
|
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,
|
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, "get_contract", fake_get_contract)
|
||||||
monkeypatch.setattr(services.crud, "update_contract", fake_update_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)
|
monkeypatch.setattr(services, "pay_period", fake_pay_period)
|
||||||
# Locks are created inside whichever loop is running; each test drives its
|
# 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.
|
# 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)))
|
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
|
assert payroll_stub.attempted == [4] # only the first in-life period paid
|
||||||
# Skipped periods still consume the schedule, so the next payday is June.
|
# Skipped periods still consume the schedule, so the next payday is June.
|
||||||
assert payroll_stub.contract.periods_done == 5
|
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)))
|
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.attempted == [0, 1, 2]
|
||||||
assert payroll_stub.contract.periods_done == 3
|
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(
|
payroll_stub.contract = make_contract(
|
||||||
start_date="2026-01-15", created_at="2026-01-01", backfill=True
|
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)))
|
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.attempted == [0] # period 1 did not jump ahead
|
||||||
assert payroll_stub.contract.periods_done == 0 # retried next tick
|
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 outcomes == []
|
||||||
assert payroll_stub.attempted == []
|
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
|
||||||
|
|
|
||||||
30
views_api.py
30
views_api.py
|
|
@ -21,7 +21,14 @@ from lnbits.utils.exchange_rates import allowed_currencies
|
||||||
|
|
||||||
from . import crud, services
|
from . import crud, services
|
||||||
from .accounts import list_directory_users, owns_wallet
|
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)])
|
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.
|
mistake" escape hatch and discards the record entirely.
|
||||||
"""
|
"""
|
||||||
return await _transition(contract_id, services.cancel)
|
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)
|
||||||
|
)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue