Files
KX-Bridge-Release/bridge_endpoints.py
viewit 6ab5c31035
All checks were successful
Testing Build / build (push) Successful in 8m44s
feat(kxgauge): add KXGauge round-display integration
Adds an optional per-printer integration with KXGauge
(https://gitea.it-drui.de/viewit/kxgauge), a small ESP32 round-face
display that shows printer status as an emotion and hotend temperature
as a color ring. KXGauge only exposes a GET-only HTTP API with no
push/websocket, so the bridge actively pushes to it from the existing
MQTT callbacks (_on_temp for the heat ring, _on_print + offline
transitions for the emotion) whenever state actually changes, with a
built-in dedupe so it doesn't spam the device every poll tick.

- kxgauge_client.py: thin synchronous HTTP client (mirrors
  spoolman_client.py's shape - called from the MQTT reader thread).
- New [kxgauge]/[kxgauge_mapping] config.ini sections; the mapping
  (kobra_state -> KXGauge emotion) is user-editable with a sane
  default and falls back per-key if only partially configured.
- Settings UI: new card under Integrations with enable/URL/target-temp
  fields, a per-state emotion mapping list, and a connection-test
  button (/api/kxgauge/test).
- Multi-printer aware: kxgauge_url/enabled/heat_peak merge per
  [printer_N] like the existing power-switch settings.

Also fixes a real bug found while testing this: _find_config_path()
resolves to the live project config/config.ini, not a sandboxed path,
so any test hitting /api/settings POST without stubbing it out will
silently overwrite the real printer config. test_settings.py already
guards against this - test_kxgauge.py now does too.
2026-08-05 15:13:43 +02:00

2689 lines
125 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""
bridge_endpoints.py - EndpointsMixin for KobraXBridge.
The bridge's own HTTP surface: KX file/profile/print/skip endpoints, printer
management (add/remove/power/restart) + version/update, file upload + print
start, camera endpoints, printer-control API (light/fan/axis/temp/ams/ace),
Moonraker DB/state, settings, log stream/download, UI/index serving, and the
gcode-script + catchall/favicon handlers. Mixed into KobraXBridge; relies on
the shared bridge state and the _json_cors/_push_status_update/_build_* spine
provided by the core class.
────────────────────────────────────────────────────────────────────────────
Copyright (C) 2026 viewit (KX-Bridge contributors)
Licensed under GPLv3 — see LICENSE in the project root. See NOTICE.md.
"""
import os
import re
import sys
import json
import copy
import time
import html
import hashlib
import pathlib
import asyncio
import logging
import subprocess
from urllib.parse import quote
try:
import config_loader as env_loader
except ImportError:
import env_loader
import aiohttp
from aiohttp import web
from gcode_meta import (
_parse_gcode_estimated_time,
_parse_gcode_layer_heights,
_extract_thumbnail,
_extract_filament_info,
)
from bridge_constants import KXGAUGE_VALID_EMOTIONS
from credentials import _kx_fetch_credentials
from camera import _find_ffmpeg
from bridge_logging import _set_verbose_http_log, _log_buffer, _log_sse_queues
log = logging.getLogger("bridge")
# Base paths (same logic as the main module): next to sys.executable in a
# PyInstaller binary, otherwise next to this file; web assets under _MEIPASS
# in a onefile binary.
_BASE = os.path.dirname(sys.executable) if getattr(sys, "frozen", False) else os.path.dirname(os.path.abspath(__file__))
_WEB_BASE = getattr(sys, "_MEIPASS", _BASE)
# Web UI: subdirectory under web/themes/<name>/index.html
_UI_THEME_NAME_RE = re.compile(r"^[a-zA-Z0-9][a-zA-Z0-9_-]{0,63}$")
# Allowed static theme files under /kx/ui/<name>
_KX_UI_ASSETS: dict[str, str] = {
"style.css": "text/css",
"app.js": "application/javascript",
}
# Files from lib/ are served based on their extension (no whitelist entry needed)
_KX_UI_LIB_TYPES: dict[str, str] = {
".js": "application/javascript",
".css": "text/css",
}
_KX_UI_TRANSLATION_RE = re.compile(r"^translations/([a-z]{2}(?:-[a-z]{2})?)\.json$")
class EndpointsMixin:
async def handle_kx_options(self, request):
return web.Response(status=204, headers=self._CORS)
async def handle_kx_files(self, request):
files = self._store.list_files()
# Backfill legacy entries without stored filament metadata
# so the dialog's left side shows GCode colors instead of AMS slots.
for f in files:
needs_refresh = not f.get("gcode_filaments")
if not needs_refresh:
try:
cached = f.get("gcode_filaments")
parsed_cached = cached if isinstance(cached, list) else json.loads(cached)
needs_refresh = any("is_used" not in item for item in (parsed_cached or []))
except Exception:
needs_refresh = True
if not needs_refresh:
continue
path = f.get("path") or ""
if not path or not os.path.isfile(path):
continue
try:
with open(path, "rb") as fh:
parsed_filaments = _extract_filament_info(fh.read())
if parsed_filaments:
f["gcode_filaments"] = json.dumps(parsed_filaments)
self._store.update_file_filaments(f["id"], parsed_filaments)
except Exception as e:
log.debug(f"Filament metadata backfill failed for {f.get('filename')}: {e}")
# Add last job status + duration per file
jobs = self._store.list_jobs(limit=500)
last_job: dict = {}
for j in reversed(jobs):
last_job[j["gcode_file_id"]] = j
for f in files:
f["web_unverified"] = bool(f.get("web_unverified"))
lj = last_job.get(f["id"])
f["last_print_status"] = lj["status"] if lj else None
f["last_print_duration"] = lj["duration_sec"] if lj else None
f["last_print_at"] = lj["started_at"] if lj else None
return self._json_cors({"result": files})
async def handle_kx_file_delete(self, request):
file_id = request.match_info["file_id"]
if self._store.delete_file(file_id):
return self._json_cors({"result": "ok"})
return self._json_cors({"error": "not found"}, status=404)
async def handle_kx_printer_files(self, request):
"""GET /kx/printer-files - lists files on the printer's OWN internal
storage (file/listLocal MQTT action), as opposed to /kx/files which
lists what the bridge itself has stored. Needed because prints
started directly from Anycubic Slicer Next (bypassing the bridge)
leave files on the printer that were previously only visible/
deletable from the printer's own display (Issue #102 context)."""
loop = asyncio.get_event_loop()
def _fetch():
return self._wait_for_file_action(
"listLocal",
lambda: self.client.publish(
"file", "listLocal",
{"page_num": 1, "page_size": 200, "path": "/"},
timeout=0,
),
timeout=8.0,
)
result = await loop.run_in_executor(None, _fetch)
if not result or result.get("code") != 200:
return self._json_cors({"error": "printer unreachable or query failed"}, status=502)
records = (result.get("data") or {}).get("records") or []
files = [r for r in records if not r.get("is_dir")]
return self._json_cors({"result": files})
async def handle_kx_printer_file_delete(self, request):
"""POST /kx/printer-files/delete - body: {"filenames": ["a.gcode", ...]}.
Single endpoint for both single and multi-select delete - the
printer's file/deleteBatch MQTT action natively accepts a list."""
try:
body = await request.json()
except Exception:
body = {}
filenames = body.get("filenames") or []
if not filenames:
return self._json_cors({"error": "no filenames given"}, status=400)
files = [{"path": "/", "filename": fn} for fn in filenames if fn]
loop = asyncio.get_event_loop()
def _delete():
return self._wait_for_file_action(
"deleteBatch",
lambda: self.client.publish(
"file", "deleteBatch",
{"root": "local", "files": files},
timeout=0,
),
timeout=8.0,
)
result = await loop.run_in_executor(None, _delete)
if not result or result.get("state") != "success":
return self._json_cors({"error": "delete failed", "detail": result}, status=502)
return self._json_cors({"result": "ok"})
async def handle_kx_printer_file_thumbnail(self, request):
"""GET /kx/printer-files/{filename}/thumbnail - fetches the embedded
GCode thumbnail for a file on the printer's own storage, via
file/fileDetails. The printer extracts and base64-encodes the
"; thumbnail begin"-block from the GCode header on demand and
returns it inline in data.file_details.thumbnail - no separate
download/presigned-URL step needed (verified live against a real
Kobra X). Cached in-memory per filename since a file's thumbnail
never changes while it exists on the printer, and re-querying on
every render/scroll would mean one MQTT roundtrip per visible card."""
filename = request.match_info.get("filename", "")
if not filename:
return self._json_cors({"error": "no filename given"}, status=400)
cached = self._printer_thumbnail_cache.get(filename)
if cached is not None:
return self._json_cors({"result": {"thumbnail": cached}})
loop = asyncio.get_event_loop()
def _fetch():
return self._wait_for_file_action(
"fileDetails",
lambda: self.client.publish(
"file", "fileDetails",
{"root": "local", "filename": filename},
timeout=0,
),
timeout=8.0,
)
result = await loop.run_in_executor(None, _fetch)
if not result or result.get("code") != 200:
return self._json_cors({"error": "printer unreachable or query failed"}, status=502)
thumb = ((result.get("data") or {}).get("file_details") or {}).get("thumbnail") or ""
self._printer_thumbnail_cache[filename] = thumb
return self._json_cors({"result": {"thumbnail": thumb}})
async def handle_kx_file_download(self, request):
file_id = request.match_info["file_id"]
f = self._store.get_file(file_id)
if not f:
return self._json_cors({"error": "not found"}, status=404)
path = f.get("path") or ""
if not path or not os.path.isfile(path):
return self._json_cors({"error": "not found"}, status=404)
filename = os.path.basename(f.get("filename") or path)
# RFC 5987: filename* with URL encoding for special chars/UTF-8,
# plus ASCII fallback (strip all " and \ from filename for the
# quoted-string-Part).
ascii_fallback = filename.encode("ascii", "replace").decode("ascii").replace('"', "").replace("\\", "")
encoded = quote(filename, safe="")
disposition = f'attachment; filename="{ascii_fallback}"; filename*=UTF-8\'\'{encoded}'
return web.FileResponse(path, headers={"Content-Disposition": disposition})
async def handle_kx_file_verify(self, request):
file_id = request.match_info["file_id"]
if self._store.clear_web_unverified(file_id):
return self._json_cors({"result": "ok"})
return self._json_cors({"error": "not found"}, status=404)
async def handle_kx_filament_slots(self, request):
slots = []
for i, s in enumerate(self._ams_slots):
gidx = int(s.get("global_index", i))
# Stale-profile guard: only show the override while its material
# family matches the loaded AMS material (else slot has no brand).
profile = self._effective_slot_profile(gidx, s.get("type", ""))
slots.append({
"slot_index": gidx,
"material": s.get("type", ""),
"color_hex": "#{:02X}{:02X}{:02X}".format(*s.get("color", [0,0,0])[:3]),
"status": "loaded" if s.get("status") == 5 else "empty",
"nozzle_temp": 0,
# Current user override from config.ini [filament_profiles]
# - (vendor,name) is unique, id is only a hint.
"filament_id": profile.get("id", ""),
"filament_vendor": profile.get("vendor", ""),
"filament_name": profile.get("name", ""),
})
return self._json_cors({"result": slots})
async def handle_kx_filament_profiles(self, request):
"""Returns the static list of OrcaSlicer filament profiles
(from bridge/data/orca_filaments.json - produced by the generator script
tools/gen_orca_filament_list.py erzeugt).
Optional Filter via ?type=PLA / ?vendor=Polymaker.
The frontend uses this for the slot profile dropdown.
"""
type_filter = request.rel_url.query.get("type", "").upper().strip()
vendor_filter = request.rel_url.query.get("vendor", "").strip()
profiles = self._load_orca_filaments()
if type_filter:
profiles = [p for p in profiles if p.get("type", "").upper() == type_filter]
if vendor_filter:
profiles = [p for p in profiles if p.get("vendor", "") == vendor_filter]
return self._json_cors({"result": profiles})
async def handle_kx_filament_profiles_user_list(self, request):
"""GET /kx/filament/profiles/user - only the user-imported profiles,
for the settings tab (management with delete buttons)."""
path = self._orca_filaments_user_path()
if not os.path.isfile(path):
return self._json_cors({"result": []})
try:
with open(path, encoding="utf-8") as f:
user_profiles = json.load(f) or []
except Exception:
user_profiles = []
return self._json_cors({"result": user_profiles})
async def handle_kx_filament_profiles_import(self, request):
"""POST /kx/filament/profiles/user - multipart upload with one
ZIP file or multiple `.json` files from
~/.config/OrcaSlicer/user/<id>/filament/.
Existing user profiles with the same (vendor, name) key are
overwritten. Parsed profiles use the same schema as
orca_filaments.json (id, name, vendor, type, color)."""
import io, zipfile
from orca_filaments import parse_profile_bytes
added: list[dict] = []
skipped: int = 0
# System index for inherits resolution: user profiles reference
# System-Parents via "inherits" (z.B. "Generic PLA @System"). Damit
# we can pull filament_id/vendor/type/color from the system parent
# when the user profile does not set them itself.
sys_idx = [p for p in self._load_orca_filaments() if not p.get("is_user")]
try:
reader = await request.multipart()
except Exception:
return self._json_cors({"error": "expected multipart"}, status=400)
async for part in reader:
if part.name not in ("file", "files", "upload"):
continue
blob = await part.read()
fn = (part.filename or "").lower()
if fn.endswith(".zip"):
try:
with zipfile.ZipFile(io.BytesIO(blob)) as zf:
for inner in zf.namelist():
if not inner.lower().endswith(".json"):
continue
try:
with zf.open(inner) as zf_in:
p = parse_profile_bytes(zf_in.read(), source_name=inner, system_index=sys_idx)
except Exception:
skipped += 1
continue
if p:
added.append(p)
else:
skipped += 1
except zipfile.BadZipFile:
return self._json_cors({"error": "bad zip"}, status=400)
elif fn.endswith(".json"):
p = parse_profile_bytes(blob, source_name=fn, system_index=sys_idx)
if p:
added.append(p)
else:
skipped += 1
if not added:
return self._json_cors({"result": "ok", "added": 0, "skipped": skipped})
# Merge with existing user JSON (same (vendor,name) -> replace)
path = self._orca_filaments_user_path()
existing: list[dict] = []
if os.path.isfile(path):
try:
with open(path, encoding="utf-8") as f:
existing = json.load(f) or []
except Exception:
existing = []
by_key = {(p.get("vendor"), p.get("name")): p for p in existing}
for p in added:
by_key[(p.get("vendor"), p.get("name"))] = p
merged = sorted(by_key.values(), key=lambda x: (x.get("vendor",""), x.get("name","")))
try:
with open(path, "w", encoding="utf-8") as f:
json.dump(merged, f, indent=2, ensure_ascii=False)
f.write("\n")
except Exception as e:
return self._json_cors({"error": f"write failed: {e}"}, status=500)
self._invalidate_filaments_cache()
return self._json_cors({"result": "ok",
"added": len(added),
"skipped": skipped,
"total_user": len(merged)})
async def handle_kx_filament_profiles_user_delete(self, request):
"""DELETE /kx/filament/profiles/user - deletes either a single
entry (?vendor=...&name=...) or all when no query is given."""
vendor = request.rel_url.query.get("vendor", "").strip()
name = request.rel_url.query.get("name", "").strip()
path = self._orca_filaments_user_path()
if not os.path.isfile(path):
return self._json_cors({"result": "ok", "removed": 0})
try:
with open(path, encoding="utf-8") as f:
existing = json.load(f) or []
except Exception:
existing = []
before = len(existing)
if vendor and name:
existing = [p for p in existing
if not (p.get("vendor") == vendor and p.get("name") == name)]
else:
existing = []
try:
with open(path, "w", encoding="utf-8") as f:
json.dump(existing, f, indent=2, ensure_ascii=False)
f.write("\n")
except Exception as e:
return self._json_cors({"error": str(e)}, status=500)
self._invalidate_filaments_cache()
return self._json_cors({"result": "ok",
"removed": before - len(existing),
"total_user": len(existing)})
def _find_orca_filaments_json(self) -> str | None:
"""Finds the static JSON file. Sits next to web/ under _WEB_BASE/data/
— in allen 3 Deployment-Modi:
• Dev: bridge/data/orca_filaments.json
* Docker: /app/data/orca_filaments.json (static in the image, NOT the
volume data/ holding runtime state - see Dockerfile)
• Onefile: sys._MEIPASS/data/orca_filaments.json
When the volume-mounted /app/data/ shadows the static data, a copy
also sits under _WEB_BASE/data/ (= /app/ in Docker = the same path).
On conflict: second lookup under ../bridge/data/ as a fallback for dev setups."""
candidates = [
# Docker: COPY bridge/data -> /app/static/ (data/ is a volume -> shadowed)
os.path.join(_WEB_BASE, "static", "orca_filaments.json"),
os.path.join(_WEB_BASE, "data", "orca_filaments.json"),
]
here = os.path.dirname(os.path.abspath(__file__))
candidates.append(os.path.join(here, "data", "orca_filaments.json"))
candidates.append(os.path.join(here, "..", "bridge", "data", "orca_filaments.json"))
for c in candidates:
if os.path.isfile(c):
return c
return None
async def handle_kx_filament_slot_profile(self, request):
"""POST /kx/filament/slots/<idx>/profile - saves or deletes
a user override mapping for a single AMS slot.
The primary selector is (vendor, name) - the ID is not unique in the Orca
data model (136 profiles share e.g. 'OGFL99'). The ID is looked up
from orca_filaments.json on save and carried along as a hint
for OrcaSlicer's `tray_info_idx`.
Body: {"vendor": "Polymaker", "name": "PolyTerra PLA"}
{"vendor": "", "name": ""} → Mapping entfernen
(Backwards compat: {"id":..., "vendor":...} is accepted,
but `name` has been the primary selector since v0.9.18.)
"""
try:
slot_idx = int(request.match_info.get("idx", "-1"))
except ValueError:
return self._json_cors({"error": "bad slot index"}, status=400)
if slot_idx < 0:
return self._json_cors({"error": "bad slot index"}, status=400)
try:
data = await request.json()
except Exception:
data = {}
new_vendor = (data.get("vendor") or "").strip()
new_name = (data.get("name") or "").strip()
new_id = (data.get("id") or "").strip() # Backwards-Kompat-Hint
if new_vendor and new_name:
# Look up the ID from JSON (not from the request body, which could
# be stale or a generic fallback).
looked_up_id = self._lookup_filament_id(new_vendor, new_name)
self._filament_profiles[slot_idx] = {
"vendor": new_vendor,
"name": new_name,
"id": looked_up_id or new_id,
}
else:
self._filament_profiles.pop(slot_idx, None)
# Persistieren in config.ini
try:
import config_loader as _cl
_cl.save_filament_profiles(self._filament_profiles, self._printer_id)
except Exception as e:
log.warning(f"save_filament_profiles failed: {e}")
return self._json_cors({"error": str(e)}, status=500)
entry = self._filament_profiles.get(slot_idx, {})
return self._json_cors({"result": "ok",
"slot_index": slot_idx,
"vendor": entry.get("vendor", ""),
"name": entry.get("name", ""),
"id": entry.get("id", "")})
async def handle_kx_visible_vendors(self, request):
"""GET/POST /kx/filament/visible_vendors — Vendor-Sichtbarkeitsfilter
for the slot profile dropdown (Issue #41 option A).
GET → {"result": ["Polymaker", "eSUN", ...]}
POST {"vendors": [...]} → speichert in config.ini [filament_profiles]
visible_vendors. Empty list = all visible. NO bridge restart
needed (display filter only)."""
if request.method == "POST":
try:
data = await request.json()
except Exception:
data = {}
vendors = data.get("vendors") or []
if not isinstance(vendors, list):
return self._json_cors({"error": "vendors must be a list"}, status=400)
self._visible_vendors = [str(v).strip() for v in vendors if str(v).strip()]
try:
import config_loader as _cl
_cl.save_visible_vendors(self._visible_vendors, self._printer_id)
except Exception as e:
log.warning(f"save_visible_vendors failed: {e}")
return self._json_cors({"error": str(e)}, status=500)
return self._json_cors({"result": self._visible_vendors})
def _load_orca_filaments(self) -> list[dict]:
"""Loads system + user profiles from the cache. System profiles come
from bridge/data/orca_filaments.json (image-embedded), user profiles
from <KX_DATA_DIR>/orca_filaments.user.json (volume-persistent -
survives image updates). User profiles get an `is_user: True`
flag so the frontend can mark them."""
if getattr(self, "_orca_filaments_cache", None) is not None:
return self._orca_filaments_cache
merged: list[dict] = []
# System
sys_path = self._find_orca_filaments_json()
if sys_path and os.path.isfile(sys_path):
try:
with open(sys_path, encoding="utf-8") as f:
merged.extend(json.load(f) or [])
except Exception as e:
log.warning(f"orca_filaments.json read error: {e}")
# User
usr_path = self._orca_filaments_user_path()
if usr_path and os.path.isfile(usr_path):
try:
with open(usr_path, encoding="utf-8") as f:
for p in (json.load(f) or []):
p["is_user"] = True
merged.append(p)
except Exception as e:
log.warning(f"orca_filaments.user.json read error: {e}")
self._orca_filaments_cache = merged
return self._orca_filaments_cache
def _orca_filaments_user_path(self) -> str:
"""Path to the user profiles JSON. Lives in the volume mount (KX_DATA_DIR)
so image updates do not destroy the data."""
data_dir = os.environ.get("KX_DATA_DIR") or os.path.join(_WEB_BASE, "data")
os.makedirs(data_dir, exist_ok=True)
return os.path.join(data_dir, "orca_filaments.user.json")
def _invalidate_filaments_cache(self):
self._orca_filaments_cache = None
def _lookup_filament_id(self, vendor: str, name: str) -> str:
"""Looks up the filament_id for a (vendor,name) tuple in
orca_filaments.json. Returns '' when not found."""
for p in self._load_orca_filaments():
if p.get("vendor") == vendor and p.get("name") == name:
return p.get("id", "")
return ""
async def handle_kx_history(self, request):
limit = int(request.rel_url.query.get("limit", 50))
offset = int(request.rel_url.query.get("offset", 0))
jobs = self._store.list_jobs(limit=limit, offset=offset)
return self._json_cors({"result": jobs})
async def handle_kx_file_objects(self, request):
"""Returns the object list + optional SVG for a file.
GET /kx/files/{id}/objects → {"names": [...], "svg_b64": "..."}
If the file has no objects yet (old entry): querying file/fileDetails
from the printer and awaiting the response is the frontend's job
(reload after upload). Only return the database state here.
"""
fid = request.match_info.get("id", "")
f = self._store.get_file(fid)
if not f:
return self._json_cors({"error": "file not found"}, status=404)
try:
names = json.loads(f.get("objects_skip_parts") or "[]")
except Exception:
names = []
# No objects in the store yet (fresh Orca/web upload): actively request
# file/fileDetails from the printer once. _on_file() backfills the store,
# the frontend polls this endpoint and receives the list on the next
# attempt (Issue #57 - skip parity outside the file browser too).
if not names:
fn = f.get("filename") or ""
if fn:
try:
self.client.publish("file", "fileDetails",
{"root": "local", "filename": fn}, timeout=0)
except Exception as e:
log.debug(f"fileDetails request failed: {e}")
return self._json_cors({
"result": {
"names": names,
"svg_b64": f.get("svg_image") or "",
}
})
async def handle_kx_skip(self, request):
"""Trigger a mid-print skip.
POST /kx/skip body={"names": ["..", ".."]}
"""
try:
body = await request.json()
except Exception:
return self._json_cors({"error": "invalid json"}, status=400)
names = body.get("names") or []
if not isinstance(names, list) or not all(isinstance(n, str) for n in names):
return self._json_cors({"error": "names must be list[str]"}, status=400)
try:
loop = asyncio.get_event_loop()
await loop.run_in_executor(None, lambda: self.client.skip_objects(names))
except Exception as e:
return self._json_cors({"error": str(e)}, status=502)
return self._json_cors({"result": "ok", "names": names})
def _build_skip_state_result(self) -> dict:
"""Builds the combined skip state for UI endpoints."""
filename = self._state.get("filename", "")
all_objects: list[str] = []
svg = ""
if filename:
try:
f = self._store.get_file_by_name(filename)
if f:
all_objects = json.loads(f.get("objects_skip_parts") or "[]")
svg = f.get("svg_image") or ""
except Exception as e:
log.warning(f"skip_state lookup failed: {e}")
return {
"objects": all_objects,
"skipped": list(self._skip_state.get("skipped", [])),
"svg_b64": svg,
"ts": self._skip_state.get("ts", 0),
"filename": filename,
}
async def handle_kx_skip_query(self, request):
"""Re-request the print object list from the printer.
POST /kx/skip/query → triggert skip/query_obj, wartet kurz auf den
async skip/report and returns the merged skip state.
"""
prev_ts = int(self._skip_state.get("ts", 0) or 0)
try:
loop = asyncio.get_event_loop()
await loop.run_in_executor(None, lambda: self.client.query_skip_objects())
except Exception as e:
return self._json_cors({"error": str(e)}, status=502)
deadline = time.time() + 1.5
while time.time() < deadline:
if int(self._skip_state.get("ts", 0) or 0) > prev_ts:
break
await asyncio.sleep(0.1)
return self._json_cors({"result": self._build_skip_state_result()})
async def handle_kx_skip_state(self, request):
"""Aktueller Skip-State.
Kombiniert:
- Full object list: from the GCode store, matched via the currently
running filename (file/report at print start populated the list).
skip/query_obj only returns the already-skipped ones,
not the full list.
- Skipped: from self._skip_state (updated by skip/report).
"""
return self._json_cors({"result": self._build_skip_state_result()})
async def handle_kx_printers(self, request):
# Collect active printers (with IP)
active = [(pid, br) for pid, br in self._all_bridges.items()
if (br._args.printer_ip or "").strip()]
# Host for bridge_url: keep the browser view, but never export "localhost" -
# otherwise browser fetches fail when the UI is opened via the LAN IP.
host = request.host.split(":")[0]
if host in ("localhost", "127.0.0.1", "::1", "0.0.0.0"):
host = ""
out = []
for pid, br in active:
port = getattr(br._args, "port", 7125)
# Only set a concrete bridge_url for multi-printer setups (cross-instance fetch).
# Single printer: empty bridge_url -> JS uses relative paths (same origin as the UI).
bridge_url = ""
if len(active) > 1 and host:
bridge_url = f"http://{host}:{port}"
out.append({
"id": pid,
"name": br._state.get("printer_name") or f"Drucker {pid}",
"bridge_url": bridge_url,
"printer_ip": br._args.printer_ip,
"device_id": br._args.device_id or "",
"has_power_control": bool(
(getattr(br._args, "power_on_url", "") or "").strip()
or (getattr(br._args, "power_off_url", "") or "").strip()
),
})
return self._json_cors({"result": out})
async def handle_kx_printer_power(self, request):
"""Toggles an external smart plug (e.g. Tasmota) for a printer that
has no MQTT-level power-off/standby command of its own (Issue #103).
Just fires a plain HTTP GET at the configured power_on_url/power_off_url -
works for Tasmota's cmnd=Power%20on/off style URLs and any other
switch that exposes a GET-triggered on/off endpoint."""
pid = str(request.match_info.get("pid", "")).strip()
br = self._all_bridges.get(pid)
if br is None:
return self._json_cors({"error": "unknown printer id"}, status=404)
try:
body = await request.json()
except Exception:
body = {}
action = str(body.get("action", "")).lower()
if action not in ("on", "off"):
return self._json_cors({"error": "action must be 'on' or 'off'"}, status=400)
url = getattr(br._args, f"power_{action}_url", "") or ""
if not url:
return self._json_cors({"error": f"no power_{action}_url configured"}, status=400)
try:
async with aiohttp.ClientSession() as session:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=5)) as resp:
ok = resp.status == 200
except Exception as e:
return self._json_cors({"error": f"power switch unreachable: {e}"}, status=502)
return self._json_cors({"result": "ok" if ok else "error", "status": "on" if action == "on" else "off"})
async def handle_kx_printer_power_status(self, request):
"""Queries the configured smart plug for its current on/off state.
Tries to parse a Tasmota-style {"POWER":"ON"/"OFF"} JSON body first,
falls back to a plain substring search for "ON"/"OFF" in the raw
response so other switch firmwares with a simpler status endpoint
still work."""
pid = str(request.match_info.get("pid", "")).strip()
br = self._all_bridges.get(pid)
if br is None:
return self._json_cors({"error": "unknown printer id"}, status=404)
url = getattr(br._args, "power_status_url", "") or ""
if not url:
return self._json_cors({"error": "no power_status_url configured"}, status=400)
try:
async with aiohttp.ClientSession() as session:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=5)) as resp:
text = await resp.text()
except Exception as e:
return self._json_cors({"error": f"power switch unreachable: {e}"}, status=502)
state = "unknown"
try:
data = json.loads(text)
power = str(data.get("POWER", "")).upper()
if power in ("ON", "OFF"):
state = power.lower()
except Exception:
pass
if state == "unknown":
up = text.upper()
if "ON" in up and "OFF" not in up:
state = "on"
elif "OFF" in up:
state = "off"
return self._json_cors({"state": state})
async def handle_kx_kxgauge_test(self, request):
"""Tests connectivity to the configured KXGauge display (GET /status)
for the "Verbindung testen" button in Settings. Uses the URL from the
request body if given (unsaved changes), else the configured one."""
try:
body = await request.json()
except Exception:
body = {}
url = str(body.get("url", "")).strip() or (getattr(self._args, "kxgauge_url", "") or "")
if not url:
return self._json_cors({"error": "no KXGauge URL configured"}, status=400)
try:
async with aiohttp.ClientSession() as session:
async with session.get(f"{url.rstrip('/')}/status", timeout=aiohttp.ClientTimeout(total=5)) as resp:
ok = resp.status == 200
data = await resp.json(content_type=None)
except Exception as e:
return self._json_cors({"error": f"KXGauge unreachable: {e}"}, status=502)
if not ok:
return self._json_cors({"error": f"KXGauge returned HTTP {resp.status}"}, status=502)
return self._json_cors({"result": "ok", "emotion": (data or {}).get("emotion", "")})
async def handle_kx_print(self, request):
"""Print start from the GCode store with optional filament assignments."""
try:
body = await request.json()
except Exception:
return self._json_cors({"error": "invalid json"}, status=400)
file_id = body.get("file_id")
if not file_id:
return self._json_cors({"error": "file_id required"}, status=400)
gcode_file = self._store.get_file(file_id)
if not gcode_file:
return self._json_cors({"error": "file not found"}, status=404)
# filament_assignments: [{slot_index, material, color_hex}, …]
assignments = body.get("filament_assignments")
# excluded_objects: ["name1","name2",...] Pre-Print Skip (v0.9.10)
excluded_objects = body.get("excluded_objects") or []
if not isinstance(excluded_objects, list):
excluded_objects = []
if assignments:
ams_box_mapping, unused_count, invalid_count = self._build_assigned_ams_box_mapping(assignments)
if unused_count:
log.debug(f"Skipped {unused_count} unused filament assignment(s) for mode={self._filament_mode}")
if invalid_count:
log.warning(f"Ignored {invalid_count} unusable filament assignment(s) for mode={self._filament_mode}")
if not ams_box_mapping:
return self._json_cors({"error": "no usable filament assignments for current filament mode"}, status=400)
else:
# No dialog -> all occupied slots as with a normal upload print
ams_box_mapping = self._build_auto_ams_box_mapping()
auto_leveling = int(body.get("auto_leveling", getattr(self._args, "auto_leveling", 1)))
filename = gcode_file["filename"]
file_path = gcode_file["path"]
# Serve the file via the internal serve endpoint
url = f"http://localhost:{self._args.port}/serve/{os.path.basename(file_path)}"
payload = self._build_print_payload(
filename, url, "", gcode_file.get("size_bytes", 0),
ams_box_mapping=ams_box_mapping,
auto_leveling=auto_leveling,
excluded_objects=excluded_objects,
)
self._reset_skip_state(excluded_objects)
log.info(f"KX store print start: {filename} ams={len(ams_box_mapping)} slots assignments={bool(assignments)} excluded={len(excluded_objects)}")
loop = asyncio.get_event_loop()
result = await loop.run_in_executor(
None, lambda: self.client.publish("print", "start", payload, timeout=15.0)
)
if result is None:
return self._json_cors({"error": "no response from printer"}, status=504)
if excluded_objects:
loop.run_in_executor(None, lambda: self._apply_preprint_skip_after_start(excluded_objects))
# Start the job in the history
self._current_job_id = self._store.start_job(
gcode_file_id=gcode_file["id"],
printer_id=getattr(self._args, "device_id", "unknown"),
filament_assignments=assignments,
)
self._current_job_filename = filename
return self._json_cors({"result": "ok", "filename": filename})
# -------------------------------------------------------------------------
# HTTP handlers
# -------------------------------------------------------------------------
async def handle_file_upload(self, request):
log.info(f"Upload-Request: {request.method} {request.path_qs} CT={request.headers.get('Content-Type','')[:60]}")
ct = request.headers.get("Content-Type", "")
if "multipart" not in ct:
return web.json_response({"error": "expected multipart"}, status=400)
auto_print = False
web_upload = False
reader = await request.multipart()
file_data = None
remote_filename = self._last_uploaded_file or "upload.gcode"
async for part in reader:
if part.name in ("file", "gcode", "upload_file"):
remote_filename = part.filename or remote_filename
file_data = await part.read()
log.info(f"Multipart-Feld '{part.name}': {remote_filename} ({len(file_data)} bytes)")
elif part.name == "path":
val = (await part.read()).decode("utf-8", errors="replace").strip()
if val:
remote_filename = val
elif part.name == "print":
val = (await part.read()).decode("utf-8", errors="replace").strip().lower()
auto_print = val == "true"
elif part.name == "web_upload":
val = (await part.read()).decode("utf-8", errors="replace").strip().lower()
web_upload = val == "true"
else:
log.debug(f"Unbekanntes Multipart-Feld: {part.name}")
if not file_data:
return web.json_response({"error": "no file received"}, status=400)
# Only allow printable files (Issue #59) - the Kobra X accepts
# only .gcode and .bgcode; .3mf uploads are not processed by the
# printer and are therefore rejected (Issue #59, @gangoke).
_allowed_ext = (".gcode", ".bgcode")
_fn_lower = (remote_filename or "").lower()
if not _fn_lower.endswith(_allowed_ext):
log.warning(f"Upload rejected (not GCode): {remote_filename}")
return web.json_response(
{"error": f"only GCode files allowed ({', '.join(_allowed_ext)})"},
status=400,
)
file_md5 = hashlib.md5(file_data).hexdigest()
file_size = len(file_data)
# Read slicer time estimate + thumbnail from GCode
est_time = _parse_gcode_estimated_time(file_data)
self._state["slicer_time"] = est_time
thumbnail_b64 = _extract_thumbnail(file_data)
gcode_filaments = _extract_filament_info(file_data)
layer_h, first_h = _parse_gcode_layer_heights(file_data)
self._state["layer_height"] = layer_h
self._state["first_layer_height"] = first_h
# Persist the file in the GCode store
self._store.save_file(
file_id=file_md5,
filename=remote_filename,
data=file_data,
est_time_sec=est_time,
thumbnail_b64=thumbnail_b64,
gcode_filaments=gcode_filaments or None,
web_unverified=web_upload,
layer_height=layer_h,
first_layer_height=first_h,
)
serve_path = os.path.join(self._serve_dir_path, os.path.basename(remote_filename))
del file_data # free RAM
self._last_uploaded_file = remote_filename
log.info(f"Upload: {remote_filename} ({file_size} bytes) md5={file_md5} -> store + printer")
# Upload the file to the printer via HTTP (serve_path is already on disk)
upload_url = self._state.get("upload_url") or None
loop = asyncio.get_event_loop()
try:
result = await loop.run_in_executor(
None, self.client.upload_gcode, serve_path, remote_filename, upload_url
)
except Exception as e:
log.error(f"Upload failed: {e}")
return web.json_response({"error": str(e)}, status=500)
log.info(f"Upload successful: {result}")
# Start the print with the full payload (incl. serve URL + md5 + size)
serve_url = f"http://{request.host}/serve/{remote_filename}"
# print=true in the multipart form (Moonraker) or query string -> start print
# print=false or missing -> upload only
if not auto_print:
auto_print = request.rel_url.query.get("print", "false").lower() == "true"
# Always request the thumbnail (printer responds async with file/report)
self._thumbnail_b64 = ""
self.client.publish("file", "fileDetails", {"root": "local", "filename": remote_filename}, timeout=0)
self._state["last_upload_url"] = serve_url
self._state["last_upload_md5"] = file_md5
self._state["last_upload_size"] = file_size
if auto_print:
mismatch = self._check_filament_mismatch(gcode_filaments)
if mismatch:
log.info(f"Upload+print blocked - filament mismatch: {mismatch}")
self._state["file_ready"] = remote_filename
self._state["filament_mismatch"] = mismatch
return self._octoprint_upload_response(
request, remote_filename,
extra={"filament_mismatch": True, "mismatch_details": mismatch},
)
log.info(f"Upload+Print (print=true): {remote_filename}")
self._state["file_ready"] = ""
loop = asyncio.get_event_loop()
loop.run_in_executor(None, lambda: self._start_print(remote_filename, serve_url, file_md5, file_size, gcode_filaments=gcode_filaments))
else:
log.info(f"Upload only (print=false): {remote_filename}")
self._state["file_ready"] = remote_filename
return self._octoprint_upload_response(request, remote_filename)
@staticmethod
def _octoprint_upload_response(request, remote_filename: str, extra: dict | None = None):
"""OctoPrint-compatible upload response (OrcaSlicer evaluates refs)."""
body = {
"done": True,
"files": {
"local": {
"name": remote_filename,
"origin": "local",
"path": remote_filename,
"refs": {
"download": f"http://{request.host}/api/files/local/{remote_filename}",
"resource": f"http://{request.host}/api/files/local/{remote_filename}",
}
}
},
"result": {
"item": {"path": remote_filename, "root": "gcodes"},
"action": "create_file",
}
}
if extra:
body.update(extra)
return web.json_response(body, status=201)
def _check_filament_mismatch(self, gcode_filaments: list | None) -> list[dict] | None:
"""Compares GCode filaments (is_used=True) with currently occupied AMS slots.
Returns a list of mismatch entries when at least one used
GCode slot has no matching material in the AMS - otherwise None.
Only triggered when AMS data is present (at least 1 occupied slot)."""
if not gcode_filaments:
return None
slots = self._ams_slots or []
occupied = {s["global_index"]: s for s in slots if s.get("type") and s.get("status") == 5}
if not occupied:
return None
mismatches = []
for f in gcode_filaments:
if not f.get("is_used"):
continue
idx = int(f.get("slot_index", -1))
gcode_mat = (f.get("material") or "").upper().strip()
if not gcode_mat:
continue
slot = occupied.get(idx)
if slot is None:
mismatches.append({
"slot_index": idx,
"gcode_material": gcode_mat,
"ams_material": None,
"reason": "empty",
})
else:
ams_mat = (slot.get("type") or "").upper().strip()
if ams_mat and ams_mat != gcode_mat:
mismatches.append({
"slot_index": idx,
"gcode_material": gcode_mat,
"ams_material": ams_mat,
"reason": "mismatch",
})
return mismatches if mismatches else None
def _build_print_payload(self, filename: str, url: str, md5: str, filesize: int,
ams_box_mapping: list, auto_leveling: int,
excluded_objects: list | None = None,
ai_type: int = 1, timelapse_type: int = 64) -> dict:
"""Builds the complete print/start MQTT payload. Single source for all
three print start paths (upload, KX store, Moonraker API)."""
return {
"taskid": "-1",
"url": url,
"filename": filename,
"md5": md5,
"filepath": None,
"filetype": 1,
"project_type": 1,
"filesize": filesize,
"ams_settings": {
"use_ams": len(ams_box_mapping) > 0,
"ams_box_mapping": ams_box_mapping,
},
"task_settings": {
"auto_leveling": auto_leveling,
"vibration_compensation": getattr(self._args, "vibration_compensation", 0),
"flow_calibration": 0,
"dry_mode": 0,
"ai_settings": {"status": 0, "count": 0, "type": ai_type},
"timelapse": {"status": 0, "count": 0, "type": timelapse_type},
"drying_settings": {"status": 0, "target_temp": 0, "duration": 0, "remain_time": 0},
"model_objects_skip_parts": excluded_objects or [],
},
}
def _reset_skip_state(self, excluded_objects: list | None = None):
"""Resets the skip state before a print start. The UI is marked as
"skipped" only after real printer confirmation."""
self._skip_state = {"skipped": [], "ts": int(time.time())}
if excluded_objects:
self._pending_preprint_skip = [str(n) for n in excluded_objects if isinstance(n, str) and n]
self._pending_preprint_skip_deadline = time.time() + 12.0
else:
self._pending_preprint_skip = []
self._pending_preprint_skip_deadline = 0.0
def _start_print(self, filename: str, url: str = "", md5: str = "", filesize: int = 0,
gcode_filaments: list | None = None):
self._state["file_ready"] = ""
loaded = self._select_loaded_slots_for_print(warn_on_empty_default=True)
# Only map the paints ACTUALLY used in the GCode to slots. OrcaSlicer
# writes all configured filaments into the header (filament_colour=...;...;...),
# but often uses only one (e.g. single color -> only T3). If we mapped all
# occupied slots, the printer would expect all colors and block
# when another (unused) slot is empty. The used paint indices
# liefert _extract_filament_info via is_used (echte T<n>-Tool-Changes).
used_paint_indices = None
if gcode_filaments:
used = [int(f["slot_index"]) for f in gcode_filaments
if f.get("is_used") and "slot_index" in f]
if used:
used_paint_indices = set(used)
if used_paint_indices is not None:
# GCode-Paint-Index N entspricht AMS-Slot N (global_index). Nur belegte
# used slots; used-but-unloaded -> a warning may follow later.
loaded = [(gidx, s) for (gidx, s) in loaded if gidx in used_paint_indices]
ams_box_mapping = self._build_auto_ams_box_mapping(loaded_slots=loaded)
log.debug(f"AMS slots: {len(loaded)} mapped (used paints: {used_paint_indices}) -> {[i for i, _ in loaded]}")
payload = self._build_print_payload(
filename, url, md5, filesize,
ams_box_mapping=ams_box_mapping,
auto_leveling=getattr(self._args, "auto_leveling", 1),
)
log.info(f"print/start → {filename} url={url} ams={len(ams_box_mapping)} slots mode={self._filament_mode}")
result = self.client.publish("print", "start", payload, timeout=15.0)
if result:
log.info(f"Print start confirmed: state={result.get('state')}")
else:
log.warning("Print start: no response from printer")
def _theme_index_path(self) -> str:
return os.path.join(_WEB_BASE, "web", "themes", self._ui_theme, "index.html")
def _load_index_template_cached(self) -> str:
path = self._theme_index_path()
mtime = os.path.getmtime(path)
key = (path, mtime)
if self._index_tpl_cache is not None and self._index_tpl_cache_key == key:
return self._index_tpl_cache
with open(path, "r", encoding="utf-8") as f:
self._index_tpl_cache = f.read()
self._index_tpl_cache_key = key
return self._index_tpl_cache
def _ui_asset_cache_buster(self) -> str:
base = os.path.join(_WEB_BASE, "web", "themes", self._ui_theme)
mt = 0.0
for fn in ("index.html", "style.css", "app.js"):
try:
mt = max(mt, os.path.getmtime(os.path.join(base, fn)))
except OSError:
pass
return str(int(mt)) if mt else "0"
async def handle_print_start(self, request):
try:
body = await request.json()
except Exception:
body = {}
filename = (request.rel_url.query.get("filename")
or body.get("filename")
or self._last_uploaded_file)
if not filename:
return web.json_response({"error": "no filename"}, status=400)
log.info(f"Starting print: {filename}")
# Optional slot selection from the filament dialog
filament_assignments = body.get("filament_assignments")
# Pre-Print Skip (v0.9.10)
excluded_objects = body.get("excluded_objects") or []
if not isinstance(excluded_objects, list):
excluded_objects = []
auto_leveling = int(body.get("auto_leveling", getattr(self._args, "auto_leveling", 1)))
url = self._state.get("last_upload_url", "")
filesize = self._state.get("last_upload_size", 0)
md5 = self._state.get("last_upload_md5", "")
if filament_assignments is not None:
# Explicit slot assignment from the filament dialog
ams_box_mapping, unused_count, invalid_count = self._build_assigned_ams_box_mapping(filament_assignments)
if unused_count:
log.debug(f"Skipped {unused_count} unused filament assignment(s) for mode={self._filament_mode}")
if invalid_count:
log.warning(f"Ignored {invalid_count} unusable filament assignment(s) for mode={self._filament_mode}")
if not ams_box_mapping:
return web.json_response({"error": "no usable filament assignments for current filament mode"}, status=400)
else:
# Dashboard reprint: load gcode_filaments from DB so the used_paint_indices
# filter applies and empty/shifted slots are not mapped incorrectly.
gcode_filaments = None
try:
db_file = self._store.get_file_by_name(filename)
if db_file and db_file.get("gcode_filaments"):
gcode_filaments = json.loads(db_file["gcode_filaments"])
except Exception as e:
log.warning(f"Could not load cached gcode_filaments for {filename}: {e} "
"- slot mapping falls back to all occupied slots")
# Set the pre-print skip before _start_print is called
self._reset_skip_state(excluded_objects)
log.info(f"print/start api=1 mode={self._filament_mode} assignments=False gcode_filaments={gcode_filaments is not None}")
loop = asyncio.get_event_loop()
loop.run_in_executor(None, lambda: self._start_print(
filename, url, md5, filesize,
gcode_filaments=gcode_filaments,
))
return web.json_response({"result": "ok"})
payload = self._build_print_payload(
filename, url, md5, filesize,
ams_box_mapping=ams_box_mapping,
auto_leveling=auto_leveling,
excluded_objects=excluded_objects,
ai_type=0, timelapse_type=0,
)
self._reset_skip_state(excluded_objects)
log.info(
f"print/start api=1 mode={self._filament_mode} "
f"ams={len(ams_box_mapping)} slots assignments=True"
)
loop = asyncio.get_event_loop()
result = await loop.run_in_executor(
None, lambda: self.client.publish("print", "start", payload, timeout=15.0)
)
if result is None:
return web.json_response({"error": "no response from printer"}, status=504)
if excluded_objects:
loop.run_in_executor(None, lambda: self._apply_preprint_skip_after_start(excluded_objects))
return web.json_response({"result": "ok"})
async def handle_print_pause(self, request):
loop = asyncio.get_event_loop()
taskid = self._state.get("taskid", "-1")
await loop.run_in_executor(None, lambda: self.client.pause_print(taskid))
return web.json_response({"result": "ok"})
async def handle_print_resume(self, request):
loop = asyncio.get_event_loop()
taskid = self._state.get("taskid", "-1")
await loop.run_in_executor(None, lambda: self.client.resume_print(taskid))
return web.json_response({"result": "ok"})
async def handle_print_cancel(self, request):
loop = asyncio.get_event_loop()
taskid = self._state.get("taskid", "-1")
await loop.run_in_executor(None, lambda: self.client.stop_print(taskid))
return web.json_response({"result": "ok"})
async def handle_api_file_ready_clear(self, request):
self._state["file_ready"] = ""
self._state["filament_mismatch"] = None
self._thumbnail_b64 = ""
self._push_status_update()
return web.json_response({"result": "ok"})
async def handle_octoprint_version(self, request):
return web.json_response({
"api": "0.1",
"server": "1.9.0",
"text": "OctoPrint (Kobra X Bridge)",
})
async def handle_kx_ui_asset(self, request):
name = request.match_info.get("name", "").lstrip("/")
ctype = _KX_UI_ASSETS.get(name)
cache_control = "public, max-age=86400"
if ctype is not None:
path = os.path.join(_WEB_BASE, "web", "themes", self._ui_theme, name)
elif name.startswith("lib/"):
ext = os.path.splitext(name)[1].lower()
ctype = _KX_UI_LIB_TYPES.get(ext)
if not ctype:
raise web.HTTPNotFound()
path = os.path.join(_WEB_BASE, "web", "themes", self._ui_theme, name)
else:
m = _KX_UI_TRANSLATION_RE.match(name)
if not m:
raise web.HTTPNotFound()
lang = m.group(1)
ctype = "application/json"
cache_control = "no-store"
path = os.path.join(_WEB_BASE, "web", "translations", f"{lang}.json")
try:
raw = pathlib.Path(path).read_text(encoding="utf-8")
except OSError:
raise web.HTTPNotFound()
if name == "app.js":
raw = raw.replace("'__VERSION__'", f"'{self._read_version()}'")
return web.Response(
text=raw,
content_type=ctype,
headers={"Cache-Control": cache_control},
)
async def handle_index(self, request):
try:
tpl = self._load_index_template_cached()
except OSError:
p = self._theme_index_path()
log.error("Web UI theme file missing or unreadable: %s (theme: %s)", p, self._ui_theme)
return web.Response(
text="<pre>KX-Bridge: index.html not found.\nExpected:\n"
+ html.escape(p, quote=True)
+ "</pre>",
status=500,
content_type="text/html; charset=utf-8",
)
page = tpl.replace("__UI_ASSETS_VER__", self._ui_asset_cache_buster())
# Embed CSS + JS INLINE instead of just linking. OrcaSlicer's
# embedded device tab webview does NOT load external <link>/<script src>
# (only the bare HTML) -> without inlining neither
# a single button works there (Issue #29). It is equally correct in a normal browser.
base = os.path.join(_WEB_BASE, "web", "themes", self._ui_theme)
# Inline vendored lib CSS/JS too — the OrcaSlicer webview loads no
# external <link>/<script src>, so GridStack (and its stylesheet) must
# be embedded like style.css/app.js. Order matters: GridStack's <script>
# sits in <head>, before app.js, so it is defined when app.js inits.
def _inline_css(rel_path: str, link_tag: str):
nonlocal page
try:
data = pathlib.Path(os.path.join(base, rel_path)).read_text(encoding="utf-8")
page = page.replace(link_tag, "<style>\n" + data + "\n</style>")
except OSError:
pass
def _inline_js(rel_path: str, script_tag: str, version_sub: bool = False):
nonlocal page
try:
data = pathlib.Path(os.path.join(base, rel_path)).read_text(encoding="utf-8")
if version_sub:
data = data.replace("'__VERSION__'", f"'{self._read_version()}'")
page = page.replace(script_tag, "<script>\n" + data + "\n</script>")
except OSError:
pass
_inline_css("lib/gridstack.min.css", '<link rel="stylesheet" href="/kx/ui/lib/gridstack.min.css">')
_inline_js("lib/gridstack-all.min.js", '<script src="/kx/ui/lib/gridstack-all.min.js"></script>')
_inline_css("style.css", '<link rel="stylesheet" href="/kx/ui/style.css">')
_inline_js("app.js", '<script src="/kx/ui/app.js"></script>', version_sub=True)
return web.Response(text=page, content_type="text/html",
headers={"Cache-Control": "no-store, no-cache, must-revalidate"})
async def handle_api_light(self, request):
try:
body = await request.json()
except Exception:
body = {}
on = bool(body.get("on", True))
brightness = int(body.get("brightness", self._state["light_brightness"]))
loop = asyncio.get_event_loop()
await loop.run_in_executor(None, lambda: self.client.publish(
"light", "control",
{"type": 3, "status": 1 if on else 0, "brightness": brightness},
timeout=0
))
self._state["light_on"] = on
self._state["light_brightness"] = brightness
return web.json_response({"result": "ok"})
async def handle_api_fan(self, request):
try:
body = await request.json()
except Exception:
body = {}
speed = int(body.get("speed", 0))
loop = asyncio.get_event_loop()
await loop.run_in_executor(None, lambda: self.client.publish(
"fan", "setSpeed", {"fan_speed_pct": speed}, timeout=0
))
self._state["fan_speed"] = speed
return web.json_response({"result": "ok"})
async def handle_api_connect(self, request):
loop = asyncio.get_event_loop()
try:
await loop.run_in_executor(None, self.client.connect)
self._state["print_state"] = "standby"
self._state["kobra_state"] = "free"
log.info("Connected manually")
return web.json_response({"result": "connected"})
except Exception as e:
return web.json_response({"error": str(e)}, status=500)
async def handle_api_disconnect(self, request):
loop = asyncio.get_event_loop()
try:
await loop.run_in_executor(None, self.client.disconnect)
except Exception:
pass
self._state["print_state"] = "error"
self._state["kobra_state"] = "offline"
log.info("Manuell getrennt")
return web.json_response({"result": "disconnected"})
async def handle_api_restart(self, request):
log.info("Restart requested via API")
response = web.json_response({"status": "restarting"})
asyncio.get_event_loop().call_later(0.3, self._restart_bridge)
return response
async def handle_api_speed(self, request):
try:
body = await request.json()
except Exception:
body = {}
mode = int(body.get("mode", 2))
loop = asyncio.get_event_loop()
taskid = self._state.get("taskid", "-1")
await loop.run_in_executor(None, lambda: self.client.publish_web(
"print", "update",
{"taskid": taskid, "settings": {"print_speed_mode": mode}},
))
self._state["print_speed_mode"] = mode
return web.json_response({"result": "ok"})
async def handle_api_ams_set_slot(self, request):
try:
body = await request.json()
except Exception:
body = {}
index = int(body.get("index", 0)) # global slot index
mat = str(body.get("type", "PLA")).upper()
color = body.get("color", [255, 255, 255])
if not (isinstance(color, list) and len(color) == 3):
return web.json_response({"error": "color must be [r,g,b]"}, status=400)
box_id, local_slot = self._global_to_box_slot(index)
loop = asyncio.get_event_loop()
self._state["last_ams_set_error"] = False
# Remembered so a later state="failed" report (which carries no slot
# info of its own, see _on_multicolor_box) can be logged alongside the
# request that triggered it - otherwise the failure is unattributable.
self._last_ams_set_request = {"global": index, "box": box_id, "local_slot": local_slot, "type": mat, "color": color}
# setInfo goes via the web/printer topic (like tempature/set). Verified via
# Workbench-Vue mqtt_setInfo verifiziert — via slicer/printer/ wurden
# slot changes are ignored by the printer and overwritten with the old
# material on the next multiColorBox/report.
def _send():
self.client.publish_web(
"multiColorBox", "setInfo",
{"multi_color_box": [{"id": box_id, "slots": [{"index": local_slot, "type": mat, "color": color}]}]},
)
log.info(f"setInfo (web) global={index} box={box_id} local_slot={local_slot} type={mat} color={color}")
await loop.run_in_executor(None, _send)
# Optimistisches Update: cached slot sofort anpassen (Drucker echoed
# gleich via multiColorBox/report — falls er den Befehl ignoriert,
# the report overwrites it again).
for s in self._ams_slots:
if s.get("global_index") == index:
s["type"] = mat
s["color"] = color
break
return web.json_response({"result": "ok"})
async def handle_api_ams_feed(self, request):
try:
body = await request.json()
except Exception:
body = {}
slot_index = int(body.get("slot_index", 0))
feed_type = int(body.get("type", 1))
if feed_type == 1:
self._pending_load_slot = slot_index
# Feed-out (type=2): if no slot was explicitly chosen, use the last loaded one
if feed_type == 2 and self._ams_loaded_slot >= 0:
slot_index = self._ams_loaded_slot
box_id, local_slot = self._global_to_box_slot(slot_index)
loop = asyncio.get_event_loop()
def _send():
resp = self.client.publish(
"multiColorBox", "feedFilament",
{"multi_color_box": [{"id": box_id, "feed_status": {"slot_index": local_slot, "type": feed_type}}]},
timeout=5
)
log.info(f"feedFilament type={feed_type} global_slot={slot_index} box={box_id} local_slot={local_slot} loaded_slot={self._ams_loaded_slot}{resp}")
await loop.run_in_executor(None, _send)
return web.json_response({"result": "ok"})
async def handle_api_ace_auto_feed(self, request):
try:
body = await request.json()
except Exception:
body = {}
ace_id_raw = body.get("ace_id", None)
on_raw = body.get("on", None)
if ace_id_raw is None or on_raw is None:
return web.json_response({"error": "ace_id and on are required"}, status=400)
try:
ace_id = int(ace_id_raw)
on = int(bool(on_raw))
except Exception:
return web.json_response({"error": "invalid parameters"}, status=400)
if not (0 <= ace_id <= 3):
return web.json_response({"error": "ace_id must be 0-3"}, status=400)
payload = {"multi_color_box": [{"id": ace_id, "auto_feed": on}]}
loop = asyncio.get_event_loop()
# Fire-and-forget: setAutoFeed ACK arrives via multiColorBox/report callback.
# Waiting for a response on that busy push topic causes false "code:0" rejections.
await loop.run_in_executor(
None,
lambda: self.client.publish("multiColorBox", "setAutoFeed", payload, timeout=0)
)
self._ace_auto_feed[ace_id] = on
self._state_dirty = True
return web.json_response({"result": "ok", "ace_id": ace_id, "auto_feed": on})
async def handle_api_ace_dry(self, request):
try:
body = await request.json()
except Exception:
body = {}
action = str(body.get("action", "start")).lower()
if action not in ("start", "stop"):
return web.json_response({"error": "action must be 'start' or 'stop'"}, status=400)
ace_ids = [i for i in self._ace_box_ids if 0 <= i <= 3]
if not ace_ids:
ace_ids = sorted({
int(s.get("box_id", -1))
for s in self._ams_slots
if 0 <= int(s.get("box_id", -1)) <= 3
})
if not ace_ids and self._state.get("filament_mode") != "toolhead":
ace_ids = [0]
if not ace_ids:
return web.json_response({"error": "ACE not detected"}, status=400)
ace_id_raw = body.get("ace_id", None)
if ace_id_raw is not None:
try:
ace_id = int(ace_id_raw)
except Exception:
return web.json_response({"error": "ace_id must be an integer"}, status=400)
if ace_id not in ace_ids:
return web.json_response({"error": f"ACE {ace_id + 1} not detected"}, status=400)
ace_ids = [ace_id]
if action == "start":
target_temp = int(body.get("target_temp", 45))
duration = int(body.get("duration", 240))
target_temp = max(30, min(80, target_temp))
duration = max(10, min(24 * 60, duration))
humidity = (self._state.get("ace_drying") or {}).get("humidity")
current_temp = (self._state.get("ace_drying") or {}).get("current_temp")
drying_status = {
"status": 1,
"target_temp": target_temp,
"duration": duration,
"remain_time": duration,
}
ui_state = {
"status": 1,
"target_temp": target_temp,
"duration": duration,
"remain_time": duration,
"humidity": humidity,
"current_temp": current_temp,
}
else:
drying_status = {"status": 0}
humidity = (self._state.get("ace_drying") or {}).get("humidity")
current_temp = (self._state.get("ace_drying") or {}).get("current_temp")
ui_state = {
"status": 0,
"target_temp": 0,
"duration": 0,
"remain_time": 0,
"humidity": humidity,
"current_temp": current_temp,
}
payload = {
"multi_color_box": [
{"id": bid, "drying_status": dict(drying_status)}
for bid in ace_ids
]
}
loop = asyncio.get_event_loop()
def _send():
return self.client.publish("multiColorBox", "setDry", payload, timeout=0)
# Fire-and-forget: setDry ACK arrives via multiColorBox/report callback.
# Waiting for a response on that busy push topic causes false "code:0" rejections.
await loop.run_in_executor(None, _send)
self._state["ace_drying"] = ui_state
self._state_dirty = True
return web.json_response({"result": "ok"})
async def handle_api_axis(self, request):
try:
body = await request.json()
except Exception:
body = {}
loop = asyncio.get_event_loop()
action = str(body.get("action", "")).lower()
if action == "turnoff":
await loop.run_in_executor(None, lambda: self.client.publish(
"axis", "turnOff", None, timeout=0
))
else:
axis = int(body.get("axis", 4))
move_type = int(body.get("move_type", 2))
distance = float(body.get("distance", 0))
await loop.run_in_executor(None, lambda: self.client.publish(
"axis", "move",
{"axis": axis, "move_type": move_type, "distance": distance},
timeout=0
))
return web.json_response({"result": "ok"})
async def handle_api_temperature(self, request):
try:
body = await request.json()
except Exception:
body = {}
nozzle = body.get("nozzle")
bed = body.get("bed")
loop = asyncio.get_event_loop()
printing = self._state.get("print_state") == "printing"
if printing:
# During print: runtime update via web/printer topic, one setting at a time
taskid = self._state.get("taskid", "-1")
if nozzle is not None:
n = int(float(nozzle))
await loop.run_in_executor(None, lambda: self.client.publish_web(
"print", "update",
{"taskid": taskid, "settings": {"target_nozzle_temp": n}},
))
if bed is not None:
b = int(float(bed))
await loop.run_in_executor(None, lambda: self.client.publish_web(
"print", "update",
{"taskid": taskid, "settings": {"target_hotbed_temp": b}},
))
else:
# Idle: tempature/set via the `web/printer` topic with a `type` field.
# Confirmed by live sniffing the Anycubic Slicer Next on 2026-05-29:
# topic = web/printer/.../tempature
# data = {"type": 0|1|2, "target_hotbed_temp": B, "target_nozzle_temp": N}
# type values (from Workbench Vue): 0=nozzle, 1=bed, 2=both.
# Ohne `type` ODER auf `slicer/printer`-Topic → Systemfehler am Drucker.
if nozzle is not None and bed is not None:
t, n, b = 2, int(float(nozzle)), int(float(bed))
elif nozzle is not None:
t, n, b = 0, int(float(nozzle)), 0
elif bed is not None:
t, n, b = 1, 0, int(float(bed))
else:
return web.json_response({"result": "ok"})
await loop.run_in_executor(None, lambda: self.client.publish_web(
"tempature", "set",
{"type": t, "target_nozzle_temp": n, "target_hotbed_temp": b},
))
return web.json_response({"result": "ok"})
async def handle_api_camera(self, request):
return web.json_response({"url": self._state["camera_url"]})
async def handle_api_camera_start(self, request):
loop = asyncio.get_event_loop()
# Wait for pushStarted confirmation before returning
result = await loop.run_in_executor(None, lambda: self.client.publish(
"video", "startCapture", None, timeout=8.0
))
state = (result or {}).get("state", "")
log.info(f"Camera startCapture: state={state}")
return web.json_response({"result": "ok", "state": state})
async def handle_api_camera_stop(self, request):
loop = asyncio.get_event_loop()
await loop.run_in_executor(None, lambda: self.client.publish(
"video", "stopCapture", None, timeout=0
))
# Prevents the auto-start guard from restarting the camera during the
# laufenden Drucks wieder einschaltet (State-Flicker-Problem).
self._camera_user_stopped = True
return web.json_response({"result": "ok"})
async def handle_api_camera_reset(self, request):
"""Reset the backoff counter and restart ffmpeg immediately.
Useful after a 429 lock (Retry-After expired) or after a printer restart."""
self.camera_cache.reset()
url = self._state.get("camera_url", "")
if not url:
log.warning("Camera reset requested but no camera_url is known yet (waiting for printer status)")
return web.json_response({
"result": "no_url",
"message": "No camera URL known yet - wait for the next printer status update, or start a print/enable the camera first.",
})
self.camera_cache.set_url(url)
await self.camera_cache.ensure_running()
return web.json_response({"result": "ok", "url": url})
async def handle_api_camera_snapshot(self, request):
"""Last JPEG frame from the CameraCache - instant from RAM,
no separate ffmpeg instance anymore (prevents the single-client 429 at the
printer and is ~1 s faster)."""
url = self._state.get("camera_url", "")
if not url:
return web.Response(status=503, text="No camera URL known")
self.camera_cache.set_url(url)
await self.camera_cache.ensure_running()
# Initial warmup: wait up to 5s for the first frame
deadline = time.time() + 5.0
while not self.camera_cache.latest_jpeg and time.time() < deadline:
await asyncio.sleep(0.1)
jpeg = self.camera_cache.latest_jpeg
if not jpeg:
return web.Response(status=503, text="No frame in cache yet")
# If the last frame is older than 10 s -> the cache ffmpeg is probably
# no longer running stably; deliver anyway but with a stale header.
age = time.time() - self.camera_cache.latest_jpeg_ts
headers = {"Cache-Control": "no-cache"}
if age > 10:
headers["X-Frame-Age"] = f"{age:.1f}"
return web.Response(body=jpeg, content_type="image/jpeg", headers=headers)
async def handle_camera_stream(self, request):
"""MJPEG live view, served as multipart/x-mixed-replace.
Fed from the central CameraCache fanout (same pattern as
handle_camera_h264) instead of spawning a dedicated ffmpeg process
per HTTP client. The printer's camera server only tolerates a very
limited number of concurrent connections (see CameraCache docstring)
- previously every consumer of this endpoint (dashboard, OrcaSlicer,
moonraker-obico, a second browser tab, ...) opened its own separate
connection, so two simultaneous viewers could already exhaust the
printer's connection limit and cause intermittent "stream
unavailable" failures. Now all consumers share one connection.
"""
url = self._state.get("camera_url", "")
if not url:
return web.Response(status=503, text="No camera URL known")
self.camera_cache.set_url(url)
await self.camera_cache.ensure_running()
q: asyncio.Queue[bytes] = asyncio.Queue(maxsize=8)
self.camera_cache.mjpeg_subscribers.add(q)
# Wait for the first frame BEFORE resp.prepare() - once prepare() sends
# the response headers the status is committed to 200, so a stalled
# source (Issue #99) must be caught here to actually return a 503
# instead of hanging the client forever with no frame ever arriving.
try:
first_frame = await asyncio.wait_for(q.get(), timeout=5.0)
except asyncio.TimeoutError:
self.camera_cache.mjpeg_subscribers.discard(q)
return web.Response(status=503, text="No frame in cache yet")
boundary = "kobraxframe"
resp = web.StreamResponse(headers={
"Content-Type": f"multipart/x-mixed-replace;boundary={boundary}",
"Cache-Control": "no-cache",
"Connection": "keep-alive",
})
await resp.prepare(request)
try:
frame = first_frame
while True:
header = (
f"--{boundary}\r\n"
f"Content-Type: image/jpeg\r\n"
f"Content-Length: {len(frame)}\r\n\r\n"
).encode()
try:
await resp.write(header + frame + b"\r\n")
except (ConnectionResetError, asyncio.CancelledError):
break
except Exception:
break
frame = await q.get()
except Exception as e:
log.warning(f"Camera stream interrupted: {e}")
finally:
self.camera_cache.mjpeg_subscribers.discard(q)
return resp
async def handle_camera_h264(self, request):
"""H.264 passthrough as MPEG-TS, fed from the central
CameraCache fanout. Allows multiple parallel consumers without an
additional FLV connection to the printer (single-client limit)."""
url = self._state.get("camera_url", "")
if not url:
return web.Response(status=503, text="No camera URL known")
self.camera_cache.set_url(url)
await self.camera_cache.ensure_running()
q: asyncio.Queue[bytes] = asyncio.Queue(maxsize=64)
self.camera_cache.h264_subscribers.add(q)
resp = web.StreamResponse(headers={
"Content-Type": "video/mp2t",
"Cache-Control": "no-cache",
"Connection": "keep-alive",
})
await resp.prepare(request)
try:
while True:
chunk = await q.get()
try:
await resp.write(chunk)
except (ConnectionResetError, asyncio.CancelledError):
break
except Exception as e:
log.warning(f"H.264-Stream unterbrochen: {e}")
finally:
self.camera_cache.h264_subscribers.discard(q)
return resp
async def handle_serve_file(self, request):
"""Serves uploaded G-code files from the temp directory (for printer download)."""
filename = os.path.basename(request.match_info.get("filename", ""))
serve_path = os.path.join(self._serve_dir_path, filename)
if not os.path.isfile(serve_path):
return web.Response(status=404, text="not found")
size = os.path.getsize(serve_path)
log.info(f"Printer downloading file: {filename} ({size} bytes)")
return web.FileResponse(serve_path, headers={
"Content-Disposition": f'attachment; filename="{filename}"'
})
async def handle_api_state(self, request):
s = self._state
# Slicer time + thumbnail are only transient in state (set during upload).
# After a browser reload or an OrcaSlicer direct print (file did not come
# through the UI upload) they are missing -> restore from the GCode store via the
# laufenden Dateinamens nachladen.
slicer_time = s["slicer_time"]
thumbnail = self._thumbnail_b64
fname = s.get("filename", "")
if fname and (not slicer_time or not thumbnail):
try:
gf = self._store.get_file_by_name(fname)
if gf:
if not slicer_time and gf.get("est_print_time_sec"):
slicer_time = int(gf["est_print_time_sec"])
if not thumbnail and gf.get("thumbnail_b64"):
thumbnail = gf["thumbnail_b64"]
except Exception:
pass
return web.json_response({
"printer_name": s["printer_name"],
"firmware_version": s["firmware_version"],
"print_state": s["print_state"],
"kobra_state": s["kobra_state"],
"nozzle_temp": s["nozzle_temp"],
"nozzle_target": s["nozzle_target"],
"bed_temp": s["bed_temp"],
"bed_target": s["bed_target"],
"progress": s["progress"],
"print_duration": s["print_duration"],
"remain_time": s["remain_time"],
"curr_layer": s["curr_layer"],
"total_layers": s["total_layers"],
"z_mm": self._estimate_current_z(),
"filename": s["filename"],
"slicer_time": slicer_time,
"camera_url": s["camera_url"],
"fan_speed": s["fan_speed"],
"print_speed_mode": s["print_speed_mode"],
"auto_leveling": getattr(self._args, "auto_leveling", 1),
"vibration_compensation": getattr(self._args, "vibration_compensation", 0),
"camera_on_print": getattr(self._args, "camera_on_print", 0),
"web_upload_warning": getattr(self._args, "web_upload_warning", 1),
"light_on": s["light_on"],
"light_brightness": s["light_brightness"],
"ams_slots": self._ams_slots,
"ams_loaded_slot": self._ams_loaded_slot,
"filament_mode": s.get("filament_mode", self._filament_mode),
"ace_drying": s.get("ace_drying", {"status": 0, "target_temp": 0, "duration": 0, "remain_time": 0, "humidity": None, "current_temp": None}),
"ace_units": list(self._ace_box_ids),
"ace_auto_feed": dict(self._ace_auto_feed),
"ace_dry_presets": self._ace_dry_presets,
"thumbnail": thumbnail,
"connection_error": s["connection_error"],
"file_ready": s["file_ready"],
"print_start_dialog": s.get("print_start_dialog", getattr(self._args, "print_start_dialog", 1)),
"version": self._read_version(),
"pause_msg": s.get("pause_msg", ""),
"error_code": s.get("error_code", 0),
"storage_total_mb": s.get("storage_total_mb", 0),
"storage_used_mb": s.get("storage_used_mb", 0),
})
async def handle_moonraker_database(self, request):
"""OrcaSlicer Filament-Sync: /server/database/item?namespace=lane_data&key=lanes (AFC-Format)"""
namespace = request.rel_url.query.get("namespace", "")
key = request.rel_url.query.get("key", "")
if namespace == "lane_data":
await asyncio.get_event_loop().run_in_executor(None, self._get_ams_slots_fresh)
lanes = self._build_lane_data()
log.info(f"AMS-Sync: {len(lanes)} Lanes an OrcaSlicer")
return web.json_response({
"result": {
"namespace": "lane_data",
"key": key or "lanes",
"value": lanes,
}
})
if namespace in ("AFC", "afc-install", "happy_hare"):
return web.json_response({
"result": {"namespace": namespace, "key": key, "value": None}
})
# mainsail/presets: Obico asks for temperature presets. The schema is evaluated in
# find_all_thermal_presets as data['value']['presets'].values(),
# so we need at least {presets: {}} to avoid a crash.
if namespace == "mainsail":
if key == "presets":
return web.json_response({
"result": {"namespace": "mainsail", "key": "presets",
"value": {"presets": {}}}
})
return web.json_response({
"result": {"namespace": "mainsail", "key": key, "value": {}}
})
# obico namespace: in-memory KV store for plugin settings (key=printer_id etc.)
if namespace == "obico":
store = self._moonraker_kv_store.setdefault("obico", {})
if key and key in store:
return web.json_response({
"result": {"namespace": "obico", "key": key, "value": store[key]}
})
return web.json_response({
"result": {"namespace": "obico", "key": key, "value": store if not key else None}
})
return web.json_response(
{"error": {"code": 404, "message": f"Namespace '{namespace}' not found"}},
status=404
)
async def handle_moonraker_database_post(self, request):
"""POST /server/database/item — KV-Store-Write (von moonraker-obico verwendet).
moonraker-obico sends namespace/key/value as form-urlencoded POST params."""
# Versuche JSON, fallback auf form-data, fallback auf Query-Params
namespace = ""
key = ""
value = None
try:
data = await request.json()
if isinstance(data, dict):
namespace = data.get("namespace", "")
key = data.get("key", "")
value = data.get("value")
except Exception:
try:
form = await request.post()
namespace = form.get("namespace", "") or ""
key = form.get("key", "") or ""
value = form.get("value")
except Exception:
pass
if not namespace:
namespace = request.rel_url.query.get("namespace", "")
if not key:
key = request.rel_url.query.get("key", "")
if namespace and key:
store = self._moonraker_kv_store.setdefault(namespace, {})
store[key] = value
return web.json_response({
"result": {"namespace": namespace, "key": key, "value": value}
})
return web.json_response({"error": {"code": 400, "message": "namespace + key required"}}, status=400)
async def handle_database_list(self, request):
"""OrcaSlicer checks which namespaces exist to detect the MMU type."""
return web.json_response({"result": {"namespaces": ["lane_data", "mainsail", "obico"]}})
def _get_ams_slots_fresh(self):
"""Frische Slot-Daten per getInfo holen, Fallback auf gecachte."""
resp = self.client.publish("multiColorBox", "getInfo", None, timeout=5)
if resp and resp.get("data"):
data = resp["data"]
self._head_tools_model = int(data.get("head_tools_model", self._head_tools_model))
boxes = data.get("multi_color_box") or []
if boxes:
self._update_ace_drying_state(data, boxes)
self._filament_mode = self._detect_filament_mode(boxes, self._head_tools_model)
self._state["filament_mode"] = self._filament_mode
global_slots, global_loaded = self._aggregate_slots(boxes, self._filament_mode)
activity_map = self._slot_activity_map(boxes, global_loaded)
for s in global_slots:
s["activity"] = activity_map.get(s.get("global_index"), "")
if global_slots:
self._ams_slots = global_slots
self._ams_loaded_slot = global_loaded
return self._ams_slots
# ─── Settings ────────────────────────────────────────────────────────────
def _find_config_path(self) -> pathlib.Path:
"""Returns the path to config.ini."""
if hasattr(env_loader, "find_config_path"):
return env_loader.find_config_path()
# Fallback for the old env_loader
script_dir = pathlib.Path(_BASE)
for base in (script_dir, script_dir.parent):
p = base / "config" / "config.ini"
if p.is_file():
return p
return script_dir / "config" / "config.ini"
async def handle_api_settings_get(self, request):
return web.json_response({
"printer_name": self._state.get("printer_name", ""),
"printer_ip": self._args.printer_ip,
"mqtt_port": self._args.mqtt_port,
"username": self._args.username,
"password": self._args.password,
"mode_id": self._args.mode_id,
"device_id": self._args.device_id,
"power_on_url": getattr(self._args, "power_on_url", "") or "",
"power_off_url": getattr(self._args, "power_off_url", "") or "",
"power_status_url": getattr(self._args, "power_status_url", "") or "",
"default_ams_slot": getattr(self._args, "default_ams_slot", "auto"),
"auto_leveling": getattr(self._args, "auto_leveling", 1),
"vibration_compensation": getattr(self._args, "vibration_compensation", 0),
"camera_on_print": getattr(self._args, "camera_on_print", 0),
"web_upload_warning": getattr(self._args, "web_upload_warning", 1),
"delete_printer_file_after_print": getattr(self._args, "delete_printer_file_after_print", 0),
"print_start_dialog": getattr(self._args, "print_start_dialog", 1),
"poll_interval": getattr(self._args, "poll_interval", 3),
"verbose_http_log": getattr(self._args, "verbose_http_log", 0),
"filament_profiles": {str(k): v for k, v in self._filament_profiles.items()},
"visible_vendors": self._visible_vendors,
"ace_dry_presets": self._ace_dry_presets,
"spoolman_server": getattr(self._args, "spoolman_server", "") or "",
"spoolman_sync_rate": getattr(self._args, "spoolman_sync_rate", 0),
"kxgauge_url": getattr(self._args, "kxgauge_url", "") or "",
"kxgauge_enabled": int(bool(getattr(self._args, "kxgauge_enabled", 0))),
"kxgauge_heat_peak": getattr(self._args, "kxgauge_heat_peak", 250),
"kxgauge_mapping": self._kxgauge_mapping,
})
async def handle_api_settings_post(self, request):
import configparser
try:
data = await request.json()
except Exception:
return self._json_cors({"error": "invalid json"}, status=400)
config_path = self._find_config_path()
config_path.parent.mkdir(parents=True, exist_ok=True)
# Read the existing config.ini (comments are lost, but values are kept)
cfg = configparser.ConfigParser(interpolation=None)
if config_path.is_file():
cfg.read(config_path, encoding="utf-8")
# Sections sicherstellen
for section in ("connection", "print", "bridge", "ace_dry_presets", "spoolman", "kxgauge", "kxgauge_mapping"):
if not cfg.has_section(section):
cfg.add_section(section)
printer_ip = str(data.get("printer_ip", self._args.printer_ip or "")).split(":")[0]
cfg.set("connection", "printer_ip", printer_ip)
cfg.set("connection", "mqtt_port", str(data.get("mqtt_port", self._args.mqtt_port or 9883)))
cfg.set("connection", "username", str(data.get("username", self._args.username or "")))
cfg.set("connection", "password", str(data.get("password", self._args.password or "")))
cfg.set("connection", "mode_id", str(data.get("mode_id", self._args.mode_id or "")))
cfg.set("connection", "device_id", str(data.get("device_id", self._args.device_id or "")))
cfg.set("connection", "power_on_url", str(data.get("power_on_url", getattr(self._args, "power_on_url", "") or "")).strip())
cfg.set("connection", "power_off_url", str(data.get("power_off_url", getattr(self._args, "power_off_url", "") or "")).strip())
cfg.set("connection", "power_status_url", str(data.get("power_status_url", getattr(self._args, "power_status_url", "") or "")).strip())
cfg.set("print", "default_ams_slot", str(data.get("default_ams_slot", getattr(self._args, "default_ams_slot", "auto"))))
cfg.set("print", "auto_leveling", str(data.get("auto_leveling", getattr(self._args, "auto_leveling", 1))))
cfg.set("print", "vibration_compensation", str(int(bool(data.get("vibration_compensation", getattr(self._args, "vibration_compensation", 0))))))
cfg.set("print", "camera_on_print", str(int(bool(data.get("camera_on_print", getattr(self._args, "camera_on_print", 0))))))
cfg.set("print", "web_upload_warning", str(int(bool(data.get("web_upload_warning", getattr(self._args, "web_upload_warning", 1))))))
cfg.set("print", "delete_printer_file_after_print", str(int(bool(data.get("delete_printer_file_after_print", getattr(self._args, "delete_printer_file_after_print", 0))))))
cfg.set("print", "print_start_dialog", str(int(bool(data.get("print_start_dialog", getattr(self._args, "print_start_dialog", 1))))))
if "poll_interval" in data:
try:
pi = max(1, min(60, int(data["poll_interval"])))
except (TypeError, ValueError):
pi = 3
cfg.set("bridge", "poll_interval", str(pi))
elif not cfg.has_option("bridge", "poll_interval"):
cfg.set("bridge", "poll_interval", "3")
verbose_http_log = int(bool(data.get("verbose_http_log", getattr(self._args, "verbose_http_log", 0))))
cfg.set("bridge", "verbose_http_log", str(verbose_http_log))
_set_verbose_http_log(bool(verbose_http_log))
self._args.verbose_http_log = verbose_http_log
printer_name = str(data.get("printer_name", "")).strip()
if printer_name:
cfg.set("bridge", "printer_name", printer_name)
elif cfg.has_option("bridge", "printer_name"):
cfg.remove_option("bridge", "printer_name")
# Spoolman
if "spoolman_server" in data:
cfg.set("spoolman", "server", str(data["spoolman_server"]).strip())
if "spoolman_sync_rate" in data:
try:
sr = max(0, int(data["spoolman_sync_rate"]))
except (TypeError, ValueError):
sr = 30
cfg.set("spoolman", "sync_rate", str(sr))
# KXGauge
if "kxgauge_url" in data:
cfg.set("kxgauge", "url", str(data["kxgauge_url"]).strip())
if "kxgauge_enabled" in data:
cfg.set("kxgauge", "enabled", str(int(bool(data["kxgauge_enabled"]))))
if "kxgauge_heat_peak" in data:
try:
hp = max(0, min(400, float(data["kxgauge_heat_peak"])))
except (TypeError, ValueError):
hp = 250.0
cfg.set("kxgauge", "heat_peak", str(hp))
incoming_mapping = data.get("kxgauge_mapping") if isinstance(data, dict) else None
if isinstance(incoming_mapping, dict):
mapping = dict(self._kxgauge_mapping)
for state, emotion in incoming_mapping.items():
state = str(state).strip().lower()
emotion = str(emotion).strip().lower()
if state and emotion in KXGAUGE_VALID_EMOTIONS:
mapping[state] = emotion
for state, emotion in mapping.items():
cfg.set("kxgauge_mapping", state, emotion)
self._kxgauge_mapping = mapping
incoming_presets = data.get("ace_dry_presets") if isinstance(data, dict) else None
presets = self._sanitize_ace_dry_presets(incoming_presets if isinstance(incoming_presets, dict) else self._ace_dry_presets)
for key, val in presets.items():
cfg.set("ace_dry_presets", f"{key}_temp", str(val["temp"]))
cfg.set("ace_dry_presets", f"{key}_duration_sec", str(val["duration_sec"]))
if key.startswith("custom_"):
cfg.set("ace_dry_presets", f"{key}_name", str(val.get("name", key.replace("_", " ").title())))
self._ace_dry_presets = presets
with open(config_path, "w", encoding="utf-8") as f:
f.write("# KX-Bridge Konfigurationsdatei\n\n")
cfg.write(f)
log.info(f"Settings saved to {config_path}")
# Send the response, then restart
response = web.json_response({"status": "restarting"})
asyncio.get_event_loop().call_later(0.3, self._restart_bridge)
return response
async def handle_kx_printer_add(self, request):
"""Adds a printer: fetches credentials via IP, writes [printer_N], restarts."""
try:
body = await request.json()
except Exception:
return self._json_cors({"error": "invalid json"}, status=400)
ip = str(body.get("printer_ip", "")).strip().split(":")[0]
name = str(body.get("name", "")).strip()
if not ip:
return self._json_cors({"error": "printer_ip required"}, status=400)
try:
creds = await _kx_fetch_credentials(ip)
except Exception as e:
return self._json_cors({"error": f"printer unreachable or error: {e}"}, status=502)
import configparser
config_path = self._find_config_path()
cfg = configparser.ConfigParser(interpolation=None)
if config_path.is_file():
cfg.read(config_path, encoding="utf-8")
# Vorhandene [printer_N]-Sektionen + belegte http_ports ermitteln
n = 1
existing_ports: set[int] = set()
while cfg.has_section(f"printer_{n}"):
p = cfg[f"printer_{n}"]
if p.get("http_port"):
try:
existing_ports.add(int(p["http_port"]))
except ValueError:
pass
n += 1
# No [printer_N], but a populated [connection]? -> migrate as printer_1
# (empty [connection] = no existing printer -> don't migrate, the new one becomes printer_1)
if n == 1 and cfg.has_section("connection") and (cfg["connection"].get("printer_ip") or "").strip():
c = cfg["connection"]
cfg.add_section("printer_1")
cfg.set("printer_1", "name", self._state.get("printer_name") or "Kobra X")
for k in ("printer_ip", "mqtt_port", "username", "password", "mode_id", "device_id"):
if c.get(k):
cfg.set("printer_1", k, c.get(k))
cfg.set("printer_1", "http_port", "7125")
existing_ports.add(7125)
n = 2
# Create the new printer as [printer_n], pick a free port
new_port = 7125 + (n - 1)
while new_port in existing_ports:
new_port += 1
sec = f"printer_{n}"
cfg.add_section(sec)
cfg.set(sec, "name", name or creds["model"])
cfg.set(sec, "printer_ip", creds["printer_ip"])
cfg.set(sec, "mqtt_port", "9883")
cfg.set(sec, "username", creds["username"])
cfg.set(sec, "password", creds["password"])
cfg.set(sec, "mode_id", creds["mode_id"])
cfg.set(sec, "device_id", creds["device_id"])
cfg.set(sec, "http_port", str(new_port))
config_path.parent.mkdir(parents=True, exist_ok=True)
with open(config_path, "w", encoding="utf-8") as f:
f.write("# KX-Bridge Konfigurationsdatei\n\n")
cfg.write(f)
log.info(f"Printer '{name or creds['model']}' added as {sec} (port {new_port})")
response = self._json_cors({"status": "restarting", "section": sec, "http_port": new_port})
asyncio.get_event_loop().call_later(0.5, self._restart_bridge)
return response
async def handle_kx_printer_remove(self, request):
"""Removes a printer from config.ini, then restarts.
- Multi mode: [printer_N] is deleted, the rest renumbered (printer_3 -> printer_2),
printer_1 bekommt immer http_port 7125.
- Single mode (no [printer_N], only [connection]): pid "1" clears the [connection] block
→ Bridge startet im Offline-Modus auf 7125, UI bleibt erreichbar.
- When the last [printer_N] is removed: all gone -> also the "empty" state.
"""
pid = str(request.match_info.get("pid", "")).strip()
if not pid:
return self._json_cors({"error": "printer id required"}, status=400)
import configparser
config_path = self._find_config_path()
cfg = configparser.ConfigParser(interpolation=None)
if config_path.is_file():
cfg.read(config_path, encoding="utf-8")
has_printer_sections = cfg.has_section("printer_1")
target = f"printer_{pid}"
if has_printer_sections:
if not cfg.has_section(target):
return self._json_cors({"error": f"{target} not found"}, status=404)
# Collect all [printer_N] (except the one being deleted), renumber
kept = []
n = 1
while cfg.has_section(f"printer_{n}"):
if str(n) != pid:
kept.append(dict(cfg[f"printer_{n}"]))
cfg.remove_section(f"printer_{n}")
n += 1
for i, sec_data in enumerate(kept, start=1):
sec = f"printer_{i}"
cfg.add_section(sec)
for k, v in sec_data.items():
cfg.set(sec, k, v)
cfg.set(sec, "http_port", str(7125 + i - 1))
remaining = len(kept)
# Was that the last printer? Then also clear [connection] -> truly "no printer"
if remaining == 0 and cfg.has_section("connection"):
for k in ("printer_ip", "username", "password", "device_id"):
cfg.set("connection", k, "")
else:
# Single mode: only pid "1" is valid (pseudo entry from handle_kx_printers)
if pid != "1":
return self._json_cors({"error": "no printer with this ID"}, status=404)
# Clear [connection] values -> bridge starts without a printer
if cfg.has_section("connection"):
for k in ("printer_ip", "username", "password", "device_id"):
cfg.set("connection", k, "")
remaining = 0
config_path.parent.mkdir(parents=True, exist_ok=True)
with open(config_path, "w", encoding="utf-8") as f:
f.write("# KX-Bridge Konfigurationsdatei\n\n")
cfg.write(f)
log.info(f"Printer {target} removed ({remaining} remaining)")
response = self._json_cors({"status": "restarting", "removed": target, "remaining": remaining})
asyncio.get_event_loop().call_later(0.5, self._restart_bridge)
return response
def _restart_bridge(self):
log.info("Restarting bridge...")
# config_loader caches config.ini values in os.environ ("only if not set").
# On restart, environ must be cleaned, otherwise the new process reads
# the old values instead of the modified config.ini. Keys are derived
# from config_loader.CONFIG_ENV_MAPPING (single source of truth) so a
# newly added setting can never be forgotten here again.
try:
import config_loader as _cl
_restart_env_keys = set(_cl.CONFIG_ENV_MAPPING.keys()) | {"FILE_READY_DIALOG"}
except Exception:
_restart_env_keys = ()
for _k in _restart_env_keys:
os.environ.pop(_k, None)
in_docker = os.path.exists("/.dockerenv") or os.environ.get("KX_IN_DOCKER")
if in_docker:
# Docker/systemd: exiting the process is enough - the supervisor restarts (fresh environ)
log.info("Container environment detected exiting for supervisor restart")
os._exit(0)
frozen = getattr(sys, "frozen", False)
# Linux: os.execv replaces the process image directly - clean even with PyInstaller onefile
# (subprocess+exit would fail there on the deleted _MEIxxxx temp directory).
if sys.platform != "win32":
exe = sys.executable
try:
if frozen:
os.execv(exe, [exe] + sys.argv[1:])
else:
os.execv(exe, [exe] + sys.argv)
except Exception as e:
log.error(f"Restart (execv) failed: {e} - please restart the bridge manually")
os._exit(1)
# Windows: os.execv is broken there (new PID, old process returns) -> subprocess
cmd = ([sys.executable] + sys.argv[1:]) if frozen else ([sys.executable] + sys.argv)
try:
subprocess.Popen(cmd, cwd=os.getcwd(),
creationflags=(subprocess.DETACHED_PROCESS
| subprocess.CREATE_NEW_PROCESS_GROUP))
except Exception as e:
log.error(f"Restart failed: {e} - please restart the bridge manually")
os._exit(0)
# ─── Update ──────────────────────────────────────────────────────────────
# limit=1 would only ever see the single newest release regardless of type -
# if that happens to be a nightly/dev prerelease (the common case, since
# those publish far more often than stable), the stable_releases filter
# below finds nothing and update checks fail with "no stable releases
# found" even though older stable releases exist (Issue #104).
STABLE_RELEASE_API = "https://gitea.it-drui.de/api/v1/repos/viewit/KX-Bridge-Release/releases?limit=20"
NIGHTLY_RELEASE_API = "https://gitea.it-drui.de/api/v1/repos/viewit/KX-Bridge-Release/releases?limit=5&pre-release=true"
DEV_RELEASE_API = "https://gitea.it-drui.de/api/v1/repos/viewit/KX-Bridge-Release/releases?limit=10&pre-release=true"
GITEA_RAW_BASE = "https://gitea.it-drui.de/viewit/KX-Bridge-Release/raw/tag"
def _read_version(self) -> str:
# PyInstaller onefile unpacks VERSION (via kx-bridge.spec datas) to
# sys._MEIPASS - therefore use _WEB_BASE instead of _BASE.
for base in (pathlib.Path(_WEB_BASE), pathlib.Path(_BASE), pathlib.Path(_BASE).parent):
p = base / "VERSION"
if p.is_file():
return p.read_text(encoding="utf-8").strip()
return "unknown"
def _write_version(self, version: str):
for base in (pathlib.Path(_BASE), pathlib.Path(_BASE).parent):
p = base / "VERSION"
if p.is_file():
p.write_text(version + "\n", encoding="utf-8")
return
(pathlib.Path(_BASE) / "VERSION").write_text(version + "\n", encoding="utf-8")
@staticmethod
def _parse_version(v: str) -> "tuple[int, ...]":
"""'v0.9.1-beta1' -> (0, 9, 1) - only numeric parts before the first '-'"""
v = v.lstrip("v").split("-")[0]
parts = re.split(r"[.\s]+", v)
result = []
for p in parts:
try:
result.append(int(p))
except ValueError:
break
return tuple(result) or (0,)
async def handle_api_log_stream(self, request):
"""SSE endpoint: streams log entries live to the browser."""
resp = web.StreamResponse(headers={
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
})
await resp.prepare(request)
# Zuerst Ring-Buffer senden
for entry in list(_log_buffer):
data = json.dumps(entry, ensure_ascii=False)
await resp.write(f"data: {data}\n\n".encode())
# Dann live streamen
q: asyncio.Queue = asyncio.Queue()
_log_sse_queues.append(q)
try:
while True:
entry = await asyncio.wait_for(q.get(), timeout=25)
data = json.dumps(entry, ensure_ascii=False)
await resp.write(f"data: {data}\n\n".encode())
except asyncio.TimeoutError:
await resp.write(b": keepalive\n\n")
except (ConnectionResetError, Exception):
pass
finally:
_log_sse_queues.remove(q) if q in _log_sse_queues else None
return resp
async def handle_api_log_download(self, request):
"""Returns all buffered log entries as plaintext for download."""
header = (f"# KX-Bridge Log | Version {self._read_version()} | "
f"{time.strftime('%Y-%m-%d %H:%M:%S')} | {len(_log_buffer)} entries\n")
lines = [f"[{e['ts']}] {e['lvl']:<7} {e['name']}: {e['msg']}" for e in _log_buffer]
text = header + "\n".join(lines) + "\n"
fname = f"kx-bridge-log_{time.strftime('%Y%m%d-%H%M%S')}.txt"
return web.Response(
body=text.encode("utf-8"),
content_type="text/plain",
headers={"Content-Disposition": f'attachment; filename="{fname}"'},
)
async def handle_api_update_check(self, request):
current = self._read_version()
# Testing channel (testing-<sha>) has no Gitea releases at all - it's
# a Docker-only channel. Report that directly instead of falling
# through to the stable path (which would wrongly offer a stable
# "update"). The :testing image is rolling, so there's nothing to
# compare a version against.
if "testing" in current:
return web.json_response({
"current": current,
"latest": current,
"update_available": False,
"tag": current,
"docker_only": True,
"changelog": "Testing channel - updates are delivered via Docker: "
"docker compose pull && docker compose up -d",
})
is_nightly = "nightly" in current
is_dev = "-dev+" in current
if is_nightly:
api_url = self.NIGHTLY_RELEASE_API
elif is_dev:
api_url = self.DEV_RELEASE_API
else:
api_url = self.STABLE_RELEASE_API
try:
async with aiohttp.ClientSession() as session:
async with session.get(api_url, timeout=aiohttp.ClientTimeout(total=10)) as resp:
if resp.status != 200:
return web.json_response({"error": f"Gitea HTTP {resp.status}"}, status=502)
releases = await resp.json(content_type=None)
if not releases:
return web.json_response({"error": "no releases found"}, status=404)
if is_nightly:
# Find the newest prerelease with a nightly tag
nightly_releases = [r for r in releases if r.get("prerelease") and "nightly" in r.get("tag_name", "")]
if not nightly_releases:
return web.json_response({"error": "no nightly releases found"}, status=404)
data = nightly_releases[0]
tag = data.get("tag_name", "")
# Tag-Format: "nightly-0.9.27-nightly4", current: "0.9.27-nightly4"
tag_version = tag[len("nightly-"):] if tag.startswith("nightly-") else tag
update_available = tag_version != current
latest = tag
return web.json_response({
"current": current,
"latest": latest,
"update_available": update_available,
"tag": tag,
"docker_only": True,
"changelog": data.get("body", ""),
})
elif is_dev:
dev_releases = [r for r in releases if "-dev+" in r.get("tag_name", "")]
if not dev_releases:
return web.json_response({"error": "no dev releases found"}, status=404)
data = dev_releases[0]
else:
# Stable: only take non-prereleases
stable_releases = [r for r in releases if not r.get("prerelease")]
if not stable_releases:
return web.json_response({"error": "no stable releases found"}, status=404)
data = stable_releases[0]
tag = data.get("tag_name", "")
latest = tag.lstrip("v")
if is_dev:
update_available = tag != f"v{current}"
else:
update_available = self._parse_version(tag) > self._parse_version(current)
download_url = f"{self.GITEA_RAW_BASE}/{tag}/kobrax_moonraker_bridge.py"
return web.json_response({
"current": current,
"latest": latest,
"update_available": update_available,
"tag": tag,
"download_url": download_url,
"docker_only": False,
"changelog": data.get("body", ""),
})
except Exception as e:
return web.json_response({"error": str(e)}, status=502)
# Bridge Python modules the self-update must include. If only the
# main file is replaced, the new version may crash with ModuleNotFoundError.
# Note: since the theme system, the frontend lives under web/themes/<name>/
# (no flat .py anymore); theme files are currently NOT included in the
# self-update - theme changes arrive via Docker image/binary updates.
_UPDATE_FILES = [
"kobrax_moonraker_bridge.py",
"kobrax_client.py",
"config_loader.py",
"env_loader.py",
]
async def handle_api_update_apply(self, request):
try:
data = await request.json()
except Exception:
return web.json_response({"error": "invalid json"}, status=400)
new_tag = data.get("tag", "")
_cur = self._read_version()
if "nightly" in _cur or "testing" in _cur:
channel = "testing" if "testing" in _cur else "nightly"
return web.json_response(
{"error": f"{channel} updates are delivered via Docker: "
"docker compose pull && docker compose up -d"}, status=400)
if getattr(sys, "frozen", False):
return web.json_response(
{"error": "self-update is not supported in binary mode - "
"please download the new binary/Docker image."}, status=400)
if not new_tag:
return web.json_response({"error": "missing tag"}, status=400)
app_dir = pathlib.Path(__file__).resolve().parent
try:
# Phase 1: ALLE Dateien herunterladen (in .new), nichts ersetzen.
downloaded: list[tuple[pathlib.Path, bytes]] = []
async with aiohttp.ClientSession() as session:
for fname in self._UPDATE_FILES:
url = f"{self.GITEA_RAW_BASE}/{new_tag}/{fname}"
async with session.get(url, timeout=aiohttp.ClientTimeout(total=30)) as resp:
if resp.status != 200:
# _web_assets.py etc. may not exist in older tags -
# the main file is mandatory, optional ones may be missing.
if fname == "kobrax_moonraker_bridge.py":
return web.json_response(
{"error": f"Download {fname}: HTTP {resp.status}"}, status=502)
log.warning(f"Update: {fname} not found in release ({resp.status}) skipped")
continue
downloaded.append((app_dir / fname, await resp.read()))
# Phase 2: replace atomically (only after a complete, successful download)
for path, content in downloaded:
tmp = path.with_suffix(path.suffix + ".new")
tmp.write_bytes(content)
os.replace(tmp, path)
self._write_version(new_tag.lstrip("v"))
log.info(f"Update to {new_tag} installed ({len(downloaded)} files), restarting...")
except Exception as e:
return web.json_response({"error": str(e)}, status=502)
response = web.json_response({"status": "updating"})
asyncio.get_event_loop().call_later(0.3, self._restart_bridge)
return response
async def handle_catchall(self, request):
body = await request.read()
log.warning(f"UNBEKANNT {request.method} {request.path_qs} body={body[:200]}")
return web.json_response({"result": {}}, status=200)
async def handle_favicon(self, request):
# Minimal 1x1 ICO so the browser doesn't log a 404
ico = bytes([
0,0,1,0,1,0,1,1,0,0,1,0,24,0,40,0,0,0,22,0,0,0,40,0,0,0,
1,0,0,0,2,0,0,0,1,0,24,0,0,0,0,0,4,0,0,0,0,0,0,0,0,0,0,0,
0,0,0,0,0,0,0,0,255,102,0,0,0,0,0,0
])
return web.Response(body=ico, content_type="image/x-icon")
# -------------------------------------------------------------------------
# Klipper G-code script emulation for moonraker-obico
# -------------------------------------------------------------------------
async def _exec_gcode_script(self, script: str) -> str:
"""Maps a Klipper or Marlin G-code line to an MQTT command
for the Kobra X. Supports:
- PAUSE / M25, RESUME / M24, CANCEL_PRINT / M0/M1/M524/ABORT
- M104 S<temp> → Nozzle-Temperatur
- M140 S<temp> → Bett-Temperatur
- SET_HEATER_TEMPERATURE HEATER=extruder TARGET=200 (Klipper)
- SET_HEATER_TEMPERATURE HEATER=heater_bed TARGET=60 (Klipper)
Unknown scripts are acknowledged with 'ok' (Obico e.g. sends G28
for homing, which the bridge silently ignores)."""
if not script:
return "ok"
s = script.strip().upper()
loop = asyncio.get_event_loop()
def _parse_marlin_temp(line: str) -> int | None:
"""Extract the temperature value from 'M104 S200' or 'M140 S60'."""
try:
return int(line.split("S", 1)[1].split()[0])
except Exception:
return None
def _parse_klipper_set_heater(line: str) -> tuple[str | None, int | None]:
"""Extract heater + target from 'SET_HEATER_TEMPERATURE HEATER=extruder TARGET=143'.
Heater ID + target. Heater is 'extruder' or
'heater_bed', target is int. Returns (None,None) on error."""
heater = None
target = None
for part in line.split():
if part.startswith("HEATER="):
heater = part.split("=", 1)[1].strip().lower()
elif part.startswith("TARGET="):
try:
target = int(float(part.split("=", 1)[1]))
except Exception:
pass
return heater, target
async def _set_temps(nozzle: int | None, bed: int | None):
"""Sets nozzle/bed temperature via the correct MQTT path -
printing: print/update with taskid, idle: tempature/set with both."""
is_printing = self._state.get("print_state") in ("printing", "paused")
if is_printing:
taskid = self._state.get("taskid", "")
if nozzle is not None:
await loop.run_in_executor(None, lambda: self.client.publish_web(
"print", "update",
{"taskid": taskid, "settings": {"target_nozzle_temp": int(nozzle)}},
))
if bed is not None:
await loop.run_in_executor(None, lambda: self.client.publish_web(
"print", "update",
{"taskid": taskid, "settings": {"target_hotbed_temp": int(bed)}},
))
else:
# Idle: tempature/set via the web/printer topic with a type field
# (Live-Sniff 2026-05-29). type: 0=Nozzle, 1=Bed, 2=beide.
if nozzle is not None and bed is not None:
t, n, b = 2, int(nozzle), int(bed)
elif nozzle is not None:
t, n, b = 0, int(nozzle), 0
elif bed is not None:
t, n, b = 1, 0, int(bed)
else:
return
await loop.run_in_executor(None, lambda: self.client.publish_web(
"tempature", "set",
{"type": t, "target_nozzle_temp": n, "target_hotbed_temp": b},
))
try:
if s in ("PAUSE", "M25"):
await loop.run_in_executor(None, self.client.pause_print)
elif s in ("RESUME", "M24"):
await loop.run_in_executor(None, self.client.resume_print)
elif s in ("CANCEL_PRINT", "M0", "M1", "M524", "ABORT"):
await loop.run_in_executor(None, self.client.stop_print)
elif s.startswith("M104 "):
t = _parse_marlin_temp(s)
if t is not None:
log.info(f"gcode.script: Nozzle-Target {t}°C (M104)")
await _set_temps(t, None)
elif s.startswith("M140 "):
t = _parse_marlin_temp(s)
if t is not None:
log.info(f"gcode.script: Bed-Target {t}°C (M140)")
await _set_temps(None, t)
elif s.startswith("SET_HEATER_TEMPERATURE"):
heater, target = _parse_klipper_set_heater(s)
if target is not None and heater:
if heater == "extruder":
log.info(f"gcode.script: Nozzle-Target {target}°C (Klipper)")
await _set_temps(target, None)
elif heater in ("heater_bed", "bed"):
log.info(f"gcode.script: Bed-Target {target}°C (Klipper)")
await _set_temps(None, target)
else:
log.debug(f"gcode.script: unbekannter Heater '{heater}' ignoriert")
else:
# Unbekanntes Script: stillschweigend OK quittieren.
log.debug(f"gcode.script ignored: {s[:60]}")
except Exception as e:
log.warning(f"gcode.script {s[:30]}: {e}")
return "ok"
async def handle_printer_gcode_script(self, request):
"""HTTP POST /printer/gcode/script — Klipper-G-Code-Wrapper (siehe _exec_gcode_script)."""
script = ""
if request.method == "POST":
try:
body = await request.json()
if isinstance(body, dict):
script = body.get("script", "") or ""
except Exception:
pass
if not script:
script = request.rel_url.query.get("script", "")
result = await self._exec_gcode_script(script)
return web.json_response({"result": result})