manager.py 92 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946
  1. """Virtual Printer Manager - coordinates SSDP, MQTT, and FTP services.
  2. Each virtual printer runs its own independent services (FTP, MQTT, SSDP, Bind)
  3. bound to its dedicated IP address, regardless of mode.
  4. """
  5. import asyncio
  6. import json
  7. import logging
  8. import time
  9. from collections.abc import Callable
  10. from datetime import datetime, timezone
  11. from pathlib import Path
  12. from typing import TYPE_CHECKING
  13. from backend.app.core.config import settings as app_settings
  14. from backend.app.models.virtual_printer import (
  15. VP_MODE_ARCHIVE,
  16. VP_MODE_PROXY,
  17. VP_MODE_QUEUE,
  18. normalize_vp_mode,
  19. )
  20. from backend.app.services.bambu_mqtt import resolve_external_spools_in_mapping
  21. from backend.app.services.virtual_printer.bind_server import BindServer
  22. from backend.app.services.virtual_printer.certificate import CertificateService
  23. from backend.app.services.virtual_printer.ftp_server import VirtualPrinterFTPServer, compute_passive_port_slice
  24. from backend.app.services.virtual_printer.mqtt_bridge import MQTTBridge
  25. from backend.app.services.virtual_printer.mqtt_server import SimpleMQTTServer
  26. from backend.app.services.virtual_printer.ssdp_server import SSDPProxy, VirtualPrinterSSDPServer
  27. from backend.app.services.virtual_printer.tcp_proxy import SlicerProxyManager, TCPProxy
  28. if TYPE_CHECKING:
  29. from backend.app.services.printer_manager import PrinterManager
  30. logger = logging.getLogger(__name__)
  31. # Mapping of SSDP model codes to display names
  32. # These are the codes that slicers expect during discovery
  33. # Sources:
  34. # - https://gist.github.com/Alex-Schaefer/72a9e2491a42da2ef99fb87601955cc3
  35. # - https://github.com/psychoticbeef/BambuLabOrcaSlicerDiscovery
  36. VIRTUAL_PRINTER_MODELS = {
  37. # X1 Series
  38. "BL-P001": "X1C", # X1 Carbon
  39. "BL-P002": "X1", # X1
  40. "C13": "X1E", # X1E
  41. # X2 Series
  42. "N6": "X2D", # X2D
  43. # A2 Series (single-FDM + integrated cutter/plotter)
  44. "N9": "A2L", # A2L
  45. # P Series
  46. "C11": "P1P", # P1P
  47. "C12": "P1S", # P1S
  48. "N7": "P2S", # P2S
  49. # A1 Series
  50. "N2S": "A1", # A1
  51. "N1": "A1 Mini", # A1 Mini
  52. # H2 Series
  53. "O1D": "H2D", # H2D
  54. "O1E": "H2D Pro", # H2D Pro
  55. "O2D": "H2D Pro", # H2D Pro
  56. "O1C": "H2C", # H2C
  57. "O1C2": "H2C", # H2C (dual nozzle variant)
  58. "O1S": "H2S", # H2S
  59. }
  60. # Serial number prefixes for each model (based on Bambu Lab serial number format)
  61. # Format: MMM??RYMDDUUUUU (15 chars total)
  62. # MMM = Model prefix (3 chars)
  63. # ?? = Unknown/revision code (2 chars)
  64. # R = Revision letter (1 char)
  65. # Y = Year digit (1 char)
  66. # M = Month (1 char, hex: 1-9, A=Oct, B=Nov, C=Dec)
  67. # DD = Day (2 chars)
  68. # UUUUU = Unit number (5 chars)
  69. MODEL_SERIAL_PREFIXES = {
  70. # X1 Series
  71. "BL-P001": "00M00A", # X1C
  72. "BL-P002": "00M00A", # X1
  73. "C13": "03W00A", # X1E
  74. # X2 Series
  75. "N6": "20P90A", # X2D (first 4 chars "20P9" match real serials)
  76. # A2 Series
  77. "N9": "26A19A", # A2L (first 5 chars "26A19" match real serials)
  78. # P Series
  79. "C11": "01S00A", # P1P
  80. "C12": "01P00A", # P1S
  81. "N7": "22E00A", # P2S
  82. # A1 Series
  83. "N2S": "03900A", # A1
  84. "N1": "03000A", # A1 Mini
  85. # H2 Series
  86. "O1D": "09400A", # H2D
  87. "O1E": "09400A", # H2D Pro (same prefix family as H2D)
  88. "O2D": "09400A", # H2D Pro
  89. "O1C": "09400A", # H2C
  90. "O1C2": "09400A", # H2C (dual nozzle variant)
  91. "O1S": "09400A", # H2S
  92. }
  93. # Reverse mapping: display name → SSDP model code (for auto-inheriting from printer model)
  94. DISPLAY_NAME_TO_MODEL_CODE = {v: k for k, v in VIRTUAL_PRINTER_MODELS.items()}
  95. # Default model
  96. DEFAULT_VIRTUAL_PRINTER_MODEL = "BL-P001" # X1C
  97. # Bound on per-instance ``_slicer_print_options`` cache size. The slicer's
  98. # project_file MQTT command stashes one dict per filename; the
  99. # corresponding ``_add_to_print_queue`` pop only fires when the file
  100. # upload completes. Failed / cancelled / non-3MF uploads orphan their
  101. # stash. The bound triggers FIFO eviction in ``on_print_command`` once
  102. # the dict fills, so a long-running VP can't leak unbounded state.
  103. _SLICER_OPTIONS_CACHE_LIMIT = 128
  104. # How long ``_add_to_print_queue`` waits for the slicer's MQTT
  105. # ``project_file`` after the FTP upload completes (#1780 round 3).
  106. # Bambu Studio sends FTP first, then MQTT immediately after — but on
  107. # wireless / loaded setups the MQTT command can land 2+ s after FTP,
  108. # which used to time the wait out and silently drop ``nozzle_mapping``
  109. # + the other slicer-driven flags. The bumped window covers the
  110. # observed worst case in the field; the late-MQTT fallback in
  111. # ``on_print_command`` covers the rest.
  112. _SLICER_OPTIONS_WAIT_TIMEOUT = 5.0
  113. # How long ``on_print_command`` will retroactively stamp slicer fields
  114. # onto a recently-committed queue item when the MQTT print command
  115. # arrives after ``_SLICER_OPTIONS_WAIT_TIMEOUT`` expired. Covers
  116. # extra-late MQTT (slow wireless slicer, NIC drop+retry) and the
  117. # scheduler tick interval before dispatch picks the item up.
  118. _RECENT_QUEUE_ITEM_TTL = 30.0
  119. # BambuStudio's tri-state calibration options (bed_leveling / flow_cali /
  120. # nozzle_offset_cali) travel on the project_file command as a bool plus an int
  121. # companion — off=0, on=1, auto=2 (getValueInt parity). The int carries the full
  122. # state; the bool is true only for "on".
  123. _TRISTATE_INT = {0: "off", 1: "on", 2: "auto"}
  124. def _tristate_from_slicer(data: dict, bool_field: str, int_field: str) -> str | None:
  125. """Reconstruct off/on/auto from a captured slicer project_file dict.
  126. Prefer the int companion (auto_bed_leveling / extrude_cali_flag / etc.) which
  127. carries all three states; fall back to the bool field (on/off only); return
  128. None when the slicer sent neither so the caller can use its own default.
  129. """
  130. if int_field in data:
  131. try:
  132. resolved = _TRISTATE_INT.get(int(data[int_field]))
  133. except (TypeError, ValueError):
  134. resolved = None
  135. if resolved is not None:
  136. return resolved
  137. if bool_field in data:
  138. return "on" if bool(data[bool_field]) else "off"
  139. return None
  140. def _extract_slicer_ams_mapping_json(data: dict, log_prefix: str, *, is_dual_nozzle: bool = False) -> str | None:
  141. """Pull the slicer's own AMS-slot pick out of a captured project_file payload.
  142. BambuStudio/OrcaSlicer resolves the physical AMS tray for each filament
  143. live, right before sending — either automatically or via the slicer's
  144. manual per-filament AMS-slot assignment dialog — and embeds the result as
  145. ``ams_mapping`` (``list[int]``, position = slot_id-1, value = global tray
  146. ID) directly in the MQTT ``project_file`` command. Confirmed by wire
  147. capture: the field is present and already in the exact shape
  148. ``PrintQueueItem.ams_mapping`` expects.
  149. The VP-queue path previously never read this — every queued print had the
  150. scheduler re-derive a mapping from just the 3MF's static type/color at
  151. dispatch time (`PrintScheduler._compute_ams_mapping_for_printer`), discarding
  152. the slicer's already-correct, live-resolved pick. That re-derivation can
  153. land on the wrong physical spool whenever the file's type+color match
  154. isn't unique (e.g. two spools of the same color) or the file's own
  155. filament-slot color wasn't what the user actually intended for that
  156. particular print. Capturing it here — mirroring the existing
  157. ``nozzle_mapping`` passthrough for H2C rack-swap models (#1780) — lets the
  158. scheduler's "already resolved, don't touch it" branch in
  159. ``_ensure_ams_mapping`` use the slicer's own choice unchanged.
  160. That branch skipping ``_compute_ams_mapping_for_printer`` is also what
  161. makes this a trade rather than a pure win: ``prefer_lowest_filament``, its
  162. AMS-filament-backup gate (#1766), the inventory-remain overrides and the
  163. per-slot force-color overrides all live inside that function. Callers are
  164. responsible for the gating — this parser only says what the slicer sent.
  165. An external spool is ``-1`` in the flat list, the same as a filament with
  166. no tray, and is named only in ``ams_mapping2``. It is resolved from there
  167. to its global tray (254/255), as the usage capture does (#3166); dispatch
  168. turns that back into the external-spool entry. Kept as ``-1`` it dispatched
  169. as unmapped and the printer stopped with 0700-8012 (#3237).
  170. ``is_dual_nozzle`` is the target printer's, since that decides which
  171. global tray an external spool is.
  172. Returns ``None`` when the field is absent, unparsable, or the classic
  173. "all -1" unresolved-race sentinel (#2589) — never worth trusting over a
  174. fresh live computation.
  175. """
  176. raw = data.get("ams_mapping")
  177. if raw is None:
  178. return None
  179. if isinstance(raw, str):
  180. try:
  181. raw = json.loads(raw)
  182. except json.JSONDecodeError:
  183. logger.warning("%s Slicer ams_mapping is unparseable JSON, dropping: %r", log_prefix, raw)
  184. return None
  185. # bool is a subclass of int in Python — isinstance(True, int) is True —
  186. # so it must be excluded explicitly, or [True, False] would pass as a
  187. # valid mapping.
  188. if not isinstance(raw, list) or not raw or not all(isinstance(v, int) and not isinstance(v, bool) for v in raw):
  189. return None
  190. detail = data.get("ams_mapping2")
  191. if isinstance(detail, str):
  192. try:
  193. detail = json.loads(detail)
  194. except json.JSONDecodeError:
  195. detail = None
  196. raw = resolve_external_spools_in_mapping(raw, detail, is_dual_nozzle)
  197. if all(v < 0 for v in raw):
  198. # #2589 sentinel — every slot unresolved. Let the scheduler compute a
  199. # fresh mapping from live AMS state instead of trusting this.
  200. return None
  201. return json.dumps(raw)
  202. def _get_serial_for_model(model: str, serial_suffix: str) -> str:
  203. """Get serial number for the given model and suffix."""
  204. prefix = MODEL_SERIAL_PREFIXES.get(model, "00M09A")
  205. return f"{prefix}{serial_suffix}"
  206. class VirtualPrinterInstance:
  207. """Per-printer state and file handling logic.
  208. Each instance represents one virtual printer with its own config,
  209. upload directory, certificates, and file handling mode.
  210. """
  211. def __init__(
  212. self,
  213. *,
  214. vp_id: int,
  215. name: str,
  216. mode: str,
  217. model: str,
  218. access_code: str,
  219. serial_suffix: str,
  220. target_printer_ip: str = "",
  221. target_printer_serial: str = "",
  222. target_printer_id: int | None = None,
  223. auto_dispatch: bool = True,
  224. queue_force_color_match: bool = False,
  225. save_ams_mapping: bool = False,
  226. gcode_injection: bool = False,
  227. bind_ip: str = "",
  228. remote_interface_ip: str = "",
  229. tailscale_disabled: bool = True,
  230. base_dir: Path,
  231. session_factory: Callable | None = None,
  232. printer_manager: "PrinterManager | None" = None,
  233. ):
  234. self.id = vp_id
  235. self.name = name
  236. # Normalize on construction so the rest of the code only compares
  237. # canonical values, even when a legacy DB row hasn't been migrated
  238. # yet (e.g. fresh-from-disk during the boot window before the
  239. # one-shot migration in `core/database.py` has executed).
  240. self.mode = normalize_vp_mode(mode) or VP_MODE_ARCHIVE
  241. self.model = model
  242. self.access_code = access_code
  243. self.serial_suffix = serial_suffix
  244. self.target_printer_ip = target_printer_ip
  245. self.target_printer_serial = target_printer_serial
  246. self.target_printer_id = target_printer_id
  247. self.auto_dispatch = auto_dispatch
  248. self.queue_force_color_match = queue_force_color_match
  249. self.save_ams_mapping = save_ams_mapping
  250. self.gcode_injection = gcode_injection
  251. self.bind_ip = bind_ip
  252. self.remote_interface_ip = remote_interface_ip
  253. self.tailscale_disabled = tailscale_disabled
  254. self._session_factory = session_factory
  255. self._printer_manager = printer_manager
  256. # Directories
  257. self.upload_dir = base_dir / "uploads" / str(vp_id)
  258. self.cert_dir = base_dir / "certs" / str(vp_id)
  259. shared_ca_dir = base_dir / "certs"
  260. # Ensure directories exist
  261. self.upload_dir.mkdir(parents=True, exist_ok=True)
  262. (self.upload_dir / "cache").mkdir(exist_ok=True)
  263. self.cert_dir.mkdir(parents=True, exist_ok=True)
  264. # Certificate service (shared CA, per-instance printer cert)
  265. self._cert_service = CertificateService(
  266. cert_dir=self.cert_dir,
  267. serial=self.serial,
  268. shared_ca_dir=shared_ca_dir,
  269. )
  270. # Pending files for MQTT correlation
  271. self._pending_files: dict[str, Path] = {}
  272. # Slicer-side print options captured from the MQTT `project_file`
  273. # command, keyed by filename. Used by `_add_to_print_queue` so the
  274. # queue item inherits the user's slicer-chosen timelapse / bed_leveling
  275. # / flow_cali / vibration_cali / layer_inspect / use_ams toggles rather
  276. # than falling back to the global `default_*` settings (#1403). FTP
  277. # completes a few hundred ms before the slicer's MQTT `project_file`
  278. # arrives, so the queue-add path waits briefly on the event below
  279. # before reading the dict. Events are popped along with the options
  280. # so the dict stays bounded.
  281. self._slicer_print_options: dict[str, dict] = {}
  282. self._slicer_print_options_events: dict[str, asyncio.Event] = {}
  283. # Queue items recently committed by `_add_to_print_queue`, keyed by
  284. # FTP filename. Used by `on_print_command` to retroactively stamp the
  285. # slicer's nozzle_mapping (and the other slicer-driven flags) onto a
  286. # queue item when the MQTT `project_file` arrives after the queue-add
  287. # wait timed out — the #1780 round-3 race. Value is
  288. # (queue_item_ids, monotonic_committed_at); entries older than
  289. # `_RECENT_QUEUE_ITEM_TTL` are evicted opportunistically on each
  290. # queue-add.
  291. self._recent_queue_items: dict[str, tuple[list[int], float]] = {}
  292. # Per-instance services
  293. self._proxy: SlicerProxyManager | None = None
  294. self._ftp: VirtualPrinterFTPServer | None = None
  295. self._mqtt: SimpleMQTTServer | None = None
  296. self._mqtt_bridge: MQTTBridge | None = None
  297. self._rtsp_proxy: TCPProxy | None = None
  298. self._bind: BindServer | None = None
  299. self._ssdp: VirtualPrinterSSDPServer | None = None
  300. self._ssdp_proxy: SSDPProxy | None = None
  301. self._tasks: list[asyncio.Task] = []
  302. # Pending timer that re-fires gcode_state=FINISH after a project_file
  303. # ack. See ``_schedule_finish_release`` for the #1658 rationale.
  304. self._finish_release_task: asyncio.Task | None = None
  305. @property
  306. def serial(self) -> str:
  307. """Full serial number for this virtual printer."""
  308. return _get_serial_for_model(self.model or DEFAULT_VIRTUAL_PRINTER_MODEL, self.serial_suffix)
  309. @property
  310. def cert_path(self) -> Path:
  311. return self._cert_service.cert_path
  312. @property
  313. def key_path(self) -> Path:
  314. return self._cert_service.key_path
  315. @property
  316. def is_proxy(self) -> bool:
  317. return self.mode == "proxy"
  318. @property
  319. def is_running(self) -> bool:
  320. return len(self._tasks) > 0 and all(not t.done() for t in self._tasks)
  321. def generate_certificates(self) -> tuple[Path, Path]:
  322. """Generate certificates for this instance."""
  323. self._cert_service.serial = self.serial if not self.is_proxy else (self.target_printer_serial or self.serial)
  324. additional_ips = [self.remote_interface_ip] if self.remote_interface_ip else None
  325. if self.bind_ip:
  326. additional_ips = additional_ips or []
  327. additional_ips.append(self.bind_ip)
  328. self._cert_service.delete_printer_certificate()
  329. return self._cert_service.generate_certificates(additional_ips=additional_ips)
  330. # -- File handling callbacks --
  331. async def on_file_received(self, file_path: Path, source_ip: str) -> None:
  332. """Handle file upload completion from FTP."""
  333. logger.info("[VP %s] Received file: %s from %s", self.name, file_path.name, source_ip)
  334. self._pending_files[file_path.name] = file_path
  335. # Accept both canonical (`archive`/`queue`) and legacy
  336. # (`immediate`/`print_queue`) wire values so a stale row that hasn't
  337. # been migrated yet still dispatches correctly. Migration in
  338. # `core/database.py` rewrites existing rows once at boot.
  339. mode = normalize_vp_mode(self.mode)
  340. if mode == VP_MODE_ARCHIVE:
  341. await self._archive_file(file_path, source_ip)
  342. elif mode == VP_MODE_QUEUE:
  343. await self._add_to_print_queue(file_path, source_ip)
  344. else:
  345. await self._queue_file(file_path, source_ip)
  346. # Signal job completion to the slicer. Send-flow slicers don't watch the
  347. # post-upload state and would be happy with anything; the Print flow
  348. # (intended for proxy-mode VPs, but users sometimes click it against
  349. # queue/immediate/review modes too — #1280) watches the gcode_state
  350. # cycle and only releases its in-flight-job lock when it sees FINISH.
  351. # Going PREPARE → IDLE wedges the slicer's UI at "Downloading...(0%)"
  352. # and blocks the next dispatch with "busy with another print job".
  353. # PREPARE → FINISH satisfies both flows. prepare_percent=100 also
  354. # unfreezes the slicer's "Downloading X%" progress bar which it ticks
  355. # against the same field during the upload window.
  356. if self._mqtt and file_path.suffix.lower() == ".3mf":
  357. self._mqtt.set_gcode_state("FINISH", filename=file_path.name, prepare_percent="100")
  358. # FINISH is the terminal state for the upload cycle per #1280
  359. # (commit 0d6171dc). The Print-flow slicer's in-flight-job lock
  360. # releases on FINISH; resetting to IDLE 2 s later would re-confuse
  361. # the slicer that just unwedged. Earlier audit suggesting the
  362. # IDLE reset was wrong — staying at FINISH is the designed
  363. # behaviour. The next upload's PREPARE→FINISH cycle starts fresh.
  364. def _target_is_dual_nozzle(self) -> bool:
  365. """Whether the target printer has two nozzles, which decides the global
  366. tray of an external spool in the slicer's mapping (#3237).
  367. Same signal dispatch uses: the live client's detection, else its model.
  368. Without a client, the VP's own model is what the slicer sliced for.
  369. """
  370. from backend.app.utils.printer_models import is_dual_nozzle_model
  371. client = (
  372. self._printer_manager.get_client(self.target_printer_id)
  373. if self._printer_manager is not None and self.target_printer_id is not None
  374. else None
  375. )
  376. if client is not None:
  377. return bool(getattr(client, "_is_dual_nozzle", False)) or is_dual_nozzle_model(
  378. getattr(client, "model", None)
  379. )
  380. return is_dual_nozzle_model(self.model)
  381. async def on_print_command(self, filename: str, data: dict) -> None:
  382. """Handle print command from MQTT.
  383. Captures the slicer's project_file options (`timelapse`, `bed_leveling`,
  384. `flow_cali`, `vibration_cali`, `layer_inspect`, `use_ams`, plus the
  385. H2C rack-pick `nozzle_mapping`) so the VP-queue path can inherit them
  386. when adding the item to the queue, rather than falling back to the
  387. global default settings (#1403, #1780).
  388. Only queue mode consumes the capture; archive / review / proxy
  389. modes ignore the print command, so we skip the stash there to keep
  390. the dict from accumulating one entry per print over the VP's
  391. uptime.
  392. Also schedules the #1658 follow-up that re-fires gcode_state=FINISH a
  393. moment after the synthetic project_file ack — for every non-proxy
  394. mode — so the slicer's "Downloading" UI releases on the slicer's
  395. FTP-first-then-MQTT send order.
  396. ``filename`` is the slicer's ``subtask_name`` (bare model name, no
  397. extension) — used verbatim for `_schedule_finish_release` because
  398. push_status echoes it back to the slicer as gcode_file / subtask_name.
  399. The queue-side stash key is derived from ``data["file"]`` (the FTP
  400. filename with extension) so `_add_to_print_queue`'s
  401. ``file_path.name`` lookup matches; falls back to ``filename`` when
  402. ``data["file"]`` is absent (legacy slicers / non-3MF uploads).
  403. Stash/lookup mismatch was the #1780 root cause — every captured field
  404. silently fell back to settings defaults on every Bambu Studio "Send".
  405. """
  406. logger.info("[VP %s] Print command for: %s", self.name, filename)
  407. mode = normalize_vp_mode(self.mode)
  408. if mode != VP_MODE_PROXY and filename and self._mqtt is not None:
  409. self._schedule_finish_release(filename)
  410. if mode != VP_MODE_QUEUE:
  411. return
  412. # Stash key must match `_add_to_print_queue`'s lookup, which uses
  413. # `file_path.name` (FTP filename WITH extension). The slicer's
  414. # `subtask_name` (== this method's `filename` arg) is the bare model
  415. # name, no extension — using it as the stash key was the #1780 root
  416. # cause.
  417. stash_key = data.get("file") or filename
  418. # Drop the oldest stash if the cache is growing — happens when the
  419. # slicer sends project_file for a filename whose FTP upload was
  420. # rejected / cancelled / non-3MF, so _add_to_print_queue's pop
  421. # never fires. With no bound, a long-running VP accumulates one
  422. # dict per such mismatch.
  423. if len(self._slicer_print_options) >= _SLICER_OPTIONS_CACHE_LIMIT:
  424. try:
  425. stale_key = next(iter(self._slicer_print_options))
  426. self._slicer_print_options.pop(stale_key, None)
  427. self._slicer_print_options_events.pop(stale_key, None)
  428. logger.debug("[VP %s] Evicted stale slicer options for %s", self.name, stale_key)
  429. except StopIteration:
  430. pass
  431. self._slicer_print_options[stash_key] = dict(data)
  432. event = self._slicer_print_options_events.get(stash_key)
  433. if event:
  434. event.set()
  435. return
  436. # No consumer waiting: `_add_to_print_queue` either already gave up
  437. # (wait_for timed out) or hasn't started yet (FTP still uploading).
  438. # If a queue item was committed within the last
  439. # `_RECENT_QUEUE_ITEM_TTL`, the wait timed out and the row holds
  440. # settings defaults instead of the slicer's pick — retroactively
  441. # stamp the slicer-driven fields so the dispatcher honours the
  442. # user's choice. Covers the #1780 round-3 race where Bambu Studio's
  443. # MQTT lands just past the bumped wait ceiling.
  444. await self._restamp_recent_queue_item(stash_key, data)
  445. async def _restamp_recent_queue_item(self, stash_key: str, data: dict) -> None:
  446. """Patch slicer-driven fields onto a queue item the MQTT command missed.
  447. ``_add_to_print_queue`` waits up to ``_SLICER_OPTIONS_WAIT_TIMEOUT``
  448. for the slicer's MQTT ``project_file`` before committing the queue
  449. item. If the MQTT command arrives after that window — observed in
  450. the field at ~2.1 s on H2C / wireless setups (#1780 round 3) — the
  451. row was already written with settings defaults. This method runs
  452. on the late MQTT path: it looks up the most recent queue items
  453. committed for this filename and patches in the slicer's
  454. ``nozzle_mapping`` + ``ams_mapping`` + workflow flags, but only
  455. while the items are still ``pending`` (scheduler hasn't dispatched
  456. them yet).
  457. """
  458. if not self._session_factory:
  459. return
  460. entry = self._recent_queue_items.get(stash_key)
  461. if entry is None:
  462. return
  463. queue_item_ids, committed_at = entry
  464. if time.monotonic() - committed_at > _RECENT_QUEUE_ITEM_TTL:
  465. self._recent_queue_items.pop(stash_key, None)
  466. return
  467. import json
  468. # Mirror the field set `_add_to_print_queue` reads off slicer_opts.
  469. # MQTT uses `bed_leveling` (single L); the column is `bed_levelling`.
  470. # `nozzles_info` is intentionally not stamped — column kept for
  471. # legacy rows but never written; see PrintQueueItem.nozzles_info.
  472. patch: dict = {}
  473. # Tri-state options (off/on/auto) — reconstruct from the int companion.
  474. for bool_field, int_field, column in (
  475. ("bed_leveling", "auto_bed_leveling", "bed_levelling"),
  476. ("flow_cali", "extrude_cali_flag", "flow_cali"),
  477. ):
  478. resolved = _tristate_from_slicer(data, bool_field, int_field)
  479. if resolved is not None:
  480. patch[column] = resolved
  481. # On/off options.
  482. for mqtt_field, column in (
  483. ("vibration_cali", "vibration_cali"),
  484. ("layer_inspect", "layer_inspect"),
  485. ("timelapse", "timelapse"),
  486. ("use_ams", "use_ams"),
  487. ):
  488. if mqtt_field in data:
  489. patch[column] = bool(data[mqtt_field])
  490. raw = data.get("nozzle_mapping")
  491. if raw is not None:
  492. if isinstance(raw, str):
  493. try:
  494. raw = json.loads(raw)
  495. except json.JSONDecodeError:
  496. logger.warning(
  497. "[VP %s] Late MQTT nozzle_mapping is unparseable JSON, dropping: %r",
  498. self.name,
  499. raw,
  500. )
  501. raw = None
  502. if raw is not None:
  503. patch["nozzle_mapping"] = json.dumps(raw)
  504. # Same two gates as the immediate path in `_add_to_print_queue`: a
  505. # model-based VP has no live AMS layout for the slicer to have resolved
  506. # tray IDs against, and taking the slicer's pick at all is the per-VP
  507. # `save_ams_mapping` opt-in (it makes the scheduler skip
  508. # `_compute_ams_mapping_for_printer`, and with it prefer-lowest and the
  509. # #1766 backup gate).
  510. ams_mapping_json = (
  511. _extract_slicer_ams_mapping_json(
  512. data, f"[VP {self.name}] Late MQTT", is_dual_nozzle=self._target_is_dual_nozzle()
  513. )
  514. if self.target_printer_id is not None and self.save_ams_mapping
  515. else None
  516. )
  517. # `Force color match` still wins for this dispatch — see the same
  518. # decision in `_add_to_print_queue`. The archive patch below is
  519. # deliberately not gated on it: persisting the pick for later reprints
  520. # is exactly what the toggle promises.
  521. if ams_mapping_json is not None and not self.queue_force_color_match:
  522. patch["ams_mapping"] = ams_mapping_json
  523. # `ams_mapping_json` alone is enough to keep going even when `patch` is
  524. # empty: with `Force color match` on it never reaches the queue item,
  525. # but it still has to be written onto the archive below.
  526. if not patch and ams_mapping_json is None:
  527. self._recent_queue_items.pop(stash_key, None)
  528. return
  529. from sqlalchemy import select, update
  530. from backend.app.models.archive import PrintArchive
  531. from backend.app.models.print_queue import PrintQueueItem
  532. try:
  533. async with self._session_factory() as db:
  534. # Only stamp items still pending; once the scheduler has
  535. # picked the row up we can't safely race the dispatcher.
  536. result = await db.execute(
  537. select(PrintQueueItem.id, PrintQueueItem.archive_id).where(
  538. PrintQueueItem.id.in_(queue_item_ids),
  539. PrintQueueItem.status == "pending",
  540. )
  541. )
  542. rows = result.all()
  543. eligible_ids = [row[0] for row in rows]
  544. if not eligible_ids:
  545. self._recent_queue_items.pop(stash_key, None)
  546. return
  547. if patch:
  548. await db.execute(update(PrintQueueItem).where(PrintQueueItem.id.in_(eligible_ids)).values(**patch))
  549. # The archive was already created (with no slicer_ams_mapping)
  550. # before this late MQTT arrived — see
  551. # `_extract_slicer_ams_mapping_json`'s docstring. Patch it here
  552. # too so a reprint later still picks up the slicer's pick, and
  553. # the "AMS mapping from slicer" badge reflects reality instead
  554. # of staying stuck on the archive's initial (empty) snapshot.
  555. # Already gated on `save_ams_mapping` above, and deliberately
  556. # NOT on `queue_force_color_match`: that toggle decides how
  557. # *this* print is matched, not whether the pick is worth
  558. # keeping for a later reprint.
  559. if ams_mapping_json is not None:
  560. archive_ids = {row[1] for row in rows if row[1] is not None}
  561. if archive_ids:
  562. archive_result = await db.execute(select(PrintArchive).where(PrintArchive.id.in_(archive_ids)))
  563. for archive in archive_result.scalars().all():
  564. extra = dict(archive.extra_data or {})
  565. extra["slicer_ams_mapping"] = {
  566. "mapping": json.loads(ams_mapping_json),
  567. "printer_id": self.target_printer_id,
  568. }
  569. archive.extra_data = extra
  570. await db.commit()
  571. logger.info(
  572. "[VP %s] Late slicer MQTT for %s — retroactively stamped %s onto queue item(s) %s%s",
  573. self.name,
  574. stash_key,
  575. sorted(patch.keys()),
  576. eligible_ids,
  577. " and saved the slicer's AMS pick onto the archive" if ams_mapping_json is not None else "",
  578. )
  579. except Exception as e:
  580. logger.error(
  581. "[VP %s] Failed to retroactively stamp queue item(s) %s for %s: %s",
  582. self.name,
  583. queue_item_ids,
  584. stash_key,
  585. e,
  586. )
  587. finally:
  588. self._recent_queue_items.pop(stash_key, None)
  589. def _schedule_finish_release(self, filename: str, delay: float = 1.5) -> None:
  590. """Re-set gcode_state=FINISH on the VP after the project_file ack.
  591. #1280 set FINISH after the FTP upload completes — that was correct
  592. for the slicer flow at the time (MQTT project_file → FTP → done).
  593. Bambu Studio 2.7.x flipped the order to FTP → FTP → MQTT project_file,
  594. which means ``_send_print_response`` runs *after* the FINISH set in
  595. ``on_file_received`` and overwrites the state back to PREPARE. The
  596. slicer's 1 Hz status stream then carries PREPARE forever and the
  597. send modal sits at "Downloading" until the VP is restarted (#1658).
  598. Re-firing FINISH after a short delay closes the gap: the slicer sees
  599. the synthetic PREPARE in the project_file ack (and likely one PREPARE
  600. push on the 1 Hz cycle), then the next push carries FINISH and the
  601. modal releases. Proxy mode is exempt — there the real printer drives
  602. the state through the bridge and a synthetic FINISH would clobber a
  603. real PREPARE/RUNNING transition coming back from the printer.
  604. Cancels any in-flight timer before scheduling a new one so a slicer
  605. that fires project_file twice in quick succession only ends in one
  606. FINISH.
  607. """
  608. if self._mqtt is None:
  609. return
  610. if self._finish_release_task is not None and not self._finish_release_task.done():
  611. self._finish_release_task.cancel()
  612. self._finish_release_task = asyncio.create_task(
  613. self._delayed_finish_release(filename, delay),
  614. name=f"vp-{self.id}-finish-release",
  615. )
  616. async def _delayed_finish_release(self, filename: str, delay: float) -> None:
  617. """Sleep, then set gcode_state=FINISH. Used by ``_schedule_finish_release``."""
  618. try:
  619. await asyncio.sleep(delay)
  620. except asyncio.CancelledError:
  621. return
  622. if self._mqtt is None:
  623. return
  624. self._mqtt.set_gcode_state("FINISH", filename=filename, prepare_percent="100")
  625. logger.debug("[VP %s] Re-set gcode_state=FINISH after project_file ack (%s)", self.name, filename)
  626. async def _archive_file(self, file_path: Path, source_ip: str) -> None:
  627. """Archive file immediately."""
  628. if not self._session_factory:
  629. logger.error("Cannot archive: no database session factory configured")
  630. return
  631. if file_path.suffix.lower() != ".3mf":
  632. logger.debug("Skipping non-3MF file: %s", file_path.name)
  633. self._pending_files.pop(file_path.name, None)
  634. try:
  635. file_path.unlink()
  636. except OSError:
  637. pass
  638. return
  639. archived = False
  640. try:
  641. from backend.app.api.routes.settings import get_setting
  642. from backend.app.services.archive import ArchiveService
  643. async with self._session_factory() as db:
  644. name_source = await get_setting(db, "virtual_printer_archive_name_source")
  645. prefer_filename = name_source == "filename"
  646. service = ArchiveService(db)
  647. archive = await service.archive_print(
  648. printer_id=None,
  649. source_file=file_path,
  650. print_data={
  651. "status": "archived",
  652. "source": "virtual_printer",
  653. "source_ip": source_ip,
  654. },
  655. prefer_filename_for_name=prefer_filename,
  656. )
  657. if archive:
  658. logger.info("[VP %s] Archived: %s - %s", self.name, archive.id, archive.print_name)
  659. await self._broadcast_archive_created(archive)
  660. archived = True
  661. else:
  662. logger.error("Failed to archive file: %s", file_path.name)
  663. except Exception as e:
  664. logger.error("Error archiving file: %s", e)
  665. finally:
  666. # Always release the in-flight marker and delete the temp file —
  667. # previously the failure paths only logged and the next upload of
  668. # the same name was silently rejected with "already uploading",
  669. # the upload_dir filled up indefinitely, and the slicer received
  670. # a clean 226 even though no archive existed (#audit-R2-1).
  671. self._pending_files.pop(file_path.name, None)
  672. if archived:
  673. try:
  674. file_path.unlink()
  675. except OSError:
  676. pass
  677. else:
  678. # Drop the failed temp file so it doesn't accumulate.
  679. try:
  680. file_path.unlink(missing_ok=True)
  681. except OSError:
  682. pass
  683. async def _queue_file(self, file_path: Path, source_ip: str) -> None:
  684. """Queue file for user review."""
  685. if not self._session_factory:
  686. logger.error("Cannot queue: no database session factory configured")
  687. return
  688. if file_path.suffix.lower() != ".3mf":
  689. self._pending_files.pop(file_path.name, None)
  690. try:
  691. file_path.unlink()
  692. except OSError:
  693. pass
  694. return
  695. # Peek at the 3MF for the embedded title BEFORE we hand it off to the
  696. # DB. Storing it now means the /pending-uploads/ list doesn't have to
  697. # reopen every 3MF on every render to keep the review card and the
  698. # eventual archive name in sync (#1152 follow-up). Failure to parse is
  699. # not fatal — the response model falls back to the filename stem.
  700. metadata_print_name: str | None = None
  701. try:
  702. from backend.app.services.archive import ThreeMFParser
  703. parsed = ThreeMFParser(file_path).parse()
  704. raw_name = parsed.get("print_name")
  705. if isinstance(raw_name, str) and raw_name.strip():
  706. metadata_print_name = raw_name.strip()[:255]
  707. except Exception as e:
  708. logger.debug("[VP %s] Metadata title peek failed for %s: %s", self.name, file_path.name, e)
  709. try:
  710. from backend.app.models.pending_upload import PendingUpload
  711. async with self._session_factory() as db:
  712. pending = PendingUpload(
  713. filename=file_path.name,
  714. file_path=str(file_path),
  715. file_size=file_path.stat().st_size,
  716. source_ip=source_ip,
  717. status="pending",
  718. uploaded_at=datetime.now(timezone.utc),
  719. metadata_print_name=metadata_print_name,
  720. )
  721. db.add(pending)
  722. await db.commit()
  723. logger.info("[VP %s] Queued: %s - %s", self.name, pending.id, file_path.name)
  724. except Exception as e:
  725. logger.error("Error queueing file: %s", e)
  726. # Queue insert failed — drop the temp file so it doesn't
  727. # accumulate. The file is unreachable without the DB row.
  728. try:
  729. file_path.unlink(missing_ok=True)
  730. except OSError:
  731. pass
  732. finally:
  733. # Always release the in-flight marker so concurrent uploads
  734. # with the same filename aren't spuriously rejected after
  735. # a queue failure.
  736. self._pending_files.pop(file_path.name, None)
  737. async def _add_to_print_queue(self, file_path: Path, source_ip: str) -> None:
  738. """Archive file and add to print queue, assigned to target printer or model."""
  739. if not self._session_factory:
  740. logger.error("Cannot add to print queue: no database session factory configured")
  741. return
  742. if file_path.suffix.lower() != ".3mf":
  743. self._pending_files.pop(file_path.name, None)
  744. try:
  745. file_path.unlink()
  746. except OSError:
  747. pass
  748. return
  749. # Wait briefly for the slicer's MQTT `project_file` command so the
  750. # queue item can inherit the slicer-side print options the user
  751. # picked (timelapse, bed_leveling, etc). Slicers send the FTP upload
  752. # first and the MQTT command immediately after, so the typical lag
  753. # is a few hundred ms. The window is generous enough to absorb
  754. # wireless / loaded-Pi jitter without making every VP-queue add
  755. # visibly slow — observed worst case in #1780 round 3 was 2.085 s,
  756. # the previous 2.0 s ceiling. Falls back to the global default_*
  757. # settings if MQTT doesn't arrive in time (legacy behaviour for
  758. # users on a slicer that doesn't send a print command). #1403.
  759. # The wait is skipped when there's no MQTT server attached — covers
  760. # unit tests that invoke `_add_to_print_queue` directly without
  761. # going through `on_print_command`, so they don't pay the wait tax.
  762. slicer_opts = self._slicer_print_options.pop(file_path.name, None)
  763. if slicer_opts is None and self._mqtt is not None:
  764. event = asyncio.Event()
  765. self._slicer_print_options_events[file_path.name] = event
  766. try:
  767. await asyncio.wait_for(event.wait(), timeout=_SLICER_OPTIONS_WAIT_TIMEOUT)
  768. slicer_opts = self._slicer_print_options.pop(file_path.name, None)
  769. except asyncio.TimeoutError:
  770. slicer_opts = None
  771. finally:
  772. self._slicer_print_options_events.pop(file_path.name, None)
  773. # If the cache still misses, queued workflow flags / nozzle pick will
  774. # silently fall back to settings defaults. Surface the missed key so a
  775. # future stash/lookup mismatch (the #1780 root cause) is obvious in
  776. # the log instead of needing a wire capture to diagnose.
  777. if slicer_opts is None:
  778. logger.debug(
  779. "[VP %s] No slicer options cached for %r (cache keys: %s); "
  780. "workflow flags + nozzle pick will fall back to settings defaults.",
  781. self.name,
  782. file_path.name,
  783. sorted(self._slicer_print_options.keys()),
  784. )
  785. try:
  786. import json
  787. from backend.app.api.routes.settings import get_setting
  788. from backend.app.models.print_queue import PrintQueueItem
  789. from backend.app.services.archive import ArchiveService
  790. from backend.app.services.filament_requirements import extract_filament_requirements
  791. from backend.app.services.print_confirmation import confirm_outcome_for_new_queue_item
  792. async with self._session_factory() as db:
  793. name_source = await get_setting(db, "virtual_printer_archive_name_source")
  794. prefer_filename = name_source == "filename"
  795. # Read workflow defaults from settings. Without this the
  796. # PrintQueueItem below would fall back to the column-level
  797. # defaults and ignore the user's workflow preferences (#1235).
  798. # Fallbacks match AppSettings defaults in schemas/settings.py.
  799. # The slicer-side options captured above (if any) take
  800. # precedence per-field over these defaults.
  801. def _bool_setting(value: str | None, default: bool) -> bool:
  802. return value.lower() == "true" if value is not None else default
  803. def _tristate_setting(value: str | None, default: str) -> str:
  804. """Tri-state workflow default, coercing legacy true/false rows."""
  805. if value is None:
  806. return default
  807. low = value.strip().lower()
  808. if low in ("on", "off", "auto"):
  809. return low
  810. if low in ("true", "1"):
  811. return "on"
  812. if low in ("false", "0"):
  813. return "off"
  814. return default
  815. def _slicer_or(field_mqtt: str, settings_default: bool) -> bool:
  816. """Slicer's MQTT value if present, else the settings default.
  817. Slicer payloads carry both bool and int (0/1) shapes
  818. depending on firmware family — coerce via bool() so
  819. `0`/`False` and `1`/`True` both work.
  820. """
  821. if slicer_opts is not None and field_mqtt in slicer_opts:
  822. return bool(slicer_opts[field_mqtt])
  823. return settings_default
  824. def _slicer_tristate(bool_field: str, int_field: str, settings_default: str) -> str:
  825. """Slicer's tri-state (off/on/auto) if present, else the default."""
  826. if slicer_opts is not None:
  827. resolved = _tristate_from_slicer(slicer_opts, bool_field, int_field)
  828. if resolved is not None:
  829. return resolved
  830. return settings_default
  831. # Note the MQTT field names differ from Bambuddy's column
  832. # names: MQTT uses `bed_leveling` (single L) while the
  833. # column / settings key use `bed_levelling` (double L).
  834. bed_levelling = _slicer_tristate(
  835. "bed_leveling",
  836. "auto_bed_leveling",
  837. _tristate_setting(await get_setting(db, "default_bed_levelling"), "auto"),
  838. )
  839. flow_cali = _slicer_tristate(
  840. "flow_cali",
  841. "extrude_cali_flag",
  842. _tristate_setting(await get_setting(db, "default_flow_cali"), "auto"),
  843. )
  844. vibration_cali = _slicer_or(
  845. "vibration_cali", _bool_setting(await get_setting(db, "default_vibration_cali"), True)
  846. )
  847. layer_inspect = _slicer_or(
  848. "layer_inspect", _bool_setting(await get_setting(db, "default_layer_inspect"), False)
  849. )
  850. timelapse = _slicer_or("timelapse", _bool_setting(await get_setting(db, "default_timelapse"), False))
  851. # "Ask for Outcome" is a default print option like the ones
  852. # above. A plate sent here from Bambu Studio is also one of the
  853. # prints `confirm_outcome_external_prints` names — but it
  854. # arrives with a queue item, so on_print_start resumes the
  855. # archive created below instead of treating it as external, and
  856. # without this the setting could never reach it (#1898).
  857. confirm_outcome = await confirm_outcome_for_new_queue_item(db, started_outside_bambuddy=True)
  858. # H2C dual-nozzle-rack slicer-pick preservation (#1780).
  859. # BambuStudio's project_file MQTT command for rack-swap models
  860. # (O1C2 today) carries `nozzle_mapping` — a per-filament array
  861. # of physical nozzle position IDs (`list[int]`). Forward it
  862. # verbatim onto the queue item so the dispatcher can replay it
  863. # in its own project_file command. Without this the H2C
  864. # firmware falls back to "last matching nozzle" auto-pick and
  865. # ignores the user's Bambu Studio choice. Every other model
  866. # has it absent from slicer_opts, so the capture is a
  867. # transparent no-op there. (`nozzles_info` was also captured
  868. # in the original fix but BambuStudio never actually sends it
  869. # — verified via wire capture on H2C — so only `nozzle_mapping`
  870. # is forwarded now.)
  871. nozzle_mapping_json: str | None = None
  872. if slicer_opts is not None:
  873. raw = slicer_opts.get("nozzle_mapping")
  874. if raw is not None:
  875. # BambuStudio's NetworkAgent embeds this as parsed
  876. # JSON in the project_file body (matching the
  877. # ams_mapping shape Bambuddy already consumes as
  878. # list[int]). Accept a JSON-encoded string defensively
  879. # in case any path arrives stringified.
  880. if isinstance(raw, str):
  881. try:
  882. raw = json.loads(raw)
  883. except json.JSONDecodeError:
  884. logger.warning(
  885. "[VP %s] Slicer nozzle_mapping is unparseable JSON, dropping: %r",
  886. self.name,
  887. raw,
  888. )
  889. raw = None
  890. if raw is not None:
  891. nozzle_mapping_json = json.dumps(raw)
  892. # Slicer's own live-resolved AMS-slot pick (see docstring on
  893. # `_extract_slicer_ams_mapping_json`). Stamped onto every plate
  894. # below, same treatment as nozzle_mapping_json above — when
  895. # present it makes `_ensure_ams_mapping` skip its own
  896. # type/color re-derivation entirely and dispatch use exactly
  897. # the tray the slicer/user picked.
  898. #
  899. # Two gates, both required:
  900. #
  901. # 1. This VP must target one fixed printer. A model-based
  902. # ("Any <model>") VP has no MQTT bridge to a real printer,
  903. # so the slicer has no live AMS layout to resolve tray IDs
  904. # against — whatever it sends here is meaningless (or,
  905. # worse, coincidentally valid for the wrong printer once
  906. # the scheduler later picks one).
  907. # 2. The per-VP `save_ams_mapping` opt-in must be on. Taking
  908. # the slicer's pick means `_ensure_ams_mapping` returns
  909. # early and `_compute_ams_mapping_for_printer` never runs —
  910. # and that function is where `prefer_lowest_filament`, its
  911. # AMS-filament-backup gate (#1766) and the inventory-remain
  912. # overrides live. Honouring the slicer unconditionally would
  913. # silently retire all of that for every existing queue-mode
  914. # VP on upgrade, so it's opt-in like every other queue-mode
  915. # behaviour toggle (#2700 review).
  916. #
  917. # Either gate failing leaves it unset, and the scheduler's
  918. # normal type/color re-derivation runs against whichever
  919. # printer actually gets the job.
  920. ams_mapping_json: str | None = None
  921. if slicer_opts is not None and self.target_printer_id is not None and self.save_ams_mapping:
  922. ams_mapping_json = _extract_slicer_ams_mapping_json(
  923. slicer_opts, f"[VP {self.name}]", is_dual_nozzle=self._target_is_dual_nozzle()
  924. )
  925. # `Force color match` is the user asking Bambuddy to do the
  926. # matching strictly, against the printer's live trays. Its only
  927. # effect on a fixed-printer item is via the per-slot
  928. # `filament_overrides` written below, which are consumed inside
  929. # `_compute_ams_mapping_for_printer` — the exact function a
  930. # stored mapping skips. So when both toggles are on, the
  931. # explicit strictness wins for *this* dispatch and the slicer's
  932. # pick is still persisted onto the archive for later reprints,
  933. # which is what `Save AMS mapping` actually promises (#2700
  934. # review).
  935. queue_ams_mapping_json = ams_mapping_json
  936. if queue_ams_mapping_json is not None and self.queue_force_color_match:
  937. logger.info(
  938. "[VP %s] Saved the slicer's AMS pick to the archive but not onto the queue item(s): "
  939. "'Force color match' is on, so the scheduler matches against live trays for this print.",
  940. self.name,
  941. )
  942. queue_ams_mapping_json = None
  943. # Parsed once for the per-plate length check in the loop below.
  944. queue_ams_mapping = json.loads(queue_ams_mapping_json) if queue_ams_mapping_json else None
  945. service = ArchiveService(db)
  946. archive = await service.archive_print(
  947. printer_id=None,
  948. source_file=file_path,
  949. print_data={
  950. "status": "archived",
  951. "source": "virtual_printer",
  952. "source_ip": source_ip,
  953. },
  954. prefer_filename_for_name=prefer_filename,
  955. # Slicer's own live AMS-slot pick -- promoted to
  956. # `extra_data.slicer_ams_mapping` by archive_print() so a
  957. # later reprint can reuse it. Already gated on the per-VP
  958. # `save_ams_mapping` opt-in above. Tagged with the printer
  959. # it was resolved against so a later reprint on a
  960. # *different* printer knows not to reuse it (#2700 review).
  961. slicer_ams_mapping=(json.loads(ams_mapping_json) if ams_mapping_json else None),
  962. slicer_ams_mapping_printer_id=self.target_printer_id,
  963. )
  964. if archive:
  965. logger.info("[VP %s] Archived: %s - %s", self.name, archive.id, archive.print_name)
  966. # Assign to specific printer if configured, otherwise use model for "Any X" scheduling
  967. target_model = None
  968. if not self.target_printer_id and self.model:
  969. target_model = VIRTUAL_PRINTER_MODELS.get(self.model)
  970. # #1733: multi-plate "Send All" uploads ship every plate in
  971. # one 3MF — `slice_info.config` lists each `<plate>` with
  972. # its own index. Enqueue one PrintQueueItem per plate so
  973. # the scheduler runs each separately. Single-plate "Send"
  974. # comes through as `[N]` (one plate index) so the loop
  975. # below runs once and the existing behaviour is preserved.
  976. plate_ids = self._extract_plate_ids(file_path)
  977. # Pick a base position the same way the manual /print-queue/
  978. # POST does, then hand consecutive positions to each plate
  979. # so a Send All keeps plate-order execution inside the
  980. # queue (#1733). Previously hardcoded to 1, which created
  981. # duplicate position=1 rows on every VP upload and made
  982. # queue execution order non-deterministic for any non-
  983. # empty queue.
  984. from sqlalchemy import func, select as _sql_select
  985. # One sequence across all pending items, not one per
  986. # printer (#3200): the plates go to the end of the queue.
  987. queue_scope = _sql_select(func.max(PrintQueueItem.position)).where(
  988. PrintQueueItem.status == "pending"
  989. )
  990. try:
  991. max_pos_raw = (await db.execute(queue_scope)).scalar()
  992. max_pos = int(max_pos_raw) if max_pos_raw is not None else 0
  993. except (TypeError, ValueError):
  994. max_pos = 0
  995. # Parse per-plate filament requirements (#1188). Each plate
  996. # has its own filament set in `slice_info.config`, so the
  997. # `required_filament_types` / `filament_overrides` columns
  998. # on each queue item reflect THAT plate, not the file's
  999. # first plate. Scoping was already plate-aware via #1697 —
  1000. # the `extract_filament_requirements(path, plate_id)` filter
  1001. # returns just the plate's filaments. required_filament_types
  1002. # is populated unconditionally — it's cheap, lets the
  1003. # scheduler reject obvious mis-matches even without
  1004. # force_color_match. filament_overrides only carries
  1005. # force_color_match=True when the per-VP setting is on, so
  1006. # upgraders keep the old behaviour by default.
  1007. queue_item_ids: list[int] = []
  1008. for offset, plate_id in enumerate(plate_ids, start=1):
  1009. required_filament_types_json: str | None = None
  1010. filament_overrides_json: str | None = None
  1011. requirements = extract_filament_requirements(file_path, plate_id)
  1012. if requirements:
  1013. types = sorted({r["type"] for r in requirements if r.get("type")})
  1014. if types:
  1015. required_filament_types_json = json.dumps(types)
  1016. if self.queue_force_color_match:
  1017. # Carry tray_info_idx so force_color_match can
  1018. # tell Bambu PLA variants apart (#2650). Bambu
  1019. # reports Basic/Matte/Silk all as tray_type
  1020. # "PLA"; the variant lives only in tray_info_idx
  1021. # (GFA00/GFA01/GFA06/...). A blank idx (custom or
  1022. # third-party spool) means "no variant
  1023. # constraint" and the scheduler falls back to
  1024. # type+colour.
  1025. overrides = [
  1026. {
  1027. "slot_id": r["slot_id"],
  1028. "type": r.get("type", ""),
  1029. "color": r.get("color", ""),
  1030. "tray_info_idx": r.get("tray_info_idx", ""),
  1031. "force_color_match": True,
  1032. }
  1033. for r in requirements
  1034. if r.get("type") and r.get("color")
  1035. ]
  1036. if overrides:
  1037. filament_overrides_json = json.dumps(overrides)
  1038. # The slicer's mapping is indexed by the 3MF's own
  1039. # file-global slot ids (position = slot_id - 1), so one
  1040. # array covers every plate of a multi-plate Send All —
  1041. # each plate just reads the entries for the slots it
  1042. # actually prints. What must be checked is that it
  1043. # reaches that far: a mapping shorter than this plate's
  1044. # highest slot id can't address the plate's own slots,
  1045. # and `_ensure_ams_mapping` would keep it anyway
  1046. # because it only rejects an all-unresolved mapping. Fall
  1047. # back to a computed mapping for that plate instead
  1048. # (#2700 review).
  1049. plate_ams_mapping_json = queue_ams_mapping_json
  1050. if queue_ams_mapping is not None and requirements:
  1051. max_slot_id = max((r.get("slot_id") or 0) for r in requirements)
  1052. if max_slot_id > len(queue_ams_mapping):
  1053. logger.warning(
  1054. "[VP %s] Slicer ams_mapping has %d entries but plate %s needs slot %d; "
  1055. "dropping it for this plate so the scheduler computes one from live AMS state.",
  1056. self.name,
  1057. len(queue_ams_mapping),
  1058. plate_id,
  1059. max_slot_id,
  1060. )
  1061. plate_ams_mapping_json = None
  1062. queue_item = PrintQueueItem(
  1063. printer_id=self.target_printer_id,
  1064. target_model=target_model,
  1065. archive_id=archive.id,
  1066. plate_id=plate_id,
  1067. position=max_pos + offset,
  1068. status="pending",
  1069. manual_start=not self.auto_dispatch,
  1070. required_filament_types=required_filament_types_json,
  1071. filament_overrides=filament_overrides_json,
  1072. bed_levelling=bed_levelling,
  1073. flow_cali=flow_cali,
  1074. vibration_cali=vibration_cali,
  1075. layer_inspect=layer_inspect,
  1076. timelapse=timelapse,
  1077. confirm_outcome=confirm_outcome,
  1078. # Per-VP opt-in for auto-print G-code injection (#1516).
  1079. # Default off; when on, the scheduler still no-ops unless
  1080. # gcode_snippets are configured for the target model, so it's
  1081. # effectively "inject when enabled AND snippets exist".
  1082. gcode_injection=self.gcode_injection,
  1083. # H2C rack-swap slicer pick (#1780). Captured above;
  1084. # stamped on every plate so a multi-plate Send All keeps
  1085. # the same nozzle pick across plates rather than only the
  1086. # first one (mirrors the #1697 / #1188 per-plate loop fix).
  1087. nozzle_mapping=nozzle_mapping_json,
  1088. # Slicer's own live AMS-slot pick, when present —
  1089. # see `_extract_slicer_ams_mapping_json`.
  1090. ams_mapping=plate_ams_mapping_json,
  1091. )
  1092. db.add(queue_item)
  1093. await db.flush() # populate queue_item.id before logging
  1094. queue_item_ids.append(queue_item.id)
  1095. await db.commit()
  1096. # Track the freshly-committed queue items so
  1097. # `on_print_command` can retroactively stamp slicer-side
  1098. # fields if the MQTT `project_file` lands AFTER the
  1099. # `_SLICER_OPTIONS_WAIT_TIMEOUT` window expired — the
  1100. # #1780 round-3 race. Eviction of stale entries here
  1101. # keeps the dict bounded; the queue path is the only
  1102. # writer, so doing it on commit is enough.
  1103. now = time.monotonic()
  1104. cutoff = now - _RECENT_QUEUE_ITEM_TTL
  1105. self._recent_queue_items = {k: v for k, v in self._recent_queue_items.items() if v[1] > cutoff}
  1106. self._recent_queue_items[file_path.name] = (list(queue_item_ids), now)
  1107. # Last-chance check: MQTT for this filename could have
  1108. # arrived during ANY await between the initial pop and
  1109. # now — wait_for itself, archive_print, db.flush,
  1110. # db.commit. In all those cases `on_print_command`
  1111. # stashed its data but neither the event-signal path nor
  1112. # the retroactive `_recent_queue_items` path was in
  1113. # place to consume it. Pop any late stash and apply
  1114. # inline so the late MQTT never leaks past the queue-add.
  1115. late_opts = self._slicer_print_options.pop(file_path.name, None)
  1116. if late_opts is not None:
  1117. logger.info(
  1118. "[VP %s] Late slicer MQTT detected for %s during queue-add — "
  1119. "applying inline (race vs commit/archive/flush yield)",
  1120. self.name,
  1121. file_path.name,
  1122. )
  1123. await self._restamp_recent_queue_item(file_path.name, late_opts)
  1124. if len(queue_item_ids) == 1:
  1125. logger.info("[VP %s] Added to queue: %s", self.name, queue_item_ids[0])
  1126. else:
  1127. logger.info(
  1128. "[VP %s] Added %d queue items for multi-plate upload (plates %s): %s",
  1129. self.name,
  1130. len(queue_item_ids),
  1131. plate_ids,
  1132. queue_item_ids,
  1133. )
  1134. await self._broadcast_archive_created(archive)
  1135. else:
  1136. logger.error("Failed to archive file: %s", file_path.name)
  1137. except Exception as e:
  1138. logger.error("Error adding to print queue: %s", e)
  1139. finally:
  1140. # Always release the marker and clean the temp file. Without this
  1141. # the same-name STOR guard would block the next upload and the
  1142. # upload_dir would accumulate failed temp files forever
  1143. # (#audit-R2-1).
  1144. self._pending_files.pop(file_path.name, None)
  1145. try:
  1146. file_path.unlink(missing_ok=True)
  1147. except OSError:
  1148. pass
  1149. async def _broadcast_archive_created(self, archive) -> None:
  1150. """Notify connected clients that a new archive exists.
  1151. Real-printer prints get this from main.py's MQTT print_start handler;
  1152. VP-uploaded prints need their own broadcast or the Archives page stays
  1153. stale until the user switches tabs (#1282).
  1154. """
  1155. try:
  1156. from backend.app.core.websocket import ws_manager
  1157. await ws_manager.send_archive_created(
  1158. {
  1159. "id": archive.id,
  1160. "printer_id": archive.printer_id,
  1161. "filename": archive.filename,
  1162. "print_name": archive.print_name,
  1163. "status": archive.status,
  1164. }
  1165. )
  1166. except Exception as e:
  1167. logger.debug("[VP %s] archive_created broadcast failed: %s", self.name, e)
  1168. @staticmethod
  1169. def _extract_plate_ids(file_path: Path) -> list[int]:
  1170. """Extract every plate index from a 3MF's slice_info.config.
  1171. A multi-plate "Send All" from BambuStudio / OrcaSlicer uploads a
  1172. single 3MF containing every plate the user selected. Each plate
  1173. has its own ``<plate>`` block with a ``<metadata key="index"
  1174. value="N"/>`` child and its own ``Metadata/plate_N.gcode`` payload
  1175. inside the same zip. Returning the full ordered list lets the VP
  1176. queue path create one queue item per plate (`_add_to_print_queue`
  1177. loops over the result), so "Send All" of a 3-plate file produces
  1178. 3 queue items sharing the same archive — one per plate to print.
  1179. Single-plate "Send" hits the same code path and returns ``[N]``
  1180. for whichever plate the user selected; the loop runs once and the
  1181. existing single-plate behaviour is preserved.
  1182. Returns ``[1]`` when the 3MF is missing ``slice_info.config``,
  1183. unparseable, or contains no plate-index metadata — the original
  1184. single-plate fallback. Production logs at debug so a non-3MF
  1185. upload doesn't spam, but the trail survives for support bundles.
  1186. """
  1187. try:
  1188. import xml.etree.ElementTree as ET
  1189. import zipfile
  1190. with zipfile.ZipFile(file_path, "r") as zf:
  1191. if "Metadata/slice_info.config" in zf.namelist():
  1192. content = zf.read("Metadata/slice_info.config").decode()
  1193. root = ET.fromstring(content) # noqa: S314 # nosec B314
  1194. plate_ids: list[int] = []
  1195. for plate in root.findall(".//plate"):
  1196. for meta in plate.findall("metadata"):
  1197. if meta.get("key") == "index" and meta.get("value"):
  1198. try:
  1199. plate_ids.append(int(meta.get("value")))
  1200. except ValueError:
  1201. continue
  1202. break
  1203. if plate_ids:
  1204. return plate_ids
  1205. except Exception as e:
  1206. logger.debug("[VP] _extract_plate_ids failed for %s: %s", file_path.name, e)
  1207. return [1]
  1208. # -- Service lifecycle --
  1209. def _resolve_cert_and_advertise(self) -> tuple[Path, Path, str]:
  1210. """Return (cert_path, key_path, advertise_address) for TLS services.
  1211. Always uses the self-signed cert chain (signed by `bbl_ca`). The user
  1212. imports `bbl_ca.crt` once into the slicer; per-VP certs validate from
  1213. there. Tailscale exposure is handled by the user picking the Tailscale
  1214. IP in the bind_ip dropdown.
  1215. """
  1216. cert_path, key_path = self.generate_certificates()
  1217. advertise = self.remote_interface_ip or self.bind_ip or ""
  1218. return cert_path, key_path, advertise
  1219. async def start_server(self) -> None:
  1220. """Start server-mode services (FTP, MQTT, SSDP, Bind) on this VP's bind_ip."""
  1221. logger.info("[VP %s] Starting server-mode services on %s", self.name, self.bind_ip)
  1222. cert_path, key_path, advertise_addr = self._resolve_cert_and_advertise()
  1223. bind_addr = self.bind_ip or "0.0.0.0" # nosec B104
  1224. async def run_with_logging(coro, svc_name):
  1225. try:
  1226. await coro
  1227. except Exception as e:
  1228. logger.error("[VP %s] %s failed: %s", self.name, svc_name, e)
  1229. self._tasks = []
  1230. # FTP server. Each VP gets a non-overlapping passive-mode port slice
  1231. # derived from its DB id so bridge-mode Docker users only have to
  1232. # expose a narrow range (#1646). Default slice is 10 ports per VP;
  1233. # see ftp_server.compute_passive_port_slice for the wrap-around
  1234. # behaviour on installs with very high VP ids.
  1235. passive_port_min, passive_port_max = compute_passive_port_slice(self.id)
  1236. self._ftp = VirtualPrinterFTPServer(
  1237. upload_dir=self.upload_dir,
  1238. access_code=self.access_code,
  1239. cert_path=cert_path,
  1240. key_path=key_path,
  1241. on_file_received=self.on_file_received,
  1242. bind_address=bind_addr,
  1243. vp_name=self.name,
  1244. passive_port_min=passive_port_min,
  1245. passive_port_max=passive_port_max,
  1246. )
  1247. self._tasks.append(
  1248. asyncio.create_task(
  1249. run_with_logging(self._ftp.start(), "FTP"),
  1250. name=f"vp_{self.id}_ftp",
  1251. )
  1252. )
  1253. # MQTT server
  1254. self._mqtt = SimpleMQTTServer(
  1255. serial=self.serial,
  1256. access_code=self.access_code,
  1257. cert_path=cert_path,
  1258. key_path=key_path,
  1259. on_print_command=self.on_print_command,
  1260. model=self.model or DEFAULT_VIRTUAL_PRINTER_MODEL,
  1261. bind_address=bind_addr,
  1262. vp_name=self.name,
  1263. )
  1264. self._tasks.append(
  1265. asyncio.create_task(
  1266. run_with_logging(self._mqtt.start(), "MQTT"),
  1267. name=f"vp_{self.id}_mqtt",
  1268. )
  1269. )
  1270. # MQTT bridge — fans out the target printer's pushes to slicers connected
  1271. # to this VP and forwards their commands back to the printer. Only meaningful
  1272. # when a target printer is configured AND printer_manager was injected (it
  1273. # always is at runtime; tests may omit it).
  1274. if self.target_printer_id is not None and self._printer_manager is not None:
  1275. self._mqtt_bridge = MQTTBridge(
  1276. vp_id=self.id,
  1277. vp_name=self.name,
  1278. vp_serial=self.serial,
  1279. target_printer_id=self.target_printer_id,
  1280. mqtt_server=self._mqtt,
  1281. printer_manager=self._printer_manager,
  1282. )
  1283. self._mqtt.set_bridge(self._mqtt_bridge)
  1284. await self._mqtt_bridge.start()
  1285. # Camera passthrough. BambuStudio / OrcaSlicer connect the "camera"
  1286. # button to the device IP they bound on (the VP), not the IP in the
  1287. # printer's `ipcam.rtsp_url`. Without a listener the slicer gets
  1288. # connection refused → "LAN connection failed" (RTSP models) or
  1289. # OrcaSlicer error `[2:-10061]` (chamber-image models, #1868).
  1290. #
  1291. # The port depends on the TARGET printer's model:
  1292. # RTSPS (X1/X2/H2/P2S) → 322
  1293. # chamber-image (A1/P1P/P1S) → 6000
  1294. #
  1295. # `get_camera_port()` is the same source of truth used by
  1296. # `routes/camera.py`, so slicer and Bambuddy UI agree.
  1297. target_client = self._printer_manager.get_client(self.target_printer_id)
  1298. target_ip = getattr(target_client, "ip_address", None) if target_client else None
  1299. target_model = getattr(target_client, "model", None) if target_client else None
  1300. if target_ip:
  1301. from backend.app.services.camera import get_camera_port
  1302. camera_port = get_camera_port(target_model)
  1303. self._rtsp_proxy = TCPProxy(
  1304. name=f"Camera-{camera_port}",
  1305. listen_port=camera_port,
  1306. target_host=target_ip,
  1307. target_port=camera_port,
  1308. bind_address=bind_addr,
  1309. )
  1310. self._tasks.append(
  1311. asyncio.create_task(
  1312. run_with_logging(self._rtsp_proxy.start(), f"Camera-{camera_port}"),
  1313. name=f"vp_{self.id}_camera",
  1314. )
  1315. )
  1316. # Bind server
  1317. self._bind = BindServer(
  1318. serial=self.serial,
  1319. model=self.model or DEFAULT_VIRTUAL_PRINTER_MODEL,
  1320. name=self.name,
  1321. bind_address=bind_addr,
  1322. cert_path=cert_path,
  1323. key_path=key_path,
  1324. )
  1325. self._tasks.append(
  1326. asyncio.create_task(
  1327. run_with_logging(self._bind.start(), "Bind"),
  1328. name=f"vp_{self.id}_bind",
  1329. )
  1330. )
  1331. # SSDP server — advertise_addr is the remote_interface_ip (Tailscale
  1332. # IP, when chosen from the bind_ip dropdown) or the bind_ip. SSDP
  1333. # Location accepts IPs only; FQDNs go in through bind_ip selection
  1334. # at the printer-IP level and resolve before reaching the SSDP
  1335. # advertisement.
  1336. self._ssdp = VirtualPrinterSSDPServer(
  1337. name=self.name,
  1338. serial=self.serial,
  1339. model=self.model or DEFAULT_VIRTUAL_PRINTER_MODEL,
  1340. advertise_ip=advertise_addr,
  1341. bind_ip=bind_addr,
  1342. )
  1343. self._tasks.append(
  1344. asyncio.create_task(
  1345. run_with_logging(self._ssdp.start(), "SSDP"),
  1346. name=f"vp_{self.id}_ssdp",
  1347. )
  1348. )
  1349. # Wait briefly for every child service to actually finish binding its
  1350. # socket so ``is_running`` doesn't lie. Without this barrier a caller
  1351. # racing the start (e.g. the diagnostic route) would see is_running=True
  1352. # while ports were still in the gap between task creation and the
  1353. # ``asyncio.start_server`` returning. Bounded timeout — if a child
  1354. # hangs we log it and move on; the existing task tracking still
  1355. # catches the failure on the next iteration.
  1356. ready_targets = [
  1357. ("FTP", self._ftp.ready),
  1358. ("MQTT", self._mqtt.ready),
  1359. ("Bind", self._bind.ready),
  1360. ("SSDP", self._ssdp.ready),
  1361. ]
  1362. try:
  1363. await asyncio.wait_for(
  1364. asyncio.gather(*(e.wait() for _, e in ready_targets)),
  1365. timeout=5.0,
  1366. )
  1367. except TimeoutError:
  1368. not_ready = [name for name, e in ready_targets if not e.is_set()]
  1369. logger.warning(
  1370. "[VP %s] Sub-service(s) didn't bind within 5s: %s — continuing anyway",
  1371. self.name,
  1372. ", ".join(not_ready) or "(none)",
  1373. )
  1374. logger.info("[VP %s] Server-mode services started on %s", self.name, bind_addr)
  1375. async def stop_server(self) -> None:
  1376. """Stop server-mode services."""
  1377. if self._finish_release_task is not None and not self._finish_release_task.done():
  1378. self._finish_release_task.cancel()
  1379. self._finish_release_task = None
  1380. if self._mqtt_bridge:
  1381. try:
  1382. await self._mqtt_bridge.stop()
  1383. except Exception:
  1384. logger.exception("[VP %s] MQTT bridge stop failed", self.name)
  1385. if self._mqtt:
  1386. self._mqtt.set_bridge(None)
  1387. self._mqtt_bridge = None
  1388. if self._rtsp_proxy:
  1389. try:
  1390. await self._rtsp_proxy.stop()
  1391. except Exception:
  1392. logger.exception("[VP %s] Camera proxy stop failed", self.name)
  1393. self._rtsp_proxy = None
  1394. if self._ftp:
  1395. await self._ftp.stop()
  1396. self._ftp = None
  1397. if self._mqtt:
  1398. await self._mqtt.stop()
  1399. self._mqtt = None
  1400. if self._bind:
  1401. await self._bind.stop()
  1402. self._bind = None
  1403. if self._ssdp:
  1404. await self._ssdp.stop()
  1405. self._ssdp = None
  1406. await self._cancel_tasks()
  1407. async def start_proxy(self) -> None:
  1408. """Start proxy mode services for this instance."""
  1409. logger.info("[VP %s] Starting proxy mode to %s", self.name, self.target_printer_ip)
  1410. cert_path, key_path, _ = self._resolve_cert_and_advertise()
  1411. self._proxy = SlicerProxyManager(
  1412. target_host=self.target_printer_ip,
  1413. cert_path=cert_path,
  1414. key_path=key_path,
  1415. on_activity=lambda n, m: logger.info("[VP %s] Proxy %s: %s", self.name, n, m),
  1416. bind_address=self.bind_ip or "0.0.0.0", # nosec B104
  1417. bind_identity={
  1418. "serial": self.target_printer_serial or self.serial,
  1419. "model": self.model or DEFAULT_VIRTUAL_PRINTER_MODEL,
  1420. "name": self.name,
  1421. "version": "01.00.00.00",
  1422. },
  1423. )
  1424. async def run_with_logging(coro, svc_name):
  1425. try:
  1426. await coro
  1427. except Exception as e:
  1428. logger.error("[VP %s] %s failed: %s", self.name, svc_name, e)
  1429. self._tasks = []
  1430. # SSDP for proxy
  1431. proxy_serial = self.target_printer_serial or self.serial
  1432. if self.remote_interface_ip:
  1433. from backend.app.services.network_utils import find_interface_for_ip
  1434. local_iface = find_interface_for_ip(self.target_printer_ip)
  1435. if local_iface:
  1436. self._ssdp_proxy = SSDPProxy(
  1437. local_interface_ip=local_iface["ip"],
  1438. remote_interface_ip=self.remote_interface_ip,
  1439. target_printer_ip=self.target_printer_ip,
  1440. name=self.name,
  1441. )
  1442. self._tasks.append(
  1443. asyncio.create_task(
  1444. run_with_logging(self._ssdp_proxy.start(), "SSDP Proxy"),
  1445. name=f"vp_{self.id}_ssdp_proxy",
  1446. )
  1447. )
  1448. else:
  1449. self._start_fallback_ssdp(proxy_serial, run_with_logging)
  1450. else:
  1451. self._start_fallback_ssdp(proxy_serial, run_with_logging)
  1452. self._tasks.append(
  1453. asyncio.create_task(
  1454. run_with_logging(self._proxy.start(), "Proxy"),
  1455. name=f"vp_{self.id}_proxy",
  1456. )
  1457. )
  1458. def _start_fallback_ssdp(self, proxy_serial: str, run_with_logging) -> None:
  1459. """Start single-interface SSDP server as fallback for proxy mode."""
  1460. self._ssdp = VirtualPrinterSSDPServer(
  1461. name=f"{self.name} (Proxy)",
  1462. serial=proxy_serial,
  1463. model=self.model or DEFAULT_VIRTUAL_PRINTER_MODEL,
  1464. advertise_ip=self.bind_ip or "",
  1465. bind_ip=self.bind_ip or "",
  1466. )
  1467. self._tasks.append(
  1468. asyncio.create_task(
  1469. run_with_logging(self._ssdp.start(), "SSDP"),
  1470. name=f"vp_{self.id}_ssdp",
  1471. )
  1472. )
  1473. async def stop_proxy(self) -> None:
  1474. """Stop proxy mode services for this instance."""
  1475. if self._proxy:
  1476. await self._proxy.stop()
  1477. self._proxy = None
  1478. if self._ssdp:
  1479. await self._ssdp.stop()
  1480. self._ssdp = None
  1481. if self._ssdp_proxy:
  1482. await self._ssdp_proxy.stop()
  1483. self._ssdp_proxy = None
  1484. await self._cancel_tasks()
  1485. async def _cancel_tasks(self) -> None:
  1486. """Cancel all running tasks and wait for cleanup."""
  1487. for task in self._tasks:
  1488. task.cancel()
  1489. if self._tasks:
  1490. try:
  1491. await asyncio.wait_for(asyncio.gather(*self._tasks, return_exceptions=True), timeout=1.0)
  1492. except TimeoutError:
  1493. pass
  1494. self._tasks = []
  1495. def get_status(self) -> dict:
  1496. """Get status for this instance."""
  1497. status: dict = {
  1498. "running": self.is_running,
  1499. "pending_files": len(self._pending_files),
  1500. }
  1501. if self.is_proxy and self._proxy:
  1502. status["proxy"] = self._proxy.get_status()
  1503. return status
  1504. class VirtualPrinterManager:
  1505. """Multi-instance virtual printer registry and orchestrator.
  1506. Every VP runs its own independent services on a dedicated bind IP.
  1507. """
  1508. def __init__(self):
  1509. self._session_factory: Callable | None = None
  1510. self._printer_manager: PrinterManager | None = None
  1511. self._instances: dict[int, VirtualPrinterInstance] = {}
  1512. # Serialize sync_from_db so concurrent PUT /vp/{id} calls can't
  1513. # race the start/stop sequence and leave duplicate sub-services
  1514. # bound to the same port. The lock is fine-grained enough that
  1515. # a single VP update completes in well under a second; if the
  1516. # user holds the lock with a long-running start they intended
  1517. # to anyway.
  1518. self._sync_lock = asyncio.Lock()
  1519. # Directories
  1520. self._base_dir = app_settings.base_dir / "virtual_printer"
  1521. # Ensure base directories exist
  1522. self._ensure_base_directories()
  1523. def _ensure_base_directories(self) -> None:
  1524. """Create base directories at startup."""
  1525. for dir_path in [self._base_dir, self._base_dir / "uploads", self._base_dir / "certs"]:
  1526. try:
  1527. dir_path.mkdir(parents=True, exist_ok=True)
  1528. except PermissionError:
  1529. logger.error(
  1530. f"Cannot create directory {dir_path}: Permission denied. "
  1531. f"For Docker: ensure the data volume is writable by the container user. "
  1532. f"For bare metal: run 'sudo chown -R $(whoami) {self._base_dir}'"
  1533. )
  1534. def set_session_factory(self, session_factory: Callable) -> None:
  1535. """Set the database session factory."""
  1536. self._session_factory = session_factory
  1537. def set_printer_manager(self, printer_manager: "PrinterManager") -> None:
  1538. """Inject the global printer_manager so non-proxy VPs can mirror their target's MQTT stream."""
  1539. self._printer_manager = printer_manager
  1540. def get_ca_certificate_info(self) -> dict:
  1541. """Return the shared virtual-printer CA certificate for slicer-trust import.
  1542. The CA is shared by every VP (one import covers all of them). It is
  1543. generated on demand here if no VP has triggered cert generation yet,
  1544. so the "copy/download certificate" UI works even before the first VP
  1545. is enabled.
  1546. """
  1547. certs_dir = self._base_dir / "certs"
  1548. cert_service = CertificateService(cert_dir=certs_dir, shared_ca_dir=certs_dir)
  1549. return cert_service.get_ca_certificate_info()
  1550. @property
  1551. def is_enabled(self) -> bool:
  1552. """Check if any virtual printer is running."""
  1553. return len(self._instances) > 0
  1554. async def sync_from_db(self) -> None:
  1555. """Load all VPs from DB, reconcile running state.
  1556. Serialised by ``self._sync_lock`` — concurrent PUT /vp/{id} routes
  1557. all call into this method; without the lock the start / stop
  1558. sequence races and can leave duplicate sub-services bound to the
  1559. same port or orphan still-running tasks.
  1560. """
  1561. if not self._session_factory:
  1562. logger.warning("Cannot sync virtual printers: no session factory")
  1563. return
  1564. async with self._sync_lock:
  1565. await self._sync_from_db_locked()
  1566. async def _sync_from_db_locked(self) -> None:
  1567. """Inner sync body — caller holds ``self._sync_lock``."""
  1568. from sqlalchemy import select
  1569. from backend.app.models.printer import Printer
  1570. from backend.app.models.virtual_printer import VirtualPrinter
  1571. async with self._session_factory() as db:
  1572. result = await db.execute(
  1573. select(VirtualPrinter).where(VirtualPrinter.enabled == True).order_by(VirtualPrinter.position) # noqa: E712
  1574. )
  1575. enabled_vps = result.scalars().all()
  1576. # Stop instances that are no longer enabled or changed mode
  1577. enabled_ids = {vp.id for vp in enabled_vps}
  1578. for vp_id in list(self._instances.keys()):
  1579. if vp_id not in enabled_ids:
  1580. await self.remove_instance(vp_id)
  1581. # Look up printer IPs for proxy VPs
  1582. proxy_vps = [vp for vp in enabled_vps if vp.mode == "proxy"]
  1583. proxy_ips: dict[int, tuple[str, str]] = {}
  1584. if proxy_vps:
  1585. async with self._session_factory() as db:
  1586. for pvp in proxy_vps:
  1587. if pvp.target_printer_id:
  1588. result = await db.execute(select(Printer).where(Printer.id == pvp.target_printer_id))
  1589. printer = result.scalar_one_or_none()
  1590. if printer:
  1591. proxy_ips[pvp.id] = (printer.ip_address, printer.serial_number)
  1592. # Detect config changes on running instances and restart if needed
  1593. for vp in enabled_vps:
  1594. instance = self._instances.get(vp.id)
  1595. if not instance:
  1596. continue
  1597. # Proxy mode: detect target printer IP / serial changes from the
  1598. # DB lookup above. Without this branch a DHCP renewal that gives
  1599. # the target printer a new IP would leave the running proxy
  1600. # forwarding to the stale IP until the user manually toggles the
  1601. # VP. The same shape covers a target-side serial change.
  1602. proxy_target_changed = False
  1603. if vp.mode == "proxy":
  1604. fresh = proxy_ips.get(vp.id)
  1605. if fresh is not None:
  1606. fresh_ip, fresh_serial = fresh
  1607. if (
  1608. getattr(instance, "target_printer_ip", None) != fresh_ip
  1609. or getattr(instance, "target_printer_serial", None) != fresh_serial
  1610. ):
  1611. proxy_target_changed = True
  1612. # Normalize the DB value before comparing — a legacy `immediate`
  1613. # row read before the migration window finishes would otherwise
  1614. # trip the "changed" branch and bounce every VP at boot.
  1615. db_mode = normalize_vp_mode(vp.mode)
  1616. changed = (
  1617. instance.mode != db_mode
  1618. or instance.model != (vp.model or DEFAULT_VIRTUAL_PRINTER_MODEL)
  1619. or instance.access_code != (vp.access_code or "")
  1620. or instance.bind_ip != (vp.bind_ip or "")
  1621. or instance.remote_interface_ip != (vp.remote_interface_ip or "")
  1622. or instance.target_printer_id != vp.target_printer_id
  1623. or instance.auto_dispatch != vp.auto_dispatch
  1624. # Queue-mode behaviour toggle — without it the running
  1625. # instance silently keeps the old value until process
  1626. # restart (#1552 follow-up family).
  1627. or instance.queue_force_color_match != vp.queue_force_color_match
  1628. or instance.save_ams_mapping != vp.save_ams_mapping
  1629. or instance.gcode_injection != vp.gcode_injection
  1630. or proxy_target_changed
  1631. )
  1632. if changed:
  1633. logger.info(
  1634. "VP %s config changed (mode: %s→%s), restarting",
  1635. instance.name,
  1636. instance.mode,
  1637. vp.mode,
  1638. )
  1639. await self.remove_instance(vp.id)
  1640. # Start instances for all enabled VPs (skip already running)
  1641. for vp in enabled_vps:
  1642. if vp.id in self._instances:
  1643. continue
  1644. if vp.mode == "proxy":
  1645. ip_info = proxy_ips.get(vp.id)
  1646. if not ip_info:
  1647. logger.warning("Proxy VP %s: target printer not found, skipping", vp.name)
  1648. continue
  1649. target_ip, target_serial = ip_info
  1650. instance = VirtualPrinterInstance(
  1651. vp_id=vp.id,
  1652. name=vp.name,
  1653. mode=vp.mode,
  1654. model=vp.model or DEFAULT_VIRTUAL_PRINTER_MODEL,
  1655. access_code=vp.access_code or "",
  1656. serial_suffix=vp.serial_suffix,
  1657. target_printer_ip=target_ip,
  1658. target_printer_serial=target_serial,
  1659. auto_dispatch=vp.auto_dispatch,
  1660. bind_ip=vp.bind_ip or "",
  1661. remote_interface_ip=vp.remote_interface_ip or "",
  1662. tailscale_disabled=vp.tailscale_disabled,
  1663. base_dir=self._base_dir,
  1664. session_factory=self._session_factory,
  1665. )
  1666. self._instances[vp.id] = instance
  1667. await instance.start_proxy()
  1668. logger.info("Started proxy VP: %s → %s (bind=%s)", instance.name, target_ip, instance.bind_ip)
  1669. else:
  1670. instance = VirtualPrinterInstance(
  1671. vp_id=vp.id,
  1672. name=vp.name,
  1673. mode=vp.mode,
  1674. model=vp.model or DEFAULT_VIRTUAL_PRINTER_MODEL,
  1675. access_code=vp.access_code or "",
  1676. serial_suffix=vp.serial_suffix,
  1677. target_printer_id=vp.target_printer_id,
  1678. auto_dispatch=vp.auto_dispatch,
  1679. queue_force_color_match=vp.queue_force_color_match,
  1680. save_ams_mapping=vp.save_ams_mapping,
  1681. gcode_injection=vp.gcode_injection,
  1682. bind_ip=vp.bind_ip or "",
  1683. remote_interface_ip=vp.remote_interface_ip or "",
  1684. tailscale_disabled=vp.tailscale_disabled,
  1685. base_dir=self._base_dir,
  1686. session_factory=self._session_factory,
  1687. printer_manager=self._printer_manager,
  1688. )
  1689. self._instances[vp.id] = instance
  1690. await instance.start_server()
  1691. logger.info("Started server-mode VP: %s on %s", instance.name, vp.bind_ip)
  1692. async def remove_instance(self, vp_id: int) -> None:
  1693. """Stop and remove a single VP instance."""
  1694. instance = self._instances.pop(vp_id, None)
  1695. if instance:
  1696. if instance.is_proxy:
  1697. await instance.stop_proxy()
  1698. else:
  1699. await instance.stop_server()
  1700. logger.info("Removed VP instance: %s", instance.name)
  1701. async def stop_all(self) -> None:
  1702. """Shutdown all virtual printer services."""
  1703. logger.info("Stopping all virtual printer services...")
  1704. for vp_id in list(self._instances.keys()):
  1705. await self.remove_instance(vp_id)
  1706. logger.info("All virtual printer services stopped")
  1707. def get_instance(self, vp_id: int) -> VirtualPrinterInstance | None:
  1708. """Get a running instance by ID."""
  1709. return self._instances.get(vp_id)
  1710. def get_all_status(self) -> list[dict]:
  1711. """Get status for all running instances."""
  1712. return [
  1713. {
  1714. "id": inst.id,
  1715. "name": inst.name,
  1716. "mode": inst.mode,
  1717. **inst.get_status(),
  1718. }
  1719. for inst in self._instances.values()
  1720. ]
  1721. # -- Legacy single-printer compat --
  1722. def get_status(self) -> dict:
  1723. """Get status for first virtual printer (backward compat)."""
  1724. if self._instances:
  1725. first = next(iter(self._instances.values()))
  1726. return {
  1727. "enabled": True,
  1728. "running": first.is_running,
  1729. "mode": first.mode,
  1730. "name": first.name,
  1731. "serial": first.serial,
  1732. "model": first.model or DEFAULT_VIRTUAL_PRINTER_MODEL,
  1733. "model_name": VIRTUAL_PRINTER_MODELS.get(
  1734. first.model or DEFAULT_VIRTUAL_PRINTER_MODEL,
  1735. first.model or DEFAULT_VIRTUAL_PRINTER_MODEL,
  1736. ),
  1737. "pending_files": first.get_status().get("pending_files", 0),
  1738. **({"target_printer_ip": first.target_printer_ip} if first.is_proxy else {}),
  1739. **({"proxy": first.get_status().get("proxy", {})} if first.is_proxy else {}),
  1740. }
  1741. return {
  1742. "enabled": False,
  1743. "running": False,
  1744. "mode": VP_MODE_ARCHIVE,
  1745. "name": "Bambuddy",
  1746. "serial": "",
  1747. "model": DEFAULT_VIRTUAL_PRINTER_MODEL,
  1748. "model_name": VIRTUAL_PRINTER_MODELS[DEFAULT_VIRTUAL_PRINTER_MODEL],
  1749. "pending_files": 0,
  1750. }
  1751. async def configure(
  1752. self,
  1753. enabled: bool,
  1754. access_code: str = "",
  1755. mode: str = VP_MODE_ARCHIVE,
  1756. model: str = "",
  1757. target_printer_ip: str = "",
  1758. target_printer_serial: str = "",
  1759. remote_interface_ip: str = "",
  1760. ) -> None:
  1761. """Legacy single-printer configure. Delegates to sync_from_db()."""
  1762. # This method is kept for backward compat with the settings endpoint.
  1763. # The actual work is done by sync_from_db() which reads from the DB.
  1764. await self.sync_from_db()
  1765. # Global instance
  1766. virtual_printer_manager = VirtualPrinterManager()