feat(distributed): run distributed mode without a message broker - #11812
feat(distributed): run distributed mode without a message broker#11812localai-org-maint-bot wants to merge 121 commits into
Conversation
Starting a Postgres and a NATS container per spec cost roughly 48 minutes of startup across the 213 specs behind SetupInfra, which is why this suite was never wired into CI. Containers move to BeforeSuite and isolation comes from CREATE DATABASE, which the dbName argument already described. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
A failed CREATE DATABASE panics out of the assertion before closeDB runs, leaking a pgx pool per attempt. With --flake-attempts 5 that exhausts postgres:16-alpine's 100 connection slots, at which point the cleanup path's own Expect fails the spec and one hiccup cascades across the suite. Scope the admin handle so the panic unwinds through defer closeDB, and let cleanup use a fallible tryAdminDB that reports rather than asserts. Register DeferCleanup immediately after CREATE so a later failure cannot leave the database behind, and warn on TestInfra that the container handles are now suite-wide. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The WebSocket log handler writes its "initial" batch before it calls Subscribe, so a line appended the instant that batch arrives lands in the circular buffer with no subscriber to receive it. Three backend-logs specs append exactly there and then wait out a 5s read deadline; once a gorilla read hits its deadline the connection is unusable, so the spec cannot retry. `--focus='Worker WebSocket log streaming' --repeat=25` failed on attempt 17 with nothing else running, which is far too often to wire into CI. Add BackendLogStore.SubscriberCount, resolving a model ID by the same exact-key and replica-prefix rules Subscribe uses, and have the specs poll it until the handler has attached. Nothing in production calls it and no assertion is weakened; the handler's own snapshot/subscribe window is left as it is, being a production streaming question rather than a test one. Verified with 60 repeats of the WebSocket specs and three consecutive --randomize-all runs of the whole distributed suite, all at --flake-attempts 1: 239 of 240 specs pass in about 80 seconds. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
… works around Three corrections from review of the previous commit. The lock-order comment on SubscriberCount claimed no path takes s.mu and a buffer lock together. Subscribe does exactly that, holding s.mu.RLock across replica registrations that take buf.mu. State the rule that is actually true — s.mu precedes any buffer lock, so counting after releasing it preserves the order — and say what follows from it: the total is a sample, not a snapshot. waitForLogSubscriber read as general-purpose but unblocks on the first registered subscription. Subscribe attaches the exact-key buffer and each replica buffer one at a time, so for a replicated model the count goes positive while later replicas are still unattached and the race survives. Rename it waitForSingleLogSubscriber, document that it holds only where Subscribe resolves to one buffer, and assert on exactly 1: misuse then fails loudly on the count rather than going quietly back to being flaky. Taking the expected count as a parameter was the alternative, but that makes callers predict a store-internal number and an under-count fails the same silent way as the original bug. The snapshot-then-subscribe race had no artifact outside a report, and review found a second site carrying it. Mark both handlers identically, including the point that swapping the two calls duplicates rather than drops and so is not the fix. The race itself is left alone; this branch stays test infrastructure. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The suite has never run in CI, so 239 specs across 32 files were verified only by hand. Path-filtered to distributed code, advisory until it earns a track record, and with flake retries at 1 rather than 5 so nondeterminism surfaces instead of being retried away. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The path allowlist covered 13 of the 99 packages the suite reaches. Commit 1dc3aee touched core/config, core/services/modeladmin and core/backend and matched no entry, so it would have merged without running the very specs that cover it. Use the paths-ignore denylist tests-e2e.yml already uses. Disable the testcontainers reaper: the runner is ephemeral, so the reaper buys nothing and its unpinned image was pulled mid-suite, defeating the pre-pull. Drop continue-on-error, which no other workflow uses and which reports a failed run as green. The job is advisory by staying out of branch protection instead. Pin Go to 1.26.0 to match go.mod, and add the tmate-on-failure step. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Runs local-ai as real child processes, one per frontend replica and one per worker, against containerised infrastructure. The in-process suites cannot express frontend-replica failure: there is no process to kill and no real HTTP boundary between a worker and the frontend it registered with. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Restarting a frontend replica must not move it: workers read LOCALAI_REGISTER_TO once at boot and never re-resolve it, so a replica that returns on a fresh port is unreachable by the workers that registered with it. startFrontend now takes the port, with <= 0 meaning "allocate". Process logs are opened for append rather than truncated, so a restarted process cannot erase the log of the instance that died, which is the log a failover post-mortem needs. The post-SIGKILL wait is bounded, so one stuck child no longer becomes a suite-wide timeout that names nothing. Stop is nil-safe because Start returns a nil cluster after stopping itself. Start's doc comment no longer claims to wait for worker registration; that needs an authenticated admin session, so it now says callers must poll /api/nodes themselves. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The register handler answers 201 both for "user created, here is your session" and for "this email already exists", so the status code cannot tell a fresh registration from a repeat one. Key on the session cookie instead and fall through to login when it is absent. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
… secret
Session rows are keyed by an HMAC of the token under a secret generated
per instance into {DataPath}/.hmac_secret. The replicas shared that
secret only because they shared a working directory, and that directory
was the source tree. Give each frontend LOCALAI_DATA_PATH under its own
baseDir and pin LOCALAI_AUTH_HMAC_SECRET, so a session minted at one
replica resolves at every other one by construction.
Assisted-by: Claude Opus 5 [claude-code]
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…ness
The point of running LocalAI as real child processes is to be able to take
one away. Add KillFrontend (SIGKILL, the lost replica), StopFrontendGracefully
(SIGTERM, the rolling update), KillWorker, RestartFrontend and FrontendAlive.
RestartFrontend pins the dead replica's original port. Workers read
LOCALAI_REGISTER_TO once at boot and never re-resolve it, so a replica that
returns on a fresh port is unreachable by exactly the workers that registered
with it and the failover under test never happens.
It also wipes the replica's data directory, so the process comes back with
empty local state and has to rehydrate node, session and job state from the
shared Postgres and NATS. Reusing the directory would model a pod with a
persistent volume and hide the class of bug these tests exist to find. That
is only safe because the harness pins LOCALAI_AUTH_HMAC_SECRET; otherwise the
wipe would take {DataPath}/.hmac_secret with it and every session minted
before the restart would 401 afterwards.
FrontendAlive consults the reaper's exited channel before signal 0: a child
that has died but has not yet been waited on is a zombie, and signal 0 to a
zombie succeeds, which would report a dead replica as alive.
The new specs cover argument validation only. Killing, stopping and
restarting a live process needs a built binary plus Postgres and NATS, so
those paths stay unexecuted until the failover suites land.
Assisted-by: Claude Opus 5 [claude-code]
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…he wipe Review round 1. Comments only, plus one guard. The note on Process.alive claimed the exited check closed the zombie window. It does not. The reaper closes exited only after Cmd.Wait returns, and Wait marks the os.Process done before returning, so exited being closed implies signal 0 already errors and the branch cannot fire earlier than the one it precedes. The window between the child exiting and waitid collecting it stays open in both versions, and the only real mitigation is for callers to poll with Eventually rather than sample once. Keep the check as hygiene, say what it actually does, and say it again on the exited field, so nobody reads the old claim and drops the Eventually. Record what the cold wipe destroys. The harness sets no LOCALAI_STORAGE_URL, so the object store is a directory under DataPath, and quantization and fine-tune outputs live there too. Postgres keeps the job row; the artifact it points at does not survive the restart. A spec that asserts otherwise will fail for a storage reason wearing a failover costume. Tell callers to let a graceful stop finish before restarting: RestartFrontend terminates with SIGKILL, so pairing it straight after StopFrontendGracefully cuts the drain short and silently converts the rolling-update case into the crash case. Refuse to wipe when the cluster has no work dir. frontendDataDir is relative when baseDir is empty, so a Cluster built by some future test helper without one would have RemoveAll walking frontend-N/data inside the source tree. The guard sits before terminate, so a refusal leaves the cluster as it was. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Tasks 4 to 6 built a harness that runs local-ai as real child processes, but none of it had ever started a process: every spec so far returned inside argument validation. These two specs are the first to run it against a real binary, a real Postgres and a real NATS. Two frontends against one database both see a worker that registered through only one of them. Every failover spec assumes this, so it is asserted first. One admin session is minted at frontend 0 and reused for both replicas rather than registering per frontend. The auth routes share a five-per-minute-per-IP limiter and all e2e traffic is 127.0.0.1, so a session per frontend would exhaust the budget as soon as a spec needs a third one. Reuse is sound because sessions live in the shared Postgres and the harness pins one HMAC secret across replicas; frontend 1 answering /api/nodes with 200 on a cookie minted at frontend 0 is what proves it. The binaries are resolved before SetupInfra so a missing build skips without first provisioning a database the skip would then have to tear down. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The Cluster label partition is these two specs and nothing else, so a missing binary skipped the entire job. Ginkgo exits 0 on skips, so a build step that broke or moved its output would have left the job reporting "0 Passed | 2 Skipped" and going green without ever starting a cluster: the silent pass this suite exists to make impossible. Skipping stays the local default, which is the right courtesy for someone who has not run `make build`, but LOCALAI_E2E_REQUIRE_BINARIES turns it into a failure that names the missing path and the target that builds it. A value that is set but unparseable counts as on, since reading it as off would restore the very skip it disables. Failures also name themselves now. The roster poll kept returning a bare nil on error, so a 401 at the second replica, a decode failure and "the worker never registered" all presented identically as an empty list. It now retains the last error and the last roster and reports whichever happened, through a lazily evaluated Gomega description that costs nothing until something fails. Finally, the two-frontend spec no longer depends on the harness to mean what it says. It asserts an unauthenticated GET /api/nodes at frontend 1 is refused, which observes the admin gate instead of assuming it, and it compares the worker's registration id across the two replicas rather than its name. A future harness that registered every worker with every frontend would have kept a name-only assertion green while it quietly stopped proving anything about shared state. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The previous round made a missing binary fail instead of skip, but only when a workflow remembered to set LOCALAI_E2E_REQUIRE_BINARIES. That leaves the silent pass one forgotten line away: the Cluster label partition is two specs, Ginkgo exits 0 on skips, and a job that skips both reports "0 Passed | 2 Skipped" and goes green having never started a cluster. So the polarity is inverted. Binaries are required whenever CI is set, which GitHub Actions always does, and the flag now exists to force the requirement OFF rather than to be remembered ON. A local developer sees no change, since CI is unset in an ordinary shell and a missing binary still skips with a message naming the path and how to build it. off, no, n and disabled are honoured as off; ParseBool rejects them, and reading a word that unambiguous as its opposite would be a worse trap than the one this removes. Also correct a claim the previous commit message got wrong. Comparing the worker's registration id across the two replicas does not pin the topology: NodeRegistry.Register looks a node up by name and preserves the existing id, and both replicas read one Postgres, so registering the worker with every frontend would yield identical ids too. The assertion is still worth keeping for what it does catch, a replica answering from its own registry or database instead of the shared one, and the comment now says that and nothing more. The topology fact moves to where someone would break it: a note on LOCALAI_REGISTER_TO recording that workers register with frontend 0 only, that the cross-replica specs depend on it, and that nothing in those specs can detect a change to it. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…plicas Four scenarios with no prior equivalent: killing a replica must not disturb a worker that never depended on it, a cold-restarted replica must rehydrate the roster from shared state and keep accepting the worker's heartbeats, a dead worker must settle to offline on every replica, and two replicas registering a worker each must converge on one roster. The timings are measured, not assumed. Node liveness is heartbeat freshness, so the only eviction path is StaleNodeThreshold (60s) plus one HealthCheckInterval tick (15s), and neither is reachable from the CLI. A worker whose registrar was killed was observed going offline at 74.2s. Every window here is sized to outlast that, because an assertion that expires before the system could have reacted proves nothing. Two assertions are deliberately unlike the obvious form. Statuses are compared for equality against a probe that returns a sentinel on error, rather than asserting a name is absent from the healthy list: the list probe returns nil on any error, and "does not contain" is satisfied by nil, so a 401 at the second replica would have passed while observing nothing. And a killed worker is required to settle to exactly offline, because it first flaps to unhealthy at ~8s and back to healthy at ~14s, which any not-healthy matcher would accept. SpreadWorkerRegistrations is new, off by default, and exists so the racing spec is a race: the harness otherwise points every worker at frontend 0, which would have left that scenario asserting on two sequential writes through one process. The default is unchanged because the baseline specs depend on it. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…r windows The two specs that assert a healthy worker stays healthy were pure negatives: they say nothing happened. A cluster whose health checking had wedged, by leaking the advisory lock the monitor takes at health.go:110, would freeze the roster and satisfy both while observing a corpse. Kill the worker once the window closes and require the roster to settle it to offline, so the preceding Consistently is a statement about behaviour rather than about a stopped clock. Applied to the cold-restart spec as well as the peer-death one: a restart is exactly the event that could leave a replacement unable to check anything. Document the hazard that can make an offline assertion hang. The staleness branch skips a node already marked unhealthy (health.go:153-155), a skip meant for nodes an operator took down, which also swallows the flap: an unhealthy mark landing after the heartbeat goes stale means MarkOffline is never called and the node stays unhealthy forever. Name the file and line at the assertion, and have the failure message say so when the roster shows a node stuck there, so a timeout sends the reader to LocalAI rather than to the harness. Stop calling the two-replica registration spec a race. Start spawns workers sequentially and the registrations land about a second apart; it is a shared-roster identity test, and saying otherwise invites someone to trust it for something it does not check. WorkerRegistrar now bound-checks its index like every other index-taking method here. It answered 0 for an out-of-range worker, and 0 is a real frontend index, so the failure mode was a spec killing the wrong replica. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Add test-e2e-cluster and a second CI job that runs it. The cluster specs spawn local-ai as real child processes and kill them, so they need a built binary; keeping them in their own job means the fast in-process suite is not held behind that build. The binary is built with a stubbed core/http/react-ui/dist. A single index.html satisfies the go:embed in core/http/app.go, and this suite drives the HTTP API only, so the job skips a Node and Vite install entirely. The job runs serial and pins --flake-attempts 1. Each Ginkgo process would otherwise get its own PostgreSQL and NATS container while every spec spawns two or three children, and a retry would hide exactly the nondeterminism the suite exists to catch. Measured at 8m39s over three runs, hence a 25 minute job timeout and a 20 minute Ginkgo timeout. LOCALAI_E2E_LOG_DIR points inside the workspace so the per-process logs upload as an artifact on failure; they are the only way to read a cluster failure. LOCALAI_E2E_REQUIRE_BINARIES is set explicitly even though CI already implies it, because a skipped cluster spec is indistinguishable from a passing one and this job's whole value is that it cannot go green without starting a cluster. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Ginkgo exits 0 when a label filter matches nothing, so a refactor that
renamed or dropped Label("Cluster") would have left the job reporting
"Test Suite Passed" having started no cluster. LOCALAI_E2E_REQUIRE_BINARIES
does not cover that case: it only fires inside a spec that is already
running. Add --fail-on-empty to both distributed targets.
Drop -r from test-e2e-cluster while here. All six Cluster specs live in the
top-level package, and the cluster subpackage contributes nothing under this
filter by design, so recursing only widened the blast radius. test-e2e-
distributed keeps -r: it must reach the eight argument-validation specs in
that subpackage.
Raise the cluster job to 45 minutes, matching its sibling. The 20 minute
Ginkgo timeout bounds the suite alone; the job timeout must also cover setup,
which is the larger and more variable half here: cold-cache module download,
protoc and protogen-go, a full build of ./cmd/local-ai and a separate test
compile, realistically 8-12 minutes on a 4-vCPU runner. At 25 minutes the
runner would have hard-killed the job before Ginkgo could report which spec
hung, which is the red-with-no-evidence outcome that gets suites disabled.
Also move upload-artifact to @v7 with the rest of the repo, and note on the
react-ui stub step that it must go if a spec ever asserts on a UI asset,
since a developer box has a real dist/ and would not catch that locally.
Assisted-by: Claude Opus 5 [claude-code]
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Two Make targets, a flake-budget variable and two environment variables landed with no way to discover them. CONTRIBUTING.md now tells a contributor how to run both suites, what each costs and which variables steer the cluster one. .agents/building-and-testing.md records the decisions that are easy to undo by accident: suite-scoped containers, the shared NATS bus and what that means for a new spec, BeforeSuite over SynchronizedBeforeSuite, the label split, --fail-on-empty, the binary gate, the flake budget of 1, the coverage exclusion, and why the cluster suite's long waits must not be shortened. .agents/ci-caching.md lists tests-e2e-distributed.yml in its paths-ignore inventory; the workflow already pointed readers there, so the cross-reference was dangling. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
--flake-attempts is total attempts, not retries: ginkgo v2.29.0 sets maxAttempts = FlakeAttempts and loops attempt < maxAttempts, and the flag's usage string reads "0 - failed tests are not retried". At 1 there is no retry at all, so "retries a failing spec once" was false in CONTRIBUTING.md and implied in .agents/building-and-testing.md. Both now say each spec runs once, and cite the source so the next reader need not re-derive it. Also restores the React-UI stub rationale, which is load-bearing because a spec asserting on a UI asset passes locally against a real dist/ and is served the stub in CI; explains why 213 and ~240 differ; records that the workflow also triggers on master pushes, where paths-ignore does not apply; and completes the LOCALAI_E2E_REQUIRE_BINARIES value table, including that any unparseable value reads as ON. In .agents/ci-caching.md the stale "13 of those 20" figure now carries its qualifier inline rather than in the following sentence. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Review of the whole branch found five comments that would send a reader to the wrong place, plus three smaller inaccuracies. Nothing here changes behaviour. The KNOWN RACE note on both backend-log WebSocket handlers said the fix needs an atomic snapshot-plus-subscribe "under the store lock". It does not: BackendLogStore.mu guards only the buffers map, and AppendLine enqueues and fans out under the per-buffer buf.mu. Whoever took the store lock would ship and the race would survive, so both notes now name buf.mu and say what s.mu does and does not exclude. Two comments in the cluster harness quoted Eventually(c.FrontendAlive) .Should(BeFalse()). FrontendAlive takes an index, so Gomega rejects that with "requested 1 arguments but received 0". Both now quote the closure form the specs actually use, and say why the closure is needed. proveHealthCheckingIsAlive claimed to prove the health monitor ran for the whole preceding window. It proves the monitor was alive at the end of it, and inferring backwards needs any wedge to be sticky. In the peer-replica-death spec that inverts: health checks are single-flighted by a session-scoped pg_try_advisory_lock, the spec SIGKILLs the replica that may hold it, and until Postgres reaps the session the survivor acquires nothing and checks nothing silently. Consistently(healthy) can then pass because nothing was checking, with the positive control still succeeding once the lock frees. The doc now states what is proven, names that gap, and says the assertion is a floor rather than a proof. The Makefile still called DISTRIBUTED_TEST_FLAKES a retry count, which is what seeded that error into the two docs just corrected against it, and the workflow called the 15s window a reconcile tick when the mechanism is HealthCheckInterval in the node health monitor. Also: the cluster suite measured 509.1s / 509.8s / 512.3s, so about 8m30s and not the 8m39s/8m40s three files claimed; the dead-worker spec title implied two independent detectors when both probes read one advisory-lock-serialised verdict out of the same row; and the sanitizeDBName length assertion used <= 50, which an empty string also satisfies, where the invariant for an over-long input is exactly 50. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The closure note in cluster/failure.go quoted a Gomega error that Gomega does not emit. Describe the argument-count failure and the Eventually().WithArguments() hint instead, so nobody greps for a string that never appears. The advisory-lock note in cluster_failover_test.go called the wedge window unbounded. A SIGKILLed local child closes its socket at once, the Postgres backend reads EOF and is reaped in milliseconds, so the mechanism bounds the window tightly. Say bounded, and keep the low probability but real framing, which was right. The workflow comment attributed HealthCheckInterval to core/services/nodes/health.go. It is declared in core/config/distributed_config.go:64; health.go only carries the ticker on the unexported checkInterval. Point a debugger at the right file. Comments only, no behaviour change. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
| // invariants still apply: non-empty, at most 72 bytes, no NUL. The | ||
| // acknowledgement is deliberate rather than incidental, so a future | ||
| // tightening of the policy cannot break every failover spec at setup time. | ||
| adminPassword = "e2e-admin-password" |
| defaultAdminEmail = "admin@e2e.local" | ||
| // testHMACSecret is shared by every frontend so a session minted at one | ||
| // replica validates at all of them. See the note in startFrontend. | ||
| testHMACSecret = "e2e-cluster-hmac-secret" |
| } | ||
| name := frontendName(i) | ||
| dir := c.frontendDir(i) | ||
| if err := os.MkdirAll(filepath.Join(dir, "models"), 0o755); err != nil { |
| if err := os.MkdirAll(filepath.Join(dir, "models"), 0o755); err != nil { | ||
| return nil, fmt.Errorf("creating %s dirs: %w", name, err) | ||
| } | ||
| if err := os.MkdirAll(filepath.Join(dir, "backends"), 0o755); err != nil { |
| cmd := exec.Command(c.opts.Binary, "run", | ||
| "--address", fmt.Sprintf("127.0.0.1:%d", port), | ||
| "--models-path", filepath.Join(dir, "models"), | ||
| "--backends-path", filepath.Join(dir, "backends"), | ||
| ) |
| logPath := filepath.Join(c.opts.LogDir, name+".log") | ||
| // Append rather than truncate: a restarted process reopens the same path, and | ||
| // the log of the instance that died is the one a failover post-mortem needs. | ||
| f, err := os.OpenFile(logPath, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o644) |
| logPath := filepath.Join(c.opts.LogDir, name+".log") | ||
| // Append rather than truncate: a restarted process reopens the same path, and | ||
| // the log of the instance that died is the one a failover post-mortem needs. | ||
| f, err := os.OpenFile(logPath, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o644) |
| } | ||
|
|
||
| func copyExecutable(src, dst string) error { | ||
| data, err := os.ReadFile(src) |
| if err != nil { | ||
| return fmt.Errorf("reading %s: %w", src, err) | ||
| } | ||
| if err := os.WriteFile(dst, data, 0o755); err != nil { |
| if err != nil { | ||
| return fmt.Errorf("reading %s: %w", src, err) | ||
| } | ||
| if err := os.WriteFile(dst, data, 0o755); err != nil { |
Replicas need to find each other to relay worker traffic, and nothing in the tree recorded a replica's address. The advertised address is discovered by opening a UDP socket toward PostgreSQL and reading back the local address, which yields the interface every replica demonstrably shares without asking an operator to configure one. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…ble addresses DiscoverAdvertisedAddr promised to return an error rather than a fallback no peer can dial, but only rejected an unspecified address. With PostgreSQL on the same host or pod as a replica, which is compose, single-node and any sidecar layout, the route to it is loopback, so every replica advertised 127.0.0.1 and a peer dialling that reached itself. Loopback, link-local and zoned source addresses are now rejected with an error naming the remedy, and a port outside 1-65535 is rejected before it becomes an undialable address. Liveness was also measured on each replica's own clock: Register and Heartbeat stamped last_seen from the Go process, and Live compared those rows against the reading replica's time.Now(). Skew therefore shrank or stretched the window by writerBehind+readerAhead, evicting healthy peers or keeping dead ones. Both sides now use the database clock, which is the one clock every replica demonstrably shares. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Returns on the first direction to finish and closes both sides so the other unblocks; a sequential copy deadlocks on any protocol where the far side speaks first. EOF and use-of-closed are normal termination, not errors. The fourth spec covers a peer that stops reading mid-body, the case where a copy is parked in Write rather than in Read. The other three tear down an idle splice and pass even against a Splice that closes only one side. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
go-yamux/v5 matches none of its errors against net.ErrClosed, so the classifier reported an ordinary teardown as a failure: when the session has gone away, the FIN that Splice's own Close writes returns ErrSessionShutdown, and a stream torn down under a live copy surfaces as ErrStreamClosed or a reset. Splice owns that Close, so it owns the errors it produces; the sentinels are named here rather than injected by the caller, which would make a forgotten classifier reintroduce the same bug silently. Cover the error half of the contract, which no in-memory pipe could reach: a scripted stream now feeds Splice a genuine transport failure and each closed-stream ending in turn. Replacing the tail of Splice with "return nil" passed every previous spec. Also assert that Splice does not return until the second direction has finished, rename a spec that promised a leak check it never made, and correct two comments that claimed more than the code did. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Matching yamux errors with errors.Is was too broad. Session.close hands every live stream ErrStreamReset wrapped around whatever killed the connection, so a keepalive timeout, a broken TCP connection or a peer that simply vanished all matched, and a relayed request that died reported a clean ending. Nothing upstream would have retried or logged it. Match the plain sentinels by identity, since only identity separates a stream that was reset from the wrapped form that means the session died. Treat a StreamError as a per-stream reset, and a GoAwayError as normal only when it carries the no-error code, read off ErrRemoteGoAway because the constant is unexported. ErrSessionShutdown needs no entry of its own; it is a GoAwayError with that code. Order matters as much as the matching: session death wraps its cause, which is routinely io.EOF or a closed socket, so the mux checks run before the generic endings. Reversing them alone puts a vanished peer back to nil. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Upgrades to a WebSocket, wraps it as a yamux server session and hands it to the caller. Rejects before upgrading so an unauthenticated dial sees a 401 rather than a WebSocket error, which is what the route-coverage test asserts. The adapter keeps the reader of a partially consumed message across Read calls. yamux reads through a 4 KiB bufio.Reader, so a small-payload test cannot see a dropped message tail; the framing specs drive the adapter directly with buffers smaller than the message. An empty configured token authorizes nobody here, unlike the worker file transfer server's check: this route is registered in every deployment, so failing open would publish an unauthenticated mux. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
agent.<name>.cancel was the last family on a message bus, and the only reason an agent worker dialled one. Its subscriber is the worker running the execution, and a worker has no database, so the family could not move to the PostgreSQL fan-out carrier: a cancel published there would reach no worker while reporting that it had been sent. It is a control verb now. An agent worker mounts workerctl.PathAgentCancel on the loopback control plane behind its tunnel and applies the cancel to the same registry the executor registers a run on. The frontend issues it through nodes.AgentControlClient.CancelAgentRun. That call is a FAN-OUT and not a pick, because nothing records which worker holds a given execution: the claim row names the claiming replica, and it is deleted when the run ends. Every agent worker a live replica can reach is asked over its own tunnel, relayed by the peer mesh when a peer holds it, and each worker answers only for itself. The answers stay apart, which is why this family was held back. A cancel a worker made is nil. A cancel some worker could not be asked is ErrAgentCancelUndelivered, which is neither a refusal nor a missing run. A cancel every reachable worker declined to own is ErrAgentRunNotOnAnyWorker. A deployment with no agent worker is ErrNoAgentWorker. Neither new sentinel wraps ErrWorkerUnroutable and neither is a worker answer, so nothing is reaped, demoted or evicted because of a cancel. A worker in the ABSENT CONNECTION condition, one whose tunnel was lost inside the reconnect grace, counts as undelivered. It is not retried in the call and not queued: a retry would spend a budget the caller did not choose, and a queue would need durable state whose only consumer is a run whose control stream went with the tunnel. A worker whose departure has outlived the grace is the one routing fact a caller may act on and is excluded, or a single retired agent node would make every cancel undelivered for ever. The fan-out reads a different node set from the pick. A draining worker takes no new work but is still finishing what it holds, so it is offered the cancel; a pending one is refused by the tunnel route on every dial and is not. With that, nothing in LocalAI connects to NATS. The agent worker's dial, its credential ladder and its refresh loop are gone, and so is the frontend's cancel carrier. LOCALAI_NATS_URL is accepted and ignored everywhere, and distributed mode no longer requires it. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Task 4 gave agent workers tunnels and deliberately left the NodeType skip in HealthMonitor.tunnelDeparted, with a spec asserting that an agent node whose presence reader answers PresenceGone is NOT marked unhealthy. That spec was scaffolding. It was true while an agent worker took its jobs and its verbs over the message bus: a departure row for one said nothing about whether it could work, and an early bug in the new tunnel client could otherwise have demoted a fleet of healthy agent workers. There is no bus. An agent worker is reachable through its tunnel and through nothing else, so a departed agent tunnel means exactly what a departed backend tunnel means: no live replica holds it, the departure has outlived the reconnect grace, and that is a routing fact the scheduler and a reaper may act on. The skip would now hide the only symptom an unreachable agent worker has. This is the deliberate removal Task 4's M6 predicted, and task-4-report.md is where that mutation already stands recorded red against the spec this commit deletes. The skip existed at ONE site. router_liveness.go has none: its candidates come from queries that already filter node_type = 'backend'. The two skips in managers_distributed.go stay, because an agent worker still runs no backend processes, so it has no backend to list and no backend op to apply. Two node types can depart now, which is why the second half exists. Before this, one type could depart and every per-node cache a departure left stale was dropped from wherever its owner happened to notice, so a reader could not tell which caches a demotion invalidated by reading the demotion path. Departure gets ONE notification point. DepartureNotifier is edge triggered, because the monitor runs on a ticker and a departed node stays departed; its subscribers are NAMED, because what has to be caught is a forgotten cache and a count can say only that one of four is missing; and NewHealthMonitor takes it as a required positional argument, so a caller that does not pass one fails to compile. Four caches subscribe: prefix-cache affinity in every model, probe freshness at every address, in-flight staging operations, and the per-node breakdown of every open gallery operation. The prefix-cache one is registered only when prefix-cache routing is enabled, so --distributed-prefix-cache=false stays a true no-op. The notification carries the node's name as well as its id, because the staging tracker keys on the name and the other two key on the id, and a subscriber should not have to read the registry from inside an eviction hook. A departure notification is an act on absence, so it fires only on the routing fact. A tunnel lost inside the grace, a worker that never dialled, a presence query that failed and a stale heartbeat all announce nothing, asserted per node type. The stale-heartbeat branch is excluded on purpose: it already marks the node offline, which deletes its rows and runs the registry's replica-removed hooks, so firing there too would double-evict and make the notification mean two different things at its subscribers. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Every carrier had already moved and no process opened a bus connection, but the surface an operator reads still described a deployment with a broker in it: a compose service, a 220-line credential-generation script, two CI steps pulling a container nothing started, two flag tables offering --nats-url, an architecture diagram with a NATS box wired to the workers, a join-command generator in the Nodes page that emitted --nats-url for agent workers, and a test suite that stood a NATS server up for specs that no longer used it. That is the one way this programme could still fail invisibly. Every test passes, every binary works, and every production deployment goes on running and paying for infrastructure that carries nothing. Nothing in this repository starts a NATS server any more. The compose file is four services, the docs say to shut the broker down and what to keep, and the e2e suite runs on one PostgreSQL container. The three LOCALAI_NATS_*_TIMEOUT env vars are KEPT, and are now documented twice as being kept. They were never broker settings: each names a control-RPC budget the frontend applies to a worker, still read and still enforced. They carry the prefix only because they arrived with the bus, and renaming them would break every existing deployment for cosmetics. The agent worker's join command was the last surface still emitting the flag, two tasks after the agent worker stopped dialling. The Playwright spec that covered it asserted the opposite of what is now true, so it is inverted rather than deleted, and it reads the rendered command string rather than the component's variables: the variables are what the fix removes, so a spec reading them would have stopped compiling instead of failing, and a compile error is not evidence about what an operator is shown. nats_jwt_test.go and its helpers are deleted. They pinned a real server ENFORCING the minted permissions. The CONTENT of those allow lists is still pinned, untouched, by pkg/natsauth's own suites, including the spec that refuses to let the agent lists go empty, since an empty allow list in NATS means unrestricted. The enforcement half is retired rather than moved: enforcement is a property of a connection, and nothing opens one. The suite's own NATS container goes with them, which the brief left for the next task. Removing the pre-pull while BeforeSuite still ran the image would have defeated the step rather than cleaned it up, and this change removes the last reader of TestInfra.NC. agent_native_executor_test.go and mcp_ci_job_test.go are moved onto infra.Bus() instead of deleted: they were the last two specs building a bridge and a dispatcher on a client nobody uses, which is exactly the drift TestInfra.Bus's own comment warns about. cluster.Options.NatsURL is now fed a deliberately dead address rather than a live container's. Frontends and agent workers still receive LOCALAI_NATS_URL, because that is the coverage for the promise that an existing command line still starts; sourcing it from a running server would have let a regression that actually dialled it pass. The control in cluster_control_test.go keeps its assertion and loses its explanation, which claimed the deployment had a bus and no longer could. One latent spec race surfaced and is fixed: the background-run spec waited for a COUNT of events and then read a snapshot for the terminal status, which is the last event of a run and therefore always arrives after the count is met. Its immediate twin had already been fixed this way. Nothing in production changed. pkg/natsauth keeps its files. It is reachable from production only through the natsauth.Config parameter thread, and that thread is the next task's. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Distributed mode has not dialled a message broker since the control plane moved onto the workers' own outward tunnels and every fan-out family moved onto PostgreSQL LISTEN/NOTIFY. What was left was the dependency itself, and the code that existed only to feed it. Dropped from go.mod: nats-io/jwt/v2, nats-io/nats.go, nats-io/nkeys, nats-io/nuid and testcontainers-go/modules/nats, along with the fourteen indirect requires that only the NATS testcontainer pulled in. go.sum carries no nats line either, so the removal is not the partial kind where the require goes and the checksum stays. Deleted with them: pkg/natsauth in full, the broker client's remaining options and TLS files, the per-node JWT minting on both the register and the approve path, and the natsauth.Config parameter threaded through the node routes. The credential manager is renamed and stripped rather than deleted, because it still holds the tunnel token that every re-registration rotates. The bus flags stay accepted and ignored, and are now hidden, on every command that had them, so an existing unit file, compose file or Helm values file still starts on the day of the upgrade. What is not kept is the validation that REQUIRED one: a distributed frontend started with no bus URL is no longer fatal. The TLS paths lose type:"existingfile" deliberately, so a certificate deleted along with the broker cannot fail a startup. One operator-visible behaviour change: --nats-require-auth no longer makes an agent worker wait through admin approval. Ask for that wait with --distributed-require-auth, which already implied it. It is documented in the migration section and pinned from both sides. A deployment now needs PostgreSQL and the frontends' own HTTP listener, and nothing else. coverage-baseline.txt moves from 54.2 to 62.0. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
… workers Tasks 1 to 17 are proven by unit and integration specs and by two e2e passes taken mid-flight. This is the pass that boots the real binaries with every carrier in place and none of the old one, and it does so on the topology the feature was built for rather than on the one-worker shape the rest of the cluster suite uses. Two frontends and two workers is the configuration that matters. With each worker's tunnel landing on a different replica, the owner path and the relay path are live at the same instant against one roster, one scheduler and one health monitor, so a routing mistake has somewhere to show up instead of hiding. It is also the only shape in which "killing a replica re-homes only ITS worker" can be stated at all. Three scenarios, all 2x2: 1. Both workers served from both replicas. No broker as a property of the ARTIFACT (debug/buildinfo reports no github.com/nats-io module, with the module count asserted non-zero so a stripped binary cannot pass vacuously), no broker in either worker's live /proc environment, and no advertised address on either worker. One completion over the owner path and one over the relay, plus the mirror image through the other replica, plus four control-plane listings covering both paths for both workers. 2. The replica owning worker 0's tunnel is killed with that tunnel blocked. Leg 1 asserts nothing and only waits for the killed instance to leave the live set, because before that it still reads as a live owner and the scenario is not yet about absence. Leg 2 then holds a window inside the reconnect grace requiring that nothing acted on the absence. Leg 3 requires the re-home and inference again. Worker 1 keeps serving throughout. 3. The suite-wide negative control. Both tunnel dials refused while registration and heartbeats flow, both workers refused at both replicas naming the routing fact and not a departure, nothing reaped and both heartbeats fresh. Then ONE tunnel is restored and exactly one worker recovers while the other stays refused. Which worker served is read back from node_models rather than assumed: the two models are pinned to one worker each through PUT /api/nodes/:id/labels and POST /api/nodes/scheduling, and every assertion requires the model to be on the expected node AND absent from the other. That the relay hit a non-owner is read from the production Owner query before the request and re-read after it. Attacks run, each alone, each reverted, each behaving as predicted: hand a worker a broker URL reddens scenario 1's environment leg; point the module check at gorm.io reddens its artifact leg; start one worker instead of two reddens all three at the topology guard; delete the relay in WorkerDialer.Dial reddens scenario 1 on exactly the request sent to the non-owner while 2 and 3 stay green; a one-nanosecond reconnect grace reddens scenario 2's leg 2 on the demotion while 1 and 3 stay green; lifting both blocks at scenario 3's differential reddens its "still unreachable" half. The brief's "restore the NatsURL validation" attack cannot be applied: DistributedConfig has no such field left to validate. Label-orphan arithmetic, counting non-skipped It nodes from --dry-run: all 256, dist 231, cluster 24, vllm 1, and 231 + 24 + 1 = 256, so no spec is orphaned by the label filters. Three test-e2e-cluster runs: 897.0s, 897.8s and 906.6s of Ginkgo time, 24 specs, 15 minutes wall. The only failure across the three was a pre-existing spec dying at cluster.Start with frontend-1 exiting status 2, which passed in the other two and is reported as a port-allocation flake rather than a regression. test-e2e-distributed is 223 plus 8 specs in 131.7s. The budget comment and .agents/building-and-testing.md move from 21 specs at 800 to 830 seconds to 24 specs at 897 to 907. The harness gains ProcessEnviron, which reads /proc for any of the three process families; WorkerEnviron and FrontendEnviron become wrappers rather than being deleted, so the specs that call them are not re-aimed for a rename. Signed-off-by: Ettore Di Giacinto <mudler@localai.io> Assisted-by: Claude Opus 5 [claude-code]
| // in the roster. | ||
| func (c *Cluster) startAgentWorker(i int) (*Process, error) { | ||
| name := agentWorkerName(i) | ||
| cmd := exec.Command(c.opts.Binary, "agent-worker") |
| // Random rather than round robin. A per-replica counter is per-replica | ||
| // state that says nothing about load, and with several replicas the | ||
| // counters agree on nothing anyway. | ||
| picked := candidates[rand.IntN(len(candidates))] |
Removing the broker left one thing carrying every broadcast family in the product: PostgreSQL LISTEN/NOTIFY, in core/services/pgbus. It is covered thoroughly in process by test-e2e-distributed, and it was covered nowhere at all by real binaries: grepping the six Cluster spec files for pgbus, bus_messages, LISTEN and NOTIFY returned zero hits. Registration, model staging over the tunnel and inference through both the owner and the relay paths were already proven by real processes; the carrier that now carries everything else was not, so a deployment whose replicas each published to themselves and heard nobody would have left every suite green. Two specs, both on two frontends and no workers against one PostgreSQL, publishing at frontend 0 and reading at frontend 1. 1. A gallery operation admitted at one replica, read out of the other, with the queued state observed before the terminal one. 2. A broadcast of about 9.3 kilobytes, which PostgreSQL refuses as a notification payload, making the round trip byte for byte through the bus_messages spill table. The family is a gallery operation for one property nothing else on this carrier has: the answer a peer gives is held in memory ALONE. GET /models/jobs/<id> reads galleryop's statuses map, which on a peer is filled by the gallery.*.progress subscriber and by nothing else, because the only other filler, Hydrate, runs once at startup and every operation here is created long afterwards. Every other family has a durable table behind it that a peer would converge through anyway, and a spec on one of those cannot separate "the broadcast arrived" from "the row was read". That is then made checkable rather than argued. The gallery_operations row is written when the gallery worker DEQUEUES an operation, so an operation still waiting in the queue has NO row, and both specs assert zero rows while the peer is already answering with the operation's own bytes. Both also read the instances table and require the reading replica to be a different live instance from the publishing one, so "the other replica" cannot decay into a spelling of "this replica". Holding the queue is what cluster.Options.Galleries is for. The gallery worker runs one operation at a time on an unbuffered channel, so an install parked inside a gated index fetch parks everything behind it; without that the admission broadcast and the terminal one are separated by two database round trips and no HTTP poller could see between them. The option also turns the startup estimate warmer off, because a second fetcher filling the process-wide index cache would leave the operation never blocking and the spec passing on an ordering nothing enforced. The spill spec is written against a failure this branch has shipped three times: a size-limit spec that cannot fail. The oversized body is an ordinary element name that the real consumer decodes and surfaces, so it is not a body the decoder would have refused at any size. The size is ABSOLUTE at 9000 bytes rather than derived from the cap, and a one-byte control operation in the same run is required to leave no spill row, so moving the 8000-byte cap in either direction reddens the spec. pgbus.FitsInline, which shares its encoder and its comparison with Publish, is asked about both payloads and must answer differently. The spilled row is then decoded and its element name compared byte for byte against what frontend 1 answers. The terminal assertion in spec 1 does not re-check the element name: a terminal status does not carry one, because updateError in galleryop.Start builds a fresh OpStatus holding only the error. It asserts the two fields that status does carry, in the relation that one place writes them. Attacks run, each alone, each reverted, each behaving as predicted. Neutering the pg_notify in pgbus.Publish so every replica knows only what it did itself reddens both specs at frontend 1, which answers 500 for an operation it was never told about; the bus_messages row assertion still passes under it, which is right, since the row is written before the notification. Releasing the queue gate reddens spec 1 at the gallery_operations count, because the operation is dequeued and the row appears. Shrinking the oversized name to 100 bytes reddens spec 2 at FitsInline; inverting that guard so the run reaches the row check reddens it there instead, with no bus_messages row written, which is what makes the row a statement about size. test-e2e-cluster is 26 specs in 933.8 seconds of Ginkgo time, 15m37s wall. The two additions cost 7.0 seconds together, 5.0s and 2.0s: they start no workers, so they pay for no registration, and what they wait on is a broadcast rather than a threshold. test-e2e-distributed is unchanged at 223 plus 8 specs, 130.6 seconds. The budget comment and .agents/building-and-testing.md move from 24 specs at 897 to 907 seconds to 26 at 933.8. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
GET /api/cluster/peer authenticated with the deployment's shared registration token and took the dialling replica's id from ?id= on trust. Every worker holds that token, so anything holding it could open a peer link as any replica: relay through it to every worker tunnel that replica owns, displace a real replica's inbound link by declaring its id, and point the roughly 31 GiB per-session receive window at one replica. Validating the id against the instances table does not fix this, because the attack declares a real replica's id. So the route now checks two credentials and needs both. The shared token still says the dialler belongs to this deployment; a new per-replica credential says which replica it is. The credential follows the per-node worker credential rather than inventing a second mechanism: crypto/rand.Text, stored only as a hex SHA-256, compared in constant time, with no fallback to the shared token. It differs in the stronger direction. A worker's credential is minted by the frontend and handed over once; a replica writes its own instances row, so it mints its own secret, publishes only the hash in the same statement that publishes its address, and never sends the plaintext anywhere but the peer dial. A peer that presents no credential is refused, not waved through. An old replica and an attacker holding the shared token send the same request, so accepting the first accepts the second; there is no safe downgrade here, only a quiet one. The refusal is made loud instead, on both sides, naming the upgrade rather than the network. On the documented frontend-first order a new replica still dials an old one; an old replica cannot dial a new one, which costs relayed requests that land on a not-yet-restarted replica and surfaces as no route, never as absence. A rejected peer gets its own sentinel, ErrPeerRejected, whose unwrap chain carries ErrPeerUnreachable as well and no absence sentinel at all. Keeping the older sentinel means no existing consumer changes behaviour; the cause stays out of the chain, so absence cannot escape through it and nothing can read an authorization failure as a worker that went away. One consequence beyond the fix: a replica with no advertised address has no instances row, so it now cannot dial out either. It was already unreachable inward. The startup error and the docs say so. Registry.Register, NewMembership, NewPeerPool, PeerHandler and RegisterClusterRoutes all gained required arguments, so the identity cannot be dropped without a compile failure. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Skills and RAG collections had no cross-replica invalidation, and the two builders that would have published one were deleted earlier in this branch because nothing called them. Wiring one now would be wrong, not merely late. Both features are derived entirely from the frontend's own state directory. A skills.Service indexes <state dir>/skills, a collections backend enumerates <state dir>/collections and holds one handle per collection it found there, and no replica reads or writes another replica's copy of either. In distributed mode PostgreSQL carries a skill's NAME and description in skills_metadata, and nothing else: Get, Search, Export and the resource verbs all read local files. So a peer told to drop a cache entry would rebuild it from a directory that does not hold the change. For a postgres-engine collection it would be worse than a no-op, since re-deriving one on a replica with no local index file yields a collection that answers with an empty file list against a populated vector store. What is missing is shared storage, not a broadcast. Recorded rather than left silent: the two cache fields say why nothing invalidates them, a distributed frontend logs the limitation once at startup, and the docs name the two deployments that avoid it. The new spec pins the premise, so a change that moved either directory onto storage every replica mounts reddens and the decision gets taken again. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The backend port allocator allocated from its own bookkeeping alone. That bookkeeping records what this worker did with a port, and the collision it cannot see is with something this worker never did: the default base port is 50051, inside Linux's default ephemeral range of 32768 to 60999, so the kernel hands ports in this range to outbound connections and to anything that binds port 0. A backend handed one of those dies on bind, and the frontend sees a backend that will not start. Every candidate is now probed by binding the exact address the backend will listen on, in all four allocation branches: the key's own port, the free pool, a grown port and a stolen one. Probing the free pool matters as much as probing a grown port, because a port this worker released is exactly as available to the kernel as one it never used. A candidate that fails the probe is quarantined rather than blacklisted, since whatever holds it is usually an ephemeral connection that gives it back, and its affinity claim is dropped so an unbindable port does not stay reserved for the key that last held it. Exhaustion now says how many candidates were skipped, which is what tells an operator "something else is in my range" from "my range is too narrow". This does not remove the race and cannot: between the probe and the child's bind the kernel can still give the port away. It removes the far larger window in which the allocator hands out a port the kernel gave away minutes ago, which was the whole of the observed one-in-three harness flake. The e2e harness comment that recorded the missing check is corrected, and the docs say how to move the range out of the ephemeral one entirely. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
HealthMonitor.misses holds one consecutive-failed-probe count per (node, model, replica) and nothing ever removed an entry whose row had gone. It is the only per-node state in a frontend that grows on model churn rather than on fleet size, so a deployment that loads and unloads models for months accumulates an integer per tuple it ever probed and gives none back. There are four ways a row stops being visible to the pass, not one. A node departs and the pass skips its probes; a node goes offline or unhealthy on a stale heartbeat and the pass skips it entirely; an operator sets a node draining; or the row is removed by an unload, a scale-down or an eviction, and nothing tells this monitor. So the bound is the pass itself, and not a subscription on the departure notifier. The notifier evicts the caches a DEPARTURE invalidates and it keeps that one meaning; this reads a different fact, that there is no longer a row to count misses against, and covers all four cases with one rule. A row the pass could not probe is marked seen before the probe, so an unreachable worker still leaves its streak exactly as it was rather than having it forgiven; a pass that could not list the fleet prunes nothing, since it observed nothing. Forgetting only ever delays a reap by up to the miss threshold and can never cause one. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
A frontend keeps two caches of one http.Client per worker: the control client's, built on the first verb issued to a node, and the HTTP file stager's, built on the first file staged to it. Both are keyed by node ID and neither was ever pruned. Their own comments said so and named what a fix would need, a signal that a node has left, which did not exist when they were written and does now. The map slot is the smaller half. Each entry holds an http.Transport whose idle connections are streams on that worker's tunnel, kept until IdleConnTimeout even after the tunnel is gone, so ForgetNode closes them rather than leaving them to the collector. Both are registered on the deployment's one departure notifier, as two subscribers and not one: a node can be in either cache without being in the other, and a single hook would say only that some client was kept. ForgetNode is on the FileStager interface rather than on the one implementation that has state to drop, so a stager that grows a per-node map later cannot be added without answering the question, and so registerDepartureEvictions can take a FileStager and still fail to compile if the registration is deleted. The S3 stager's is a documented no-op that deliberately does not forward to the control client, which registers itself: forwarding would evict a cache it does not own, twice per departure, and the second drop would not appear in the subscriber names the wiring spec reads. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The claim queue's attempts counter grew without bound and nothing read it. At the default two-second poll a permanently undispatchable row cost about 43000 UPDATEs a day, and it cost more than writes: rows are claimed oldest first, so the oldest stuck row was re-claimed ahead of every newer one on every tick and held a dispatch slot while it failed. One poison row starved the queue behind it. No dead letter, and that is the decision rather than the omission. Read settleClaim: the only outcome that releases a claim is one where NOTHING was learned about the work. No agent worker was connected, the tunnel broke, a peer could not be reached, the stream was refused before the request body left this replica. Not one of those is a worker saying it ran the job and it failed, and an attempt ceiling would turn "the fleet was away long enough" into a job failure nobody reported, which is the collapse this whole design exists to prevent pointed at work instead of at nodes. The one verdict available here, that no build of any worker serves this kind, is already settled as an answer. So the retry stays unbounded and the RATE does not. Each release stamps the row with the earliest it may be claimed again, doubling from two seconds to a cap of sixty, computed in the release statement from the row's own attempts count and stamped on the DATABASE clock, because that is the clock competing replicas order the queue on. Queued work becomes claimable again within one cap of the fleet returning, and a stuck row no longer holds the head of the queue. A claim released by the reap carries no delay at all: that work was never handed to anyone, so there is nothing to back off from. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
None was introduced by this branch and all three are in test code, which is what made them survive: every suite passed on every run and only the race detector said otherwise. A known-failing -race run is worse than a noisy one, because a real race raised by production code lands in the same report and is read as one of these. galleryop: gatedModelManager guarded the recorded names and not the gate channel itself. A spec frees the parked worker by closing the gate and installing a fresh one, on the spec goroutine, while the worker goroutine reads the field to park on it. The channel is now read and replaced under the same mutex, and cleanup closes idempotently. pkg/model: two specs swapped xlog's package logger to capture output and swapped it back on cleanup. xlog.SetLogger writes an unsynchronised global, so the restore raced with the backend process watcher, which logs while a process is stopping; the captured bytes.Buffer was written by that goroutine and read by an Eventually at the same time. SetLogger is now called once for the whole test binary, from init, before a goroutine exists to race with, and a spec swaps the DESTINATION under a mutex through a routing slog.Handler. Per-spec level filtering is preserved deliberately: one of these specs asserts that a debug emission is filtered OUT and would pass vacuously against a handler that recorded everything. openai: fakeTransport appended to its event and audio logs from the response and turn coordinators' goroutines while a spec ranged over them. Both are behind a mutex and are read through snapshot accessors; the fields are renamed so a raw read from another spec file does not compile. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Probing every candidate before handing it out is right, and it made sixteen specs that were never about the kernel depend on it. They build a supervisor directly, name the ports they expect literally, and those literals sit inside Linux's default ephemeral range, so with the real probe each one asks this host whether 50051 is bindable at that instant. The first full -race run over ./core/... and ./pkg/... after the probe landed went red on four of them, and holding 50051, 50052, 50060 and 50061 from another process turns eleven red deterministically. Nothing was wrong with the allocator in either case: something else on the machine held a port, which is the situation the probe exists to survive. So the specs that assert bookkeeping now inject a probe that always says yes, and say why once. The two specs that are about the probe leave the field unset and keep asking real sockets, which is what still fails if the probe is removed. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The mutex round the realtime transport double added a snapshot accessor for each
recorded slice. Only the event one has a caller, so make lint refuses the build:
realtime_doubles_test.go:64:25: func (*fakeTransport).recordedAudio is unused (unused)
No spec has ever read the audio log, before the mutex or after it, so the
accessor is deleted rather than nolinted and the struct comment says where the
next one comes from. audioLog stays written, because a double that silently
discarded what a coordinator sent it would be a different double.
Assisted-by: Claude Opus 5 [claude-code]
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Preserve heartbeat checkpoints and backend readiness across the tunnel transport changes. Update incoming tests for the renamed worker address fields and health monitor arguments. Assisted-by: Codex:gpt-6
|
@localai-org-maint-bot fix tests |
Assisted-by: Codex:GPT-6
Apply the existing Docker skip at the container entry point so the maintenance-connection regression cannot bypass it. Assisted-by: Codex:GPT-6
The general e2e target builds a container, not the host binary that cluster specs require. Run only its own package; dedicated distributed and cluster targets retain the full subpackage coverage. Assisted-by: Codex:GPT-6
|
Pushed two test fixes on top of the latest master:
Verification: the maintenance test fails before the change and skips afterward with a Go overlay selecting the macOS path; four FakeBus specs and Still unresolved: the original Linux timeout in the Update 2026-09-08: resumed this request and pushed Reproduced the exact failure on Linux with a symlinked TMPDIR before the fix. Afterward, all 13 file-stager round-trip specs and four new containment cases pass with that TMPDIR, including rejection of sibling paths and symlinks escaping the allowed root. The focused worker file-staging tests, incoming virtual-model companion tests, The existing Linux, distributed, cluster, and VibeVoice test runs were cancelled, so full verification remains incomplete. No CI jobs were restarted or polled, and this request remains in progress. |
Merge current master and adapt incoming tests to the distributed APIs. The sound mock now implements NewClientForNode and uses WorkerLocalAddress. The upload finalization test supplies the existing direct test dialer. Assisted-by: Codex:gpt-6
Include current gallery and virtual-model path fixes before continuing the requested test repairs. Assisted-by: Codex:gpt-6
Resolve allowed directories before comparing them with resolved files. Otherwise staging rejects valid files under macOS temporary paths. Cover aliased roots, sibling paths, and symlinks escaping the root. Assisted-by: Codex:gpt-6
Resolve the existing model directory before choosing a missing file. On macOS, resolving the missing leaf fails and otherwise compares an aliased path with a canonical allowed directory. Keep this test on the Upload error branch instead of the containment rejection branch. Assisted-by: Codex:GPT-6
Preserve staging capacity limits and request cleanup over worker tunnels. Adapt the release handlers and tests to the HTTP control transport. Verified 33 control-path, 148 node staging, and 123 worker specs. PostgreSQL integration specs require Docker, which is unavailable here. Assisted-by: Codex:GPT-6
Keep the worker control request test and canonicalize its temporary directory. Assisted-by: Codex:GPT-6
|
This is a large PR (100 files) with tests-linux and tests-apple CANCELLED (required checks), plus a gosec failure and backend build failures (magpie-tts-cpp, face-detect). Could you check if the cancelled tests are due to a cascading failure or a real issue? The gosec failure also needs investigation. |
Distributed mode without a message broker
Running LocalAI distributed used to mean running and managing a NATS cluster, and giving every worker a routable address with inbound ports open back to it. Both are gone.
A worker now dials the frontend load balancer over ordinary HTTP and holds one multiplexed yamux tunnel that carries gRPC, HTTP and websocket traffic together. That tunnel lands on exactly one frontend replica; every other replica reaches the worker by relaying through the owner. Workers bind loopback only and advertise nothing. A deployment needs PostgreSQL and the frontends' own HTTP listener, and nothing else.
github.com/nats-io/...is out ofgo.modandgo.sum.What a worker needs now
One outbound HTTPS connection to the same load balancer a browser would use. No routable address, no inbound ports, no broker credentials.
How it works
GET /api/cluster/connectand holds a yamux session. Control RPCs, file staging and inference all ride it.LISTEN/NOTIFY. Payloads above the 8000-byte notification cap become a row plus a notify-with-id, never a truncated payload.SELECT ... FOR UPDATE SKIP LOCKEDclaim queue. A claim is released when the replica holding it stops being live, with no age term, so a slow claimant is never robbed.The invariant this rests on
Four conditions must never be reported as each other: a routing fact, a connection absent within the reconnect grace, an unreachable peer, and the worker's own answer. Only the first and last may be acted on. Absence makes the scheduler reap rows and evict models, and one of those paths runs during inference, so a collapse between them is not a cosmetic bug.
Review found and fixed fourteen instances of that collapse, including one created by the fix for another, one where an expired caller deadline was blamed on the peer under contention, and one where a worker's own refusal was read as "no route", which left permanently unreapable ghost rows.
Operator impact
LOCALAI_ADVERTISE_ADDRandLOCALAI_ADVERTISE_HTTP_ADDRare no longer used.--nats-urlis accepted and ignored so existing command lines keep starting.--nats-require-authno longer gates agent-worker approval.Testing
A real-process end-to-end suite runs a two-frontend, two-worker cluster of live
local-aibinaries, which is the configuration where the owner path and the relay path are live at once:The suite previously could not fail honestly: with no binary present it ran zero specs, exited zero and printed "Test Suite Passed". It now builds the binary and refuses to run against one older than the tree. That guard has since caught three real mistakes.
The race detector runs clean across 98 suites and 7095 specs with no known-failure allowance.
Bugs found in existing code and fixed
Follow-ups closed before merge
The gaps found during the work were closed rather than deferred:
/api/cluster/peertook a self-declared replica id, so anything holding the shared registration token could relay to every worker a replica owned and repeatedly evict its inbound link. Each replica now mints its own secret at startup and publishes only the hash; the plaintext never leaves the process that minted it. A mixed-version rollout refuses rather than falling back, because an old replica and an attacker send byte-identical requests.Known remaining
DistributedManager.Listtruncates skill content to 500 bytes for agents, a separate pre-existing defect.🤖 Generated with Claude Code
https://claude.ai/code/session_01Y2TjpXdY7SszRrM5PhSp1e