Files
FabledCurator/backend/app/services/worker_control.py
T
bvandeusenandClaude Opus 5 ffcd13096a
CI / lint (push) Successful in 2s
CI / extension-version (push) Successful in 2s
Build images / sign-extension (push) Successful in 3s
Build images / build-agent (push) Successful in 6s
CI / frontend-build (push) Successful in 20s
extension / lint (push) Successful in 23s
CI / backend-lint-and-test (push) Failing after 31s
CI / integration (push) Successful in 2m16s
Build images / build-ml (push) Successful in 3m8s
Build images / build-web (push) Successful in 3m16s
Build images / smoke-web (push) Skipped
Build images / promote (push) Skipped
feat: one image for every lane, with the model fetch gated on enabling (4296)
Milestone 422 step 6. Dockerfile.ml is gone; the main image carries torch,
torchvision, transformers, onnxruntime and opencv, and serves every lane.

WHY IT HAD TO MERGE: step 5 runs every lane in one process tree, so a second
image would mean the `ml` lane could never be enabled from the UI — there
would be no worker in that container to enable. The switch needs something to
switch.

THE MODEL NO LONGER DOWNLOADS AT BOOT. `entrypoint.sh`'s ml-worker role ran
download_models before celery started, so every boot of that role reached
HuggingFace for ~3.5GB — a startup dependency on a third party for a feature
the operator may never use. Rule 164 permits a runtime fetch only for
something "optional and clearly off", so the fetch is now a TASK, enqueued
the moment the lane is ENABLED.

Being a task is what makes it visible: it gets a TaskRun row, so the download
shows in Activity with a duration and a status, and a failure is something an
operator can see and retry rather than a container that quietly never became
useful. Idempotent, so re-enabling a provisioned lane costs one no-op.

Enqueued only when the lane actually came ON (`enabled is True`, not the
resolved value) so re-saving slots does not re-fetch, and only when the
consumer change landed — a task queued onto a queue nothing consumes would
sit pending with no explanation.

`fabledcurator-ml` KEEPS PUBLISHING, from the merged Dockerfile. The
operator's Swarm stack references that name and lives outside this repo;
dropping it would not break their deploy, it would freeze it silently at the
last publish — the exact failure class this milestone keeps finding. Retiring
the NAME is its own task, gated on that stack moving. Same two-phase shape
#406 used for pixiv.

THREE LIVE BREAKAGES from deleting the file, found by grepping for it rather
than assuming the build was the only consumer:

- `docker-compose.override.yml` built the ml service from it (contributor
  path would have failed at `docker compose build`).
- `tests/test_artifact_paths.py` pins the ml path set.
- `scripts/artifacts.sh` ML_PATHS named it. A path set naming a deleted file
  silently stops contributing to the derived revision — which the reuse check
  and the version string both read. That is #3202's recorded shape.

The `--with-ml` flag is gone from the generator and the healthcheck rather
than left defaulting to true. One image carries every lane now, so a flag
that can only be passed one way is a branch pretending to be a choice.

The advisory shipped in ecbd325 is what makes this honest to an adopter: the
lane says it is optional, names the model, and gives its download and
per-slot RAM before the switch is thrown.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR
2026-09-22 08:51:44 -04:00

538 lines
22 KiB
Python

