| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518 |
- """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
|