test_scheduler_global_queue_order_3200.py 9.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251
  1. """One queue order for pinned and "Any <model>" jobs (#3200).
  2. Three P2S printers. The user dragged a job pinned to printer 1 above two
  3. "Any P2S" jobs, printer 1 finished, and it started one of the "Any" jobs from
  4. further down the queue. The pinned job sat at the top saying "Busy".
  5. The scheduler read the queue ``ORDER BY printer_id, position``, which put the
  6. lane ahead of the position. A model-based job has ``printer_id`` NULL, which
  7. SQLite sorts first and PostgreSQL sorts last, so which job won a printer both
  8. wanted was decided by the database, never by where the user put it.
  9. These run on SQLite, where the "Any" jobs used to win.
  10. """
  11. from contextlib import ExitStack
  12. from types import SimpleNamespace
  13. from unittest.mock import AsyncMock, MagicMock, patch
  14. import pytest
  15. from sqlalchemy import select
  16. from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
  17. import backend.app.models # noqa: F401 - populate Base.metadata
  18. from backend.app.core.database import Base
  19. from backend.app.models.library import LibraryFile
  20. from backend.app.models.print_queue import PrintQueueItem
  21. from backend.app.models.printer import Printer
  22. from backend.app.models.settings import Settings
  23. from backend.app.services.print_scheduler import PrintScheduler
  24. @pytest.fixture
  25. async def ctx():
  26. engine = create_async_engine("sqlite+aiosqlite:///:memory:", echo=False)
  27. async with engine.begin() as conn:
  28. await conn.run_sync(Base.metadata.create_all)
  29. session_maker = async_sessionmaker(engine, expire_on_commit=False)
  30. async with session_maker() as db:
  31. for pid in (1, 2, 3):
  32. db.add(
  33. Printer(
  34. id=pid,
  35. name=f"3D printer 0{pid}",
  36. serial_number=f"P2S000{pid}",
  37. ip_address=f"10.0.0.{pid}",
  38. access_code="x",
  39. model="P2S",
  40. is_active=True,
  41. )
  42. )
  43. await db.commit()
  44. try:
  45. yield SimpleNamespace(session_maker=session_maker)
  46. finally:
  47. await engine.dispose()
  48. async def _add(ctx, *, position, printer_id=None, target_model=None, print_time=None, manual_start=False):
  49. async with ctx.session_maker() as db:
  50. lib = LibraryFile(
  51. filename="job.gcode.3mf",
  52. file_path="/library/job.gcode.3mf",
  53. file_size=10,
  54. file_type="gcode.3mf",
  55. file_metadata={"sliced_for_model": "P2S"},
  56. )
  57. db.add(lib)
  58. await db.flush()
  59. item = PrintQueueItem(
  60. status="pending",
  61. position=position,
  62. printer_id=printer_id,
  63. target_model=target_model,
  64. library_file_id=lib.id,
  65. print_time_seconds=print_time,
  66. manual_start=manual_start,
  67. )
  68. db.add(item)
  69. await db.commit()
  70. return item.id
  71. async def _pinned(ctx, position, printer_id=1, **kw):
  72. return await _add(ctx, position=position, printer_id=printer_id, **kw)
  73. async def _any(ctx, position, **kw):
  74. return await _add(ctx, position=position, target_model="P2S", **kw)
  75. async def _set(ctx, key, value):
  76. async with ctx.session_maker() as db:
  77. db.add(Settings(key=key, value=value))
  78. await db.commit()
  79. async def _item(ctx, item_id):
  80. async with ctx.session_maker() as db:
  81. return (await db.execute(select(PrintQueueItem).where(PrintQueueItem.id == item_id))).scalar_one()
  82. async def _run(ctx, *, idle_printers):
  83. """One check_queue pass; returns {item_id: printer_id} for what went out."""
  84. scheduler = PrintScheduler()
  85. launched = MagicMock()
  86. patches = [
  87. patch("backend.app.services.print_scheduler.async_session", ctx.session_maker),
  88. patch("backend.app.core.database.async_session", ctx.session_maker),
  89. patch("backend.app.services.print_scheduler.printer_manager.is_connected", MagicMock(return_value=True)),
  90. patch("backend.app.services.print_scheduler.printer_manager.get_status", MagicMock(return_value=None)),
  91. patch(
  92. "backend.app.services.print_scheduler.printer_manager.is_awaiting_plate_clear",
  93. MagicMock(return_value=False),
  94. ),
  95. patch(
  96. "backend.app.services.print_scheduler.ha_sensor_manager.blocked_printers",
  97. AsyncMock(return_value={}),
  98. ),
  99. patch(
  100. "backend.app.services.notification_service.notification_service.on_queue_job_waiting",
  101. AsyncMock(),
  102. ),
  103. patch(
  104. "backend.app.services.notification_service.notification_service.on_queue_job_assigned",
  105. AsyncMock(),
  106. ),
  107. patch.object(
  108. scheduler,
  109. "_is_printer_idle",
  110. MagicMock(side_effect=lambda pid, *_a, **_k: pid in idle_printers),
  111. ),
  112. patch.object(scheduler, "_check_auto_drying", AsyncMock()),
  113. patch.object(scheduler, "_ensure_ams_mapping", AsyncMock(return_value=None)),
  114. patch.object(scheduler, "_block_on_filament_deficit", AsyncMock(return_value=False)),
  115. patch.object(scheduler, "_get_smart_plugs", AsyncMock(return_value=[])),
  116. patch.object(scheduler, "_launch_uploads", launched),
  117. ]
  118. with ExitStack() as stack:
  119. for p in patches:
  120. stack.enter_context(p)
  121. await scheduler.check_queue()
  122. if not launched.called:
  123. return {}
  124. dispatched = {}
  125. for item_id in launched.call_args[0][0]:
  126. dispatched[item_id] = (await _item(ctx, item_id)).printer_id
  127. return dispatched
  128. class TestPositionDecidesWhoGetsThePrinter:
  129. @pytest.mark.asyncio
  130. async def test_the_reporters_queue(self, ctx):
  131. """Pinned job dragged to the top, two "Any P2S" jobs below it, only
  132. printer 1 free: the pinned job starts."""
  133. shuttle_a = await _any(ctx, 2)
  134. shuttle_b = await _any(ctx, 3)
  135. coral = await _pinned(ctx, 1)
  136. assert await _run(ctx, idle_printers={1}) == {coral: 1}
  137. assert (await _item(ctx, shuttle_a)).status == "pending"
  138. assert (await _item(ctx, shuttle_b)).status == "pending"
  139. @pytest.mark.asyncio
  140. async def test_an_any_job_above_a_pinned_one_goes_first(self, ctx):
  141. """The other direction: position still decides, so a model-based job
  142. the user put higher takes the printer."""
  143. pinned = await _pinned(ctx, 2)
  144. any_job = await _any(ctx, 1)
  145. assert await _run(ctx, idle_printers={1}) == {any_job: 1}
  146. assert (await _item(ctx, pinned)).status == "pending"
  147. @pytest.mark.asyncio
  148. async def test_a_pinned_job_that_cannot_start_does_not_hold_up_other_printers(self, ctx):
  149. """The top job is pinned to a busy printer; the "Any" job below it
  150. still takes the free one."""
  151. await _pinned(ctx, 1, printer_id=2)
  152. any_job = await _any(ctx, 2)
  153. assert await _run(ctx, idle_printers={1}) == {any_job: 1}
  154. @pytest.mark.asyncio
  155. async def test_several_free_printers_each_get_the_next_job(self, ctx):
  156. coral = await _pinned(ctx, 1)
  157. shuttle_a = await _any(ctx, 2)
  158. shuttle_b = await _any(ctx, 3)
  159. dispatched = await _run(ctx, idle_printers={1, 2, 3})
  160. assert dispatched[coral] == 1
  161. assert {dispatched[shuttle_a], dispatched[shuttle_b]} == {2, 3}
  162. @pytest.mark.asyncio
  163. async def test_a_staged_job_at_the_top_does_not_block_the_printer(self, ctx):
  164. """A manual-start job waits for the user, so the next job takes the
  165. printer rather than the queue stalling behind it."""
  166. await _pinned(ctx, 1, manual_start=True)
  167. any_job = await _any(ctx, 2)
  168. assert await _run(ctx, idle_printers={1}) == {any_job: 1}
  169. class TestShortestJobFirstAcrossLanes:
  170. @pytest.mark.asyncio
  171. async def test_a_shorter_any_job_beats_a_longer_pinned_one(self, ctx):
  172. await _set(ctx, "queue_shortest_first", "true")
  173. pinned = await _pinned(ctx, 1, print_time=7200)
  174. short_any = await _any(ctx, 2, print_time=600)
  175. assert await _run(ctx, idle_printers={1}) == {short_any: 1}
  176. assert (await _item(ctx, pinned)).status == "pending"
  177. @pytest.mark.asyncio
  178. async def test_the_pinned_job_it_jumped_is_marked_and_goes_next(self, ctx):
  179. """The starvation guard has to look across lanes, or a stream of short
  180. "Any" jobs holds the pinned one back forever."""
  181. await _set(ctx, "queue_shortest_first", "true")
  182. pinned = await _pinned(ctx, 1, print_time=7200)
  183. await _any(ctx, 2, print_time=600)
  184. await _run(ctx, idle_printers={1})
  185. assert (await _item(ctx, pinned)).been_jumped is True
  186. await _any(ctx, 3, print_time=300)
  187. async with ctx.session_maker() as db:
  188. for row in (await db.execute(select(PrintQueueItem).where(PrintQueueItem.status != "pending"))).scalars():
  189. row.status = "completed"
  190. await db.commit()
  191. assert await _run(ctx, idle_printers={1}) == {pinned: 1}
  192. @pytest.mark.asyncio
  193. async def test_a_jumped_any_job_is_marked_when_a_pinned_job_goes_first(self, ctx):
  194. await _set(ctx, "queue_shortest_first", "true")
  195. long_any = await _any(ctx, 1, print_time=7200)
  196. short_pinned = await _pinned(ctx, 2, print_time=600)
  197. assert await _run(ctx, idle_printers={1}) == {short_pinned: 1}
  198. assert (await _item(ctx, long_any)).been_jumped is True
  199. @pytest.mark.asyncio
  200. async def test_an_any_job_for_another_model_is_not_marked(self, ctx):
  201. await _set(ctx, "queue_shortest_first", "true")
  202. other_model = await _add(ctx, position=1, target_model="X1C", print_time=7200)
  203. await _pinned(ctx, 2, print_time=600)
  204. await _run(ctx, idle_printers={1})
  205. assert (await _item(ctx, other_model)).been_jumped is False