Build images / sign-extension (push) Successful in 3s
CI / lint (push) Successful in 2s
Build images / build-agent (push) Successful in 5s
CI / extension-version (push) Successful in 2s
CI / frontend-build (push) Successful in 23s
CI / backend-lint-and-test (push) Successful in 31s
Build images / build-web (push) Successful in 1m21s
Build images / smoke-web (push) Skipped
Build images / build-ml (push) Successful in 2m12s
Build images / promote (push) Skipped
CI / integration (push) Successful in 2m39s
Operator, 2026-09-22: "since the ml-worker is optional it should be shown as such in the UI and have a warning about what it does and that it pulls the models and what models and their projected size and ram requirements to run." The card previously said "a few GB, once" — a number sourced from nothing, which is exactly the hand-wave I had flagged in this step's own survey log as something that should be measured rather than asserted. ONE FACT CORRECTED WHILE WRITING THE COPY. I had named the lane "ML tagging". It downloads an EMBEDDER: google/siglip-so400m-patch14-384. WD14 tagging is the GPU agent's job — celery_app.py:5 still names both, but that has been stale since B3 (#1238), when the agent took over and this lane was left as the CPU embed fallback for stacks running no agent (see MLSettings.cpu_embed_enabled). Telling someone the lane "does tagging" would have been wrong in exactly the way this request exists to prevent. The facts are structured data on the lane, not prose in a component: ModelRequirement(repo, approx_download_bytes, approx_resident_bytes, measured). The API carries them; the card renders them. Numbers come from the system, wording from the UI. ML_BYTES_PER_SLOT IS NOW DERIVED from that requirement rather than stated separately. They have to be one number: the figure quoted to the operator before they enable the lane and the figure the cap enforces. Two copies could disagree, and the UI would promise a slot the cap then refuses. `measured=False` travels with the numbers and the card renders "about". They are estimates from the checkpoint's parameter count and dtype — ~877M params at fp32 is ~3.5GB of weights — not from a build. This decides whether someone's server survives, so it is labelled rather than rounded into something that reads like a fact. A test asserts the flag is false, to be flipped in the same commit that records a real measurement. The card now shows: an "optional" chip in the row itself (someone scanning the table should not have to enable a lane to learn it was never required), and before the switch, what the lane does, that you only need it if you are NOT running the GPU agent, the repo id, the download size, the per-slot RAM, and why the ceiling is what it is — including saying plainly when a box has too little memory to run it at all. Keyed on the lane's own `optional` flag, not on the name 'ml', so a second optional lane gets the same treatment without anyone remembering to add it. A test asserts no REQUIRED lane declares a model: if one ever needs a download, it stops being required. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR
499 lines
20 KiB
Python
499 lines
20 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)
|
|
|
|
return {
|
|
"name": lane.name,
|
|
"slots": row.slots,
|
|
"slots_cap": row.slots_cap,
|
|
"ceiling": ceiling,
|
|
"enabled": row.enabled,
|
|
"applied": applied,
|
|
"apply_error": error,
|
|
}
|
|
|
|
|
|
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}
|