Finance-Accounts/ar-aging-app/backend/app/services/jobs.py

433 lines
20 KiB
Python

"""Background processing job: run the engine over a session's files and persist results."""
from __future__ import annotations
import logging
import time
import traceback
from ..config import FX_AUTO_DAILY
from ..core.pipeline import process
from ..core.i18n import CURRENCY_BY_REGION, DEFAULT_FX_USD, currency_for_region, default_fx_for_region
from ..db import models
from ..db.database import SessionLocal
from .store import TransactionSink, persist_aggregates, clear_session_results
logger = logging.getLogger(__name__)
def recover_stale_jobs() -> None:
"""Called once at startup. Jobs run in-process (single worker), so a session still
marked processing/exporting at boot was killed mid-run by a restart or deploy — without
this it stays stuck forever and the 409 "already processing" guard blocks every re-run."""
db = SessionLocal()
try:
stale = db.query(models.Session).filter(
models.Session.status.in_(("processing", "exporting"))).all()
for s in stale:
was = s.status
s.status = "error" if was == "processing" else "processed"
s.error = ("Interrupted by a server restart before it finished — run it again."
if was == "processing"
else "Export was interrupted by a server restart — export again.")
s.progress_stage = "Interrupted"
logger.warning("recovered stale job: session %s (%s) was '%s'", s.id, s.name, was)
if stale:
db.commit()
except Exception: # noqa: BLE001 — recovery must never prevent startup
db.rollback()
logger.exception("stale-job recovery failed")
finally:
db.close()
def load_mapping_rules(db) -> dict[str, str]:
"""Admin-saved header rules (normalized header -> canonical field)."""
try:
return {r.normalized_header: r.field for r in db.query(models.MappingRule)}
except Exception: # noqa: BLE001 — table may not exist on very old DBs
return {}
def run_processing(session_id: int) -> None:
"""Executed in a background thread. Owns its own DB session."""
db = SessionLocal()
sink: TransactionSink | None = None
try:
session = db.get(models.Session, session_id)
if session is None:
return
files = db.query(models.SessionFile).filter(
models.SessionFile.session_id == session_id).all()
paths = [f.stored_path for f in files]
if not paths:
_fail(db, session, "No files uploaded.")
return
reserves = {(r.marketplace, r.account_type): r.amount
for r in db.query(models.Reserve).filter(
models.Reserve.session_id == session_id)}
fx_rows = db.query(models.FxRate).filter(models.FxRate.session_id == session_id).all()
# Marketplace defaults (seeded from the Jan-26 workbook) under session overrides.
fx_rates = {**DEFAULT_FX_USD, **{r.marketplace: r.rate for r in fx_rows}}
currencies = {**CURRENCY_BY_REGION, **{r.marketplace: r.currency for r in fx_rows}}
mapping_rules = load_mapping_rules(db)
# Bank receipts: Finance's record of when each payout actually reached the bank.
# A receipt overrides the clearing-lag heuristic for its payout — received iff the
# BANK date is on/before month-end. In manual mode the heuristic is off entirely
# and a payout without a receipt is not received.
receipts = db.query(models.PayoutReceipt).filter(
models.PayoutReceipt.session_id == session_id).all()
received_overrides = {
(r.marketplace, r.account_type, r.settlement_id):
bool(r.bank_date and session.month_end_date
and r.bank_date <= session.month_end_date)
for r in receipts
}
manual_payouts = (session.payout_mode or "auto") == "manual"
session.status = "processing"
session.error = ""
session.progress_stage = "Reading workbook"
session.progress_pct = 0.02
db.commit()
clear_session_results(db, session_id)
t0 = time.time()
def progress(stage: str, pct: float, rows_done: int = 0, rows_total: int = 0) -> None:
elapsed = time.time() - t0
eta = int(elapsed / pct - elapsed) if pct > 0.03 else 0
session.progress_stage = stage
session.progress_pct = pct
session.progress_rows_done = rows_done
session.progress_rows_total = rows_total
session.eta_seconds = max(eta, 0)
db.commit()
sink = TransactionSink(session_id)
result = process(
paths,
month_end=session.month_end_date,
clearing_lag_days=session.clearing_lag_days or 2,
reserves=reserves, fx_rates=fx_rates, currencies=currencies,
received_overrides=received_overrides, manual_payouts=manual_payouts,
manual_adjustments=session.manual_adjustment or 0.0,
tolerance=session.rounding_tolerance or 0.01,
saved_column_overrides=mapping_rules,
record_sink=sink.add, progress=progress,
)
progress("Persisting transactions", 0.94)
sink.finish()
sink.apply_classification(result)
progress("Saving results", 0.96)
persist_aggregates(db, session_id, result, files)
# Seed editable FX rows for any marketplace that appeared without one.
have_fx = {r.marketplace for r in fx_rows}
for mkt in (result.receivable.marketplaces if result.receivable else {}):
if mkt not in have_fx:
db.add(models.FxRate(
session_id=session_id, marketplace=mkt,
currency=currency_for_region(mkt), rate=default_fx_for_region(mkt),
source="default (Jan-26 workbook)"))
db.commit()
# Daily FX from the provider, covering the span of dates the files actually
# contain, so every dated movement converts at the rate effective on ITS OWN
# transaction date (ledger / fx-daily). Advisory: a provider outage never blocks
# the close — conversion falls back to the last available fixing, then the month
# rate, and the shortfall is surfaced below as an exception.
if FX_AUTO_DAILY:
progress("Fetching daily FX rates", 0.97)
from .fx_service import auto_seed_daily_fx
fx_daily_out = auto_seed_daily_fx(db, session)
if fx_daily_out.get("error"):
db.add(models.Exception_(
session_id=session_id, category="fx_daily_unavailable",
severity="warning",
detail=(f"Daily exchange rates could not be fetched from the provider "
f"({fx_daily_out['error']}). Dated movements convert at "
f"previously fetched daily rates or the month rate until "
f"'Fetch daily rates' on the AR Ledger tab succeeds."),
source="fx provider"))
db.commit()
# Journal-entry decomposition (separate pass; part of the close).
try:
progress("Building journal entry", 0.98)
import json
from ..core.journal import journal_payload
payload = journal_payload(paths, mapping_rules)
old = db.query(models.JournalEntry).filter(
models.JournalEntry.session_id == session_id).first()
# The entry number survives a re-process; review/approval deliberately do NOT —
# the numbers just changed, so any sign-off no longer attests to what's stored.
entry_no = old.entry_no if old else ""
# Bulk delete executes immediately — an ORM delete+add pair can flush the INSERT
# before the DELETE and trip the unique(session_id) constraint, silently keeping
# the OLD journal (and its stale sign-off) via the except below.
db.query(models.JournalEntry).filter(
models.JournalEntry.session_id == session_id).delete(synchronize_session=False)
db.add(models.JournalEntry(session_id=session_id, data=json.dumps(payload),
entry_no=entry_no))
db.commit()
except Exception: # noqa: BLE001 — journal is supplementary; never fail the close over it
db.rollback()
# Bank-receipt sanity: a receipt whose amount differs from Amazon's payout, or whose
# key matches no payout in the files, is surfaced — never silently absorbed.
try:
_receipt_exceptions(db, session_id, receipts, result)
except Exception: # noqa: BLE001 — advisory only; never fail the close over it
db.rollback()
session.status = "processed"
session.progress_stage = "Done"
session.progress_pct = 1.0
session.needs_reprocess = False # this run reflects the receipts as of now
db.commit()
# Month-end controls run LAST, over everything that was just persisted, and set
# status to "blocked" if any of them fails with error severity. A close that cannot
# be trusted must not publish a receivable figure.
progress("Running month-end controls", 0.99)
from .controls_run import run_and_persist
run_and_persist(db, session_id, result)
except Exception as exc: # noqa: BLE001
db.rollback()
session = db.get(models.Session, session_id)
if session:
_fail(db, session, f"{type(exc).__name__}: {exc}\n{traceback.format_exc()[-1500:]}")
finally:
if sink is not None:
sink.close()
db.close()
def _receipt_exceptions(db, session_id: int, receipts, result) -> None:
"""Warn on bank receipts that disagree with the files (wrong amount / no such payout)."""
agg = result.aggregation
if agg is None:
return
# Amazon payout per receipt key = the bucket's transfer total.
by_key = {(m, a, s): st.transfer_total for (m, a, s), st in agg.settlements.items()
if st.transfer_count}
for r in receipts:
key = (r.marketplace, r.account_type, r.settlement_id)
amazon = by_key.get(key)
if amazon is None:
db.add(models.Exception_(
session_id=session_id, category="unmatched_bank_receipt", severity="warning",
detail=(f"{r.marketplace}: a bank receipt dated {r.bank_date} is entered for "
f"settlement {r.settlement_id}, but no payout with that settlement "
f"exists in the uploaded files — check the settlement id"),
source=r.marketplace))
elif r.bank_amount is not None and abs(abs(r.bank_amount) - abs(amazon)) > 0.01:
db.add(models.Exception_(
session_id=session_id, category="bank_amount_variance", severity="warning",
detail=(f"{r.marketplace} settlement {r.settlement_id}: bank received "
f"{r.bank_amount:,.2f} on {r.bank_date} but Amazon's payout is "
f"{amazon:,.2f} — difference {abs(r.bank_amount) - abs(amazon):+,.2f} "
f"(bank fee or partial payment; the ledger uses Amazon's amount)"),
source=r.marketplace))
db.commit()
def _fail(db, session, message: str) -> None:
session.status = "error"
session.error = message
session.progress_stage = "Error"
db.commit()
def _finalize_export(db, session_id: int, out_path: str, kind: str) -> None:
import hashlib
import os
h = hashlib.sha256()
with open(out_path, "rb") as fh:
for chunk in iter(lambda: fh.read(1 << 20), b""):
h.update(chunk)
db.add(models.ExportRecord(session_id=session_id, kind=kind, path=out_path,
sha256=h.hexdigest(), size_bytes=os.path.getsize(out_path)))
def run_summary_export(session_id: int) -> None:
"""Background: write the compact Finance summary pack (reads stored results — no re-parse)."""
import json
from ..config import EXPORT_DIR
from ..core.summary_export import export_summary_workbook
from ..api.routes.ar import build_finance_summary
from ..api.routes.control import get_control
db = SessionLocal()
try:
session = db.get(models.Session, session_id)
if session is None:
return
session.status = "exporting"
session.progress_stage = "Building Finance summary"
session.progress_pct = 0.25
db.commit()
summary = build_finance_summary(db, session_id)
if not summary.get("available"):
raise RuntimeError("Process the closing before exporting the summary.")
control = get_control(session_id, db)
j = db.query(models.JournalEntry).filter(
models.JournalEntry.session_id == session_id).first()
journal = json.loads(j.data) if j and j.data else {}
files = db.query(models.SessionFile).filter(
models.SessionFile.session_id == session_id).all()
meta = {
"session_name": session.name,
"entry_no": j.entry_no if j else "",
"files": [{"filename": f.filename, "rows": f.imported_rows,
"dates": f"{f.min_date}..{f.max_date}", "sha256": f.sha256 or ""}
for f in files],
}
session.progress_stage = "Writing summary workbook"
session.progress_pct = 0.7
db.commit()
EXPORT_DIR.mkdir(parents=True, exist_ok=True)
month = session.reporting_month or "output"
out_path = str(EXPORT_DIR / f"AR_Summary_{month}_session{session_id}.xlsx")
from .controls_run import payload as controls_payload
export_summary_workbook(out_path, summary, control, journal, meta,
controls=controls_payload(db, session_id))
_finalize_export(db, session_id, out_path, "summary")
session.status = "processed"
session.progress_stage = "Summary ready"
session.progress_pct = 1.0
db.commit()
except Exception as e: # noqa: BLE001
db.rollback()
s = db.get(models.Session, session_id)
if s:
s.status = "processed"
s.error = f"{e}\n{traceback.format_exc(limit=3)}"
db.commit()
finally:
db.close()
def run_export(session_id: int) -> None:
"""Background: regenerate a fresh result and write the full A/R Aging workbook to disk."""
import hashlib
import os
from ..config import EXPORT_DIR
from ..core.excel_export import export_workbook
from ..db.database import ENGINE
db = SessionLocal()
try:
session = db.get(models.Session, session_id)
if session is None:
return
files = db.query(models.SessionFile).filter(
models.SessionFile.session_id == session_id).all()
paths = [f.stored_path for f in files]
reserves = {(r.marketplace, r.account_type): r.amount
for r in db.query(models.Reserve).filter(
models.Reserve.session_id == session_id)}
fx_rows = db.query(models.FxRate).filter(models.FxRate.session_id == session_id).all()
# Marketplace defaults (seeded from the Jan-26 workbook) under session overrides.
fx_rates = {**DEFAULT_FX_USD, **{r.marketplace: r.rate for r in fx_rows}}
currencies = {**CURRENCY_BY_REGION, **{r.marketplace: r.currency for r in fx_rows}}
mapping_rules = load_mapping_rules(db)
session.status = "exporting"
session.progress_stage = "Recomputing for export"
session.progress_pct = 0.02
db.commit()
t0 = time.time()
# Phase 1 (0.00..0.45): re-aggregate from source. Phase 2 (0.45..0.99): write workbook.
def progress(stage, pct, rows_done=0, rows_total=0):
elapsed = time.time() - t0
p = pct * 0.45
session.progress_stage = stage
session.progress_pct = p
session.progress_rows_done = rows_done
session.progress_rows_total = rows_total
session.eta_seconds = int(elapsed / p - elapsed) if p > 0.03 else 0
db.commit()
# Same bank-receipt overrides as run_processing — the exported workbook must show
# the identical paid/receivable split the dashboard shows.
receipts = db.query(models.PayoutReceipt).filter(
models.PayoutReceipt.session_id == session_id).all()
received_overrides = {
(r.marketplace, r.account_type, r.settlement_id):
bool(r.bank_date and session.month_end_date
and r.bank_date <= session.month_end_date)
for r in receipts
}
result = process(paths, month_end=session.month_end_date,
clearing_lag_days=session.clearing_lag_days or 2,
reserves=reserves, fx_rates=fx_rates, currencies=currencies,
received_overrides=received_overrides,
manual_payouts=(session.payout_mode or "auto") == "manual",
manual_adjustments=session.manual_adjustment or 0.0,
tolerance=session.rounding_tolerance or 0.01,
saved_column_overrides=mapping_rules, progress=progress)
session.progress_stage = "Writing workbook"
session.progress_pct = 0.45
db.commit()
def write_progress(frac, rows_done, rows_total):
elapsed = time.time() - t0
p = 0.45 + 0.54 * frac
session.progress_stage = f"Writing workbook ({rows_done:,}/{rows_total:,} rows)"
session.progress_pct = p
session.progress_rows_done = rows_done
session.progress_rows_total = rows_total
session.eta_seconds = int(elapsed / p - elapsed) if p > 0.03 else 0
db.commit()
# Finance summary + journal travel with the full workbook as their own sheets.
import json as _json
from ..api.routes.ar import build_finance_summary
try:
_summary = build_finance_summary(db, session_id)
_j = db.query(models.JournalEntry).filter(
models.JournalEntry.session_id == session_id).first()
_journal = _json.loads(_j.data) if _j and _j.data else {}
except Exception: # noqa: BLE001
_summary, _journal = {"available": False}, {}
EXPORT_DIR.mkdir(parents=True, exist_ok=True)
month = session.reporting_month or "output"
out_path = str(EXPORT_DIR / f"AR_Aging_{month}_session{session_id}.xlsx")
# Bank dates for the "Settled Settlements" sheet, so the workbook records WHY each
# excluded settlement was excluded and who said so.
receipt_notes = {
(r.marketplace, r.account_type, r.settlement_id):
f"{r.bank_date}" + (f" · {r.entered_by}" if r.entered_by else "")
for r in receipts if r.bank_date
}
export_workbook(result, paths, out_path, reserves=reserves,
allowance_for_returns=session.allowance_for_returns or 0.0,
saved_column_overrides=mapping_rules,
progress=write_progress, summary=_summary, journal=_journal,
payout_receipts=receipt_notes)
_finalize_export(db, session_id, out_path, "full")
session.status = "processed"
session.progress_stage = "Export ready"
session.progress_pct = 1.0
db.commit()
except Exception as exc: # noqa: BLE001
db.rollback()
session = db.get(models.Session, session_id)
if session:
_fail(db, session, f"Export failed: {type(exc).__name__}: {exc}")
finally:
db.close()