| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707270827092710271127122713271427152716271727182719272027212722272327242725272627272728272927302731273227332734273527362737273827392740274127422743274427452746274727482749275027512752275327542755275627572758275927602761276227632764276527662767276827692770277127722773277427752776277727782779278027812782278327842785278627872788278927902791279227932794279527962797279827992800280128022803280428052806280728082809281028112812281328142815281628172818281928202821282228232824282528262827282828292830283128322833283428352836283728382839284028412842284328442845284628472848284928502851285228532854285528562857285828592860286128622863286428652866286728682869287028712872287328742875287628772878287928802881288228832884288528862887288828892890289128922893289428952896289728982899290029012902290329042905290629072908290929102911291229132914291529162917291829192920292129222923292429252926292729282929293029312932293329342935293629372938293929402941294229432944294529462947294829492950295129522953295429552956295729582959296029612962296329642965296629672968296929702971297229732974297529762977297829792980298129822983298429852986298729882989299029912992299329942995299629972998299930003001300230033004300530063007300830093010301130123013301430153016301730183019302030213022302330243025302630273028302930303031303230333034303530363037303830393040304130423043304430453046304730483049305030513052305330543055305630573058305930603061306230633064306530663067306830693070307130723073307430753076307730783079308030813082308330843085308630873088308930903091309230933094309530963097309830993100310131023103310431053106310731083109311031113112311331143115311631173118311931203121312231233124312531263127312831293130313131323133313431353136313731383139314031413142314331443145314631473148314931503151315231533154315531563157315831593160316131623163316431653166316731683169317031713172317331743175317631773178317931803181318231833184318531863187318831893190319131923193319431953196319731983199320032013202320332043205320632073208320932103211321232133214321532163217321832193220322132223223322432253226322732283229323032313232323332343235323632373238323932403241324232433244324532463247324832493250325132523253325432553256325732583259326032613262326332643265326632673268326932703271327232733274327532763277327832793280328132823283328432853286328732883289329032913292329332943295329632973298329933003301330233033304330533063307330833093310331133123313331433153316331733183319332033213322332333243325332633273328332933303331333233333334333533363337333833393340334133423343334433453346334733483349335033513352335333543355335633573358335933603361336233633364336533663367336833693370337133723373337433753376337733783379338033813382338333843385338633873388338933903391339233933394339533963397339833993400340134023403340434053406340734083409341034113412341334143415341634173418341934203421342234233424342534263427342834293430343134323433343434353436343734383439344034413442344334443445344634473448344934503451345234533454345534563457345834593460346134623463346434653466346734683469347034713472347334743475347634773478347934803481348234833484348534863487348834893490349134923493349434953496349734983499350035013502350335043505350635073508350935103511351235133514351535163517351835193520352135223523352435253526352735283529353035313532353335343535353635373538353935403541354235433544354535463547354835493550355135523553355435553556355735583559356035613562356335643565356635673568356935703571357235733574357535763577357835793580358135823583358435853586358735883589359035913592359335943595359635973598359936003601360236033604360536063607360836093610361136123613361436153616361736183619362036213622362336243625362636273628362936303631363236333634363536363637363836393640364136423643364436453646364736483649365036513652365336543655365636573658365936603661366236633664366536663667366836693670367136723673367436753676367736783679368036813682368336843685368636873688368936903691369236933694369536963697369836993700370137023703370437053706370737083709371037113712371337143715371637173718371937203721372237233724372537263727372837293730373137323733373437353736373737383739374037413742374337443745374637473748374937503751375237533754375537563757375837593760376137623763376437653766376737683769377037713772377337743775377637773778377937803781378237833784378537863787378837893790379137923793379437953796379737983799380038013802380338043805380638073808380938103811381238133814381538163817381838193820382138223823382438253826382738283829383038313832383338343835383638373838383938403841384238433844384538463847384838493850385138523853385438553856385738583859386038613862386338643865386638673868386938703871387238733874387538763877387838793880388138823883388438853886388738883889389038913892389338943895389638973898389939003901390239033904390539063907390839093910391139123913391439153916391739183919392039213922392339243925392639273928392939303931393239333934393539363937393839393940394139423943394439453946394739483949395039513952395339543955395639573958395939603961396239633964396539663967396839693970397139723973397439753976397739783979398039813982398339843985398639873988398939903991399239933994399539963997399839994000400140024003400440054006400740084009401040114012401340144015401640174018401940204021402240234024402540264027402840294030403140324033403440354036403740384039404040414042404340444045404640474048404940504051405240534054405540564057405840594060406140624063406440654066406740684069407040714072407340744075407640774078407940804081408240834084408540864087408840894090409140924093409440954096409740984099410041014102410341044105410641074108410941104111411241134114411541164117411841194120412141224123412441254126412741284129413041314132 |
- """Print scheduler service - processes the print queue."""
- import asyncio
- import json
- import logging
- import time
- from datetime import datetime, timezone
- from pathlib import Path
- from sqlalchemy import func, select, update
- from sqlalchemy.ext.asyncio import AsyncSession
- from sqlalchemy.orm import selectinload
- from backend.app.core.config import settings
- from backend.app.core.database import async_session, run_with_retry
- from backend.app.core.tasks import spawn_background_task
- from backend.app.core.websocket import ws_manager
- from backend.app.models.archive import PrintArchive
- from backend.app.models.library import LibraryFile
- from backend.app.models.print_queue import PrintQueueItem
- from backend.app.models.printer import Printer
- from backend.app.models.settings import Settings
- from backend.app.models.smart_plug import SmartPlug
- from backend.app.models.spool_assignment import SpoolAssignment
- from backend.app.models.spoolman_slot_assignment import SpoolmanSlotAssignment
- from backend.app.services import print_dispatch_context
- from backend.app.services.bambu_ftp import (
- UploadCancelled,
- cache_3mf_download,
- delete_file_async,
- get_ftp_retry_settings,
- upload_file_async,
- with_ftp_retry,
- )
- from backend.app.services.bambu_mqtt import HMS_MQTT_VERIFY_FAILED
- from backend.app.services.filament_deficit import compute_deficit_for_queue_item
- from backend.app.services.notification_service import notification_service
- from backend.app.services.printer_manager import (
- printer_manager,
- supports_airduct,
- supports_chamber_heater,
- supports_chamber_temp,
- supports_drying,
- supports_drying_while_printing,
- )
- from backend.app.services.smart_plug_manager import smart_plug_manager
- from backend.app.utils.filename import derive_remote_filename
- from backend.app.utils.printer_models import is_gcode_compatible, normalize_printer_model
- logger = logging.getLogger(__name__)
- # Dispatch-toast progress throttling (#1625 follow-up). Mirrors the legacy
- # background_dispatch.py upload_progress_callback (200 ms time gate + 256 KB
- # byte gate) from before the scheduler unification. Time gate keeps small
- # files from going silent (a single 8 KB chunk fires once and that's it);
- # byte gate caps the broadcast rate on slow LAN where 200 ms covers many
- # chunks. uploaded >= total always emits so the bar closes cleanly even on
- # sub-200 ms files.
- _DISPATCH_PROGRESS_BYTE_STEP = 256 * 1024
- _DISPATCH_PROGRESS_MIN_INTERVAL_SECS = 0.2
- class _UploadProgressBridge:
- """Thread-safe bridge from ``upload_file_async`` to the WS broadcaster.
- ``upload_file_async`` runs the FTP transfer in an executor thread and
- invokes its ``progress_callback`` from that thread, so the callback
- body cannot ``await`` directly. This bridge captures the asyncio loop
- at construction (on the scheduler thread) and uses
- ``run_coroutine_threadsafe`` to hop back. The byte/time throttle
- matches the legacy background_dispatch.py path 1:1 so the toast feels
- identical to the pre-#1625 experience.
- Failures inside the emit are swallowed — progress is a UX nicety, the
- upload itself must not fail because of a WS hiccup.
- """
- def __init__(self, user_id: int | None, queue_item_id: int):
- self._user_id = user_id
- self._queue_item_id = queue_item_id
- try:
- self._loop = asyncio.get_running_loop()
- except RuntimeError:
- self._loop = None
- self._last_emit_bytes = 0
- self._last_emit_monotonic = 0.0
- self._has_emitted = False
- def __call__(self, bytes_transferred: int, total_bytes: int) -> None:
- if self._loop is None or total_bytes <= 0:
- return
- now = time.monotonic()
- # Mirrors legacy bg-dispatch: emit if first call OR upload complete
- # OR 200 ms elapsed OR ≥256 KB transferred since last emit. Two of
- # the four matter most: first-call so the user sees something even
- # for sub-chunk-size files; uploaded >= total so the bar locks at
- # 100% even when the throttle would otherwise eat it.
- should_emit = (
- not self._has_emitted
- or bytes_transferred >= total_bytes
- or now - self._last_emit_monotonic >= _DISPATCH_PROGRESS_MIN_INTERVAL_SECS
- or bytes_transferred - self._last_emit_bytes >= _DISPATCH_PROGRESS_BYTE_STEP
- )
- if not should_emit:
- return
- self._has_emitted = True
- self._last_emit_bytes = bytes_transferred
- self._last_emit_monotonic = now
- try:
- asyncio.run_coroutine_threadsafe(
- ws_manager.send_queue_item_upload_progress(
- user_id=self._user_id,
- queue_item_id=self._queue_item_id,
- bytes_transferred=bytes_transferred,
- total_bytes=total_bytes,
- ),
- self._loop,
- )
- except Exception:
- pass # progress is best-effort, never block the upload
- # Bambu firmware states that mean the project_file has actually been accepted
- # and the printer is now processing / running / paused mid-print. Used by the
- # dispatch watchdog (#1370): a transition into one of these states means the
- # print landed, anything else (e.g. FINISH -> IDLE after the user dismisses
- # a post-print prompt) is NOT a valid "command landed" signal even though the
- # state value did change. SLICING is included because some firmwares park
- # briefly in SLICING between PREPARE and RUNNING while parsing the g-code.
- _ACTIVE_PRINT_STATES: frozenset[str] = frozenset({"PREPARE", "SLICING", "RUNNING", "PAUSE"})
- # How many times the start-watchdog may revert an item to 'pending' before it
- # gives up and fails the row instead (#2555). Each attempt costs a full 3MF
- # re-upload plus the watchdog's wait, so a wedged printer left to retry forever
- # both never recovers and starves the other printers of dispatch slots. Three
- # is chosen to clear the transient causes the watchdog already recovers from —
- # a lost MQTT publish on a half-broken session (#887/#936) is fixed by the
- # force-reconnect on the very next attempt — while still bounding the loop.
- DISPATCH_MAX_ATTEMPTS = 3
- # Filament type equivalence groups — types within the same group are
- # interchangeable on the printer side (Bambu Lab firmware treats them as compatible).
- _FILAMENT_TYPE_GROUPS: list[list[str]] = [
- ["PA-CF", "PA12-CF", "PAHT-CF"],
- ]
- _FILAMENT_EQUIV_MAP: dict[str, str] = {}
- for _group in _FILAMENT_TYPE_GROUPS:
- _canonical = _group[0].upper()
- for _t in _group:
- _FILAMENT_EQUIV_MAP[_t.upper()] = _canonical
- def _canonical_filament_type(ftype: str) -> str:
- """Return canonical type for equivalence matching."""
- upper = ftype.upper()
- return _FILAMENT_EQUIV_MAP.get(upper, upper)
- def _mapping_is_all_unresolved(mapping: list | None) -> bool:
- """True if ``mapping`` is a non-empty list whose every entry is the
- unresolved sentinel (-1 / None) — i.e. no required slot ever matched a tray.
- Such a mapping is a bug artifact: a frontend status-load race can serialize
- ``[-1]`` before the printer's AMS trays are known (#2589). It must be
- recomputed from live status at dispatch rather than trusted, otherwise it
- reaches the print command and is silently downgraded to external-spool mode.
- A partially-resolved mapping (``[-1, -1, 5]`` where slot 3 matched, or a
- padding ``-1`` for a slot this plate does not print) is NOT unresolved. An
- explicit external selection (``>= 254``) is NOT unresolved either — those
- keep their meaning.
- """
- if not isinstance(mapping, list) or not mapping:
- return False
- return all(t is None or (isinstance(t, int) and t < 0) for t in mapping)
- def _mqtt_commands_rejected(status) -> bool:
- """True when the printer is currently reporting that it refused a command.
- ``HMS_MQTT_VERIFY_FAILED`` means the firmware's authorization check rejected
- a control command it could not verify. Queries still answer, so the printer
- looks connected and idle while project_file, gcode_line and
- ams_change_filament are all dropped — no amount of waiting or re-uploading
- changes that (#2732).
- Tolerates a missing status and errors without a ``full_code`` (the 8-char
- ``print_error`` path builds HMSError differently), so this is safe to call on
- every watchdog poll.
- """
- for err in getattr(status, "hms_errors", None) or []:
- if getattr(err, "full_code", "") == HMS_MQTT_VERIFY_FAILED:
- return True
- return False
- def _installed_nozzle_diameters(status) -> list[float]:
- """Parse the installed nozzle diameters from a PrinterState (#1899).
- Returns the diameters the printer actually reports (e.g. [0.4] single-nozzle,
- [0.4, 0.6] dual-nozzle), skipping the empty-string defaults that populate a
- NozzleInfo before MQTT fills it in. An empty list means "the printer hasn't
- told us its nozzle hardware" — callers must treat that as unknown, not as a
- mismatch, so we never block a print on missing data.
- """
- diameters: list[float] = []
- for nozzle in getattr(status, "nozzles", None) or []:
- raw = getattr(nozzle, "nozzle_diameter", "") or ""
- try:
- value = float(raw)
- except (TypeError, ValueError):
- continue
- if value > 0:
- diameters.append(value)
- return diameters
- def _nozzle_mismatch_message(sliced_nozzle: float | None, installed: list[float]) -> str | None:
- """Return an actionable error message when the sliced nozzle can't be
- printed on any installed nozzle, else None (#1899).
- Fail-safe: returns None whenever we lack the data to judge — no sliced
- diameter, or the printer reported no nozzles — so a print is only ever
- blocked on a POSITIVE mismatch. On dual-nozzle printers a match against
- EITHER installed nozzle passes (a 0.6 slice is fine if one hotend is 0.6).
- The 0.05 tolerance absorbs float noise while staying well inside the 0.2
- gap between adjacent nozzle sizes (0.2/0.4/0.6/0.8).
- """
- if not sliced_nozzle or not installed:
- return None
- if any(abs(d - sliced_nozzle) < 0.05 for d in installed):
- return None
- installed_str = " / ".join(f"{d:g}mm" for d in installed)
- return (
- f"File sliced for a {sliced_nozzle:g}mm nozzle, but the printer has "
- f"{installed_str} installed. Re-slice for the installed nozzle, or "
- f"install the matching nozzle before printing."
- )
- class PrintScheduler:
- """Background scheduler that processes the print queue."""
- # Built-in drying presets per filament type (from BambuStudio filament profiles)
- # Format: { n3f_temp, n3s_temp, n3f_hours, n3s_hours }
- DEFAULT_DRYING_PRESETS: dict[str, dict[str, int]] = {
- "PLA": {"n3f": 45, "n3s": 45, "n3f_hours": 12, "n3s_hours": 12},
- "PETG": {"n3f": 65, "n3s": 65, "n3f_hours": 12, "n3s_hours": 12},
- "TPU": {"n3f": 65, "n3s": 75, "n3f_hours": 12, "n3s_hours": 18},
- "ABS": {"n3f": 65, "n3s": 80, "n3f_hours": 12, "n3s_hours": 8},
- "ASA": {"n3f": 65, "n3s": 80, "n3f_hours": 12, "n3s_hours": 8},
- "PA": {"n3f": 65, "n3s": 85, "n3f_hours": 12, "n3s_hours": 12},
- "PC": {"n3f": 65, "n3s": 80, "n3f_hours": 12, "n3s_hours": 8},
- "PVA": {"n3f": 65, "n3s": 85, "n3f_hours": 12, "n3s_hours": 18},
- }
- def __init__(self):
- self._running = False
- self._check_interval = 30 # seconds
- # After a pass that actually dispatched something, loop again almost
- # immediately instead of sleeping the full interval (#2555). A dispatch
- # changes printer state — a batch launch fans out over several passes as
- # printers free up, a wedged head-of-line job reverts to pending, an
- # upload slot opens — and the next batch of ready work should not have to
- # wait 30 s behind an idle sleep. When a pass dispatches nothing (all
- # pending items are behind printers that are genuinely busy printing),
- # there is nothing to react to, so we fall back to the normal interval;
- # that also means this can never tight-loop, since fast ticks only
- # continue while dispatches keep happening and the queue is draining.
- self._fast_check_interval = 3 # seconds
- self._power_on_wait_time = 180 # seconds to wait for printer after power on (3 min)
- self._power_on_check_interval = 10 # seconds between connection checks
- # Track which printers are currently auto-drying (printer_id -> start timestamp)
- self._drying_in_progress: dict[int, float] = {}
- # Defensive in-memory dispatch hold (#1157): a printer that just received
- # a project_file command must not get a second dispatch until either it
- # transitions out of pre_state OR the hard timeout expires. The H2D Pro
- # can take 80–210 s to flip FINISH→PREPARE after project_file, and
- # during that window the DB busy_printers seed is empirically unreliable
- # (multi-plate batches double-/triple-dispatched onto the same printer
- # 30 s apart). Keyed by printer_id; cleared by the watchdog on success
- # or revert.
- # printer_id -> (monotonic_started_at, pre_state, pre_subtask_id)
- self._dispatch_holds: dict[int, tuple[float, str, str | None]] = {}
- # Minimum cooldown between dispatches to the same printer (covers the
- # H2D's project_file digestion window).
- self._dispatch_min_cooldown = 60.0
- # Hard timeout — drop the hold even if we never observed a transition,
- # so a lost MQTT session can't lock a printer out of the queue forever.
- # Matches the watchdog timeout (90 s) plus a safety margin so the
- # watchdog runs first on the unhappy path.
- self._dispatch_max_hold = 180.0
- # Refillable upload pool (#2602). Items whose FTP upload was launched by
- # an earlier pass and is still running. `_start_print` flips the row
- # pending -> printing only *after* the upload completes, so until then
- # the row stays `pending`: each tick, check_queue excludes these
- # item_ids from re-selection and their printers from new dispatch /
- # auto-drying, and launches only `limit - len(_inflight)` new uploads so
- # freed slots refill on the next fast tick. check_queue is the sole,
- # sequential caller and the prune done-callbacks run in the same
- # event-loop thread, so this dict needs no lock.
- # item_id -> (task, printer_id)
- self._inflight: dict[int, tuple[asyncio.Task, int | None]] = {}
- # Expected prints registered by `_start_print` that have not yet had a
- # print command sent. Populated at registration, dropped once
- # `start_print()` succeeds, and rolled back by `_dispatch_one` on every
- # other exit. Same threading argument as `_inflight` above: one
- # sequential caller, callbacks on the same loop, so no lock.
- # item_id -> (printer_id, remote_filename, archive_id)
- self._unconfirmed_expected_print: dict[int, tuple[int, str, int]] = {}
- async def run(self):
- """Main loop - check queue every interval."""
- self._running = True
- logger.info("Print scheduler started")
- await self._clear_stale_dispatch_claims(at_startup=True)
- while self._running:
- dispatched = False
- try:
- # No-op while any upload is in flight; on a quiet tick it releases
- # a claim whose best-effort clear failed (e.g. the database was
- # briefly unreachable), instead of leaving the row wedged until
- # the next restart.
- await self._clear_stale_dispatch_claims()
- dispatched = await self.check_queue()
- except Exception as e:
- logger.error("Scheduler error: %s", e)
- # Re-check quickly after a productive pass so a draining batch does
- # not stall behind the idle interval; otherwise sleep normally (#2555).
- await asyncio.sleep(self._fast_check_interval if dispatched else self._check_interval)
- async def _clear_stale_dispatch_claims(self, *, at_startup: bool = False) -> None:
- """Clear dispatch claims with no live dispatch coroutine behind them (#2615).
- A claim is only ever held by a live dispatch coroutine, so when this
- process has nothing in ``_inflight`` every ``dispatching_at`` in the table
- is stale. At startup that is trivially true — no coroutine survives a
- restart. It is equally true on any later tick where no upload is running,
- which is what makes this safe to repeat rather than only run once.
- Repeating it matters because ``_clear_dispatch_claim`` is best-effort: if
- the database is briefly unreachable at exactly the moment dispatch ends,
- the claim survives and the row is wedged out of the selection query. That
- used to last until the next restart (#2702 follow-up, seen when
- PostgreSQL refused a connection mid-dispatch).
- ``_inflight`` is populated when the task is spawned, before the coroutine
- claims its row, and pruned by a done-callback that cannot run before the
- coroutine's own ``finally`` — so "claim present, nothing in flight" has no
- race window and needs no age threshold. A size-derived upload deadline
- (``max(600s, size/25KB/s)``) has no safe fixed bound anyway.
- """
- if self._inflight:
- return
- try:
- async with async_session() as db:
- res = await db.execute(
- update(PrintQueueItem).where(PrintQueueItem.dispatching_at.is_not(None)).values(dispatching_at=None)
- )
- await db.commit()
- if res.rowcount:
- logger.info(
- "Cleared %d orphaned dispatch claim(s)%s (#2615)",
- res.rowcount,
- " at startup" if at_startup else "",
- )
- except Exception as exc:
- logger.error("Failed to clear orphaned dispatch claims: %s", exc)
- def stop(self):
- """Stop the scheduler."""
- self._running = False
- logger.info("Print scheduler stopped")
- async def check_queue(self) -> bool:
- """Check for prints ready to start.
- Returns True if this pass dispatched at least one item, so the caller
- can loop again quickly instead of sleeping the full interval (#2555).
- """
- async with async_session() as db:
- # Check if shortest-job-first scheduling is enabled
- sjf_enabled = await self._get_bool_setting(db, "queue_shortest_first")
- # Get all pending items, ordered by printer and position (or SJF order)
- if sjf_enabled:
- # SJF: group by printer (and target_model for model-based jobs),
- # then items already jumped get top priority (starvation guard),
- # then sort by print_time ascending. Items with no print time go last.
- result = await db.execute(
- select(PrintQueueItem)
- .where(PrintQueueItem.status == "pending")
- # Never re-select a row a dispatch worker has already claimed
- # (#2615) — belt-and-suspenders with the _inflight exclusion
- # below, and the guard that lets an orphaned claim be ignored
- # until startup reconciliation clears it.
- .where(PrintQueueItem.dispatching_at.is_(None))
- # archive/library_file are read by the cross-model gate
- # (#2578); eager-load once per pass instead of a lazy-load
- # (which would raise in async) per item.
- .options(
- selectinload(PrintQueueItem.archive),
- selectinload(PrintQueueItem.library_file),
- )
- .order_by(
- PrintQueueItem.printer_id,
- PrintQueueItem.target_model,
- PrintQueueItem.been_jumped.desc(),
- PrintQueueItem.print_time_seconds.asc().nullslast(),
- PrintQueueItem.position,
- )
- )
- else:
- result = await db.execute(
- select(PrintQueueItem)
- .where(PrintQueueItem.status == "pending")
- # Skip rows already claimed by a dispatch worker (#2615).
- .where(PrintQueueItem.dispatching_at.is_(None))
- .options(
- selectinload(PrintQueueItem.archive),
- selectinload(PrintQueueItem.library_file),
- )
- .order_by(PrintQueueItem.printer_id, PrintQueueItem.position)
- )
- items = list(result.scalars().all())
- # Drop rows whose upload is still in flight from an earlier pass
- # (#2602). They stay `pending` until the upload finishes, so without
- # this a fast tick would re-select and re-dispatch the same row.
- # Belt-and-suspenders with the printer exclusion below.
- if self._inflight:
- inflight_ids = set(self._inflight)
- items = [it for it in items if it.id not in inflight_ids]
- # Read plate-clear setting once per queue check. Default MUST be
- # False to match the schema (SettingsSchema.require_plate_clear
- # defaults False) and the frontend (toggle + card badge both treat a
- # missing value as off). When no settings row exists, a True default
- # here re-enabled the plate-clear gate the UI showed as disabled,
- # blocking dispatch to FINISH-state printers forever with no UI path
- # to clear it (#1865).
- require_plate_clear = await self._get_bool_setting(db, "require_plate_clear", default=False)
- if not items:
- # No dispatchable pending items — still check auto-drying on idle
- # printers, but keep any printer with an upload still in flight
- # from an earlier pass out of it (#2602): its print is imminent,
- # so it must not be auto-dried in the gap before the row flips to
- # printing. Report the pass as productive while uploads run so the
- # loop stays on the fast interval.
- inflight_printers = {pid for (_task, pid) in self._inflight.values() if pid is not None}
- await self._check_auto_drying(db, [], inflight_printers, require_plate_clear=require_plate_clear)
- return bool(self._inflight)
- logger.info(
- "Queue check: found %d pending items: %s",
- len(items),
- [(i.id, i.printer_id, i.archive_id, i.library_file_id) for i in items],
- )
- # Seed busy_printers with printers that already have an item in 'printing'
- # status. _is_printer_idle() alone is not sufficient as a dispatch gate —
- # on H2D / P1 series the MQTT state transition from IDLE to RUNNING can
- # lag several seconds behind the print command, so the next check_queue
- # tick still sees IDLE and would double-dispatch onto the same printer.
- # Without this guard, two pending items targeting the same printer
- # (e.g. a batch with quantity>1) both end up in 'printing' status —
- # surfaced via the "BUG: Multiple queue items" warning in on_print_complete.
- busy_result = await db.execute(
- select(PrintQueueItem.printer_id)
- .where(PrintQueueItem.status == "printing")
- .where(PrintQueueItem.printer_id.is_not(None))
- )
- busy_printers: set[int] = {pid for (pid,) in busy_result.all() if pid is not None}
- # Defense-in-depth (#1157): augment busy_printers with any printer
- # still in its post-dispatch hold window. Empirically, the DB seed
- # above can miss in-flight items in a multi-plate batch — same-file
- # plates were being dispatched 30 s apart while the H2D was still
- # digesting the first project_file. The hold is keyed in-memory and
- # released by the watchdog on the success path, so it adds a layer
- # that doesn't depend on DB row visibility or completion-callback
- # timing.
- for held_printer_id in list(self._dispatch_holds.keys()):
- if self._printer_in_dispatch_hold(held_printer_id):
- busy_printers.add(held_printer_id)
- # Exclude printers whose upload is still in flight from an earlier
- # pass (#2602). The row is `pending` until the upload finishes and
- # the printing-state seed / dispatch hold above only arm once the
- # upload completes, so this is what holds the printer (and, via
- # busy_printers, its auto-drying) out of the pass during the upload.
- for _task, inflight_pid in self._inflight.values():
- if inflight_pid is not None:
- busy_printers.add(inflight_pid)
- # Log skip reasons once per queue check (not per item)
- skip_reasons: dict[str, int] = {}
- # Items selected for dispatch in this pass, one per printer. The
- # loop below only *decides* — the uploads happen afterwards, in
- # parallel (#2555). See _dispatch_selected().
- dispatch_ids: list[int] = []
- # Library rows queued with `cleanup_library_after_dispatch` (the
- # printer-card "upload and print" flow) are CONSUMED by the dispatch
- # that prints them: the row is deleted and the 3MF is unlinked from
- # disk. That was safe only because dispatch was serial. Run two of
- # them against the same row at once and the second DELETE matches no
- # row (StaleDataError), and the winner's unlink can pull the file out
- # from under the loser's in-flight upload.
- #
- # Only the cleanup flag mutates the row. An ordinary library print
- # just reads it, so the common fan-out — one file, many printers,
- # which is exactly the reporter's workload — still goes out fully in
- # parallel. Narrow the guard to the mutating case; do not serialise
- # the case the whole fix exists for.
- dispatch_libs: set[int] = set()
- consumed_libs: set[int] = set()
- def _library_row_conflict(candidate: PrintQueueItem) -> bool:
- """True if dispatching `candidate` now would race another item's cleanup."""
- lib_id = candidate.library_file_id
- if lib_id is None:
- return False
- if candidate.cleanup_library_after_dispatch:
- # We would delete a row someone else in this pass is reading.
- return lib_id in dispatch_libs
- # Someone else in this pass will delete the row out from under us.
- return lib_id in consumed_libs
- def _claim_library_row(candidate: PrintQueueItem) -> None:
- lib_id = candidate.library_file_id
- if lib_id is None:
- return
- dispatch_libs.add(lib_id)
- if candidate.cleanup_library_after_dispatch:
- consumed_libs.add(lib_id)
- for item in items:
- # Check scheduled time first (scheduled_time is stored in UTC from ISO string)
- if item.scheduled_time:
- sched = item.scheduled_time
- if sched.tzinfo is None:
- sched = sched.replace(tzinfo=timezone.utc)
- if sched > datetime.now(timezone.utc):
- skip_reasons["scheduled_future"] = skip_reasons.get("scheduled_future", 0) + 1
- continue
- # Skip items that require manual start
- if item.manual_start:
- skip_reasons["manual_start"] = skip_reasons.get("manual_start", 0) + 1
- continue
- if item.printer_id:
- # Specific printer assignment (existing behavior)
- if item.printer_id in busy_printers:
- continue
- # Check if printer is idle
- printer_idle = self._is_printer_idle(item.printer_id, require_plate_clear)
- printer_connected = printer_manager.is_connected(item.printer_id)
- # If printer not connected, try to power on via smart plug
- if not printer_connected:
- plugs = await self._get_smart_plugs(db, item.printer_id)
- auto_on_plugs = [p for p in plugs if p.auto_on and p.enabled]
- if auto_on_plugs:
- logger.info("Printer %s offline, attempting to power on via smart plug(s)", item.printer_id)
- # Power on using the plug that actually feeds the printer, and
- # wait for it to boot on that one only (#2629).
- primary_plug = self._pick_power_plug(auto_on_plugs)
- powered_on = await self._power_on_and_wait(primary_plug, item.printer_id, db)
- if powered_on:
- # Also turn on any remaining auto_on plugs (e.g., filter)
- for extra_plug in [p for p in auto_on_plugs if p.id != primary_plug.id]:
- try:
- service = await smart_plug_manager.get_service_for_plug(extra_plug, db)
- await service.turn_on(extra_plug)
- logger.info(
- "Also powered on plug '%s' for printer %s", extra_plug.name, item.printer_id
- )
- except Exception as e:
- logger.warning("Failed to power on extra plug '%s': %s", extra_plug.name, e)
- printer_connected = True
- printer_idle = self._is_printer_idle(item.printer_id, require_plate_clear)
- else:
- logger.warning("Could not power on printer %s via smart plug", item.printer_id)
- busy_printers.add(item.printer_id)
- continue
- else:
- # No plug or auto_on disabled
- busy_printers.add(item.printer_id)
- continue
- # Check if printer is idle (busy with another print)
- if not printer_idle:
- # If printer is drying (not truly busy), handle based on queue_drying_block
- if self._drying_in_progress.get(item.printer_id):
- block_for_drying = await self._get_bool_setting(db, "queue_drying_block")
- if block_for_drying:
- # Drying blocks queue — skip this printer
- busy_printers.add(item.printer_id)
- continue
- else:
- # Print takes priority — stop drying
- await self._stop_drying(item.printer_id)
- # Re-check idle after stopping drying
- printer_idle = self._is_printer_idle(item.printer_id, require_plate_clear)
- if not printer_idle:
- busy_printers.add(item.printer_id)
- continue
- else:
- busy_printers.add(item.printer_id)
- continue
- # Check condition (previous print success)
- if item.require_previous_success:
- if not await self._check_previous_success(db, item):
- item.status = "skipped"
- item.error_message = "Previous print failed or was aborted"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.info("Skipped queue item %s - previous print failed", item.id)
- # Send notification
- job_name = await self._get_job_name(db, item)
- printer = await self._get_printer(db, item.printer_id)
- await notification_service.on_queue_job_skipped(
- job_name=job_name,
- printer_id=item.printer_id,
- printer_name=printer.name if printer else "Unknown",
- reason="Previous print failed or was aborted",
- db=db,
- )
- continue
- # Resolve the AMS mapping when it's missing OR unresolved
- # (all -1). A stored all-[-1] mapping is a bug artifact — a
- # frontend status-load race can persist [-1] (#2589) — and
- # must be recomputed from live trays rather than trusted.
- await self._ensure_ams_mapping(db, item.printer_id, item)
- # Filament-deficit pre-dispatch check (#1496). If the
- # assigned spool can't satisfy any required slot grams,
- # promote the item to manual_start so the user must
- # acknowledge via the ▶ button (which re-checks live).
- if await self._block_on_filament_deficit(db, item):
- continue
- # Hold this item back for the next pass rather than racing
- # another dispatch over the same transient library row. The
- # printer is still marked busy so a later item does not jump
- # its place in this printer's queue.
- if _library_row_conflict(item):
- skip_reasons["library_row_in_use"] = skip_reasons.get("library_row_in_use", 0) + 1
- busy_printers.add(item.printer_id)
- continue
- # Queue the dispatch instead of running it here — see
- # _dispatch_selected(). busy_printers still gets the printer
- # immediately, so nothing else in this pass can target it.
- _claim_library_row(item)
- dispatch_ids.append(item.id)
- busy_printers.add(item.printer_id)
- # SJF starvation guard: mark items that were jumped
- if sjf_enabled and item.print_time_seconds is not None:
- for other in items:
- if (
- other.id != item.id
- and other.status == "pending"
- and other.printer_id == item.printer_id
- and not other.been_jumped
- and other.position < item.position
- and (
- other.print_time_seconds is None
- or other.print_time_seconds > item.print_time_seconds
- )
- ):
- other.been_jumped = True
- await db.commit()
- elif item.target_model:
- # Model-based assignment - find any idle printer of matching model
- # Parse required filament types if present
- required_types = None
- if item.required_filament_types:
- try:
- required_types = json.loads(item.required_filament_types)
- except json.JSONDecodeError:
- pass # Ignore malformed filament types; treat as no constraint
- # Parse filament overrides if present
- filament_overrides = None
- if item.filament_overrides:
- try:
- filament_overrides = json.loads(item.filament_overrides)
- except json.JSONDecodeError:
- pass
- # If overrides exist, use override types for validation instead
- effective_types = required_types
- if filament_overrides:
- override_types = sorted({o["type"] for o in filament_overrides if "type" in o})
- if override_types:
- # Merge: keep original types for non-overridden slots, add override types
- effective_types = sorted(set(required_types or []) | set(override_types))
- # Cross-model safety gate (#2578): never hand a 3MF sliced
- # for an incompatible model to a printer, no matter how the
- # row got into the DB (old rows, direct API writes). Held
- # as pending with an actionable waiting_reason — the user
- # fixes it by editing the item's target model.
- sliced_for = None
- if item.archive:
- sliced_for = item.archive.sliced_for_model
- elif item.library_file and item.library_file.file_metadata:
- sliced_for = item.library_file.file_metadata.get("sliced_for_model")
- if not is_gcode_compatible(sliced_for, item.target_model):
- printer_id = None
- waiting_reason = (
- f"File was sliced for {sliced_for}, which is not compatible with "
- f"{item.target_model} — edit the item and fix its target model"
- )
- skip_reasons["sliced_model_mismatch"] = skip_reasons.get("sliced_model_mismatch", 0) + 1
- else:
- printer_id, waiting_reason = await self._find_idle_printer_for_model(
- db,
- item.target_model,
- busy_printers,
- effective_types,
- item.target_location,
- filament_overrides=filament_overrides,
- require_plate_clear=require_plate_clear,
- )
- # Update waiting_reason if changed and send notification when first waiting
- if item.waiting_reason != waiting_reason:
- was_waiting = item.waiting_reason is not None
- item.waiting_reason = waiting_reason
- await db.commit()
- # Send waiting notification only when transitioning to waiting state
- # and the reason requires user action (not just "all printers busy")
- if waiting_reason and not was_waiting and not self._is_busy_only(waiting_reason):
- job_name = await self._get_job_name(db, item)
- await notification_service.on_queue_job_waiting(
- job_name=job_name,
- target_model=item.target_model,
- waiting_reason=waiting_reason,
- db=db,
- )
- if printer_id:
- # Before claiming the printer: hold back rather than race
- # another dispatch over the same transient library row.
- # Checked here so a held item does not get a printer
- # assigned and then sit on it. See _library_row_conflict().
- #
- # No busy_printers.add() here, unlike the fixed-printer
- # branch above: that one protects its printer's own queue
- # ordering, but this item was never assigned to `printer_id`
- # — the matcher merely offered it. Marking it busy would
- # strand an idle printer for the rest of the pass.
- if _library_row_conflict(item):
- skip_reasons["library_row_in_use"] = skip_reasons.get("library_row_in_use", 0) + 1
- continue
- # Check condition (previous print success) before assigning
- if item.require_previous_success:
- if not await self._check_previous_success(db, item):
- item.status = "skipped"
- item.error_message = "Previous print failed or was aborted"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.info("Skipped queue item %s - previous print failed", item.id)
- # Send notification
- job_name = await self._get_job_name(db, item)
- printer = await self._get_printer(db, printer_id)
- await notification_service.on_queue_job_skipped(
- job_name=job_name,
- printer_id=printer_id,
- printer_name=printer.name if printer else "Unknown",
- reason="Previous print failed or was aborted",
- db=db,
- )
- continue
- # Assign printer and start - clear waiting reason
- item.printer_id = printer_id
- item.waiting_reason = None
- logger.info("Model-based assignment: queue item %s assigned to printer %s", item.id, printer_id)
- # Send assignment notification
- job_name = await self._get_job_name(db, item)
- printer = await self._get_printer(db, printer_id)
- await notification_service.on_queue_job_assigned(
- job_name=job_name,
- printer_id=printer_id,
- printer_name=printer.name if printer else "Unknown",
- target_model=item.target_model,
- db=db,
- )
- # Resolve the AMS mapping for the assigned printer when it's
- # missing OR unresolved (all -1). Critical for model-based
- # jobs where mapping wasn't computed upfront, and it also
- # self-heals a bogus stored [-1] (#2589).
- await self._ensure_ams_mapping(db, printer_id, item)
- # Filament-deficit pre-dispatch check (#1496).
- if await self._block_on_filament_deficit(db, item):
- continue
- _claim_library_row(item)
- dispatch_ids.append(item.id)
- busy_printers.add(printer_id)
- # SJF starvation guard: mark model-based items that were jumped
- if sjf_enabled and item.print_time_seconds is not None:
- for other in items:
- if (
- other.id != item.id
- and other.status == "pending"
- and other.printer_id is None
- and other.target_model
- and other.target_model.upper() == item.target_model.upper()
- and not other.been_jumped
- and other.position < item.position
- and (
- other.print_time_seconds is None
- or other.print_time_seconds > item.print_time_seconds
- )
- ):
- other.been_jumped = True
- await db.commit()
- # Log the decisions BEFORE dispatching. The dispatch below blocks for
- # as long as the slowest upload takes (minutes on a big 3MF), and a
- # skip summary that only lands after the transfers have finished is
- # useless for working out why an item did not go out.
- if skip_reasons:
- logger.info("Queue skip summary: %s", skip_reasons)
- if busy_printers:
- # Log why each printer was busy (first time it was checked)
- for pid in busy_printers:
- state = printer_manager.get_status(pid)
- connected = printer_manager.is_connected(pid)
- awaiting = printer_manager.is_awaiting_plate_clear(pid)
- state_name = state.state if state else "NO_STATUS"
- logger.info(
- "Queue: printer %d not available — connected=%s, state=%s, awaiting_plate_clear=%s",
- pid,
- connected,
- state_name,
- awaiting,
- )
- # Read the concurrency limit BEFORE the commit below, not inside
- # _dispatch_selected(). A SELECT on this session after the commit
- # implicitly opens a fresh transaction that nothing then closes, and
- # it would stay open for the whole dispatch — minutes of "idle in
- # transaction" on Postgres (pinned MVCC snapshot, vacuum blocked),
- # and on SQLite a pinned WAL read snapshot that stops the WAL being
- # checkpointed while every dispatch is writing to it.
- upload_limit = max(1, await self._get_int_setting(db, "queue_max_concurrent_uploads", default=4))
- # Selection is done; every decision above is recorded on `db`
- # (model-based printer assignment, computed ams_mapping). Flush it
- # before the dispatch tasks open their own sessions, or they will
- # read a row that still says printer_id=None. This also releases the
- # connection back to the pool for the duration of the dispatch.
- await db.commit()
- if dispatch_ids:
- item_printers = {it.id: it.printer_id for it in items}
- self._launch_uploads(dispatch_ids, item_printers, upload_limit)
- # Auto-drying: start drying on idle printers that have no pending queue items
- await self._check_auto_drying(db, items, busy_printers, require_plate_clear=require_plate_clear)
- # Keep the loop on the fast interval while any upload is in flight so
- # a slot freed mid-tick refills within seconds rather than after the
- # 30 s idle sleep (#2602). Selecting anything this pass (launched or
- # deferred because the pool was full) also counts as productive.
- return bool(dispatch_ids) or bool(self._inflight)
- def _launch_uploads(self, item_ids: list[int], item_printers: dict[int, int | None], limit: int) -> None:
- """Launch selected uploads as a refillable pool, capped at ``limit`` (#2602).
- Dispatch used to happen inline in the selection loop: ``await
- _start_print(db, item)`` per item in turn. Since ``_start_print``
- performs the FTP upload, that serialized every printer behind every
- other printer's transfer even though the printers are independent
- machines; #2555 moved it to a parallel ``asyncio.gather()``. But that
- gather was awaited before ``check_queue`` returned, so the run loop
- stayed blocked until the *slowest* upload in the batch finished — a
- 513 s upload left 15 of 16 configured slots idle for 8.5 minutes on a
- 93-printer farm even as other printers came free (#2602).
- Each upload now runs as an independent background task tracked in
- ``self._inflight``. check_queue excludes in-flight item_ids (still
- `pending` until their upload completes) and their printers from the
- next pass's selection, and this method launches at most
- ``limit - len(self._inflight)`` new uploads, so a freed slot refills on
- the next fast tick instead of waiting out the whole batch. The bound
- exists because the printers are independent but the host is not: each
- in-flight upload holds a thread in the FTP pool, a TLS session and a
- file handle.
- The no-overlapping-dispatch invariant the batch-await used to provide
- is now carried by the in-flight exclusion in check_queue. Everything
- else — the pending->printing CAS, the busy-printer guard (#2598), the
- per-printer hold, and each item's independent failure handling — still
- lives in ``_start_print`` and runs per task exactly as before.
- Synchronous on purpose: it registers every launched task into
- ``self._inflight`` before returning, so the next (sequential) tick sees
- an accurate in-flight count with no interleaving await.
- """
- free = limit - len(self._inflight)
- if free <= 0:
- logger.info(
- "Upload pool full (%d/%d in flight) — deferring %d item(s) to a later tick: %s",
- len(self._inflight),
- limit,
- len(item_ids),
- item_ids,
- )
- return
- to_launch = item_ids[:free]
- deferred = item_ids[free:]
- logger.info(
- "Launching %d upload(s) (pool %d/%d in flight)%s",
- len(to_launch),
- len(self._inflight),
- limit,
- f" — deferring {deferred} to a later tick" if deferred else "",
- )
- for item_id in to_launch:
- task = spawn_background_task(self._dispatch_one(item_id), name=f"queue-upload-{item_id}")
- self._inflight[item_id] = (task, item_printers.get(item_id))
- # Prune on completion so the freed slot is refillable next tick.
- # spawn_background_task already logs any uncaught exception; this
- # only reclaims the pool slot (fires on success, failure, or cancel).
- task.add_done_callback(lambda _t, iid=item_id: self._inflight.pop(iid, None))
- async def _dispatch_one(self, item_id: int) -> None:
- """Upload + start one queue item in its own session (pool worker, #2602).
- Its own session: pool workers run concurrently and an AsyncSession is
- not safe to share across tasks; it also keeps a slow upload from pinning
- the scheduler's session (and, on SQLite, its transaction) open for the
- transfer's duration.
- """
- async with async_session() as item_db:
- # Claim the row for dispatch BEFORE reading the printer snapshot or
- # touching any slow I/O (#2615). The claim is an atomic CAS on
- # (status='pending', dispatching_at IS NULL); while it's held the edit
- # routes reject reassignment (409), so printer_id can't change out from
- # under the in-flight upload and split the queue row from the
- # archive/expected-print/physical command.
- if not await self._claim_for_dispatch(item_db, item_id):
- logger.info(
- "Queue item %s not claimable for dispatch (cancelled, removed, or already claimed) — skipping",
- item_id,
- )
- return
- try:
- item = await item_db.get(PrintQueueItem, item_id)
- if not item:
- logger.info("Queue item %s vanished after claim — skipping", item_id)
- return
- await self._start_print(item_db, item)
- finally:
- # Undo an expected-print registration whose print command never
- # went out. One choke point covers every way `_start_print` can
- # end without sending: a raised exception (a DB failure mid-
- # dispatch is the reported case), an early return, a cancel
- # winning the #1853 CAS, or `start_print()` returning False.
- # A confirmed send removes the entry itself, so this is a no-op
- # on the happy path.
- self._rollback_unconfirmed_expected_print(item_id)
- # Release the claim on every exit. Once dispatch has finished the
- # row's status carries the lock (printing/failed/cancelled are all
- # != pending), so the token is only needed for the duration of the
- # upload. A row left pending (e.g. busy-printer deferral) becomes
- # dispatchable again on the next tick.
- await self._clear_dispatch_claim(item_db, item_id)
- def _rollback_unconfirmed_expected_print(self, item_id: int) -> None:
- """Drop an expectation for a print command that was never sent.
- Best-effort and never raises: this runs in the ``finally`` of dispatch,
- where the interesting exception is usually the one already propagating.
- """
- pending = self._unconfirmed_expected_print.pop(item_id, None)
- if pending is None:
- return
- printer_id, remote_filename, archive_id = pending
- try:
- from backend.app.main import unregister_expected_print
- unregister_expected_print(printer_id, remote_filename, archive_id)
- except Exception:
- logger.warning(
- "Queue item %s: failed to unregister expected print (printer=%s, file=%s, archive=%s)",
- item_id,
- printer_id,
- remote_filename,
- archive_id,
- exc_info=True,
- )
- async def _claim_for_dispatch(self, db: AsyncSession, item_id: int) -> bool:
- """Atomically stamp ``dispatching_at`` on a still-pending, unclaimed row.
- Returns True if this call won the claim, False if the row was already
- claimed, no longer pending (cancelled mid-tick), or removed. The CAS is
- the load-bearing guard against reassign-during-dispatch (#2615)."""
- res = await db.execute(
- update(PrintQueueItem)
- .where(PrintQueueItem.id == item_id)
- .where(PrintQueueItem.status == "pending")
- .where(PrintQueueItem.dispatching_at.is_(None))
- .values(dispatching_at=datetime.now(timezone.utc))
- )
- await db.commit()
- return res.rowcount > 0
- async def _clear_dispatch_claim(self, db: AsyncSession, item_id: int) -> None:
- """Clear the dispatch claim (#2615). Best-effort: a failure here must not
- mask the dispatch outcome.
- Retried, because the failure mode in practice is transient and narrow: a
- database that is momentarily unreachable — PostgreSQL out of connection
- slots is the observed case — refuses this write for a second or two while
- the dispatch that just ended is still holding the row out of the selection
- query. One attempt was enough to wedge the item; a couple of spaced
- attempts clear it. Each attempt rolls back first, since a failed write
- leaves the session needing it before it can be reused.
- If every attempt fails, ``_clear_stale_dispatch_claims`` picks the row up
- on the next quiet tick.
- """
- for attempt in range(1, 4):
- try:
- await db.execute(update(PrintQueueItem).where(PrintQueueItem.id == item_id).values(dispatching_at=None))
- await db.commit()
- return
- except Exception as exc:
- try:
- await db.rollback()
- except Exception:
- pass
- if attempt == 3:
- logger.warning(
- "Queue item %s: failed to clear dispatch claim after %d attempts: %s "
- "— a later quiet tick will release it",
- item_id,
- attempt,
- exc,
- )
- return
- await asyncio.sleep(0.5 * attempt)
- async def _find_idle_printer_for_model(
- self,
- db: AsyncSession,
- model: str,
- exclude_ids: set[int],
- required_filament_types: list[str] | None = None,
- target_location: str | None = None,
- filament_overrides: list[dict] | None = None,
- require_plate_clear: bool = True,
- ) -> tuple[int | None, str | None]:
- """Find an idle, connected printer matching the model with compatible filaments.
- Args:
- db: Database session
- model: Printer model to match (e.g., "X1C", "P1S")
- exclude_ids: Printer IDs to exclude (already busy)
- required_filament_types: Optional list of filament types needed (e.g., ["PLA", "PETG"])
- If provided, only printers with all required types loaded will match.
- target_location: Optional location filter. If provided, only printers in this location are considered.
- filament_overrides: Optional list of override dicts. Each entry may include
- ``force_color_match: true`` to require an exact type+color match
- on the printer for that slot. Without the flag the existing
- colour-preference logic applies.
- Returns:
- Tuple of (printer_id, waiting_reason):
- - (printer_id, None) if a matching printer was found
- - (None, reason) if no printer is available, with explanation
- """
- # Normalize model name and use case-insensitive matching
- normalized_model = normalize_printer_model(model) or model
- query = (
- select(Printer)
- .where(func.lower(Printer.model) == normalized_model.lower())
- .where(Printer.is_active == True) # noqa: E712
- )
- # Add location filter if specified
- if target_location:
- query = query.where(Printer.location == target_location)
- result = await db.execute(query)
- printers = list(result.scalars().all())
- location_suffix = f" in {target_location}" if target_location else ""
- if not printers:
- return None, f"No active {normalized_model} printers{location_suffix} configured"
- # Separate force-matched overrides from preference-only overrides
- force_overrides = [o for o in (filament_overrides or []) if o.get("force_color_match")]
- pref_overrides = [o for o in (filament_overrides or []) if not o.get("force_color_match")]
- # Track reasons for skipping printers
- printers_busy = []
- printers_offline = []
- printers_missing_filament: list[tuple[str, list[str]]] = []
- candidates: list[tuple[int, int]] = [] # (printer_id, color_match_count)
- for printer in printers:
- if printer.id in exclude_ids:
- # Printer is already claimed by another job in this scheduling run.
- # For force-color jobs, still check if the color would match — if not,
- # report it as a color mismatch rather than plain "Busy" so the user
- # knows the job needs a filament change, not just to wait for availability.
- if force_overrides and not pref_overrides:
- missing_colors = self._get_missing_force_color_slots(printer.id, force_overrides)
- if missing_colors:
- printers_missing_filament.append((printer.name, missing_colors))
- continue
- printers_busy.append(printer.name)
- continue
- is_connected = printer_manager.is_connected(printer.id)
- is_idle = self._is_printer_idle(printer.id, require_plate_clear) if is_connected else False
- if not is_connected:
- printers_offline.append(printer.name)
- continue
- if not is_idle:
- # Printer is currently printing. For force-color jobs, check whether the
- # loaded color would satisfy the requirement — if not, surface it as a
- # color-mismatch reason rather than plain "Busy" so the user understands
- # that the job is waiting for a filament change, not just printer availability.
- if force_overrides and not pref_overrides:
- missing_colors = self._get_missing_force_color_slots(printer.id, force_overrides)
- if missing_colors:
- printers_missing_filament.append((printer.name, missing_colors))
- logger.debug(
- "Printer %s (%s) is busy but also has wrong force-color: %s",
- printer.id,
- printer.name,
- missing_colors,
- )
- continue
- printers_busy.append(printer.name)
- continue
- # Validate filament compatibility if required types are specified
- if required_filament_types:
- missing = self._get_missing_filament_types(printer.id, required_filament_types)
- if missing:
- # When force_overrides are present, enrich missing entries with color info
- # so the "Waiting on" message includes "TYPE (color)" instead of just "TYPE"
- if force_overrides:
- force_color_map = {
- (o.get("type") or "").upper(): o.get("color_name") or o.get("color", "?")
- for o in force_overrides
- }
- missing_enriched = [
- f"{t} ({force_color_map[t_upper]})" if (t_upper := t.upper()) in force_color_map else t
- for t in missing
- ]
- printers_missing_filament.append((printer.name, missing_enriched))
- else:
- printers_missing_filament.append((printer.name, missing))
- logger.debug("Skipping printer %s (%s) - missing filaments: %s", printer.id, printer.name, missing)
- continue
- # Force color match: ALL flagged slots must have an exact type+color match
- if force_overrides:
- missing_colors = self._get_missing_force_color_slots(printer.id, force_overrides)
- if missing_colors:
- printers_missing_filament.append((printer.name, missing_colors))
- logger.debug(
- "Skipping printer %s (%s) - missing force-matched colors: %s",
- printer.id,
- printer.name,
- missing_colors,
- )
- continue
- # If preference-only overrides exist, rank by color matches (existing behaviour)
- if pref_overrides:
- color_matches = self._count_override_color_matches(printer.id, pref_overrides)
- if color_matches > 0:
- candidates.append((printer.id, color_matches))
- else:
- override_colors = [f"{o.get('type', '?')} ({o.get('color', '?')})" for o in pref_overrides]
- printers_missing_filament.append((printer.name, override_colors))
- logger.debug("Skipping printer %s (%s) - no matching override colors", printer.id, printer.name)
- continue
- elif force_overrides:
- # Passed all force checks — immediately eligible (no preference ordering needed)
- return printer.id, None
- else:
- # No overrides at all - take first available (existing behavior)
- return printer.id, None
- # If we have candidates from preference override matching, pick the one with most color matches
- if candidates:
- candidates.sort(key=lambda c: c[1], reverse=True)
- return candidates[0][0], None
- # Build waiting reason from what we found
- reasons = []
- if printers_missing_filament:
- # Filament/color mismatch is most actionable - show first
- if force_overrides and not pref_overrides:
- # All mismatches are force-color failures — use descriptive message only;
- # but only if there are no busy printers that DO have the matching color.
- # If a printer has the right color but is busy, surface "Busy" instead so
- # the user knows the job will start automatically once that printer is free.
- if not printers_busy:
- all_missing = sorted({c for _, cols in printers_missing_filament for c in cols})
- return None, f"No matching material/color. Waiting on {', '.join(all_missing)}"
- # else: fall through — printers_busy will be appended below
- else:
- names_and_missing = [
- f"{name} (needs {', '.join(missing)})" for name, missing in printers_missing_filament
- ]
- reasons.append(f"Waiting for filament: {'; '.join(names_and_missing)}")
- if printers_busy:
- reasons.append(f"Busy: {', '.join(printers_busy)}")
- if printers_offline:
- reasons.append(f"Offline: {', '.join(printers_offline)}")
- return None, " | ".join(reasons) if reasons else f"No available {model} printers{location_suffix}"
- @staticmethod
- def _is_busy_only(waiting_reason: str) -> bool:
- """Check if the waiting reason only contains 'Busy' entries.
- When all matching printers are simply busy printing, the queued job
- will start automatically once a printer finishes — no user action
- is required, so we skip the notification.
- """
- parts = [p.strip() for p in waiting_reason.split(" | ")]
- return all(p.startswith("Busy:") for p in parts)
- def _get_missing_force_color_slots(self, printer_id: int, force_overrides: list[dict]) -> list[str]:
- """Return descriptive strings for force_color_match slots not satisfied by the printer.
- Each entry in ``force_overrides`` must have ``type`` and ``color`` fields and is expected
- to carry ``force_color_match: True``. The printer must have **every** such slot loaded
- with an exact type+color match.
- When both the override and a candidate tray carry a ``tray_info_idx``, they must also
- match on it: Bambu reports every PLA variant as ``tray_type == "PLA"``, so the
- Basic/Matte/Silk distinction lives only in ``tray_info_idx`` (GFA00/GFA01/GFA06/...).
- Without this, a job sliced for PLA Matte matched every white PLA regardless of variant
- (#2650). If either side lacks an idx (custom/third-party spools report a blank one, and
- older 3MFs carry none) we fall back to the historical type+colour behaviour so those
- setups are unaffected.
- Returns:
- List of ``"TYPE (color)"`` strings for unmatched slots (empty list means all match).
- """
- status = printer_manager.get_status(printer_id)
- if not status:
- return [f"{o.get('type', '?')} ({o.get('color_name') or o.get('color', '?')})" for o in force_overrides]
- # Build loaded (type, colour, tray_info_idx) triples from AMS and external spool.
- loaded: list[tuple[str, str, str]] = []
- for ams_unit in status.raw_data.get("ams", []):
- for tray in ams_unit.get("tray", []):
- tray_type = tray.get("tray_type")
- if tray_type:
- color_norm = (tray.get("tray_color", "") or "").replace("#", "").lower()[:6]
- loaded.append(
- (_canonical_filament_type(tray_type), color_norm, tray.get("tray_info_idx", "") or "")
- )
- for vt in status.raw_data.get("vt_tray") or []:
- vt_type = vt.get("tray_type")
- if vt_type:
- color_norm = (vt.get("tray_color", "") or "").replace("#", "").lower()[:6]
- loaded.append((_canonical_filament_type(vt_type), color_norm, vt.get("tray_info_idx", "") or ""))
- missing = []
- for o in force_overrides:
- o_type = _canonical_filament_type(o.get("type") or "")
- o_color = (o.get("color") or "").replace("#", "").lower()[:6]
- o_idx = o.get("tray_info_idx") or ""
- satisfied = any(
- t_type == o_type and t_color == o_color and (not o_idx or not t_idx or o_idx == t_idx)
- for t_type, t_color, t_idx in loaded
- )
- if not satisfied:
- color_label = o.get("color_name") or o.get("color", "?")
- missing.append(f"{o_type} ({color_label})")
- return missing
- def _get_missing_filament_types(self, printer_id: int, required_types: list[str]) -> list[str]:
- """Get the list of required filament types that are not loaded on the printer.
- Args:
- printer_id: The printer ID
- required_types: List of filament types needed (e.g., ["PLA", "PETG"])
- Returns:
- List of missing filament types (empty if all are loaded)
- """
- status = printer_manager.get_status(printer_id)
- if not status:
- return required_types # Can't determine, assume all missing
- # Collect all filament types loaded on this printer (AMS units + external spool)
- # Use canonical types so equivalence groups (e.g. PA-CF/PA12-CF/PAHT-CF) match.
- loaded_types: set[str] = set()
- # Check AMS units (stored in raw_data["ams"])
- ams_data = status.raw_data.get("ams", [])
- if ams_data:
- for ams_unit in ams_data:
- for tray in ams_unit.get("tray", []):
- tray_type = tray.get("tray_type")
- if tray_type:
- loaded_types.add(_canonical_filament_type(tray_type))
- # Check external spool(s) (virtual tray, stored in raw_data["vt_tray"] as list)
- for vt in status.raw_data.get("vt_tray") or []:
- vt_type = vt.get("tray_type")
- if vt_type:
- loaded_types.add(_canonical_filament_type(vt_type))
- # Find which required types are missing (using canonical type for equivalence)
- missing = []
- for req_type in required_types:
- if _canonical_filament_type(req_type) not in loaded_types:
- missing.append(req_type)
- return missing
- def _count_override_color_matches(self, printer_id: int, overrides: list[dict]) -> int:
- """Count how many filament overrides have an exact color match on the printer.
- Used to prefer printers that already have the desired override colors loaded.
- """
- status = printer_manager.get_status(printer_id)
- if not status:
- return 0
- # Collect loaded filaments' type+color pairs
- loaded: set[tuple[str, str]] = set()
- for ams_unit in status.raw_data.get("ams", []):
- for tray in ams_unit.get("tray", []):
- tray_type = tray.get("tray_type")
- tray_color = tray.get("tray_color", "")
- if tray_type:
- color_norm = tray_color.replace("#", "").lower()[:6]
- loaded.add((tray_type.upper(), color_norm))
- for vt in status.raw_data.get("vt_tray") or []:
- vt_type = vt.get("tray_type")
- if vt_type:
- color_norm = (vt.get("tray_color", "") or "").replace("#", "").lower()[:6]
- loaded.add((vt_type.upper(), color_norm))
- matches = 0
- for o in overrides:
- o_type = (o.get("type") or "").upper()
- o_color = (o.get("color") or "").replace("#", "").lower()[:6]
- if (o_type, o_color) in loaded:
- matches += 1
- return matches
- async def _ensure_ams_mapping(self, db: AsyncSession, printer_id: int, item: PrintQueueItem) -> None:
- """Ensure the queue item carries a usable AMS mapping before dispatch.
- Recomputes from live printer status when the stored mapping is missing OR
- unresolved (all -1). A stored all-[-1] mapping is a bug artifact — a
- frontend status-load race can serialize [-1] before the printer's AMS
- trays are known (#2589) — and must not be trusted: downstream it would be
- silently downgraded to external-spool mode and print against an empty
- feed. A resolved mapping (including manual overrides, or a partially
- padded one) is left untouched.
- When recompute cannot resolve it either (no compatible tray loaded), the
- bogus [-1] is cleared to None so it is not later mistaken for an explicit
- external selection; the print command then keeps use_ams=True and the
- firmware surfaces a clear AMS-mapping error instead of silently printing
- to the empty external feed.
- """
- stored_mapping: list | None = None
- if item.ams_mapping:
- try:
- stored_mapping = json.loads(item.ams_mapping)
- except (json.JSONDecodeError, TypeError):
- stored_mapping = None
- # Already resolved (present and not all-unresolved) — keep as-is so a
- # user's manual mapping is never overwritten.
- if item.ams_mapping and not _mapping_is_all_unresolved(stored_mapping):
- return
- computed_mapping = await self._compute_ams_mapping_for_printer(db, printer_id, item)
- if computed_mapping and not _mapping_is_all_unresolved(computed_mapping):
- item.ams_mapping = json.dumps(computed_mapping)
- logger.info(
- "Queue item %s: Computed AMS mapping for printer %s: %s",
- item.id,
- printer_id,
- computed_mapping,
- )
- await db.commit()
- elif _mapping_is_all_unresolved(stored_mapping):
- logger.warning(
- "Queue item %s: stored ams_mapping %s is unresolved and could not be recomputed "
- "from live status on printer %s; clearing it so dispatch does not treat it as external",
- item.id,
- stored_mapping,
- printer_id,
- )
- item.ams_mapping = None
- await db.commit()
- async def _compute_ams_mapping_for_printer(
- self, db: AsyncSession, printer_id: int, item: PrintQueueItem
- ) -> list[int] | None:
- """Compute AMS mapping for a printer based on filament requirements.
- Called when a queue item has no ams_mapping set — either for model-based
- items after printer assignment, or printer-specific items (e.g. from VP).
- Args:
- db: Database session
- printer_id: The assigned printer ID
- item: The queue item (contains archive_id or library_file_id)
- Returns:
- AMS mapping array or None if no mapping needed/possible
- """
- # Get printer status
- status = printer_manager.get_status(printer_id)
- if not status:
- logger.warning("Cannot compute AMS mapping: printer %s status unavailable", printer_id)
- return None
- # Filament Track Switch (FTS): when installed it routes any AMS slot to
- # either extruder, so the per-nozzle hard filter below must NOT apply.
- # Otherwise a print on one nozzle can't use a spool physically loaded in
- # an AMS on the *other* nozzle, and the matcher falls through to a
- # same-type wrong-colour spool on the target nozzle — the H2C + FTS
- # wrong-filament bug (#2186). Mirrors the frontend skip added for #1162.
- fts_installed = bool(getattr(getattr(status, "fila_switch", None), "installed", False))
- # Get filament requirements from source file
- filament_reqs = await self._get_filament_requirements(db, item)
- if not filament_reqs:
- # When the 3MF can't be read but force-color overrides are present, build a
- # direct mapping from the overrides so the printer uses the correct AMS slot.
- if item.filament_overrides:
- try:
- overrides = json.loads(item.filament_overrides)
- force_overrides = [o for o in overrides if o.get("force_color_match")]
- if force_overrides:
- logger.info(
- "Queue item %s: No filament reqs from 3MF; building AMS mapping from %d "
- "force-color override(s)",
- item.id,
- len(force_overrides),
- )
- return self._build_override_direct_mapping(force_overrides, status)
- except (json.JSONDecodeError, KeyError, TypeError) as e:
- logger.warning("Queue item %s: Force-color fallback mapping failed: %s", item.id, e)
- logger.debug("No filament requirements found for queue item %s", item.id)
- return None
- # Apply filament overrides if present
- if item.filament_overrides:
- try:
- overrides = json.loads(item.filament_overrides)
- override_map = {o["slot_id"]: o for o in overrides}
- for req in filament_reqs:
- if req["slot_id"] in override_map:
- override = override_map[req["slot_id"]]
- req["type"] = override["type"]
- req["color"] = override["color"]
- # A manual/preference override SWAPS the slot's filament, so the
- # 3MF's original tray_info_idx now points at the old spool and must
- # be cleared — matching then falls back to type+colour. A
- # force_color_match override is not a swap: it carries the 3MF's
- # intended variant (Basic GFA00 / Matte GFA01 / Silk GFA06), so keep
- # it here too, letting the matcher pin the correct variant slot on a
- # printer holding two same-colour spools of different variants (#2650).
- # If that variant isn't loaded the matcher falls back to type+colour,
- # so an eligible printer never fails to map.
- req["tray_info_idx"] = (
- override.get("tray_info_idx", "") if override.get("force_color_match") else ""
- )
- logger.debug(
- "Queue item %s: Override slot %d -> %s %s",
- item.id,
- req["slot_id"],
- override["type"],
- override["color"],
- )
- except (json.JSONDecodeError, KeyError, TypeError) as e:
- logger.warning("Failed to apply filament overrides for queue item %s: %s", item.id, e)
- # Build loaded filaments from printer status
- loaded_filaments = self._build_loaded_filaments(status)
- if not loaded_filaments:
- logger.debug("No filaments loaded on printer %s", printer_id)
- return None
- # Check if user prefers lowest remaining filament when multiple spools match
- prefer_lowest = await self._get_bool_setting(db, "prefer_lowest_filament")
- # Gate prefer_lowest on the printer's AMS Filament Backup state (#1766).
- # Without backup, the printer will not switch to a second spool when the
- # picked one runs out — so sorting toward the lowest leaves the print
- # at risk of running dry mid-job. None (unknown / A1 family) preserves
- # today's behaviour intentionally.
- if prefer_lowest and status.ams_filament_backup is False:
- logger.info("[prefer-lowest] skipped (AMS Backup OFF on printer %s)", printer_id)
- prefer_lowest = False
- # When the preference is on, surface Bambuddy's inventory-side
- # remaining for each slot that's bound to a tracked spool, so the
- # sort beats the MQTT-only blind spot (#1508). Skip the lookup
- # entirely when the preference is off — no behaviour change for
- # users who haven't opted in.
- inventory_remain_overrides: dict[int, float] | None = None
- if prefer_lowest:
- inventory_remain_overrides = await self._build_inventory_remain_overrides(db, printer_id, loaded_filaments)
- # Compute mapping: match required filaments to available slots
- return self._match_filaments_to_slots(
- filament_reqs, loaded_filaments, prefer_lowest, inventory_remain_overrides, fts_installed
- )
- def _build_override_direct_mapping(self, force_overrides: list[dict], status) -> list[int] | None:
- """Build an AMS mapping directly from force-color overrides without a 3MF.
- Used when ``_get_filament_requirements`` returns nothing (e.g. the 3MF's
- slice_info is missing or unreadable) but ``force_color_match`` overrides
- are present. Each override's ``slot_id``, ``type``, and ``color`` are
- treated as the filament requirement for that slot and matched against the
- current AMS state of the printer.
- Returns the same format as ``_match_filaments_to_slots``, or None when
- the AMS has no loaded filaments.
- """
- loaded = self._build_loaded_filaments(status)
- if not loaded:
- return None
- reqs = [
- {
- "slot_id": o["slot_id"],
- "type": o.get("type", ""),
- "color": o.get("color", ""),
- # These are all force_color_match overrides, so the idx (when the
- # 3MF carried one) is the intended variant, not a stale swap —
- # keep it so the matcher pins the right variant slot, falling back
- # to type+colour when it isn't loaded (#2650).
- "tray_info_idx": o.get("tray_info_idx", ""),
- }
- for o in force_overrides
- ]
- return self._match_filaments_to_slots(reqs, loaded)
- async def _get_filament_requirements(self, db: AsyncSession, item: PrintQueueItem) -> list[dict] | None:
- """Resolve the queue item's source 3MF and parse the per-slot
- filament requirements out of it. Thin DB-resolver wrapper around
- ``filament_requirements.extract_filament_requirements`` so the VP
- queue-mode write path (#1188) can reuse the same parser at upload
- time.
- """
- from backend.app.services.filament_requirements import extract_filament_requirements
- file_path: Path | None = None
- if item.archive_id:
- result = await db.execute(select(PrintArchive).where(PrintArchive.id == item.archive_id))
- archive = result.scalar_one_or_none()
- if archive:
- file_path = settings.base_dir / archive.file_path
- elif item.library_file_id:
- result = await db.execute(LibraryFile.active().where(LibraryFile.id == item.library_file_id))
- library_file = result.scalar_one_or_none()
- if library_file:
- lib_path = Path(library_file.file_path)
- file_path = lib_path if lib_path.is_absolute() else settings.base_dir / library_file.file_path
- if not file_path or not file_path.exists():
- return None
- filaments = extract_filament_requirements(file_path, plate_id=item.plate_id)
- return filaments if filaments else None
- def _build_loaded_filaments(self, status) -> list[dict]:
- """Build list of loaded filaments from printer status.
- Args:
- status: PrinterState from printer_manager
- Returns:
- List of loaded filament dicts with type, color, ams_id, tray_id, global_tray_id
- """
- filaments = []
- # Get ams_extruder_map for dual-nozzle printers (H2D, H2D Pro)
- ams_extruder_map = status.raw_data.get("ams_extruder_map", {})
- # Parse AMS units from raw_data
- ams_data = status.raw_data.get("ams", [])
- for ams_unit in ams_data:
- ams_id = int(ams_unit.get("id", 0))
- trays = ams_unit.get("tray", [])
- is_ht = len(trays) == 1 # AMS-HT has single tray
- for tray in trays:
- tray_type = tray.get("tray_type")
- if tray_type:
- tray_id = int(tray.get("id", 0))
- tray_color = tray.get("tray_color", "")
- # tray_info_idx identifies the specific spool (e.g., "GFA00", "P4d64437")
- tray_info_idx = tray.get("tray_info_idx", "")
- # Normalize color: remove alpha, add hash
- color = self._normalize_color(tray_color)
- # Calculate global tray ID
- # AMS-HT units have IDs starting at 128 with a single tray
- global_tray_id = ams_id if ams_id >= 128 else ams_id * 4 + tray_id
- filaments.append(
- {
- "type": tray_type,
- "color": color,
- "tray_info_idx": tray_info_idx,
- "ams_id": ams_id,
- "tray_id": tray_id,
- "is_ht": is_ht,
- "is_external": False,
- "global_tray_id": global_tray_id,
- "extruder_id": ams_extruder_map.get(str(ams_id)),
- "remain": tray.get("remain", -1),
- }
- )
- # Check external spool(s) (vt_tray is a list)
- for idx, vt in enumerate(status.raw_data.get("vt_tray") or []):
- if vt.get("tray_type"):
- color = self._normalize_color(vt.get("tray_color", ""))
- tray_id = int(vt.get("id", 254))
- filaments.append(
- {
- "type": vt["tray_type"],
- "color": color,
- "tray_info_idx": vt.get("tray_info_idx", ""),
- "ams_id": -1,
- "tray_id": idx,
- "is_ht": False,
- "is_external": True,
- "global_tray_id": tray_id,
- "extruder_id": (255 - tray_id) if ams_extruder_map else None,
- "remain": vt.get("remain", -1),
- }
- )
- return filaments
- def _normalize_color(self, color: str | None) -> str:
- """Normalize color to #RRGGBB format."""
- if not color:
- return "#808080"
- hex_color = color.replace("#", "")[:6]
- return f"#{hex_color}"
- def _normalize_color_for_compare(self, color: str | None) -> str:
- """Normalize color for comparison (lowercase, no hash)."""
- if not color:
- return ""
- return color.replace("#", "").lower()[:6]
- def _colors_are_similar(self, color1: str | None, color2: str | None, threshold: int = 40) -> bool:
- """Check if two colors are visually similar within a threshold."""
- hex1 = self._normalize_color_for_compare(color1)
- hex2 = self._normalize_color_for_compare(color2)
- if not hex1 or not hex2 or len(hex1) < 6 or len(hex2) < 6:
- return False
- try:
- r1 = int(hex1[0:2], 16)
- g1 = int(hex1[2:4], 16)
- b1 = int(hex1[4:6], 16)
- r2 = int(hex2[0:2], 16)
- g2 = int(hex2[2:4], 16)
- b2 = int(hex2[4:6], 16)
- return abs(r1 - r2) <= threshold and abs(g1 - g2) <= threshold and abs(b1 - b2) <= threshold
- except ValueError:
- return False
- async def _build_inventory_remain_overrides(
- self, db: AsyncSession, printer_id: int, loaded: list[dict]
- ) -> dict[int, float]:
- """Return ``{global_tray_id: remaining_grams}`` for AMS slots the user
- has bound to an inventory spool — Bambuddy-side or Spoolman-side.
- The MQTT ``remain`` field on a tray is the printer firmware's
- RFID-decremented value, which has two limitations the "Prefer Lowest
- Remaining Filament" feature has been ignoring (#1508):
- - it's only meaningful for Bambu RFID spools; everything else reports
- ``-1`` (then clamped to a sentinel), so multiple non-RFID trays
- compare equal and the sort collapses to AMS-slot order — the user
- who's curating inventory weights gets the lower-slot pick instead
- of the lower-remaining pick;
- - even when set, it's the *printer's* counter, not Bambuddy's
- ``label_weight - weight_used`` (internal mode) or Spoolman's
- ``remaining_weight`` (Spoolman mode) — the two diverge any time the
- user re-spools, swaps cardboard, or runs a print outside Bambuddy.
- When the user has bound a spool to a slot, their own inventory
- tracking is authoritative; this helper surfaces that value so the
- sort can prefer it. Slots without a binding are absent from the
- returned map — the caller then falls back to MQTT ``remain`` for
- those, preserving the pre-#1508 behaviour for un-tracked spools.
- Returns an empty map on any failure (no inventory bindings, DB
- error, Spoolman unreachable). A best-effort lookup; "Prefer Lowest"
- is a preference, not a guarantee.
- """
- if not loaded:
- return {}
- # External / virtual-tray slots are tracked separately from AMS — skip
- # them so a VT-loaded spool doesn't accidentally inherit a tracked
- # AMS binding (the tables use ams_id 254/255 for VT, but the cross
- # match is fiddly and out of scope for this fix).
- tracked_slots = [(f["ams_id"], f["tray_id"], f["global_tray_id"]) for f in loaded if not f.get("is_external")]
- if not tracked_slots:
- return {}
- is_spoolman = await self._is_spoolman_mode(db)
- overrides: dict[int, float] = {}
- if is_spoolman:
- result = await db.execute(
- select(SpoolmanSlotAssignment).where(SpoolmanSlotAssignment.printer_id == printer_id)
- )
- assignments = list(result.scalars().all())
- by_slot = {(a.ams_id, a.tray_id): a.spoolman_spool_id for a in assignments}
- from backend.app.services.filament_deficit import _spoolman_remaining_grams
- for ams_id, tray_id, gtid in tracked_slots:
- spoolman_id = by_slot.get((ams_id, tray_id))
- if spoolman_id is None:
- continue
- grams = await _spoolman_remaining_grams(spoolman_id)
- if grams is not None:
- overrides[gtid] = grams
- return overrides
- # Internal inventory mode (default). selectinload matches the pattern
- # used elsewhere (inventory.py, spoolman.py routes) — a single query
- # plus an eager-loaded relationship rather than an explicit join, so
- # the row-attribute shape is exactly what those routes already rely on.
- result = await db.execute(
- select(SpoolAssignment)
- .options(selectinload(SpoolAssignment.spool))
- .where(SpoolAssignment.printer_id == printer_id)
- )
- assignments = list(result.scalars().all())
- by_slot = {(a.ams_id, a.tray_id): a.spool for a in assignments}
- for ams_id, tray_id, gtid in tracked_slots:
- spool = by_slot.get((ams_id, tray_id))
- if spool is None:
- continue
- label = float(spool.label_weight or 0)
- used = float(spool.weight_used or 0)
- overrides[gtid] = max(0.0, label - used)
- return overrides
- @staticmethod
- async def _is_spoolman_mode(db: AsyncSession) -> bool:
- """Mirror of ``filament_deficit._is_spoolman_mode`` — kept private
- here to avoid making this module import-dependent on that private
- helper's signature."""
- try:
- from backend.app.api.routes.settings import get_setting
- v = await get_setting(db, "spoolman_enabled")
- return bool(v) and v.lower() == "true"
- except Exception:
- return False
- @staticmethod
- def _slot_priority(ams_id: int | None, tray_id: int | None) -> int:
- """Deterministic slot-position tie-breaker for the prefer-lowest sort.
- Three bands, matched to the emission order in ``_build_loaded_filaments``
- so a tied sort produces the same physical-position order the pre-#1508
- stable sort did (preserves the regression-free baseline):
- - Regular AMS (``ams_id`` 0..7): ``ams_id * 4 + tray_id`` → 0..31
- - AMS-HT (``ams_id`` >= 128, single tray): ``1000 + (ams_id - 128) * 4``
- - External / VT (``ams_id`` < 0, or ``None``): ``10_000``
- Banding ensures regular AMS < AMS-HT < external on ties, regardless of
- what the raw ``ams_id`` happens to be (in particular, ``ams_id = -1``
- for VT must NOT sort to a negative number or it would beat AMS slot 0).
- """
- if ams_id is None or ams_id < 0:
- return 10_000
- if ams_id >= 128:
- return 1_000 + (ams_id - 128) * 4 + (tray_id or 0)
- return ams_id * 4 + (tray_id or 0)
- @staticmethod
- def _prefer_lowest_sort_key(f: dict, overrides: dict[int, float] | None) -> tuple[int, float, int]:
- """Sort key for the "Prefer Lowest Remaining Filament" preference.
- Two-tier ordering: inventory-tracked spools always sort BEFORE
- non-tracked spools (the user has told us they care about these
- specifically), then ascending by remaining within each tier, then
- ascending by AMS slot position as the deterministic tie-breaker.
- Tiers are flagged by the first tuple element (0 = inventory-tracked,
- 1 = MQTT-only / unknown). Cross-tier value comparisons never run
- because the tier flag dominates — which is what lets us mix grams
- (inventory) and percent (MQTT) without a unit conversion.
- Within the MQTT tier ``remain = -1`` (unknown) is mapped to 101 so
- spools the printer DOES know something about sort ahead of those
- it knows nothing about — preserves pre-#1508 behaviour for the
- no-inventory-binding case.
- Slot tie-breaker via ``_slot_priority`` so regular AMS < AMS-HT <
- external on ties, matching the legacy emission-order stable sort.
- """
- gtid = f.get("global_tray_id")
- slot_order = PrintScheduler._slot_priority(f.get("ams_id"), f.get("tray_id"))
- if overrides and gtid in overrides:
- return (0, overrides[gtid], slot_order)
- remain = f.get("remain", -1)
- return (1, float(remain) if remain is not None and remain >= 0 else 101.0, slot_order)
- def _match_filaments_to_slots(
- self,
- required: list[dict],
- loaded: list[dict],
- prefer_lowest: bool = False,
- inventory_remain_overrides: dict[int, float] | None = None,
- fts_installed: bool = False,
- ) -> list[int] | None:
- """Match required filaments to loaded filaments and build AMS mapping.
- Priority: unique tray_info_idx match > exact color match > similar color match > type-only match
- The tray_info_idx is a filament type identifier stored in the 3MF file when the user
- slices (e.g., "GFA00" for generic PLA, "P4d64437" for custom presets). If the same
- tray_info_idx appears in only ONE available tray, we use that tray. If multiple trays
- have the same tray_info_idx (e.g., two spools of generic PLA), we fall back to color
- matching among those trays.
- Args:
- required: List of required filaments with slot_id, type, color, tray_info_idx
- loaded: List of loaded filaments with type, color, tray_info_idx, global_tray_id
- Returns:
- AMS mapping array (position = slot_id - 1, value = global_tray_id or -1)
- """
- if not required:
- return None
- # Track used trays to avoid duplicate assignment
- used_tray_ids: set[int] = set()
- comparisons = []
- for req in required:
- req_type = (req.get("type") or "").upper()
- req_color = req.get("color", "")
- req_tray_info_idx = req.get("tray_info_idx", "")
- # Find best match: unique tray_info_idx > exact color > similar color > type-only
- idx_match = None
- exact_match = None
- similar_match = None
- type_only_match = None
- # Get available trays (not already used)
- available = [f for f in loaded if f["global_tray_id"] not in used_tray_ids]
- # Nozzle-aware filtering: restrict to trays on the correct nozzle.
- # Hard filter — cross-nozzle assignment causes print failures
- # ("position of left hotend is abnormal"), so never fall back.
- # Skipped when an FTS is installed: it routes any AMS slot to either
- # extruder, so restricting to one nozzle would wrongly exclude the
- # correct spool sitting in the other nozzle's AMS (#2186).
- req_nozzle_id = req.get("nozzle_id")
- if req_nozzle_id is not None and not fts_installed:
- available = [f for f in available if f.get("extruder_id") == req_nozzle_id]
- # Sort by remaining filament (ascending) so lowest-remain spool wins .find().
- # Inventory-tracked spools sort before MQTT-only ones (#1508); see
- # _prefer_lowest_sort_key for the full rationale.
- if prefer_lowest:
- available.sort(key=lambda f: self._prefer_lowest_sort_key(f, inventory_remain_overrides))
- # INFO-level decision trace for "Prefer Lowest Filament" #1766.
- # One line per filament req so a bug report can be diagnosed
- # without enabling debug logging: shows what the matcher saw
- # (req shape + sorted candidate trays with their remain values
- # and any inventory override that was applied). Mirrored by
- # the picked-match log at the bottom of the loop.
- logger.info(
- "[prefer-lowest] req slot=%s type=%r color=%r tii=%r nozzle=%s; available (sorted lowest-first): %s",
- req.get("slot_id"),
- req_type,
- req_color,
- req_tray_info_idx,
- req_nozzle_id,
- [
- {
- "gtid": f.get("global_tray_id"),
- "type": f.get("type"),
- "color": f.get("color"),
- "tii": f.get("tray_info_idx"),
- "remain": f.get("remain"),
- "inv_g": (
- inventory_remain_overrides.get(f.get("global_tray_id"))
- if inventory_remain_overrides
- else None
- ),
- }
- for f in available
- ],
- )
- # Check if tray_info_idx is unique among available trays
- if req_tray_info_idx:
- idx_matches = [f for f in available if f.get("tray_info_idx") == req_tray_info_idx]
- if len(idx_matches) == 1:
- # Unique tray_info_idx - use it as definitive match
- idx_match = idx_matches[0]
- logger.debug(
- f"Matched filament slot {req.get('slot_id')} by unique tray_info_idx={req_tray_info_idx} "
- f"-> tray {idx_match['global_tray_id']}"
- )
- elif len(idx_matches) > 1:
- # Multiple trays with same tray_info_idx - use color matching among them
- logger.debug(
- f"Non-unique tray_info_idx={req_tray_info_idx} found in {len(idx_matches)} trays, "
- f"using color matching among trays: {[f['global_tray_id'] for f in idx_matches]}"
- )
- if prefer_lowest:
- idx_matches.sort(key=lambda f: self._prefer_lowest_sort_key(f, inventory_remain_overrides))
- # Use color matching within this subset
- for f in idx_matches:
- f_color = f.get("color", "")
- if self._normalize_color_for_compare(f_color) == self._normalize_color_for_compare(req_color):
- if not exact_match:
- exact_match = f
- elif self._colors_are_similar(f_color, req_color):
- if not similar_match:
- similar_match = f
- elif not type_only_match:
- type_only_match = f
- # If no idx_match yet, do standard type/color matching on all available trays
- if not idx_match and not exact_match and not similar_match and not type_only_match:
- for f in available:
- f_type = (f.get("type") or "").upper()
- if _canonical_filament_type(f_type) != _canonical_filament_type(req_type):
- continue
- # Type matches - check color
- f_color = f.get("color", "")
- if self._normalize_color_for_compare(f_color) == self._normalize_color_for_compare(req_color):
- if not exact_match:
- exact_match = f
- elif self._colors_are_similar(f_color, req_color):
- if not similar_match:
- similar_match = f
- elif not type_only_match:
- type_only_match = f
- match = idx_match or exact_match or similar_match or type_only_match
- if match:
- used_tray_ids.add(match["global_tray_id"])
- comparisons.append({"slot_id": req.get("slot_id", 0), "global_tray_id": match["global_tray_id"]})
- else:
- comparisons.append({"slot_id": req.get("slot_id", 0), "global_tray_id": -1})
- if prefer_lowest:
- # Pair with the "available (sorted)" log above so the reporter
- # bundle shows BOTH what the matcher saw AND which match bucket
- # won — fast triage when "Prefer Lowest Filament" picks the
- # wrong slot (#1766).
- if match:
- bucket = (
- "idx"
- if idx_match is not None
- else "exact_color"
- if exact_match is not None
- else "similar_color"
- if similar_match is not None
- else "type_only"
- )
- logger.info(
- "[prefer-lowest] picked gtid=%s via %s for req slot=%s",
- match["global_tray_id"],
- bucket,
- req.get("slot_id"),
- )
- else:
- logger.info(
- "[prefer-lowest] NO MATCH for req slot=%s (type=%r color=%r tii=%r)",
- req.get("slot_id"),
- req_type,
- req_color,
- req_tray_info_idx,
- )
- # Build mapping array
- if not comparisons:
- return None
- max_slot_id = max(c["slot_id"] for c in comparisons)
- if max_slot_id <= 0:
- return None
- mapping = [-1] * max_slot_id
- for c in comparisons:
- slot_id = c["slot_id"]
- if slot_id and slot_id > 0:
- mapping[slot_id - 1] = c["global_tray_id"]
- return mapping
- def _mark_printer_dispatched(
- self,
- printer_id: int,
- pre_state: str | None,
- pre_subtask_id: str | None,
- ) -> None:
- """Record that a print command was just sent to ``printer_id``.
- Held until either the watchdog observes a state/subtask transition
- (success path) or the hard timeout expires. See ``_dispatch_holds``.
- """
- if not pre_state:
- # No pre_state means we can't detect a transition — fall back to a
- # pure time-based hold using empty string as a sentinel that won't
- # match any real printer state.
- pre_state = ""
- self._dispatch_holds[printer_id] = (time.monotonic(), pre_state, pre_subtask_id)
- def _release_dispatch_hold(self, printer_id: int) -> None:
- """Drop the dispatch hold for ``printer_id`` (called by the watchdog)."""
- self._dispatch_holds.pop(printer_id, None)
- def _printer_in_dispatch_hold(self, printer_id: int) -> bool:
- """True if ``printer_id`` is still inside its post-dispatch hold window.
- Returns False (and clears the hold) once any of these are true:
- - hard timeout (``_dispatch_max_hold``) has elapsed
- - the printer has transitioned out of pre_state and we're past the
- minimum cooldown
- - the printer's subtask_id has advanced past pre_subtask_id and we're
- past the minimum cooldown
- Otherwise the printer is held — caller should treat it as busy.
- """
- entry = self._dispatch_holds.get(printer_id)
- if not entry:
- return False
- started_at, pre_state, pre_subtask_id = entry
- elapsed = time.monotonic() - started_at
- if elapsed >= self._dispatch_max_hold:
- self._dispatch_holds.pop(printer_id, None)
- return False
- # Without a pre_state we can't detect a transition — fall back to the
- # min cooldown alone, then drop the hold.
- if not pre_state:
- if elapsed >= self._dispatch_min_cooldown:
- self._dispatch_holds.pop(printer_id, None)
- return False
- return True
- status = printer_manager.get_status(printer_id)
- current_state = getattr(status, "state", None) if status else None
- current_subtask_id = getattr(status, "subtask_id", None) if status else None
- transitioned = (current_state is not None and current_state != pre_state) or (
- pre_subtask_id is not None and current_subtask_id is not None and current_subtask_id != pre_subtask_id
- )
- if transitioned and elapsed >= self._dispatch_min_cooldown:
- self._dispatch_holds.pop(printer_id, None)
- return False
- return True
- def _is_printer_idle(self, printer_id: int, require_plate_clear: bool = True) -> bool:
- """Check if a printer is connected and idle."""
- if not printer_manager.is_connected(printer_id):
- logger.debug("Printer %d: not connected", printer_id)
- return False
- state = printer_manager.get_status(printer_id)
- if not state:
- logger.debug("Printer %d: no status available", printer_id)
- return False
- # Plate-clear gate: if the printer finished/failed a previous print and the user
- # hasn't acknowledged the plate was cleared, the queue must not dispatch the next
- # job — even if the printer currently reports IDLE. After Auto Off cycles the
- # printer, it boots back into IDLE with no memory of the previous finish; without
- # the persisted awaiting flag we'd bypass the confirmation prompt (#961).
- if require_plate_clear and printer_manager.is_awaiting_plate_clear(printer_id):
- logger.debug(
- "Printer %d: not idle — awaiting plate-clear acknowledgment (state=%s)",
- printer_id,
- state.state,
- )
- return False
- idle = state.state in ("IDLE", "FINISH", "FAILED")
- if not idle:
- logger.debug("Printer %d: not idle — state=%s", printer_id, state.state)
- return idle
- async def _get_setting(self, db: AsyncSession, key: str) -> str | None:
- """Read a setting value from the database."""
- result = await db.execute(select(Settings).where(Settings.key == key))
- setting = result.scalar_one_or_none()
- return setting.value if setting else None
- async def _get_bool_setting(self, db: AsyncSession, key: str, default: bool = False) -> bool:
- """Read a boolean setting from the database."""
- result = await db.execute(select(Settings).where(Settings.key == key))
- setting = result.scalar_one_or_none()
- if setting:
- return setting.value.lower() == "true"
- return default
- async def _get_int_setting(self, db: AsyncSession, key: str, default: int) -> int:
- """Read an int setting; falls back to default on missing/unparseable rows."""
- result = await db.execute(select(Settings).where(Settings.key == key))
- setting = result.scalar_one_or_none()
- if setting and setting.value:
- try:
- return int(setting.value)
- except ValueError:
- pass
- return default
- async def _get_drying_presets(self, db: AsyncSession) -> dict[str, dict[str, int]]:
- """Get drying presets (user-configured or built-in defaults)."""
- result = await db.execute(select(Settings).where(Settings.key == "drying_presets"))
- setting = result.scalar_one_or_none()
- if setting and setting.value:
- try:
- presets = json.loads(setting.value)
- if isinstance(presets, dict) and presets:
- return presets
- except json.JSONDecodeError:
- pass
- return self.DEFAULT_DRYING_PRESETS
- async def _get_humidity_thresholds(self, db: AsyncSession) -> dict[str, int]:
- """Per-filament humidity thresholds (#1605).
- Returns the user-configured overrides map keyed by normalized filament
- type (uppercase base, e.g. ``PLA``, ``ASA``) plus a ``default`` key for
- unknown / unmapped types. Empty / unset → empty dict, in which case
- callers fall back to ``ams_humidity_fair``.
- """
- result = await db.execute(select(Settings).where(Settings.key == "ams_humidity_thresholds"))
- setting = result.scalar_one_or_none()
- if not setting or not setting.value:
- return {}
- try:
- data = json.loads(setting.value)
- except json.JSONDecodeError:
- return {}
- if not isinstance(data, dict):
- return {}
- out: dict[str, int] = {}
- for key, value in data.items():
- try:
- out[str(key).upper() if key != "default" else "default"] = int(value)
- except (TypeError, ValueError):
- continue
- return out
- @staticmethod
- def resolve_humidity_threshold(trays: list[dict], thresholds: dict[str, int], fallback: int) -> int:
- """Resolve the effective humidity threshold for an AMS unit (#1605).
- For mixed filament types loaded into one AMS, returns the most
- restrictive (lowest) threshold across all loaded tray types — matches
- the conservative-params strategy already used for drying temp/hours.
- Empty / unloaded trays contribute no constraint. Unknown types use the
- ``default`` key, falling through to ``fallback`` (= ``ams_humidity_fair``)
- when no per-type map is configured at all.
- """
- default = thresholds.get("default", fallback)
- if not thresholds:
- return fallback
- candidates: list[int] = []
- for tray in trays:
- tray_type = str(tray.get("tray_type") or "").strip()
- if not tray_type:
- continue
- base_type = tray_type.split()[0].upper()
- candidates.append(thresholds.get(base_type, default))
- if not candidates:
- return default
- return min(candidates)
- def _get_conservative_drying_params(
- self, trays: list[dict], module_type: str, presets: dict[str, dict[str, int]]
- ) -> tuple[int, int, str] | None:
- """Get the most conservative drying params for mixed filament types in an AMS unit.
- Returns (temp, duration_hours, filament_type) or None if no drying-eligible filaments.
- """
- temp_key = module_type if module_type in ("n3f", "n3s") else "n3f"
- hours_key = f"{temp_key}_hours"
- min_temp = None
- max_hours = None
- filament_type = ""
- for tray in trays:
- tray_type = tray.get("tray_type", "")
- if not tray_type:
- continue
- # Normalize filament type for preset lookup (e.g., "PLA Basic" -> "PLA")
- base_type = tray_type.split()[0].upper()
- preset = presets.get(base_type)
- if not preset:
- continue
- temp = preset.get(temp_key, 55)
- hours = preset.get(hours_key, 12)
- # Conservative: lowest temp, longest duration
- if min_temp is None or temp < min_temp:
- min_temp = temp
- if max_hours is None or hours > max_hours:
- max_hours = hours
- if not filament_type:
- filament_type = base_type
- if min_temp is None:
- return None
- return (min_temp, max_hours or 12, filament_type)
- async def _check_auto_drying(
- self,
- db: AsyncSession,
- queue_items: list[PrintQueueItem],
- busy_printers: set[int],
- *,
- require_plate_clear: bool = True,
- ):
- """Start drying on idle printers based on humidity.
- Three modes (can all be enabled independently):
- - queue_drying_enabled: Dry between scheduled queue prints
- - ambient_drying_enabled: Dry any idle printer when humidity is high, regardless of queue
- - print_drying_enabled: Also evaluate printers that are currently printing,
- when model+firmware supports "Print While Drying" (gated by
- supports_drying_while_printing). Drying temperature is capped at
- max(40, preset_temp - 5) to protect spools mid-print.
- """
- queue_drying_enabled = await self._get_bool_setting(db, "queue_drying_enabled")
- ambient_drying_enabled = await self._get_bool_setting(db, "ambient_drying_enabled")
- print_drying_enabled = await self._get_bool_setting(db, "print_drying_enabled")
- if not queue_drying_enabled and not ambient_drying_enabled:
- # Stop active drying on all printers if both features disabled
- if self._drying_in_progress:
- for pid in list(self._drying_in_progress):
- logger.info("Auto-drying: printer %d — stopping, auto-drying disabled", pid)
- await self._stop_drying(pid)
- return
- # Update drying state from printer status (handles backend restart)
- self._sync_drying_state()
- # Find printers with scheduled items (for queue drying mode)
- printers_with_scheduled: set[int] = set()
- printers_with_items: set[int] = set()
- for item in queue_items:
- if item.printer_id:
- printers_with_items.add(item.printer_id)
- if item.scheduled_time and not item.manual_start:
- printers_with_scheduled.add(item.printer_id)
- # If only queue mode is on and no printers have scheduled items, stop drying
- # (but skip this short-circuit when print_drying_enabled is on — busy printers
- # may still be eligible for mid-print drying regardless of queue state).
- if not ambient_drying_enabled and not printers_with_scheduled and not print_drying_enabled:
- for pid in list(self._drying_in_progress):
- logger.info("Auto-drying: printer %d — stopping, no scheduled prints in queue", pid)
- await self._stop_drying(pid)
- return
- # Get humidity threshold (global fallback)
- result = await db.execute(select(Settings).where(Settings.key == "ams_humidity_fair"))
- setting = result.scalar_one_or_none()
- global_humidity_threshold = int(setting.value) if setting else 60
- # Per-filament humidity threshold overrides (#1605). Empty → fall back
- # to the global threshold for every AMS unit.
- per_type_thresholds = await self._get_humidity_thresholds(db)
- # Get drying presets
- presets = await self._get_drying_presets(db)
- # Determine if drying should be skipped for printers with pending items
- block_for_drying = await self._get_bool_setting(db, "queue_drying_block")
- # Get all active printers
- all_printers = await db.execute(select(Printer).where(Printer.is_active.is_(True)))
- for printer in all_printers.scalars():
- pid = printer.id
- # Resolve model+firmware up front — needed to decide whether this printer
- # qualifies for mid-print drying (busy printer on capable hardware).
- state = printer_manager.get_status(pid)
- if not state:
- logger.debug("Auto-drying: printer %d skipped — no state", pid)
- continue
- model = printer_manager.get_model(pid)
- firmware = state.firmware_version
- mid_print = (
- pid in busy_printers and print_drying_enabled and supports_drying_while_printing(model, firmware)
- )
- if pid in busy_printers and not mid_print:
- logger.debug("Auto-drying: printer %d skipped — busy", pid)
- continue
- if not mid_print:
- # In queue-only mode, only dry printers that have scheduled prints
- if not ambient_drying_enabled and pid not in printers_with_scheduled:
- if self._drying_in_progress.get(pid):
- logger.info("Auto-drying: printer %d — stopping, no scheduled prints for this printer", pid)
- await self._stop_drying(pid)
- logger.debug("Auto-drying: printer %d skipped — no scheduled prints", pid)
- continue
- # When block mode is on, don't START new drying on printers with pending items.
- # But allow already-drying printers through so humidity auto-stop logic still runs.
- if block_for_drying and pid in printers_with_items and not self._drying_in_progress.get(pid):
- logger.debug("Auto-drying: printer %d skipped — has pending items (block mode)", pid)
- continue
- if not printer_manager.is_connected(pid):
- logger.debug("Auto-drying: printer %d skipped — not connected", pid)
- continue
- if not mid_print and not self._is_printer_idle(pid, require_plate_clear):
- logger.debug("Auto-drying: printer %d skipped — not idle", pid)
- continue
- # Check drying capability. For mid-print path, supports_drying_while_printing
- # was already verified when computing mid_print above.
- if not mid_print and not supports_drying(model, firmware):
- logger.debug("Auto-drying: printer %d skipped — model %s does not support drying", pid, model)
- continue
- # Check each AMS unit from raw_data
- ams_list = state.raw_data.get("ams", [])
- logger.debug("Auto-drying: printer %d — checking %d AMS units", pid, len(ams_list))
- for ams_data in ams_list:
- module_type = str(ams_data.get("module_type") or "")
- ams_id = int(ams_data.get("id", 0))
- # Only n3f/n3s support drying
- if module_type not in ("n3f", "n3s"):
- logger.debug("Auto-drying: printer %d AMS %d skipped — module_type=%s", pid, ams_id, module_type)
- continue
- # Resolve per-filament humidity threshold for this AMS unit (#1605).
- # Most-restrictive of all loaded tray types; falls back to the
- # global threshold when no overrides are configured.
- trays = ams_data.get("tray", []) or []
- humidity_threshold = self.resolve_humidity_threshold(
- trays, per_type_thresholds, global_humidity_threshold
- )
- dry_time = int(ams_data.get("dry_time") or 0)
- # Read humidity — prefer humidity_raw (actual %) over humidity (index 1-5)
- humidity = None
- h_raw = ams_data.get("humidity_raw")
- if h_raw is not None:
- try:
- humidity = int(h_raw)
- except (ValueError, TypeError):
- pass
- if humidity is None:
- h_idx = ams_data.get("humidity")
- if h_idx is not None:
- try:
- humidity = int(h_idx)
- except (ValueError, TypeError):
- pass
- # Already drying — let it run to its configured duration (#1892).
- #
- # We deliberately do NOT stop drying from a humidity re-check here.
- # Relative humidity drops steeply in heated air, so the AMS sensor
- # reads ~15-20% within minutes of the dryer starting even while the
- # filament is still saturated. A humidity-based early-stop therefore
- # always fires at the minimum-time floor, truncating both user-started
- # manual cycles and Bambuddy's own preset-duration dries to ~30 min.
- # The firmware stops when the configured duration elapses; scheduling
- # stops (print takes priority, queue no longer needs drying) are
- # handled separately via _stop_drying().
- if dry_time > 0:
- if pid not in self._drying_in_progress:
- # Drying we didn't start (manual or from before restart) —
- # track it so scheduling stops still apply; never auto-stop it.
- self._drying_in_progress[pid] = time.monotonic()
- logger.debug(
- "Auto-drying: printer %d AMS %d — drying (%dm left, humidity %s%%), letting it run",
- pid,
- ams_id,
- dry_time,
- humidity,
- )
- continue
- # Humidity below threshold — no need to start drying
- if humidity is None or humidity <= humidity_threshold:
- logger.debug(
- "Auto-drying: printer %d AMS %d skipped — humidity %s <= threshold %d",
- pid,
- ams_id,
- humidity,
- humidity_threshold,
- )
- continue
- # Check cannot-dry reasons (power constraints etc.)
- sf_reasons = ams_data.get("dry_sf_reason", [])
- if sf_reasons:
- logger.debug(
- "Auto-drying: printer %d AMS %d skipped — cannot dry reasons: %s",
- pid,
- ams_id,
- sf_reasons,
- )
- continue
- # Get conservative drying params for mixed filaments
- params = self._get_conservative_drying_params(trays, module_type, presets)
- if not params:
- logger.debug(
- "Auto-drying: printer %d AMS %d skipped — no drying-eligible filaments in trays", pid, ams_id
- )
- continue
- temp, duration_hours, filament_type = params
- # Mid-print drying: cap drying temperature to protect spools (Bambu warns
- # "drying temperature must not exceed the filament's softening temperature"
- # for Print While Drying). Floor at 40 degC — below that the dryer is
- # ineffective and firmware will reject anyway.
- if mid_print:
- temp = max(40, temp - 5)
- # Start drying
- logger.info(
- "Auto-drying: printer %d AMS %d — humidity %d%% > threshold %d%%, "
- "starting %s drying at %d°C for %dh%s",
- pid,
- ams_id,
- humidity,
- humidity_threshold,
- filament_type,
- temp,
- duration_hours,
- " (mid-print)" if mid_print else "",
- )
- success = printer_manager.send_drying_command(
- pid, ams_id, temp, duration_hours, mode=1, filament=filament_type
- )
- if success:
- self._drying_in_progress[pid] = time.monotonic()
- def _sync_drying_state(self):
- """Sync in-memory drying state with actual printer status.
- Handles backend restart — if a printer is drying but we don't know about it,
- update our state. If we think it's drying but it's not, clear it.
- """
- to_remove = []
- for pid in self._drying_in_progress:
- state = printer_manager.get_status(pid)
- if not state:
- to_remove.append(pid)
- continue
- # Check if any AMS unit is still drying
- ams_list = state.raw_data.get("ams", [])
- any_drying = any(int(a.get("dry_time") or 0) > 0 for a in ams_list)
- if not any_drying:
- to_remove.append(pid)
- for pid in to_remove:
- self._drying_in_progress.pop(pid, None)
- async def _stop_drying(self, printer_id: int):
- """Stop all active drying on a printer (print takes priority)."""
- state = printer_manager.get_status(printer_id)
- if not state:
- self._drying_in_progress.pop(printer_id, None)
- return
- ams_list = state.raw_data.get("ams", [])
- for ams_data in ams_list:
- dry_time = int(ams_data.get("dry_time") or 0)
- if dry_time > 0:
- ams_id = int(ams_data.get("id", 0))
- logger.info(
- "Auto-drying: stopping drying on printer %d AMS %d — print takes priority",
- printer_id,
- ams_id,
- )
- printer_manager.send_drying_command(printer_id, ams_id, 0, 0, mode=0)
- self._drying_in_progress.pop(printer_id, None)
- async def _get_smart_plugs(self, db: AsyncSession, printer_id: int) -> list[SmartPlug]:
- """Get all smart plugs associated with a printer."""
- result = await db.execute(select(SmartPlug).where(SmartPlug.printer_id == printer_id))
- return list(result.scalars().all())
- @staticmethod
- def _pick_power_plug(auto_on_plugs: list[SmartPlug]) -> SmartPlug:
- """Pick the plug to power-cycle a printer back online with (#2629).
- Only a plug flagged ``controls_printer_power`` can actually bring the
- printer back; waiting for a boot on an accessory (filter fan, lights)
- just burns the power-on timeout and fails the dispatch. Falls back to
- the first plug when none is flagged, which is the pre-#2629 behaviour.
- Callers must pass a non-empty list.
- """
- for plug in auto_on_plugs:
- if plug.controls_printer_power:
- return plug
- return auto_on_plugs[0]
- # Bundled defaults for preheat_filament_targets (#1468). Values are the
- # chamber-temperature recommendations BambuStudio ships for the matching
- # filament profile; users can override via Settings → Workflow → Preheat
- # card. "default" applies when a loaded tray's normalised type isn't in
- # the map (rare — Bambu RFID-tagged spools always carry a known type).
- DEFAULT_PREHEAT_FILAMENT_TARGETS: dict[str, int] = {
- "PLA": 0,
- "PETG": 0,
- "PETG-CF": 40,
- "ABS": 45,
- "ASA": 45,
- "PA": 50,
- "PA-CF": 55,
- "PC": 50,
- "PC-FR": 50,
- "TPU": 0,
- "PVA": 0,
- "default": 0,
- }
- async def _get_preheat_filament_targets(self, db: AsyncSession) -> dict[str, int]:
- """Parse the user-configured filament→chamber-target map, falling back
- to DEFAULT_PREHEAT_FILAMENT_TARGETS on missing / malformed JSON. Keys
- are uppercased and the 'default' fallback is always present in the
- returned dict so the resolution loop can index it unconditionally."""
- raw = await self._get_setting(db, "preheat_filament_targets")
- if not raw:
- return dict(self.DEFAULT_PREHEAT_FILAMENT_TARGETS)
- try:
- parsed = json.loads(raw)
- if not isinstance(parsed, dict):
- raise ValueError("not an object")
- except (json.JSONDecodeError, ValueError) as exc:
- logger.warning("preheat_filament_targets unparseable, using defaults: %s", exc)
- return dict(self.DEFAULT_PREHEAT_FILAMENT_TARGETS)
- # Coerce values to int; drop unparseable rows so a stray string
- # doesn't crash the loop.
- out: dict[str, int] = {}
- for key, value in parsed.items():
- try:
- out[str(key).upper()] = int(value)
- except (TypeError, ValueError):
- continue
- if "DEFAULT" not in out:
- out["DEFAULT"] = self.DEFAULT_PREHEAT_FILAMENT_TARGETS["default"]
- return out
- @staticmethod
- def _normalize_filament_type(tray_type: str) -> str:
- """Reduce the printer's tray_type to a preset-lookup key. Mirrors the
- existing drying-preset normalisation (split-at-space, upper-case) so
- the two maps share vocabulary — "PLA Basic" → "PLA", "PA-CF" stays
- "PA-CF" (no space to split on)."""
- return tray_type.split()[0].upper() if tray_type else ""
- def _derive_chamber_target(
- self,
- printer: Printer,
- targets: dict[str, int],
- ) -> int:
- """Look up the chamber target for each loaded AMS tray and return the
- max. Returns 0 when no AMS data is available (e.g. external-spool
- prints) or when every loaded slot maps to 0 — the chamber phase then
- short-circuits in the main loop.
- Reads from `printer_manager.get_status(...).raw_data['ams']`, which is
- the same source the dispatcher uses for AMS slot mapping. Empty / RFID-
- less slots have empty `tray_type` and contribute nothing."""
- state = printer_manager.get_status(printer.id)
- if state is None:
- return 0
- ams_list = (state.raw_data or {}).get("ams") if state.raw_data else None
- # Older Bambu firmware nests AMS as {"ams": {"ams": [...]}} — try both.
- if isinstance(ams_list, dict):
- ams_list = ams_list.get("ams") or []
- if not isinstance(ams_list, list):
- return 0
- best = 0
- for ams in ams_list:
- for tray in (ams.get("tray") or []) if isinstance(ams, dict) else []:
- normalised = self._normalize_filament_type(tray.get("tray_type") or "")
- if not normalised:
- continue
- target = targets.get(normalised, targets.get("DEFAULT", 0))
- if target > best:
- best = target
- return best
- async def _preheat_and_soak(
- self,
- db: AsyncSession,
- item: PrintQueueItem,
- printer: Printer,
- archive: PrintArchive | None,
- ) -> None:
- """Run the per-printer preheat + heat-soak stage before FTP upload (#1468).
- Resolution order:
- 1. `item.preheat_override` — 'off' skips entirely; 'inherit' falls back
- to the global `preheat_enabled` setting; 'on' forces the stage on
- even if the global is off.
- 2. Chamber target — `item.preheat_chamber_target_override` if non-null;
- else max of `preheat_filament_targets[normalize(t.tray_type)]`
- across loaded AMS slots; else 0 (skips chamber phase, keeps bed
- phase + soak timer).
- 3. Three hardware tiers branch the wait loop:
- - Chamber heater (H2C/H2D/H2DPro/H2S/X2D/X1E via supports_chamber_heater):
- send M141 to the resolved target, then wait for the chamber sensor
- to reach it (or the max-wait timeout to elapse).
- - Chamber sensor only (X1C/P2S via supports_chamber_temp ∧ ¬supports_chamber_heater):
- no M141; the bed is the only heat source, so we wait for the chamber
- sensor to rise via bed radiation OR fall through on timeout.
- - No chamber sensor (P1S/P1P/A1/A1 Mini): no way to verify chamber
- temperature; the function just heats the bed and holds for the
- configured soak duration.
- The bed target comes from the archive's parsed metadata
- (`bed_temperature`); if missing the preheat stage logs and returns
- without dispatching anything, rather than guessing at a default that
- might wreck filament setup.
- Failures are logged but never re-raised — preheat is best-effort. A
- printer that goes offline mid-soak, a refused gcode command, or a
- missing temperature reading must not turn into a failed queue item; the
- normal upload + start path runs immediately after this method returns.
- """
- override = (getattr(item, "preheat_override", None) or "inherit").lower()
- if override == "off":
- return
- if override == "inherit":
- enabled = await self._get_bool_setting(db, "preheat_enabled", default=False)
- if not enabled:
- return
- # override == "on" forces the stage on regardless of the global setting.
- max_wait = await self._get_int_setting(db, "preheat_max_wait_seconds", default=900)
- soak_seconds = await self._get_int_setting(db, "preheat_soak_seconds", default=300)
- # Chamber target resolution:
- # 1. Explicit per-item override beats everything (user knows best).
- # 2. Otherwise derive from loaded AMS filament types via the per-
- # filament target map. PLA-only print derives 0 → chamber phase
- # auto-skips without the user touching anything.
- explicit_target = getattr(item, "preheat_chamber_target_override", None)
- if explicit_target is not None and explicit_target > 0:
- chamber_target = int(explicit_target)
- chamber_source = "item-override"
- elif explicit_target == 0:
- chamber_target = 0 # explicit 0 means "no chamber, even if filament wants it"
- chamber_source = "item-override-zero"
- else:
- targets = await self._get_preheat_filament_targets(db)
- chamber_target = self._derive_chamber_target(printer, targets)
- chamber_source = "filament-map"
- bed_target = int(archive.bed_temperature) if archive and archive.bed_temperature else 0
- if bed_target <= 0:
- logger.info(
- "Queue item %s: preheat skipped — archive has no bed_temperature metadata",
- item.id,
- )
- return
- client = printer_manager.get_client(printer.id)
- if client is None:
- logger.warning("Queue item %s: preheat skipped — printer client unavailable", item.id)
- return
- model = printer.model or ""
- has_heater = supports_chamber_heater(model)
- has_sensor = supports_chamber_temp(model)
- do_chamber = chamber_target > 0 and (has_heater or has_sensor)
- logger.info(
- "Queue item %s: preheat starting — bed=%d°C chamber_target=%d°C (source=%s override=%s "
- "model=%s has_heater=%s has_sensor=%s) max_wait=%ds soak=%ds",
- item.id,
- bed_target,
- chamber_target if do_chamber else 0,
- chamber_source,
- override,
- model,
- has_heater,
- has_sensor,
- max_wait,
- soak_seconds,
- )
- # Dispatch heaters. set_bed_temperature / set_chamber_temperature already
- # cache the target locally so the polling reads below see consistent
- # state (firmware MQTT echoes lag by ~1s).
- try:
- client.set_bed_temperature(bed_target)
- except Exception as exc:
- logger.warning("Queue item %s: preheat bed M140 failed: %s", item.id, exc)
- return
- # Airduct mode (#1468 follow-up). Models with the cooling/heating flap
- # (H2C/H2D/H2D Pro/H2S/X2D/P2S) keep the flap whatever the user last
- # left it on, regardless of M141. Default cooling actively vents the
- # chamber, so a `chamber_target > 0` print with the flap stuck in
- # cooling never converges — the heater fights the open exhaust. We
- # flip the flap BEFORE M141 to "heating" when the preheat wants
- # chamber heat, and back to "cooling" when it doesn't (PLA-only print
- # on an H2D that was previously running ABS would otherwise stay in
- # heating mode and overheat PLA). The current-state read keeps the
- # command idempotent — no MQTT chatter when the flap is already where
- # we want it.
- if supports_airduct(model):
- desired_airduct = "heating" if chamber_target > 0 else "cooling"
- desired_id = 1 if desired_airduct == "heating" else 0
- current_state = printer_manager.get_status(printer.id)
- current_airduct = getattr(current_state, "airduct_mode", None) if current_state else None
- if current_airduct != desired_id:
- try:
- client.set_airduct_mode(desired_airduct)
- except Exception as exc:
- logger.warning(
- "Queue item %s: preheat airduct %s mode failed: %s",
- item.id,
- desired_airduct,
- exc,
- )
- if do_chamber and has_heater:
- try:
- client.set_chamber_temperature(chamber_target)
- except Exception as exc:
- logger.warning("Queue item %s: preheat chamber M141 failed: %s", item.id, exc)
- # Release the pooled DB connection before the (potentially many-minute)
- # heat-soak wait below (#2572). Every setting this method needs is read
- # above; the wait/soak loop only polls printer_manager state and sleeps —
- # it never touches the DB. Without this the caller's transaction sat
- # "idle in transaction" for the whole soak, pinning one pooled connection
- # per preheating printer. expire_on_commit=False keeps item/printer
- # readable afterwards; there are no pending writes to lose here.
- await db.commit()
- # Wait for convergence. Bed warm-up is fast (~5 min from cold); chamber
- # via M141 takes a few minutes; chamber via bed radiation can take 20+.
- # Poll every 3s — frequent enough for responsive logging without
- # spamming the MQTT state stream. The "converged" predicate is:
- # bed reached target (within 2°C tolerance for floating-point + heater hysteresis),
- # AND
- # chamber phase satisfied (no chamber phase, no sensor, or sensor reached target).
- BED_TOLERANCE = 2.0
- CHAMBER_TOLERANCE = 2.0
- POLL_INTERVAL = 3.0
- deadline = asyncio.get_event_loop().time() + max_wait
- while True:
- state = printer_manager.get_status(printer.id)
- if state is None:
- logger.warning("Queue item %s: preheat lost state during wait", item.id)
- break
- temps = state.temperatures or {}
- bed_now = float(temps.get("bed", 0) or 0)
- chamber_now = float(temps.get("chamber", 0) or 0)
- bed_ok = bed_now >= bed_target - BED_TOLERANCE
- if not do_chamber:
- chamber_ok = True # phase disabled or model has neither sensor nor heater
- elif not has_sensor:
- chamber_ok = True # P1S etc — can't read, rely on soak timer only
- else:
- chamber_ok = chamber_now >= chamber_target - CHAMBER_TOLERANCE
- if bed_ok and chamber_ok:
- logger.info(
- "Queue item %s: preheat target reached (bed=%.1f chamber=%.1f) — entering soak",
- item.id,
- bed_now,
- chamber_now,
- )
- break
- if asyncio.get_event_loop().time() >= deadline:
- logger.info(
- "Queue item %s: preheat max_wait reached (bed=%.1f/%d chamber=%.1f/%d) — falling through to soak",
- item.id,
- bed_now,
- bed_target,
- chamber_now,
- chamber_target if do_chamber else 0,
- )
- break
- await asyncio.sleep(POLL_INTERVAL)
- if soak_seconds > 0:
- logger.info("Queue item %s: preheat soak — holding for %ds", item.id, soak_seconds)
- await asyncio.sleep(soak_seconds)
- logger.info("Queue item %s: preheat complete — proceeding to upload", item.id)
- async def _power_on_and_wait(self, plug: SmartPlug, printer_id: int, db: AsyncSession) -> bool:
- """Turn on smart plug and wait for printer to connect.
- Returns True if printer connected successfully within timeout.
- """
- # Get the appropriate service for the plug type (Tasmota or Home Assistant)
- service = await smart_plug_manager.get_service_for_plug(plug, db)
- # Check current plug state
- status = await service.get_status(plug)
- if not status.get("reachable"):
- logger.warning("Smart plug '%s' is not reachable", plug.name)
- return False
- # Turn on if not already on
- if status.get("state") != "ON":
- success = await service.turn_on(plug)
- if not success:
- logger.warning("Failed to turn on smart plug '%s'", plug.name)
- return False
- logger.info("Powered on smart plug '%s' for printer %s", plug.name, printer_id)
- # Get printer from database for connection
- result = await db.execute(select(Printer).where(Printer.id == printer_id))
- printer = result.scalar_one_or_none()
- if not printer:
- logger.error("Printer %s not found in database", printer_id)
- return False
- # Wait for printer to boot (give it some time before trying to connect)
- logger.info("Waiting 30s for printer %s to boot...", printer_id)
- await asyncio.sleep(30)
- # Try to connect to the printer periodically
- elapsed = 30 # Already waited 30s
- while elapsed < self._power_on_wait_time:
- # Try to connect
- logger.info("Attempting to connect to printer %s...", printer_id)
- try:
- connected = await printer_manager.connect_printer(printer)
- if connected:
- logger.info("Printer %s connected after %ss", printer_id, elapsed)
- # Give it a moment to stabilize and get status
- await asyncio.sleep(5)
- return True
- except Exception as e:
- logger.debug("Connection attempt failed: %s", e)
- await asyncio.sleep(self._power_on_check_interval)
- elapsed += self._power_on_check_interval
- logger.debug("Waiting for printer %s to connect... (%ss)", printer_id, elapsed)
- logger.warning("Printer %s did not connect within %ss after power on", printer_id, self._power_on_wait_time)
- return False
- async def _check_previous_success(self, db: AsyncSession, item: PrintQueueItem) -> bool:
- """Check if the previous print on this printer succeeded.
- A user-cancelled predecessor is treated as neutral — `cancelled` is a
- deliberate action, not a failure, so subsequent items should still
- dispatch (#1667). `skipped` is excluded from the lookback entirely:
- a skip isn't an actual print attempt, so it must not gate downstream
- items — counting it as a failed predecessor was the cascade bug that
- let a single cancellation block 18 items over 3 days for the reporter.
- Only `failed` and `aborted` — real print-attempt failures — block.
- Failures with `gate_acknowledged=True` (set by the per-printer Resume
- action — #1818) are also excluded from the lookback so the user can
- clear the gate after fixing the physical issue without having to
- re-queue every downstream job.
- """
- result = await db.execute(
- select(PrintQueueItem)
- .where(PrintQueueItem.printer_id == item.printer_id)
- .where(PrintQueueItem.id != item.id)
- .where(PrintQueueItem.status.in_(["completed", "failed", "cancelled", "aborted"]))
- .where(PrintQueueItem.gate_acknowledged == False) # noqa: E712
- .order_by(PrintQueueItem.completed_at.desc())
- .limit(1)
- )
- prev_item = result.scalar_one_or_none()
- # If no previous item, assume success (first in queue)
- if not prev_item:
- return True
- return prev_item.status in ("completed", "cancelled")
- async def _power_off_if_needed(self, db: AsyncSession, item: PrintQueueItem):
- """Schedule power-off if the queue item enabled auto_off_after.
- Delegates to the smart-plug manager so the off honours each plug's
- configured strategy (time delay or temperature threshold), is cancelled
- if the printer starts printing again, and never cuts power on a loaded
- print (#1890). Previously this hardcoded a 50°C / 600s cooldown wait and
- powered off on the timeout regardless of print state.
- """
- if not item.auto_off_after:
- return
- try:
- await smart_plug_manager.schedule_off_after_queue_job(item.printer_id, db)
- except Exception as e:
- logger.warning("Auto-off: Failed to schedule power-off for printer %s: %s", item.printer_id, e)
- async def _get_job_name(self, db: AsyncSession, item: PrintQueueItem) -> str:
- """Get a human-readable name for a queue item."""
- if item.archive_id:
- result = await db.execute(select(PrintArchive).where(PrintArchive.id == item.archive_id))
- archive = result.scalar_one_or_none()
- if archive:
- return archive.filename.replace(".gcode.3mf", "").replace(".3mf", "")
- if item.library_file_id:
- result = await db.execute(LibraryFile.active().where(LibraryFile.id == item.library_file_id))
- library_file = result.scalar_one_or_none()
- if library_file:
- return library_file.filename.replace(".gcode.3mf", "").replace(".3mf", "")
- return f"Job #{item.id}"
- async def _get_printer(self, db: AsyncSession, printer_id: int) -> Printer | None:
- """Get printer by ID."""
- result = await db.execute(select(Printer).where(Printer.id == printer_id))
- return result.scalar_one_or_none()
- async def _notify_dispatch_gave_up(
- self,
- queue_item_id: int,
- printer_id: int,
- created_by_id: int | None,
- reason: str = "Printer accepted the file but never started printing",
- ) -> None:
- """Tell the user the queue item was failed after exhausting its dispatch retries.
- Called from the watchdog, which is a background task with no session of
- its own — hence the fresh one here. Best-effort throughout: the row is
- already marked failed and that is the load-bearing part; a notification
- provider being down must not resurrect the retry loop we just stopped.
- ``reason`` defaults to the exhausted-retries wording. The command-rejected
- path passes its own, because "accepted the file but never started" is the
- opposite of what happened there — the printer refused it outright (#2732).
- """
- try:
- async with async_session() as db:
- item = await db.get(PrintQueueItem, queue_item_id)
- if not item:
- return
- job_name = await self._get_job_name(db, item)
- printer = await self._get_printer(db, printer_id)
- await notification_service.on_queue_job_failed(
- job_name=job_name,
- printer_id=printer_id,
- printer_name=printer.name if printer else "Unknown",
- reason=reason,
- db=db,
- )
- except Exception as e:
- logger.warning("Queue item %s: give-up notification failed: %s", queue_item_id, e)
- try:
- await ws_manager.send_queue_item_failed(
- user_id=created_by_id,
- queue_item_id=queue_item_id,
- printer_id=printer_id,
- reason="never_started",
- )
- except Exception:
- pass # toast is best-effort
- async def _block_on_filament_deficit(
- self,
- db: AsyncSession,
- item: PrintQueueItem,
- ) -> bool:
- """Promote the item to manual_start when the assigned spool is short (#1496).
- Returns True when this dispatch attempt was blocked, False when the
- item is clear to start. A previously-flagged item whose spool has
- since been swapped to one with enough material clears the flag here
- so the next scheduler tick dispatches it.
- """
- # User has explicitly acknowledged the deficit ("Print Anyway") —
- # don't re-flag, don't even compute. Without this short-circuit the
- # scheduler bounces between "user said anyway" (route clears
- # manual_start) and "scheduler re-blocked" (this method re-flags it
- # on identical spool state) (#1698-followup).
- if item.skip_filament_check:
- # #1762 diagnostic: surface the short-circuit at INFO so a
- # future "Print Anyway didn't work" report (e.g. issue #1762
- # comment 3) has actionable evidence in the support bundle
- # without needing DEBUG enabled.
- logger.info(
- "Queue item %s honouring user's Print Anyway acknowledgement — skipping deficit check",
- item.id,
- )
- return False
- try:
- deficit = await compute_deficit_for_queue_item(db, item)
- except Exception as e:
- # Never let a flaky deficit check wedge the queue — log and let
- # dispatch proceed. The PrintModal-side check still runs on the
- # manual paths.
- logger.warning("Filament deficit check failed for item %s: %s", item.id, e)
- return False
- if deficit:
- item.filament_short = True
- item.manual_start = True
- await db.commit()
- job_name = await self._get_job_name(db, item)
- printer = await self._get_printer(db, item.printer_id) if item.printer_id else None
- logger.info(
- "Queue item %s blocked on filament deficit (%d slot(s)) — promoted to manual_start",
- item.id,
- len(deficit),
- )
- try:
- await notification_service.on_queue_job_waiting(
- job_name=job_name,
- target_model=(printer.model if printer else "") or "",
- waiting_reason="filament_short",
- db=db,
- )
- except Exception as e:
- logger.debug("filament_short notification failed for item %s: %s", item.id, e)
- return True
- # No deficit — clear any stale flag from a previous tick.
- if item.filament_short:
- item.filament_short = False
- await db.commit()
- return False
- async def _propagate_owner_to_printer_manager(self, db: AsyncSession, item: PrintQueueItem) -> None:
- """Hand the queue item's owner to printer_manager so the
- print-complete callback can credit the user in PrintLogEntry (#1670).
- No-ops when the item has no `created_by_id` or the referenced user
- row is missing (e.g. user deleted between queue-add and dispatch —
- in that case the print log row falls back to the existing un-credited
- behaviour rather than crashing the dispatch).
- """
- if not item.created_by_id:
- return
- from backend.app.models.user import User
- owner = await db.get(User, item.created_by_id)
- if owner:
- printer_manager.set_current_print_user(item.printer_id, owner.id, owner.username)
- async def _start_print(self, db: AsyncSession, item: PrintQueueItem):
- """Upload file and start print for a queue item.
- Supports two sources:
- - archive_id: Print from an existing archive
- - library_file_id: Print from a library file (file manager)
- """
- logger.info("Starting queue item %s", item.id)
- # Get printer first (needed for both paths)
- result = await db.execute(select(Printer).where(Printer.id == item.printer_id))
- printer = result.scalar_one_or_none()
- if not printer:
- item.status = "failed"
- item.error_message = "Printer not found"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.error("Queue item %s: Printer %s not found", item.id, item.printer_id)
- await self._power_off_if_needed(db, item)
- return
- # Check printer is connected
- if not printer_manager.is_connected(item.printer_id):
- item.status = "failed"
- item.error_message = "Printer not connected"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.error("Queue item %s: Printer %s not connected", item.id, item.printer_id)
- await self._power_off_if_needed(db, item)
- return
- # Cancel-while-dispatching race (#1853): the scheduler's snapshot of
- # `items` was taken at the top of check_queue, but the user can /cancel
- # any pending row in the gap before we reach this point. Re-read the
- # row and bail out cleanly instead of starting an FTP upload for a row
- # that's already cancelled. The atomic CAS at the pending→printing
- # transition (below, before start_print) is the load-bearing guard;
- # this is the early-exit optimisation that avoids wasted FTP I/O.
- await db.refresh(item)
- if item.status != "pending":
- logger.info(
- "Queue item %s no longer pending (status=%s) — aborting dispatch",
- item.id,
- item.status,
- )
- return
- # Busy-printer guard (#2598). check_queue gates dispatch on
- # _is_printer_idle(), but that treats FINISH as idle and a printer can
- # keep reporting FINISH for tens of seconds *after* it accepted a
- # project_file (see the watchdog's phase-B note). A watchdog revert
- # (#2555) also releases the dispatch hold, so a re-selected item can
- # reach here while its printer has actually started printing. Uploading
- # and dispatching then collides with the live job — the firmware answers
- # 0500_4004 and, on an A1 mini, cancels the running print. Re-check the
- # live state right before the expensive FTP upload: if the printer is
- # busy, leave the item pending and let a later tick dispatch it once the
- # printer is genuinely idle. No wasted upload, no collision.
- pre_dispatch_state = getattr(printer_manager.get_status(item.printer_id), "state", None)
- if pre_dispatch_state in _ACTIVE_PRINT_STATES:
- logger.info(
- "Queue item %s: printer %s is busy (state=%s) — deferring dispatch, "
- "leaving item pending for a later tick (#2598)",
- item.id,
- item.printer_id,
- pre_dispatch_state,
- )
- return
- # Determine source: archive or library file
- archive = None
- library_file = None
- file_path = None
- filename = None
- cleanup_disk_paths: list[Path] = []
- if item.archive_id:
- # Print from archive
- result = await db.execute(select(PrintArchive).where(PrintArchive.id == item.archive_id))
- archive = result.scalar_one_or_none()
- if not archive:
- item.status = "failed"
- item.error_message = "Archive not found"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.error("Queue item %s: Archive %s not found", item.id, item.archive_id)
- await self._power_off_if_needed(db, item)
- return
- # Persist the queue item's selected plate onto the archive so Print
- # History can show the actual plate after cancel/fail/complete (#2603).
- # Only when the archive doesn't already carry one, so a reprint of a
- # plate-specific archive isn't relabelled by a differently-plated
- # queue row.
- if archive.plate_id is None and item.plate_id is not None:
- archive.plate_id = item.plate_id
- file_path = settings.base_dir / archive.file_path
- filename = archive.filename
- elif item.library_file_id:
- # Print from library file (file manager)
- result = await db.execute(LibraryFile.active().where(LibraryFile.id == item.library_file_id))
- library_file = result.scalar_one_or_none()
- if not library_file:
- item.status = "failed"
- item.error_message = "Library file not found"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.error("Queue item %s: Library file %s not found", item.id, item.library_file_id)
- await self._power_off_if_needed(db, item)
- return
- # Library files store absolute paths
- lib_path = Path(library_file.file_path)
- file_path = lib_path if lib_path.is_absolute() else settings.base_dir / library_file.file_path
- filename = library_file.filename
- # Create archive from library file so usage tracking has access to the 3MF
- queue_item_id = item.id
- try:
- from backend.app.services.archive import ArchiveService
- archive_service = ArchiveService(db)
- archive = await archive_service.archive_print(
- printer_id=item.printer_id,
- source_file=file_path,
- original_filename=filename,
- created_by_id=item.created_by_id,
- project_id=item.project_id,
- library_file_id=item.library_file_id, # per-file project progress (#1897)
- plate_id=item.plate_id, # selected plate → Print History (#2603)
- )
- if archive:
- item.archive_id = archive.id
- if item.cleanup_library_after_dispatch and not library_file.is_external:
- item.library_file_id = None
- cleanup_disk_paths.append(file_path)
- if library_file.thumbnail_path:
- thumb_path = Path(library_file.thumbnail_path)
- if not thumb_path.is_absolute():
- thumb_path = settings.base_dir / library_file.thumbnail_path
- cleanup_disk_paths.append(thumb_path)
- await db.delete(library_file)
- file_path = settings.base_dir / archive.file_path
- filename = archive.filename
- # Commit, not flush — flush opens the SQLite write
- # transaction (item.archive_id update + library_file
- # delete) and would hold the WAL writer lock through the
- # FTP upload below, causing "database is locked" cascades
- # for sensor history + concurrent cancels (#1853).
- await db.commit()
- logger.info(
- "Queue item %s: Created archive %s from library file %s",
- item.id,
- archive.id,
- item.library_file_id,
- )
- except Exception as e:
- logger.warning(
- "Queue item %s: Failed to create archive from library file: %s",
- queue_item_id,
- e,
- exc_info=True,
- )
- await db.rollback()
- item = await db.get(PrintQueueItem, queue_item_id)
- if item:
- item.status = "failed"
- item.error_message = "Failed to create archive from library file"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- await self._power_off_if_needed(db, item)
- return
- if not archive:
- item.status = "failed"
- item.error_message = "Failed to create archive from library file"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.error("Queue item %s: Archive creation from library file returned no archive", item.id)
- await self._power_off_if_needed(db, item)
- return
- else:
- # Neither archive nor library file specified
- item.status = "failed"
- item.error_message = "No source file specified"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.error("Queue item %s: No archive_id or library_file_id specified", item.id)
- await self._power_off_if_needed(db, item)
- return
- # Check file exists on disk
- if not file_path.exists():
- item.status = "failed"
- item.error_message = "Source file not found on disk"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.error("Queue item %s: File not found: %s", item.id, file_path)
- await self._power_off_if_needed(db, item)
- return
- # Nozzle-diameter mismatch guard (#1899). A file sliced for one nozzle
- # size dispatched to a printer with a different nozzle installed is
- # rejected by the firmware with a cryptic HMS ("Failed to get AMS mapping
- # table" 0700_8012, or "nozzle diameter … not consistent" 0500_4038) that
- # gives the user no idea what went wrong. Catch it here, before we spend
- # time preheating and uploading, and fail with an actionable message.
- # Fail-safe by construction: only a POSITIVE mismatch blocks — when the
- # slice carries no nozzle diameter (archive.nozzle_diameter is None) or
- # the printer hasn't reported its nozzles yet, we fall through and let the
- # print proceed exactly as before. On dual-nozzle printers (H2D) a match
- # against EITHER installed nozzle passes, so a 0.6 slice is fine as long
- # as one of the two hotends is a 0.6.
- sliced_nozzle = archive.nozzle_diameter if archive else None
- if sliced_nozzle:
- installed = _installed_nozzle_diameters(printer_manager.get_status(item.printer_id))
- mismatch_msg = _nozzle_mismatch_message(sliced_nozzle, installed)
- if mismatch_msg:
- item.status = "failed"
- item.error_message = mismatch_msg
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.warning("Queue item %s: nozzle mismatch — %s", item.id, mismatch_msg)
- await notification_service.on_queue_job_failed(
- job_name=filename.replace(".gcode.3mf", "").replace(".3mf", ""),
- printer_id=printer.id,
- printer_name=printer.name,
- reason=mismatch_msg,
- db=db,
- )
- try:
- await ws_manager.send_queue_item_failed(
- user_id=item.created_by_id,
- queue_item_id=item.id,
- printer_id=item.printer_id,
- reason="nozzle_mismatch",
- )
- except Exception:
- pass
- await self._power_off_if_needed(db, item)
- return
- # Preheat / heat-soak (#1468) — fires before upload so the printer's
- # bed (and chamber, if applicable) is at temperature when the firmware
- # starts the actual print routine. Best-effort: any failure logs and
- # falls through to the normal upload+start path rather than turning a
- # configuration issue into a failed queue item.
- await self._preheat_and_soak(db, item, printer, archive)
- # G-code injection for auto-print systems (#422)
- injected_path = None
- # #2547: tracked separately from `injected_path`, which is also set when
- # only a START snippet was injected. Only an END snippet changes what the
- # camera sees at print completion.
- end_gcode_injected = False
- if item.gcode_injection:
- try:
- snippets_raw = await self._get_setting(db, "gcode_snippets")
- if snippets_raw:
- snippets = json.loads(snippets_raw)
- model_snippets = snippets.get(printer.model, {})
- start_gc = (model_snippets.get("start_gcode") or "").strip()
- end_gc = (model_snippets.get("end_gcode") or "").strip()
- if start_gc or end_gc:
- from backend.app.utils.threemf_tools import inject_gcode_into_3mf
- injected_path = inject_gcode_into_3mf(
- file_path, item.plate_id or 1, start_gc or None, end_gc or None
- )
- if injected_path:
- file_path = injected_path
- end_gcode_injected = bool(end_gc)
- logger.info("Queue item %s: G-code injected for model %s", item.id, printer.model)
- else:
- logger.warning(
- "Queue item %s: G-code injection returned no result, using original", item.id
- )
- except Exception as e:
- logger.warning("Queue item %s: G-code injection failed, using original: %s", item.id, e)
- # #2547: the finish-photo path can't learn from telemetry that this print
- # ends with user End G-code — which means the plate may be gone by the
- # time FINISH arrives (#1867). Flag it here; `on_print_start` binds it to
- # the print once the printer confirms it running.
- if end_gcode_injected:
- print_dispatch_context.mark_pending(printer.id)
- # Upload to root directory (not /cache/) - the start_print command references
- # files by name only (ftp://{filename}), so they must be in the root
- remote_filename = derive_remote_filename(filename)
- remote_path = f"/{remote_filename}"
- # Get FTP retry settings
- ftp_retry_enabled, ftp_retry_count, ftp_retry_delay, ftp_timeout = await get_ftp_retry_settings()
- logger.info(
- f"Queue item {item.id}: FTP upload starting - printer={printer.name} ({printer.model}), "
- f"ip={printer.ip_address}, file={remote_filename}, local_path={file_path}, "
- f"retry_enabled={ftp_retry_enabled}, retry_count={ftp_retry_count}, timeout={ftp_timeout}"
- )
- # Release the pooled DB connection before the FTP delete/upload (#2572).
- # Every read this method needs (printer, archive/library, preheat) is
- # done, and the library-file branch already committed its archive
- # creation. Without this the transaction opened by the first SELECT above
- # stays "idle in transaction" for the entire upload — multiple seconds
- # for a large 3MF — pinning one pooled connection per in-flight dispatch;
- # a farm dispatching many jobs at once then exhausts the pool. This was
- # correlated to an exact idle-in-transaction session on a 93-printer farm
- # (reporter @Jostxxl). expire_on_commit=False keeps item/printer/archive
- # readable; the status writes below (upload-failure path and the
- # pending->printing CAS) transparently open a fresh transaction.
- await db.commit()
- # Delete existing file if present (avoids 553 error on overwrite)
- try:
- logger.debug("Queue item %s: Deleting existing file %s if present...", item.id, remote_path)
- delete_result = await delete_file_async(
- printer.ip_address,
- printer.access_code,
- remote_path,
- socket_timeout=ftp_timeout,
- printer_model=printer.model,
- )
- logger.debug("Queue item %s: Delete result: %s", item.id, delete_result)
- except Exception as e:
- logger.debug("Queue item %s: Delete failed (may not exist): %s", item.id, e)
- # Dispatch toast — announce the upload start with the total byte
- # count so the frontend can render an honest progress bar.
- toast_uid = item.created_by_id
- toast_file_name = filename.replace(".gcode.3mf", "").replace(".3mf", "")
- try:
- total_bytes = file_path.stat().st_size
- except OSError:
- total_bytes = 0
- try:
- await ws_manager.send_queue_item_uploading(
- user_id=toast_uid,
- queue_item_id=item.id,
- printer_id=item.printer_id,
- printer_name=printer.name,
- file_name=toast_file_name,
- total_bytes=total_bytes,
- )
- except Exception:
- pass # toast is best-effort
- progress_bridge = _UploadProgressBridge(toast_uid, item.id)
- # A deadline expiry gets its own message: "check your SD card" is the
- # wrong advice for a link that was simply too slow to finish (#2529).
- upload_error: str | None = None
- try:
- if ftp_retry_enabled:
- uploaded = await with_ftp_retry(
- upload_file_async,
- printer.ip_address,
- printer.access_code,
- file_path,
- remote_path,
- socket_timeout=ftp_timeout,
- printer_model=printer.model,
- progress_callback=progress_bridge,
- max_retries=ftp_retry_count,
- retry_delay=ftp_retry_delay,
- operation_name=f"Upload print to {printer.name}",
- )
- else:
- uploaded = await upload_file_async(
- printer.ip_address,
- printer.access_code,
- file_path,
- remote_path,
- socket_timeout=ftp_timeout,
- printer_model=printer.model,
- progress_callback=progress_bridge,
- )
- except UploadCancelled as e:
- uploaded = False
- upload_error = (
- "Upload was too slow to finish and was cancelled. The printer's connection could not sustain "
- "the transfer — check its Wi-Fi signal, or move it closer to the access point."
- )
- logger.error("Queue item %s: upload deadline exceeded: %s", item.id, e)
- except Exception as e:
- uploaded = False
- logger.error("Queue item %s: FTP error: %s (type: %s)", item.id, e, type(e).__name__)
- # Clean up injected temp file after upload attempt
- if injected_path and injected_path.exists():
- injected_path.unlink(missing_ok=True)
- if not uploaded:
- error_msg = upload_error or (
- "Failed to upload file to printer. Check if SD card is inserted and properly formatted (FAT32/exFAT). "
- "See server logs for detailed diagnostics."
- )
- item.status = "failed"
- item.error_message = error_msg
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.error(
- f"Queue item {item.id}: FTP upload failed - printer={printer.name}, model={printer.model}, "
- f"ip={printer.ip_address}. Check logs above for storage diagnostics and specific error codes."
- )
- # Send failure notification
- await notification_service.on_queue_job_failed(
- job_name=filename.replace(".gcode.3mf", "").replace(".3mf", ""),
- printer_id=printer.id,
- printer_name=printer.name,
- reason="Failed to upload file to printer",
- db=db,
- )
- try:
- await ws_manager.send_queue_item_failed(
- user_id=toast_uid,
- queue_item_id=item.id,
- printer_id=item.printer_id,
- reason="upload_failed",
- )
- except Exception:
- pass
- await self._power_off_if_needed(db, item)
- return
- # Parse AMS mapping if stored
- ams_mapping = None
- if item.ams_mapping:
- try:
- ams_mapping = json.loads(item.ams_mapping)
- except json.JSONDecodeError:
- logger.warning("Queue item %s: Invalid AMS mapping JSON, ignoring", item.id)
- # Register as expected print so we don't create a duplicate archive
- # Only applicable for archive-based prints
- if archive:
- from backend.app.main import register_expected_print
- register_expected_print(
- item.printer_id,
- remote_filename,
- archive.id,
- ams_mapping=ams_mapping,
- created_by_id=item.created_by_id,
- plate_id=item.plate_id,
- )
- # Registration happens before the print command by necessity (the
- # printer can report the print before the send returns), so record
- # what to undo if we never get as far as sending. `_dispatch_one`
- # rolls back anything still pending here on every exit — exception,
- # early return, or cancel winning the CAS below.
- self._unconfirmed_expected_print[item.id] = (item.printer_id, remote_filename, archive.id)
- # Propagate the queue item's owner into printer_manager so the
- # print-complete callback can credit the user in the PrintLogEntry
- # (#1670). `created_by_id` is set either at queue-add time (UI-added
- # items) or when the user clicks the manual-start button.
- await self._propagate_owner_to_printer_manager(db, item)
- # IMPORTANT: Set status to "printing" BEFORE sending the print command.
- # This prevents phantom reprints if the backend crashes/restarts after the
- # print command is sent but before the status update is committed.
- # If we crash after this commit but before start_print(), the item will be
- # in "printing" status without actually printing - but that's safer than
- # accidentally reprinting the same file hours later.
- #
- # Atomic CAS (#1853): a user pressing /cancel mid-dispatch (between the
- # initial pending read at the top of check_queue and this point) flips
- # the row to "cancelled" in a separate session. Without the WHERE
- # status='pending' clause, the unconditional update here would silently
- # overwrite that cancellation and we'd ship the MQTT start_print below
- # — printer obeys, user sees "I pressed cancel and the print started".
- # rowcount==0 means the user won the race; bail out, best-effort delete
- # the file we just uploaded, do NOT send start_print.
- now_utc = datetime.now(timezone.utc)
- cas = await db.execute(
- update(PrintQueueItem)
- .where(PrintQueueItem.id == item.id)
- .where(PrintQueueItem.status == "pending")
- .values(status="printing", started_at=now_utc)
- )
- await db.commit()
- if cas.rowcount == 0:
- logger.info(
- "Queue item %s no longer pending at print-command time "
- "(cancelled or removed mid-dispatch) — aborting before MQTT send (#1853)",
- item.id,
- )
- try:
- await delete_file_async(
- printer.ip_address,
- printer.access_code,
- remote_path,
- socket_timeout=ftp_timeout,
- printer_model=printer.model,
- )
- except Exception as cleanup_err:
- logger.debug(
- "Queue item %s: best-effort cleanup of uploaded file failed: %s",
- item.id,
- cleanup_err,
- )
- try:
- await ws_manager.send_queue_item_failed(
- user_id=toast_uid,
- queue_item_id=item.id,
- printer_id=item.printer_id,
- reason="cancelled_mid_dispatch",
- )
- except Exception:
- pass
- return
- # Sync the in-memory item so subsequent code that reads item.status /
- # item.started_at sees the values we just persisted.
- item.status = "printing"
- item.started_at = now_utc
- for cleanup_path in cleanup_disk_paths:
- try:
- if cleanup_path.exists():
- cleanup_path.unlink()
- except OSError as cleanup_err:
- logger.warning(
- "TRANSIENT_LIBRARY_FILE_ORPHAN %s",
- json.dumps(
- {
- "queue_item_id": item.id,
- "path": str(cleanup_path),
- "error": str(cleanup_err),
- },
- sort_keys=True,
- ),
- )
- # Clear the awaiting-plate-clear flag now that we're starting a new print
- printer_manager.set_awaiting_plate_clear(item.printer_id, False)
- logger.info("Queue item %s: Status set to 'printing', sending print command...", item.id)
- # Capture state before dispatch so the watchdog can detect whether the
- # printer actually transitioned (#967). Also capture subtask_id so the
- # watchdog can recognise "command landed but state hasn't flipped yet"
- # on slow H2D transitions (#1078).
- pre_status = printer_manager.get_status(item.printer_id)
- pre_state = getattr(pre_status, "state", None) if pre_status else None
- pre_subtask_id = getattr(pre_status, "subtask_id", None) if pre_status else None
- pre_gcode_file = getattr(pre_status, "gcode_file", None) if pre_status else None
- # #1721: respect the user's explicit timelapse choice. The #1397
- # force-on at dispatch was removed because it caused per-layer nozzle
- # parking on slicer profiles with Timelapse Type = Smooth. Finish-photo
- # capture is now driven by the stg_cur=22 transition in bambu_mqtt.py
- # ("Filament unloading", toolhead parked, bed not yet dropped) with a
- # FINISH-state fallback — no need to force a video.
- effective_timelapse = bool(item.timelapse)
- # Start the print with AMS mapping, plate_id and print options.
- # nozzle_mapping rides through verbatim — JSON string captured from
- # Bambu Studio's project_file on VP intake (#1780); the MQTT layer
- # parses + injects it only for dual-nozzle models so a null on every
- # other model is a transparent pass-through.
- started = printer_manager.start_print(
- item.printer_id,
- remote_filename,
- plate_id=item.plate_id or 1,
- ams_mapping=ams_mapping,
- bed_levelling=item.bed_levelling,
- flow_cali=item.flow_cali,
- vibration_cali=item.vibration_cali,
- layer_inspect=item.layer_inspect,
- timelapse=effective_timelapse,
- use_ams=item.use_ams,
- nozzle_offset_cali=item.nozzle_offset_cali,
- nozzle_mapping=item.nozzle_mapping,
- )
- if started:
- # The command is away, so the expectation is now legitimate and must
- # survive. Anything still in this dict when _dispatch_one exits gets
- # rolled back.
- self._unconfirmed_expected_print.pop(item.id, None)
- logger.info("Queue item %s: Print started successfully - %s", item.id, filename)
- # No dispatch-toast event here: the legacy bg-dispatch path kept
- # status='processing' from upload start until the printer acked
- # (or timed out). The frontend derives "Awaiting printer…" purely
- # from upload_progress_pct >= 99.9; an explicit 'dispatched' WS
- # event would push the status chip out of 'PROCESSING' prematurely
- # — which is exactly what the screenshot at #1625-followup
- # complained about.
- # Register the local 3MF in the cover-cache so /cover skips FTP
- # (#1166 follow-up). file_path was resolved earlier from either the
- # archive or the library file row.
- if file_path is not None:
- cache_3mf_download(item.printer_id, remote_filename, file_path)
- # Hold the printer against further dispatches until the watchdog
- # confirms the printer transitioned (or until the hard timeout).
- # Prevents multi-plate batches from triple-dispatching onto the
- # same H2D Pro while it digests the first project_file (#1157).
- self._mark_printer_dispatched(item.printer_id, pre_state, pre_subtask_id)
- # Watchdog: if the printer never transitions out of pre_state AND
- # never advances subtask_id, the MQTT publish was accepted locally but
- # didn't reach the printer (half-broken session — same shape as
- # #887/#936). Revert the queue item so the next dispatch can pick it
- # up instead of leaving it stuck in "printing" (#967). subtask_id
- # check avoids false reverts on slow H2D FINISH→PREPARE transitions
- # that would otherwise cause the item to re-dispatch as a reprint
- # of the just-finished job (#1078).
- if pre_state:
- spawn_background_task(
- self._watchdog_print_start(
- item.id,
- item.printer_id,
- pre_state,
- pre_subtask_id,
- pre_gcode_file,
- created_by_id=toast_uid,
- ),
- name=f"watchdog-print-start-{item.id}",
- )
- # Get estimated time for notification.
- #
- # This used to fall back to `library_file.print_time_seconds`, a column
- # LibraryFile does not have — the print time it knows about lives in
- # `file_metadata`. So a library print whose archive carried no parseable
- # print time (a plain .gcode, or a 3MF the parser could not read) raised
- # AttributeError right here, *after* the printer had already been sent
- # the job: the started-notification never fired, and the exception
- # unwound the whole queue pass, so every other printer still waiting to
- # be dispatched on that tick silently missed its turn.
- #
- # The queue item caches the print time at creation ("Cached from
- # archive/library"), which is the value this was reaching for.
- estimated_time = None
- if archive and archive.print_time_seconds:
- estimated_time = archive.print_time_seconds
- elif item.print_time_seconds:
- estimated_time = item.print_time_seconds
- # Send job started notification
- await notification_service.on_queue_job_started(
- job_name=filename.replace(".gcode.3mf", "").replace(".3mf", ""),
- printer_id=printer.id,
- printer_name=printer.name,
- db=db,
- estimated_time=estimated_time,
- )
- # MQTT relay - publish queue job started
- try:
- from backend.app.services.mqtt_relay import mqtt_relay
- await mqtt_relay.on_queue_job_started(
- job_id=item.id,
- filename=filename,
- printer_id=printer.id,
- printer_name=printer.name,
- printer_serial=printer.serial_number,
- )
- except Exception:
- pass # Don't fail if MQTT fails
- else:
- # Clean up uploaded file from SD card to prevent phantom prints
- try:
- await delete_file_async(
- printer.ip_address,
- printer.access_code,
- remote_path,
- printer_model=printer.model,
- )
- except Exception:
- pass # Best-effort — don't fail the error handler
- # Busy-refusal is a deferral, not a failure (#2598). The printer's
- # state can flip from idle to active in the window between the
- # pre-dispatch check above and this publish (the FTP upload takes
- # seconds); start_print() then refuses to send project_file to the
- # now-busy printer and returns False. Failing the item here would be
- # wrong — the printer is fine, it is simply busy — so revert to
- # pending and let a later tick dispatch it once the printer is idle,
- # exactly like the pre-dispatch guard. Only a start_print() False on
- # an idle/unknown printer is a genuine command failure.
- post_dispatch_state = getattr(printer_manager.get_status(item.printer_id), "state", None)
- if post_dispatch_state in _ACTIVE_PRINT_STATES:
- logger.info(
- "Queue item %s: printer %s became busy (state=%s) before the start "
- "command was sent — deferring, reverting item to pending (#2598)",
- item.id,
- item.printer_id,
- post_dispatch_state,
- )
- item.status = "pending"
- item.started_at = None
- await db.commit()
- return
- # Print command failed - revert status
- item.status = "failed"
- item.error_message = "Failed to send print command to printer"
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- logger.error(
- f"Queue item {item.id}: Failed to start print on {printer.name} ({printer.model}) - "
- f"printer_manager.start_print() returned False. "
- f"This may indicate: printer not connected, MQTT error, unsupported model configuration, or firmware issue. "
- f"Check printer status and backend logs for details."
- )
- # Send failure notification
- await notification_service.on_queue_job_failed(
- job_name=filename.replace(".gcode.3mf", "").replace(".3mf", ""),
- printer_id=printer.id,
- printer_name=printer.name,
- reason="Failed to send print command to printer - check printer connection and status",
- db=db,
- )
- try:
- await ws_manager.send_queue_item_failed(
- user_id=toast_uid,
- queue_item_id=item.id,
- printer_id=item.printer_id,
- reason="start_command_failed",
- )
- except Exception:
- pass
- await self._power_off_if_needed(db, item)
- @staticmethod
- async def _watchdog_print_start(
- queue_item_id: int,
- printer_id: int,
- pre_state: str,
- pre_subtask_id: str | None = None,
- pre_gcode_file: str | None = None,
- timeout: float = 90.0,
- phase_b_timeout: float = 180.0,
- poll_interval: float = 3.0,
- created_by_id: int | None = None,
- ) -> None:
- """Revert a queue item if the printer never acknowledges the start command.
- Bambuddy optimistically marks the queue item as "printing" right after the
- MQTT project_file publish succeeds locally. The watchdog runs in two phases:
- Phase A (up to ``timeout``): wait for either an active-state transition
- or a ``subtask_id`` advance past ``pre_subtask_id``. State alone is the
- primary signal; subtask_id advance handles the H2D case where state can
- sit at FINISH for ~50 s after the printer accepted ``project_file``
- before flipping to PREPARE (#1078). If neither happens, the MQTT publish
- was lost on a half-broken session (#887/#936) — revert and force
- reconnect (the #967 recovery path).
- Phase B (up to ``phase_b_timeout``, only if Phase A exited on subtask_id
- alone): keep watching for the active-state transition. subtask_id alone
- proves the file landed but not that the printer started — and a printer
- that accepts the command but stays at IDLE/FINISH indefinitely (e.g.
- cloud+LAN re-auth dance after a power cycle on old firmware, #1678)
- used to leave the queue item stuck in 'printing' forever because the
- old watchdog returned success as soon as subtask_id advanced. If Phase
- B times out, revert the queue item so the user can retry without
- restarting Bambuddy. Skip ``force_reconnect`` here: the file landed and
- a forced reconnect mid-parse triggers 0500_4003 (#1150).
- Phase A timeout raised from 45 s → 90 s as belt-and-braces for slow
- transitions that also don't emit an early subtask_id tick.
- Both phases also watch for ``HMS_MQTT_VERIFY_FAILED``. A printer that
- refuses to verify our commands will never start this job or any other,
- so waiting out the full 270 s and re-uploading the 3MF twice more only
- burns an upload slot the rest of the farm is queued behind — that path
- is for a printer that might still come good, which this one cannot
- (#2732). It fails the item on the spot with the actual reason instead.
- """
- last_status = None
- landed_on_subtask = False
- # Latched, not level-tested: state.hms_errors is rebuilt from scratch on
- # every push carrying an `hms` key, so the fault can come and go between
- # 3-second polls. Seeing it once inside the dispatch window is enough.
- command_rejected = False
- deadline = time.monotonic() + timeout
- while time.monotonic() < deadline:
- await asyncio.sleep(poll_interval)
- status = printer_manager.get_status(printer_id)
- if not status:
- # Printer disconnected — don't mess with the DB. Drop the
- # in-memory dispatch hold too so a fresh dispatch can retry
- # once the printer comes back; the hard timeout would
- # otherwise hold the printer unnecessarily.
- scheduler._release_dispatch_hold(printer_id)
- return
- last_status = status
- if status.state in _ACTIVE_PRINT_STATES:
- # Printer is actively processing the job — release the
- # post-dispatch hold so the next pending item for this printer
- # can be evaluated normally. We do NOT accept arbitrary state
- # transitions: a printer going FINISH -> IDLE (user dismissed
- # the post-print prompt without accepting our project_file)
- # would otherwise look like "command landed" and leave the
- # queue item stuck in 'printing' forever (#1370).
- scheduler._release_dispatch_hold(printer_id)
- try:
- await ws_manager.send_queue_item_acked(
- user_id=created_by_id,
- queue_item_id=queue_item_id,
- printer_id=printer_id,
- )
- except Exception:
- pass
- return
- # Checked only after the active-state exit above: a stale HMS left
- # over from an earlier job must never abort a print that is visibly
- # running. An actually-refused command leaves the printer idle, so
- # this ordering costs the detection nothing.
- if _mqtt_commands_rejected(status):
- command_rejected = True
- break
- if pre_subtask_id is not None and status.subtask_id is not None and status.subtask_id != pre_subtask_id:
- # Phase A exit — printer accepted the file (subtask_id flipped
- # to our submission id). Don't return yet: the printer may
- # have accepted the command but never actually start (e.g.
- # cloud+LAN re-auth dance after a power cycle, #1678). Phase
- # B watches for the active-state transition.
- landed_on_subtask = True
- break
- if landed_on_subtask and not command_rejected:
- phase_b_deadline = time.monotonic() + phase_b_timeout
- while time.monotonic() < phase_b_deadline:
- await asyncio.sleep(poll_interval)
- status = printer_manager.get_status(printer_id)
- if not status:
- scheduler._release_dispatch_hold(printer_id)
- return
- last_status = status
- if status.state in _ACTIVE_PRINT_STATES:
- scheduler._release_dispatch_hold(printer_id)
- try:
- await ws_manager.send_queue_item_acked(
- user_id=created_by_id,
- queue_item_id=queue_item_id,
- printer_id=printer_id,
- )
- except Exception:
- pass
- return
- # Same ordering rule as Phase A: a running print wins over a
- # lingering HMS.
- if _mqtt_commands_rejected(status):
- command_rejected = True
- break
- # No active-state transition. Revert the item so the scheduler can retry.
- # Drop the in-memory hold so the retry isn't blocked by it.
- scheduler._release_dispatch_hold(printer_id)
- # Four outcomes from the revert attempt, each routed differently:
- # "reverted": row flipped from printing -> pending, run recovery
- # "gave_up": same, but the retry budget is spent — row failed
- # rather than pending, so it stops going round again
- # "already_moved_on": item.status != 'printing' (completed/cancelled by
- # on_print_complete or user). Skip recovery entirely
- # — the print clearly landed somewhere even if the
- # watchdog didn't see the active-state transition.
- # "revert_failed": SQLite contention exhausted retries. Still run
- # recovery so the MQTT session gets a fresh client_id
- # on the half-broken-session path.
- #
- # The retry budget (#2555): reverting to 'pending' hands the item straight
- # back to the next queue pass, which re-uploads the whole 3MF and waits out
- # the watchdog again. For a printer that is genuinely wedged that loop never
- # ends — the reporter had one printer "since this morning still not launch"
- # — and each lap also consumes an upload slot that the other printers in the
- # farm are waiting on. Retrying is right; retrying forever is not.
- async def _do_revert(db):
- item = await db.get(PrintQueueItem, queue_item_id)
- if not item or item.status != "printing":
- return "already_moved_on"
- item.dispatch_attempts = (item.dispatch_attempts or 0) + 1
- item.started_at = None
- if command_rejected:
- # No retry budget for this one: the printer refused to verify the
- # command, and re-uploading the same 3MF to the same printer will
- # be refused the same way. Fail now with the fix rather than after
- # three laps of a message about SD cards (#2732).
- item.status = "failed"
- item.error_message = (
- "The printer rejected the print command: MQTT command verification failed "
- "(HMS 0500-0500-0001-0007). Enable Developer Mode on the printer, restart it, "
- "then start the job again."
- )
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- return "command_rejected"
- if item.dispatch_attempts >= DISPATCH_MAX_ATTEMPTS:
- item.status = "failed"
- item.error_message = (
- f"The printer accepted the file but never started printing, after "
- f"{item.dispatch_attempts} attempts. Check the printer's screen for a "
- f"prompt or error, confirm its SD card is readable, and start the job again."
- )
- item.completed_at = datetime.now(timezone.utc)
- await db.commit()
- return "gave_up"
- item.status = "pending"
- await db.commit()
- return "reverted"
- try:
- revert_outcome = await run_with_retry(_do_revert, label=f"watchdog revert item={queue_item_id}")
- except Exception as e:
- logger.warning(
- "Queue item %s: failed to revert to 'pending' (printer %d): %s — "
- "scheduler may keep treating this item as in-flight",
- queue_item_id,
- printer_id,
- e,
- )
- revert_outcome = "revert_failed"
- if revert_outcome == "already_moved_on":
- # Preserves the pre-#1370 early-return: if on_print_complete (or any
- # other path) already moved the item past 'printing', don't run the
- # MQTT session-recovery logic below — a forced reconnect on a healthy
- # session breaks ongoing prints on the same printer.
- return
- total_timeout = timeout + (phase_b_timeout if landed_on_subtask else 0.0)
- if revert_outcome == "command_rejected":
- logger.error(
- "Queue item %s: printer %d reported HMS %s (MQTT command verification "
- "failed) — the print command was rejected, not lost. Failing the item "
- "without retrying; enable Developer Mode on the printer and restart it (#2732)",
- queue_item_id,
- printer_id,
- HMS_MQTT_VERIFY_FAILED,
- )
- await scheduler._notify_dispatch_gave_up(
- queue_item_id,
- printer_id,
- created_by_id,
- reason="Printer rejected the print command (MQTT command verification failed)",
- )
- # Same reasoning as the landed_on_subtask path below: the file is on
- # the printer and a forced reconnect would only add 0500_4003 to a
- # problem that has nothing to do with the MQTT session (#1150).
- return
- if revert_outcome == "gave_up":
- logger.error(
- "Queue item %s: printer %d never started the print after %d dispatch "
- "attempts (last one waited %.0fs) — marking the item failed instead of "
- "re-uploading it again (#2555)",
- queue_item_id,
- printer_id,
- DISPATCH_MAX_ATTEMPTS,
- total_timeout,
- )
- await scheduler._notify_dispatch_gave_up(queue_item_id, printer_id, created_by_id)
- elif revert_outcome == "reverted":
- if landed_on_subtask:
- logger.warning(
- "Queue item %s: printer %d accepted project_file (subtask_id "
- "advanced) but never transitioned to an active state within "
- "%.0fs — printer wedged post-acceptance; reverted to 'pending' "
- "for retry (#1678)",
- queue_item_id,
- printer_id,
- total_timeout,
- )
- else:
- logger.warning(
- "Queue item %s: printer %d did not respond to print command within "
- "%.0fs (state still %s, subtask_id still %s) — reverted to 'pending' "
- "for retry (#967)",
- queue_item_id,
- printer_id,
- timeout,
- pre_state,
- pre_subtask_id,
- )
- # Phase B was entered iff subtask_id advanced, which means the
- # project_file landed on the printer. A forced reconnect at this point
- # would interrupt the printer's parse and trigger 0500_4003 (#1150) —
- # skip the recovery entirely.
- if landed_on_subtask:
- return
- # Phase A timeout path: if the printer's gcode_file changed since
- # pre-dispatch, the project_file command landed and the printer is
- # parsing — a forced reconnect mid-parse triggers 0500_4003 (#1150).
- # If gcode_file is unchanged, the publish was silently swallowed
- # (#887/#936) and force_reconnect recovery is what we want.
- client = printer_manager.get_client(printer_id)
- current_gcode_file = getattr(last_status, "gcode_file", None) if last_status else None
- publish_landed = current_gcode_file is not None and current_gcode_file != pre_gcode_file
- if publish_landed:
- logger.warning(
- "Queue item %s: gcode_file changed to %r (was %r) — printer "
- "received the command and is parsing slowly. Skipping forced "
- "MQTT reconnect to avoid 0500_4003 mid-parse (#1150).",
- queue_item_id,
- current_gcode_file,
- pre_gcode_file,
- )
- elif client and hasattr(client, "force_reconnect_stale_session"):
- client.force_reconnect_stale_session(
- f"queue print command unacknowledged after {timeout:.0f}s "
- f"(state still {pre_state}, gcode_file {current_gcode_file!r})"
- )
- # Global scheduler instance
- scheduler = PrintScheduler()
|