Files
Conflict_Checker/backend/jobs.py
T
John Wilganowski afa1089311
Docker Release / build-and-push (push) Successful in 55s
Docker Release / release (push) Skipped
Add job run logs, OpenRouter model picker, and discipline grouping.
- Job logs: each job's stdout/stderr is teed into outputs/<id>/job.log
  (survives restarts) and served at GET /jobs/{id}/log as text/plain, so
  full run logs can be shared for debugging and refinement.
- Model picker: GET /models proxies OpenRouter's public model list with
  per-1M-token pricing (1h cache, 502 on failure); the UI shows a model
  dropdown with costs when OpenRouter compute is selected, and the pick
  overrides vision+text models for that job (Classic and Agent modes).
- Conflicts in the report view are grouped by discipline pair
  (collapsible sections, severity-ordered within groups) instead of one
  flat severity-only list.
2026-07-28 21:23:42 +00:00

282 lines
11 KiB
Python

"""
jobs.py - Lightweight async job registry for conflict checks.
A conflict check takes minutes, so the HTTP request must not block on it. Each
upload becomes a job that runs on a background thread; the client gets a job_id
immediately and can either poll GET /jobs/{id} or just close the page and wait
for the completion email.
State is in-memory (fine for a single-user tool); the report is also persisted
to outputs/<job_id>/ so results survive a restart even though live status does
not. No external queue/DB.
"""
import json
import contextlib
import os
import sys
import time
import uuid
import shutil
import threading
from typing import Dict, Optional
from backend import config
from backend import llm
from backend.agents.runner import run_agent_pipeline
from backend.pipeline.runner import run_pipeline
from backend.email_sender import send_conflict_report, send_review_required
_jobs: Dict[str, Dict] = {}
_lock = threading.Lock()
PIPELINE_MODES = {"classic", "agent"}
class _Tee:
"""Write to both the real stream and the job log file."""
def __init__(self, stream, log_file) -> None:
self._stream = stream
self._log = log_file
def write(self, data):
self._stream.write(data)
self._log.write(data)
def flush(self):
self._stream.flush()
self._log.flush()
@contextlib.contextmanager
def _tee_log(log_path: str, header: str):
"""Mirror stdout/stderr into a per-job log file for the duration of a run.
sys.stdout is process-global, so two concurrent jobs would interleave in
each other's logs - acceptable for this single-user tool (same tradeoff as
the LLM cost globals in llm.py).
"""
with open(log_path, "a", encoding="utf-8") as log_file:
log_file.write(header + "\n")
real_out, real_err = sys.stdout, sys.stderr
sys.stdout, sys.stderr = _Tee(real_out, log_file), _Tee(real_err, log_file)
try:
yield
finally:
sys.stdout, sys.stderr = real_out, real_err
def _set(job_id: str, **fields) -> None:
with _lock:
_jobs[job_id].update(fields)
def create_job(pdf_path: str, source_filename: str, email: Optional[str] = None,
project_input: Optional[Dict] = None, text_local: bool = False,
pipeline_mode: str = "classic", model: Optional[str] = None) -> str:
"""Register a job and kick off its background thread. Returns the job_id."""
pipeline_mode = pipeline_mode.strip().lower()
if pipeline_mode not in PIPELINE_MODES:
raise ValueError(f"Unsupported pipeline mode: {pipeline_mode!r}")
# Agent mode v1 is OpenRouter-only.
text_local = bool(text_local and pipeline_mode == "classic")
model = (model or "").strip() or None
job_id = uuid.uuid4().hex[:12]
with _lock:
_jobs[job_id] = {
"job_id": job_id,
"status": "queued", # queued -> running -> done | needs_review | error
"source": source_filename,
"email": email or None,
"project_input": project_input or {},
"text_local": text_local,
"pipeline_mode": pipeline_mode,
"model": model,
"stage": None,
"created_at": time.time(),
"finished_at": None,
"report": None,
"error": None,
}
threading.Thread(target=_run, args=(
job_id, pdf_path, project_input, text_local, pipeline_mode, model,
),
daemon=True).start()
return job_id
def _run(job_id: str, pdf_path: str, project_input: Optional[Dict] = None,
text_local: bool = False, pipeline_mode: str = "classic",
model: Optional[str] = None) -> None:
out_dir = os.path.join(config.OUTPUT_DIR, job_id)
try:
_set(job_id, status="running")
# Keep a copy of the source PDF so its sheets can be viewed later.
os.makedirs(out_dir, exist_ok=True)
header = (f"=== Job {job_id} | {pipeline_mode} | {_jobs[job_id].get('source')} | "
f"model={model or 'default'} | "
f"started {time.strftime('%Y-%m-%d %H:%M:%S %Z', time.gmtime())} UTC ===")
with _tee_log(os.path.join(out_dir, "job.log"), header):
_run_pipeline(job_id, pdf_path, out_dir, project_input, text_local,
pipeline_mode, model)
except Exception as e:
print(f"[Jobs] Job {job_id} failed: {e}")
_set(job_id, status="error", error=str(e), finished_at=time.time())
_notify_error(job_id)
finally:
try:
os.remove(pdf_path)
except OSError:
pass
def _run_pipeline(job_id: str, pdf_path: str, out_dir: str,
project_input: Optional[Dict], text_local: bool,
pipeline_mode: str, model: Optional[str]) -> None:
"""The body of a job run; executes inside the job's tee'd log capture."""
# Persist minimal job metadata so the disk fallback in get_job can
# recover the recipient email / pipeline mode after a server restart
# (plain json.dump, matching the _dump style used elsewhere).
with open(os.path.join(out_dir, "job.json"), "w", encoding="utf-8") as f:
json.dump({
"job_id": job_id,
"email": _jobs[job_id].get("email"),
"pipeline_mode": pipeline_mode,
"source": _jobs[job_id].get("source"),
}, f, indent=2)
shutil.copy2(pdf_path, os.path.join(out_dir, "source.pdf"))
runner = run_agent_pipeline if pipeline_mode == "agent" else run_pipeline
runner_kwargs = {
"out_dir": out_dir,
"on_stage": lambda name: _set(job_id, stage=name),
"project_input": project_input,
"source_name": _jobs[job_id].get("source"),
}
if pipeline_mode == "classic":
runner_kwargs["text_local"] = text_local
else:
runner_kwargs["require_review"] = config.AGENT_REQUIRE_REVIEW
if model:
print(f"[Jobs] Model override for this run: {model}")
llm.set_model_override(model)
try:
report = runner(pdf_path, **runner_kwargs)
finally:
if model:
llm.set_model_override(None)
report.setdefault("summary", {})["pipeline_mode"] = pipeline_mode
if report["summary"].get("agent_status") == "needs_review":
# Human-review gate: hold the job, don't email the unreviewed report.
_set(job_id, status="needs_review", report=report,
finished_at=time.time(), stage=None)
email = _jobs[job_id].get("email")
if email:
review_url = f"{config.APP_BASE_URL.rstrip('/')}/?job={job_id}"
send_review_required(email, report, review_url)
else:
_set(job_id, status="done", report=report, finished_at=time.time(), stage=None)
_notify(job_id, report, out_dir)
def _notify(job_id: str, report: Dict, out_dir: str) -> None:
email = _jobs[job_id].get("email")
if not email:
return
results_url = f"{config.APP_BASE_URL.rstrip('/')}/?job={job_id}"
attachments = [
os.path.join(out_dir, "report.md"),
os.path.join(out_dir, "conflicts.json"),
os.path.join(out_dir, "validated_issues.json"),
os.path.join(out_dir, "rfis.json"),
]
send_conflict_report(email, report, results_url=results_url, attachments=attachments)
def _notify_error(job_id: str) -> None:
job = _jobs[job_id]
email = job.get("email")
if not email:
return
# Reuse the report mailer with a minimal error-shaped payload.
err_report = {
"source": job.get("source", ""),
"summary": {"conflicts_found": 0, "by_severity": {}, "disciplines": []},
}
try:
from backend.email_sender import _smtp_ready, _send
from email.message import EmailMessage
if not _smtp_ready():
print(f"[Email] SMTP not configured - skipping error notice to {email}")
return
msg = EmailMessage()
msg["Subject"] = f"Conflict Checker - {job.get('source','')} - run FAILED"
msg["From"] = config.SMTP_FROM or config.SMTP_USER
msg["To"] = email
msg.set_content(
"Your conflict check did not complete.\n\n"
f"Drawing set: {job.get('source','')}\n"
f"Error: {job.get('error','unknown')}\n\n"
"Generated by Conflict Checker"
)
_send(msg)
except Exception as e:
print(f"[Email] Failed to send error notice: {e}")
def get_job(job_id: str) -> Optional[Dict]:
"""Public job view. Includes the full report only when done.
Falls back to the on-disk conflicts.json when the job isn't in the
in-memory registry (e.g. after a server restart).
"""
with _lock:
job = _jobs.get(job_id)
if job:
return dict(job)
# Try loading from disk
report_path = os.path.join(config.OUTPUT_DIR, job_id, "conflicts.json")
if not os.path.isfile(report_path):
return None
try:
with open(report_path, encoding="utf-8") as f:
report = json.load(f)
summary = report.get("summary", {})
# Recover the job's real state: a job that stopped at the review gate
# must come back as needs_review (not done) or it can never finalize.
status = "needs_review" if summary.get("agent_status") == "needs_review" else "done"
# job.json (written at job start) carries the recipient email and
# pipeline mode so the final notification still fires after a restart.
# Missing/corrupt job.json degrades to the previous derivations.
meta: Dict = {}
meta_path = os.path.join(config.OUTPUT_DIR, job_id, "job.json")
try:
with open(meta_path, encoding="utf-8") as f:
loaded = json.load(f)
if isinstance(loaded, dict):
meta = loaded
except (OSError, json.JSONDecodeError):
pass
source_pdf = os.path.join(config.OUTPUT_DIR, job_id, "source.pdf")
job = {
"job_id": job_id,
"status": status,
"source": meta.get("source") or report.get("source", os.path.basename(report_path)),
"email": meta.get("email"),
"project_input": report.get("project_input", {}),
"text_local": summary.get("text_backend") == "local",
"pipeline_mode": meta.get("pipeline_mode") or summary.get("pipeline_mode", "classic"),
"stage": None,
"created_at": os.path.getmtime(source_pdf) if os.path.isfile(source_pdf) else None,
"finished_at": os.path.getmtime(report_path),
"report": report,
"error": None,
}
# Hydrate the in-memory registry so _set(...) transitions (reviewing,
# finalizing, done) work for restart-recovered jobs.
with _lock:
return dict(_jobs.setdefault(job_id, job))
except Exception as e:
print(f"[Jobs] Failed to load job {job_id} from disk: {e}")
return None