feat: recurring payout scheduler
Turns a contract into money moving. One permanent task ticks every five minutes and settles whatever each active contract owes; the actual transfer is a plain internal LNbits invoice on the employee's wallet, paid from the source wallet. services.py is split into a pure half and an effectful half on purpose. Paydays are the part of payroll that is easy to get subtly wrong and expensive to get wrong in production, so the schedule math has no DB, no wallets and no clock of its own, and is covered by tests. Decisions worth reviewing: - The n-th payday is a function of start_date and n alone. Advancing a stored date would drift on every late tick and would pin a month-end contract to the 28th forever; anchoring means 31 Jan pays 28 Feb and then 31 Mar. Tested both ways round. - A failed period does not advance the contract's position, and a backlog halts at the first failure so paydays cannot settle out of order. - The sat amount is derived exactly once, by create_invoice, and the value it returns is what gets recorded — never recomputed from amount x rate. - Back-dated start dates skip rather than back-pay by default; a mistyped start date is far more likely than a genuine back-pay request. Explicit `backfill` opts in, and a single tick is capped at 12 periods either way. - Per-contract asyncio lock, with the row re-read under it. Not needed by the scheduler alone, but off-cycle payout paths land mid-tick and double-paying is the worst thing this extension could do. Known gap, addressed by the payout-ledger commit that follows: a failure is retried indefinitely, once per tick, with nothing but a log line to show for it. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_018jy52j9GRZ6XKa1Zt21LLj
This commit is contained in:
parent
55f3dfac91
commit
99f2131474
8 changed files with 829 additions and 2 deletions
363
services.py
Normal file
363
services.py
Normal file
|
|
@ -0,0 +1,363 @@
|
|||
"""Payroll scheduling and payout.
|
||||
|
||||
Two separable things live here, and the split is deliberate:
|
||||
|
||||
* **Schedule math** (`occurrence_on`, `due_period_indices`) is pure — no
|
||||
DB, no wallets, no clock of its own. Paydays are the part of payroll that
|
||||
is easy to get subtly wrong and expensive to get wrong in production, so
|
||||
it is kept testable in isolation.
|
||||
* **Payout** (`run_due_periods`) is the effectful half: it turns a due
|
||||
period into an actual wallet-to-wallet transfer.
|
||||
|
||||
Everything is anchored on `start_date`. The *n*-th payday is a function of
|
||||
the start date and `n` alone — never of the previous payday, and never of
|
||||
when the scheduler happened to run. That is what stops a missed tick, a
|
||||
restart, or a short month from shifting the whole remaining schedule.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
from calendar import monthrange
|
||||
from dataclasses import dataclass
|
||||
from datetime import date, datetime, timedelta, timezone
|
||||
|
||||
from lnbits.core.crud import get_wallet
|
||||
from lnbits.core.services import create_invoice, pay_invoice
|
||||
from loguru import logger
|
||||
|
||||
from . import crud
|
||||
from .models import Contract, ContractStatus, Frequency
|
||||
|
||||
# Frequencies that are an exact number of days: no calendar involved, so no
|
||||
# clamping is possible or needed.
|
||||
_DAY_STEP = {
|
||||
Frequency.daily: 1,
|
||||
Frequency.weekly: 7,
|
||||
Frequency.biweekly: 14,
|
||||
}
|
||||
|
||||
# Frequencies that step whole months and therefore clamp into short months.
|
||||
_MONTH_STEP = {
|
||||
Frequency.monthly: 1,
|
||||
Frequency.quarterly: 3,
|
||||
Frequency.yearly: 12,
|
||||
}
|
||||
|
||||
# A single tick will not fire more than this many periods for one contract.
|
||||
# Protects against a mistyped start date turning into a hundred transfers.
|
||||
MAX_CATCH_UP_PERIODS = 12
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Schedule math (pure)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def add_months(anchor: date, months: int) -> date:
|
||||
"""`anchor` shifted by `months`, clamped to the target month's last day.
|
||||
|
||||
31 Jan + 1 month is 28 Feb (29 in a leap year), because 31 Feb does not
|
||||
exist. Note this is applied to the *anchor*, not to a previous result:
|
||||
31 Jan + 2 months is 31 Mar, not 28 Mar. Payroll anchored on month-end
|
||||
should keep paying on month-end.
|
||||
"""
|
||||
total = anchor.year * 12 + (anchor.month - 1) + months
|
||||
year, month_index = divmod(total, 12)
|
||||
month = month_index + 1
|
||||
return date(year, month, min(anchor.day, monthrange(year, month)[1]))
|
||||
|
||||
|
||||
def occurrence_on(start: date, frequency: Frequency, index: int) -> date:
|
||||
"""The payday for period `index`, counting the start date as period 0."""
|
||||
if frequency in _DAY_STEP:
|
||||
return start + timedelta(days=_DAY_STEP[frequency] * index)
|
||||
return add_months(start, _MONTH_STEP[frequency] * index)
|
||||
|
||||
|
||||
def parse_start_date(contract: Contract) -> date:
|
||||
return datetime.strptime(contract.start_date, "%Y-%m-%d").date()
|
||||
|
||||
|
||||
def next_payday(contract: Contract) -> date | None:
|
||||
"""When this contract pays next, or None if it has no periods left."""
|
||||
remaining = contract.periods_remaining
|
||||
if remaining is not None and remaining <= 0:
|
||||
return None
|
||||
return occurrence_on(
|
||||
parse_start_date(contract), contract.frequency, contract.periods_done
|
||||
)
|
||||
|
||||
|
||||
def upcoming_paydays(contract: Contract, count: int) -> list[date]:
|
||||
"""The next `count` paydays, truncated by any period cap."""
|
||||
start = parse_start_date(contract)
|
||||
last = contract.total_periods if contract.total_periods is not None else None
|
||||
indices = range(contract.periods_done, contract.periods_done + count)
|
||||
return [
|
||||
occurrence_on(start, contract.frequency, i)
|
||||
for i in indices
|
||||
if last is None or i < last
|
||||
]
|
||||
|
||||
|
||||
def due_period_indices(contract: Contract, today: date) -> list[int]:
|
||||
"""Every period index that is payable as of `today`, in order.
|
||||
|
||||
Normally this is zero or one entry. It is a list because a contract can
|
||||
legitimately have a backlog: a back-dated start date, or an instance
|
||||
that was down over a payday. Returning them all lets the caller decide
|
||||
what to do with each — `run_due_periods` pays or skips them one at a
|
||||
time rather than collapsing a backlog into a single surprise transfer.
|
||||
"""
|
||||
if contract.status != ContractStatus.active:
|
||||
return []
|
||||
|
||||
start = parse_start_date(contract)
|
||||
cap = contract.total_periods
|
||||
indices: list[int] = []
|
||||
index = contract.periods_done
|
||||
while cap is None or index < cap:
|
||||
if occurrence_on(start, contract.frequency, index) > today:
|
||||
break
|
||||
indices.append(index)
|
||||
index += 1
|
||||
# A daily contract left unpaid for years would otherwise build an
|
||||
# unbounded list. Anything beyond this is an operator problem, not
|
||||
# something to silently drain in one tick.
|
||||
if len(indices) >= MAX_CATCH_UP_PERIODS:
|
||||
break
|
||||
return indices
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Payout (effectful)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
@dataclass
|
||||
class PeriodOutcome:
|
||||
"""What happened to one period. `amount_msat` is the canonical settled
|
||||
figure, taken from the invoice LNbits actually priced — never
|
||||
recomputed from amount x rate."""
|
||||
|
||||
index: int
|
||||
payday: date
|
||||
status: str # "paid" | "skipped" | "failed"
|
||||
detail: str = ""
|
||||
amount_msat: int | None = None
|
||||
payment_hash: str | None = None
|
||||
|
||||
@property
|
||||
def consumed(self) -> bool:
|
||||
"""Whether this period should advance the contract's position.
|
||||
|
||||
Paid and skipped both consume the period. Failed does not — that is
|
||||
what makes the next tick retry the same payday rather than dropping
|
||||
it.
|
||||
"""
|
||||
return self.status in ("paid", "skipped")
|
||||
|
||||
|
||||
# One lock per contract id. The scheduler is a single task, but an operator
|
||||
# can trigger an off-cycle payout at any moment; without this, a manual run
|
||||
# landing mid-tick could pay the same period twice. Double-paying is the
|
||||
# worst thing this extension can do, so the guard is cheap insurance.
|
||||
_contract_locks: dict[str, asyncio.Lock] = {}
|
||||
|
||||
|
||||
def _lock_for(contract_id: str) -> asyncio.Lock:
|
||||
lock = _contract_locks.get(contract_id)
|
||||
if lock is None:
|
||||
lock = asyncio.Lock()
|
||||
_contract_locks[contract_id] = lock
|
||||
return lock
|
||||
|
||||
|
||||
def payout_memo(contract: Contract, payday: date) -> str:
|
||||
label = contract.memo or contract.label or "Payroll"
|
||||
return f"{label} — {payday.isoformat()}"
|
||||
|
||||
|
||||
async def pay_period(contract: Contract, index: int, payday: date) -> PeriodOutcome:
|
||||
"""Move one period's money from the source wallet to the employee wallet.
|
||||
|
||||
An internal invoice on the destination wallet, paid from the source
|
||||
wallet — the standard LNbits wallet-to-wallet transfer. Creating the
|
||||
invoice is also what prices a fiat contract in sats, so the amount is
|
||||
converted exactly once, here, and the result is what gets recorded.
|
||||
|
||||
A failure after the invoice exists leaves an unpaid internal invoice on
|
||||
the employee's wallet. That is harmless — it expires on its own and was
|
||||
never settled — and it is the price of not pre-computing the sat amount
|
||||
a second time just to run a balance check.
|
||||
"""
|
||||
memo = payout_memo(contract, payday)
|
||||
extra = {
|
||||
"tag": "payroll",
|
||||
"contract_id": contract.id,
|
||||
"period": index,
|
||||
"payday": payday.isoformat(),
|
||||
}
|
||||
|
||||
try:
|
||||
invoice = await create_invoice(
|
||||
wallet_id=contract.employee_wallet,
|
||||
amount=contract.amount,
|
||||
currency=contract.currency,
|
||||
memo=memo,
|
||||
internal=True,
|
||||
extra=extra,
|
||||
)
|
||||
except Exception as exc:
|
||||
return PeriodOutcome(
|
||||
index=index,
|
||||
payday=payday,
|
||||
status="failed",
|
||||
detail=f"could not raise invoice: {exc}",
|
||||
)
|
||||
|
||||
# invoice.amount is msat and is now the canonical figure for this period.
|
||||
amount_msat = invoice.amount
|
||||
|
||||
source = await get_wallet(contract.source_wallet)
|
||||
if not source:
|
||||
return PeriodOutcome(
|
||||
index=index,
|
||||
payday=payday,
|
||||
status="failed",
|
||||
detail="source wallet not found",
|
||||
amount_msat=amount_msat,
|
||||
)
|
||||
if source.balance_msat < amount_msat:
|
||||
return PeriodOutcome(
|
||||
index=index,
|
||||
payday=payday,
|
||||
status="failed",
|
||||
detail=(
|
||||
f"insufficient balance in source wallet: "
|
||||
f"{source.balance_msat // 1000} sat available, "
|
||||
f"{amount_msat // 1000} sat required"
|
||||
),
|
||||
amount_msat=amount_msat,
|
||||
)
|
||||
|
||||
try:
|
||||
await pay_invoice(
|
||||
wallet_id=contract.source_wallet,
|
||||
payment_request=invoice.bolt11,
|
||||
description=memo,
|
||||
tag="payroll",
|
||||
extra=extra,
|
||||
)
|
||||
except Exception as exc:
|
||||
return PeriodOutcome(
|
||||
index=index,
|
||||
payday=payday,
|
||||
status="failed",
|
||||
detail=f"payment failed: {exc}",
|
||||
amount_msat=amount_msat,
|
||||
payment_hash=invoice.payment_hash,
|
||||
)
|
||||
|
||||
return PeriodOutcome(
|
||||
index=index,
|
||||
payday=payday,
|
||||
status="paid",
|
||||
amount_msat=amount_msat,
|
||||
payment_hash=invoice.payment_hash,
|
||||
)
|
||||
|
||||
|
||||
async def run_due_periods(
|
||||
contract: Contract, today: date | None = None
|
||||
) -> list[PeriodOutcome]:
|
||||
"""Settle every period this contract owes as of `today`.
|
||||
|
||||
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
|
||||
contract's position.
|
||||
"""
|
||||
today = today or datetime.now(timezone.utc).date()
|
||||
|
||||
async with _lock_for(contract.id):
|
||||
# Re-read under the lock: a manual run may have moved the position
|
||||
# since the caller loaded this row.
|
||||
fresh = await crud.get_contract(contract.id)
|
||||
if not fresh:
|
||||
return []
|
||||
contract = fresh
|
||||
|
||||
outcomes: list[PeriodOutcome] = []
|
||||
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)
|
||||
if not outcome.consumed:
|
||||
break
|
||||
contract.periods_done = index + 1
|
||||
|
||||
if outcomes:
|
||||
_maybe_complete(contract)
|
||||
await crud.update_contract(contract)
|
||||
|
||||
return outcomes
|
||||
|
||||
|
||||
async def _settle(contract: Contract, index: int, payday: date) -> PeriodOutcome:
|
||||
"""Pay a period, or skip it as back-dated.
|
||||
|
||||
A contract created with a start date in the past is normally an
|
||||
*anchoring* choice — "we pay on the 1st" — not a request for back-pay.
|
||||
Periods whose payday fell before the contract existed are therefore
|
||||
skipped unless the operator explicitly asked to backfill, which keeps a
|
||||
new contract from firing months of transfers on its first tick.
|
||||
"""
|
||||
if not contract.backfill and payday < contract.created_at.date():
|
||||
logger.info(
|
||||
f"payroll: contract {contract.id} period {index} ({payday}) "
|
||||
"predates the contract; skipping (backfill is off)"
|
||||
)
|
||||
return PeriodOutcome(
|
||||
index=index,
|
||||
payday=payday,
|
||||
status="skipped",
|
||||
detail="payday predates the contract and backfill is off",
|
||||
)
|
||||
|
||||
outcome = await pay_period(contract, index, payday)
|
||||
if outcome.status == "paid":
|
||||
logger.success(
|
||||
f"payroll: contract {contract.id} period {index} paid "
|
||||
f"{(outcome.amount_msat or 0) // 1000} sat to "
|
||||
f"{contract.employee_wallet}"
|
||||
)
|
||||
else:
|
||||
logger.warning(
|
||||
f"payroll: contract {contract.id} period {index} "
|
||||
f"{outcome.status}: {outcome.detail}"
|
||||
)
|
||||
return outcome
|
||||
|
||||
|
||||
def _maybe_complete(contract: Contract) -> None:
|
||||
if (
|
||||
contract.total_periods is not None
|
||||
and contract.periods_done >= contract.total_periods
|
||||
):
|
||||
contract.status = ContractStatus.completed
|
||||
logger.info(
|
||||
f"payroll: contract {contract.id} completed "
|
||||
f"({contract.periods_done} periods)"
|
||||
)
|
||||
|
||||
|
||||
async def tick(today: date | None = None) -> None:
|
||||
"""One scheduler pass over every active contract."""
|
||||
contracts = await crud.get_contracts_by_status(ContractStatus.active)
|
||||
for contract in contracts:
|
||||
try:
|
||||
await run_due_periods(contract, today)
|
||||
except Exception as exc:
|
||||
logger.error(f"payroll: contract {contract.id} tick failed: {exc}")
|
||||
Loading…
Add table
Add a link
Reference in a new issue