"""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}} for uid in box.unchecked_uids(): msg = box.fetch(uid) subject = decode_subject(msg) html, text = mailparse.bodies(msg) 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 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 box.move_to_trash(uid) box.close() return summary finally: release_lock(lock)