468 lines
18 KiB
Python
468 lines
18 KiB
Python
"""
|
||
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 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 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))
|
||
|
||
print(f"[Agent] gestartet für Maschine {machine_id} gegen {cfg['directus_url']}")
|
||
|
||
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 –
|
||
# 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:
|
||
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()
|
||
beat("fehler")
|
||
time.sleep(poll)
|
||
continue
|
||
|
||
# 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):
|
||
for ctx in laufend["contexts"]:
|
||
_abschluss(dx, ctx)
|
||
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.
|
||
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)
|
||
|
||
time.sleep(poll)
|
||
|
||
except requests.HTTPError as e:
|
||
print(f"[Agent] HTTP-Fehler: {e}")
|
||
time.sleep(poll)
|
||
except Exception as e:
|
||
print(f"[Agent] Fehler: {e}")
|
||
time.sleep(poll)
|
||
|
||
|
||
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 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(
|
||
"/items/jobs",
|
||
**{
|
||
"filter[machine][_eq]": machine_id,
|
||
"filter[status][_eq]": "queued",
|
||
# 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,
|
||
"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 []
|
||
|
||
|
||
def _effektives_template(job):
|
||
"""Template des Jobs inkl. Pro-Druck-Übersteuerung (params) aus dem Popup."""
|
||
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 _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 (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.
|
||
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 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
|
||
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
|
||
|
||
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()}
|
||
|
||
|
||
def _abschluss(dx, ctx):
|
||
"""Schließt einen fertig gedruckten Job ab: Status, Indikatoren, Kuvertierung.
|
||
|
||
Wird erst aufgerufen, wenn die Maschine nach dem Druck wieder „Idle" meldet.
|
||
"""
|
||
job_id = ctx["job_id"]
|
||
order_id = ctx.get("order_id")
|
||
order_nr = ctx.get("order_nr")
|
||
typ = ctx.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()
|