""" Skrift – schlanker Produktions-Agent ==================================== Läuft lokal (Windows) neben dem Plotter. Er macht bewusst nur drei Dinge: 1. Heartbeat – meldet den Maschinenzustand an Directus (machines.status/last_seen) 2. Poll – holt den nächsten Druck-Job seiner Maschine (jobs, status=queued) 3. Drucken – lädt die SVGs vom Directus-Datei-Endpoint, schickt sie an den Plotter und meldet den Status zurück (printing → printed/failed) Alles andere (Aufträge, Warteschlange, Templates, Maschinenpflege) steuerst du in Directus / der Webapp. Der Agent kennt nur seine eine `machine`-ID. Konfiguration: config.json (siehe config.example.json). Start: python agent.py """ import io import json import os import sys import time import tempfile import subprocess import datetime as dt import requests HIER = os.path.dirname(os.path.abspath(__file__)) def lade_config(): pfad = os.path.join(HIER, "config.json") if not os.path.exists(pfad): print("config.json fehlt – bitte config.example.json kopieren und ausfüllen.") sys.exit(1) with open(pfad, "r", encoding="utf-8") as f: return json.load(f) class Directus: """Minimaler Directus-Client (Service-Token).""" def __init__(self, base_url, token): self.base = base_url.rstrip("/") self.s = requests.Session() self.s.headers["Authorization"] = f"Bearer {token}" def get(self, path, **params): """Für /items/… – Directus verpackt die Nutzdaten in `data`.""" r = self.s.get(self.base + path, params=params, timeout=30) r.raise_for_status() return r.json().get("data") def get_raw(self, path, **params): """Für Custom-Endpunkte (z. B. /skrift-orders/files) ohne `data`-Hülle.""" r = self.s.get(self.base + path, params=params, timeout=30) r.raise_for_status() return r.json() def patch(self, path, body): r = self.s.patch(self.base + path, json=body, timeout=30) r.raise_for_status() return r.json().get("data") def post(self, path, body): r = self.s.post(self.base + path, json=body, timeout=30) r.raise_for_status() return r.json().get("data") def download(self, path, params): r = self.s.get(self.base + path, params=params, timeout=120) r.raise_for_status() return r.content class Plotter: """Sendet SVGs an die Maschinen-API (Auszug aus der bisherigen App).""" def __init__(self, machine): self.base = str(machine.get("base_url", "")).rstrip("/") self.mid = machine.get("mid") or "" self.user = machine.get("username") or "" self.passwd = machine.get("password") or "" self.s = requests.Session() def login(self): if not self.user: return True r = self.s.get(f"{self.base}/login", params={"uin": self.user, "passwd": self.passwd}, timeout=10) return r.json().get("error", -1) == 0 def status(self): params = {"mids": self.mid} if self.mid else {} try: r = self.s.get(f"{self.base}/machine_status", params=params, timeout=10) data = r.json() return data[0] if isinstance(data, list) and data else None except Exception: return None def recover(self): params = {"mid": self.mid} if self.mid else {} try: r = self.s.get(f"{self.base}/recover_write", params=params, timeout=10) return r.json().get("code", -1) == 0 except Exception: return False @staticmethod def _fmt(v): v = float(v or 0) return str(int(v)) if v == int(v) else str(v) def write(self, dateien, tpl): """dateien: [(name, bytes)]; tpl: Template-Dict mit width/height/... (Zoll).""" basis = { "width": f"{self._fmt(tpl.get('width'))}in", "height": f"{self._fmt(tpl.get('height'))}in", "xpos": f"{self._fmt(tpl.get('xpos'))}in", "ypos": f"{self._fmt(tpl.get('ypos'))}in", "rotation": self._fmt(tpl.get("rotation")), "scale": self._fmt(tpl.get("scale") if tpl.get("scale") is not None else 1), "repeat": "1", } if self.mid: basis["mid"] = self.mid letzte = len(dateien) - 1 antwort = {} for i, (name, inhalt) in enumerate(dateien): data = {**basis, "clear": "1" if i == 0 else "0", "start": "1" if i == letzte else "0"} files = [("file", (name, io.BytesIO(inhalt), "image/svg+xml"))] r = self.s.post(f"{self.base}/write_svg", data=data, files=files, timeout=60) antwort = r.json() if antwort.get("error", -1) != 0: raise RuntimeError(f"Maschine meldet Fehler {antwort.get('error')} bei {name}") return antwort def jetzt(): return dt.datetime.now(dt.timezone.utc).isoformat() def parse_nummern(spec): """„1-3,5,7-8" → {1,2,3,5,7,8}. Leer/None → None (= alle drucken).""" if not spec or not str(spec).strip(): return None treffer = set() for teil in str(spec).split(","): teil = teil.strip() if not teil: continue if "-" in teil: try: a, b = teil.split("-", 1) a, b = int(a), int(b) treffer.update(range(min(a, b), max(a, b) + 1)) except ValueError: continue else: try: treffer.add(int(teil)) except ValueError: continue return treffer or None def seiten_reihenfolge(files): """Seiten je Brief ABSTEIGEND (…p3, p2, p1), Briefe aufsteigend. Der Plotter stapelt face-up: zuletzt gedruckt liegt oben. Damit am Ende Seite 1 obenauf liegt, muss je Brief die letzte Seite zuerst gedruckt werden. Ergebnis: 1.3, 1.2, 1.1, 2.3, 2.2, 2.1, … Einseitige Briefe (ohne _pN) und andere Typen (Kuvert/Signatur) bleiben unverändert. """ def key(f): name = str(f.get("name", "")).lower() brief, seite = 0, 1 if name.startswith("letter_"): rest = name[7:] if rest.endswith(".svg"): rest = rest[:-4] if "_p" in rest: a, _, b = rest.partition("_p") brief = int(a) if a.isdigit() else 0 seite = int(b) if b.isdigit() else 1 elif rest.isdigit(): brief = int(rest) return (brief, -seite) return sorted(files, key=key) def main(): cfg = lade_config() dx = Directus(cfg["directus_url"], cfg["directus_token"]) machine_id = cfg["machine_id"] poll = int(cfg.get("poll_interval_seconds", 5)) log(dx, machine_id, "info", "agent", f"Agent gestartet für Maschine {machine_id} gegen {cfg['directus_url']}") # Plotter-Fehler nur EINMAL loggen (nicht bei jedem Poll erneut) – erst wenn er # wieder verschwindet, ist der nächste Fehler wieder log-würdig. fehler_gemeldet = {"an": False} KIND = {"brief": "schriftstueck", "umschlag": "umschlag", "signatur": "signatur"} # Sicherheitsfenster: Sehen wir nach dem Senden binnen dieser Zeit kein # „Printing" (Job war sehr kurz), gilt er als fertig – verhindert Deadlock. start_grace = max(60, poll * 6) # Kontext des aktuell an die Maschine gesendeten Jobs (oder None). Der # nächste Job wird ERST gesendet, wenn die Maschine wieder „Idle" ist – # sonst würde ein neuer write_svg (clear=1) den laufenden Druck verwerfen. laufend = None # Heartbeat gedrosselt: nur schreiben, wenn sich der Status ändert oder # spätestens alle 60 s – sonst flutet der PATCH bei jedem Poll das Log. beat_state = {"status": None, "ts": 0.0} def beat(status): now = time.monotonic() if status != beat_state["status"] or now - beat_state["ts"] >= 60: _heartbeat(dx, machine_id, status) beat_state["status"] = status beat_state["ts"] = now while True: try: # Offene Befehle (Programm-Neustart) abarbeiten – auch wenn die Maschine # gerade offline/inaktiv ist. pruefe_befehle(dx, machine_id, cfg) maschine = dx.get(f"/items/machines/{machine_id}", fields="base_url,mid,username,password,active,current_format") if not maschine or not maschine.get("active", True): beat("offline") time.sleep(poll) continue plotter = Plotter(maschine) st = plotter.status() or {} s = str(st.get("status", "")).strip().lower() # Fehler (z. B. Papierende): fortsetzen, aber nichts Neues senden. if s in ("error", "fehler"): plotter.recover() if not fehler_gemeldet["an"]: log(dx, machine_id, "error", "plotter", f"Plotter meldet Fehler ({st.get('status')}) – Recover ausgelöst.") fehler_gemeldet["an"] = True beat("fehler") time.sleep(poll) continue fehler_gemeldet["an"] = False # Maschine druckt noch → warten (auf keinen Fall einen neuen Job senden). if s == "printing": if laufend is not None: laufend["gestartet"] = True beat("druckt") time.sleep(poll) continue # Nur bei sicherem „Idle" handeln. Unbekannt/leer → abwarten. if s != "idle": beat("bereit" if s else "offline") time.sleep(poll) continue beat("bereit") # Ein zuvor gesendeter Job ist fertig gedruckt (Idle nach Printing) → # abschließen. Das start_grace fängt sehr kurze Jobs ab, bei denen wir # das „Printing" nie zu Gesicht bekommen. if laufend is not None: if laufend.get("gestartet") or (time.monotonic() - laufend["ts"] > start_grace): _abschluss(dx, laufend) laufend = None else: time.sleep(poll) continue # Nächsten Job nur senden, wenn nichts (mehr) läuft. if laufend is None: current_format = (maschine.get("current_format") or "").strip() # Format-Warteschlange: nur passendes Format, kein Auto-Anschluss. if current_format: jobs = _naechste_jobs(dx, machine_id, current_format) if jobs: laufend = _sende_batch(dx, plotter, jobs, KIND) time.sleep(poll) except requests.HTTPError as e: log(dx, machine_id, "error", "agent", f"HTTP-Fehler: {e}") time.sleep(poll) except Exception as e: log(dx, machine_id, "error", "agent", f"Fehler: {e}") time.sleep(poll) def log(dx, machine_id, level, source, message): """Loggt lokal UND nach Directus (agent_logs). Sende-Fehler nie fatal.""" print(f"[{source}:{level}] {message}") try: dx.post("/items/agent_logs", {"machine": machine_id, "level": level, "source": source, "message": str(message)[:2000]}) except Exception: pass def restart_program(cfg, target): """Startet ein Windows-Programm neu: erst beenden (taskkill), dann starten. Config-Form: "restart": { "process": "plotter.exe", "start": "C:/Pfad/plotter.exe" } `target` (aus dem Befehl) überschreibt: sieht es nach einer .exe mit Pfad aus, gilt es als Startbefehl, sonst als zu beendender Prozessname.""" rc = cfg.get("restart") or {} process = rc.get("process") start = rc.get("start") if target: t = str(target).strip() if t.lower().endswith(".exe") and ("\\" in t or "/" in t): start = t process = process or os.path.basename(t) else: process = t if not process and not start: raise RuntimeError("Kein Ziel: weder cfg['restart'] noch Befehl-target gesetzt.") schritte = [] if process: try: out = subprocess.run(["taskkill", "/IM", process, "/F"], capture_output=True, text=True, timeout=30) schritte.append(f"taskkill {process}: rc={out.returncode} {(out.stdout or out.stderr).strip()}") except Exception as e: schritte.append(f"taskkill {process} Fehler: {e}") time.sleep(1.5) if start: cwd = os.path.dirname(start) if (os.path.sep in start and start.lower().endswith(".exe")) else None subprocess.Popen(start, shell=True, cwd=cwd) schritte.append(f"start: {start}") return " | ".join(schritte) def pruefe_befehle(dx, machine_id, cfg): """Arbeitet offene agent_commands (status=pending) dieser Maschine ab.""" try: cmds = dx.get("/items/agent_commands", **{ "filter[status][_eq]": "pending", "filter[machine][_eq]": machine_id, "sort": "date_created", "limit": 5, "fields": "id,type,target", }) or [] except Exception: return for c in cmds: cid = c.get("id") if c.get("type") != "restart": continue try: ergebnis = restart_program(cfg, c.get("target")) dx.patch(f"/items/agent_commands/{cid}", {"status": "done", "result": ergebnis[:2000]}) log(dx, machine_id, "info", "restart", f"Neustart ausgeführt (#{cid}): {ergebnis}") except Exception as e: dx.patch(f"/items/agent_commands/{cid}", {"status": "error", "result": str(e)[:2000]}) log(dx, machine_id, "error", "restart", f"Neustart fehlgeschlagen (#{cid}): {e}") def _heartbeat(dx, machine_id, status): try: dx.patch(f"/items/machines/{machine_id}", {"status": status, "last_seen": jetzt()}) except Exception: pass def _naechste_jobs(dx, machine_id, current_format): """ALLE eingereihten Jobs des aktuell eingelegten Formats – in stabiler Warteschlangen-Reihenfolge (Priorität, dann Erstellzeit, dann ID als eindeutiger Tiebreak). Diese Reihenfolge = Erstell-/Auftragsnummern-Reihenfolge (enqueue-bulk legt die Jobs nach Auftragsnummer an) und MUSS exakt erhalten bleiben, damit Schriftstücke, Kuverts und die Sammel-PDF zueinander passen.""" return dx.get( "/items/jobs", **{ "filter[machine][_eq]": machine_id, "filter[status][_eq]": "queued", "filter[format][_eq]": current_format, "sort": "-priority,date_created,id", "limit": -1, "fields": "id,type,output_ref,params,numbers,spaltenweise," "order.id,order.order_number,order.needs_envelope," "template.width,template.height," "template.xpos,template.ypos,template.scale,template.rotation", }, ) or [] def _eff_tpl(job): """Effektives Maschinen-Template: Template-Werte + Pro-Druck-Übersteuerung.""" tpl = dict(job.get("template") or {}) params = job.get("params") if isinstance(params, dict): for k in ("width", "height", "xpos", "ypos", "scale", "rotation"): if params.get(k) is not None: tpl[k] = params[k] return tpl def _batch_key(job): """Jobs mit gleichem Typ UND gleicher Platzierung dürfen zusammen in EINEN write_svg-Batch (die Maschine kennt nur eine Platzierung je Sendung).""" return (job.get("type") or "brief", json.dumps(_eff_tpl(job), sort_keys=True, default=str)) def _job_dateien(dx, job, KIND): """Lädt die SVG-Dateien EINES Jobs (Typ-passend, Nummern-Filter, Seiten- Reihenfolge, ggf. spaltenweise). Gibt [(name, bytes)] zurück; wirft bei Fehler.""" order = job.get("order") or {} order_nr = job.get("output_ref") or order.get("order_number") typ = job.get("type") or "brief" if not order_nr: raise RuntimeError("Auftragsnummer fehlt am Job.") liste = dx.get_raw(f"/skrift-orders/files/{order_nr}") or {} gewuenscht = KIND.get(typ, "schriftstueck") passende = sorted( [f for f in (liste.get("files") or []) if f.get("kind") == gewuenscht], key=lambda f: f.get("name", ""), ) # Nur bestimmte Briefnummern (1-basiert, positionsbasiert), falls angegeben. nummern = parse_nummern(job.get("numbers")) if nummern is not None: passende = [f for i, f in enumerate(passende, start=1) if i in nummern] if not passende: raise RuntimeError(f"Keine Dateien vom Typ '{gewuenscht}' für {order_nr} gefunden.") # Mehrseitige Briefe: Seiten je Brief absteigend (Stapel-Reihenfolge). passende = seiten_reihenfolge(passende) # Rollenschonende Zeichenreihenfolge (spaltenweise) nur fürs Schriftstück. reorder = bool(job.get("spaltenweise")) and gewuenscht == "schriftstueck" dateien = [] for f in passende: p = {"rel": f["rel"]} if reorder: p["reorder"] = "column" dateien.append((f["name"], dx.download(f"/skrift-orders/files/{order_nr}/download", p))) return dateien def _sende_batch(dx, plotter, jobs, KIND): """Sendet ALLE Jobs mit gleichem Typ+Platzierung wie der erste als EINEN zusammenhängenden Druck an die Maschine (clear beim ersten, start beim letzten Dokument). Markiert alle als „printing". Gibt einen Sammel-Kontext zum späteren Abschluss zurück (oder None, wenn nichts gesendet wurde).""" ref = _batch_key(jobs[0]) batch = [j for j in jobs if _batch_key(j) == ref] tpl = _eff_tpl(batch[0]) dateien = [] infos = [] for job in batch: jid = job["id"] try: teil = _job_dateien(dx, job, KIND) except Exception as e: dx.patch(f"/items/jobs/{jid}", {"status": "failed", "error": str(e)[:1000]}) print(f"[Agent] Job {jid} übersprungen: {e}") continue order = job.get("order") or {} oid = order.get("id") dx.patch(f"/items/jobs/{jid}", {"status": "printing", "error": None}) if oid: try: dx.patch(f"/items/orders/{oid}", {"production_status": "im_druck"}) except Exception: pass dateien.extend(teil) infos.append({"job_id": jid, "order_id": oid, "order_nr": job.get("output_ref") or order.get("order_number"), "typ": job.get("type") or "brief"}) if not dateien: return None try: plotter.login() plotter.write(dateien, tpl) except Exception as e: for info in infos: dx.patch(f"/items/jobs/{info['job_id']}", {"status": "failed", "error": str(e)[:1000]}) print(f"[Agent] Batch fehlgeschlagen: {e}") return None print(f"[Agent] Batch: {len(infos)} Job(s) / {len(dateien)} Dateien in einem Zug an die Maschine – warte auf Fertigstellung.") return {"jobs": infos, "gestartet": False, "ts": time.monotonic()} def _abschluss(dx, ctx): """Schließt ALLE Jobs eines fertig gedruckten Batches ab: Status, Indikatoren, Kuvertierung. Wird erst aufgerufen, wenn die Maschine nach dem Druck wieder „Idle" meldet (der ganze Batch war EINE Sendung).""" for info in ctx.get("jobs", []): job_id = info["job_id"] order_id = info.get("order_id") order_nr = info.get("order_nr") typ = info.get("typ") try: dx.patch(f"/items/jobs/{job_id}", {"status": "printed"}) except Exception as e: print(f"[Agent] Konnte Job {job_id} nicht auf 'printed' setzen: {e}") # Indikator am Auftrag setzen (Schriftstück bzw. Umschlag gedruckt) und # bei Vollständigkeit automatisch in die Kuvertierung schieben. if order_id: feld = "kuvert_gedruckt" if typ == "umschlag" else "brief_gedruckt" try: dx.patch(f"/items/orders/{order_id}", {feld: True}) stand = dx.get(f"/items/orders/{order_id}", fields="brief_gedruckt,kuvert_gedruckt,needs_envelope,production_status") or {} braucht_kuvert = bool(stand.get("needs_envelope")) fertig = stand.get("brief_gedruckt") and (stand.get("kuvert_gedruckt") or not braucht_kuvert) if fertig and stand.get("production_status") == "im_druck": dx.patch(f"/items/orders/{order_id}", {"production_status": "kuvertieren"}) print(f"[Agent] Auftrag {order_nr} vollständig gedruckt → Kuvertierung.") except Exception: pass print(f"[Agent] Job {job_id} fertig gedruckt.") if __name__ == "__main__": main()