Root cause chain (evidence: 25+ failed traces on nproj/3022414): straight quote in page text ends the guided-decoding JSON string early -> grammar deadlock -> whitespace loop until context/token limit -> chat_json fails 3x (~35 min) -> IMAP conn idles out -> move_to_trash raises -> mail never trashed -> reprocessed forever. - ingest/mailer: MailBox.ensure_alive() reconnect before finalize; move_to_trash failure falls back to flag_checked, run continues - stages: page_text() feeds only div.content minus skill-tag cloud to the LLM and maps quote chars to apostrophes (the actual degeneration trigger, live-verified: extract now OK in 21 s) - llm: requirements maxItems 40; max_tokens 8192 caps degenerate runs at ~100 s instead of ~10 min Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
154 lines
6.0 KiB
Python
154 lines
6.0 KiB
Python
"""Flow 1: poll inbox, split projects, dedup, dispatch Flow 2, alert, trash."""
|
|
from __future__ import annotations
|
|
|
|
import email.header
|
|
import fcntl
|
|
import json
|
|
import os
|
|
|
|
import requests
|
|
|
|
from . import espocrm, mailer, mailparse
|
|
|
|
SUMMARY_KEYS = ("created", "rejected", "duplicate", "failed")
|
|
LOCK_NAME = "projekt-matching.lock"
|
|
|
|
|
|
def acquire_lock(data_dir):
|
|
path = os.path.join(data_dir, LOCK_NAME)
|
|
handle = open(path, "w") # noqa: SIM115
|
|
try:
|
|
fcntl.flock(handle, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
|
return handle
|
|
except OSError:
|
|
handle.close()
|
|
return None
|
|
|
|
|
|
def release_lock(handle):
|
|
fcntl.flock(handle, fcntl.LOCK_UN)
|
|
handle.close()
|
|
|
|
|
|
def decode_subject(msg):
|
|
raw = msg["Subject"] or ""
|
|
try:
|
|
parts = email.header.decode_header(raw)
|
|
return "".join(p.decode(enc or "utf-8", "replace")
|
|
if isinstance(p, bytes) else p for p, enc in parts)
|
|
except Exception:
|
|
return raw
|
|
|
|
|
|
def dispatch_flow2(cfg, item):
|
|
resp = requests.post(
|
|
f"{cfg.langflow_base}/api/v1/run/{cfg.flow2_id}?stream=false",
|
|
headers={"x-api-key": cfg.langflow_api_key},
|
|
json={"input_value": json.dumps(
|
|
{"canonical": item["canonical"], "title": item["title"]},
|
|
ensure_ascii=False),
|
|
"input_type": "text", "output_type": "text"},
|
|
timeout=1800)
|
|
resp.raise_for_status()
|
|
return parse_run_result(resp.json())
|
|
|
|
|
|
def parse_run_result(payload):
|
|
"""Depth-first search for the stage summary JSON in the run response."""
|
|
stack = [payload]
|
|
while stack:
|
|
node = stack.pop()
|
|
if isinstance(node, dict):
|
|
stack.extend(node.values())
|
|
elif isinstance(node, list):
|
|
stack.extend(node)
|
|
elif isinstance(node, str) and node.lstrip().startswith("{"):
|
|
try:
|
|
data = json.loads(node)
|
|
except json.JSONDecodeError:
|
|
continue
|
|
if isinstance(data, dict) and "status" in data:
|
|
return data
|
|
raise ValueError("Kein Status-JSON in der Run-Antwort gefunden")
|
|
|
|
|
|
def run_ingest(cfg, mailbox=None, dispatch=None, espo=None, resolve=None):
|
|
lock = acquire_lock(cfg.data_dir)
|
|
if lock is None:
|
|
return {"skipped": "locked"}
|
|
try:
|
|
box = mailbox or mailer.MailBox(cfg.imap_user, cfg.imap_password)
|
|
crm = espo or espocrm.EspoClient(cfg.espo_base, cfg.espo_api_key)
|
|
send = dispatch or (lambda item: dispatch_flow2(cfg, item))
|
|
resolve = resolve or mailparse.canonical_url
|
|
summary = {"mails": 0, "projects": 0,
|
|
**{k: 0 for k in SUMMARY_KEYS}}
|
|
try:
|
|
for uid in box.unchecked_uids():
|
|
try:
|
|
msg = box.fetch(uid)
|
|
subject = decode_subject(msg)
|
|
html, text = mailparse.bodies(msg)
|
|
except Exception: # noqa: BLE001
|
|
box.flag_checked(uid) # nie wieder anfassen, Pipeline läuft weiter
|
|
summary["unreadable"] = summary.get("unreadable", 0) + 1
|
|
continue
|
|
items = [] if mailparse.is_own_mail(subject) else \
|
|
mailparse.split_projects(html, text)
|
|
if not items:
|
|
box.flag_checked(uid)
|
|
continue
|
|
summary["mails"] += 1
|
|
results = []
|
|
for item in items:
|
|
summary["projects"] += 1
|
|
result = dict(item)
|
|
try:
|
|
result["canonical"] = resolve(item["url"])
|
|
dup = crm.find_opportunity_by_link(result["canonical"])
|
|
if dup:
|
|
result["status"] = "duplicate"
|
|
else:
|
|
result.update(send(result))
|
|
except Exception as exc: # noqa: BLE001
|
|
result.update(status="failed", error=str(exc))
|
|
if result.get("status") not in SUMMARY_KEYS:
|
|
result.update(
|
|
status="failed",
|
|
error=result.get("error")
|
|
or f"unbekannter Status: {result.get('status')!r}")
|
|
results.append(result)
|
|
summary[result["status"]] += 1
|
|
try:
|
|
box.ensure_alive() # Flow-2-Dispatches können >30 min dauern — Server-Idle-Timeout
|
|
except Exception: # noqa: BLE001
|
|
pass # Alert/Trash unten haben eigene Fallbacks
|
|
failed = [r for r in results if r["status"] == "failed"]
|
|
if failed:
|
|
lines = "\n\n".join(
|
|
f"- {f.get('title', '?')}\n"
|
|
f" {f.get('canonical', f.get('url', '?'))}\n"
|
|
f" Fehler: {f.get('error', '?')}" for f in failed)
|
|
try:
|
|
mailer.send_mail(
|
|
cfg.imap_user, cfg.imap_password, cfg.alert_to,
|
|
f"[Projekt-Match-Fehler] Freelancermap — "
|
|
f"{len(failed)} Projekt(e)",
|
|
"Folgende Projekte konnten nicht verarbeitet "
|
|
"werden:\n\n" + lines + "\n")
|
|
except Exception: # noqa: BLE001
|
|
summary["alertFailed"] = summary.get("alertFailed", 0) + 1
|
|
try:
|
|
box.move_to_trash(uid)
|
|
except Exception: # noqa: BLE001
|
|
summary["trashFailed"] = summary.get("trashFailed", 0) + 1
|
|
try:
|
|
box.flag_checked(uid) # nie erneut verarbeiten — Terminal-Zustand
|
|
except Exception: # noqa: BLE001
|
|
pass # nächster Lauf räumt auf
|
|
return summary
|
|
finally:
|
|
box.close()
|
|
finally:
|
|
release_lock(lock)
|