From 7ee7745ac0f98ce43957c8c7e18770f4f1e722b8 Mon Sep 17 00:00:00 2001 From: s4luorth Date: Fri, 11 Sep 2026 19:58:02 +0200 Subject: [PATCH] Neustart-Befehle von Directus abarbeiten + Logs/Fehler nach Directus MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Directus.post; agent_commands (type=restart) pollen und Programm neu starten (taskkill /IM + start), Ergebnis + Fehler zurückschreiben. - Agent-Log & Plotter-Fehler (throttled) an agent_logs melden. - config.example.json: restart-Block (process/start). Co-Authored-By: Claude Opus 4.8 --- agent.py | 307 ++++++++++++++++++++++---------------------- config.example.json | 6 +- 2 files changed, 159 insertions(+), 154 deletions(-) diff --git a/agent.py b/agent.py index 60cb6ac..645f77b 100644 --- a/agent.py +++ b/agent.py @@ -21,6 +21,7 @@ import os import sys import time import tempfile +import subprocess import datetime as dt import requests @@ -62,6 +63,11 @@ class Directus: 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() @@ -184,40 +190,25 @@ def seiten_reihenfolge(files): return sorted(files, key=key) -def brief_nummer(name): - """Brief-Nummer aus dem Dateinamen: letter_001_p2.svg → 1, envelope_003.svg → 3. - - None, wenn nicht parsebar. Die Seite (_pN) wird ignoriert – alle Seiten eines - Briefs teilen sich dieselbe Nummer. - """ - rest = str(name or "").lower() - for pre in ("letter_", "envelope_", "kuvert_"): - if rest.startswith(pre): - rest = rest[len(pre):] - break - rest = rest.split(".")[0] - if "_p" in rest: - rest = rest.split("_p")[0] - return int(rest) if rest.isdigit() else None - - 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)) - print(f"[Agent] gestartet für Maschine {machine_id} gegen {cfg['directus_url']}") + 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 Batches (oder None): - # {"contexts": [ctx, …], "gestartet": bool, "ts": monotonic}. Ein Batch fasst - # ALLE queued-Jobs gleichen Typs+Templates zu EINEM Durchlauf zusammen. Der - # nächste Batch wird ERST gesendet, wenn die Maschine wieder „Idle" ist – + # 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 @@ -234,6 +225,10 @@ def main(): 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): @@ -248,9 +243,14 @@ def main(): # 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": @@ -273,35 +273,97 @@ def main(): # das „Printing" nie zu Gesicht bekommen. if laufend is not None: if laufend.get("gestartet") or (time.monotonic() - laufend["ts"] > start_grace): - for ctx in laufend["contexts"]: - _abschluss(dx, ctx) + _abschluss(dx, laufend) laufend = None else: time.sleep(poll) continue - # Nichts läuft → alle queued-Jobs des aktuellen Formats holen und ALLE - # vom selben Typ (gleiches Template) als EINEN Durchlauf drucken. Danach - # holt die nächste Runde automatisch den Rest desselben Typs (falls noch - # queued) bzw. den nächsten Typ. + # 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, _batch_gleicher_typ(jobs), KIND) + job = _naechster_job(dx, machine_id, current_format) + if job: + laufend = _sende(dx, plotter, job, KIND) time.sleep(poll) except requests.HTTPError as e: - print(f"[Agent] HTTP-Fehler: {e}") + log(dx, machine_id, "error", "agent", f"HTTP-Fehler: {e}") time.sleep(poll) except Exception as e: - print(f"[Agent] Fehler: {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()}) @@ -309,10 +371,8 @@ def _heartbeat(dx, machine_id, status): pass -def _naechste_jobs(dx, machine_id, current_format): - """ALLE queued-Jobs der Maschine im aktuellen Format (sortiert nach Priorität, - dann Alter). Die Auswahl fürs gemeinsame Drucken trifft _batch_gleicher_typ.""" - return dx.get( +def _naechster_job(dx, machine_id, current_format): + jobs = dx.get( "/items/jobs", **{ "filter[machine][_eq]": machine_id, @@ -320,139 +380,80 @@ def _naechste_jobs(dx, machine_id, current_format): # Nur Jobs des aktuell eingelegten Formats. "filter[format][_eq]": current_format, "sort": "-priority,date_created", - # Praktisch „alle" – nur eine Obergrenze, damit die Abfrage begrenzt bleibt. - "limit": 500, + "limit": 1, "fields": "id,type,output_ref,params,numbers," "order.id,order.order_number,order.needs_envelope," "template.width,template.height," "template.xpos,template.ypos,template.scale,template.rotation", }, - ) or [] + ) + return jobs[0] if jobs else None -def _effektives_template(job): - """Template des Jobs inkl. Pro-Druck-Übersteuerung (params) aus dem Popup.""" +def _sende(dx, plotter, job, KIND): + """Sendet einen Job an die Maschine und markiert ihn als „printing". + + Gibt einen Kontext zum späteren Abschluss zurück (oder None bei Fehler). + Es wird NICHT auf das Druckende gewartet – das übernimmt die Hauptschleife + (erst wenn die Maschine wieder „Idle" ist, wird der nächste Job gesendet). + """ + job_id = job["id"] + order = job.get("order") or {} + order_id = order.get("id") + # Ordner: entweder Auftragsnummer oder ein output_ref (z. B. Testdruck). + order_nr = job.get("output_ref") or order.get("order_number") tpl = dict(job.get("template") or {}) + # Pro-Druck-Übersteuerung der Maße (aus dem Druck-Popup) hat Vorrang. 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 _template_sig(job): - t = _effektives_template(job) - return tuple(t.get(k) for k in ("width", "height", "xpos", "ypos", "scale", "rotation")) - - -def _batch_gleicher_typ(jobs): - """Aus der sortierten queued-Liste den Kopf-Job nehmen und ALLE queued-Jobs - gleichen Typs UND gleichen Templates dazu sammeln – ein Format hat i. d. R. - genau ein Template, also faktisch „alles von einem Typ auf einmal". Jobs mit - abweichendem Template laufen im nächsten Durchgang. Reihenfolge (Priorität/ - Alter) bleibt erhalten.""" - kopf = jobs[0] - typ = kopf.get("type") or "brief" - sig = _template_sig(kopf) - return [j for j in jobs - if (j.get("type") or "brief") == typ and _template_sig(j) == sig] - - -def _job_dateien(dx, job, KIND): - """Lädt die zu druckenden SVGs eines Jobs in Stapel-Reihenfolge → [(name, bytes)].""" - order = job.get("order") or {} - # Ordner: entweder Auftragsnummer oder ein output_ref (z. B. Testdruck). - order_nr = job.get("output_ref") or order.get("order_number") - if not order_nr: - raise RuntimeError("Auftragsnummer fehlt am Job.") - typ = job.get("type") or "brief" - 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 drucken, falls im Job angegeben. - # Positionsbasiert über DISTINKTE Brief-Nummern (nicht über die Dateiposition!), - # damit mehrseitige Briefe komplett bleiben: 24 zweiseitige Briefe = 48 Dateien, - # numbers="1-24" wählt alle 24 Briefe → alle 48 Seiten. Robust gegen 0-/1-basierte - # Namen. Fallback auf die alte Dateiposition, falls keine Nummer parsebar ist. nummern = parse_nummern(job.get("numbers")) - if nummern is not None: - distinkt = sorted({n for n in (brief_nummer(f.get("name")) for f in passende) if n is not None}) - if distinkt: - erlaubt = {distinkt[i - 1] for i in nummern if 1 <= i <= len(distinkt)} - passende = [f for f in passende if brief_nummer(f.get("name")) in erlaubt] - else: - 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.") + typ = job.get("type") or "brief" + print(f"[Agent] Job {job_id}: {typ} für {order_nr} → an Maschine senden") - # Mehrseitige Briefe: Seiten je Brief absteigend an den Plotter (Stapel- - # Reihenfolge, damit Seite 1 am Ende oben liegt). Einseitige unverändert. - passende = seiten_reihenfolge(passende) - dateien = [] - for f in passende: - inhalt = dx.download(f"/skrift-orders/files/{order_nr}/download", {"rel": f["rel"]}) - dateien.append((f["name"], inhalt)) - return dateien - - -def _sende_batch(dx, plotter, jobs, KIND): - """Sendet mehrere gleichartige Jobs als EINEN Druckdurchlauf an die Maschine. - - Alle Jobs teilen Typ und Template (siehe _batch_gleicher_typ). Ihre Dokumente - werden nacheinander in EINEM Lauf gedruckt (clear nur am Anfang, start am Ende) – - so hält die Maschine zwischen den Aufträgen nicht an. Ein einzelner Job, dessen - Dateien fehlen/fehlerhaft sind, wird als „failed" markiert und übersprungen; der - Rest wird trotzdem gedruckt. Gibt den Batch-Kontext zurück (oder None, wenn nichts - Druckbares übrig bleibt). Es wird NICHT auf das Druckende gewartet – das übernimmt - die Hauptschleife (nächster Batch erst bei „Idle"). - """ - contexts = [] - dateien = [] - tpl_final = None - for job in jobs: - job_id = job["id"] - order = job.get("order") or {} - order_id = order.get("id") - order_nr = job.get("output_ref") or order.get("order_number") - typ = job.get("type") or "brief" - dx.patch(f"/items/jobs/{job_id}", {"status": "printing", "error": None}) - # Produktionsstatus fürs Kanban: sobald gedruckt wird → „im_druck". - if order_id: - try: dx.patch(f"/items/orders/{order_id}", {"production_status": "im_druck"}) - except Exception: pass - try: - job_dateien = _job_dateien(dx, job, KIND) - except Exception as e: - dx.patch(f"/items/jobs/{job_id}", {"status": "failed", "error": str(e)[:1000]}) - print(f"[Agent] Job {job_id} fehlgeschlagen: {e}") - continue - if tpl_final is None: - tpl_final = _effektives_template(job) - dateien.extend(job_dateien) - contexts.append({"job_id": job_id, "order_id": order_id, - "order_nr": order_nr, "typ": typ}) - print(f"[Agent] Job {job_id}: {typ} für {order_nr} ({len(job_dateien)} Dateien) → Batch") - - if not dateien: - return None + dx.patch(f"/items/jobs/{job_id}", {"status": "printing", "error": None}) + # Produktionsstatus fürs Kanban: sobald gedruckt wird → „im_druck". + if order_id: + try: dx.patch(f"/items/orders/{order_id}", {"production_status": "im_druck"}) + except Exception: pass try: - plotter.login() - plotter.write(dateien, tpl_final) - except Exception as e: - # Ganzer Durchlauf fehlgeschlagen → alle beteiligten Jobs auf „failed". - for ctx in contexts: - try: dx.patch(f"/items/jobs/{ctx['job_id']}", {"status": "failed", "error": str(e)[:1000]}) - except Exception: pass - print(f"[Agent] Batch-Druck fehlgeschlagen: {e}") - return None + 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 drucken, falls im Job angegeben. + # Positionsbasiert (1 = erstes Dokument), NICHT über die Zahl im + # Dateinamen: Alt-Aufträge sind 0-basiert (letter_000.svg), neue + # 1-basiert (letter_001.svg). Die Position stimmt in beiden Fällen. + 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.") - print(f"[Agent] Batch an Maschine übergeben: {len(contexts)} Job(s), " - f"{len(dateien)} Dateien – warte auf Fertigstellung.") - return {"contexts": contexts, "gestartet": False, "ts": time.monotonic()} + # Mehrseitige Briefe: Seiten je Brief absteigend an den Plotter (Stapel- + # Reihenfolge, damit Seite 1 am Ende oben liegt). Einseitige unverändert. + passende = seiten_reihenfolge(passende) + + dateien = [] + for f in passende: + inhalt = dx.download(f"/skrift-orders/files/{order_nr}/download", {"rel": f["rel"]}) + dateien.append((f["name"], inhalt)) + + plotter.login() + plotter.write(dateien, tpl) + print(f"[Agent] Job {job_id} an Maschine übergeben ({len(dateien)} Dateien) – warte auf Fertigstellung.") + return {"job_id": job_id, "order_id": order_id, "order_nr": order_nr, + "typ": typ, "gestartet": False, "ts": time.monotonic()} + except Exception as e: + dx.patch(f"/items/jobs/{job_id}", {"status": "failed", "error": str(e)[:1000]}) + print(f"[Agent] Job {job_id} fehlgeschlagen: {e}") + return None def _abschluss(dx, ctx): diff --git a/config.example.json b/config.example.json index e812591..a8d56e9 100644 --- a/config.example.json +++ b/config.example.json @@ -2,5 +2,9 @@ "directus_url": "https://admin.skrift.de", "directus_token": "SERVICE-TOKEN-HIER-EINTRAGEN", "machine_id": 1, - "poll_interval_seconds": 5 + "poll_interval_seconds": 5, + "restart": { + "process": "PROGRAMMNAME.exe", + "start": "C:/Pfad/zum/PROGRAMM.exe" + } }