payroll/services.py
Padreug e77f431d47 fix: resuming an auto-paused contract no longer discards the backlog
Found by tracing what happens when a back-dated backfill contract cannot
fetch a historical rate. The failure handling itself was fine — period 0
fails, the backlog halts so nothing settles out of order, five ledger rows
record the reason, no money moves, and the contract auto-pauses once the
retry budget is spent. The recovery was not.

The operator fixes the cause (switches to a stated rate, or to current),
clicks Resume, and periods_done jumps 0 -> 6: every unpaid payday silently
written off, contract back to looking healthy, employee never paid. The
confirm dialog even asserted the missed paydays "are written off" — true of
one kind of pause and a lie about the other.

Two features colliding. "Do not backfill a deliberate pause" is right when
the operator paused: the pause *was* the decision not to pay. It is wrong
when payroll paused, because nobody decided anything — the money is still
owed and the operator has just removed whatever blocked it.

Contracts now carry `paused_reason`, set only when payroll pauses them and
cleared by a deliberate pause. Resume infers from it, and an explicit
`catch_up` still overrides either way. The console asks a different question
for each, quoting the reason, and flags a payroll-paused contract in the
table so the distinction is visible before anyone clicks.

Verified end to end: five failing ticks leave periods_done at 0 and pause
with "period 0 (2026-08-01) failed 5 times: no historical EUR rate
available for 2026-08-01"; resuming after switching to a manual rate keeps
the position at 0, and the next tick settles all seven owed periods.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018jy52j9GRZ6XKa1Zt21LLj
2026-08-31 23:02:45 +02:00

776 lines
28 KiB
Python

