Commit Graph
3 Commits
Author SHA1 Message Date
Ettore Di Giacinto 1a6952ef39 feat(distributed): serve file staging over the worker tunnel
The four nodes.<id>.files.* subjects were the last commands a
serve-backend worker took off the bus. They are now HTTP routes under
workerctl.Prefix, on the same loopback server and behind the same bearer
check as the ten lifecycle verbs, so the frontend reaches them through
the worker's tunnel.

files.listdir is the verb this matters most for. Its reply had to fit a
payload the bus would carry, which put a wide model directory close to
the limit; a response body has no such ceiling, so nothing truncates the
listing at either end. A short listing reads to the frontend as files
the worker does not have.

S3NATSFileStager becomes S3FileStager and calls ControlClient, which
means every failure now lands in the bucket phase 3 exists to keep
straight: a route this frontend could not use is unroutable and nothing
may act on it, while the worker's own answer, including "that file is
not there", is evidence a caller may act on. Each RPC's deadline is
DERIVED FROM the caller's context rather than started fresh, at every
one of the five call sites, so a caller that gave up stops the RPC too.

A worker started without an object store mounts no file verb at all and
answers 404, which is the same answer a build too old to know them
gives. The subjects and the backend worker's files.> publish grant go
with them; a backend worker now publishes nowhere but its own inbox.

Assisted-by: Claude Opus 5 [claude-code]
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
2026-09-27 03:05:12 +00:00
Ettore Di Giacinto ab79d3130b fix(worker): pin the rune cut, answer unload honestly, drop the dead publisher
Review fix round 1. Seven non-blocking findings; the blocking one is a
merge gate for Task 4 rather than anything in this diff, and the report's
concern about it is corrected: until Task 4 lands, PingNode probes two
subjects no serve-backend worker subscribes to any more, so every healthy
worker reads as absent and is marked unhealthy on the scheduling path.

The rune-boundary cut in truncate was true behaviour with nothing
holding it: a byte-wise mutation survived all 201 specs. isRuneStart is
replaced by utf8.RuneStart, the same predicate the cluster package uses
for this rule, and two specs pin it, one with a rune straddling the
bound and one with a rune ending exactly on it so the fix cannot be
"always walk back".

unloadModel answered Success:true whatever Free did. That is the worker
saying "done" about work it did not do, and the frontend's only caller
is EvictLRU, so a false yes told the scheduler VRAM had been released
and let it place the next model on a node still holding the old one. It
now reports the failure, following stopModelExact, which is the honest
pattern already in this package. Still a 200: the worker answered, only
its verdict is negative. An address with nothing loaded still answers
success, which is a true answer rather than a claim about work done.

NewDebouncedInstallProgressPublisher had no production caller after the
last commit, only its own spec. Deleted rather than wired: wiring it
would publish every event on two carriers, which is what the carrier
decision exists to avoid. Its specs now run against the sink, plus one
that pins the identity stamped on each event, since the subject used to
carry the op and node id and now nothing but the body does.

The install progress wiring was exercised by no spec, because with no
gallery nothing ever invokes the download callback. The guard moves into
startProgress, shared by install and upgrade, which also emits one
resolving event before any gallery work. That is worth having on its
own: a cold install spends minutes on a manifest and a progress stream
with nothing on it is indistinguishable from a broken one. It also makes
the wiring observable end to end, and four specs now drive the real
installBackend and upgradeBackend over HTTP with no override.

model/stop and backend/stop keep taking Background rather than the
caller's context, and the sites now say why. model/stop is the
acknowledged stop path: it reserves the process, frees it, kills it,
waits for exit and releases the port, and abandoning that because the
caller hung up would leave a process marked stopping, a port not
returned to the allocator and a row nothing reconciles. In
stopBackendExact the Free is a courtesy before a kill that happens
anyway. model/unload differs because Free IS the operation there.

A route set with no prefix or no registrar is now a startup error rather
than a silent no-op: a server that comes up healthy while every route
the caller registered answers 404 is, through a tunnel, indistinguishable
from a version skew. And the AllPaths spec no longer claims to catch a
constant that was never added to the set, which it cannot; it asserts
the whole set instead, which catches a verb dropped from it.

Assisted-by: Claude Opus 5 [claude-code]
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
2026-09-27 03:05:12 +00:00
Ettore Di Giacinto 22dfa35fc3 feat(worker): serve the control plane over the tunnel, not over NATS
Ten NATS subscriptions on the worker become ten HTTP routes under
/v1/control/, served on the loopback HTTP server the worker already runs
and reached only through the tunnel's existing `http` stream tag.

The carrier is the tag that already exists rather than a new one. A new
tag would have had to invent correlation, per-request deadlines,
unbounded payloads and a progress stream, and each of those is a place
this branch has already put a defect. It would also have added a fifth
entry to the worker's stream-refusal vocabulary, which decides what a
frontend reaps on and took eight fixes to settle. Riding `http` means a
control RPC to a worker another replica holds takes the same relay the
inference path takes, which is the path that has been measured.

The request and reply DTOs are untouched, so a body on a control route
is byte-for-byte what the corresponding subject carried. No subject was
deleted: agent workers still subscribe to nodes.<id>.backend.stop.

Install and upgrade stream. They answer application/x-ndjson: zero or
more {"progress":...} lines carrying the same event the per-op NATS
subject carried, then exactly one {"reply":...} line, always last. That
deletes the 8000-byte notification cap structurally instead of
reproducing it on a new carrier: a progress line is written into the
response the caller is already reading, so there is nothing to size and
no subscribe-before-request window. The debouncer is shared with the
NATS publisher rather than forked, so the ~4/s tick bound is one fact.

A verb's own failure is a 200 with Error set, never a 5xx. The frontend
maps a transport failure onto "no route to that worker", which nothing
may act on, and the worker's answer onto evidence a reap guard may act
on; answering 500 for a failed install would put the worker's verdict
in the bucket reserved for a broken link. Only a request that could not
be read or routed is non-2xx.

Control RPCs carry the caller's budget. r.Context() replaces four
context.Background() calls at the gallery-install sites, and the one
pre-existing fixed timeout on model.unload is now derived from the
caller's context so a shorter budget is honoured. No timeout is invented.

The inner `go func()` in the install and upgrade handlers is deleted
rather than nested: it existed because one subscription served every
install, and over HTTP each request already has its own goroutine.
Per-backend serialization stays lockBackend, which is what actually
prevented two requests racing the gallery directory.

Bounds against a boundary the worker now serves: every body is capped at
8 MiB before any decode; the 404 echoes at most 128 bytes of the request
path, cut on a rune boundary so a half rune cannot travel downstream as
a replacement character; non-POST is refused before the body is read so
a probe cannot fire a command; the streaming responses set nosniff.

The routes mount through nodes.AuthenticatedRoutes, which hands the
registrar a private mux and puts the whole prefix behind the same
constant-time bearer check as the file routes. The worker's HTTP server
now takes the supervisor as a required parameter, so there is no way to
start it without the control plane mounted.

Assisted-by: Claude Opus 5 [claude-code]
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
2026-09-27 03:05:12 +00:00