|
|
@@ -0,0 +1,518 @@
|
|
|
+"""A failed upload must not eat the queue (#3210).
|
|
|
+
|
|
|
+The reporter queued 20 jobs for "any P2S". One P2S's file service was out of
|
|
|
+connection slots -- it answered port 990 with ``421 There are too many
|
|
|
+connections from your internet address`` -- so every upload to it failed. Each
|
|
|
+failure marked the item ``failed``, which left that printer idle, so the next
|
|
|
+pass handed it the next item. One item died about every nine seconds: 43 jobs
|
|
|
+gone in ten minutes, none printed on that printer.
|
|
|
+
|
|
|
+Two changes, pinned here:
|
|
|
+
|
|
|
+* An upload whose file never reached the printer (handshake refused, cool-off,
|
|
|
+ timeout, dropped connection) puts the item back in the queue. Failures that
|
|
|
+ would repeat on every retry -- a rejected access code, a full card -- still
|
|
|
+ fail it.
|
|
|
+* The printer is out of dispatch for a backoff window, so nothing else is
|
|
|
+ sent to it, and the item that was put back waits there with a reason.
|
|
|
+"""
|
|
|
+
|
|
|
+import asyncio
|
|
|
+import time
|
|
|
+from contextlib import ExitStack, asynccontextmanager
|
|
|
+from pathlib import Path
|
|
|
+from types import SimpleNamespace
|
|
|
+from unittest.mock import AsyncMock, MagicMock, patch
|
|
|
+
|
|
|
+import pytest
|
|
|
+from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
|
|
|
+
|
|
|
+import backend.app.models # noqa: F401 - populate Base.metadata
|
|
|
+import backend.app.services.archive as archive_module
|
|
|
+import backend.app.services.print_scheduler as scheduler_module
|
|
|
+from backend.app.core.database import Base
|
|
|
+from backend.app.models.archive import PrintArchive
|
|
|
+from backend.app.models.print_queue import PrintQueueItem
|
|
|
+from backend.app.models.printer import Printer
|
|
|
+from backend.app.services.bambu_ftp import BambuFTPClient, FtpFailure, FtpFailureKind, UploadCancelled
|
|
|
+from backend.app.services.print_scheduler import (
|
|
|
+ UPLOAD_FAILURE_BACKOFF_MAX_SECONDS,
|
|
|
+ UPLOAD_FAILURE_BACKOFF_SECONDS,
|
|
|
+ PrintScheduler,
|
|
|
+)
|
|
|
+
|
|
|
+pytestmark = pytest.mark.unit
|
|
|
+
|
|
|
+REFUSING_IP = "10.0.0.8"
|
|
|
+HEALTHY_IP = "10.0.0.9"
|
|
|
+
|
|
|
+
|
|
|
+@pytest.fixture
|
|
|
+async def farm(tmp_path):
|
|
|
+ """Two printers -- one whose file service refuses everything -- and an item factory."""
|
|
|
+ engine = create_async_engine("sqlite+aiosqlite:///:memory:", echo=False)
|
|
|
+ async with engine.begin() as conn:
|
|
|
+ await conn.run_sync(Base.metadata.create_all)
|
|
|
+ session_maker = async_sessionmaker(engine, expire_on_commit=False)
|
|
|
+
|
|
|
+ base_dir = tmp_path / "farm"
|
|
|
+ (base_dir / "archives").mkdir(parents=True, exist_ok=True)
|
|
|
+
|
|
|
+ async with session_maker() as db:
|
|
|
+ refusing = Printer(
|
|
|
+ name="P2S-8", serial_number="S8", ip_address=REFUSING_IP, access_code="12345678", model="P2S"
|
|
|
+ )
|
|
|
+ healthy = Printer(name="P2S-9", serial_number="S9", ip_address=HEALTHY_IP, access_code="12345678", model="P2S")
|
|
|
+ db.add_all([refusing, healthy])
|
|
|
+ await db.commit()
|
|
|
+ ids = SimpleNamespace(refusing=refusing.id, healthy=healthy.id)
|
|
|
+
|
|
|
+ counter = iter(range(1000))
|
|
|
+
|
|
|
+ async def add_item(*, printer_id: int | None = None, target_model: str | None = None) -> int:
|
|
|
+ n = next(counter)
|
|
|
+ async with session_maker() as db:
|
|
|
+ archive_rel = Path("archives") / f"job-{n}.3mf"
|
|
|
+ (base_dir / archive_rel).write_bytes(b"archive payload")
|
|
|
+ archive = PrintArchive(
|
|
|
+ printer_id=printer_id,
|
|
|
+ filename=f"job-{n}.3mf",
|
|
|
+ file_path=str(archive_rel),
|
|
|
+ file_size=15,
|
|
|
+ status="completed",
|
|
|
+ )
|
|
|
+ db.add(archive)
|
|
|
+ await db.flush()
|
|
|
+ item = PrintQueueItem(
|
|
|
+ printer_id=printer_id,
|
|
|
+ target_model=target_model,
|
|
|
+ archive_id=archive.id,
|
|
|
+ status="pending",
|
|
|
+ position=n,
|
|
|
+ )
|
|
|
+ db.add(item)
|
|
|
+ await db.commit()
|
|
|
+ return item.id
|
|
|
+
|
|
|
+ try:
|
|
|
+ yield SimpleNamespace(session_maker=session_maker, base_dir=base_dir, ids=ids, add_item=add_item)
|
|
|
+ finally:
|
|
|
+ await engine.dispose()
|
|
|
+
|
|
|
+
|
|
|
+def _upload_refused_by(ip: str, failure: FtpFailure):
|
|
|
+ """An upload that fails with *failure* on *ip* and succeeds everywhere else."""
|
|
|
+
|
|
|
+ async def _upload(ip_address, *_args, **kwargs):
|
|
|
+ if ip_address != ip:
|
|
|
+ return True
|
|
|
+ if kwargs.get("failure") is not None:
|
|
|
+ kwargs["failure"].failure = failure
|
|
|
+ return False
|
|
|
+
|
|
|
+ return _upload
|
|
|
+
|
|
|
+
|
|
|
+@asynccontextmanager
|
|
|
+async def _scheduler(
|
|
|
+ ctx,
|
|
|
+ upload,
|
|
|
+ *,
|
|
|
+ busy: set[int] | None = None,
|
|
|
+ scheduler: PrintScheduler | None = None,
|
|
|
+ waiting_notify: AsyncMock | None = None,
|
|
|
+):
|
|
|
+ """A scheduler wired to *ctx*'s database, with every printer idle unless in *busy*.
|
|
|
+
|
|
|
+ Pass the same *waiting_notify* to several passes to count notifications
|
|
|
+ across them.
|
|
|
+ """
|
|
|
+ scheduler = scheduler or PrintScheduler()
|
|
|
+ busy = busy if busy is not None else set()
|
|
|
+ failed_notify = AsyncMock()
|
|
|
+ waiting_notify = waiting_notify or AsyncMock()
|
|
|
+
|
|
|
+ def _real_spawn(coro, *, name=None):
|
|
|
+ return asyncio.create_task(coro, name=name)
|
|
|
+
|
|
|
+ patches = [
|
|
|
+ patch.object(scheduler_module.settings, "base_dir", ctx.base_dir),
|
|
|
+ patch.object(archive_module.settings, "base_dir", ctx.base_dir),
|
|
|
+ patch.object(archive_module.settings, "archive_dir", ctx.base_dir / "archive"),
|
|
|
+ patch("backend.app.services.print_scheduler.async_session", ctx.session_maker),
|
|
|
+ patch("backend.app.core.database.async_session", ctx.session_maker),
|
|
|
+ patch("backend.app.services.print_scheduler.printer_manager.is_connected", MagicMock(return_value=True)),
|
|
|
+ patch(
|
|
|
+ "backend.app.services.print_scheduler.printer_manager.get_status",
|
|
|
+ MagicMock(return_value=SimpleNamespace(state="IDLE", subtask_id=None, gcode_file=None, raw_data={})),
|
|
|
+ ),
|
|
|
+ patch(
|
|
|
+ "backend.app.services.print_scheduler.printer_manager.is_awaiting_plate_clear",
|
|
|
+ MagicMock(return_value=False),
|
|
|
+ ),
|
|
|
+ patch("backend.app.services.print_scheduler.printer_manager.start_print", MagicMock(return_value=True)),
|
|
|
+ patch("backend.app.services.print_scheduler.printer_manager.set_awaiting_plate_clear", MagicMock()),
|
|
|
+ patch("backend.app.services.print_scheduler.upload_file_async", upload),
|
|
|
+ patch("backend.app.services.print_scheduler.delete_file_async", AsyncMock(return_value=True)),
|
|
|
+ patch(
|
|
|
+ "backend.app.services.print_scheduler.get_ftp_retry_settings",
|
|
|
+ AsyncMock(return_value=(False, 0, 0, 1.0)),
|
|
|
+ ),
|
|
|
+ patch("backend.app.services.print_scheduler.cache_3mf_download", MagicMock()),
|
|
|
+ patch("backend.app.services.print_scheduler.spawn_background_task", _real_spawn),
|
|
|
+ patch("backend.app.services.notification_service.notification_service.on_queue_job_started", AsyncMock()),
|
|
|
+ patch("backend.app.services.notification_service.notification_service.on_queue_job_failed", failed_notify),
|
|
|
+ patch("backend.app.services.notification_service.notification_service.on_queue_job_assigned", AsyncMock()),
|
|
|
+ patch("backend.app.services.notification_service.notification_service.on_queue_job_waiting", waiting_notify),
|
|
|
+ patch("backend.app.services.mqtt_relay.mqtt_relay.on_queue_job_started", AsyncMock()),
|
|
|
+ patch.object(scheduler, "_is_printer_idle", MagicMock(side_effect=lambda pid, *_a, **_k: pid not in busy)),
|
|
|
+ patch.object(scheduler, "_ensure_ams_mapping", AsyncMock(return_value=None)),
|
|
|
+ patch.object(scheduler, "_block_on_filament_deficit", AsyncMock(return_value=False)),
|
|
|
+ patch.object(scheduler, "_propagate_owner_to_printer_manager", AsyncMock()),
|
|
|
+ patch.object(scheduler, "_power_off_if_needed", AsyncMock()),
|
|
|
+ patch.object(scheduler, "_preheat_and_soak", AsyncMock()),
|
|
|
+ patch.object(scheduler, "_check_auto_drying", AsyncMock()),
|
|
|
+ patch.object(scheduler, "_watchdog_print_start", AsyncMock()),
|
|
|
+ ]
|
|
|
+ with ExitStack() as stack:
|
|
|
+ for patcher in patches:
|
|
|
+ stack.enter_context(patcher)
|
|
|
+ scheduler.failed_notify = failed_notify
|
|
|
+ yield scheduler
|
|
|
+ tasks = [task for (task, _pid) in scheduler._inflight.values()]
|
|
|
+ if tasks:
|
|
|
+ await asyncio.gather(*tasks, return_exceptions=True)
|
|
|
+
|
|
|
+
|
|
|
+async def _item(ctx, item_id: int) -> PrintQueueItem:
|
|
|
+ async with ctx.session_maker() as db:
|
|
|
+ return await db.get(PrintQueueItem, item_id)
|
|
|
+
|
|
|
+
|
|
|
+HANDSHAKE = FtpFailure(FtpFailureKind.HANDSHAKE, "WRONG_VERSION_NUMBER (printer answered in cleartext: 421 ...)")
|
|
|
+
|
|
|
+
|
|
|
+# ---------------------------------------------------------------------------
|
|
|
+# What one failed upload does to its item
|
|
|
+# ---------------------------------------------------------------------------
|
|
|
+class TestOneFailedUpload:
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ @pytest.mark.parametrize(
|
|
|
+ "kind",
|
|
|
+ [FtpFailureKind.HANDSHAKE, FtpFailureKind.COOLOFF, FtpFailureKind.TIMEOUT, FtpFailureKind.NETWORK],
|
|
|
+ )
|
|
|
+ async def test_a_file_that_never_arrived_keeps_the_item(self, farm, kind):
|
|
|
+ item_id = await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+ upload = _upload_refused_by(REFUSING_IP, FtpFailure(kind, "detail"))
|
|
|
+
|
|
|
+ async with _scheduler(farm, upload) as scheduler:
|
|
|
+ await scheduler.check_queue()
|
|
|
+
|
|
|
+ item = await _item(farm, item_id)
|
|
|
+ assert item.status == "pending"
|
|
|
+ assert item.printer_id == farm.ids.refusing
|
|
|
+ assert item.error_message is None
|
|
|
+ assert item.dispatch_attempts == 0, "the start-watchdog budget is for a printer that took the file"
|
|
|
+ scheduler.failed_notify.assert_not_awaited()
|
|
|
+ assert farm.ids.refusing in scheduler._upload_backoff
|
|
|
+
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ @pytest.mark.parametrize(
|
|
|
+ "failure",
|
|
|
+ [
|
|
|
+ FtpFailure(FtpFailureKind.AUTH, "530 Login incorrect.", "530"),
|
|
|
+ FtpFailure(FtpFailureKind.STORAGE, "553 Could not create file.", "553"),
|
|
|
+ FtpFailure(FtpFailureKind.NOT_FOUND, "550 Permission denied.", "550"),
|
|
|
+ FtpFailure(FtpFailureKind.UNKNOWN, "500 what"),
|
|
|
+ None,
|
|
|
+ ],
|
|
|
+ ids=["auth", "storage", "not_found", "unknown", "unreported"],
|
|
|
+ )
|
|
|
+ async def test_a_failure_that_would_repeat_still_fails_the_item(self, farm, failure):
|
|
|
+ """Retrying a wrong access code or a full card every five minutes helps no one."""
|
|
|
+ item_id = await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+
|
|
|
+ async def upload(*_args, **kwargs):
|
|
|
+ if failure is not None:
|
|
|
+ kwargs["failure"].failure = failure
|
|
|
+ return False
|
|
|
+
|
|
|
+ async with _scheduler(farm, upload) as scheduler:
|
|
|
+ await scheduler.check_queue()
|
|
|
+
|
|
|
+ item = await _item(farm, item_id)
|
|
|
+ assert item.status == "failed"
|
|
|
+ assert item.error_message
|
|
|
+ scheduler.failed_notify.assert_awaited_once()
|
|
|
+ assert farm.ids.refusing not in scheduler._upload_backoff
|
|
|
+
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ async def test_an_upload_that_overran_its_deadline_still_fails(self, farm):
|
|
|
+ """A link too slow to finish would be just as slow next time (#2529)."""
|
|
|
+ item_id = await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+
|
|
|
+ async def upload(*_args, **kwargs):
|
|
|
+ kwargs["failure"].failure = FtpFailure(FtpFailureKind.TIMEOUT, "deadline")
|
|
|
+ raise UploadCancelled("too slow")
|
|
|
+
|
|
|
+ async with _scheduler(farm, upload) as scheduler:
|
|
|
+ await scheduler.check_queue()
|
|
|
+
|
|
|
+ assert (await _item(farm, item_id)).status == "failed"
|
|
|
+
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ async def test_a_cancel_during_the_upload_is_not_undone(self, farm):
|
|
|
+ """The row is never written back to pending, so a cancel that won stays won."""
|
|
|
+ item_id = await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+
|
|
|
+ async def upload(*_args, **kwargs):
|
|
|
+ async with farm.session_maker() as other:
|
|
|
+ row = await other.get(PrintQueueItem, item_id)
|
|
|
+ row.status = "cancelled"
|
|
|
+ await other.commit()
|
|
|
+ kwargs["failure"].failure = HANDSHAKE
|
|
|
+ return False
|
|
|
+
|
|
|
+ async with _scheduler(farm, upload) as scheduler:
|
|
|
+ await scheduler.check_queue()
|
|
|
+
|
|
|
+ assert (await _item(farm, item_id)).status == "cancelled"
|
|
|
+
|
|
|
+
|
|
|
+# ---------------------------------------------------------------------------
|
|
|
+# The report: a refusing printer must not drain the queue
|
|
|
+# ---------------------------------------------------------------------------
|
|
|
+class TestTheQueueIsNotDrained:
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ async def test_any_model_items_stop_going_to_the_refusing_printer(self, farm):
|
|
|
+ """#3210's shape: "any P2S" items, one P2S refusing every upload.
|
|
|
+
|
|
|
+ Pass 1 sends one item to each printer; the refusing one keeps its item.
|
|
|
+ Pass 2 has the healthy printer busy printing and the refusing one in
|
|
|
+ backoff, so the rest wait instead of being fed to it one by one.
|
|
|
+ """
|
|
|
+ ids = [await farm.add_item(target_model="P2S") for _ in range(5)]
|
|
|
+ upload = _upload_refused_by(REFUSING_IP, HANDSHAKE)
|
|
|
+ scheduler = PrintScheduler()
|
|
|
+
|
|
|
+ async with _scheduler(farm, upload, scheduler=scheduler):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ async with _scheduler(farm, upload, scheduler=scheduler, busy={farm.ids.healthy}):
|
|
|
+ for _ in range(3):
|
|
|
+ await scheduler.check_queue()
|
|
|
+
|
|
|
+ items = [await _item(farm, i) for i in ids]
|
|
|
+ assert [i.status for i in items].count("failed") == 0, [(i.id, i.status, i.error_message) for i in items]
|
|
|
+ assert [i.status for i in items].count("printing") == 1
|
|
|
+ held = [i for i in items if i.printer_id == farm.ids.refusing]
|
|
|
+ assert len(held) == 1, "exactly one item waits on the refusing printer"
|
|
|
+ assert held[0].status == "pending"
|
|
|
+ waiting = [i for i in items if i.status == "pending" and i.printer_id is None]
|
|
|
+ assert len(waiting) == 3, "the rest stay unassigned, free for any printer that comes free"
|
|
|
+
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ async def test_the_item_kept_on_the_printer_says_why(self, farm):
|
|
|
+ item_id = await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+ upload = _upload_refused_by(REFUSING_IP, HANDSHAKE)
|
|
|
+ scheduler = PrintScheduler()
|
|
|
+
|
|
|
+ async with _scheduler(farm, upload, scheduler=scheduler):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ async with _scheduler(farm, upload, scheduler=scheduler):
|
|
|
+ await scheduler.check_queue()
|
|
|
+
|
|
|
+ item = await _item(farm, item_id)
|
|
|
+ assert item.status == "pending"
|
|
|
+ assert item.waiting_reason == "P2S-8 is not accepting files — Bambuddy will retry automatically"
|
|
|
+
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ async def test_the_printer_is_tried_again_once_the_backoff_ends(self, farm):
|
|
|
+ item_id = await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+ scheduler = PrintScheduler()
|
|
|
+
|
|
|
+ async with _scheduler(farm, _upload_refused_by(REFUSING_IP, HANDSHAKE), scheduler=scheduler):
|
|
|
+ await scheduler.check_queue()
|
|
|
+
|
|
|
+ retry_at = scheduler._upload_backoff[farm.ids.refusing]
|
|
|
+ assert retry_at - time.monotonic() == pytest.approx(UPLOAD_FAILURE_BACKOFF_SECONDS, abs=5)
|
|
|
+
|
|
|
+ # The printer recovered and the window has passed.
|
|
|
+ scheduler._upload_backoff[farm.ids.refusing] = time.monotonic() - 1
|
|
|
+ async with _scheduler(farm, AsyncMock(return_value=True), scheduler=scheduler):
|
|
|
+ await scheduler.check_queue()
|
|
|
+
|
|
|
+ assert (await _item(farm, item_id)).status == "printing"
|
|
|
+ assert farm.ids.refusing not in scheduler._upload_backoff
|
|
|
+
|
|
|
+
|
|
|
+# ---------------------------------------------------------------------------
|
|
|
+# A printer that stays broken: retries thin out, and say so once
|
|
|
+# ---------------------------------------------------------------------------
|
|
|
+def _expire_backoff(scheduler: PrintScheduler, printer_id: int) -> None:
|
|
|
+ """Jump past the printer's backoff window, as if the time had passed."""
|
|
|
+ scheduler._upload_backoff[printer_id] = time.monotonic() - 1
|
|
|
+
|
|
|
+
|
|
|
+def _remaining(scheduler: PrintScheduler, printer_id: int) -> float:
|
|
|
+ return scheduler._upload_backoff[printer_id] - time.monotonic()
|
|
|
+
|
|
|
+
|
|
|
+async def _finish(ctx, scheduler: PrintScheduler, item_id: int) -> None:
|
|
|
+ """Mark a dispatched item's print done, so its printer can take the next one."""
|
|
|
+ async with ctx.session_maker() as db:
|
|
|
+ row = await db.get(PrintQueueItem, item_id)
|
|
|
+ assert row.status == "printing"
|
|
|
+ row.status = "completed"
|
|
|
+ await db.commit()
|
|
|
+ scheduler._release_dispatch_hold(row.printer_id)
|
|
|
+
|
|
|
+
|
|
|
+class TestAPrinterThatStaysBroken:
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ async def test_each_refusal_in_a_row_doubles_the_wait_up_to_the_cap(self, farm):
|
|
|
+ """#3210's printer refused for about 40 hours.
|
|
|
+
|
|
|
+ At a flat five minutes that is ~480 retries, each one a preheat cycle
|
|
|
+ where preheat is on, five connection attempts and a page of log.
|
|
|
+ """
|
|
|
+ await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+ upload = _upload_refused_by(REFUSING_IP, HANDSHAKE)
|
|
|
+ scheduler = PrintScheduler()
|
|
|
+
|
|
|
+ waits = []
|
|
|
+ for _ in range(6):
|
|
|
+ async with _scheduler(farm, upload, scheduler=scheduler):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ waits.append(_remaining(scheduler, farm.ids.refusing))
|
|
|
+ _expire_backoff(scheduler, farm.ids.refusing)
|
|
|
+
|
|
|
+ expected = [300, 600, 1200, 2400, 3600, 3600]
|
|
|
+ assert expected[0] == UPLOAD_FAILURE_BACKOFF_SECONDS
|
|
|
+ assert expected[-1] == UPLOAD_FAILURE_BACKOFF_MAX_SECONDS
|
|
|
+ assert waits == [pytest.approx(w, abs=5) for w in expected]
|
|
|
+
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ async def test_a_successful_upload_resets_the_wait(self, farm):
|
|
|
+ first = await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+ refused = _upload_refused_by(REFUSING_IP, HANDSHAKE)
|
|
|
+ scheduler = PrintScheduler()
|
|
|
+
|
|
|
+ for _ in range(3):
|
|
|
+ async with _scheduler(farm, refused, scheduler=scheduler):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ _expire_backoff(scheduler, farm.ids.refusing)
|
|
|
+ async with _scheduler(farm, AsyncMock(return_value=True), scheduler=scheduler):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ assert farm.ids.refusing not in scheduler._upload_refusals
|
|
|
+ await _finish(farm, scheduler, first)
|
|
|
+
|
|
|
+ # It breaks again later: back to the first window, not the fourth.
|
|
|
+ await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+ async with _scheduler(farm, refused, scheduler=scheduler, busy=set()):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ assert _remaining(scheduler, farm.ids.refusing) == pytest.approx(UPLOAD_FAILURE_BACKOFF_SECONDS, abs=5)
|
|
|
+
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ async def test_one_outage_sends_one_waiting_notification(self, farm):
|
|
|
+ """Every retry clears the waiting reason and the next refusal sets it again.
|
|
|
+
|
|
|
+ `hold_item` reads that as a new reason each time, and "Job Waiting" is
|
|
|
+ on by default -- so without a guard a 40-hour outage notifies ~480 times.
|
|
|
+ """
|
|
|
+ await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+ upload = _upload_refused_by(REFUSING_IP, HANDSHAKE)
|
|
|
+ scheduler = PrintScheduler()
|
|
|
+ waiting = AsyncMock()
|
|
|
+
|
|
|
+ for _ in range(4):
|
|
|
+ # The pass that dispatches and is refused...
|
|
|
+ async with _scheduler(farm, upload, scheduler=scheduler, waiting_notify=waiting):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ # ...and the pass that holds the item while the printer waits.
|
|
|
+ async with _scheduler(farm, upload, scheduler=scheduler, waiting_notify=waiting):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ _expire_backoff(scheduler, farm.ids.refusing)
|
|
|
+
|
|
|
+ assert waiting.await_count == 1
|
|
|
+ assert "not accepting files" in waiting.await_args.kwargs["waiting_reason"]
|
|
|
+
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ async def test_a_new_outage_after_a_recovery_notifies_again(self, farm):
|
|
|
+ first = await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+ refused = _upload_refused_by(REFUSING_IP, HANDSHAKE)
|
|
|
+ scheduler = PrintScheduler()
|
|
|
+ waiting = AsyncMock()
|
|
|
+
|
|
|
+ async def outage():
|
|
|
+ async with _scheduler(farm, refused, scheduler=scheduler, waiting_notify=waiting):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ async with _scheduler(farm, refused, scheduler=scheduler, waiting_notify=waiting):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ _expire_backoff(scheduler, farm.ids.refusing)
|
|
|
+
|
|
|
+ await outage()
|
|
|
+ async with _scheduler(farm, AsyncMock(return_value=True), scheduler=scheduler, waiting_notify=waiting):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ await _finish(farm, scheduler, first)
|
|
|
+ await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+ await outage()
|
|
|
+
|
|
|
+ assert waiting.await_count == 2
|
|
|
+
|
|
|
+ @pytest.mark.asyncio
|
|
|
+ async def test_keep_warm_does_not_hold_a_bed_for_a_printer_in_backoff(self, farm):
|
|
|
+ """There is no print coming to keep the bed warm for.
|
|
|
+
|
|
|
+ The printer is in busy_printers and has a pending item, which is
|
|
|
+ exactly what keep-warm looks for -- and each retry would restart its
|
|
|
+ time cap, so the bed could stay hot for the whole outage.
|
|
|
+ """
|
|
|
+ await farm.add_item(printer_id=farm.ids.refusing)
|
|
|
+ upload = _upload_refused_by(REFUSING_IP, HANDSHAKE)
|
|
|
+ scheduler = PrintScheduler()
|
|
|
+
|
|
|
+ async with _scheduler(farm, upload, scheduler=scheduler):
|
|
|
+ await scheduler.check_queue()
|
|
|
+ keep_warm = AsyncMock()
|
|
|
+ async with _scheduler(farm, upload, scheduler=scheduler):
|
|
|
+ with patch.object(scheduler, "_apply_keep_warm", keep_warm):
|
|
|
+ await scheduler.check_queue()
|
|
|
+
|
|
|
+ busy_seen = keep_warm.await_args.args[3]
|
|
|
+ assert farm.ids.refusing not in busy_seen
|
|
|
+
|
|
|
+
|
|
|
+# ---------------------------------------------------------------------------
|
|
|
+# 452 is the card, not the network
|
|
|
+# ---------------------------------------------------------------------------
|
|
|
+class TestA452IsAStorageReply:
|
|
|
+ def _upload_raising(self, error, tmp_path):
|
|
|
+ local = tmp_path / "job.3mf"
|
|
|
+ local.write_bytes(b"x" * 16)
|
|
|
+ client = BambuFTPClient(REFUSING_IP, "12345678")
|
|
|
+ client._ftp = MagicMock()
|
|
|
+ client._ftp.transfercmd.side_effect = error
|
|
|
+ assert client.upload_file(local, "/job.3mf") is False
|
|
|
+ return client.last_failure
|
|
|
+
|
|
|
+ def test_452_is_storage(self, tmp_path):
|
|
|
+ """ftplib raises it as error_temp, but it says the card is full.
|
|
|
+
|
|
|
+ Read as NETWORK it would be retried forever instead of failing with
|
|
|
+ the advice to check the card.
|
|
|
+ """
|
|
|
+ import ftplib # nosec B402 -- tests construct real ftplib error types
|
|
|
+
|
|
|
+ failure = self._upload_raising(ftplib.error_temp("452 Insufficient storage space."), tmp_path)
|
|
|
+ assert failure.kind is FtpFailureKind.STORAGE
|
|
|
+ assert failure.code == "452"
|
|
|
+
|
|
|
+ def test_other_transient_replies_are_still_network(self, tmp_path):
|
|
|
+ import ftplib # nosec B402 -- tests construct real ftplib error types
|
|
|
+
|
|
|
+ failure = self._upload_raising(ftplib.error_temp("421 There are too many connections."), tmp_path)
|
|
|
+ assert failure.kind is FtpFailureKind.NETWORK
|
|
|
+
|
|
|
+ def test_a_socket_error_is_never_read_as_a_reply_code(self, tmp_path):
|
|
|
+ failure = self._upload_raising(OSError("452 looks like a code but is not a reply"), tmp_path)
|
|
|
+ assert failure.kind is FtpFailureKind.NETWORK
|