Praxis-Telefon-Agent: Anrufannahme, Transkription, Triage — vollständig on-prem

Nimmt Anrufe einer Hausarztpraxis entgegen, transkribiert sie lokal,
kategorisiert das Anliegen und legt es strukturiert ab. Keine Patientendaten
verlassen den Rechner: Spracherkennung, Sprachmodell und Sprachsynthese laufen
lokal, im Datenpfad steht kein Cloud-Dienst.

Aufbau:
- FreeSWITCH nimmt über einen SIP-Trunk an, spielt die Ansage und nimmt auf.
- pipe.watch erkennt fertige Aufnahmen und schiebt sie durch die Pipe.
- whisper.cpp transkribiert, ein lokales Sprachmodell (Ollama) ordnet in
  Kategorie und Dringlichkeit ein — mit erzwungenem JSON-Schema.
- Ablage lokal, optional Nextcloud (WebDAV) und ein Deck-Board.
- Leitstand (pipe.server): Board mit Eingang/In Bearbeitung/Rückfragen/Erledigt,
  Aufnahmen zum Anhören, Verlauf und Ampeln für alle beteiligten Dienste.

Zugangsdaten und Anrufdaten liegen außerhalb des Repositorys (.env, ablage/,
telefon/). Konfiguration siehe .env.example, Einrichtung siehe README.
This commit is contained in:
Jeuner 2026-08-22 22:37:17 +02:00
commit 1a1280e79d
24 changed files with 2253 additions and 0 deletions

0
pipe/__init__.py Normal file
View file

116
pipe/categorize.py Normal file
View file

@ -0,0 +1,116 @@
"""Kategorisierung: Transkript -> strukturierte Felder (Ollama/Qwen, lokal).
Nutzt den Chat-Endpunkt von Ollama mit erzwungenem JSON-Schema. Denk-Tokens des
Modells landen separat im "thinking"-Feld und werden ignoriert; zurück kommt
reines JSON.
"""
from __future__ import annotations
import json
import urllib.error
import urllib.request
from functools import lru_cache
from . import config
KATEGORIEN = [
"termin", "rezept", "ueberweisung", "befund",
"verwaltung", "beschwerden", "notfall", "rueckruf", "sonstiges",
]
DRINGLICHKEITEN = ["niedrig", "normal", "hoch", "notfall"]
# JSON-Schema, das Ollama dem Modell als Ausgabeformat aufzwingt.
_SCHEMA = {
"type": "object",
"properties": {
"kategorie": {"type": "string", "enum": KATEGORIEN},
"dringlichkeit": {"type": "string", "enum": DRINGLICHKEITEN},
"anrufer_name": {"type": ["string", "null"]},
"rueckrufnummer": {"type": ["string", "null"]},
"rueckruf_gewuenscht": {"type": "boolean"},
"sprache": {"type": "string"},
"anliegen_kurz": {"type": "string"},
"stichworte": {"type": "array", "items": {"type": "string"}},
},
# Alle Felder erzwingen: was hier fehlt, lässt das Modell einfach weg statt
# null zu liefern - Name und Rückrufnummer gingen so still verloren, obwohl
# sie im Transkript standen. Der Typ erlaubt weiterhin null, wenn nichts da ist.
"required": [
"kategorie", "dringlichkeit", "anrufer_name", "rueckrufnummer",
"rueckruf_gewuenscht", "sprache", "anliegen_kurz", "stichworte",
],
}
@lru_cache(maxsize=1)
def _prompt() -> str:
return config.PROMPT_DATEI.read_text(encoding="utf-8")
class KategorisierungsFehler(RuntimeError):
pass
def kategorisiere(transkript: str) -> dict:
"""Ordnet ein Transkript ein. Wirft KategorisierungsFehler bei Problemen."""
if not transkript.strip():
raise KategorisierungsFehler("Leeres Transkript kann nicht kategorisiert werden.")
anfrage = {
"model": config.OLLAMA_MODELL,
"messages": [
{"role": "system", "content": _prompt()},
{"role": "user", "content": f"Transkript der Sprachnachricht:\n\n{transkript}"},
],
"stream": False,
"think": False, # Klassifikation braucht kein Reasoning -> deutlich schneller
"format": _SCHEMA,
"options": {"temperature": 0},
}
daten = json.dumps(anfrage).encode("utf-8")
req = urllib.request.Request(
f"{config.OLLAMA_URL}/api/chat",
data=daten,
headers={"Content-Type": "application/json"},
method="POST",
)
try:
with urllib.request.urlopen(req, timeout=120) as antwort:
roh = json.loads(antwort.read().decode("utf-8"))
except urllib.error.URLError as fehler:
raise KategorisierungsFehler(
f"Ollama nicht erreichbar unter {config.OLLAMA_URL}: {fehler}"
) from fehler
inhalt = roh.get("message", {}).get("content", "").strip()
if not inhalt:
raise KategorisierungsFehler("Ollama lieferte keine Antwort.")
try:
ergebnis = json.loads(inhalt)
except json.JSONDecodeError as fehler:
raise KategorisierungsFehler(f"Antwort war kein gültiges JSON: {inhalt[:200]}") from fehler
return _bereinige(ergebnis)
def _bereinige(d: dict) -> dict:
"""Absichern gegen Ausreißer trotz Schema (defensiv)."""
if d.get("kategorie") not in KATEGORIEN:
d["kategorie"] = "sonstiges"
if d.get("dringlichkeit") not in DRINGLICHKEITEN:
d["dringlichkeit"] = "normal"
# Das Modell vergibt gelegentlich die Kategorie "notfall" und stuft die
# Dringlichkeit trotzdem nur auf "hoch" ein. Das ist widersprüchlich und
# führt zu einer harmloseren Farbe auf dem Board, als der Anruf verdient.
# In der Praxis eskaliert man im Zweifel, statt abzuschwächen.
if d["kategorie"] == "notfall":
d["dringlichkeit"] = "notfall"
d.setdefault("anrufer_name", None)
d.setdefault("rueckrufnummer", None)
d.setdefault("rueckruf_gewuenscht", False)
d.setdefault("sprache", "de")
d.setdefault("anliegen_kurz", "")
d.setdefault("stichworte", [])
if not isinstance(d["stichworte"], list):
d["stichworte"] = []
return d

128
pipe/config.py Normal file
View file

@ -0,0 +1,128 @@
"""Zentrale Konfiguration der Anruf-Pipeline.
Alle Werte kommen aus Umgebungsvariablen (optional aus einer .env-Datei im
Projektwurzelverzeichnis). Keine Geheimnisse im Code die Nextcloud-Zugangs-
daten werden erst gesetzt, wenn sie vorliegen; bis dahin bleibt die Ablage rein
lokal.
"""
from __future__ import annotations
import os
from pathlib import Path
WURZEL = Path(__file__).resolve().parent.parent
def _lade_dotenv(pfad: Path) -> None:
"""Minimaler .env-Loader ohne externe Abhängigkeit."""
if not pfad.is_file():
return
for zeile in pfad.read_text(encoding="utf-8").splitlines():
zeile = zeile.strip()
if not zeile or zeile.startswith("#") or "=" not in zeile:
continue
schluessel, _, wert = zeile.partition("=")
schluessel = schluessel.strip()
wert = wert.strip().strip('"').strip("'")
os.environ.setdefault(schluessel, wert)
_lade_dotenv(WURZEL / ".env")
def _bool(name: str, standard: bool) -> bool:
wert = os.environ.get(name)
if wert is None:
return standard
return wert.strip().lower() in {"1", "true", "yes", "ja", "on"}
# --- Spracherkennung (lokal) -------------------------------------------------
# Backend: "whispercpp" (Apple-GPU, nutzt vorhandenes ggml-Modell, kein Download)
# "faster" (faster-whisper / CTranslate2, lädt Modell bei Bedarf)
STT_BACKEND = os.environ.get("STT_BACKEND", "whispercpp").strip().lower()
# Feste Sprache erzwingen (z. B. "de") oder "auto" für Auto-Erkennung.
WHISPER_SPRACHE = os.environ.get("WHISPER_SPRACHE", "de").strip() or "auto"
# whisper.cpp
WHISPERCPP_BIN = os.environ.get("WHISPERCPP_BIN", "whisper-cli")
WHISPERCPP_MODELL = os.environ.get(
"WHISPERCPP_MODELL", str(Path.home() / "whisper-models" / "ggml-large-v3-turbo.bin")
)
WHISPER_THREADS = os.environ.get("WHISPER_THREADS", "8")
# Kontext für die Spracherkennung. Whisper erkennt Namen, Rufnummern und
# Praxis-Vokabular deutlich zuverlässiger, wenn es weiß, worum es geht - bei
# der schlechten Tonqualität einer Telefonleitung macht das den Unterschied.
WHISPER_PROMPT = os.environ.get(
"WHISPER_PROMPT",
"Anruf auf dem Anrufbeantworter einer Hausarztpraxis. Der Anrufer nennt "
"seinen Namen, sein Anliegen und eine Rückrufnummer. Häufig geht es um "
"Termin, Rezept, Überweisung, Befund, Krankschreibung oder Schmerzen.",
).strip()
# faster-whisper (nur bei STT_BACKEND=faster)
WHISPER_MODELL = os.environ.get("WHISPER_MODELL", "medium")
WHISPER_COMPUTE = os.environ.get("WHISPER_COMPUTE", "int8")
# --- Kategorisierung (Ollama, lokal) -----------------------------------------
OLLAMA_URL = os.environ.get("OLLAMA_URL", "http://127.0.0.1:11434")
OLLAMA_MODELL = os.environ.get("OLLAMA_MODELL", "qwen3.5:latest")
PROMPT_DATEI = WURZEL / "prompts" / os.environ.get(
"PROMPT_DATEI", "categorize_de.txt"
)
# --- Ablage ------------------------------------------------------------------
ABLAGE_LOKAL = Path(os.environ.get("ABLAGE_LOKAL", str(WURZEL / "ablage")))
# --- Leitstand (Weboberfläche) ----------------------------------------------
LEITSTAND_HOST = os.environ.get("LEITSTAND_HOST", "127.0.0.1")
LEITSTAND_PORT = int(os.environ.get("LEITSTAND_PORT", "8088"))
# Zugangsschutz. Pflicht, sobald der Leitstand nicht nur lokal erreichbar ist -
# auf den Seiten stehen Transkripte und Aufnahmen von Patienten.
LEITSTAND_USER = os.environ.get("LEITSTAND_USER", "praxis")
LEITSTAND_PASS = os.environ.get("LEITSTAND_PASS", "")
def leitstand_oeffentlich() -> bool:
return LEITSTAND_HOST not in {"127.0.0.1", "localhost", "::1"}
# --- Telefonie ---------------------------------------------------------------
# Die Anlage legt Aufnahmen im Eingang ab; der Watcher (pipe.watch) räumt sie
# nach der Verarbeitung weg.
TELEFON = Path(os.environ.get("TELEFON_ORDNER", str(WURZEL / "telefon")))
TELEFON_EINGANG = TELEFON / "eingang"
TELEFON_VERARBEITET = TELEFON / "verarbeitet"
TELEFON_FEHLER = TELEFON / "fehler"
# SIP-Zugang der Telefonanlage (Werte in .env, nie im Code).
SIP_USER = os.environ.get("SIP_USER", "")
SIP_DOMAIN = os.environ.get("SIP_DOMAIN", "")
SIP_PASS = os.environ.get("SIP_PASS", "")
# Nextcloud (steckbar): nur aktiv, wenn URL + Nutzer + Passwort gesetzt sind.
NEXTCLOUD_URL = os.environ.get("NEXTCLOUD_URL", "").rstrip("/")
NEXTCLOUD_USER = os.environ.get("NEXTCLOUD_USER", "") # Login (kann E-Mail sein)
NEXTCLOUD_PASS = os.environ.get("NEXTCLOUD_PASS", "") # App-Passwort
# WebDAV-Pfad braucht die interne User-ID (weicht bei E-Mail-Login ab).
NEXTCLOUD_USERID = os.environ.get("NEXTCLOUD_USERID", "") or NEXTCLOUD_USER
NEXTCLOUD_ORDNER = os.environ.get("NEXTCLOUD_ORDNER", "Anrufe").strip("/")
# Deck-Triage-Board (optional): pro Anruf eine Karte, Stapel = Kategorie.
NEXTCLOUD_DECK = os.environ.get("NEXTCLOUD_DECK", "").strip().lower() in {
"1", "true", "yes", "ja", "on"
}
DECK_BOARD = os.environ.get("DECK_BOARD", "Anrufe")
# Stapel des Boards = Arbeitsablauf, nicht Kategorie. Neue Anrufe landen immer
# im ersten Stapel; die Kategorie steht auf der Karte.
DECK_STAPEL = [s.strip() for s in os.environ.get(
"DECK_STAPEL", "Eingang,In Bearbeitung,Rückfragen,Erledigt").split(",") if s.strip()]
# Nutzer, die das Board sehen sollen (kommagetrennt). Ohne Freigabe sieht nur
# das Konto der Pipe die Karten - die Praxis schaut in ein leeres Deck.
DECK_TEILEN = [n.strip() for n in os.environ.get("DECK_TEILEN", "").split(",") if n.strip()]
def nextcloud_aktiv() -> bool:
return bool(NEXTCLOUD_URL and NEXTCLOUD_USER and NEXTCLOUD_PASS)

