mirror of
https://github.com/maziggy/bambuddy.git
synced 2026-10-04 05:01:37 +02:00
fix(asyncio): track strong refs on orphan create_task sites
asyncio holds only a weak reference to tasks returned by ``create_task``. Fire-and-forget callers that discard the return value let the event loop GC the task before it finishes, logging ``Task was destroyed but it is pending!`` with no traceback. The #1648 support-bundle review surfaced 94 such warnings in 8 days of v0.2.4.5 -- the silently-vanished exceptions reach support bundles as opaque GC notices instead of actionable errors. New backend/app/core/tasks.py::spawn_background_task(coro, *, name=None) is the one place in the codebase that calls asyncio.create_task. It stores the task in a module-level set, attaches a done-callback that auto-removes on completion AND surfaces any uncaught exception via the logger with the originating traceback, and accepts name= so a leak source is traceable through /tracebacks and the log line. Cancelled tasks don't log (a shutting-down service is not an error). Migrated the 16 truly-orphan create_task call sites to the helper: main.py (8): reconcile-stale, cooldown-poweroff, energy calc, smart-plug, maintenance-check, photo-then-notify, layer-timelapse, scan-timelapse, print-scheduler, notify-no-archive (the last one was hand-rolling the same pattern with task + no-op done_callback) printers.py:3123 apply-pa-after-refresh print_queue.py:1034 queue cooldown-poweroff firmware_update.py:261 firmware upload archive.py:1514 timelapse mp4 convert print_scheduler.py:2199 watchdog print-start library.py:1614 STL backfill smart_plugs.py:259 tasmota scan discovery.py:159 subnet scan smart_plug_manager.py x3 plug auto-off-pending background_dispatch.py x2 (lambda-wrapped inside loop.call_soon_threadsafe) upload progress Sites that already kept strong refs are unchanged: self._tasks.append(asyncio.create_task(...)) -- VP manager, tcp_proxy, mqtt_server self._x_task = asyncio.create_task(...) on service instances -- mqtt_bridge, obico_detection, github_backup, archive_purge, local_backup, library_trash, discovery service Locally assigned + awaited/gathered -- tcp_proxy bidirectional pumps, camera_fanout, slice_dispatch, slicer_api progress_task, manager._finish_release_task, main.py module-level cleanup loops
This commit is contained in:
@@ -8,6 +8,7 @@ All notable changes to Bambuddy will be documented in this file.
|
||||
- **VP access code is now auto-derived from the target printer in non-proxy modes (Discord report)** — A user on Discord set up a Queue-mode VP with a different access code than the real target printer and couldn't get the slicer to connect, even after the cert-trust path was sorted. Root cause: the live target-printer mirror that landed earlier in the 0.2.5 cycle forwards the slicer's MQTT/RTSPS auth bytes through to the real printer — the slicer holds **one** code in its profile (the one it bound the VP with), and that code has to pass two checks (VP listener, then real printer). If the codes diverge the bridge silently fails at the second hop and the slicer abandons the connection (e.g. opens 8883, FINs before sending a ClientHello). The wiki *did* document a code-match requirement but framed it as a camera-only concern (`MQTT and FTP work either way; only the camera path needs the match`) — wrong, all bridged protocols inherit. **The fix removes the foot-gun rather than re-document it.** When a target printer is selected on a non-proxy VP (Archive / Review / Queue), the access-code field in the VP card switches to a read-only display showing the target's code with an Eye-toggle reveal, and the backend auto-inherits the value on every `create` / `update` (any explicit `access_code` submitted alongside a target is silently overridden — belt-and-braces for non-UI clients). When no target is set, the field stays editable as before. The same `inheritsAccessCodeFromTarget` predicate gates a small "Inherited from target" badge in place of the existing `isSet` / `notSet` status pill. Changing the target after the slicer has already bound triggers an info toast ("Access code now matches the new target — re-add this device in your slicer") because the slicer's stored code is now stale. **One-shot startup migration** in `core/database.py` corrects any pre-existing mismatched VPs on first boot after the upgrade: SELECTs the diverged rows for an INFO log per VP (`VP 'Workshop Queue' (id=3) access code synced from target printer 'X1C #2'` — audit trail for anyone digging through logs), then UPDATEs via correlated subquery (idempotent — the WHERE clause excludes already-synced rows, so re-running is a no-op; portable across SQLite and Postgres). No user-facing banner because there's no action for the user to take — the fix is done, and a previously-stuck bridge now works. **Wiki**: `features/virtual-printer.md` line 1189 flipped from the wrong MQTT/FTP-work-either-way claim to "the bridge forwards slicer auth bytes through; Bambuddy auto-derives so the codes can't diverge", the line-84 tip's "for camera" framing replaced with the broader rule, and the port-table row for RTSP `:322` annotated with "transparent passthrough to the real printer's `:322`, same end-to-end TLS as proxy mode" so the dedicated-bind-IP-vs-passthrough-to-printer apparent contradiction reads as one consistent model. **i18n**: 5 new keys (`accessCode.inheritedFromTarget`, `accessCode.derivedFromTargetHint`, `accessCode.reveal`, `accessCode.hide`, `toast.targetCodeChangedRebind`) translated in all 11 locales (de/en/es/fr/it/ja/ko/pt-BR/tr/zh-CN/zh-TW), no English fallbacks per the project's hard rule.
|
||||
|
||||
### Fixed
|
||||
- **Background asyncio tasks no longer get garbage-collected mid-flight (#1648 follow-up)** — Support-bundle review under #1648 surfaced 94 `Task was destroyed but it is pending!` warnings in 8 days of v0.2.4.5. **Root cause:** asyncio holds only a weak reference to the result of `create_task` — any "fire and forget" call site that doesn't store the returned task lets the event loop GC the task before it finishes. The warning gives no traceback, so the originating exception (if any) vanishes silently into a support bundle that looks scary but isn't actionable. **Fix:** new `backend/app/core/tasks.py::spawn_background_task(coro, *, name=None)` helper that stores a strong reference in a module-level set, attaches a done-callback that auto-removes on completion AND surfaces any uncaught exception via the logger with the originating traceback, and accepts a `name=` argument so a leak source is traceable in `/tracebacks` and the log line. **Migration:** the 16 truly-orphan `asyncio.create_task(...)` call sites — across `main.py` (8), `printers.py`, `print_queue.py`, `firmware_update.py`, `archive.py`, `print_scheduler.py`, `library.py`, `smart_plugs.py`, `discovery.py`, `smart_plug_manager.py` (3), and `background_dispatch.py` (2 lambda-wrapped) — switched to `spawn_background_task`. Other `create_task` sites already kept strong refs via `self._tasks.append(...)`, `self._x_task = ...`, or local `await`/`gather` and stay unchanged. **Tests:** 5 unit cases in `test_tasks.py` pin the contract — strong-ref retention through completion, set-shrinkage after done, uncaught exception logged at WARNING with `exc_info`, cancellation does not log (a shutting-down service is not an error), and named tasks propagate `name=`. Net result: support bundles stop showing the opaque GC warnings, and any silent fire-and-forget exception now reaches the logger with a traceback attached. Severity reclassification of unrelated noise (the 791 "Failed to get cloud preset 400" spam, the `bambu_cloud.Login failed` mis-ERRORs, etc.) is a separate follow-up.
|
||||
- **Home-page filament assign no longer leaves the slicer unaware of PFCN cloud presets (#1648, reported by @ferch-G)** — Reporter on an H2D with a Polymaker spool noticed that assigning the spool from the Dashboard left the slicer's filament dropdown showing "unknown" / generic, but clicking Configure right after made the slicer recognize it correctly — "Configure" felt like a mandatory follow-up step rather than a refinement. **Root cause: PFCN-prefix cloud preset IDs were never handled.** Bambu's cloud uses three preset-ID shapes: `GFS…` (official Bambu), `PFUS…` (cloud user-created), and `PFCN…` (cloud shared / partner-uploaded — e.g. Polymaker's "(Custom)" Bambu Lab H2D variants like the reporter's `PFCN80e80c1f79db85`). `apply_spool_to_slot_via_mqtt` only routed `GFS` and `PFUS` through the cloud-detail lookup that extracts the real `filament_id`. PFCN slipped past the cloud-lookup branch, fell into the local-preset `int()` parse path, raised ValueError, dropped into `normalize_slicer_filament` which returns any `P`-prefix unchanged, and the raw PFCN landed in `tray_info_idx` — which the printer's calibration table can't index, so the slicer rendered "unknown". The Configure modal rescued each assign because it does its own `getCloudSettingDetail` lookup and writes the resolved `filament_id`. **Fix:** extend the cloud-detail-lookup branch (`inventory.py:129`) and the discard safety net (`inventory.py:223`) to include `PFCN` alongside `GFS`/`PFUS`. After the fix, the same three paths work: cloud-authenticated → real `filament_id` from `detail["filament_id"]` ships as `tray_info_idx` (Polymaker PLA Matte resolves to `GFL05`); cloud unavailable → raw PFCN discarded, slot reuses an existing valid P-prefix preset if material matches; nothing else available → falls through to the spool's generic material id (`PLA → GFL99`). Source comment now lists all three cloud-ID shapes so the next time Bambu invents a new prefix (PFXX, PFYY, …) the maintainer doesn't have to re-derive the structure from a bug report. **Tests:** 3 new integration cases in `test_inventory_assign.py::TestAssignSpoolPfcnCloudPreset` — falls back to generic when cloud unavailable (and pins the no-PFCN-leak invariant), reuses an existing slot's valid P-prefix preset when material matches, and the happy-path cloud lookup that produces a resolved `filament_id` while preserving the original PFCN as `setting_id`. Existing 28 assign-flow tests stay green.
|
||||
- **Bambu cloud A1 Mini filament / process profiles no longer hidden in AMS slot picker (#1649, root-caused by @technopaw)** — Reporter on an A1 Mini observed that the AMS slot Configure dropdown showed no Bambu / Generic filament profiles; only user-authored profiles surfaced. Mirror in the Profiles tab: filtering by "A1 Mini" left only A1 (non-mini) results. **Root cause: Bambu rolled out a profile rename mid-2026.** The `@BBL <code>` suffix on cloud profiles shifted from the long display form to a terse model code — `Bambu PLA Basic @BBL A1 Mini ...` is now `Bambu PLA Basic @BBL A1M ...` across 106 cloud profiles. User-authored profiles still use the long form (which is why the reporter's custom A1 Mini profile worked, and Bambu PLA Basic happened to render via the `localPreset` always-shown path). Bambuddy's filter compared the extracted token verbatim against the display name (`"A1M".toUpperCase() === "A1 MINI"` is false), so the rename silently stripped every newly-renamed profile from the picker. **Fix: centralized alias-aware match in `frontend/src/utils/slicerPrinterMatch.ts`.** New `PRINTER_MODEL_SUFFIX_ALIASES` table holds the bidirectional `A1 Mini` ⇄ `A1M` mapping (uppercase-normalised, narrow on purpose — wide-net aliasing like `X1` ⇄ `X1C` would silently group truly distinct printers); exported `matchesPrinterModelSuffix(presetSuffix, printerModel)` helper does the case-insensitive compare with alias fallback. Both consumer sites switched to the helper: `ConfigureAmsSlotModal.tsx:586,607` (the AMS slot picker, hit directly + reached from SpoolBuddy's AMS page via `mapModelCode(printer?.model)`), and `slicerPrinterMatch.ts:classifyByBambuName` (the SliceModal Process / Filament compatibility check). Backend `PRINTER_MODEL_MAP` also gains a `Bambu Lab A1M` → `A1 Mini` entry so server-side 3MF printer-model normalization stays consistent if a future 3MF embeds the short form. The structure stays open: when Bambu introduces the next rename, it's a single new row in the alias table — `/api/v1/cloud/settings` is the place to grep, called out in the source comment. **Tests**: 7 new unit cases in `slicerPrinterMatch.test.ts` pin the alias helper (canonical, case-insensitive both directions, A1M ↔ A1 Mini in both orientations, A1M does NOT collapse to A1, A1 does NOT collapse to A1 Mini, unrelated models reject) plus 3 integration cases in `presetCompatibility` (cloud filament `@BBL A1M` matches A1 Mini, cloud process `@BBL A1M` matches A1 Mini, `@BBL A1M` does NOT match A1). 2 new component-level cases in `ConfigureAmsSlotModal.test.tsx`: `@BBL A1M` cloud preset surfaces when picker is for A1 Mini (with `@BBL A1` correctly filtered out), and `@BBL X1C` stays filtered out when picker is for A1 Mini (sanity check against accidental widening). All existing 2062 vitest cases stay green.
|
||||
- **VP Queue / Archive / Review: Bambu Studio 2.7.x stayed stuck at "Downloading" after Send (#1658, reported by @IndividualGhost1905)** — Reporter on Bambu Studio 2.7.1.57 + X1C reported that sending a model to a Queue-mode VP (with Auto-Dispatch off) left the slicer's send modal stuck at "Downloading" forever; clicking Delete on the queued item didn't release it, and even Auto-Dispatch ON + a successful real print didn't release it. Only toggling the VP off/on cleared the slicer. The deleted-from-queue framing is a red herring — the slicer was stuck *before* deletion, the user just noticed it most when they deleted. **Root cause: the #1280 fix assumed the wrong event order.** The original assumption was MQTT `project_file` → FTP upload → set `gcode_state=FINISH`, and the slicer's "Downloading" UI releases on FINISH. Bambu Studio 2.7.x flipped the Send sequence to FTP `verify_job` → FTP `.3mf` → MQTT `project_file`, so on_file_received's `set_gcode_state("FINISH", …)` fires *first*, then the synthetic `_send_print_response` ack runs and overwrites `_gcode_state` back to `"PREPARE"`. From that point the 1 Hz cached-as-base push stream carries PREPARE forever, the slicer waits for the FINISH transition it'll never see, and the modal sits stuck. **Auto-Dispatch ON is the same bug**: the real printer's gcode_state goes PREPARE → RUNNING → FINISH on its bridge, but `_send_status_report` overrides the cached push's `gcode_state` with the local `_gcode_state` (still PREPARE), so the real state changes never reach the slicer. The fix re-fires `set_gcode_state("FINISH", filename, prepare_percent="100")` from `on_print_command` 1.5 s after the synthetic ack, for every non-proxy mode (queue / archive / review). The 1.5 s window is long enough for the slicer's modal to see at least one PREPARE push on the 1 Hz cycle (so the transition reads as PREPARE → FINISH, matching what the slicer expects) and short enough that the modal feels responsive. Proxy mode is exempt — there the real printer drives the bridge state and a synthetic FINISH would clobber a real PREPARE/RUNNING transition. The scheduler cancels any in-flight timer when a new project_file lands so a slicer that retries doesn't end with two competing FINISH timers. **Tests**: 6 new cases in `test_virtual_printer.py` — schedules on archive (and by extension queue/review), proxy mode does NOT schedule, no-MQTT skip is silent, second `project_file` cancels the first timer, delayed run sets the expected `(state, filename, prepare_percent)` triple, empty filename does not schedule.
|
||||
|
||||
@@ -12,6 +12,7 @@ from pydantic import BaseModel
|
||||
|
||||
from backend.app.core.auth import RequirePermissionIfAuthEnabled
|
||||
from backend.app.core.permissions import Permission
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.models.user import User
|
||||
from backend.app.services.discovery import (
|
||||
discovery_service,
|
||||
@@ -154,9 +155,10 @@ async def start_subnet_scan(
|
||||
request: Subnet to scan in CIDR notation (e.g., "192.168.1.0/24")
|
||||
"""
|
||||
# Start scan in background
|
||||
import asyncio
|
||||
|
||||
asyncio.create_task(subnet_scanner.scan_subnet(request.subnet, request.timeout))
|
||||
spawn_background_task(
|
||||
subnet_scanner.scan_subnet(request.subnet, request.timeout),
|
||||
name=f"subnet-scan-{request.subnet}",
|
||||
)
|
||||
|
||||
# Return immediate status
|
||||
scanned, total = subnet_scanner.progress
|
||||
|
||||
@@ -1,6 +1,5 @@
|
||||
"""API routes for File Manager (Library) functionality."""
|
||||
|
||||
import asyncio
|
||||
import base64
|
||||
import binascii
|
||||
import contextlib
|
||||
@@ -30,6 +29,7 @@ from backend.app.core.auth import (
|
||||
from backend.app.core.config import settings as app_settings
|
||||
from backend.app.core.database import async_session, get_db
|
||||
from backend.app.core.permissions import Permission
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.models.archive import PrintArchive
|
||||
from backend.app.models.library import LibraryFile, LibraryFolder
|
||||
from backend.app.models.print_queue import PrintQueueItem
|
||||
@@ -1611,7 +1611,7 @@ async def scan_external_folder(
|
||||
# folder_cache.values() covers the root + every pre-existing subfolder
|
||||
# + every subfolder created during this scan. all_folder_ids on its own
|
||||
# would miss the newly-created ones (it's snapshotted before the walk).
|
||||
asyncio.create_task(
|
||||
spawn_background_task(
|
||||
_backfill_external_stl_thumbnails(list(set(folder_cache.values()))),
|
||||
name=f"stl-backfill-folder-{folder_id}",
|
||||
)
|
||||
|
||||
@@ -16,6 +16,7 @@ from backend.app.core.auth import RequirePermissionIfAuthEnabled, require_owners
|
||||
from backend.app.core.config import settings
|
||||
from backend.app.core.database import get_db
|
||||
from backend.app.core.permissions import Permission
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.models.archive import PrintArchive
|
||||
from backend.app.models.library import LibraryFile
|
||||
from backend.app.models.print_batch import PrintBatch
|
||||
@@ -961,7 +962,6 @@ async def stop_queue_item(
|
||||
_: User | None = RequirePermissionIfAuthEnabled(Permission.QUEUE_UPDATE_ALL),
|
||||
):
|
||||
"""Stop an actively printing queue item."""
|
||||
import asyncio
|
||||
|
||||
from backend.app.models.smart_plug import SmartPlug
|
||||
from backend.app.services.printer_manager import printer_manager
|
||||
@@ -1031,7 +1031,7 @@ async def stop_queue_item(
|
||||
logger.info("Auto-off: Powering off printer %s", printer_id)
|
||||
await tasmota_service.turn_off(plug)
|
||||
|
||||
asyncio.create_task(cooldown_and_poweroff())
|
||||
spawn_background_task(cooldown_and_poweroff(), name=f"queue-cooldown-poweroff-{printer_id}")
|
||||
|
||||
return {"message": "Print stopped" if stop_sent else "Queue item cancelled (printer was offline)"}
|
||||
|
||||
|
||||
@@ -12,6 +12,7 @@ from backend.app.core.auth import RequireCameraStreamTokenIfAuthEnabled, Require
|
||||
from backend.app.core.config import settings
|
||||
from backend.app.core.database import get_db
|
||||
from backend.app.core.permissions import Permission
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.models.ams_label import AmsLabel
|
||||
from backend.app.models.printer import Printer
|
||||
from backend.app.models.slot_preset import SlotPresetMapping
|
||||
@@ -3120,7 +3121,10 @@ async def refresh_ams_slot(
|
||||
raise HTTPException(400, message)
|
||||
|
||||
# Apply PA profile after delay (RFID re-read takes a few seconds)
|
||||
asyncio.create_task(_apply_pa_after_refresh(printer_id, ams_id, slot_id))
|
||||
spawn_background_task(
|
||||
_apply_pa_after_refresh(printer_id, ams_id, slot_id),
|
||||
name=f"apply-pa-after-refresh-{printer_id}-{ams_id}-{slot_id}",
|
||||
)
|
||||
|
||||
return {"success": True, "message": message}
|
||||
|
||||
|
||||
@@ -12,6 +12,7 @@ from backend.app.api.routes.settings import get_setting
|
||||
from backend.app.core.auth import RequirePermissionIfAuthEnabled
|
||||
from backend.app.core.database import get_db
|
||||
from backend.app.core.permissions import Permission
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.models.printer import Printer
|
||||
from backend.app.models.smart_plug import SmartPlug
|
||||
from backend.app.models.user import User
|
||||
@@ -249,14 +250,16 @@ async def start_tasmota_scan(
|
||||
|
||||
Auto-detects local network if no IP range provided.
|
||||
"""
|
||||
import asyncio
|
||||
|
||||
# Auto-detect network
|
||||
from_ip, to_ip = get_local_network_range()
|
||||
timeout = request.timeout if request else 1.0
|
||||
|
||||
# Start scan in background
|
||||
asyncio.create_task(tasmota_scanner.scan_range(from_ip, to_ip, timeout))
|
||||
spawn_background_task(
|
||||
tasmota_scanner.scan_range(from_ip, to_ip, timeout),
|
||||
name="tasmota-scan",
|
||||
)
|
||||
|
||||
# Return immediate status
|
||||
scanned, total = tasmota_scanner.progress
|
||||
|
||||
@@ -0,0 +1,85 @@
|
||||
"""Background-task helper that keeps a strong reference to fire-and-forget tasks.
|
||||
|
||||
asyncio holds only a weak reference to tasks returned by ``create_task`` --
|
||||
when the caller discards the return value (the "fire and forget" pattern),
|
||||
the task can be garbage-collected mid-execution and the event loop logs
|
||||
``Task was destroyed but it is pending!`` with no traceback. A support
|
||||
bundle review under #1648 surfaced 94 such warnings in 8 days of v0.2.4.5.
|
||||
|
||||
``spawn_background_task`` is the one place in the codebase that calls
|
||||
``asyncio.create_task``. It stores the task in a module-level set, removes
|
||||
it when the task completes, and surfaces any uncaught exception through
|
||||
the logger so a silently-swallowed error becomes a visible WARNING with
|
||||
the originating traceback instead of an opaque GC warning.
|
||||
|
||||
Use this for any work that should run in the background without being
|
||||
awaited inline. For tasks that the service owns and needs to cancel on
|
||||
shutdown, store the returned ``asyncio.Task`` on the service instance
|
||||
instead (the helper still adds the strong reference, so storing it twice
|
||||
is redundant but harmless).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from collections.abc import Coroutine
|
||||
from typing import Any
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Strong-reference holder. Tasks live here from creation through completion.
|
||||
# Module-level so the set survives across spawn calls; the done-callback
|
||||
# removes each task as it finishes so the set doesn't grow without bound
|
||||
# (the event loop's GC can't reap an entry the callback still holds, but
|
||||
# the discard breaks the cycle immediately).
|
||||
_background_tasks: set[asyncio.Task[Any]] = set()
|
||||
|
||||
|
||||
def spawn_background_task(
|
||||
coro: Coroutine[Any, Any, Any],
|
||||
*,
|
||||
name: str | None = None,
|
||||
) -> asyncio.Task[Any]:
|
||||
"""Schedule ``coro`` on the running loop without losing the task reference.
|
||||
|
||||
Args:
|
||||
coro: The coroutine to run. Must not already be a Task.
|
||||
name: Optional task name surfaced in /tracebacks and the
|
||||
done-callback log line so a leaked task is traceable to its
|
||||
spawn site.
|
||||
|
||||
Returns:
|
||||
The created ``asyncio.Task``. Most callers ignore it -- the helper
|
||||
keeps its own strong reference. Callers that need to ``await`` or
|
||||
cancel later can store it on a service instance.
|
||||
"""
|
||||
task = asyncio.create_task(coro, name=name)
|
||||
_background_tasks.add(task)
|
||||
task.add_done_callback(_on_task_done)
|
||||
return task
|
||||
|
||||
|
||||
def _on_task_done(task: asyncio.Task[Any]) -> None:
|
||||
"""Discard the strong reference and surface any uncaught exception.
|
||||
|
||||
Without this, an exception raised inside a fire-and-forget task is
|
||||
silently retrieved by ``Task.__del__`` and never reaches the logger.
|
||||
Surface it here as a WARNING with the task name so support bundles
|
||||
capture the originating error instead of an opaque GC notice.
|
||||
"""
|
||||
_background_tasks.discard(task)
|
||||
if task.cancelled():
|
||||
return
|
||||
exc = task.exception()
|
||||
if exc is not None:
|
||||
logger.warning(
|
||||
"Background task %r raised an uncaught exception",
|
||||
task.get_name(),
|
||||
exc_info=exc,
|
||||
)
|
||||
|
||||
|
||||
def active_task_count() -> int:
|
||||
"""Number of background tasks currently in flight. Used by tests."""
|
||||
return len(_background_tasks)
|
||||
+21
-12
@@ -72,6 +72,7 @@ from backend.app.api.routes.maintenance import _get_printer_maintenance_internal
|
||||
from backend.app.api.routes.support import init_debug_logging
|
||||
from backend.app.core.config import APP_VERSION, settings as app_settings
|
||||
from backend.app.core.database import async_session, engine, init_db
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.core.websocket import ws_manager
|
||||
from backend.app.models.smart_plug import SmartPlug
|
||||
from backend.app.services.archive import ArchiveService, peek_plate_index_in_3mf, swap_plate_suffix
|
||||
@@ -823,7 +824,10 @@ async def on_printer_status_change(printer_id: int, state: PrinterState):
|
||||
# the same connection don't re-trigger reconciliation.
|
||||
if state.connected and not _printer_reconciled_since_connect.get(printer_id, False):
|
||||
_printer_reconciled_since_connect[printer_id] = True
|
||||
asyncio.create_task(reconcile_stale_active_prints(printer_id))
|
||||
spawn_background_task(
|
||||
reconcile_stale_active_prints(printer_id),
|
||||
name=f"reconcile-stale-prints-{printer_id}",
|
||||
)
|
||||
elif not state.connected and _printer_reconciled_since_connect.get(printer_id, False):
|
||||
# Re-arm so the next reconnect triggers reconciliation again.
|
||||
_printer_reconciled_since_connect[printer_id] = False
|
||||
@@ -3839,7 +3843,10 @@ async def on_print_complete(printer_id: int, data: dict):
|
||||
except Exception as e:
|
||||
logger.warning("Failed to power off plug %s for printer %s: %s", plug_id, pid, e)
|
||||
|
||||
asyncio.create_task(cooldown_and_poweroff(printer_id, [p.id for p in enabled_plugs]))
|
||||
spawn_background_task(
|
||||
cooldown_and_poweroff(printer_id, [p.id for p in enabled_plugs]),
|
||||
name=f"cooldown-poweroff-{printer_id}",
|
||||
)
|
||||
except Exception as e:
|
||||
logging.getLogger(__name__).warning(f"Queue item update failed: {e}")
|
||||
|
||||
@@ -4024,8 +4031,7 @@ async def on_print_complete(printer_id: int, data: dict):
|
||||
except Exception as e:
|
||||
logger.warning("[NOTIFY-BG] Failed to send notification without archive: %s", e, exc_info=True)
|
||||
|
||||
task = asyncio.create_task(_notify_no_archive())
|
||||
task.add_done_callback(lambda _t: None)
|
||||
spawn_background_task(_notify_no_archive(), name="notify-no-archive")
|
||||
return
|
||||
|
||||
log_timing("Archive lookup")
|
||||
@@ -4363,9 +4369,9 @@ async def on_print_complete(printer_id: int, data: dict):
|
||||
logger.warning("[PHOTO-BG] Failed: %s", e)
|
||||
return None
|
||||
|
||||
asyncio.create_task(_background_energy_calculation())
|
||||
spawn_background_task(_background_energy_calculation(), name="background-energy-calc")
|
||||
# Photo capture task - result will be used by notifications
|
||||
photo_task = asyncio.create_task(_background_finish_photo())
|
||||
photo_task = spawn_background_task(_background_finish_photo(), name="background-finish-photo")
|
||||
log_timing("Background tasks scheduled (energy, photo)")
|
||||
|
||||
# Also run smart plug, notifications, and maintenance as background tasks
|
||||
@@ -4544,8 +4550,8 @@ async def on_print_complete(printer_id: int, data: dict):
|
||||
except Exception as e:
|
||||
logger.warning("[MAINT-BG] Failed: %s", e)
|
||||
|
||||
asyncio.create_task(_background_smart_plug())
|
||||
asyncio.create_task(_background_maintenance_check())
|
||||
spawn_background_task(_background_smart_plug(), name="background-smart-plug")
|
||||
spawn_background_task(_background_maintenance_check(), name="background-maintenance-check")
|
||||
|
||||
# Notification task waits for photo capture to complete first (with timeout).
|
||||
# When a timelapse was recording, photo sourcing polls the per-print
|
||||
@@ -4572,7 +4578,7 @@ async def on_print_complete(printer_id: int, data: dict):
|
||||
except Exception as e:
|
||||
logger.error("[PHOTO-NOTIFY] Notification sending failed: %s", e, exc_info=True)
|
||||
|
||||
asyncio.create_task(_photo_then_notify())
|
||||
spawn_background_task(_photo_then_notify(), name="photo-then-notify")
|
||||
|
||||
# Stitch external camera layer timelapse if session was active
|
||||
print_status = data.get("status", "completed")
|
||||
@@ -4611,7 +4617,7 @@ async def on_print_complete(printer_id: int, data: dict):
|
||||
except Exception:
|
||||
pass # Best-effort timelapse session cancellation on error
|
||||
|
||||
asyncio.create_task(_background_layer_timelapse())
|
||||
spawn_background_task(_background_layer_timelapse(), name="background-layer-timelapse")
|
||||
|
||||
log_timing("All background tasks scheduled")
|
||||
|
||||
@@ -4621,7 +4627,10 @@ async def on_print_complete(printer_id: int, data: dict):
|
||||
# Schedule timelapse scan as background task with retries
|
||||
# The printer needs time to encode the video after print completion
|
||||
baseline = _timelapse_baselines.pop(printer_id, None)
|
||||
asyncio.create_task(_scan_for_timelapse_with_retries(archive_id, baseline))
|
||||
spawn_background_task(
|
||||
_scan_for_timelapse_with_retries(archive_id, baseline),
|
||||
name=f"scan-timelapse-{archive_id}",
|
||||
)
|
||||
log_timing("Timelapse scan scheduled")
|
||||
|
||||
logger.info("[CALLBACK] on_print_complete finished for printer %s, archive %s", printer_id, archive_id)
|
||||
@@ -5442,7 +5451,7 @@ async def lifespan(app: FastAPI):
|
||||
logging.warning("Failed to auto-connect to Spoolman: %s", e)
|
||||
|
||||
# Start the print scheduler
|
||||
asyncio.create_task(print_scheduler.run())
|
||||
spawn_background_task(print_scheduler.run(), name="print-scheduler")
|
||||
|
||||
# Start background dispatch worker for send/start operations
|
||||
await background_dispatch.start()
|
||||
|
||||
@@ -13,6 +13,7 @@ from sqlalchemy import and_, or_, select, text
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from backend.app.core.config import settings
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.models.archive import PrintArchive
|
||||
from backend.app.models.filament import Filament
|
||||
from backend.app.models.printer import Printer
|
||||
@@ -1511,7 +1512,7 @@ class ArchiveService:
|
||||
|
||||
# For non-MP4 videos (e.g. AVI from P1S), kick off background conversion
|
||||
if not filename.lower().endswith(".mp4"):
|
||||
asyncio.create_task(
|
||||
spawn_background_task(
|
||||
_convert_timelapse_to_mp4(archive_id, timelapse_file),
|
||||
name=f"timelapse-convert-{archive_id}",
|
||||
)
|
||||
|
||||
@@ -20,6 +20,7 @@ from sqlalchemy import select
|
||||
|
||||
from backend.app.core.config import settings
|
||||
from backend.app.core.database import async_session
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.core.websocket import ws_manager
|
||||
from backend.app.models.library import LibraryFile
|
||||
from backend.app.models.printer import Printer
|
||||
@@ -629,7 +630,10 @@ class BackgroundDispatchService:
|
||||
progress_state["last_emit"] = now
|
||||
progress_state["last_bytes"] = uploaded
|
||||
loop.call_soon_threadsafe(
|
||||
lambda u=uploaded, t=total: asyncio.create_task(self._set_active_upload_progress(job, u, t))
|
||||
lambda u=uploaded, t=total: spawn_background_task(
|
||||
self._set_active_upload_progress(job, u, t),
|
||||
name=f"upload-progress-{job.id}",
|
||||
)
|
||||
)
|
||||
|
||||
if ftp_retry_enabled:
|
||||
@@ -828,7 +832,10 @@ class BackgroundDispatchService:
|
||||
progress_state["last_emit"] = now
|
||||
progress_state["last_bytes"] = uploaded
|
||||
loop.call_soon_threadsafe(
|
||||
lambda u=uploaded, t=total: asyncio.create_task(self._set_active_upload_progress(job, u, t))
|
||||
lambda u=uploaded, t=total: spawn_background_task(
|
||||
self._set_active_upload_progress(job, u, t),
|
||||
name=f"upload-progress-{job.id}",
|
||||
)
|
||||
)
|
||||
|
||||
if ftp_retry_enabled:
|
||||
|
||||
@@ -16,6 +16,7 @@ from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from backend.app.core.compat import StrEnum
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.core.websocket import ws_manager
|
||||
from backend.app.models.printer import Printer
|
||||
from backend.app.services.bambu_ftp import (
|
||||
@@ -258,14 +259,15 @@ class FirmwareUpdateService:
|
||||
await self._broadcast_progress(printer_id, state)
|
||||
|
||||
# Run the upload in background
|
||||
asyncio.create_task(
|
||||
spawn_background_task(
|
||||
self._do_upload(
|
||||
printer_id=printer_id,
|
||||
ip_address=printer.ip_address,
|
||||
access_code=printer.access_code,
|
||||
model=model,
|
||||
target_version=target_version,
|
||||
)
|
||||
),
|
||||
name=f"firmware-upload-{printer_id}",
|
||||
)
|
||||
|
||||
return True
|
||||
|
||||
@@ -13,6 +13,7 @@ from sqlalchemy.orm import selectinload
|
||||
|
||||
from backend.app.core.config import settings
|
||||
from backend.app.core.database import async_session, run_with_retry
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.models.archive import PrintArchive
|
||||
from backend.app.models.library import LibraryFile
|
||||
from backend.app.models.print_queue import PrintQueueItem
|
||||
@@ -2196,14 +2197,15 @@ class PrintScheduler:
|
||||
# that would otherwise cause the item to re-dispatch as a reprint
|
||||
# of the just-finished job (#1078).
|
||||
if pre_state:
|
||||
asyncio.create_task(
|
||||
spawn_background_task(
|
||||
self._watchdog_print_start(
|
||||
item.id,
|
||||
item.printer_id,
|
||||
pre_state,
|
||||
pre_subtask_id,
|
||||
pre_gcode_file,
|
||||
)
|
||||
),
|
||||
name=f"watchdog-print-start-{item.id}",
|
||||
)
|
||||
|
||||
# Get estimated time for notification
|
||||
|
||||
@@ -8,6 +8,7 @@ from typing import TYPE_CHECKING
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from backend.app.core.tasks import spawn_background_task
|
||||
from backend.app.services.homeassistant import homeassistant_service
|
||||
from backend.app.services.printer_manager import printer_manager
|
||||
from backend.app.services.rest_smart_plug import rest_smart_plug_service
|
||||
@@ -340,7 +341,7 @@ class SmartPlugManager:
|
||||
logger.info("Scheduling turn-off for plug '%s' in %s seconds", plug.name, delay_seconds)
|
||||
|
||||
# Mark as pending in database (survives restarts)
|
||||
asyncio.create_task(self._mark_auto_off_pending(plug.id, True))
|
||||
spawn_background_task(self._mark_auto_off_pending(plug.id, True), name=f"plug-auto-off-pending-{plug.id}")
|
||||
|
||||
task = asyncio.create_task(
|
||||
self._delayed_off(
|
||||
@@ -419,7 +420,7 @@ class SmartPlugManager:
|
||||
logger.info("Scheduling temperature-based turn-off for plug '%s' (threshold: %s°C)", plug.name, temp_threshold)
|
||||
|
||||
# Mark as pending in database (survives restarts)
|
||||
asyncio.create_task(self._mark_auto_off_pending(plug.id, True))
|
||||
spawn_background_task(self._mark_auto_off_pending(plug.id, True), name=f"plug-auto-off-pending-{plug.id}")
|
||||
|
||||
task = asyncio.create_task(
|
||||
self._temp_based_off(
|
||||
@@ -579,7 +580,7 @@ class SmartPlugManager:
|
||||
self._pending_off[plug_id].cancel()
|
||||
del self._pending_off[plug_id]
|
||||
# Clear pending state in database
|
||||
asyncio.create_task(self._mark_auto_off_pending(plug_id, False))
|
||||
spawn_background_task(self._mark_auto_off_pending(plug_id, False), name=f"plug-auto-off-pending-{plug_id}")
|
||||
|
||||
def cancel_all_pending(self):
|
||||
"""Cancel all pending turn-off tasks."""
|
||||
|
||||
@@ -0,0 +1,104 @@
|
||||
"""Unit tests for the spawn_background_task helper (#1648 follow-up).
|
||||
|
||||
asyncio holds only weak references to tasks, so a fire-and-forget
|
||||
create_task whose return value is discarded can be GC'd mid-flight and
|
||||
log ``Task was destroyed but it is pending!`` with no traceback.
|
||||
``spawn_background_task`` is the central helper that fixes this: it
|
||||
stores a strong reference until completion, surfaces uncaught exceptions
|
||||
through the logger, and auto-removes finished tasks.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
|
||||
import pytest
|
||||
|
||||
from backend.app.core.tasks import active_task_count, spawn_background_task
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_holds_strong_ref_until_completion():
|
||||
"""Discarding the returned task must not let asyncio reap it mid-flight.
|
||||
Pre-fix, ``asyncio.create_task(coro)`` with no caller-side reference
|
||||
would let GC swallow short tasks before they finished."""
|
||||
finished = asyncio.Event()
|
||||
|
||||
async def work() -> None:
|
||||
await asyncio.sleep(0)
|
||||
finished.set()
|
||||
|
||||
# Note: NOT storing the returned task -- this is exactly the
|
||||
# pattern the helper exists to support.
|
||||
spawn_background_task(work())
|
||||
await asyncio.wait_for(finished.wait(), timeout=1.0)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_removes_from_strong_ref_set_after_completion():
|
||||
"""The strong-ref set must shrink as tasks complete; otherwise a
|
||||
long-running process accumulates one entry per spawned task and the
|
||||
helper itself becomes a leak."""
|
||||
before = active_task_count()
|
||||
|
||||
async def work() -> None:
|
||||
await asyncio.sleep(0)
|
||||
|
||||
spawn_background_task(work())
|
||||
spawn_background_task(work())
|
||||
# Yield enough times for both tasks + their done-callbacks to run.
|
||||
for _ in range(5):
|
||||
await asyncio.sleep(0)
|
||||
assert active_task_count() == before
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_uncaught_exception_logged_as_warning(caplog):
|
||||
"""A fire-and-forget task that raises must surface the exception via
|
||||
the logger with the traceback attached -- otherwise the error vanishes
|
||||
silently and only an opaque ``Task was destroyed`` notice reaches the
|
||||
support bundle."""
|
||||
|
||||
async def boom() -> None:
|
||||
raise RuntimeError("synthetic failure for test")
|
||||
|
||||
with caplog.at_level(logging.WARNING, logger="backend.app.core.tasks"):
|
||||
spawn_background_task(boom(), name="boom-task")
|
||||
for _ in range(5):
|
||||
await asyncio.sleep(0)
|
||||
|
||||
# One WARNING with the task name and the exception info.
|
||||
boom_records = [r for r in caplog.records if "boom-task" in r.message]
|
||||
assert len(boom_records) == 1
|
||||
assert boom_records[0].levelno == logging.WARNING
|
||||
assert boom_records[0].exc_info is not None
|
||||
assert isinstance(boom_records[0].exc_info[1], RuntimeError)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_cancelled_task_does_not_log_exception(caplog):
|
||||
"""Explicit cancellation isn't an error -- a service shutting down
|
||||
its background loops should not be reported as 'uncaught exception'."""
|
||||
|
||||
async def long_running() -> None:
|
||||
await asyncio.sleep(10.0)
|
||||
|
||||
with caplog.at_level(logging.WARNING, logger="backend.app.core.tasks"):
|
||||
task = spawn_background_task(long_running(), name="cancel-me")
|
||||
await asyncio.sleep(0) # Let it start.
|
||||
task.cancel()
|
||||
try:
|
||||
await task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
assert not any("cancel-me" in r.message for r in caplog.records)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_task_name_propagates():
|
||||
"""Named tasks make the leak source visible in tracebacks and the
|
||||
done-callback log line. Pin that ``name`` reaches the underlying
|
||||
Task so support bundles surface the spawn site."""
|
||||
task = spawn_background_task(asyncio.sleep(0), name="named-spawn-test")
|
||||
assert task.get_name() == "named-spawn-test"
|
||||
await task
|
||||
Reference in New Issue
Block a user