CI / lint (push) Successful in 3s
CI / extension-version (push) Successful in 2s
CI / frontend-build (push) Successful in 19s
CI / backend-lint-and-test (push) Successful in 31s
CI / integration (push) Successful in 2m6s
Build images / sign-extension (push) Successful in 3s
Build images / build-agent (push) Successful in 6s
Build images / build-web (push) Successful in 1m52s
Build images / smoke-web (push) Successful in 56s
Build images / promote (push) Skipped
Found by the all-role smoke on its very first execution (run 7319), which is
the whole argument for having added it one commit ago.
Celery's default node name is `celery@<hostname>`. In the single-container
layout all four lanes share one hostname, so all four registered as the SAME
node. Celery says so itself:
DuplicateNodenameWarning: Received multiple replies from node name:
celery@72adc5b706a7
`inspect` collapses four replies into one dict key and the last one wins, so
three lanes read as absent — and WHICH three varies between calls:
lanes not answering: maintenance_long, ml, worker
lanes not answering: maintenance_long, scheduler, worker
Fatal twice over:
* The composite healthcheck can never pass. In Swarm that is a container
that never goes healthy — restart loop, then an automatic rollback of a
deploy whose image was fine.
* `pool_grow`/`pool_shrink` take a `destination` of node names. The UI dial
and the autoscaler would have resized whichever lane happened to answer
rather than the one asked for — silently, and differently each time.
Every celery role now starts with `-n "${CELERY_NODENAME:-celery}@%h"`, and
the generated supervisord config sets that per lane. The lanes become
worker@<cid>, scheduler@<cid>, maintenance_long@<cid>, ml@<cid> — distinct,
so inspect keeps four entries and `destination` addresses what it names.
`inspect_lanes_sync` maps hostname to lane by QUEUES, so nothing there
changes; it just stops having three of its four entries overwritten.
Unset, it falls back to `celery` — exactly celery's own default — so every
service in the multi-service stack is byte-identical to before, including the
`celery@$HOSTNAME` healthcheck in docker-compose.yml and in the operator's
Swarm stack.
The test asserts DISTINCTNESS across the whole lane table rather than a fixed
string per lane. The property that broke is that no two collide, and stating
it that way keeps holding when a lane is added.
This is the bug I said a live deploy was needed to find, found in CI instead
for the price of one `docker run` — and it would have met the operator as a
rollback loop on their first consolidated deploy.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR
186 lines
7.8 KiB
Python
186 lines
7.8 KiB
Python
"""Emit a supervisord config for the single-container layout.
|
|
|
|
Milestone 422 step 5. Writes to stdout; `entrypoint.sh all` redirects it to a
|
|
file and execs supervisord against it.
|
|
|
|
## Why this is generated and not a checked-in .conf
|
|
|
|
A static config would spell out each lane's `-Q` list, and that would be a
|
|
FIFTH hand-kept copy of the queue names — after `celery_app.task_routes`, and
|
|
the three collapsed in steps 1, 2 and 4 (`service_roster.ROLE_NAMES`,
|
|
`system_activity._QUEUE_NAMES`, and the Activity filter). Every one of those
|
|
had already drifted by the time it was found.
|
|
|
|
Generating from `worker_lanes.LANES` makes a stronger guarantee than "they
|
|
match today": the processes this container runs and the lanes the application
|
|
believes in are the same list, so a lane added to `LANES` gets a process
|
|
without anyone remembering to add one, and a queue can never end up with no
|
|
consumer because a config file was missed.
|
|
|
|
## Why supervisord
|
|
|
|
It is one pip dependency on an image that is already Python, and it does the
|
|
four things this needs without being clever: restart a program that exits,
|
|
give each one its OWN stop timeout, signal the process GROUP rather than the
|
|
leader, and put every program's output on one stdout.
|
|
|
|
The process-group part is not a detail. Celery's prefork pool forks children,
|
|
and a TERM delivered only to the parent leaves them running — which is how a
|
|
"graceful" shutdown turns into orphaned workers holding tasks. `stopasgroup`
|
|
and `killasgroup` are both set for every program.
|
|
|
|
s6-overlay is the other standard answer and would work; it needs a build-time
|
|
download and a second mental model, and its advantage (correct PID-1 signal
|
|
and zombie handling) is available here from `init: true` in compose, which
|
|
puts tini in front of supervisord. Neither choice reaches the application —
|
|
nothing in FC talks to the supervisor — so this is reversible without touching
|
|
a line of product code.
|
|
|
|
## Every lane, including ml
|
|
|
|
Step 6 merged the images, so this one carries torch and the ML requirements
|
|
and the `ml` lane gets a program like any other. It starts at one slot with
|
|
its consumers CANCELLED — `enabled=false` in the seeded settings — so it
|
|
holds a process and no model. That matters: `add_consumer` needs a running
|
|
worker to reach, and without one the UI switch would have nothing to switch.
|
|
|
|
Nothing is downloaded by starting it. The model fetch is enqueued when the
|
|
lane is enabled, which is what lets rule 164 permit a runtime fetch at all.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import shlex
|
|
import sys
|
|
|
|
from ..services.worker_lanes import LANES, Lane
|
|
|
|
# One number for the whole container, and it must cover the SLOWEST lane —
|
|
# docker gives the container a single stop timeout, where compose today gives
|
|
# each service its own (90/60/180/120s). `maintenance_long` is the 180s one:
|
|
# DB backups, library audits and translation backfill. Anything less turns a
|
|
# routine restart into a SIGKILL mid-backup.
|
|
#
|
|
# Per-program values below are the old per-service ones, preserved: supervisord
|
|
# waits `stopwaitsecs` for each, and they stop in parallel, so the container's
|
|
# own timeout needs to cover the max rather than the sum.
|
|
STOP_WAIT_SECONDS: dict[str, int] = {
|
|
"worker": 90,
|
|
"scheduler": 60,
|
|
"maintenance_long": 180,
|
|
"ml": 120,
|
|
}
|
|
DEFAULT_STOP_WAIT = 60
|
|
|
|
def _program(lane: Lane, *, slots: int) -> str:
|
|
"""One [program:x] block.
|
|
|
|
`stdout_logfile=/dev/fd/1` with maxbytes 0 puts the lane's output straight
|
|
on the container's stdout unbuffered, so `docker logs` shows every lane
|
|
interleaved rather than supervisord swallowing them into rotated files.
|
|
|
|
The output is prefixed through `sed` so a line can be attributed to a lane
|
|
— four celery workers and hypercorn on one stream are otherwise
|
|
indistinguishable. The shell that the pipe requires is exactly why
|
|
`stopasgroup` matters: the signal has to reach the celery process, not the
|
|
`sh` holding the pipeline.
|
|
"""
|
|
inner = f"./entrypoint.sh {lane.entrypoint_role}"
|
|
prefixed = f"{inner} 2>&1 | sed -u 's/^/[{lane.name}] /'"
|
|
stop_wait = STOP_WAIT_SECONDS.get(lane.name, DEFAULT_STOP_WAIT)
|
|
return "\n".join([
|
|
f"[program:{lane.name}]",
|
|
f"command=sh -c {shlex.quote(prefixed)}",
|
|
# QUOTED, and that is load-bearing. supervisord parses `environment`
|
|
# as a COMMA-separated KEY=VALUE list, so an unquoted queue list reads
|
|
# as CELERY_QUEUES=default followed by three malformed entries — and
|
|
# the lane would consume only its first queue. Silent: the worker
|
|
# starts, reports healthy, and simply never picks up `import`.
|
|
f'environment=CELERY_QUEUES="{",".join(lane.queues)}",'
|
|
f"CELERY_CONCURRENCY={slots},"
|
|
# A UNIQUE celery node name per lane, and the reason is not cosmetic.
|
|
# These processes share one hostname, so celery's default
|
|
# `celery@<hostname>` made all four the SAME node: inspect collapsed
|
|
# their replies, three lanes read as absent, and which three varied
|
|
# per call (run 7319). The healthcheck could never pass, and
|
|
# pool_grow's `destination` would have addressed an arbitrary lane.
|
|
f"CELERY_NODENAME={lane.name}",
|
|
"autostart=true",
|
|
"autorestart=true",
|
|
# A lane that dies instantly and repeatedly is a broken image, not a
|
|
# transient fault. Backing off stops it burning a core in a restart
|
|
# loop while still recovering from a one-off crash.
|
|
"startretries=3",
|
|
"startsecs=5",
|
|
f"stopwaitsecs={stop_wait}",
|
|
"stopasgroup=true",
|
|
"killasgroup=true",
|
|
"stdout_logfile=/dev/fd/1",
|
|
"stdout_logfile_maxbytes=0",
|
|
"redirect_stderr=true",
|
|
"",
|
|
])
|
|
|
|
|
|
def _web_program() -> str:
|
|
"""hypercorn. Started FIRST (priority) because its role runs
|
|
`alembic upgrade head`, and a worker that boots against an un-migrated
|
|
schema fails in a way that looks like application breakage."""
|
|
prefixed = "./entrypoint.sh web 2>&1 | sed -u 's/^/[web] /'"
|
|
return "\n".join([
|
|
"[program:web]",
|
|
f"command=sh -c {shlex.quote(prefixed)}",
|
|
"priority=1",
|
|
"autostart=true",
|
|
"autorestart=true",
|
|
"startretries=3",
|
|
"startsecs=5",
|
|
# Short: HTTP requests and the occasional file download. Matches the
|
|
# 30s the operator's production stack gives the web service.
|
|
"stopwaitsecs=30",
|
|
"stopasgroup=true",
|
|
"killasgroup=true",
|
|
"stdout_logfile=/dev/fd/1",
|
|
"stdout_logfile_maxbytes=0",
|
|
"redirect_stderr=true",
|
|
"",
|
|
])
|
|
|
|
|
|
def render() -> str:
|
|
parts = [
|
|
"\n".join([
|
|
"[supervisord]",
|
|
# PID 1 in the container, so it must not daemonise.
|
|
"nodaemon=true",
|
|
# supervisord's OWN log. /dev/fd/1 keeps it on the container's
|
|
# stdout beside the programs rather than in a file nobody reads.
|
|
"logfile=/dev/fd/1",
|
|
"logfile_maxbytes=0",
|
|
"loglevel=info",
|
|
"",
|
|
]),
|
|
_web_program(),
|
|
]
|
|
# Lanes after web, in LANES order, so the log reads in a stable sequence.
|
|
for lane in LANES:
|
|
# A lane configured at zero slots still gets a PROCESS, at one slot
|
|
# with its consumers cancelled by the reconcile. Without a running
|
|
# worker there is nothing for `add_consumer` to reach, so enabling the
|
|
# lane from the UI could not work at all — the process has to exist for
|
|
# the switch to have something to switch.
|
|
parts.append(_program(lane, slots=max(1, lane.default_slots)))
|
|
return "\n".join(parts)
|
|
|
|
|
|
def main(argv: list[str] | None = None) -> int:
|
|
ap = argparse.ArgumentParser(description=__doc__)
|
|
ap.parse_args(argv)
|
|
sys.stdout.write(render())
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|