"""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 lnbits.helpers import urlsafe_short_hash
from lnbits.utils.exchange_rates import fiat_amount_as_satoshis
from loguru import logger
from . import crud
from .models import (
TERMINAL_STATUSES,
Contract,
ContractStatus,
Frequency,
Payout,
PayoutStatus,
PricingMode,
)
from .rates import historical_btc_rate, sats_for
# 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
# 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)
# ---------------------------------------------------------------------------
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 paydays_from(
start: date,
frequency: Frequency,
first_index: int,
count: int,
total_periods: int | None = None,
) -> list[date]:
"""`count` paydays starting at period `first_index`, truncated by any cap.
Takes loose parameters rather than a Contract so the operator can preview
a schedule *before* the contract exists — which is the point at which a
mistyped start date or the wrong frequency is cheap to fix.
"""
return [
occurrence_on(start, frequency, i)
for i in range(first_index, first_index + count)
if total_periods is None or i < total_periods
]
def upcoming_paydays(contract: Contract, count: int) -> list[date]:
"""The next `count` paydays for a live contract, from its position."""
return paydays_from(
parse_start_date(contract),
contract.frequency,
contract.periods_done,
count,
contract.total_periods,
)
def final_payday(
start: date, frequency: Frequency, total_periods: int | None
) -> date | None:
"""When an open-ended contract would end — None, by definition."""
if total_periods is None:
return None
return occurrence_on(start, frequency, total_periods - 1)
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: PayoutStatus
detail: str = ""
amount_msat: int | None = None
payment_hash: str | None = None
rate: float | None = None
rate_source: str = ""
@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 (PayoutStatus.paid, PayoutStatus.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 _today() -> date:
return datetime.now(timezone.utc).date()
def _fmt(value: float) -> str:
"""Trim a float for display without lying about precision."""
return f"{value:,.2f}".rstrip("0").rstrip(".")
def payout_memo(contract: Contract, payday: date, rate: float | None = None) -> str:
"""What both wallets show for this payout.
Names the rate applied, because for a back-dated period the sat figure
on its own is unexplainable — 18,353 sat and 14,794 sat are both "ten
euro", and only the rate says why they differ. The rate is always the
one payroll actually converted at, never a lookup made for display.
"""
label = contract.memo or contract.label or "Payroll"
parts = [f"{label} — {payday.isoformat()}"]
if contract.currency.lower() != "sat":
money = f"{_fmt(contract.amount)} {contract.currency}"
if rate:
money += f" @ {_fmt(rate)} {contract.currency}/BTC"
parts.append(money)
return " · ".join(parts)
class PricingError(ValueError):
"""A period's sat amount could not be established.
Deliberately fatal for that period. The alternative — falling back to
today's rate when the payday's rate is unavailable — pays a materially
different amount than intended and hides it, which is exactly the class
of error nobody finds until an audit.
"""
@dataclass
class Price:
"""What to hand `create_invoice`, plus the provenance to record.
When `currency` is the contract's own, LNbits does the conversion and
`rate` is filled in afterwards from the invoice it priced — read back
rather than recomputed, so the ledger records what actually happened.
"""
amount: float
currency: str
rate: float | None
source: str
async def resolve_price(
contract: Contract,
payday: date,
today: date,
mode: PricingMode | None = None,
manual_rate: float | None = None,
) -> Price:
"""Decide what a period is worth, and say where the number came from.
`mode`/`manual_rate` override the contract's own settings for a single
call, which is how an operator prices one off-cycle payout differently
without editing the contract.
"""
if contract.currency.lower() == "sat":
return Price(contract.amount, "sat", None, "")
mode = mode or contract.pricing_mode
manual_rate = manual_rate or contract.manual_rate
# An explicit rate wins outright, whatever the calendar says. It is an
# instruction rather than a lookup — "pay at the rate we agreed" is as
# meaningful for a payday next week as for one last month, and a
# contract pegged to a fixed rate should honour it on every period.
if mode == PricingMode.manual:
if not manual_rate:
raise PricingError("manual pricing selected but no rate given")
return Price(
sats_for(contract.amount, manual_rate), "sat", manual_rate, "manual"
)
# A payday that is today or ahead has no history to consult, so `payday`
# mode collapses into `current` — a future rate is not knowable, and the
# live one is the only defensible answer.
#
# Payroll still does the conversion rather than handing create_invoice a
# fiat amount, because knowing the rate *before* the invoice exists is
# what lets the memo state it. Same single conversion either way —
# fiat_amount_as_satoshis is the function create_invoice would have
# called — not a second opinion.
if payday >= today or mode == PricingMode.current:
amount_sat = await fiat_amount_as_satoshis(contract.amount, contract.currency)
if amount_sat <= 0:
raise PricingError(
f"{contract.amount} {contract.currency} rounds to zero sats"
)
# Derived exactly as lnbits.core.services.calculate_fiat_amounts does.
rate = (contract.amount / amount_sat) * 100_000_000
return Price(amount_sat, "sat", rate, "current")
historical = await historical_btc_rate(payday, contract.currency)
if historical is None:
raise PricingError(
f"no historical {contract.currency} rate available for {payday}"
)
return Price(sats_for(contract.amount, historical), "sat", historical, "payday")
async def pay_period(
contract: Contract,
index: int,
payday: date,
mode: PricingMode | None = None,
manual_rate: float | None = None,
) -> 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.
"""
extra = {
"tag": "payroll",
"contract_id": contract.id,
"period": index,
"payday": payday.isoformat(),
}
try:
price = await resolve_price(contract, payday, _today(), mode, manual_rate)
except (PricingError, ValueError) as exc:
return PeriodOutcome(
index=index,
payday=payday,
status=PayoutStatus.failed,
# The memo cannot be built without a rate, so name the contract
# plainly in the failure instead.
detail=str(exc),
)
memo = payout_memo(contract, payday, price.rate)
if price.source:
extra["rate_source"] = price.source
if price.rate:
# Payroll converted, so create_invoice is handed sats and will not
# stamp these itself. Same keys and formulas as calculate_fiat_amounts,
# so a payroll payment reads like any other fiat-priced one.
extra.update(
{
"fiat_currency": contract.currency,
"fiat_amount": round(contract.amount, 3),
"fiat_rate": price.amount / contract.amount,
"btc_rate": price.rate,
}
)
try:
invoice = await create_invoice(
wallet_id=contract.employee_wallet,
amount=price.amount,
currency=price.currency,
memo=memo,
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=PayoutStatus.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
rate = price.rate
source = await get_wallet(contract.source_wallet)
if not source:
return PeriodOutcome(
index=index,
payday=payday,
status=PayoutStatus.failed,
detail="source wallet not found",
amount_msat=amount_msat,
rate=rate,
rate_source=price.source,
)
if source.balance_msat < amount_msat:
return PeriodOutcome(
index=index,
payday=payday,
status=PayoutStatus.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,
rate=rate,
rate_source=price.source,
)
try:
await pay_invoice(
wallet_id=contract.source_wallet,
payment_request=invoice.bolt11,
description=memo,
tag="payroll",
extra=extra,
)
except Exception as exc: # a failed payment is an expected outcome
return PeriodOutcome(
index=index,
payday=payday,
status=PayoutStatus.failed,
detail=f"payment failed: {exc}",
amount_msat=amount_msat,
payment_hash=invoice.payment_hash,
rate=rate,
rate_source=price.source,
)
return PeriodOutcome(
index=index,
payday=payday,
status=PayoutStatus.paid,
amount_msat=amount_msat,
payment_hash=invoice.payment_hash,
rate=rate,
rate_source=price.source,
)
async def run_due_periods(
contract: Contract, today: date | None = None
) -> 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
contract's position.
"""
today = today or _today()
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
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)
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 payouts:
_maybe_complete(contract)
await crud.update_contract(contract)
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,
rate=outcome.rate,
rate_source=outcome.rate_source,
)
)
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
contract.paused_reason = (
f"period {payout.period_index} ({payout.payday}) failed "
f"{payout.attempt} times: {payout.detail}"
)
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:
"""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=PayoutStatus.skipped,
detail="payday predates the contract and backfill is off",
)
outcome = await pay_period(contract, index, payday)
if outcome.status == PayoutStatus.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.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
):
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: # one bad row must not stop the rest of payroll
logger.error(f"payroll: contract {contract.id} tick failed: {exc}")
# ---------------------------------------------------------------------------
# Lifecycle
# ---------------------------------------------------------------------------
class LifecycleError(ValueError):
"""An illegal status transition. Mapped to 409 at the API boundary."""
def fast_forward_index(contract: Contract, today: date) -> int:
"""The first period index whose payday is not already in the past."""
start = parse_start_date(contract)
cap = contract.total_periods
index = contract.periods_done
while (cap is None or index < cap) and occurrence_on(
start, contract.frequency, index
) < today:
index += 1
return index
def pause(contract: Contract) -> Contract:
if contract.status != ContractStatus.active:
raise LifecycleError(
f"Only an active contract can be paused (is {contract.status.value})."
)
contract.status = ContractStatus.paused
# A deliberate pause carries no reason, which is what makes resume treat
# the missed paydays as written off rather than still owed.
contract.paused_reason = ""
return contract
def resume(
contract: Contract, today: date | None = None, catch_up: bool | None = None
) -> Contract:
"""Put a paused contract back to work.
What happens to the paydays that fell during the pause depends on who
paused it, because the two cases mean opposite things:
* **The operator paused it.** Those paydays are written off — the pause
*was* the decision not to pay them, and resuming into a surprise
multi-period transfer is the opposite of what "resume" implies. The
position fast-forwards to the next payday on or after today.
* **Payroll paused it** — a period exhausted its retries, because the
wallet was empty or the date could not be priced. Nobody decided
anything, the money is still owed, and the operator has just fixed
whatever blocked it. The backlog is kept and settles on the next tick.
Sharing one path between those two silently destroys money that is owed,
so the inference is on `paused_reason`. `catch_up` overrides it.
"""
if contract.status != ContractStatus.paused:
raise LifecycleError(
f"Only a paused contract can be resumed (is {contract.status.value})."
)
if catch_up is None:
catch_up = bool(contract.paused_reason)
contract.status = ContractStatus.active
contract.paused_reason = ""
if not catch_up:
today = today or _today()
skipped_to = fast_forward_index(contract, today)
if skipped_to != contract.periods_done:
logger.info(
f"payroll: contract {contract.id} resumed, skipping "
f"{skipped_to - contract.periods_done} payday(s) missed while paused"
)
contract.periods_done = skipped_to
_maybe_complete(contract)
return contract
def cancel(contract: Contract) -> Contract:
"""Stop a contract for good, keeping the row and its history.
This — not DELETE — is how a running payroll is stopped: deleting throws
away the schedule position and the record that the contract ever existed.
"""
if contract.status in TERMINAL_STATUSES:
raise LifecycleError(f"Contract is already {contract.status.value}.")
contract.status = ContractStatus.cancelled
return contract
# ---------------------------------------------------------------------------
# Off-cycle payout
# ---------------------------------------------------------------------------
async def pay_now(
contract: Contract,
mode: PricingMode | None = None,
manual_rate: float | None = None,
) -> Payout:
"""Settle the contract's next period immediately, whatever the calendar
says.
The operator's manual override: "run it now" when waiting five minutes
for the tick is not acceptable, and "pay it early" when a payday needs to
land before it is due. Both are the same operation — the next period
settles and is consumed — so there is one endpoint rather than two that
differ only in whether today happens to be the payday.
It calls `pay_period` directly rather than going through `_settle`,
because `_settle`'s back-dated skip exists to stop a *new* contract
firing surprise back-pay. An operator explicitly asking to pay a period
is not a surprise.
`mode`/`manual_rate` price this one payout differently without editing
the contract — the case an operator hits when entering a payment that
happened weeks ago at a rate they already know.
"""
async with _lock_for(contract.id):
fresh = await crud.get_contract(contract.id)
if not fresh:
raise LifecycleError("Contract not found.")
contract = fresh
if contract.status != ContractStatus.active:
raise LifecycleError(
f"Only an active contract can be paid (is {contract.status.value})."
)
remaining = contract.periods_remaining
if remaining is not None and remaining <= 0:
raise LifecycleError("Contract has no periods left to pay.")
index = contract.periods_done
payday = occurrence_on(parse_start_date(contract), contract.frequency, index)
today = _today()
outcome = await pay_period(contract, index, payday, mode, manual_rate)
if payday > today and outcome.status == PayoutStatus.paid:
outcome.detail = f"paid early on {today.isoformat()} (due {payday})"
payout = await _record(contract, outcome)
if outcome.consumed:
contract.periods_done = index + 1
_maybe_complete(contract)
else:
_maybe_pause_after_repeated_failure(contract, payout)
await crud.update_contract(contract)
return payout