105
pipe/dashboard.py Normal file
View file

@ -0,0 +1,105 @@
"""Erzeugt aus der lokalen Ablage eine HTML-Triage-Übersicht.
python3 -m pipe.dashboard [ausgabe.html]
Rein lokal, liest nur die meta.json-Dateien. Notfälle stehen oben.
"""
from __future__ import annotations
import html
import json
import sys
from pathlib import Path
from . import config
RANG = {"notfall": 0, "hoch": 1, "normal": 2, "niedrig": 3}
FARBE = {"notfall": "#E9322D", "hoch": "#E0A339", "normal": "#C9A227", "niedrig": "#31CC7C"}
KAT_LABEL = {
"termin": "Termin", "rezept": "Rezept", "ueberweisung": "Überweisung",
"befund": "Befund", "verwaltung": "Verwaltung", "beschwerden": "Beschwerden", "notfall": "Notfall",
"rueckruf": "Rückruf", "sonstiges": "Sonstiges",
}
def _anrufe() -> list[dict]:
daten = []
if not config.ABLAGE_LOKAL.is_dir():
return daten
for meta in config.ABLAGE_LOKAL.glob("*/*/meta.json"):
try:
daten.append(json.loads(meta.read_text(encoding="utf-8")))
except Exception:
continue
daten.sort(key=lambda d: (RANG.get(d["auswertung"]["dringlichkeit"], 9), d["empfangen"]))
return daten
def _karte(d: dict) -> str:
e = d["auswertung"]
farbe = FARBE.get(e["dringlichkeit"], "#888")
wer = html.escape(e.get("anrufer_name") or "")
nummer = html.escape(d.get("anrufer_nummer") or e.get("rueckrufnummer") or "")
stich = " ".join(f"<span class='tag'>{html.escape(s)}</span>" for s in (e.get("stichworte") or []))
return f"""
<article class="karte" style="--akzent:{farbe}">
<div class="kopf">
<span class="kat">{html.escape(KAT_LABEL.get(e['kategorie'], e['kategorie']))}</span>
<span class="dring" style="background:{farbe}">{html.escape(e['dringlichkeit'])}</span>
<span class="zeit">{html.escape(d['empfangen'].replace('T', ' '))}</span>
</div>
<div class="wer"><strong>{wer}</strong> · <span class="nr">{nummer}</span>
{"· 📞 Rückruf" if e.get('rueckruf_gewuenscht') else ""}</div>
<p class="anliegen">{html.escape(e.get('anliegen_kurz') or '')}</p>
<div class="tags">{stich}</div>
<details><summary>Transkript</summary><p class="tr">{html.escape(d.get('transkript') or '')}</p></details>
</article>"""
def baue(ausgabe: Path) -> Path:
anrufe = _anrufe()
zahl = {k: sum(1 for d in anrufe if d["auswertung"]["dringlichkeit"] == k) for k in RANG}
zusammenfassung = " · ".join(
f"<b style='color:{FARBE[k]}'>{zahl[k]}</b> {k}" for k in RANG if zahl[k]
)
karten = "\n".join(_karte(d) for d in anrufe)
doc = f"""<!doctype html><html lang="de"><head><meta charset="utf-8">
<meta name="viewport" content="width=device-width,initial-scale=1">
<title>Anruf-Übersicht Praxis Musterhausen</title>
<style>
:root{{--bg:#f6f7f9;--fg:#1a2230;--karte:#fff;--rand:#e4e8ee;--muted:#6b7688}}
@media(prefers-color-scheme:dark){{:root{{--bg:#12151b;--fg:#e7ecf3;--karte:#1b2029;--rand:#2a313d;--muted:#93a0b4}}}}
*{{box-sizing:border-box}}body{{margin:0;background:var(--bg);color:var(--fg);
font:15px/1.5 -apple-system,BlinkMacSystemFont,'Segoe UI',Roboto,sans-serif}}
.wrap{{max-width:840px;margin:0 auto;padding:28px 18px}}
h1{{font-size:22px;margin:0 0 4px}}.sub{{color:var(--muted);margin:0 0 20px;font-size:14px}}
.karte{{background:var(--karte);border:1px solid var(--rand);border-left:5px solid var(--akzent);
border-radius:12px;padding:14px 16px;margin:12px 0}}
.kopf{{display:flex;align-items:center;gap:10px;flex-wrap:wrap;margin-bottom:6px}}
.kat{{font-weight:700}}.dring{{color:#fff;font-size:11px;font-weight:700;text-transform:uppercase;
padding:2px 8px;border-radius:999px;letter-spacing:.03em}}
.zeit{{margin-left:auto;color:var(--muted);font-size:13px;font-variant-numeric:tabular-nums}}
.wer{{margin:2px 0 6px}}.nr{{color:var(--muted)}}
.anliegen{{margin:6px 0}}.tags{{margin:6px 0}}
.tag{{display:inline-block;background:var(--bg);border:1px solid var(--rand);border-radius:6px;
padding:1px 7px;font-size:12px;color:var(--muted);margin:2px 4px 2px 0}}
details{{margin-top:6px}}summary{{cursor:pointer;color:var(--muted);font-size:13px}}
.tr{{color:var(--muted);font-size:13px;margin:6px 0 0}}
</style></head><body><div class="wrap">
<h1>Anruf-Übersicht</h1>
<p class="sub">{len(anrufe)} Anrufe · {zusammenfassung or ''} · lokal erzeugt aus der Ablage</p>
{karten}
</div></body></html>"""
ausgabe.write_text(doc, encoding="utf-8")
return ausgabe
def main() -> int:
ziel = Path(sys.argv[1]) if len(sys.argv) > 1 else config.WURZEL / "dashboard.html"
baue(ziel)
print(f"Dashboard: {ziel}")
return 0
if __name__ == "__main__":
raise SystemExit(main())

162
pipe/deck.py Normal file
View file

