"""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 ..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() # 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()