"""Read and change a lane's live pool, over the broker.
Milestone 422 step 2. The half of the milestone that does something.
## No docker socket is involved, and that is the point
Milestone 365 put "acting on the state" out of scope because restarting a
dead worker needs a docker socket the web container deliberately does not
have. That is true of RESTARTING a container. It is not true of changing how
much work a RUNNING worker does: celery's remote control sends a message over
the broker and the worker resizes its own pool. Same Redis the app already
uses, no new privilege, no new surface.
pool_grow / pool_shrink how many slots a lane runs
add_consumer / cancel_consumer whether it consumes its queues at all
The operator ruled the socket out independently (2026-09-22: *"this feature
is a very invasive idea in my mind and I'd like to avoid it"*), and nothing
here raises the question.
## The setting is PER PROCESS, not per lane total
`pool_grow(n, destination=[...])` adds n slots to EACH destination it names.
While the stack still runs several containers per lane — the operator's
production `worker` is `replicas: 2` — a single delta applied to a lane's
total would be wrong for every replica.
So `slots` means what `CELERY_CONCURRENCY` means: the pool size of one
process. The reconcile below drives EACH replica to that number
independently, computing its own delta from that replica's current pool, so
replicas that have drifted apart (one restarted, one was grown) converge
rather than being moved in lockstep from a shared baseline.
After step 5 there is one process per lane and the distinction disappears.
It matters now, and getting it wrong now would be invisible — the totals
would simply be double what the UI claimed.
## Why reserved() is read alongside the queue depth
Celery PREFETCHES: a worker pulls more messages than it can run and holds
them in memory. Those have already left the Redis list, so `LLEN` — which is
what `/api/system/activity/queues` reports — can read 0 while thirty tasks
are waiting inside a worker. Any judgement about backlog that uses only LLEN
under-reports, which matters for the UI and is disqualifying for step 7's
autoscaler.
"""
from __future__ import annotations
import asyncio
import logging
from dataclasses import dataclass, field
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from ..models import WorkerLane
from .worker_lanes import LANES, LANES_BY_QUEUE_KEY, Lane, derived_ceiling
log = logging.getLogger(__name__)
# celery control is a broker round trip on a request path, so it gets a
# deadline (rule 156) — the same reasoning and the same budget as
# service_roster's inspect. A broker that stopped answering must make this
# report "not present", which is true, rather than hang the page.
CONTROL_TIMEOUT_SECONDS = 2.0
@dataclass
class LaneLiveState:
"""What `celery inspect` says about one lane right now.
`present=False` is NOT "zero slots" — it is "nothing answered". A lane
whose worker is restarting, or whose broker is unreachable, must read as
unknown rather than as stopped: an unswept absence is not a verdict
(snippet #3969). The reconcile in step 3 skips an absent lane rather than
correcting it, which is only safe because this distinction is kept.
"""
present: bool = False
replicas: int = 0
active: int = 0
reserved: int = 0
hostnames: list[str] = field(default_factory=list)
# The queues this lane is actually consuming right now, across replicas.
# Distinct from the lane's CONFIGURED queues: `cancel_consumer` stops a
# worker consuming one without changing what it was started with, which
# is how `enabled=false` is implemented. The reconcile needs this to tell
# "already disabled" from "needs disabling" — without it, it would re-send
# add_consumer for every queue on every tick forever (lesson #4183).
consuming: set[str] = field(default_factory=set)
# Pool size PER HOSTNAME, not aggregated. The resize below computes each
# replica's own delta from its own current pool, so replicas that have
# drifted apart converge instead of being moved in lockstep from a shared
# baseline — which is what an aggregate here would silently reintroduce.
pools: dict[str, int] = field(default_factory=dict)
@property
def pool(self) -> int | None:
"""One number for the UI. `max` rather than a sum: `slots` means the
pool size of ONE process (see the module docstring), so the largest
replica is the honest answer to "what is this lane set to". None when
no replica reported — unknown, never zero."""
return max(self.pools.values()) if self.pools else None
def _lane_for_queues(queues: tuple[str, ...]) -> Lane | None:
return LANES_BY_QUEUE_KEY.get(tuple(sorted(queues)))
def inspect_lanes_sync() -> dict[str, LaneLiveState]:
"""Live state per lane name. Sync — callers wrap in asyncio.to_thread.
Never raises. Every lane is present in the result; ones nothing answered
for carry `present=False`, so a caller cannot accidentally read a missing
lane as an empty one by iterating only what came back.
"""
out = {lane.name: LaneLiveState() for lane in LANES}
try:
from ..celery_app import celery as celery_app
insp = celery_app.control.inspect(timeout=CONTROL_TIMEOUT_SECONDS)
active_queues = insp.active_queues() or {}
stats = insp.stats() or {}
active = insp.active() or {}
reserved = insp.reserved() or {}
except Exception:
log.warning("worker_control: celery inspect failed", exc_info=True)
return out
for hostname, queues in active_queues.items():
lane = _lane_for_queues(tuple(q["name"] for q in queues))
if lane is None:
# A deployment slicing CELERY_QUEUES differently. Reported by the
# roster under its raw queue list; it simply has no lane row to
# control, which is honest rather than an error.
continue
state = out[lane.name]
state.present = True
state.replicas += 1
state.hostnames.append(hostname)
state.active += len(active.get(hostname, []))
state.reserved += len(reserved.get(hostname, []))
state.consuming.update(q["name"] for q in queues)
# `pool.max-concurrency` is the number pool_grow/pool_shrink move and
# the number the UI shows. Absent on a worker whose stats did not
# answer, which leaves pool=None — unknown, not zero.
pool = (stats.get(hostname) or {}).get("pool", {}).get("max-concurrency")
if isinstance(pool, int):
state.pools[hostname] = pool
for state in out.values():
state.hostnames.sort()
return out
def set_lane_slots_sync(
lane: Lane, target: int, live: LaneLiveState | None = None,
) -> tuple[bool, str | None]:
"""Drive every replica of `lane` to `target` slots. Returns (applied, err).
Per-replica deltas rather than one shared delta: see the module docstring.
A replica already at the target is issued nothing at all, which is what
makes step 3's periodic reconcile converge instead of re-sending a grow of
zero forever (lesson #4183 — an enforcer without a reachable fixed point
re-does its own work every tick).
`applied=False` is not a failure of the SETTING. The caller has already
stored the value; this says only that the live push did not land, and the
reconcile will carry it when the lane answers again.
"""
try:
from ..celery_app import celery as celery_app
if live is None:
live = inspect_lanes_sync()[lane.name]
if not live.present:
return False, "lane is not running"
if not live.pools:
return False, "worker did not report its pool size"
control = celery_app.control
unreported = [h for h in live.hostnames if h not in live.pools]
for hostname, current in live.pools.items():
delta = target - current
if delta > 0:
control.pool_grow(delta, destination=[hostname])
elif delta < 0:
control.pool_shrink(-delta, destination=[hostname])
if unreported:
# Resized what could be resized, and said which could not. Silence
# here would leave a replica running at a size the UI claims it is
# not, with nothing anywhere recording the gap.
return False, f"no pool size reported by {', '.join(sorted(unreported))}"
return True, None
except Exception as exc: # noqa: BLE001 — reported, never raised at a caller
log.warning("worker_control: could not resize %s", lane.name, exc_info=True)
return False, str(exc)
def set_lane_enabled_sync(
lane: Lane, enabled: bool, live: LaneLiveState | None = None,
) -> tuple[bool, str | None]:
"""Start or stop `lane` consuming its queues, without killing the process.
`cancel_consumer` rather than a shutdown: a stopped consumer keeps its
worker alive and answering `inspect`, so a disabled lane stays visible and
can be turned back on. A killed worker would read as absent, which is the
same signal as a crash — and the whole point of the roster (#365) is that
those two must not look alike.
"""
try:
from ..celery_app import celery as celery_app
if live is None:
live = inspect_lanes_sync()[lane.name]
if not live.present:
return False, "lane is not running"
control = celery_app.control
for queue in lane.queues:
if enabled:
control.add_consumer(queue, destination=live.hostnames)
else:
control.cancel_consumer(queue, destination=live.hostnames)
return True, None
except Exception as exc: # noqa: BLE001
log.warning(
"worker_control: could not %s %s",
"enable" if enabled else "disable", lane.name, exc_info=True,
)
return False, str(exc)
# --- the settings half, which is async ----------------------------------------
#
# Sync celery control above, async DB below, in one module. Same split
# `service_roster` already runs (`_inspect_celery_sync` beside `touch_service`)
# — the boundary is the transport, not the concern, and "control the workers"
# is one concern.
async def _rows_by_name(session: AsyncSession) -> dict[str, WorkerLane]:
"""Every lane's row, creating any that are missing from its LANES defaults.
Self-heals rather than depending on a migration having run for a lane
added later: alembic 0103 seeded the four that existed on 2026-09-22, and
a fifth added to LANES afterwards gets its row the first time anything
asks. Without this, a new lane would read as absent and the UI would
simply not show it.
"""
rows = {
row.name: row
for row in (await session.execute(select(WorkerLane))).scalars()
}
missing = [lane for lane in LANES if lane.name not in rows]
for lane in missing:
row = WorkerLane(
name=lane.name,
slots=lane.default_slots,
slots_cap=lane.default_slots_cap,
enabled=lane.default_enabled,
)
session.add(row)
rows[lane.name] = row
if missing:
await session.commit()
return rows
async def lane_view(session: AsyncSession) -> list[dict]:
"""Every lane: what is configured, what is live, what it may grow to.
One call rather than making the UI join three sources. `pending` is the
honest backlog — Redis depth PLUS reserved — because celery prefetches and
LLEN alone reads 0 while a worker holds tasks in memory.
"""
rows = await _rows_by_name(session)
live = await asyncio.to_thread(inspect_lanes_sync)
depths = await asyncio.to_thread(_queue_depths_sync)
out = []
for lane in LANES:
row = rows[lane.name]
state = live[lane.name]
# None for a queue the broker did not answer for, which must not be
# silently summed as zero — an unknown depth is not an empty one.
known = [depths.get(q) for q in lane.queues]
depth = sum(d for d in known if d is not None) if any(
d is not None for d in known
) else None
out.append({
"name": lane.name,
"display_name": lane.display_name,
"queues": list(lane.queues),
"slots": row.slots,
"slots_cap": row.slots_cap,
"ceiling": derived_ceiling(lane),
"enabled": row.enabled,
"memory_bound": lane.memory_bound,
"optional": lane.optional,
# What enabling this lane will download, so the UI can say WHICH
# model and how big BEFORE the switch is thrown rather than after
# a multi-GB fetch has started. `measured` travels with the
# numbers: the card must not present an estimate as a fact.
"models": [
{
"repo": m.repo,
"download_bytes": m.approx_download_bytes,
"resident_bytes": m.approx_resident_bytes,
"measured": m.measured,
}
for m in lane.models
],
"live": {
"present": state.present,
"replicas": state.replicas,
"pool": state.pool,
"active": state.active,
"reserved": state.reserved,
},
"queue_depth": depth,
"pending": None if depth is None else depth + state.reserved,
})
return out
def _queue_depths_sync() -> dict[str, int | None]:
"""Redis LLEN per queue. None for one that did not answer — see lane_view.
Sync; the caller threads it. A per-queue try/except so one bad queue does
not cost the whole report, matching `api/system_activity._read_queues_sync`.
"""
import redis
from ..config import get_config
out: dict[str, int | None] = {}
try:
client = redis.Redis.from_url(get_config().celery_broker_url)
except Exception:
log.warning("worker_control: no broker for queue depths", exc_info=True)
return {q: None for lane in LANES for q in lane.queues}
for lane in LANES:
for queue in lane.queues:
try:
out[queue] = int(client.llen(queue))
except Exception: # noqa: BLE001 — a hiccup must not break the UI
out[queue] = None
return out
class LaneUpdateRefused(ValueError):
"""A requested value is outside what the lane may hold. Carries the reason
the UI shows — a greyed control with no explanation reads as a bug."""
async def set_lane(
session: AsyncSession,
lane: Lane,
*,
slots: int | None = None,
slots_cap: int | None = None,
enabled: bool | None = None,
) -> dict:
"""Store the operator's choice, then push it to the running lane.
BOTH, in one call, and the order matters. `pool_grow`/`pool_shrink` are
not durable — a restart drops every lane back to its env concurrency — so
a UI that only pushed would have its setting evaporate on the next deploy
with nothing to show for it (lesson #4202: the live change does not
survive, and nothing says so). Storing alone would be a number that
describes nothing until something restarts.
A failed PUSH is not a failed setting. The value is saved either way and
step 3's reconcile carries it when the lane answers again; the result says
`applied: false` with a reason so the UI can say "saved, not yet live"
rather than "that didn't work".
"""
rows = await _rows_by_name(session)
row = rows[lane.name]
new_cap = row.slots_cap if slots_cap is None else slots_cap
new_slots = row.slots if slots is None else slots
new_enabled = row.enabled if enabled is None else enabled
ceiling = derived_ceiling(lane)
if new_cap < 0 or new_slots < 0:
raise LaneUpdateRefused("slots and cap cannot be negative")
if new_cap > ceiling:
raise LaneUpdateRefused(
f"cap {new_cap} is above what this container can hold "
f"({ceiling} for {lane.display_name})"
)
if new_slots > new_cap:
raise LaneUpdateRefused(f"slots {new_slots} is above the cap {new_cap}")
row.slots_cap = new_cap
row.slots = new_slots
row.enabled = new_enabled
await session.commit()
applied, error = True, None
if enabled is not None:
applied, error = await asyncio.to_thread(
set_lane_enabled_sync, lane, new_enabled,
)
if applied and slots is not None:
applied, error = await asyncio.to_thread(set_lane_slots_sync, lane, new_slots)
# Enabling a lane that needs models is what triggers the fetch (milestone
# 422 step 6). Never at boot: that made every start of the ML role reach
# HuggingFace for ~3.5GB, and rule 164 permits a runtime fetch only for a
# feature that is optional and clearly OFF.
#
# Only when the lane actually came on — `enabled is True` rather than
# `new_enabled`, so re-saving slots on an already-enabled lane does not
# re-enqueue. And only when the consumer change landed: enqueueing a task
# onto a queue nothing is consuming would leave it pending with no
# explanation until the lane returns.
fetching = False
if enabled is True and lane.models and applied:
fetching = _enqueue_model_fetch()
return {
"name": lane.name,
"slots": row.slots,
"slots_cap": row.slots_cap,
"ceiling": ceiling,
"enabled": row.enabled,
"applied": applied,
"apply_error": error,
# Tells the card to say a download has started rather than leaving the
# operator to wonder why a freshly enabled lane is busy.
"fetching_models": fetching,
}
def _enqueue_model_fetch() -> bool:
"""Queue the model download. Returns whether it was accepted.
Import inside the function: `backend.app.tasks.ml` pulls in torch, and web
must not pay that import cost on a module that every settings request
touches.
Never raises. A broker that will not take the task is worth reporting, but
the SETTING has already been stored and the lane is already enabled — so
failing the whole request here would roll back nothing and tell the
operator their change did not happen when it did.
"""
try:
from ..tasks.ml import ensure_models
ensure_models.delay()
return True
except Exception: # noqa: BLE001 — reported, never raised at a caller
log.warning("worker_control: could not enqueue the model fetch", exc_info=True)
return False
def reconcile_lanes_sync(desired: dict[str, tuple[int, bool]]) -> dict:
"""Drive every RUNNING lane to its stored slots and enabled flag.
`desired` is lane name -> (slots, enabled), read from the database by the
caller. This function touches no database: the celery task that schedules
it owns the sync session, and keeping the DB out of here is what lets the
same code be called from anywhere that already knows the target.
## Why this exists at all
`pool_grow` is not durable. A worker that dies and is restarted by its
supervisor comes back at its ENV concurrency — silently below whatever the
operator set — and nothing in step 2's path would ever notice. Storing the
value made it survivable; this is what makes it actually survive.
## It must converge and then go quiet
One `inspect` for all lanes, and `set_lane_slots_sync` issues nothing at
all to a replica already at its target. So a settled system performs one
broker round trip per tick and sends no control messages — the reachable
fixed point lesson #4183 is about. An enforcer that re-sent a grow of zero
every tick would churn forever and bury a real correction in its own noise,
which is why `changed` below counts only lanes that actually moved.
## An absent lane is SKIPPED, not corrected
`present=False` means nothing answered — a worker restarting, or a broker
that is unreachable. It does NOT mean zero slots. Correcting an absence
would be drawing a conclusion from an unswept read (snippet #3969), and
here it would be worse than useless: there is nothing to send the message
to. The lane is reported as skipped and picked up on a later tick.
"""
live = inspect_lanes_sync()
changed: list[str] = []
skipped: list[str] = []
failed: dict[str, str] = {}
for lane in LANES:
target = desired.get(lane.name)
if target is None:
continue
slots, enabled = target
state = live[lane.name]
if not state.present:
skipped.append(lane.name)
continue
# Enabled first: a lane being turned on should be consuming before
# its pool is sized, so the slots it gains have work to pick up.
#
# Only when it DISAGREES. Calling this unconditionally would send
# add_consumer for every queue on every tick of a settled system —
# the exact churn lesson #4183 describes, and invisible because
# add_consumer on a queue already consumed is harmless.
consuming_all = state.consuming.issuperset(lane.queues)
if enabled != consuming_all:
ok, err = set_lane_enabled_sync(lane, enabled, live=state)
if not ok:
failed[lane.name] = err or "could not set consumers"
continue
changed.append(lane.name)
current = state.pool
if current is not None and current == slots:
continue
ok, err = set_lane_slots_sync(lane, slots, live=state)
if ok:
if lane.name not in changed:
changed.append(lane.name)
log.info(
"worker_control: %s reconciled %s -> %s slots",
lane.name, current, slots,
)
else:
failed[lane.name] = err or "could not resize"
return {"changed": changed, "skipped": skipped, "failed": failed}