@ -0,0 +1,162 @@
"""Nextcloud-Deck-Anbindung: pro Anruf eine Triage-Karte.
Aufbau:
- Ein Board (Standard: "Anrufe").
- Ein Stapel je Kategorie (termin, rezept, ) wird bei Bedarf angelegt.
- Ein Label je Dringlichkeit (farbig) wird bei Bedarf angelegt und der Karte
zugewiesen.
- Karte je Anruf: Titel kompakt, Beschreibung = die Zusammenfassung (Markdown).
Nur aktiv, wenn NEXTCLOUD_DECK=true. Alles über die Deck-REST-API (v1.1).
"""
from __future__ import annotations
import base64
import json
import ssl
import urllib.error
import urllib.request
from functools import lru_cache
from . import config
_API = "/index.php/apps/deck/api/v1.1"
# Farben (6-stelliges Hex ohne #) je Dringlichkeit.
DRINGLICHKEIT_FARBE = {
"niedrig": "31CC7C", # grün
"normal": "F1DB50", # gelb
"hoch": "E0A339", # orange
"notfall": "E9322D", # rot
}
# Lesbare Namen der Kategorien für den Kartentitel - der Stapel sagt jetzt,
# wie weit die Bearbeitung ist, nicht mehr worum es geht.
KATEGORIE_LABEL = {
"termin": "Termin", "rezept": "Rezept", "ueberweisung": "Überweisung",
"befund": "Befund", "verwaltung": "Verwaltung", "beschwerden": "Beschwerden", "notfall": "Notfall",
"rueckruf": "Rückruf", "sonstiges": "Sonstiges",
}
@lru_cache(maxsize=1)
def _ssl_kontext() -> ssl.SSLContext:
try:
import certifi
return ssl.create_default_context(cafile=certifi.where())
except Exception:
return ssl.create_default_context()
def _api(method: str, pfad: str, koerper: dict | None = None):
roh = f"{config.NEXTCLOUD_USER}:{config.NEXTCLOUD_PASS}".encode("utf-8")
headers = {
"Authorization": "Basic " + base64.b64encode(roh).decode("ascii"),
"OCS-APIRequest": "true",
"Content-Type": "application/json",
"User-Agent": "praxis-telefon-agent/1.0",
}
daten = json.dumps(koerper).encode("utf-8") if koerper is not None else None
req = urllib.request.Request(
f"{config.NEXTCLOUD_URL}{_API}{pfad}", data=daten, headers=headers, method=method
)
with urllib.request.urlopen(req, timeout=30, context=_ssl_kontext()) as antwort:
text = antwort.read().decode("utf-8")
return json.loads(text) if text.strip() else None
# --- Board / Stapel / Labels sicherstellen -----------------------------------
def _freigabe(board_id: int, acl: list | None) -> None:
"""Teilt das Board mit den Nutzern aus DECK_TEILEN, falls noch nicht geschehen.
Das Board gehört dem Konto, mit dem die Pipe arbeitet. Ohne Freigabe legt
sie zwar Karten an, die Praxis sieht aber ein leeres Deck.
"""
if not config.DECK_TEILEN:
return
schon = {(a.get("participant") or {}).get("uid") for a in (acl or [])}
for nutzer in config.DECK_TEILEN:
if nutzer in schon:
continue
try:
_api("POST", f"/boards/{board_id}/acl",
{"type": 0, "participant": nutzer, "permissionEdit": True,
"permissionShare": False, "permissionManage": False})
print(f" → Deck: Board mit '{nutzer}' geteilt")
except Exception as fehler:
print(f" → Deck: Freigabe für '{nutzer}' fehlgeschlagen ({fehler})")
def _board_id() -> int:
for board in _api("GET", "/boards") or []:
if board.get("title") == config.DECK_BOARD and not board.get("archived"):
_freigabe(board["id"], board.get("acl"))
return board["id"]
neu = _api("POST", "/boards", {"title": config.DECK_BOARD, "color": "0082C9"})
_freigabe(neu["id"], [])
return neu["id"]
def _stapel(board_id: int) -> dict[str, int]:
"""Stellt die Stapel des Arbeitsablaufs sicher und gibt sie zurück."""
vorhanden = {s["title"]: s["id"] for s in (_api("GET", f"/boards/{board_id}/stacks") or [])}
for i, name in enumerate(config.DECK_STAPEL):
if name not in vorhanden:
neu = _api("POST", f"/boards/{board_id}/stacks", {"title": name, "order": i})
vorhanden[name] = neu["id"]
return vorhanden
def _labels(board_id: int) -> dict[str, int]:
board = _api("GET", f"/boards/{board_id}")
vorhanden = {l["title"]: l["id"] for l in (board.get("labels") or [])}
for name, farbe in DRINGLICHKEIT_FARBE.items():
if name not in vorhanden:
neu = _api("POST", f"/boards/{board_id}/labels",
{"title": name, "color": farbe})
vorhanden[name] = neu["id"]
return vorhanden
# --- Karte anlegen ------------------------------------------------------------
def _kartentitel(datensatz: dict) -> str:
"""Zeit, Kategorie, Anrufer - in dieser Reihenfolge lesbar auf der Karte."""
e = datensatz["auswertung"]
wer = e.get("anrufer_name") or datensatz.get("anrufer_nummer") or "unbekannt"
zeit = datensatz["empfangen"][11:16]
kategorie = KATEGORIE_LABEL.get(e["kategorie"], e["kategorie"])
return f"{zeit} · {kategorie} · {wer}"
def karte_anlegen(datensatz: dict, beschreibung: str) -> int | None:
"""Legt eine Deck-Karte für den Anruf an. Gibt die Karten-ID zurück.
Wirft nie Deck ist ein Zusatz, kein Muss.
"""
if not (config.nextcloud_aktiv() and config.NEXTCLOUD_DECK):
return None
try:
board_id = _board_id()
stapel = _stapel(board_id)
labels = _labels(board_id)
e = datensatz["auswertung"]
eingang = config.DECK_STAPEL[0]
stack_id = stapel[eingang]
karte = _api(
"POST", f"/boards/{board_id}/stacks/{stack_id}/cards",
{"title": _kartentitel(datensatz), "type": "plain", "order": 0,
"description": beschreibung},
)
karten_id = karte["id"]
label_id = labels.get(e.get("dringlichkeit"))
if label_id is not None:
_api("PUT",
f"/boards/{board_id}/stacks/{stack_id}/cards/{karten_id}/assignLabel",
{"labelId": label_id})
return karten_id
except Exception as fehler:
print(f" → Deck: FEHLGESCHLAGEN ({fehler})")
return None

72
pipe/process_call.py Normal file
View file

@ -0,0 +1,72 @@
"""Orchestrator der Anruf-Pipeline: Audio -> Transkript -> Kategorie -> Ablage.
Aufruf:
python3 -m pipe.process_call <audiodatei> [--nummer 021031234567]
Gibt am Ende den Ablageort aus. Vollständig lokal, sofern keine Nextcloud-
Zugangsdaten konfiguriert sind.
"""
from __future__ import annotations
import argparse
import sys
from datetime import datetime
from . import categorize, protokoll, store, transcribe
def verarbeite(audio: str, anrufer_nummer: str | None = None,
empfangen: datetime | None = None) -> dict:
empfangen = empfangen or datetime.now()
print(f"▸ Transkribiere {audio}")
protokoll.schreibe("anruf", "Transkribiere …", nummer=anrufer_nummer)
t = transcribe.transkribiere(audio)
protokoll.schreibe("anruf", (t.text[:200] or "(kein Text erkannt)"),
nummer=anrufer_nummer, dauer=t.dauer)
print(f" Sprache: {t.sprache} ({t.sprache_wahrscheinlichkeit}), Dauer: {t.dauer}s")
print(f" Text: {t.text[:160]}{'' if len(t.text) > 160 else ''}")
print("▸ Kategorisiere …")
auswertung = categorize.kategorisiere(t.text)
print(f" Kategorie: {auswertung['kategorie']} · Dringlichkeit: {auswertung['dringlichkeit']}")
protokoll.schreibe(
"notfall" if auswertung["dringlichkeit"] == "notfall" else "anruf",
f"{auswertung['kategorie']} · {auswertung['dringlichkeit']}"
+ (f" · {auswertung['anrufer_name']}" if auswertung.get("anrufer_name") else ""),
nummer=anrufer_nummer, kategorie=auswertung["kategorie"],
dringlichkeit=auswertung["dringlichkeit"])
datensatz = {
"empfangen": empfangen.isoformat(timespec="seconds"),
"anrufer_nummer": anrufer_nummer,
"transkript": t.text,
"erkannte_sprache": t.sprache,
"audio_dauer_s": t.dauer,
"auswertung": auswertung,
}
print("▸ Lege ab …")
ziel = store.speichere(datensatz, audio)
print(f"✔ Fertig: {ziel}")
protokoll.schreibe("fertig", f"abgelegt: {ziel.parent.name}/{ziel.name}",
nummer=anrufer_nummer)
protokoll.kuerze()
return datensatz
def main(argv: list[str] | None = None) -> int:
p = argparse.ArgumentParser(description="Anruf verarbeiten und ablegen.")
p.add_argument("audio", help="Pfad zur Audiodatei (WAV/OPUS/…)")
p.add_argument("--nummer", default=None, help="Rufnummer des Anrufers (CLIP), falls bekannt")
args = p.parse_args(argv)
try:
verarbeite(args.audio, anrufer_nummer=args.nummer)
except Exception as fehler:
print(f"✖ Fehler: {fehler}", file=sys.stderr)
return 1
return 0
if __name__ == "__main__":
raise SystemExit(main())

56
pipe/protokoll.py Normal file
View file

@ -0,0 +1,56 @@
"""Gemeinsames Ereignisprotokoll der Pipe.
Jede Zeile ist ein JSON-Objekt (JSONL) maschinenlesbar für das Dashboard und
trotzdem mit ``tail -f`` lesbar. Bewusst schlank: anhängen, nie sperren, nie
scheitern. Ein kaputtes Protokoll darf keinen Anruf kosten.
"""
from __future__ import annotations
import json
from datetime import datetime
from . import config
DATEI = config.TELEFON / "verlauf.log"
MAX_ZEILEN = 2000 # darüber wird beim Schreiben gekürzt
def schreibe(art: str, text: str, **felder) -> None:
"""Hängt ein Ereignis an. Wirft nie."""
try:
DATEI.parent.mkdir(parents=True, exist_ok=True)
eintrag = {"zeit": datetime.now().isoformat(timespec="seconds"),
"art": art, "text": text, **felder}
with DATEI.open("a", encoding="utf-8") as f:
f.write(json.dumps(eintrag, ensure_ascii=False) + "\n")
except Exception:
pass
def lies(anzahl: int = 200) -> list[dict]:
"""Liest die letzten Ereignisse. Defekte Zeilen werden übersprungen."""
if not DATEI.is_file():
return []
try:
zeilen = DATEI.read_text(encoding="utf-8").splitlines()[-anzahl:]
except Exception:
return []
eintraege = []
for z in zeilen:
try:
eintraege.append(json.loads(z))
except json.JSONDecodeError:
continue
return eintraege
def kuerze() -> None:
"""Hält das Protokoll auf MAX_ZEILEN. Wirft nie."""
try:
if not DATEI.is_file():
return
zeilen = DATEI.read_text(encoding="utf-8").splitlines()
if len(zeilen) > MAX_ZEILEN:
DATEI.write_text("\n".join(zeilen[-MAX_ZEILEN:]) + "\n", encoding="utf-8")
except Exception:
pass

