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
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
538 lines
22 KiB
Python
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}
|