pipeline_runs.py 44 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054
  1. """API routes for Slicer Pipeline runs (#1425 PR B + PR C).
  2. PR B implemented single-target dispatch: one Run-pipeline click =
  3. slice the source once → enqueue ONE print on ``target_printer_id``.
  4. PR C extends this with:
  5. * ``copies > 1`` — slice once, enqueue N copies.
  6. * ``target_kind='printer_class'`` — pipeline targets a Bambu model code
  7. (X1C / P1S / H2D / …); orchestrator distributes copies across matching
  8. printers using the pipeline's ``fanout_strategy``.
  9. * Retry-failed runs that re-attempt only the failed/cancelled copies of
  10. a partial-failure run.
  11. * Dashboard list endpoint (``GET /pipeline-runs``) with status + pipeline
  12. filters and pagination.
  13. * WebSocket ``pipeline_run_updated`` events on state transitions so the
  14. dashboard refreshes live without polling.
  15. The slice itself runs through ``slice_dispatch`` (same path as the manual
  16. SliceModal), so the ``Slicing X — Generating G-code 75%`` toast renders
  17. end-to-end. The slice job's id rides on the run response so the frontend
  18. can call ``trackJob`` directly.
  19. """
  20. from __future__ import annotations
  21. import json
  22. import logging
  23. from datetime import datetime, timezone
  24. from pathlib import Path
  25. from typing import Literal
  26. from fastapi import APIRouter, Depends, HTTPException
  27. from sqlalchemy import delete, desc, func, select
  28. from sqlalchemy.ext.asyncio import AsyncSession
  29. from backend.app.api.routes.cloud import resolve_api_key_cloud_owner
  30. from backend.app.core.auth import QueueReviewRequired, RequestPrinterScope, RequirePermissionIfAuthEnabled
  31. from backend.app.core.config import settings as app_settings
  32. from backend.app.core.database import async_session, get_db
  33. from backend.app.core.permissions import Permission
  34. from backend.app.core.printer_scope import ALL_PRINTERS, PrinterScope, ensure_model_target_allowed
  35. from backend.app.core.websocket import ws_manager
  36. from backend.app.models.archive import PrintArchive
  37. from backend.app.models.library import LibraryFile
  38. from backend.app.models.pipeline_run import PipelineJob, PipelineRun
  39. from backend.app.models.print_queue import PrintQueueItem
  40. from backend.app.models.printer import Printer
  41. from backend.app.models.slicer_pipeline import SlicerPipeline
  42. from backend.app.models.user import User
  43. from backend.app.schemas.pipeline_run import (
  44. CheckEligibilityRequest,
  45. EligibilityIssueResponse,
  46. EligibilityReportResponse,
  47. PerPrinterReport as PerPrinterReportResponse,
  48. PipelineJobResponse,
  49. PipelineRunCreateRequest,
  50. PipelineRunListResponse,
  51. PipelineRunResponse,
  52. )
  53. from backend.app.schemas.slicer import PresetRef, SliceRequest
  54. from backend.app.services.pipeline_eligibility import (
  55. EligibilityReport,
  56. check_pipeline_eligibility,
  57. )
  58. from backend.app.services.print_confirmation import confirm_outcome_for_new_queue_item
  59. logger = logging.getLogger(__name__)
  60. pipeline_run_create_router = APIRouter(prefix="/slicer-pipelines", tags=["Slicer Pipelines"])
  61. pipeline_run_router = APIRouter(prefix="/pipeline-runs", tags=["Slicer Pipelines"])
  62. # ---------------------------------------------------------------------------
  63. # Helpers
  64. # ---------------------------------------------------------------------------
  65. def _serialise_status(report: EligibilityReport) -> EligibilityReportResponse:
  66. return EligibilityReportResponse(
  67. ok=report.ok,
  68. target_kind=report.target_kind,
  69. target_printer_id=report.target_printer_id,
  70. target_printer_name=report.target_printer_name,
  71. target_model_class=report.target_model_class,
  72. issues=[
  73. EligibilityIssueResponse(
  74. kind=issue.kind,
  75. slot_index=issue.slot_index,
  76. expected=issue.expected,
  77. actual=issue.actual,
  78. )
  79. for issue in report.issues
  80. ],
  81. printer_reports=[
  82. PerPrinterReportResponse(
  83. printer_id=r.printer_id,
  84. printer_name=r.printer_name,
  85. ok=r.ok,
  86. issues=[
  87. EligibilityIssueResponse(
  88. kind=i.kind,
  89. slot_index=i.slot_index,
  90. expected=i.expected,
  91. actual=i.actual,
  92. )
  93. for i in r.issues
  94. ],
  95. )
  96. for r in report.printer_reports
  97. ],
  98. )
  99. async def _load_pipeline(db: AsyncSession, pipeline_id: int) -> SlicerPipeline:
  100. pipeline = (
  101. await db.execute(
  102. select(SlicerPipeline).where(
  103. SlicerPipeline.id == pipeline_id,
  104. SlicerPipeline.is_deleted.is_(False),
  105. )
  106. )
  107. ).scalar_one_or_none()
  108. if pipeline is None:
  109. raise HTTPException(404, "Pipeline not found")
  110. return pipeline
  111. async def _load_printer_status(printer_id: int | None) -> dict | None:
  112. """Snapshot the printer_manager's live PrinterState for the eligibility
  113. matcher. Returns ``None`` when the printer has no MQTT client."""
  114. if printer_id is None:
  115. return None
  116. from backend.app.services.printer_manager import printer_manager
  117. state = printer_manager.get_status(printer_id)
  118. if state is None:
  119. return None
  120. return {"connected": state.connected, "raw_data": state.raw_data}
  121. def _make_status_lookup():
  122. """Closure that snapshots the printer_manager once per printer_id call.
  123. Passed to the matcher's class-targeting branch so it can read live state
  124. for every candidate printer."""
  125. def _lookup(printer_id: int) -> dict | None:
  126. from backend.app.services.printer_manager import printer_manager
  127. state = printer_manager.get_status(printer_id)
  128. if state is None:
  129. return None
  130. return {"connected": state.connected, "raw_data": state.raw_data}
  131. return _lookup
  132. def _slice_request_from_pipeline(pipeline: SlicerPipeline) -> SliceRequest:
  133. try:
  134. raw_filaments = json.loads(pipeline.filament_presets_json or "[]")
  135. except (json.JSONDecodeError, TypeError):
  136. raw_filaments = []
  137. filament_presets = [
  138. PresetRef(source=r["source"], id=r["id"])
  139. for r in raw_filaments
  140. if isinstance(r, dict) and "source" in r and "id" in r
  141. ]
  142. return SliceRequest(
  143. printer_preset=PresetRef(source=pipeline.printer_preset_source, id=pipeline.printer_preset_id),
  144. process_preset=PresetRef(source=pipeline.process_preset_source, id=pipeline.process_preset_id),
  145. filament_presets=filament_presets,
  146. bed_type=pipeline.bed_type,
  147. export_3mf=True,
  148. )
  149. def _compute_job_status(
  150. persisted: str,
  151. queue_entry: PrintQueueItem | None,
  152. ) -> str:
  153. if persisted in ("failed", "cancelled", "completed"):
  154. return persisted
  155. if queue_entry is None:
  156. return persisted
  157. qs = queue_entry.status
  158. if qs == "completed":
  159. return "completed"
  160. if qs in ("failed", "aborted"):
  161. return "failed"
  162. if qs == "cancelled":
  163. return "cancelled"
  164. if qs == "printing":
  165. return "printing"
  166. return "queued"
  167. def _roll_up_run_status(
  168. persisted: str,
  169. job_statuses: list[str],
  170. ) -> str:
  171. """Compute the run-level status from the per-job statuses.
  172. Terminal-persisted always wins for explicit cancels / hard failures so
  173. the dashboard doesn't flicker when one job's queue entry hasn't caught
  174. up. Otherwise:
  175. - all completed → completed
  176. - any in_progress / printing / queued / dispatching → in_progress
  177. - any failed alongside any completed → partial_failure
  178. - all failed/cancelled → failed
  179. """
  180. if persisted in ("cancelled",):
  181. return persisted
  182. if not job_statuses:
  183. return persisted
  184. completed = sum(1 for s in job_statuses if s == "completed")
  185. failed = sum(1 for s in job_statuses if s == "failed")
  186. cancelled = sum(1 for s in job_statuses if s == "cancelled")
  187. in_flight = sum(1 for s in job_statuses if s in ("printing", "queued", "awaiting_printer", "pending"))
  188. total = len(job_statuses)
  189. if completed == total:
  190. return "completed"
  191. if in_flight > 0:
  192. return "in_progress" if persisted not in ("queued", "slicing", "dispatching") else persisted
  193. # All copies are in terminal states.
  194. if failed == 0 and cancelled == total:
  195. return "cancelled"
  196. if completed > 0 and (failed > 0 or cancelled > 0):
  197. return "partial_failure"
  198. if failed > 0:
  199. return "failed"
  200. return persisted
  201. async def _materialise_run(db: AsyncSession, run: PipelineRun) -> PipelineRunResponse:
  202. pipeline_name: str | None = None
  203. target_kind = None
  204. target_printer_id = None
  205. target_model_class = None
  206. fanout_strategy = None
  207. if run.pipeline_id:
  208. pipeline = (
  209. await db.execute(select(SlicerPipeline).where(SlicerPipeline.id == run.pipeline_id))
  210. ).scalar_one_or_none()
  211. if pipeline:
  212. pipeline_name = pipeline.name
  213. target_kind = pipeline.target_kind # type: ignore[assignment]
  214. target_printer_id = pipeline.target_printer_id
  215. target_model_class = pipeline.target_model_class
  216. fanout_strategy = pipeline.fanout_strategy # type: ignore[assignment]
  217. source_filename: str | None = None
  218. if run.source_library_file_id:
  219. src = (
  220. await db.execute(select(LibraryFile).where(LibraryFile.id == run.source_library_file_id))
  221. ).scalar_one_or_none()
  222. source_filename = src.filename if src else None
  223. elif run.source_archive_id:
  224. arc = (
  225. await db.execute(select(PrintArchive).where(PrintArchive.id == run.source_archive_id))
  226. ).scalar_one_or_none()
  227. source_filename = (arc.print_name or arc.filename) if arc else None
  228. job_rows = (
  229. (
  230. await db.execute(
  231. select(PipelineJob).where(PipelineJob.pipeline_run_id == run.id).order_by(PipelineJob.copy_index)
  232. )
  233. )
  234. .scalars()
  235. .all()
  236. )
  237. job_responses: list[PipelineJobResponse] = []
  238. job_live_statuses: list[str] = []
  239. for job in job_rows:
  240. queue_entry = None
  241. if job.queue_entry_id:
  242. queue_entry = (
  243. await db.execute(select(PrintQueueItem).where(PrintQueueItem.id == job.queue_entry_id))
  244. ).scalar_one_or_none()
  245. printer_name: str | None = None
  246. if job.assigned_printer_id:
  247. p = (await db.execute(select(Printer).where(Printer.id == job.assigned_printer_id))).scalar_one_or_none()
  248. printer_name = p.name if p else None
  249. live_job_status = _compute_job_status(job.status, queue_entry)
  250. # If the job WAS dispatched (had a queue_entry_id) but the entry has
  251. # since been deleted from the queue page, the user's intent was
  252. # cancellation. Otherwise the run would stay forever showing as
  253. # ``queued`` because the persisted job.status hasn't been updated.
  254. if (
  255. job.queue_entry_id is not None
  256. and queue_entry is None
  257. and live_job_status not in ("completed", "failed", "cancelled")
  258. ):
  259. live_job_status = "cancelled"
  260. job_live_statuses.append(live_job_status)
  261. job_responses.append(
  262. PipelineJobResponse(
  263. id=job.id,
  264. pipeline_run_id=job.pipeline_run_id,
  265. copy_index=job.copy_index,
  266. assigned_printer_id=job.assigned_printer_id,
  267. assigned_printer_name=printer_name,
  268. queue_entry_id=job.queue_entry_id,
  269. status=live_job_status, # type: ignore[arg-type]
  270. error_message=job.error_message,
  271. dispatched_at=job.dispatched_at,
  272. completed_at=job.completed_at,
  273. )
  274. )
  275. rolled_up = _roll_up_run_status(run.status, job_live_statuses)
  276. return PipelineRunResponse(
  277. id=run.id,
  278. pipeline_id=run.pipeline_id,
  279. pipeline_name=pipeline_name,
  280. source_library_file_id=run.source_library_file_id,
  281. source_archive_id=run.source_archive_id,
  282. source_filename=source_filename,
  283. parent_run_id=run.parent_run_id,
  284. copies=run.copies,
  285. copies_completed=sum(1 for s in job_live_statuses if s == "completed"),
  286. copies_failed=sum(1 for s in job_live_statuses if s == "failed"),
  287. copies_cancelled=sum(1 for s in job_live_statuses if s == "cancelled"),
  288. copies_in_progress=sum(
  289. 1 for s in job_live_statuses if s in ("printing", "queued", "awaiting_printer", "pending")
  290. ),
  291. status=rolled_up, # type: ignore[arg-type]
  292. slice_job_id=run.slice_job_id,
  293. sliced_library_file_id=run.sliced_library_file_id,
  294. eligibility_overridden=run.eligibility_overridden,
  295. error_message=run.error_message,
  296. created_by=run.created_by,
  297. created_at=run.created_at,
  298. started_at=run.started_at,
  299. completed_at=run.completed_at,
  300. jobs=job_responses,
  301. target_kind=target_kind,
  302. target_printer_id=target_printer_id,
  303. target_model_class=target_model_class,
  304. fanout_strategy=fanout_strategy,
  305. )
  306. async def _publish_run_event(db: AsyncSession, run: PipelineRun) -> None:
  307. """Broadcast a ``pipeline_run_updated`` event with the full materialised
  308. run. Per-user routing via ``broadcast_to_user`` falls back to a global
  309. broadcast when ``created_by`` is None (auth-disabled installs)."""
  310. try:
  311. payload = await _materialise_run(db, run)
  312. await ws_manager.broadcast_to_user(
  313. run.created_by,
  314. {
  315. "type": "pipeline_run_updated",
  316. "run": payload.model_dump(mode="json"),
  317. },
  318. )
  319. except Exception:
  320. logger.exception("Failed to broadcast pipeline_run_updated for run %d", run.id)
  321. # ---------------------------------------------------------------------------
  322. # Source resolution + orchestration
  323. # ---------------------------------------------------------------------------
  324. SourceKind = Literal["library_file", "archive"]
  325. async def _resolve_source(
  326. db: AsyncSession,
  327. *,
  328. library_file_id: int | None,
  329. archive_id: int | None,
  330. user: User | None,
  331. printer_scope: PrinterScope,
  332. ) -> tuple[SourceKind, int, str, Path]:
  333. # Per-row ownership gate (IDOR fix): a caller may only run a pipeline on a
  334. # source they can see. Without this a READ_OWN caller could reference
  335. # another user's library file / archive by raw id and have it sliced (and,
  336. # via /run, printed) even though a direct GET on that id returned 404.
  337. # Auth-disabled and API-key callers (user is None) keep can_read_all=True —
  338. # no per-row identity, matching the library/archive read helpers.
  339. from backend.app.api.routes.archives import _ensure_archive_visible
  340. from backend.app.api.routes.library import _ensure_library_file_visible
  341. if library_file_id is not None:
  342. lib = (await db.execute(select(LibraryFile).where(LibraryFile.id == library_file_id))).scalar_one_or_none()
  343. can_read_all = user is None or user.has_permission(Permission.LIBRARY_READ_ALL.value)
  344. lib = _ensure_library_file_visible(lib, user, can_read_all)
  345. src_path = (
  346. Path(app_settings.base_dir) / lib.file_path
  347. ) # SEC-PATH-OK: lib.file_path is a LibraryFile DB column set only by the upload route, which writes a UUID-named file under base_dir/library_files/.
  348. if not src_path.exists():
  349. raise HTTPException(404, "Source library file missing on disk")
  350. return ("library_file", lib.id, lib.filename, src_path)
  351. assert archive_id is not None
  352. arc = (await db.execute(select(PrintArchive).where(PrintArchive.id == archive_id))).scalar_one_or_none()
  353. can_read_all = user is None or user.has_permission(Permission.ARCHIVES_READ_ALL.value)
  354. arc = _ensure_archive_visible(arc, user, can_read_all, printer_scope)
  355. rel = arc.source_3mf_path or arc.file_path
  356. if not rel:
  357. raise HTTPException(400, "Archive has no source file to slice")
  358. src_path = (
  359. Path(app_settings.base_dir) / rel
  360. ) # SEC-PATH-OK: rel is archive.source_3mf_path / archive.file_path, both set by upload-time validators that already do resolve+relative_to containment.
  361. if not src_path.exists():
  362. raise HTTPException(404, "Archive source file missing on disk")
  363. name = arc.filename or arc.print_name or src_path.name
  364. return ("archive", arc.id, name, src_path)
  365. async def _pick_assignments(
  366. db: AsyncSession,
  367. pipeline: SlicerPipeline,
  368. copies: int,
  369. printer_scope: PrinterScope = ALL_PRINTERS,
  370. ) -> list[tuple[int | None, str | None]]:
  371. """Return ``[(printer_id_or_None, target_model_or_None), ...]`` of length
  372. ``copies`` per the pipeline's fanout strategy. ``target_model_class``
  373. items leave ``printer_id`` None so the scheduler picks any free matching
  374. printer; specific assignments fill ``printer_id``. Class targeting only
  375. pins copies to printers in the runner's ``printer_scope`` (#1727)."""
  376. target_kind = pipeline.target_kind or "specific_printer"
  377. if target_kind == "specific_printer" or pipeline.target_printer_id is not None:
  378. assert pipeline.target_printer_id is not None
  379. return [(pipeline.target_printer_id, None)] * copies
  380. # Class-targeting. Enumerate matching printers + apply the strategy.
  381. matching = (
  382. (
  383. await db.execute(
  384. select(Printer)
  385. .where(Printer.model == pipeline.target_model_class)
  386. .where(Printer.is_active.is_(True))
  387. .order_by(Printer.id)
  388. )
  389. )
  390. .scalars()
  391. .all()
  392. )
  393. matching = [p for p in matching if printer_scope.allows(p.id)]
  394. if not matching:
  395. # Shouldn't reach here when eligibility passes, but failing gracefully
  396. # is better than a TypeError on next-slot pick.
  397. return [(None, pipeline.target_model_class)] * copies
  398. strategy = pipeline.fanout_strategy or "max_parallel"
  399. if strategy == "fill_one_first":
  400. # Pin every copy to the first match. Scheduler dispatches them serially
  401. # to that printer. If the printer breaks, copies wait; that's the
  402. # documented trade-off.
  403. return [(matching[0].id, None)] * copies
  404. if strategy == "round_robin":
  405. # Cycle through eligible printers — copy ``i`` lands on
  406. # ``matching[i % len(matching)]``. Each item gets a fixed printer_id.
  407. return [(matching[i % len(matching)].id, None) for i in range(copies)]
  408. # max_parallel — leave printer_id=None, set target_model so the scheduler
  409. # picks any free X1C / P1S / … for each item independently.
  410. return [(None, pipeline.target_model_class)] * copies
  411. def _make_orchestration_callable(
  412. *,
  413. run_id: int,
  414. pipeline_id: int,
  415. src_kind: SourceKind,
  416. src_id: int,
  417. src_filename: str,
  418. src_path: Path,
  419. creator_user_id: int | None,
  420. copies: int,
  421. printer_scope: PrinterScope = ALL_PRINTERS,
  422. review_required: bool = False,
  423. ):
  424. """Returns the async callable that ``slice_dispatch.enqueue`` runs as the
  425. background slice job. Wraps slice + multi-copy enqueue + state update."""
  426. async def _orchestrate(slice_job_id: int) -> dict:
  427. from backend.app.api.routes.library import slice_and_persist
  428. async with async_session() as session:
  429. run = (await session.execute(select(PipelineRun).where(PipelineRun.id == run_id))).scalar_one_or_none()
  430. pipeline = (
  431. await session.execute(select(SlicerPipeline).where(SlicerPipeline.id == pipeline_id))
  432. ).scalar_one_or_none()
  433. if run is None or pipeline is None:
  434. logger.warning("pipeline_run %d or pipeline %d disappeared mid-orchestration", run_id, pipeline_id)
  435. return {}
  436. # Honour a cancel that landed between ``POST /run`` returning and
  437. # this background task starting. If the run was cancelled while
  438. # still in ``queued`` we must NOT flip it back to ``slicing`` —
  439. # the operator's intent was to stop, and overwriting status here
  440. # was the bug that left runs stuck at ``dispatching`` after a
  441. # user-side cancel (#1425 PR C bug report).
  442. if run.status == "cancelled":
  443. logger.info("pipeline_run %d was cancelled before slicing started", run_id)
  444. return {}
  445. run.status = "slicing"
  446. run.started_at = datetime.now(timezone.utc)
  447. await session.commit()
  448. await _publish_run_event(session, run)
  449. slice_request = _slice_request_from_pipeline(pipeline)
  450. model_bytes = src_path.read_bytes()
  451. folder_id: int | None = None
  452. if src_kind == "library_file":
  453. lib = (await session.execute(select(LibraryFile).where(LibraryFile.id == src_id))).scalar_one_or_none()
  454. if lib is not None:
  455. folder_id = lib.folder_id
  456. try:
  457. slice_response = await slice_and_persist(
  458. session,
  459. model_bytes=model_bytes,
  460. model_filename=src_filename,
  461. folder_id=folder_id,
  462. extra_metadata={
  463. f"sliced_from_{src_kind}_id": src_id,
  464. "sliced_via_pipeline_id": pipeline.id,
  465. "sliced_via_pipeline_run_id": run.id,
  466. },
  467. request=slice_request,
  468. current_user_id=creator_user_id,
  469. job_id=slice_job_id,
  470. )
  471. except HTTPException as exc:
  472. run.status = "failed"
  473. run.error_message = f"Slice failed: {exc.detail}"
  474. run.completed_at = datetime.now(timezone.utc)
  475. await session.commit()
  476. await _publish_run_event(session, run)
  477. raise
  478. except Exception as exc:
  479. logger.exception("Pipeline run %d slice raised unexpectedly", run_id)
  480. run.status = "failed"
  481. run.error_message = f"Slice failed: {exc}"
  482. run.completed_at = datetime.now(timezone.utc)
  483. await session.commit()
  484. await _publish_run_event(session, run)
  485. raise
  486. run.sliced_library_file_id = slice_response.library_file_id
  487. # Re-check cancellation: the slice can take minutes, and the
  488. # operator may have hit Cancel during that window. Refresh from
  489. # the DB rather than trusting our in-memory `run` (the cancel
  490. # route writes via a separate session). When cancelled, don't
  491. # enqueue print queue items — that's the whole point of cancel.
  492. await session.refresh(run)
  493. if run.status == "cancelled":
  494. logger.info("pipeline_run %d cancelled mid-slice; skipping queue enqueue", run_id)
  495. await session.commit()
  496. return slice_response.model_dump()
  497. # PR C: enqueue N copies per the picked assignment strategy.
  498. assignments = await _pick_assignments(session, pipeline, copies, printer_scope)
  499. jobs = (
  500. (
  501. await session.execute(
  502. select(PipelineJob)
  503. .where(PipelineJob.pipeline_run_id == run_id)
  504. .order_by(PipelineJob.copy_index)
  505. )
  506. )
  507. .scalars()
  508. .all()
  509. )
  510. if len(jobs) != copies:
  511. logger.warning("pipeline_run %d expected %d jobs, found %d", run_id, copies, len(jobs))
  512. # No per-job ask-for-outcome toggle on a pipeline run either (#1898).
  513. confirm_outcome = await confirm_outcome_for_new_queue_item(session)
  514. for job, (printer_id, target_model) in zip(jobs, assignments, strict=False):
  515. queue_item = PrintQueueItem(
  516. printer_id=printer_id,
  517. target_model=target_model,
  518. library_file_id=slice_response.library_file_id,
  519. created_by_id=creator_user_id,
  520. status="pending",
  521. confirm_outcome=confirm_outcome,
  522. # The copies wait for review like any other job of theirs (#1620)
  523. manual_start=review_required,
  524. )
  525. session.add(queue_item)
  526. await session.flush()
  527. job.queue_entry_id = queue_item.id
  528. job.assigned_printer_id = printer_id # may be None for max_parallel
  529. # Don't write job.status yet — final cancellation check below
  530. # may flip it to 'cancelled' instead. dispatched_at is fine to
  531. # set unconditionally since the orchestration actually got here.
  532. job.dispatched_at = datetime.now(timezone.utc)
  533. # Final cancellation check before committing 'dispatching'. The
  534. # cancel route writes via a separate session so we have to refresh
  535. # to see the latest. If the cancel landed in this narrow window —
  536. # AFTER the post-slice refresh but BEFORE this commit — the queue
  537. # entries we just created would otherwise pick up and print. Mark
  538. # them + the per-copy jobs cancelled so the user's intent sticks.
  539. await session.refresh(run)
  540. if run.status == "cancelled":
  541. logger.info(
  542. "pipeline_run %d cancelled in the dispatch window; cancelling its %d queue entries",
  543. run_id,
  544. len(jobs),
  545. )
  546. for job in jobs:
  547. if job.queue_entry_id:
  548. qe = (
  549. await session.execute(select(PrintQueueItem).where(PrintQueueItem.id == job.queue_entry_id))
  550. ).scalar_one_or_none()
  551. if qe is not None and qe.status in ("pending", "queued"):
  552. qe.status = "cancelled"
  553. if job.status not in ("completed", "failed", "cancelled"):
  554. job.status = "cancelled"
  555. job.completed_at = datetime.now(timezone.utc)
  556. await session.commit()
  557. await _publish_run_event(session, run)
  558. return slice_response.model_dump()
  559. for job in jobs:
  560. job.status = "queued"
  561. run.status = "dispatching"
  562. await session.commit()
  563. await _publish_run_event(session, run)
  564. return slice_response.model_dump()
  565. return _orchestrate
  566. # ---------------------------------------------------------------------------
  567. # /slicer-pipelines/{id}/check-eligibility
  568. # ---------------------------------------------------------------------------
  569. @pipeline_run_create_router.post("/{pipeline_id}/check-eligibility", response_model=EligibilityReportResponse)
  570. async def check_eligibility(
  571. pipeline_id: int,
  572. body: CheckEligibilityRequest,
  573. current_user: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_READ),
  574. printer_scope: PrinterScope = RequestPrinterScope,
  575. db: AsyncSession = Depends(get_db),
  576. ):
  577. pipeline = await _load_pipeline(db, pipeline_id)
  578. printer_scope.ensure(pipeline.target_printer_id)
  579. await _resolve_source(
  580. db,
  581. library_file_id=body.source_library_file_id,
  582. archive_id=body.source_archive_id,
  583. user=current_user,
  584. printer_scope=printer_scope,
  585. )
  586. if pipeline.target_kind == "printer_class" and pipeline.target_printer_id is None:
  587. report = await check_pipeline_eligibility(db, pipeline, status_lookup=_make_status_lookup())
  588. else:
  589. status = await _load_printer_status(pipeline.target_printer_id)
  590. report = await check_pipeline_eligibility(db, pipeline, status)
  591. return _serialise_status(report)
  592. # ---------------------------------------------------------------------------
  593. # /slicer-pipelines/{id}/run
  594. # ---------------------------------------------------------------------------
  595. @pipeline_run_create_router.post("/{pipeline_id}/run", response_model=PipelineRunResponse, status_code=202)
  596. async def run_pipeline(
  597. pipeline_id: int,
  598. body: PipelineRunCreateRequest,
  599. current_user: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_RUN),
  600. api_key_cloud_owner: User | None = Depends(resolve_api_key_cloud_owner),
  601. printer_scope: PrinterScope = RequestPrinterScope,
  602. review_required: bool = QueueReviewRequired,
  603. db: AsyncSession = Depends(get_db),
  604. ):
  605. from backend.app.api.routes.settings import get_setting
  606. from backend.app.services.slice_dispatch import slice_dispatch
  607. pipeline = await _load_pipeline(db, pipeline_id)
  608. # The pipeline is shared config; running it is limited to what the caller
  609. # may print on (#1727)
  610. printer_scope.ensure(pipeline.target_printer_id)
  611. # ``user=current_user`` deliberately, not the cloud owner below: an API-key
  612. # caller has no per-row identity and must keep can_read_all, the same as
  613. # every other read helper.
  614. src_kind, src_id, src_filename, src_path = await _resolve_source(
  615. db,
  616. library_file_id=body.source_library_file_id,
  617. archive_id=body.source_archive_id,
  618. user=current_user,
  619. printer_scope=printer_scope,
  620. )
  621. # The permission gate answers an API-keyed request with current_user=None,
  622. # so a pipeline built on Bambu/Orca Cloud presets would have nobody whose
  623. # stored cloud token could resolve them. Fall back to the key's owner, the
  624. # same fallback POST /library/files/{id}/slice makes (#1182 follow-up).
  625. # Only keys with the cloud scope resolve to an owner here; everything else
  626. # stays None and slices against local presets exactly as before.
  627. creator = current_user or api_key_cloud_owner
  628. # Copies left to the scheduler run within their creator's printers. A
  629. # limited API key can't be held to its own printers that way -- its cloud
  630. # owner may see more -- and even a pinning strategy falls back to "any
  631. # printer of the class" when none of the class is in scope, so a limited
  632. # key may only run pipelines aimed at one printer (#1727).
  633. if pipeline.target_printer_id is None:
  634. ensure_model_target_allowed(current_user, printer_scope)
  635. # Cap copies against the configured ceiling.
  636. raw_cap = await get_setting(db, "pipeline_max_copies")
  637. try:
  638. cap = int(raw_cap) if raw_cap else 50
  639. except (TypeError, ValueError):
  640. cap = 50
  641. if body.copies > cap:
  642. raise HTTPException(
  643. 422,
  644. f"copies={body.copies} exceeds pipeline_max_copies setting ({cap})",
  645. )
  646. # Eligibility pre-flight.
  647. if pipeline.target_kind == "printer_class" and pipeline.target_printer_id is None:
  648. report = await check_pipeline_eligibility(db, pipeline, status_lookup=_make_status_lookup())
  649. else:
  650. status = await _load_printer_status(pipeline.target_printer_id)
  651. report = await check_pipeline_eligibility(db, pipeline, status)
  652. if not report.ok and not body.force:
  653. raise HTTPException(status_code=409, detail=_serialise_status(report).model_dump())
  654. # Need a target — specific or class — to dispatch.
  655. if pipeline.target_printer_id is None and not pipeline.target_model_class:
  656. raise HTTPException(
  657. 400,
  658. "Pipeline has no target. Open the pipeline in Settings → Workflow → Pipelines and choose a target printer or printer class.",
  659. )
  660. run = PipelineRun(
  661. pipeline_id=pipeline.id,
  662. source_library_file_id=src_id if src_kind == "library_file" else None,
  663. source_archive_id=src_id if src_kind == "archive" else None,
  664. copies=body.copies,
  665. status="queued",
  666. eligibility_overridden=(not report.ok and body.force),
  667. created_by=creator.id if creator else None,
  668. )
  669. db.add(run)
  670. await db.flush()
  671. # One PipelineJob per copy. PR B was copies=1, PR C generalises.
  672. for i in range(body.copies):
  673. db.add(
  674. PipelineJob(
  675. pipeline_run_id=run.id,
  676. copy_index=i,
  677. status="pending",
  678. )
  679. )
  680. await db.commit()
  681. await db.refresh(run)
  682. await _publish_run_event(db, run)
  683. orchestrate = _make_orchestration_callable(
  684. run_id=run.id,
  685. pipeline_id=pipeline.id,
  686. src_kind=src_kind,
  687. src_id=src_id,
  688. src_filename=src_filename,
  689. src_path=src_path,
  690. creator_user_id=creator.id if creator else None,
  691. copies=body.copies,
  692. printer_scope=printer_scope,
  693. review_required=review_required,
  694. )
  695. slice_job = await slice_dispatch.enqueue(
  696. kind="library_file" if src_kind == "library_file" else "archive",
  697. source_id=src_id,
  698. source_name=src_filename,
  699. owner_id=creator.id if creator else None,
  700. run=orchestrate,
  701. )
  702. run.slice_job_id = slice_job.id
  703. await db.commit()
  704. await db.refresh(run)
  705. return await _materialise_run(db, run)
  706. # ---------------------------------------------------------------------------
  707. # Lists, reads, cancel, retry-failed
  708. # ---------------------------------------------------------------------------
  709. def _run_scope_clause(printer_scope: PrinterScope):
  710. """Runs the caller may see (#1727): their pipeline targets no printer out of
  711. scope, and no copy was assigned to one. None when unrestricted."""
  712. if printer_scope.is_unrestricted:
  713. return None
  714. allowed = printer_scope.printer_ids
  715. hidden_by_job = select(PipelineJob.pipeline_run_id).where(
  716. PipelineJob.assigned_printer_id.is_not(None), PipelineJob.assigned_printer_id.not_in(allowed)
  717. )
  718. hidden_pipelines = select(SlicerPipeline.id).where(
  719. SlicerPipeline.target_printer_id.is_not(None), SlicerPipeline.target_printer_id.not_in(allowed)
  720. )
  721. return PipelineRun.id.not_in(hidden_by_job) & (
  722. PipelineRun.pipeline_id.is_(None) | PipelineRun.pipeline_id.not_in(hidden_pipelines)
  723. )
  724. async def _load_run_in_scope(db: AsyncSession, run_id: int, printer_scope: PrinterScope) -> PipelineRun:
  725. query = select(PipelineRun).where(PipelineRun.id == run_id)
  726. if (clause := _run_scope_clause(printer_scope)) is not None:
  727. query = query.where(clause)
  728. run = (await db.execute(query)).scalar_one_or_none()
  729. if run is None:
  730. raise HTTPException(404, "Pipeline run not found")
  731. return run
  732. @pipeline_run_create_router.get("/{pipeline_id}/runs", response_model=PipelineRunListResponse)
  733. async def list_runs_for_pipeline(
  734. pipeline_id: int,
  735. limit: int = 10,
  736. _: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_READ),
  737. db: AsyncSession = Depends(get_db),
  738. printer_scope: PrinterScope = RequestPrinterScope,
  739. ):
  740. limit = max(1, min(limit, 100))
  741. conditions = [PipelineRun.pipeline_id == pipeline_id]
  742. if (clause := _run_scope_clause(printer_scope)) is not None:
  743. conditions.append(clause)
  744. rows = (
  745. (await db.execute(select(PipelineRun).where(*conditions).order_by(PipelineRun.id.desc()).limit(limit)))
  746. .scalars()
  747. .all()
  748. )
  749. total = (await db.execute(select(func.count()).select_from(PipelineRun).where(*conditions))).scalar() or 0
  750. return PipelineRunListResponse(
  751. runs=[await _materialise_run(db, r) for r in rows],
  752. total=total,
  753. )
  754. @pipeline_run_router.get("", response_model=PipelineRunListResponse)
  755. async def list_all_runs(
  756. limit: int = 25,
  757. offset: int = 0,
  758. pipeline_id: int | None = None,
  759. status: str | None = None,
  760. target_printer_id: int | None = None,
  761. target_model_class: str | None = None,
  762. _: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_READ),
  763. db: AsyncSession = Depends(get_db),
  764. printer_scope: PrinterScope = RequestPrinterScope,
  765. ):
  766. """Dashboard list. Newest first; filters on pipeline_id + status +
  767. target_printer_id + target_model_class. The ``status`` filter matches
  768. the persisted snapshot, not the live roll-up — in-progress runs may
  769. appear under ``dispatching`` until the next state transition writes
  770. through. ``target_*`` filters JOIN to the pipeline so runs whose
  771. pipeline currently points at the printer / class are returned."""
  772. limit = max(1, min(limit, 100))
  773. offset = max(0, offset)
  774. stmt = select(PipelineRun)
  775. count_stmt = select(func.count()).select_from(PipelineRun)
  776. if (clause := _run_scope_clause(printer_scope)) is not None:
  777. stmt = stmt.where(clause)
  778. count_stmt = count_stmt.where(clause)
  779. if pipeline_id is not None:
  780. stmt = stmt.where(PipelineRun.pipeline_id == pipeline_id)
  781. count_stmt = count_stmt.where(PipelineRun.pipeline_id == pipeline_id)
  782. if status:
  783. stmt = stmt.where(PipelineRun.status == status)
  784. count_stmt = count_stmt.where(PipelineRun.status == status)
  785. if target_printer_id is not None or target_model_class is not None:
  786. stmt = stmt.join(SlicerPipeline, SlicerPipeline.id == PipelineRun.pipeline_id)
  787. count_stmt = count_stmt.join(SlicerPipeline, SlicerPipeline.id == PipelineRun.pipeline_id)
  788. if target_printer_id is not None:
  789. stmt = stmt.where(SlicerPipeline.target_printer_id == target_printer_id)
  790. count_stmt = count_stmt.where(SlicerPipeline.target_printer_id == target_printer_id)
  791. if target_model_class is not None:
  792. stmt = stmt.where(SlicerPipeline.target_model_class == target_model_class)
  793. count_stmt = count_stmt.where(SlicerPipeline.target_model_class == target_model_class)
  794. rows = (await db.execute(stmt.order_by(desc(PipelineRun.id)).offset(offset).limit(limit))).scalars().all()
  795. total = (await db.execute(count_stmt)).scalar() or 0
  796. return PipelineRunListResponse(
  797. runs=[await _materialise_run(db, r) for r in rows],
  798. total=total,
  799. )
  800. _TERMINAL_RUN_STATUSES = ("completed", "failed", "cancelled", "partial_failure")
  801. @pipeline_run_router.post("/clear")
  802. async def clear_terminal_runs(
  803. _: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_WRITE),
  804. db: AsyncSession = Depends(get_db),
  805. printer_scope: PrinterScope = RequestPrinterScope,
  806. ):
  807. """Delete every terminal pipeline run (completed / failed / cancelled /
  808. partial_failure). In-flight runs (queued / slicing / dispatching /
  809. in_progress) are preserved — clearing those mid-flight would lose the
  810. operator's intent. Cascades to PipelineJob via the ondelete='CASCADE'
  811. relationship; the linked PrintQueueItem rows stay (they have their own
  812. lifecycle on the queue page)."""
  813. # Count first so the response can report how many got cleared. Done
  814. # under the same session/transaction as the delete so the numbers can't
  815. # drift if another caller races in.
  816. # Runs on printers the caller can't see are left alone (#1727)
  817. conditions = [PipelineRun.status.in_(_TERMINAL_RUN_STATUSES)]
  818. if (clause := _run_scope_clause(printer_scope)) is not None:
  819. conditions.append(clause)
  820. count_stmt = select(func.count()).select_from(PipelineRun).where(*conditions)
  821. n = (await db.execute(count_stmt)).scalar() or 0
  822. if n > 0:
  823. await db.execute(delete(PipelineRun).where(*conditions))
  824. await db.commit()
  825. return {"deleted": n}
  826. @pipeline_run_router.get("/{run_id}", response_model=PipelineRunResponse)
  827. async def get_run(
  828. run_id: int,
  829. _: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_READ),
  830. db: AsyncSession = Depends(get_db),
  831. printer_scope: PrinterScope = RequestPrinterScope,
  832. ):
  833. run = await _load_run_in_scope(db, run_id, printer_scope)
  834. return await _materialise_run(db, run)
  835. @pipeline_run_router.post("/{run_id}/cancel", response_model=PipelineRunResponse)
  836. async def cancel_run(
  837. run_id: int,
  838. _: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_RUN),
  839. db: AsyncSession = Depends(get_db),
  840. printer_scope: PrinterScope = RequestPrinterScope,
  841. ):
  842. """Cancel a queued / in-flight run. Cascades to all non-terminal queue
  843. entries; in-flight prints continue on the printer (operator must Stop)."""
  844. run = await _load_run_in_scope(db, run_id, printer_scope)
  845. if run.status in ("completed", "failed", "cancelled", "partial_failure"):
  846. return await _materialise_run(db, run)
  847. run.status = "cancelled"
  848. run.completed_at = datetime.now(timezone.utc)
  849. if not run.error_message:
  850. run.error_message = "Cancelled by user"
  851. job_rows = (await db.execute(select(PipelineJob).where(PipelineJob.pipeline_run_id == run.id))).scalars().all()
  852. for job in job_rows:
  853. if job.queue_entry_id:
  854. queue_entry = (
  855. await db.execute(select(PrintQueueItem).where(PrintQueueItem.id == job.queue_entry_id))
  856. ).scalar_one_or_none()
  857. if queue_entry is not None and queue_entry.status in ("pending", "queued"):
  858. queue_entry.status = "cancelled"
  859. if job.status not in ("completed", "failed", "cancelled"):
  860. job.status = "cancelled"
  861. job.completed_at = datetime.now(timezone.utc)
  862. await db.commit()
  863. await db.refresh(run)
  864. await _publish_run_event(db, run)
  865. return await _materialise_run(db, run)
  866. @pipeline_run_router.post("/{run_id}/retry-failed", response_model=PipelineRunResponse, status_code=202)
  867. async def retry_failed(
  868. run_id: int,
  869. current_user: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_RUN),
  870. api_key_cloud_owner: User | None = Depends(resolve_api_key_cloud_owner),
  871. printer_scope: PrinterScope = RequestPrinterScope,
  872. review_required: bool = QueueReviewRequired,
  873. db: AsyncSession = Depends(get_db),
  874. ):
  875. """Create a new run with copies = (failed + cancelled count) from the
  876. parent. Same pipeline, same source. Eligibility re-checked at run time
  877. (it might pass this time — operator may have fixed the issue)."""
  878. parent = await _load_run_in_scope(db, run_id, printer_scope)
  879. if parent.pipeline_id is None:
  880. raise HTTPException(400, "Original pipeline was deleted; cannot retry")
  881. if parent.source_library_file_id is None and parent.source_archive_id is None:
  882. raise HTTPException(400, "Original source was deleted; cannot retry")
  883. # Count the parent's failed + cancelled jobs.
  884. parent_jobs = (
  885. (await db.execute(select(PipelineJob).where(PipelineJob.pipeline_run_id == parent.id))).scalars().all()
  886. )
  887. fail_count = 0
  888. for j in parent_jobs:
  889. queue_entry = None
  890. if j.queue_entry_id:
  891. queue_entry = (
  892. await db.execute(select(PrintQueueItem).where(PrintQueueItem.id == j.queue_entry_id))
  893. ).scalar_one_or_none()
  894. live = _compute_job_status(j.status, queue_entry)
  895. if live in ("failed", "cancelled"):
  896. fail_count += 1
  897. if fail_count == 0:
  898. raise HTTPException(400, "No failed copies to retry")
  899. # Build the request payload the same way the user would have via /run.
  900. body = PipelineRunCreateRequest(
  901. source_library_file_id=parent.source_library_file_id,
  902. source_archive_id=parent.source_archive_id,
  903. copies=fail_count,
  904. force=True, # operator already accepted eligibility on the parent
  905. )
  906. # Reuse the run_pipeline route logic via a direct call — keeps the
  907. # orchestration single-sourced. The result inherits parent_run_id. Every
  908. # dependency it declares has to be forwarded explicitly: FastAPI resolves
  909. # those only for a routed request, so an omitted one would arrive as the
  910. # Depends() marker object itself rather than as None.
  911. new_run_response = await run_pipeline(
  912. parent.pipeline_id,
  913. body,
  914. current_user=current_user,
  915. api_key_cloud_owner=api_key_cloud_owner,
  916. printer_scope=printer_scope,
  917. review_required=review_required,
  918. db=db,
  919. )
  920. # Stamp parent_run_id on the freshly-created run.
  921. new_row = (await db.execute(select(PipelineRun).where(PipelineRun.id == new_run_response.id))).scalar_one_or_none()
  922. if new_row is not None:
  923. new_row.parent_run_id = parent.id
  924. await db.commit()
  925. await db.refresh(new_row)
  926. return await _materialise_run(db, new_row)
  927. return new_run_response