53
pipe/resync.py Normal file
View file

@ -0,0 +1,53 @@
"""Nachsync: lädt lokal abgelegte Anrufe nach Nextcloud, die noch nicht oben sind.
Nützlich nach einem Nextcloud-Ausfall (z. B. Wartungsmodus): die Pipe speichert
in dem Fall nur lokal; dieser Befehl holt das Hochladen nach.
python3 -m pipe.resync
Erkennt bereits hochgeladene Ordner an der Marker-Datei ``.nextcloud_ok``.
"""
from __future__ import annotations
import sys
from . import config, store
def _anruf_ordner():
"""Alle Anruf-Ordner (JJJJ-MM-TT/HH-MM-SS_*) unter der lokalen Ablage."""
if not config.ABLAGE_LOKAL.is_dir():
return
for tag in sorted(config.ABLAGE_LOKAL.iterdir()):
if tag.is_dir():
for ordner in sorted(tag.iterdir()):
if ordner.is_dir() and (ordner / "meta.json").is_file():
yield ordner
def nachsync() -> tuple[int, int, int]:
"""Gibt (hochgeladen, bereits_oben, fehlgeschlagen) zurück."""
hoch = schon = fehler = 0
for ordner in _anruf_ordner():
if (ordner / store.MARKER).exists():
schon += 1
continue
if store.hochladen(ordner):
print(f"{ordner.parent.name}/{ordner.name}")
hoch += 1
else:
fehler += 1
return hoch, schon, fehler
def main() -> int:
if not config.nextcloud_aktiv():
print("Nextcloud ist nicht konfiguriert — nichts zu tun.", file=sys.stderr)
return 1
hoch, schon, fehler = nachsync()
print(f"\nNachsync fertig: {hoch} hochgeladen, {schon} bereits oben, {fehler} fehlgeschlagen.")
return 0 if fehler == 0 else 2
if __name__ == "__main__":
raise SystemExit(main())

526
pipe/server.py Normal file
View file

