pipeline_runs.py 33 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855
  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 desc, func, select
  28. from sqlalchemy.ext.asyncio import AsyncSession
  29. from backend.app.core.auth import RequirePermissionIfAuthEnabled
  30. from backend.app.core.config import settings as app_settings
  31. from backend.app.core.database import async_session, get_db
  32. from backend.app.core.permissions import Permission
  33. from backend.app.core.websocket import ws_manager
  34. from backend.app.models.archive import PrintArchive
  35. from backend.app.models.library import LibraryFile
  36. from backend.app.models.pipeline_run import PipelineJob, PipelineRun
  37. from backend.app.models.print_queue import PrintQueueItem
  38. from backend.app.models.printer import Printer
  39. from backend.app.models.slicer_pipeline import SlicerPipeline
  40. from backend.app.models.user import User
  41. from backend.app.schemas.pipeline_run import (
  42. CheckEligibilityRequest,
  43. EligibilityIssueResponse,
  44. EligibilityReportResponse,
  45. PerPrinterReport as PerPrinterReportResponse,
  46. PipelineJobResponse,
  47. PipelineRunCreateRequest,
  48. PipelineRunListResponse,
  49. PipelineRunResponse,
  50. )
  51. from backend.app.schemas.slicer import PresetRef, SliceRequest
  52. from backend.app.services.pipeline_eligibility import (
  53. EligibilityReport,
  54. check_pipeline_eligibility,
  55. )
  56. logger = logging.getLogger(__name__)
  57. pipeline_run_create_router = APIRouter(prefix="/slicer-pipelines", tags=["Slicer Pipelines"])
  58. pipeline_run_router = APIRouter(prefix="/pipeline-runs", tags=["Slicer Pipelines"])
  59. # ---------------------------------------------------------------------------
  60. # Helpers
  61. # ---------------------------------------------------------------------------
  62. def _serialise_status(report: EligibilityReport) -> EligibilityReportResponse:
  63. return EligibilityReportResponse(
  64. ok=report.ok,
  65. target_kind=report.target_kind,
  66. target_printer_id=report.target_printer_id,
  67. target_printer_name=report.target_printer_name,
  68. target_model_class=report.target_model_class,
  69. issues=[
  70. EligibilityIssueResponse(
  71. kind=issue.kind,
  72. slot_index=issue.slot_index,
  73. expected=issue.expected,
  74. actual=issue.actual,
  75. )
  76. for issue in report.issues
  77. ],
  78. printer_reports=[
  79. PerPrinterReportResponse(
  80. printer_id=r.printer_id,
  81. printer_name=r.printer_name,
  82. ok=r.ok,
  83. issues=[
  84. EligibilityIssueResponse(
  85. kind=i.kind,
  86. slot_index=i.slot_index,
  87. expected=i.expected,
  88. actual=i.actual,
  89. )
  90. for i in r.issues
  91. ],
  92. )
  93. for r in report.printer_reports
  94. ],
  95. )
  96. async def _load_pipeline(db: AsyncSession, pipeline_id: int) -> SlicerPipeline:
  97. pipeline = (
  98. await db.execute(
  99. select(SlicerPipeline).where(
  100. SlicerPipeline.id == pipeline_id,
  101. SlicerPipeline.is_deleted.is_(False),
  102. )
  103. )
  104. ).scalar_one_or_none()
  105. if pipeline is None:
  106. raise HTTPException(404, "Pipeline not found")
  107. return pipeline
  108. async def _load_printer_status(printer_id: int | None) -> dict | None:
  109. """Snapshot the printer_manager's live PrinterState for the eligibility
  110. matcher. Returns ``None`` when the printer has no MQTT client."""
  111. if printer_id is None:
  112. return None
  113. from backend.app.services.printer_manager import printer_manager
  114. state = printer_manager.get_status(printer_id)
  115. if state is None:
  116. return None
  117. return {"connected": state.connected, "raw_data": state.raw_data}
  118. def _make_status_lookup():
  119. """Closure that snapshots the printer_manager once per printer_id call.
  120. Passed to the matcher's class-targeting branch so it can read live state
  121. for every candidate printer."""
  122. def _lookup(printer_id: int) -> dict | None:
  123. from backend.app.services.printer_manager import printer_manager
  124. state = printer_manager.get_status(printer_id)
  125. if state is None:
  126. return None
  127. return {"connected": state.connected, "raw_data": state.raw_data}
  128. return _lookup
  129. def _slice_request_from_pipeline(pipeline: SlicerPipeline) -> SliceRequest:
  130. try:
  131. raw_filaments = json.loads(pipeline.filament_presets_json or "[]")
  132. except (json.JSONDecodeError, TypeError):
  133. raw_filaments = []
  134. filament_presets = [
  135. PresetRef(source=r["source"], id=r["id"])
  136. for r in raw_filaments
  137. if isinstance(r, dict) and "source" in r and "id" in r
  138. ]
  139. return SliceRequest(
  140. printer_preset=PresetRef(source=pipeline.printer_preset_source, id=pipeline.printer_preset_id),
  141. process_preset=PresetRef(source=pipeline.process_preset_source, id=pipeline.process_preset_id),
  142. filament_presets=filament_presets,
  143. bed_type=pipeline.bed_type,
  144. export_3mf=True,
  145. )
  146. def _compute_job_status(
  147. persisted: str,
  148. queue_entry: PrintQueueItem | None,
  149. ) -> str:
  150. if persisted in ("failed", "cancelled", "completed"):
  151. return persisted
  152. if queue_entry is None:
  153. return persisted
  154. qs = queue_entry.status
  155. if qs == "completed":
  156. return "completed"
  157. if qs in ("failed", "aborted"):
  158. return "failed"
  159. if qs == "cancelled":
  160. return "cancelled"
  161. if qs == "printing":
  162. return "printing"
  163. return "queued"
  164. def _roll_up_run_status(
  165. persisted: str,
  166. job_statuses: list[str],
  167. ) -> str:
  168. """Compute the run-level status from the per-job statuses.
  169. Terminal-persisted always wins for explicit cancels / hard failures so
  170. the dashboard doesn't flicker when one job's queue entry hasn't caught
  171. up. Otherwise:
  172. - all completed → completed
  173. - any in_progress / printing / queued / dispatching → in_progress
  174. - any failed alongside any completed → partial_failure
  175. - all failed/cancelled → failed
  176. """
  177. if persisted in ("cancelled",):
  178. return persisted
  179. if not job_statuses:
  180. return persisted
  181. completed = sum(1 for s in job_statuses if s == "completed")
  182. failed = sum(1 for s in job_statuses if s == "failed")
  183. cancelled = sum(1 for s in job_statuses if s == "cancelled")
  184. in_flight = sum(1 for s in job_statuses if s in ("printing", "queued", "awaiting_printer", "pending"))
  185. total = len(job_statuses)
  186. if completed == total:
  187. return "completed"
  188. if in_flight > 0:
  189. return "in_progress" if persisted not in ("queued", "slicing", "dispatching") else persisted
  190. # All copies are in terminal states.
  191. if failed == 0 and cancelled == total:
  192. return "cancelled"
  193. if completed > 0 and (failed > 0 or cancelled > 0):
  194. return "partial_failure"
  195. if failed > 0:
  196. return "failed"
  197. return persisted
  198. async def _materialise_run(db: AsyncSession, run: PipelineRun) -> PipelineRunResponse:
  199. pipeline_name: str | None = None
  200. target_kind = None
  201. target_printer_id = None
  202. target_model_class = None
  203. fanout_strategy = None
  204. if run.pipeline_id:
  205. pipeline = (
  206. await db.execute(select(SlicerPipeline).where(SlicerPipeline.id == run.pipeline_id))
  207. ).scalar_one_or_none()
  208. if pipeline:
  209. pipeline_name = pipeline.name
  210. target_kind = pipeline.target_kind # type: ignore[assignment]
  211. target_printer_id = pipeline.target_printer_id
  212. target_model_class = pipeline.target_model_class
  213. fanout_strategy = pipeline.fanout_strategy # type: ignore[assignment]
  214. source_filename: str | None = None
  215. if run.source_library_file_id:
  216. src = (
  217. await db.execute(select(LibraryFile).where(LibraryFile.id == run.source_library_file_id))
  218. ).scalar_one_or_none()
  219. source_filename = src.filename if src else None
  220. elif run.source_archive_id:
  221. arc = (
  222. await db.execute(select(PrintArchive).where(PrintArchive.id == run.source_archive_id))
  223. ).scalar_one_or_none()
  224. source_filename = (arc.print_name or arc.filename) if arc else None
  225. job_rows = (
  226. (
  227. await db.execute(
  228. select(PipelineJob).where(PipelineJob.pipeline_run_id == run.id).order_by(PipelineJob.copy_index)
  229. )
  230. )
  231. .scalars()
  232. .all()
  233. )
  234. job_responses: list[PipelineJobResponse] = []
  235. job_live_statuses: list[str] = []
  236. for job in job_rows:
  237. queue_entry = None
  238. if job.queue_entry_id:
  239. queue_entry = (
  240. await db.execute(select(PrintQueueItem).where(PrintQueueItem.id == job.queue_entry_id))
  241. ).scalar_one_or_none()
  242. printer_name: str | None = None
  243. if job.assigned_printer_id:
  244. p = (await db.execute(select(Printer).where(Printer.id == job.assigned_printer_id))).scalar_one_or_none()
  245. printer_name = p.name if p else None
  246. live_job_status = _compute_job_status(job.status, queue_entry)
  247. job_live_statuses.append(live_job_status)
  248. job_responses.append(
  249. PipelineJobResponse(
  250. id=job.id,
  251. pipeline_run_id=job.pipeline_run_id,
  252. copy_index=job.copy_index,
  253. assigned_printer_id=job.assigned_printer_id,
  254. assigned_printer_name=printer_name,
  255. queue_entry_id=job.queue_entry_id,
  256. status=live_job_status, # type: ignore[arg-type]
  257. error_message=job.error_message,
  258. dispatched_at=job.dispatched_at,
  259. completed_at=job.completed_at,
  260. )
  261. )
  262. rolled_up = _roll_up_run_status(run.status, job_live_statuses)
  263. return PipelineRunResponse(
  264. id=run.id,
  265. pipeline_id=run.pipeline_id,
  266. pipeline_name=pipeline_name,
  267. source_library_file_id=run.source_library_file_id,
  268. source_archive_id=run.source_archive_id,
  269. source_filename=source_filename,
  270. parent_run_id=run.parent_run_id,
  271. copies=run.copies,
  272. copies_completed=sum(1 for s in job_live_statuses if s == "completed"),
  273. copies_failed=sum(1 for s in job_live_statuses if s == "failed"),
  274. copies_cancelled=sum(1 for s in job_live_statuses if s == "cancelled"),
  275. copies_in_progress=sum(
  276. 1 for s in job_live_statuses if s in ("printing", "queued", "awaiting_printer", "pending")
  277. ),
  278. status=rolled_up, # type: ignore[arg-type]
  279. slice_job_id=run.slice_job_id,
  280. sliced_library_file_id=run.sliced_library_file_id,
  281. eligibility_overridden=run.eligibility_overridden,
  282. error_message=run.error_message,
  283. created_by=run.created_by,
  284. created_at=run.created_at,
  285. started_at=run.started_at,
  286. completed_at=run.completed_at,
  287. jobs=job_responses,
  288. target_kind=target_kind,
  289. target_printer_id=target_printer_id,
  290. target_model_class=target_model_class,
  291. fanout_strategy=fanout_strategy,
  292. )
  293. async def _publish_run_event(db: AsyncSession, run: PipelineRun) -> None:
  294. """Broadcast a ``pipeline_run_updated`` event with the full materialised
  295. run. Per-user routing via ``broadcast_to_user`` falls back to a global
  296. broadcast when ``created_by`` is None (auth-disabled installs)."""
  297. try:
  298. payload = await _materialise_run(db, run)
  299. await ws_manager.broadcast_to_user(
  300. run.created_by,
  301. {
  302. "type": "pipeline_run_updated",
  303. "run": payload.model_dump(mode="json"),
  304. },
  305. )
  306. except Exception:
  307. logger.exception("Failed to broadcast pipeline_run_updated for run %d", run.id)
  308. # ---------------------------------------------------------------------------
  309. # Source resolution + orchestration
  310. # ---------------------------------------------------------------------------
  311. SourceKind = Literal["library_file", "archive"]
  312. async def _resolve_source(
  313. db: AsyncSession,
  314. *,
  315. library_file_id: int | None,
  316. archive_id: int | None,
  317. ) -> tuple[SourceKind, int, str, Path]:
  318. if library_file_id is not None:
  319. lib = (await db.execute(select(LibraryFile).where(LibraryFile.id == library_file_id))).scalar_one_or_none()
  320. if lib is None:
  321. raise HTTPException(404, "Source library file not found")
  322. src_path = (
  323. Path(app_settings.base_dir) / lib.file_path
  324. ) # 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/.
  325. if not src_path.exists():
  326. raise HTTPException(404, "Source library file missing on disk")
  327. return ("library_file", lib.id, lib.filename, src_path)
  328. assert archive_id is not None
  329. arc = (await db.execute(select(PrintArchive).where(PrintArchive.id == archive_id))).scalar_one_or_none()
  330. if arc is None:
  331. raise HTTPException(404, "Source archive not found")
  332. rel = arc.source_3mf_path or arc.file_path
  333. if not rel:
  334. raise HTTPException(400, "Archive has no source file to slice")
  335. src_path = (
  336. Path(app_settings.base_dir) / rel
  337. ) # 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.
  338. if not src_path.exists():
  339. raise HTTPException(404, "Archive source file missing on disk")
  340. name = arc.filename or arc.print_name or src_path.name
  341. return ("archive", arc.id, name, src_path)
  342. async def _pick_assignments(
  343. db: AsyncSession,
  344. pipeline: SlicerPipeline,
  345. copies: int,
  346. ) -> list[tuple[int | None, str | None]]:
  347. """Return ``[(printer_id_or_None, target_model_or_None), ...]`` of length
  348. ``copies`` per the pipeline's fanout strategy. ``target_model_class``
  349. items leave ``printer_id`` None so the scheduler picks any free matching
  350. printer; specific assignments fill ``printer_id``."""
  351. target_kind = pipeline.target_kind or "specific_printer"
  352. if target_kind == "specific_printer" or pipeline.target_printer_id is not None:
  353. assert pipeline.target_printer_id is not None
  354. return [(pipeline.target_printer_id, None)] * copies
  355. # Class-targeting. Enumerate matching printers + apply the strategy.
  356. matching = (
  357. (
  358. await db.execute(
  359. select(Printer)
  360. .where(Printer.model == pipeline.target_model_class)
  361. .where(Printer.is_active.is_(True))
  362. .order_by(Printer.id)
  363. )
  364. )
  365. .scalars()
  366. .all()
  367. )
  368. if not matching:
  369. # Shouldn't reach here when eligibility passes, but failing gracefully
  370. # is better than a TypeError on next-slot pick.
  371. return [(None, pipeline.target_model_class)] * copies
  372. strategy = pipeline.fanout_strategy or "max_parallel"
  373. if strategy == "fill_one_first":
  374. # Pin every copy to the first match. Scheduler dispatches them serially
  375. # to that printer. If the printer breaks, copies wait; that's the
  376. # documented trade-off.
  377. return [(matching[0].id, None)] * copies
  378. if strategy == "round_robin":
  379. # Cycle through eligible printers — copy ``i`` lands on
  380. # ``matching[i % len(matching)]``. Each item gets a fixed printer_id.
  381. return [(matching[i % len(matching)].id, None) for i in range(copies)]
  382. # max_parallel — leave printer_id=None, set target_model so the scheduler
  383. # picks any free X1C / P1S / … for each item independently.
  384. return [(None, pipeline.target_model_class)] * copies
  385. def _make_orchestration_callable(
  386. *,
  387. run_id: int,
  388. pipeline_id: int,
  389. src_kind: SourceKind,
  390. src_id: int,
  391. src_filename: str,
  392. src_path: Path,
  393. creator_user_id: int | None,
  394. copies: int,
  395. ):
  396. """Returns the async callable that ``slice_dispatch.enqueue`` runs as the
  397. background slice job. Wraps slice + multi-copy enqueue + state update."""
  398. async def _orchestrate(slice_job_id: int) -> dict:
  399. from backend.app.api.routes.library import slice_and_persist
  400. async with async_session() as session:
  401. run = (await session.execute(select(PipelineRun).where(PipelineRun.id == run_id))).scalar_one_or_none()
  402. pipeline = (
  403. await session.execute(select(SlicerPipeline).where(SlicerPipeline.id == pipeline_id))
  404. ).scalar_one_or_none()
  405. if run is None or pipeline is None:
  406. logger.warning("pipeline_run %d or pipeline %d disappeared mid-orchestration", run_id, pipeline_id)
  407. return {}
  408. run.status = "slicing"
  409. run.started_at = datetime.now(timezone.utc)
  410. await session.commit()
  411. await _publish_run_event(session, run)
  412. slice_request = _slice_request_from_pipeline(pipeline)
  413. model_bytes = src_path.read_bytes()
  414. folder_id: int | None = None
  415. if src_kind == "library_file":
  416. lib = (await session.execute(select(LibraryFile).where(LibraryFile.id == src_id))).scalar_one_or_none()
  417. if lib is not None:
  418. folder_id = lib.folder_id
  419. try:
  420. slice_response = await slice_and_persist(
  421. session,
  422. model_bytes=model_bytes,
  423. model_filename=src_filename,
  424. folder_id=folder_id,
  425. extra_metadata={
  426. f"sliced_from_{src_kind}_id": src_id,
  427. "sliced_via_pipeline_id": pipeline.id,
  428. "sliced_via_pipeline_run_id": run.id,
  429. },
  430. request=slice_request,
  431. current_user_id=creator_user_id,
  432. job_id=slice_job_id,
  433. )
  434. except HTTPException as exc:
  435. run.status = "failed"
  436. run.error_message = f"Slice failed: {exc.detail}"
  437. run.completed_at = datetime.now(timezone.utc)
  438. await session.commit()
  439. await _publish_run_event(session, run)
  440. raise
  441. except Exception as exc:
  442. logger.exception("Pipeline run %d slice raised unexpectedly", run_id)
  443. run.status = "failed"
  444. run.error_message = f"Slice failed: {exc}"
  445. run.completed_at = datetime.now(timezone.utc)
  446. await session.commit()
  447. await _publish_run_event(session, run)
  448. raise
  449. run.sliced_library_file_id = slice_response.library_file_id
  450. # PR C: enqueue N copies per the picked assignment strategy.
  451. assignments = await _pick_assignments(session, pipeline, copies)
  452. jobs = (
  453. (
  454. await session.execute(
  455. select(PipelineJob)
  456. .where(PipelineJob.pipeline_run_id == run_id)
  457. .order_by(PipelineJob.copy_index)
  458. )
  459. )
  460. .scalars()
  461. .all()
  462. )
  463. if len(jobs) != copies:
  464. logger.warning("pipeline_run %d expected %d jobs, found %d", run_id, copies, len(jobs))
  465. for job, (printer_id, target_model) in zip(jobs, assignments, strict=False):
  466. queue_item = PrintQueueItem(
  467. printer_id=printer_id,
  468. target_model=target_model,
  469. library_file_id=slice_response.library_file_id,
  470. created_by_id=creator_user_id,
  471. status="pending",
  472. )
  473. session.add(queue_item)
  474. await session.flush()
  475. job.queue_entry_id = queue_item.id
  476. job.assigned_printer_id = printer_id # may be None for max_parallel
  477. job.status = "queued"
  478. job.dispatched_at = datetime.now(timezone.utc)
  479. run.status = "dispatching"
  480. await session.commit()
  481. await _publish_run_event(session, run)
  482. return slice_response.model_dump()
  483. return _orchestrate
  484. # ---------------------------------------------------------------------------
  485. # /slicer-pipelines/{id}/check-eligibility
  486. # ---------------------------------------------------------------------------
  487. @pipeline_run_create_router.post("/{pipeline_id}/check-eligibility", response_model=EligibilityReportResponse)
  488. async def check_eligibility(
  489. pipeline_id: int,
  490. body: CheckEligibilityRequest,
  491. _: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_READ),
  492. db: AsyncSession = Depends(get_db),
  493. ):
  494. pipeline = await _load_pipeline(db, pipeline_id)
  495. await _resolve_source(
  496. db,
  497. library_file_id=body.source_library_file_id,
  498. archive_id=body.source_archive_id,
  499. )
  500. if pipeline.target_kind == "printer_class" and pipeline.target_printer_id is None:
  501. report = await check_pipeline_eligibility(db, pipeline, status_lookup=_make_status_lookup())
  502. else:
  503. status = await _load_printer_status(pipeline.target_printer_id)
  504. report = await check_pipeline_eligibility(db, pipeline, status)
  505. return _serialise_status(report)
  506. # ---------------------------------------------------------------------------
  507. # /slicer-pipelines/{id}/run
  508. # ---------------------------------------------------------------------------
  509. @pipeline_run_create_router.post("/{pipeline_id}/run", response_model=PipelineRunResponse, status_code=202)
  510. async def run_pipeline(
  511. pipeline_id: int,
  512. body: PipelineRunCreateRequest,
  513. current_user: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_RUN),
  514. db: AsyncSession = Depends(get_db),
  515. ):
  516. from backend.app.api.routes.settings import get_setting
  517. from backend.app.services.slice_dispatch import slice_dispatch
  518. pipeline = await _load_pipeline(db, pipeline_id)
  519. src_kind, src_id, src_filename, src_path = await _resolve_source(
  520. db,
  521. library_file_id=body.source_library_file_id,
  522. archive_id=body.source_archive_id,
  523. )
  524. # Cap copies against the configured ceiling.
  525. raw_cap = await get_setting(db, "pipeline_max_copies")
  526. try:
  527. cap = int(raw_cap) if raw_cap else 50
  528. except (TypeError, ValueError):
  529. cap = 50
  530. if body.copies > cap:
  531. raise HTTPException(
  532. 422,
  533. f"copies={body.copies} exceeds pipeline_max_copies setting ({cap})",
  534. )
  535. # Eligibility pre-flight.
  536. if pipeline.target_kind == "printer_class" and pipeline.target_printer_id is None:
  537. report = await check_pipeline_eligibility(db, pipeline, status_lookup=_make_status_lookup())
  538. else:
  539. status = await _load_printer_status(pipeline.target_printer_id)
  540. report = await check_pipeline_eligibility(db, pipeline, status)
  541. if not report.ok and not body.force:
  542. raise HTTPException(status_code=409, detail=_serialise_status(report).model_dump())
  543. # Need a target — specific or class — to dispatch.
  544. if pipeline.target_printer_id is None and not pipeline.target_model_class:
  545. raise HTTPException(
  546. 400,
  547. "Pipeline has no target. Open the pipeline in Settings → Workflow → Pipelines and choose a target printer or printer class.",
  548. )
  549. run = PipelineRun(
  550. pipeline_id=pipeline.id,
  551. source_library_file_id=src_id if src_kind == "library_file" else None,
  552. source_archive_id=src_id if src_kind == "archive" else None,
  553. copies=body.copies,
  554. status="queued",
  555. eligibility_overridden=(not report.ok and body.force),
  556. created_by=current_user.id if current_user else None,
  557. )
  558. db.add(run)
  559. await db.flush()
  560. # One PipelineJob per copy. PR B was copies=1, PR C generalises.
  561. for i in range(body.copies):
  562. db.add(
  563. PipelineJob(
  564. pipeline_run_id=run.id,
  565. copy_index=i,
  566. status="pending",
  567. )
  568. )
  569. await db.commit()
  570. await db.refresh(run)
  571. await _publish_run_event(db, run)
  572. orchestrate = _make_orchestration_callable(
  573. run_id=run.id,
  574. pipeline_id=pipeline.id,
  575. src_kind=src_kind,
  576. src_id=src_id,
  577. src_filename=src_filename,
  578. src_path=src_path,
  579. creator_user_id=current_user.id if current_user else None,
  580. copies=body.copies,
  581. )
  582. slice_job = await slice_dispatch.enqueue(
  583. kind="library_file" if src_kind == "library_file" else "archive",
  584. source_id=src_id,
  585. source_name=src_filename,
  586. run=orchestrate,
  587. )
  588. run.slice_job_id = slice_job.id
  589. await db.commit()
  590. await db.refresh(run)
  591. return await _materialise_run(db, run)
  592. # ---------------------------------------------------------------------------
  593. # Lists, reads, cancel, retry-failed
  594. # ---------------------------------------------------------------------------
  595. @pipeline_run_create_router.get("/{pipeline_id}/runs", response_model=PipelineRunListResponse)
  596. async def list_runs_for_pipeline(
  597. pipeline_id: int,
  598. limit: int = 10,
  599. _: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_READ),
  600. db: AsyncSession = Depends(get_db),
  601. ):
  602. limit = max(1, min(limit, 100))
  603. rows = (
  604. (
  605. await db.execute(
  606. select(PipelineRun)
  607. .where(PipelineRun.pipeline_id == pipeline_id)
  608. .order_by(PipelineRun.id.desc())
  609. .limit(limit)
  610. )
  611. )
  612. .scalars()
  613. .all()
  614. )
  615. total = (
  616. await db.execute(select(func.count()).select_from(PipelineRun).where(PipelineRun.pipeline_id == pipeline_id))
  617. ).scalar() or 0
  618. return PipelineRunListResponse(
  619. runs=[await _materialise_run(db, r) for r in rows],
  620. total=total,
  621. )
  622. @pipeline_run_router.get("", response_model=PipelineRunListResponse)
  623. async def list_all_runs(
  624. limit: int = 25,
  625. offset: int = 0,
  626. pipeline_id: int | None = None,
  627. status: str | None = None,
  628. _: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_READ),
  629. db: AsyncSession = Depends(get_db),
  630. ):
  631. """Dashboard list. Newest first; filters on pipeline_id + status. The
  632. `status` filter matches the persisted snapshot, not the live roll-up —
  633. in-progress runs may appear under `dispatching` until the next state
  634. transition writes through."""
  635. limit = max(1, min(limit, 100))
  636. offset = max(0, offset)
  637. stmt = select(PipelineRun)
  638. count_stmt = select(func.count()).select_from(PipelineRun)
  639. if pipeline_id is not None:
  640. stmt = stmt.where(PipelineRun.pipeline_id == pipeline_id)
  641. count_stmt = count_stmt.where(PipelineRun.pipeline_id == pipeline_id)
  642. if status:
  643. stmt = stmt.where(PipelineRun.status == status)
  644. count_stmt = count_stmt.where(PipelineRun.status == status)
  645. rows = (await db.execute(stmt.order_by(desc(PipelineRun.id)).offset(offset).limit(limit))).scalars().all()
  646. total = (await db.execute(count_stmt)).scalar() or 0
  647. return PipelineRunListResponse(
  648. runs=[await _materialise_run(db, r) for r in rows],
  649. total=total,
  650. )
  651. @pipeline_run_router.get("/{run_id}", response_model=PipelineRunResponse)
  652. async def get_run(
  653. run_id: int,
  654. _: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_READ),
  655. db: AsyncSession = Depends(get_db),
  656. ):
  657. run = (await db.execute(select(PipelineRun).where(PipelineRun.id == run_id))).scalar_one_or_none()
  658. if run is None:
  659. raise HTTPException(404, "Pipeline run not found")
  660. return await _materialise_run(db, run)
  661. @pipeline_run_router.post("/{run_id}/cancel", response_model=PipelineRunResponse)
  662. async def cancel_run(
  663. run_id: int,
  664. _: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_RUN),
  665. db: AsyncSession = Depends(get_db),
  666. ):
  667. """Cancel a queued / in-flight run. Cascades to all non-terminal queue
  668. entries; in-flight prints continue on the printer (operator must Stop)."""
  669. run = (await db.execute(select(PipelineRun).where(PipelineRun.id == run_id))).scalar_one_or_none()
  670. if run is None:
  671. raise HTTPException(404, "Pipeline run not found")
  672. if run.status in ("completed", "failed", "cancelled", "partial_failure"):
  673. return await _materialise_run(db, run)
  674. run.status = "cancelled"
  675. run.completed_at = datetime.now(timezone.utc)
  676. if not run.error_message:
  677. run.error_message = "Cancelled by user"
  678. job_rows = (await db.execute(select(PipelineJob).where(PipelineJob.pipeline_run_id == run.id))).scalars().all()
  679. for job in job_rows:
  680. if job.queue_entry_id:
  681. queue_entry = (
  682. await db.execute(select(PrintQueueItem).where(PrintQueueItem.id == job.queue_entry_id))
  683. ).scalar_one_or_none()
  684. if queue_entry is not None and queue_entry.status in ("pending", "queued"):
  685. queue_entry.status = "cancelled"
  686. if job.status not in ("completed", "failed", "cancelled"):
  687. job.status = "cancelled"
  688. job.completed_at = datetime.now(timezone.utc)
  689. await db.commit()
  690. await db.refresh(run)
  691. await _publish_run_event(db, run)
  692. return await _materialise_run(db, run)
  693. @pipeline_run_router.post("/{run_id}/retry-failed", response_model=PipelineRunResponse, status_code=202)
  694. async def retry_failed(
  695. run_id: int,
  696. current_user: User | None = RequirePermissionIfAuthEnabled(Permission.PIPELINES_RUN),
  697. db: AsyncSession = Depends(get_db),
  698. ):
  699. """Create a new run with copies = (failed + cancelled count) from the
  700. parent. Same pipeline, same source. Eligibility re-checked at run time
  701. (it might pass this time — operator may have fixed the issue)."""
  702. parent = (await db.execute(select(PipelineRun).where(PipelineRun.id == run_id))).scalar_one_or_none()
  703. if parent is None:
  704. raise HTTPException(404, "Pipeline run not found")
  705. if parent.pipeline_id is None:
  706. raise HTTPException(400, "Original pipeline was deleted; cannot retry")
  707. if parent.source_library_file_id is None and parent.source_archive_id is None:
  708. raise HTTPException(400, "Original source was deleted; cannot retry")
  709. # Count the parent's failed + cancelled jobs.
  710. parent_jobs = (
  711. (await db.execute(select(PipelineJob).where(PipelineJob.pipeline_run_id == parent.id))).scalars().all()
  712. )
  713. fail_count = 0
  714. for j in parent_jobs:
  715. queue_entry = None
  716. if j.queue_entry_id:
  717. queue_entry = (
  718. await db.execute(select(PrintQueueItem).where(PrintQueueItem.id == j.queue_entry_id))
  719. ).scalar_one_or_none()
  720. live = _compute_job_status(j.status, queue_entry)
  721. if live in ("failed", "cancelled"):
  722. fail_count += 1
  723. if fail_count == 0:
  724. raise HTTPException(400, "No failed copies to retry")
  725. # Build the request payload the same way the user would have via /run.
  726. body = PipelineRunCreateRequest(
  727. source_library_file_id=parent.source_library_file_id,
  728. source_archive_id=parent.source_archive_id,
  729. copies=fail_count,
  730. force=True, # operator already accepted eligibility on the parent
  731. )
  732. # Reuse the run_pipeline route logic via a direct call — keeps the
  733. # orchestration single-sourced. The result inherits parent_run_id.
  734. new_run_response = await run_pipeline(parent.pipeline_id, body, current_user=current_user, db=db)
  735. # Stamp parent_run_id on the freshly-created run.
  736. new_row = (await db.execute(select(PipelineRun).where(PipelineRun.id == new_run_response.id))).scalar_one_or_none()
  737. if new_row is not None:
  738. new_row.parent_run_id = parent.id
  739. await db.commit()
  740. await db.refresh(new_row)
  741. return await _materialise_run(db, new_row)
  742. return new_run_response