Merge main: dual model dropdowns + richer job logs, adapted for agent-mode.
Docker Release / build-and-push (push) Successful in 1m0s
Docker Release / release (push) Skipped

- llm.py: set_model_overrides(vision, text) replaces the single job override;
  UI picks still beat per-call agent model args, but never name the hybrid
  local model (avoids main's hybrid footgun); local->cloud fallback uses the
  text pick.
- jobs.py: timestamped line-split tee (job_log.py), in-memory log + log_tail
  polls, full log on terminal states (done/error/needs_review/finalization_error),
  log-only disk recovery, error email links to the run log, and failed runs now
  append the full traceback to job.log. Keeps pipeline_mode, job.json, and the
  review gate.
- models.py: vision/text split via architecture modalities, pricing kept;
  /models returns {vision, text, defaults}; /check takes vision_model/text_model
  (replacing model); /health adds text_model. models_catalog.py dropped.
- UI: two priced dropdowns (OpenRouter compute only) + live run-log panel.
- Tests updated for dual overrides and the /models shape; new coverage for
  traceback capture and local-model immunity.
This commit is contained in:
2026-08-02 09:55:09 -05:00
11 changed files with 699 additions and 150 deletions
+143 -78
View File
@@ -9,61 +9,33 @@ 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.
A teed stdout/stderr log is kept in memory and written to outputs/<job_id>/job.log
so failed or suspicious runs can be reviewed after the fact.
"""
import json
import contextlib
import os
import sys
import time
import traceback
import uuid
import shutil
import threading
from typing import Dict, Optional
from typing import Dict, List, Optional
from backend import config
from backend import llm
from backend.job_log import capture_stdio, read_log_file, stamp_line
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()
_LOG_TAIL = 80
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
# States where the job will produce no more log output; polls get the full log.
_TERMINAL_STATES = {"done", "error", "needs_review", "finalization_error"}
def _set(job_id: str, **fields) -> None:
@@ -71,16 +43,39 @@ def _set(job_id: str, **fields) -> None:
_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:
def _append_log(job_id: str, raw_line: str, log_path: str) -> None:
"""Stamp, store, and append one captured stdout/stderr line."""
entry = stamp_line(raw_line)
with _lock:
job = _jobs.get(job_id)
if job is not None:
job.setdefault("log", []).append(entry)
try:
os.makedirs(os.path.dirname(log_path), exist_ok=True)
with open(log_path, "a", encoding="utf-8") as f:
f.write(entry + "\n")
except OSError:
# Don't fail the job over log I/O; avoid print() here — it would
# re-enter the stdio tee while a job is capturing.
pass
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",
vision_model: Optional[str] = None,
text_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] = {
@@ -91,35 +86,61 @@ def create_job(pdf_path: str, source_filename: str, email: Optional[str] = None,
"project_input": project_input or {},
"text_local": text_local,
"pipeline_mode": pipeline_mode,
"model": model,
"vision_model": (vision_model or "").strip() or None,
"text_model": (text_model or "").strip() or None,
"stage": None,
"created_at": time.time(),
"finished_at": None,
"report": None,
"error": None,
"log": [],
}
threading.Thread(target=_run, args=(
job_id, pdf_path, project_input, text_local, pipeline_mode, model,
),
daemon=True).start()
threading.Thread(
target=_run,
args=(job_id, pdf_path, project_input, text_local, pipeline_mode,
vision_model, text_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:
def _run(
job_id: str,
pdf_path: str,
project_input: Optional[Dict] = None,
text_local: bool = False,
pipeline_mode: str = "classic",
vision_model: Optional[str] = None,
text_model: Optional[str] = None,
) -> None:
out_dir = os.path.join(config.OUTPUT_DIR, job_id)
log_path = os.path.join(out_dir, "job.log")
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)
# Truncate any leftover log if job_id somehow collided (shouldn't).
with open(log_path, "w", encoding="utf-8"):
pass
header = (f"=== Job {job_id} | {pipeline_mode} | {_jobs[job_id].get('source')} | "
f"model={model or 'default'} | "
f"vision={vision_model or 'default'} text={text_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):
_append_log(job_id, header, log_path)
def on_line(raw: str) -> None:
_append_log(job_id, raw, log_path)
with capture_stdio(on_line):
_run_pipeline(job_id, pdf_path, out_dir, project_input, text_local,
pipeline_mode, model)
pipeline_mode, vision_model, text_model)
except Exception as e:
# Land the failure AND its traceback in the job log so failed runs can
# be diagnosed from the log alone (the stdio tee is already torn down).
try:
_append_log(job_id, f"[Jobs] Job {job_id} failed: {e}", log_path)
for ln in traceback.format_exc().rstrip().splitlines():
_append_log(job_id, ln, log_path)
except Exception:
pass
print(f"[Jobs] Job {job_id} failed: {e}")
_set(job_id, status="error", error=str(e), finished_at=time.time())
_notify_error(job_id)
@@ -132,7 +153,8 @@ def _run(job_id: str, pdf_path: str, project_input: Optional[Dict] = None,
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:
pipeline_mode: str, vision_model: Optional[str],
text_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
@@ -143,7 +165,10 @@ def _run_pipeline(job_id: str, pdf_path: str, out_dir: str,
"email": _jobs[job_id].get("email"),
"pipeline_mode": pipeline_mode,
"source": _jobs[job_id].get("source"),
"vision_model": vision_model,
"text_model": text_model,
}, f, indent=2)
# Keep a copy of the source PDF so its sheets can be viewed later.
shutil.copy2(pdf_path, os.path.join(out_dir, "source.pdf"))
runner = run_agent_pipeline if pipeline_mode == "agent" else run_pipeline
runner_kwargs = {
@@ -153,17 +178,22 @@ def _run_pipeline(job_id: str, pdf_path: str, out_dir: str,
"source_name": _jobs[job_id].get("source"),
}
if pipeline_mode == "classic":
# run_pipeline takes the picks as params and clears them in finally.
runner_kwargs["text_local"] = text_local
runner_kwargs["vision_model"] = vision_model
runner_kwargs["text_model"] = text_model
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)
# The agent runner has no override params; set them module-level.
if vision_model or text_model:
print(f"[Jobs] Model overrides for this run: "
f"vision={vision_model or '(default)'} text={text_model or '(default)'}")
llm.set_model_overrides(vision_model, text_model)
try:
report = runner(pdf_path, **runner_kwargs)
finally:
if model:
llm.set_model_override(None)
if pipeline_mode == "agent":
llm.set_model_overrides(None, 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.
@@ -197,11 +227,6 @@ def _notify_error(job_id: str) -> None:
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
@@ -216,6 +241,7 @@ def _notify_error(job_id: str) -> None:
"Your conflict check did not complete.\n\n"
f"Drawing set: {job.get('source','')}\n"
f"Error: {job.get('error','unknown')}\n\n"
f"Review the run log at: {config.APP_BASE_URL.rstrip('/')}/?job={job_id}\n\n"
"Generated by Conflict Checker"
)
_send(msg)
@@ -223,28 +249,63 @@ def _notify_error(job_id: str) -> None:
print(f"[Email] Failed to send error notice: {e}")
def _log_from_disk(job_id: str) -> List[str]:
return read_log_file(os.path.join(config.OUTPUT_DIR, job_id, "job.log"))
def get_job_log(job_id: str) -> Optional[List[str]]:
"""Full job log lines, from memory or disk. None if job unknown."""
with _lock:
job = _jobs.get(job_id)
if job is not None:
return list(job.get("log") or [])
log = _log_from_disk(job_id)
# Job exists on disk if we have a log or a report artifact.
report_path = os.path.join(config.OUTPUT_DIR, job_id, "conflicts.json")
if log or os.path.isfile(report_path):
return log
return None
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).
Falls back to the on-disk artifacts (conflicts.json / job.log) 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)
out = dict(job)
log = list(job.get("log") or [])
out["log_tail"] = log[-_LOG_TAIL:]
# Full log on terminal states so the UI can show it without a
# second fetch; keep polls light while running.
if out.get("status") in _TERMINAL_STATES:
out["log"] = log
else:
out.pop("log", None)
return out
# Try loading from disk
report_path = os.path.join(config.OUTPUT_DIR, job_id, "conflicts.json")
if not os.path.isfile(report_path):
log = _log_from_disk(job_id)
if not os.path.isfile(report_path) and not log:
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"
report = None
if os.path.isfile(report_path):
with open(report_path, encoding="utf-8") as f:
report = json.load(f)
summary = (report or {}).get("summary", {})
if report is None:
# Crashed before writing a report; the log is the only artifact.
status = "error"
else:
# 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.
@@ -261,16 +322,20 @@ def get_job(job_id: str) -> Optional[Dict]:
job = {
"job_id": job_id,
"status": status,
"source": meta.get("source") or report.get("source", os.path.basename(report_path)),
"source": meta.get("source") or (report or {}).get("source", os.path.basename(report_path)),
"email": meta.get("email"),
"project_input": report.get("project_input", {}),
"project_input": (report or {}).get("project_input", {}),
"text_local": summary.get("text_backend") == "local",
"pipeline_mode": meta.get("pipeline_mode") or summary.get("pipeline_mode", "classic"),
"vision_model": meta.get("vision_model"),
"text_model": meta.get("text_model"),
"stage": None,
"created_at": os.path.getmtime(source_pdf) if os.path.isfile(source_pdf) else None,
"finished_at": os.path.getmtime(report_path),
"finished_at": os.path.getmtime(report_path) if os.path.isfile(report_path) else None,
"report": report,
"error": None,
"error": None if report is not None else "Report missing; see job log",
"log": log,
"log_tail": log[-_LOG_TAIL:],
}
# Hydrate the in-memory registry so _set(...) transitions (reviewing,
# finalizing, done) work for restart-recovered jobs.