@ -0,0 +1,526 @@
"""Live-Dashboard: zeigt Anrufe, Systemzustand und das laufende Protokoll.
python3 -m pipe.server [--port 8088]
Rein lokal, stdlib. Liest die Ablage, das Ereignisprotokoll und den Zustand der
beteiligten Dienste. Die Seite fragt alle paar Sekunden nach kein WebSocket,
kein Framework, nichts, was im Praxisbetrieb kaputtgehen kann.
"""
from __future__ import annotations
import argparse
import base64
import hmac
import json
import re
import subprocess
import sys
import time
import urllib.error
import urllib.request
from datetime import datetime
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from urllib.parse import unquote, urlparse
from . import config, dashboard, protokoll, stapel, store
# Statusabfragen sind teuer (Netz, Unterprozesse) - kurz zwischenspeichern,
# damit haeufiges Nachfragen den Rechner nicht belastet.
_CACHE: dict[str, tuple[float, dict]] = {}
CACHE_S = 4.0
def _laeuft(muster: str) -> bool:
try:
return subprocess.run(["pgrep", "-f", muster], capture_output=True).returncode == 0
except Exception:
return False
def _http_ok(url: str, timeout: float = 5.0) -> bool:
"""Erreichbarkeitstest. Mit eigenem User-Agent - ein WAF vor Nextcloud
blockt den Standardnamen von urllib und meldete fälschlich 'nicht
erreichbar'."""
req = urllib.request.Request(url, headers={"User-Agent": "praxis-telefon-agent/1.0"})
try:
with urllib.request.urlopen(req, timeout=timeout, context=store._ssl_kontext()):
return True
except urllib.error.HTTPError:
return True # antwortet - reicht als Lebenszeichen
except Exception:
return False
def _trunk() -> tuple[str, str]:
"""(zustand, text) des SIP-Trunks."""
if not _laeuft("[f]reeswitch"):
return "aus", "FreeSWITCH läuft nicht"
try:
roh = subprocess.run(
["fs_cli", "-P", "8099", "-x", "sofia status gateway plusnet"],
capture_output=True, text=True, timeout=8).stdout
for zeile in roh.splitlines():
if zeile.startswith("Status"):
wert = zeile.split()[-1]
return ("gut", "registriert") if wert == "UP" else ("schlecht", f"Trunk {wert}")
return "unklar", "Gateway unbekannt"
except Exception:
return "unklar", "fs_cli nicht erreichbar"
def status() -> dict:
jetzt = time.time()
gepuffert = _CACHE.get("status")
if gepuffert and jetzt - gepuffert[0] < CACHE_S:
return gepuffert[1]
trunk_zustand, trunk_text = _trunk()
watcher = _laeuft("[p]ipe.watch")
ollama = _http_ok(f"{config.OLLAMA_URL}/api/version")
nextcloud = (_http_ok(f"{config.NEXTCLOUD_URL}/status.php")
if config.nextcloud_aktiv() else None)
offen = 0
if config.TELEFON_EINGANG.is_dir():
offen = sum(1 for d in config.TELEFON_EINGANG.iterdir()
if d.is_file() and not d.name.startswith("."))
fehler = 0
if config.TELEFON_FEHLER.is_dir():
fehler = sum(1 for d in config.TELEFON_FEHLER.iterdir() if d.is_file())
daten = {
"dienste": [
{"name": "Telefon", "zustand": trunk_zustand, "text": trunk_text},
{"name": "Verarbeitung", "zustand": "gut" if watcher else "aus",
"text": "beobachtet Eingang" if watcher else "Watcher läuft nicht"},
{"name": "Spracherkennung", "zustand": "gut" if ollama else "schlecht",
"text": "Ollama bereit" if ollama else "Ollama nicht erreichbar"},
{"name": "Nextcloud", "zustand": "unklar" if nextcloud is None
else ("gut" if nextcloud else "schlecht"),
"text": "nicht konfiguriert" if nextcloud is None
else ("erreichbar" if nextcloud else "nicht erreichbar")},
],
"eingang_offen": offen,
"fehler": fehler,
"stand": datetime.now().strftime("%H:%M:%S"),
}
_CACHE["status"] = (jetzt, daten)
return daten
def anrufe() -> list[dict]:
liste = []
for d in dashboard._anrufe():
e = d["auswertung"]
zeit = d["empfangen"]
ordner = None
try:
z = datetime.fromisoformat(zeit)
tag = f"{z:%Y-%m-%d}"
for kandidat in (config.ABLAGE_LOKAL / tag).iterdir():
if kandidat.is_dir() and kandidat.name.startswith(f"{z:%H-%M-%S}"):
ordner = f"{tag}/{kandidat.name}"
break
except Exception:
pass
audio = None
if ordner:
for endung in (".wav", ".mp3", ".opus", ".m4a"):
if (config.ABLAGE_LOKAL / ordner / f"aufnahme{endung}").is_file():
audio = f"/audio/{ordner}/aufnahme{endung}"
break
liste.append({
"id": ordner or "",
"stapel": stapel.lies(config.ABLAGE_LOKAL / ordner) if ordner else stapel.erlaubt()[0],
"empfangen": zeit,
"kategorie": e["kategorie"],
"dringlichkeit": e["dringlichkeit"],
"name": e.get("anrufer_name"),
"nummer": d.get("anrufer_nummer") or e.get("rueckrufnummer"),
"rueckruf": bool(e.get("rueckruf_gewuenscht")),
"anliegen": e.get("anliegen_kurz") or "",
"stichworte": e.get("stichworte") or [],
"transkript": d.get("transkript") or "",
"audio": audio,
})
return liste
FS_LOG = Path("/opt/homebrew/var/log/freeswitch/freeswitch.log")
_ZEIT = re.compile(r"\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}")
# Was im FreeSWITCH-Log für die Praxis interessant ist - und wie es heißen soll.
_FS_MUSTER = [
("New Channel sofia/external/", "anruf", "Anruf von {nummer}"),
("has been answered", "anruf", "abgenommen, Ansage läuft"),
("Hangup sofia/external/", "anruf", "Anrufer hat aufgelegt"),
("Failed Registration", "fehler", "Telefon: Registrierung fehlgeschlagen"),
("Ping failed", "fehler", "Telefon: Provider antwortet nicht"),
]
def _telefon_ereignisse(anzahl: int = 60) -> list[dict]:
"""Liest Anruf-Ereignisse aus dem FreeSWITCH-Log.
So sieht man den Anruf schon beim Klingeln, nicht erst wenn die fertige
Aufnahme durch die Pipe gelaufen ist.
"""
if not FS_LOG.is_file():
return []
try:
with FS_LOG.open("rb") as f: # nur das Ende lesen
f.seek(0, 2)
f.seek(max(0, f.tell() - 300_000))
zeilen = f.read().decode("utf-8", "replace").splitlines()
except Exception:
return []
ereignisse = []
for zeile in zeilen[-4000:]:
for schnipsel, art, vorlage in _FS_MUSTER:
if schnipsel not in zeile:
continue
# Bei Kanal-Ereignissen steht die Channel-UUID vor dem Zeitstempel,
# die Position ist also nicht verlässlich - daher per Muster suchen.
treffer = _ZEIT.search(zeile)
zeit = treffer.group(0) if treffer else ""
nummer = ""
if "sofia/external/" in zeile:
rest = zeile.split("sofia/external/", 1)[1]
nummer = rest.split("@", 1)[0].strip()
ereignisse.append({
"zeit": zeit.replace(" ", "T"),
"art": art,
"text": vorlage.format(nummer=nummer or "unbekannt"),
})
break
return ereignisse[-anzahl:]
def verlauf(anzahl: int = 150) -> list[dict]:
"""Pipe-Protokoll und Telefonie-Ereignisse, zeitlich zusammengeführt."""
alle = protokoll.lies(anzahl) + _telefon_ereignisse()
alle.sort(key=lambda e: e.get("zeit") or "")
return alle[-anzahl:]
class Handler(BaseHTTPRequestHandler):
def log_message(self, *_): # kein Zugriffslog auf der Konsole
pass
def _angemeldet(self) -> bool:
"""Prüft die Zugangsdaten, sofern eines gesetzt ist.
Ohne Passwort bleibt der Leitstand offen - das ist nur zulässig, wenn er
an 127.0.0.1 hängt; darauf besteht main() beim Start.
"""
if not config.LEITSTAND_PASS:
return True
kopf = self.headers.get("Authorization", "")
if not kopf.startswith("Basic "):
return False
try:
nutzer, _, wort = base64.b64decode(kopf[6:]).decode("utf-8").partition(":")
except Exception:
return False
# compare_digest: Vergleichsdauer verrät nichts über das Passwort.
return (hmac.compare_digest(nutzer, config.LEITSTAND_USER)
and hmac.compare_digest(wort, config.LEITSTAND_PASS))
def _anmeldung_verlangen(self) -> None:
self.send_response(401)
self.send_header("WWW-Authenticate", 'Basic realm="Anruf-Leitstand"')
self.send_header("Content-Length", "0")
self.end_headers()
def _sende(self, inhalt: bytes, typ: str, code: int = 200) -> None:
self.send_response(code)
self.send_header("Content-Type", typ)
self.send_header("Content-Length", str(len(inhalt)))
self.send_header("Cache-Control", "no-store")
self.end_headers()
self.wfile.write(inhalt)
def _json(self, daten) -> None:
self._sende(json.dumps(daten, ensure_ascii=False).encode("utf-8"),
"application/json; charset=utf-8")
def do_GET(self) -> None:
if not self._angemeldet():
self._anmeldung_verlangen()
return
pfad = unquote(urlparse(self.path).path)
if pfad == "/":
self._sende(SEITE.encode("utf-8"), "text/html; charset=utf-8")
elif pfad == "/api/status":
self._json(status())
elif pfad == "/api/anrufe":
self._json({"stapel": stapel.erlaubt(), "anrufe": anrufe()})
elif pfad == "/api/verlauf":
self._json(verlauf(150))
elif pfad.startswith("/audio/"):
self._audio(pfad[len("/audio/"):])
else:
self._sende(b"nicht gefunden", "text/plain; charset=utf-8", 404)
def do_POST(self) -> None:
if not self._angemeldet():
self._anmeldung_verlangen()
return
if unquote(urlparse(self.path).path) != "/api/stapel":
self._sende(b"nicht gefunden", "text/plain; charset=utf-8", 404)
return
try:
laenge = int(self.headers.get("Content-Length") or 0)
if laenge > 10_000:
raise ValueError("Anfrage zu groß")
wunsch = json.loads(self.rfile.read(laenge).decode("utf-8"))
kennung, ziel = wunsch["id"], wunsch["stapel"]
except Exception:
self._json({"ok": False, "fehler": "ungültige Anfrage"})
return
ordner = stapel.ordner_zu(kennung)
if ordner is None:
self._json({"ok": False, "fehler": "Anruf nicht gefunden"})
return
if not stapel.setze(ordner, ziel):
self._json({"ok": False, "fehler": f"unbekannter Stapel: {ziel}"})
return
protokoll.schreibe("stapel", f"{ordner.name}{ziel}")
self._json({"ok": True})
def _audio(self, rest: str) -> None:
"""Liefert eine Aufnahme aus der Ablage - nur von dort."""
ziel = (config.ABLAGE_LOKAL / rest).resolve()
wurzel = config.ABLAGE_LOKAL.resolve()
if not str(ziel).startswith(str(wurzel) + "/") or not ziel.is_file():
self._sende(b"nicht gefunden", "text/plain; charset=utf-8", 404)
return
typ = {".wav": "audio/wav", ".mp3": "audio/mpeg",
".opus": "audio/ogg", ".m4a": "audio/mp4"}.get(ziel.suffix, "application/octet-stream")
self._sende(ziel.read_bytes(), typ)
SEITE = r"""<!doctype html><html lang="de"><head><meta charset="utf-8">
<meta name="viewport" content="width=device-width,initial-scale=1">
<title>Anruf-Leitstand Praxis Musterhausen</title>
<style>
:root{--bg:#f4f6f9;--fg:#18202c;--karte:#fff;--spalte:#eceff4;--rand:#dde3ec;--muted:#68748a;
--gut:#1f9d5c;--schlecht:#d93a34;--unklar:#c8901f;--aus:#8a93a3;--log:#f8fafc}
@media(prefers-color-scheme:dark){:root{--bg:#0f1218;--fg:#e6ecf4;--karte:#1b212b;
--spalte:#151a22;--rand:#28303d;--muted:#8d99ad;--log:#10141b}}
*{box-sizing:border-box}
body{margin:0;background:var(--bg);color:var(--fg);
font:15px/1.5 -apple-system,BlinkMacSystemFont,'Segoe UI',Roboto,sans-serif}
.wrap{max-width:1600px;margin:0 auto;padding:18px 16px 30px}
header{display:flex;align-items:baseline;gap:12px;flex-wrap:wrap;margin-bottom:12px}
h1{font-size:19px;margin:0}
.stand{color:var(--muted);font-size:13px;font-variant-numeric:tabular-nums;margin-left:auto}
.dienste{display:grid;grid-template-columns:repeat(auto-fit,minmax(180px,1fr));gap:9px}
.dienst{background:var(--karte);border:1px solid var(--rand);border-radius:9px;padding:9px 11px;
display:flex;align-items:center;gap:9px}
.punkt{width:9px;height:9px;border-radius:50%;flex:none}
.gut{background:var(--gut)}.schlecht{background:var(--schlecht)}
.unklar{background:var(--unklar)}.aus{background:var(--aus)}
.dienst b{display:block;font-size:13px}.dienst span{color:var(--muted);font-size:12px}
.hinweis{background:#fdf1e7;border:1px solid #f0c9a0;color:#8a4b12;border-radius:8px;
padding:7px 11px;margin:9px 0;font-size:13px}
@media(prefers-color-scheme:dark){.hinweis{background:#2b1e12;border-color:#5a3a18;color:#f0c08a}}
.board{display:grid;grid-template-columns:repeat(4,1fr);gap:12px;margin-top:12px;align-items:start}
@media(max-width:1100px){.board{grid-template-columns:repeat(2,1fr)}}
@media(max-width:660px){.board{grid-template-columns:1fr}}
.spalte{background:var(--spalte);border:1px solid var(--rand);border-radius:12px;padding:10px;
min-height:130px;transition:background .12s,border-color .12s}
.spalte.ueber{border-color:var(--gut);background:color-mix(in srgb,var(--gut) 8%,var(--spalte))}
.spalte h2{font-size:12px;text-transform:uppercase;letter-spacing:.06em;color:var(--muted);
margin:0 0 9px;display:flex;gap:7px;align-items:center}
.zahl{background:var(--rand);color:var(--fg);border-radius:999px;padding:0 7px;font-size:11px}
a.waehlen{color:inherit;text-decoration:underline;text-decoration-style:dotted;
text-underline-offset:2px;cursor:pointer}
a.waehlen:hover{color:var(--gut);text-decoration-style:solid}
.karte{background:var(--karte);border:1px solid var(--rand);border-left:4px solid var(--akzent,#888);
border-radius:9px;padding:9px 11px;margin-bottom:9px;cursor:grab}
.karte:active{cursor:grabbing}
.karte.zieht{opacity:.4}
.kopf{display:flex;align-items:center;gap:7px;flex-wrap:wrap}
.titel{font-weight:650;font-size:14px}
.dring{color:#fff;font-size:9px;font-weight:700;text-transform:uppercase;padding:2px 6px;
border-radius:999px;letter-spacing:.04em;margin-left:auto}
.meta{color:var(--muted);font-size:12px;margin:3px 0 5px;font-variant-numeric:tabular-nums}
.anliegen{margin:5px 0;font-size:13.5px}
.tag{display:inline-block;background:var(--bg);border:1px solid var(--rand);border-radius:5px;
padding:1px 6px;font-size:11px;color:var(--muted);margin:2px 3px 2px 0}
audio{width:100%;height:32px;margin-top:6px;display:block}
details{margin-top:4px}summary{cursor:pointer;color:var(--muted);font-size:12px}
.tr{color:var(--muted);font-size:12.5px;margin:4px 0 0}
.schieber{display:flex;gap:5px;margin-top:7px;flex-wrap:wrap}
.schieber button{font:inherit;font-size:11px;padding:2px 8px;border-radius:6px;cursor:pointer;
border:1px solid var(--rand);background:var(--bg);color:var(--muted)}
.schieber button:hover{color:var(--fg);border-color:var(--muted)}
.leer{color:var(--muted);font-size:12.5px;padding:5px 2px}
.dienste-kopf{font-size:12px;text-transform:uppercase;letter-spacing:.06em;color:var(--muted);
margin:22px 0 7px}
#verlauf{background:var(--log);border:1px solid var(--rand);border-radius:11px;padding:10px;
height:180px;overflow-y:auto;font:12px/1.55 ui-monospace,SFMono-Regular,Menlo,monospace}
.z{display:flex;gap:9px;padding:1px 0}
.z time{color:var(--muted);flex:none;font-variant-numeric:tabular-nums}
.z .a{flex:none;width:64px;font-weight:600}
.a-anruf{color:#3b86d8}.a-eingang{color:#1f9d5c}.a-fertig{color:#1f9d5c}
.a-fehler{color:var(--schlecht)}.a-notfall{color:var(--schlecht);font-weight:700}
.a-system{color:var(--muted)}.a-stapel{color:#9b6bd8}
</style></head><body><div class="wrap">
<header><h1>Anruf-Leitstand</h1><span class="stand" id="stand"></span></header>
<div id="verlauf"></div>
<div id="hinweise"></div>
<div class="board" id="board"></div>
<div class="dienste-kopf">Dienste</div>
<div class="dienste" id="dienste"></div>
</div>
<script>
const FARBE={notfall:'#d93a34',hoch:'#e08a2a',normal:'#c8a227',niedrig:'#1f9d5c'};
const KAT={termin:'Termin',rezept:'Rezept',ueberweisung:'Überweisung',befund:'Befund',
verwaltung:'Verwaltung',beschwerden:'Beschwerden',notfall:'Notfall',rueckruf:'Rückruf',sonstiges:'Sonstiges'};
const esc=s=>String(s??'').replace(/[&<>"]/g,c=>({'&':'&amp;','<':'&lt;','>':'&gt;','"':'&quot;'}[c]));
let STAPEL=[], ANRUFE=[], offeneTranskripte=new Set(), pausiert=false;
async function hole(p){try{const r=await fetch(p);return r.ok?await r.json():null}catch(e){return null}}
async function verschiebe(id,ziel){
// sofort anzeigen, dann speichern - fuehlt sich unmittelbar an
const a=ANRUFE.find(x=>x.id===id); if(!a||a.stapel===ziel)return;
a.stapel=ziel; zeichneBoard();
const r=await fetch('/api/stapel',{method:'POST',headers:{'Content-Type':'application/json'},
body:JSON.stringify({id:id,stapel:ziel})});
const d=await r.json().catch(()=>null);
if(!d||!d.ok){takt()} // hat nicht geklappt: echten Stand neu holen
}
// Nur Ziffern und ein fuehrendes Plus taugen fuer tel:
const waehlbar=n=>String(n||'').replace(/[^\d+]/g,'').replace(/(?!^)\+/g,'');
function nummerLink(d){
const n=d.nummer, w=waehlbar(n);
if(!w)return '';
return `<a class="waehlen" href="tel:${esc(w)}" data-id="${esc(d.id)}"
title="Rückruf an ${esc(n)} — verschiebt die Karte in ${esc(STAPEL[1]||'')}">${esc(n)}</a>`;
}
function karte(d){
// Notfall zaehlt als Notfall, auch wenn aeltere Datensaetze die Dringlichkeit
// niedriger gesetzt haben.
const f=(d.kategorie==='notfall')?FARBE.notfall:(FARBE[d.dringlichkeit]||'#888');
const i=STAPEL.indexOf(d.stapel);
const knoepfe=STAPEL.map((s,j)=>j===i?'':
`<button data-id="${esc(d.id)}" data-ziel="${esc(s)}">${j<i?'':''} ${esc(s)}</button>`).join('');
const offen=offeneTranskripte.has(d.id)?' open':'';
return `<article class="karte" draggable="true" data-id="${esc(d.id)}" style="--akzent:${f}">
<div class="kopf"><span class="titel">${esc(KAT[d.kategorie]||d.kategorie)} · ${esc(d.name||'unbekannt')}</span>
<span class="dring" style="background:${f}">${esc(d.kategorie==='notfall'?'notfall':d.dringlichkeit)}</span></div>
<div class="meta">${esc(d.empfangen.replace('T',' '))} · ${nummerLink(d)}${d.rueckruf?' · 📞 Rückruf':''}</div>
<p class="anliegen">${esc(d.anliegen)}</p>
<div>${(d.stichworte||[]).map(s=>`<span class="tag">${esc(s)}</span>`).join('')}</div>
${d.audio?`<audio controls preload="none" src="${d.audio}"></audio>`:''}
<details data-id="${esc(d.id)}"${offen}><summary>Transkript</summary><p class="tr">${esc(d.transkript)}</p></details>
<div class="schieber">${knoepfe}</div>
</article>`;
}
function zeichneBoard(){
const el=document.getElementById('board');
el.innerHTML=STAPEL.map(s=>{
const drin=ANRUFE.filter(a=>a.stapel===s);
return `<section class="spalte" data-stapel="${esc(s)}">
<h2>${esc(s)}<span class="zahl">${drin.length}</span></h2>
${drin.map(karte).join('')||'<p class="leer">leer</p>'}
</section>`}).join('');
el.querySelectorAll('.schieber button').forEach(b=>
b.onclick=e=>{e.stopPropagation();verschiebe(b.dataset.id,b.dataset.ziel)});
el.querySelectorAll('a.waehlen').forEach(a=>
a.onclick=()=>{ // der Link waehlt, wir ziehen nur nach
const d=ANRUFE.find(x=>x.id===a.dataset.id);
if(d&&STAPEL[1]&&d.stapel===STAPEL[0])verschiebe(d.id,STAPEL[1]);
});
el.querySelectorAll('details').forEach(d=>{
d.ontoggle=()=>{d.open?offeneTranskripte.add(d.dataset.id):offeneTranskripte.delete(d.dataset.id)};
});
el.querySelectorAll('.karte').forEach(k=>{
k.ondragstart=e=>{pausiert=true;k.classList.add('zieht');e.dataTransfer.setData('text/plain',k.dataset.id)};
k.ondragend=()=>{pausiert=false;k.classList.remove('zieht')};
});
el.querySelectorAll('.spalte').forEach(sp=>{
sp.ondragover=e=>{e.preventDefault();sp.classList.add('ueber')};
sp.ondragleave=()=>sp.classList.remove('ueber');
sp.ondrop=e=>{e.preventDefault();sp.classList.remove('ueber');pausiert=false;
verschiebe(e.dataTransfer.getData('text/plain'),sp.dataset.stapel)};
});
}
function zeichneStatus(s){
if(!s)return;
document.getElementById('stand').textContent='Stand '+s.stand;
document.getElementById('dienste').innerHTML=s.dienste.map(d=>
`<div class="dienst"><span class="punkt ${d.zustand}"></span>
<span><b>${esc(d.name)}</b><span>${esc(d.text)}</span></span></div>`).join('');
const h=[];
// Die Ampeln stehen unten - eine Stoerung muss trotzdem sofort auffallen.
s.dienste.filter(d=>d.zustand==='schlecht'||d.zustand==='aus')
.forEach(d=>h.push(`${d.name}: ${d.text}`));
if(s.eingang_offen>0)h.push(`${s.eingang_offen} Aufnahme(n) warten auf Verarbeitung.`);
if(s.fehler>0)h.push(`${s.fehler} Aufnahme(n) in telefon/fehler nicht verarbeitet.`);
document.getElementById('hinweise').innerHTML=h.map(t=>`<div class="hinweis">${esc(t)}</div>`).join('');
}
function zeichneVerlauf(v){
if(!v)return;
const el=document.getElementById('verlauf');
const unten=el.scrollHeight-el.scrollTop-el.clientHeight<40;
el.innerHTML=v.map(e=>`<div class="z"><time>${esc((e.zeit||'').slice(11,19))}</time>
<span class="a a-${esc(e.art)}">${esc(e.art)}</span><span>${esc(e.text)}</span></div>`).join('')
||'<p class="leer">Noch nichts passiert.</p>';
if(unten)el.scrollTop=el.scrollHeight;
}
async function takt(){
if(pausiert)return; // nicht neu zeichnen, waehrend jemand zieht
const [s,a,v]=await Promise.all([hole('/api/status'),hole('/api/anrufe'),hole('/api/verlauf')]);
zeichneStatus(s);
if(a){STAPEL=a.stapel;ANRUFE=a.anrufe;zeichneBoard()}
zeichneVerlauf(v);
}
takt();setInterval(takt,3000);
</script></body></html>"""
def main(argv: list[str] | None = None) -> int:
p = argparse.ArgumentParser(description="Live-Dashboard der Anruf-Pipe.")
p.add_argument("--port", type=int, default=config.LEITSTAND_PORT)
p.add_argument("--host", default=config.LEITSTAND_HOST,
help="0.0.0.0 macht den Leitstand im Netz erreichbar (dann ist "
"LEITSTAND_PASS Pflicht)")
args = p.parse_args(argv)
oeffentlich = args.host not in {"127.0.0.1", "localhost", "::1"}
if oeffentlich and not config.LEITSTAND_PASS:
print("Abbruch: Der Leitstand soll über das Netz erreichbar sein, aber es ist\n"
"kein LEITSTAND_PASS gesetzt. Auf den Seiten stehen Transkripte und\n"
"Aufnahmen von Patienten — ohne Passwort läge das für jeden im Netz offen.\n"
"Passwort in die .env eintragen und erneut starten.", file=sys.stderr)
return 2
server = ThreadingHTTPServer((args.host, args.port), Handler)
schutz = "mit Passwort" if config.LEITSTAND_PASS else "ohne Passwort (nur lokal)"
print(f"Leitstand: http://{args.host}:{args.port} [{schutz}] — Abbruch mit Strg-C")
try:
server.serve_forever()
except KeyboardInterrupt:
print("\nLeitstand beendet.")
return 0
if __name__ == "__main__":
raise SystemExit(main())

