294 lines
12 KiB
Python
294 lines
12 KiB
Python
"""Background processing job: run the engine over a session's files and persist results."""
|
|
from __future__ import annotations
|
|
|
|
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
|
|
|
|
|
|
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)
|
|
|
|
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,
|
|
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)
|
|
db.query(models.JournalEntry).filter(
|
|
models.JournalEntry.session_id == session_id).delete()
|
|
db.add(models.JournalEntry(session_id=session_id, data=json.dumps(payload)))
|
|
db.commit()
|
|
except Exception: # noqa: BLE001 — journal is supplementary; never fail the close over it
|
|
db.rollback()
|
|
|
|
session.status = "processed"
|
|
session.progress_stage = "Done"
|
|
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"{type(exc).__name__}: {exc}\n{traceback.format_exc()[-1500:]}")
|
|
finally:
|
|
if sink is not None:
|
|
sink.close()
|
|
db.close()
|
|
|
|
|
|
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")
|
|
export_summary_workbook(out_path, summary, control, journal, meta)
|
|
_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()
|
|
|
|
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,
|
|
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")
|
|
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)
|
|
|
|
_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()
|