Neustart-Befehle von Directus abarbeiten + Logs/Fehler nach Directus

- 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 <noreply@anthropic.com>
This commit is contained in:
s4luorth
2026-09-11 19:58:02 +02:00
parent dfaf5ccc58
commit 7ee7745ac0
2 changed files with 159 additions and 154 deletions

307
agent.py
View File

@@ -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):