58
pipe/stapel.py Normal file
View file

@ -0,0 +1,58 @@
"""Bearbeitungsstand eines Anrufs — wo er auf dem Board liegt.
Der Stand gehört nicht in ``meta.json``: dort steht, was der Anruf war, und das
ändert sich nie. Wo er in der Bearbeitung steht, ändert sich dauernd. Deshalb
eine eigene kleine Datei im Anrufordner so bleibt beides unabhängig, und ein
verlorener Stand kostet nie die Anrufdaten.
"""
from __future__ import annotations
import json
from datetime import datetime
from . import config
DATEI = "stapel.json"
def erlaubt() -> list[str]:
return config.DECK_STAPEL
def lies(ordner) -> str:
"""Stapel eines Anrufs. Unbekanntes oder Fehlendes gilt als Eingang."""
standard = erlaubt()[0]
pfad = ordner / DATEI
if not pfad.is_file():
return standard
try:
wert = json.loads(pfad.read_text(encoding="utf-8")).get("stapel")
except Exception:
return standard
return wert if wert in erlaubt() else standard
def setze(ordner, stapel: str) -> bool:
"""Schreibt den Stapel. False, wenn der Name nicht vorgesehen ist."""
if stapel not in erlaubt():
return False
try:
(ordner / DATEI).write_text(json.dumps(
{"stapel": stapel, "geaendert": datetime.now().isoformat(timespec="seconds")},
ensure_ascii=False, indent=2), encoding="utf-8")
return True
except Exception:
return False
def ordner_zu(kennung: str):
"""Wandelt eine Anruf-Kennung (JJJJ-MM-TT/HH-MM-SS_nummer) in einen Pfad.
Prüft, dass der Pfad wirklich in der Ablage liegt die Kennung kommt aus
dem Browser und ist damit nichts, worauf man sich verlässt.
"""
ziel = (config.ABLAGE_LOKAL / kennung).resolve()
wurzel = config.ABLAGE_LOKAL.resolve()
if not str(ziel).startswith(str(wurzel) + "/") or not ziel.is_dir():
return None
return ziel

208
pipe/store.py Normal file
View file

@ -0,0 +1,208 @@
"""Ablage eines verarbeiteten Anrufs.
Immer lokal (Ordner pro Anruf mit meta.json, zusammenfassung.md, Audiokopie).
Zusätzlich Nextcloud per WebDAV, sobald Zugangsdaten in der Konfiguration
stehen sonst wird der Nextcloud-Schritt übersprungen (steckbar).
"""
from __future__ import annotations
import base64
import json
import re
import shutil
import ssl
import time
import urllib.error
import urllib.request
from datetime import datetime
from functools import lru_cache
from pathlib import Path
from . import config
@lru_cache(maxsize=1)
def _ssl_kontext() -> ssl.SSLContext:
"""SSL-Kontext mit CA-Bundle (macOS-Python bringt keinen System-Store mit)."""
try:
import certifi
return ssl.create_default_context(cafile=certifi.where())
except Exception:
return ssl.create_default_context()
DRINGLICHKEIT_SYMBOL = {"niedrig": "", "normal": "", "hoch": "", "notfall": ""}
def _ordnername(zeit: datetime, nummer: str | None) -> str:
"""Ordnername aus Uhrzeit und Rufnummer, dateisystemsicher.
Die alte Fassung schnitt auf die letzten 12 Zeichen zu und zerstörte damit
internationale Nummern: aus +49000000000 wurde 915223062462. Vorwahlen
werden daher in 00-Schreibweise gebracht und die Nummer bleibt vollstaendig.
"""
roh = (nummer or "").strip()
if roh.startswith("+"):
roh = "00" + roh[1:]
sauber = re.sub(r"[^0-9A-Za-z]", "", roh)
return f"{zeit:%H-%M-%S}_{sauber[:24] or 'unbekannt'}"
def _markdown(datensatz: dict) -> str:
e = datensatz["auswertung"]
symbol = DRINGLICHKEIT_SYMBOL.get(e["dringlichkeit"], "")
zeilen = [
f"# Anruf {datensatz['empfangen']}{e['kategorie']}",
"",
f"- **Dringlichkeit:** {symbol} {e['dringlichkeit']}",
f"- **Anrufer:** {e.get('anrufer_name') or ''}",
f"- **Rufnummer (CLIP):** {datensatz.get('anrufer_nummer') or ''}",
f"- **Rückrufnummer (genannt):** {e.get('rueckrufnummer') or ''}",
f"- **Rückruf gewünscht:** {'ja' if e.get('rueckruf_gewuenscht') else 'nein'}",
f"- **Sprache:** {e.get('sprache')}",
f"- **Stichworte:** {', '.join(e.get('stichworte') or []) or ''}",
"",
"## Anliegen",
e.get("anliegen_kurz") or "",
"",
"## Transkript",
datensatz.get("transkript") or "",
"",
]
return "\n".join(zeilen)
def speichere(datensatz: dict, audio: str | Path | None) -> Path:
"""Legt den Anruf lokal ab (und in Nextcloud, falls konfiguriert).
Gibt den lokalen Ordnerpfad zurück.
"""
zeit = datetime.fromisoformat(datensatz["empfangen"])
tages_ordner = config.ABLAGE_LOKAL / f"{zeit:%Y-%m-%d}"
# Die Rufnummer der Anlage (CLIP) ist verlässlich, die aus dem Transkript
# gelesene nur geraten - daher hat CLIP Vorrang im Ordnernamen.
nummer = (datensatz.get("anrufer_nummer")
or datensatz["auswertung"].get("rueckrufnummer"))
ziel = tages_ordner / _ordnername(zeit, nummer)
ziel.mkdir(parents=True, exist_ok=True)
markdown = _markdown(datensatz)
(ziel / "meta.json").write_text(
json.dumps(datensatz, ensure_ascii=False, indent=2), encoding="utf-8"
)
(ziel / "zusammenfassung.md").write_text(markdown, encoding="utf-8")
audio_ziel = None
if audio and Path(audio).is_file():
audio_ziel = ziel / f"aufnahme{Path(audio).suffix or '.wav'}"
shutil.copy2(audio, audio_ziel)
if config.nextcloud_aktiv():
if hochladen(ziel):
print(f" → Nextcloud: hochgeladen ({config.NEXTCLOUD_ORDNER}/{zeit:%Y-%m-%d}/{ziel.name})")
else:
print(" → Nextcloud: nicht erreichbar, nur lokal gespeichert (Nachsync via pipe.resync)")
if config.NEXTCLOUD_DECK:
from . import deck # späte Einbindung: Deck ist optional
if deck.karte_anlegen(datensatz, markdown):
print(f" → Deck: Karte in '{config.DECK_STAPEL[0]}' angelegt")
return ziel
MARKER = ".nextcloud_ok"
def hochladen(ordner: Path) -> bool:
"""Lädt einen Anruf-Ordner nach Nextcloud. True bei Erfolg.
Legt bei Erfolg eine Marker-Datei an, damit der Nachsync weiß, was schon
oben ist. Wirft nie die lokale Ablage darf nie an Nextcloud scheitern.
"""
if not config.nextcloud_aktiv():
return False
try:
zeit = datetime.strptime(ordner.parent.name, "%Y-%m-%d")
_nach_nextcloud(ordner, zeit)
(ordner / MARKER).write_text("", encoding="utf-8")
return True
except Exception as fehler:
print(f" Upload {ordner.name}: {fehler}")
return False
# --- Nextcloud WebDAV (steckbar, stdlib) -------------------------------------
def _webdav_basis() -> str:
return f"{config.NEXTCLOUD_URL}/remote.php/dav/files/{config.NEXTCLOUD_USERID}"
def _auth_header() -> dict:
roh = f"{config.NEXTCLOUD_USER}:{config.NEXTCLOUD_PASS}".encode("utf-8")
return {
"Authorization": "Basic " + base64.b64encode(roh).decode("ascii"),
# Ein WAF/Proxy vor Nextcloud blockt den Standard-User-Agent von urllib.
"User-Agent": "praxis-telefon-agent/1.0",
}
# Nextcloud antwortet gelegentlich träge (PHP-Kaltstart, Hintergrund-Jobs). Ein
# einzelner Aussetzer darf einen Anruf nicht in den manuellen Nachsync schicken.
WEBDAV_TIMEOUT = 60
WEBDAV_VERSUCHE = 3
def _sende(req: urllib.request.Request, timeout: int = WEBDAV_TIMEOUT):
"""Führt eine WebDAV-Anfrage aus, mit Wiederholung bei Aussetzern.
Wiederholt nur, was sich durch Wiederholen lösen lässt: Zeitüberschreitung,
Verbindungsabbruch, Serverfehler (5xx). MKCOL und PUT sind idempotent, ein
zweiter Versuch kann also nichts doppelt anlegen.
"""
for versuch in range(1, WEBDAV_VERSUCHE + 1):
try:
return urllib.request.urlopen(req, timeout=timeout, context=_ssl_kontext())
except urllib.error.HTTPError as fehler:
if fehler.code < 500 or versuch == WEBDAV_VERSUCHE:
raise
grund = f"HTTP {fehler.code}"
except OSError as fehler: # umfasst URLError und Zeitüberschreitung
if versuch == WEBDAV_VERSUCHE:
raise
grund = str(fehler) or type(fehler).__name__
wartezeit = 2 ** versuch
print(f" {req.get_method()} {req.selector.rsplit('/', 1)[-1]}: {grund}"
f" - neuer Versuch in {wartezeit}s")
time.sleep(wartezeit)
def _mkcol(pfad: str) -> None:
req = urllib.request.Request(
f"{_webdav_basis()}/{pfad}", method="MKCOL", headers=_auth_header()
)
try:
_sende(req)
except urllib.error.HTTPError as e:
if e.code != 405: # 405 = Ordner existiert bereits
raise
def _put(lokal: Path, pfad: str) -> None:
req = urllib.request.Request(
f"{_webdav_basis()}/{pfad}",
data=lokal.read_bytes(),
method="PUT",
headers=_auth_header(),
)
_sende(req, timeout=180)
def _nach_nextcloud(ordner: Path, zeit: datetime) -> None:
basis = config.NEXTCLOUD_ORDNER
tag = f"{zeit:%Y-%m-%d}"
_mkcol(basis)
_mkcol(f"{basis}/{tag}")
_mkcol(f"{basis}/{tag}/{ordner.name}")
for datei in sorted(ordner.iterdir()):
if datei.is_file() and datei.name != MARKER:
_put(datei, f"{basis}/{tag}/{ordner.name}/{datei.name}")

141
pipe/transcribe.py Normal file
View file

@ -0,0 +1,141 @@
"""Spracherkennung: Audiodatei -> Transkript (lokal, kein Cloud-Aufruf).
Zwei Backends:
- "whispercpp": nutzt das CLI `whisper-cli` mit einem vorhandenen ggml-Modell
(Apple-GPU-beschleunigt, kein Modell-Download).
- "faster": faster-whisper (CTranslate2); lädt sein Modell bei Bedarf.
Beliebige Eingabeformate/Abtastraten werden vor der Erkennung per ffmpeg auf
16 kHz mono normalisiert passt auch für 8-kHz-Telefonaufnahmen.
"""
from __future__ import annotations
import json
import subprocess
import tempfile
from dataclasses import dataclass, field
from functools import lru_cache
from pathlib import Path
from . import config
@dataclass
class Transkript:
text: str
sprache: str
sprache_wahrscheinlichkeit: float
dauer: float
segmente: list[dict] = field(default_factory=list)
def _letzte_zeile(ausgabe: str | None) -> str:
"""Letzte nicht-leere Zeile einer Fehlerausgabe (für knappe Meldungen)."""
zeilen = [z for z in (ausgabe or "").strip().splitlines() if z.strip()]
return zeilen[-1].strip() if zeilen else "keine Fehlerausgabe"
def _nach_16k_mono(audio: Path) -> Path:
"""Normalisiert beliebiges Audio auf 16 kHz mono WAV (Tempdatei)."""
ziel = Path(tempfile.mkstemp(suffix=".wav", prefix="stt_")[1])
try:
subprocess.run(
["ffmpeg", "-y", "-i", str(audio), "-ar", "16000", "-ac", "1", str(ziel)],
check=True, capture_output=True, text=True,
)
except FileNotFoundError as fehler:
ziel.unlink(missing_ok=True)
raise RuntimeError("ffmpeg wurde nicht gefunden (brew install ffmpeg).") from fehler
except subprocess.CalledProcessError as fehler:
ziel.unlink(missing_ok=True)
raise RuntimeError(
f"ffmpeg konnte {audio.name} nicht lesen: {_letzte_zeile(fehler.stderr)}"
) from fehler
return ziel
def transkribiere(audio: str | Path) -> Transkript:
"""Wandelt eine Audiodatei in Text um. FileNotFoundError, wenn sie fehlt."""
pfad = Path(audio)
if not pfad.is_file():
raise FileNotFoundError(f"Audiodatei nicht gefunden: {pfad}")
if config.STT_BACKEND == "faster":
return _mit_faster_whisper(pfad)
return _mit_whispercpp(pfad)
# --- Backend: whisper.cpp -----------------------------------------------------
def _mit_whispercpp(audio: Path) -> Transkript:
wav = _nach_16k_mono(audio)
ausgabe_basis = Path(tempfile.mkstemp(suffix="", prefix="stt_out_")[1])
try:
cmd = [
config.WHISPERCPP_BIN,
"-m", config.WHISPERCPP_MODELL,
"-f", str(wav),
"-l", config.WHISPER_SPRACHE,
"-t", str(config.WHISPER_THREADS),
"-oj", "-of", str(ausgabe_basis),
"-np",
]
if config.WHISPER_PROMPT:
cmd += ["--prompt", config.WHISPER_PROMPT]
try:
subprocess.run(cmd, check=True, capture_output=True, text=True)
except FileNotFoundError as fehler:
raise RuntimeError(
f"Spracherkennung '{config.WHISPERCPP_BIN}' nicht gefunden "
f"(brew install whisper-cpp, oder STT_BACKEND=faster setzen)."
) from fehler
except subprocess.CalledProcessError as fehler:
raise RuntimeError(
f"whisper-cli brach ab: {_letzte_zeile(fehler.stderr)}"
) from fehler
roh = json.loads(Path(f"{ausgabe_basis}.json").read_text(encoding="utf-8"))
finally:
wav.unlink(missing_ok=True)
Path(f"{ausgabe_basis}.json").unlink(missing_ok=True)
ausgabe_basis.unlink(missing_ok=True)
segmente = []
for s in roh.get("transcription", []):
text = (s.get("text") or "").strip()
if text:
off = s.get("offsets", {})
segmente.append({
"start": round(off.get("from", 0) / 1000, 2),
"ende": round(off.get("to", 0) / 1000, 2),
"text": text,
})
text_ganz = " ".join(s["text"] for s in segmente).strip()
sprache = roh.get("result", {}).get("language", config.WHISPER_SPRACHE)
dauer = segmente[-1]["ende"] if segmente else 0.0
return Transkript(text=text_ganz, sprache=sprache,
sprache_wahrscheinlichkeit=1.0, dauer=dauer, segmente=segmente)
# --- Backend: faster-whisper --------------------------------------------------
@lru_cache(maxsize=1)
def _faster_modell():
from faster_whisper import WhisperModel
return WhisperModel(config.WHISPER_MODELL, device="cpu",
compute_type=config.WHISPER_COMPUTE)
def _mit_faster_whisper(audio: Path) -> Transkript:
sprache = None if config.WHISPER_SPRACHE == "auto" else config.WHISPER_SPRACHE
segmente, info = _faster_modell().transcribe(str(audio), language=sprache, vad_filter=True)
teile, seg_liste = [], []
for s in segmente:
t = s.text.strip()
if t:
teile.append(t)
seg_liste.append({"start": round(s.start, 2), "ende": round(s.end, 2), "text": t})
return Transkript(
text=" ".join(teile).strip(), sprache=info.language,
sprache_wahrscheinlichkeit=round(info.language_probability, 3),
dauer=round(info.duration, 2), segmente=seg_liste,
)

133
pipe/watch.py Normal file
View file

@ -0,0 +1,133 @@
"""Beobachtet den Eingangsordner und schiebt neue Aufnahmen durch die Pipe.
python3 -m pipe.watch [--einmal]
Die Telefonanlage legt Aufnahmen in ``telefon/eingang`` ab. Der Dateiname trägt
die Rufnummer des Anrufers (CLIP), damit sie nicht verloren geht:
JJJJMMTT-HHMMSS_<nummer>.wav
Verarbeitete Aufnahmen wandern nach ``telefon/verarbeitet``, gescheiterte nach
``telefon/fehler`` dort bleiben sie liegen, statt still verloren zu gehen.
Bewusst mit Polling statt Dateisystem-Ereignissen: eine Praxis bekommt Anrufe im
Minutenabstand, nicht im Millisekundenabstand, und Polling übersteht Neustarts
und Netzlaufwerke ohne Sonderfälle.
"""
from __future__ import annotations
import argparse
import re
import shutil
import sys
import time
from datetime import datetime
from pathlib import Path
from . import config, process_call, protokoll
TAKT_S = 3 # Wartezeit zwischen zwei Durchläufen
RUHE_S = 2 # So lange muss die Dateigröße stabil sein (Aufnahme fertig)
ENDUNGEN = {".wav", ".mp3", ".opus", ".ogg", ".m4a", ".aiff"}
# JJJJMMTT-HHMMSS_<nummer>.<ext> - beides optional, damit auch handverlesene
# Dateien verarbeitet werden.
_NAME = re.compile(r"^(?:(\d{8})-(\d{6}))?_?(.*)$")
def _nummer_und_zeit(datei: Path) -> tuple[str | None, datetime | None]:
"""Liest Rufnummer und Aufnahmezeit aus dem Dateinamen."""
treffer = _NAME.match(datei.stem)
if not treffer:
return None, None
tag, uhrzeit, rest = treffer.groups()
zeit = None
if tag and uhrzeit:
try:
zeit = datetime.strptime(f"{tag}{uhrzeit}", "%Y%m%d%H%M%S")
except ValueError:
zeit = None
nummer = rest.strip() or None
if nummer in {"unbekannt", "anonymous", "restricted", "0"}:
nummer = None
return nummer, zeit
def _fertig(datei: Path) -> bool:
"""Wahr, wenn die Datei nicht mehr wächst - die Aufnahme also steht."""
try:
groesse = datei.stat().st_size
except FileNotFoundError:
return False
if groesse == 0:
return False
time.sleep(RUHE_S)
try:
return datei.stat().st_size == groesse
except FileNotFoundError:
return False
def _neue_dateien(eingang: Path):
for datei in sorted(eingang.iterdir()):
if datei.is_file() and datei.suffix.lower() in ENDUNGEN and not datei.name.startswith("."):
yield datei
def verarbeite_datei(datei: Path) -> bool:
"""Schiebt eine Aufnahme durch die Pipe und räumt sie weg. True bei Erfolg."""
nummer, zeit = _nummer_und_zeit(datei)
print(f"\n▸ Neue Aufnahme: {datei.name}"
f"{f' (Rufnummer {nummer})' if nummer else ' (Rufnummer unbekannt)'}")
protokoll.schreibe("eingang", f"Neue Aufnahme {datei.name}", nummer=nummer)
try:
process_call.verarbeite(str(datei), anrufer_nummer=nummer, empfangen=zeit)
except Exception as fehler:
print(f"✖ Verarbeitung fehlgeschlagen: {fehler}", file=sys.stderr)
protokoll.schreibe("fehler", f"{datei.name}: {fehler}", nummer=nummer)
config.TELEFON_FEHLER.mkdir(parents=True, exist_ok=True)
shutil.move(str(datei), config.TELEFON_FEHLER / datei.name)
print(f" Aufnahme liegt in {config.TELEFON_FEHLER}/{datei.name} und geht nicht verloren.")
return False
config.TELEFON_VERARBEITET.mkdir(parents=True, exist_ok=True)
shutil.move(str(datei), config.TELEFON_VERARBEITET / datei.name)
return True
def durchlauf(eingang: Path) -> int:
"""Verarbeitet alles, was gerade fertig im Eingang liegt."""
zahl = 0
for datei in _neue_dateien(eingang):
if _fertig(datei):
verarbeite_datei(datei)
zahl += 1
return zahl
def main(argv: list[str] | None = None) -> int:
p = argparse.ArgumentParser(description="Eingangsordner beobachten und Anrufe verarbeiten.")
p.add_argument("--einmal", action="store_true",
help="nur einmal durchlaufen statt dauerhaft beobachten")
args = p.parse_args(argv)
eingang = config.TELEFON_EINGANG
eingang.mkdir(parents=True, exist_ok=True)
if args.einmal:
zahl = durchlauf(eingang)
print(f"{zahl} Aufnahme(n) verarbeitet.")
return 0
print(f"Beobachte {eingang} — Abbruch mit Strg-C")
protokoll.schreibe("system", "Watcher gestartet")
try:
while True:
durchlauf(eingang)
time.sleep(TAKT_S)
except KeyboardInterrupt:
print("\nBeobachtung beendet.")
return 0
if __name__ == "__main__":
raise SystemExit(main())