Repository navigation
feat(distributed): add carrier state and switchable holders - #12549
Draft
mudler-agent wants to merge 148 commits into
Draft
mudler-agent wants to merge 148 commits into
mudler-agent wants to merge 148 commits into
Conversation
Add table cluster_carrier and a store for it. The row names the carrier that the whole cluster uses, an epoch that grows with every change, and the fields a later change of carrier needs: state, target and the carrier that is draining. Seed inserts the row only when it is absent, so replicas that start together agree on one writer. Transition is a compare-and-set on the epoch, so two writers cannot interleave. A check constraint keeps the table at one row. Nothing reads the row yet. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The backend and model managers and the replica reconciler took the NATS sender by its concrete type. Add NodeControl, the interface of every control call they make, and take it instead. The sender implements it, so nothing changes at run time. NodeControl names the three unload interfaces that the model loader reads by type assertion. A sender that lacked one would still compile and would silently lose that behaviour. Export the install-with-force fallback so a sender in another package can implement it. Add a ClientFactory option to the reconciler, so its default probe dials the way the router does. Add NewTokenClientFactory for the direct factory. Narrow the agent event bridge from MessagingClient to Broadcaster. It only publishes and subscribes. Pin that the NATS client returns a client for a server that is not up yet, which startup relies on. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
A carrier.Set is the implementation of every seam for one carrier. The holders implement the same interfaces as the seams (fan-out, work queue, control verbs, file staging, backend clients, worker dialer, agent RPC) and forward each call to the set that an atomic pointer names. They take no lock and add no allocation. The fan-out holder keeps a registry of its subscriptions. Listen attaches every subscription to a second set, so each handler listens on both carriers while a swap runs. Release drops the first set again. Publish always uses the set the pointer names. The holder also forwards the reconnect hook that syncstate and the failover sync register. NewNATSSet builds the set for the NATS carrier. A spec fails when any other non-test file imports the NATS libraries. Nothing stores a second set yet, so no behaviour changes. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…lders Read the cluster carrier row at startup, or insert it when it is absent, and build the set it names. An existing deployment has a NATS URL, so it seeds nats and runs as before. A row that names a carrier this build does not run stops the start with a message, and does not fall back. DistributedServices no longer holds the NATS client. The router, the health monitor, the reconciler, the job dispatcher, the agent bridge, the prefix cache and the other users of the bus take holders over the active set. The file manager is built before the set because the NATS file stager needs it. The NATS URL keeps its meaning. A URL the client cannot use still stops the start. A server that is not up does not: the client retries, as it did before. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
make test-e2e-distributed was not run by any workflow. Add one, with continue-on-error so a failure is visible on the pull request and does not block a merge until the job has a track record. The suite starts its own PostgreSQL and NATS servers, so the job pre-pulls both images. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Add a section on the holders, the carrier set, the subscription registry and the cluster_carrier row. Record the import boundary of the NATS libraries and the startup rule for an unreachable NATS server. Remove the open item for a reconciler client factory, which now exists. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The fan-out holder is a Broadcaster, so it has to pass the suite that every carrier passes. Run it on a holder as built, and on a holder whose carrier was swapped. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
New admin endpoint GET /api/models/storage reports disk usage of the installed models: each model's files and sizes, files shared between models counted once, and configured files that are missing from disk. The models WebUI page shows a usage summary, each model's file table, and the full file list with per-file status. Assisted-by: Claude Code:claude-fable-5 Signed-off-by: Plamen K. Kosseff <p.kosseff@gmail.com>
…nnections Each frontend replica records itself in the instances table and refreshes the row from a membership loop. The database clock stamps every row, so no replica decides on the clock of another. The row also carries the version and the epoch of the carrier that the replica has built, with a reason when it could not. The node_connections table records which live replica holds the connection of a worker. Every claim takes a fresh epoch from a sequence, so a replica that lost a worker cannot clear the claim of the replica that won it. A released or orphaned row is kept as a departure, and Presence tells a worker that is reconnecting from one that is gone only after a grace period. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
A worker that follows a change of carrier is attached to two carriers for a while. Add the columns that hold what the worker reports and the epoch it saw when it reported. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…istry Run the membership loop in distributed mode, on every carrier, and remove the row when the replica stops. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…yloads Move the subject matcher out of the test helper into messaging, so a carrier and the in-memory doubles use one rule. Add MaxBroadcastBytes, checked by the NATS client and by the fake bus, so a message that one carrier accepts is not refused by another. The conformance suite now takes two ends of a carrier and runs every broadcast root, the sizes around the notify cap and the bound. It also asserts what a carrier does with the control roots, and reads the drop counter when the carrier has one. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
… string The postgres image declares its data directory as a volume. A removed container leaves that anonymous volume behind, so every spec that started a database left tens of megabytes on the disk of the host. Put the data directory on a tmpfs. Also return the connection string, for code that needs a connection of its own, such as a LISTEN connection that cannot come from a pool. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…t and reply A carrier that has fan-out and no request and reply serves the broadcast roots only. ValidateBroadcastSubject applies the shared rules and then refuses the nodes and mcp roots with the same class of error as a root that is not served. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The second implementation of messaging.Broadcaster. A message goes out with pg_notify. When its encoded notification is 8000 bytes or more, the payload is written to the bus_messages table and the notification carries the row id. Every replica deletes the rows that are older than the retention, by the clock of the database. Delivery is at-most-once, as with NATS. The listener copies each notification into a bounded queue and never blocks, so the server is not held up by a slow consumer. A lost connection is dialled again and every channel is listened again, then the reconnect callbacks run. The subject rules are the shared ones, and the carrier refuses the control roots. It checks the shared payload bound. Counters for the messages published by path, the messages lost by stage, the spilled rows deleted and the reconnections go to OpenTelemetry, because a lost broadcast is otherwise visible only in a log. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
… order A broadcast of 64 KiB is written to a row, and every replica reads that row with a SELECT. One reader at a time limited a replica to about a thousand such broadcasts per second, and a burst above that was dropped. Read up to eight rows at the same time and hand the notifications to the handlers in the order they arrived, so the order of one subject does not change. In a burst of 4000 broadcasts of 64 KiB from 16 publishers, the replica with one reader delivered 1860, and the replica with eight readers delivered all of them. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
NewPgbusFanout opens the LISTEN connection and returns the broadcaster, the reconnect hook, a ready check and the close that a set needs. It opens the connection when it is called, so a deployment on NATS holds no LISTEN session until a change of carrier asks for one. Run the broadcaster conformance suite on holders that were swapped from one carrier to this one on a real database, and check that a message published after the swap reaches a subscriber made before it. The import-boundary spec now also keeps the PostgreSQL driver in the pgbus package, and the pgbus package behind the carrier builder. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
⬆️ Checksum updates in gallery/index.yaml Signed-off-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> Co-authored-by: mudler <2420543+mudler@users.noreply.github.com>
A worker holds one outbound websocket to a frontend and a yamux session runs on it. The frontend opens a stream for every request. The first frame of a stream names the local service of the worker, and the worker always answers with a reply frame before the carried protocol starts. The four refusals stay apart, so that a caller acts on the three that are evidence about a backend and does nothing about the one that says the worker learned nothing. The upgrader and the dialer come from one constructor each, with 64 KiB buffers. With the 4 KiB default of gorilla, a stream that sends from the worker to the frontend ran at 560 to 600 MB/s on a loopback socket and the other direction at 3,200 to 3,800 MB/s. With 64 KiB the first direction runs at 2,900 to 3,000 MB/s. BenchmarkTransfer repeats the measurement. A session has a lane. The inference lane keeps the yamux windows. The bulk lane uses 4 MiB and 32 MiB, which a transfer needs on a link with a long round trip, and which would let more data queue ahead of a small call on a session that is shared. Only the files of the tunnel package import yamux. The import boundary spec lists it. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Splice used to close both streams as soon as one direction ended. A protocol that closes its sending side and then waits for the reply lost the reply that was still in flight. Now a direction that ends clean passes the end on with CloseWrite and the other direction goes on until it ends too. When the destination cannot half-close, or a direction fails or ends because a stream was closed, both streams are closed at once, as before. This is what wakes a copy that is parked in a read or in a write. A copy into a TCP connection wraps the error in a net.OpError. The classifier of the endings of yamux now removes that layer, so a stream that this side reset is not reported as a failure of the transport. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
A proxy in the spec works as one link of 1 Gbit/s with 10 ms of delay in each direction and one queue that every connection shares. A transfer of 256 MiB runs from the worker to the frontend while a small call goes every 5 ms in the other direction, and its reply has to wait in the queue of the loaded direction. The spec runs three set-ups: a probe on its own TCP connection, a probe on the session of the transfer, and a probe on the inference session with the transfer on the bulk session. It requires the bulk lane to stay within 10 ms of the TCP baseline at the 99th percentile, and the shared session to be worse than the bulk lane. A build with the race detector moves 128 MiB. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The registry keeps the sessions that workers opened to this replica and keeps the node_connections table in agreement with them. Attach writes the claim first and stores the session second, so a failed claim leaves the replica as it was, and the claim and its record are one step for each node. Detach matches the token of the attachment by equality, so a session that was replaced cannot remove the one that replaced it. Open returns ErrNotOwner only when the node is not held here. A session that ended under a held entry, and a caller whose budget ran out, are reported as themselves. A node has two lanes. The inference lane owns the claim. The bulk lane is stored under the claim of the inference lane, is refused with ErrNotOwner on a replica that does not hold it, and is closed with it. Open on the bulk lane uses the inference lane when the node has no bulk session, so a transfer on an older worker or during a reconnect works. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
When a peer sweeps a replica that stalled, it deletes the instance row of that replica and records a departure for every connection it owned. The replica registers again and rebuilds the instance row, but the sockets are still open on it. The other replicas then answer "not connected" for workers that are connected. The membership loop now takes a Reclaimer and calls it after a successful registration that follows a sweep. The interface keeps the package free of the code that holds the sockets. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
A node gets its own secret for the tunnel. Registration mints a new random token with 128 bits or more, returns it once, and stores only its SHA-256 in a new column. The registration token of the deployment never opens a tunnel, so a leak of it does not let its holder be any worker. Every registration mints a new token, and the worker reads it each time it dials. A node waits for approval with a token that the connect route refuses with 403 until the approval. A node type that holds no credential has the column cleared. The answer to a registration now names the active carrier and its epoch when the endpoint has a carrier reader, and hands over the credential for that carrier only: the tunnel token when the tunnel is active, the NATS user JWT otherwise. Without a reader, or when the row cannot be read, the answer is the one that a deployment with only NATS always got. NATS still requires the address of a backend worker. With the tunnel active the address is optional, and a worker that says that its address is not routable has it dropped and cleared. When the tunnel is active and registration is open, the frontend logs one warning. It does not refuse. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
GET /api/cluster/connect takes the dial of a worker. It reads the bearer token before it checks anything else, so an anonymous dial gets 401. A frontend with no registry answers 503 to a caller that has a token. An unknown node gets 401 and not 404, a lookup that fails gets 500, and a node that waits for approval gets 403. The token is compared in constant time with the hash on the row of the node. The registration token of the deployment is never accepted. The upgrade uses the shared upgrader with 64 KiB buffers. The query parameter lane picks the session. A bulk dial to a replica that does not hold the inference session of the node gets 409 before the upgrade, so the worker can dial again at once and land on the owner. The route is registered on every replica. The global auth middleware lets exactly this path through to its own check. The prefix /api/cluster/ is not delegated. The route coverage test walks the route and a new spec fails if the route is removed. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The proxy that works as a link with a rate and a delay is now a package, so that the specs of the worker can use it on the real endpoint. It also passes the end of a connection on behind the data that is still queued. Before, a connection that ended was closed at once and the data in the queue was lost. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
A worker that is attached to the tunnel is ready when it holds a tunnel session. The probe is the counterpart of the one for NATS, and its error is a different one, so that /readyz says which link is down. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…amed The credential manager is no longer only for NATS. The answer to a registration names the carrier and its epoch. When the carrier is the tunnel, the manager waits for the tunnel token and not for the NATS JWT, and it does not wait for the approval of the node: the frontend mints the token at registration and refuses the dial of a node that is still pending, so the tunnel client waits for the approval itself. The token of the latest registration is the one that the client presents, because every registration mints a new one. For NATS, or for a frontend that names no carrier, the manager works as before. RegisterWithRetry now uses RegisterFullWithRetry, which returns the whole answer. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The worker dials the frontend and serves the streams that the frontend opens. It holds two sessions. The inference lane carries model calls and health. The bulk lane carries file transfers and lives only as long as the inference lane. A dial of the bulk lane that lands on a replica which does not hold the inference lane gets 409, and the worker tries again soon. A frontend that refuses the bulk lane does not cost the worker its inference lane. A stream is routed by the tag in its first frame. The host that the frontend names is dropped, and the port must be in the range of the allocator of the worker, so the tunnel cannot reach the LAN of the worker. A refusal says which of four things is true, and a late frame, a timeout and a cancelled context are never evidence about a backend. The wait between two attempts grows from half a second to thirty seconds with jitter, and it goes back to the floor only after a session that lasted. The token is read at every dial. The specs run the client against the real connect endpoint and a real PostgreSQL. Behind a link of 1 Gbit/s with 10 ms of delay in each direction, a transfer of 256 MiB delays a small call on the same session to a p99 of 88 ms, and on the bulk lane to 34 ms, against 31 ms for a call on its own TCP connection. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
… its changes A frontend no longer refuses a cluster that runs on the tunnel. It reads the row, seeds it when no replica has (the tunnel for a deployment with only PostgreSQL, NATS for one with a NATS URL, which keeps the behaviour of an existing deployment), and builds the set the row names. For the tunnel that is pgbus, the claim queue, the control client and the dialers behind the holders. For NATS it is what it was. A NATS URL that cannot be parsed stops the start, and one whose server is not up does not. The NATS URL becomes a setting of the cluster. The flag is copied there when none is stored, and a stored URL wins, so that every replica uses the same one. While the tunnel is the carrier in use or is draining, each replica claims queued work and drives it on agent workers, through the dialer of the tunnel set and not through the holder, so that a run that started on the tunnel finishes there. When the tunnel is released, the rows that are still pending are taken out of the table and published to the queue of the carrier that took over. Every replica now follows the row: a swapper that polls it and wakes on a hint, a leader loop that drives the protocol under an advisory lock, and a report of which carriers the replica could build. The report is refreshed when a dry run asks for it. The routing of calls to workers during a change is wired into every holder and into the dialer, and the draining carrier forgets a node when it is removed. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
GET /api/cluster/carrier reports the active carrier, the epoch, the state of a change, each live replica with its version and readiness, the workers that could not follow a change to the other carrier and why, the work in flight, the waits of a change and the live replica list that the admin confirms before a change. POST /api/cluster/carrier takes a target with dry_run and force, or abort. A dry run asks every replica to look at what it can build, and answers from what they wrote, and a request runs the same preflight and starts the change. The answer is 202 because the replicas carry out the change and one of them leads it. A blocked request is 422 with the blockers and the report, a change already under way is 409. GET and PUT /api/cluster/settings store the NATS address that frontends use, the one workers are told to use, and the waits of a change. Saving a NATS address makes NATS available and checks that this replica reaches it with its own credentials. It does not switch, and the switch stays an explicit request. A password in an address is masked when it is read back. No credential is stored. The routes sit under /api/cluster/, which is not a public prefix. Only the exact paths of the connect and peer routes skip the global authentication, and these four answer to an admin, so an anonymous caller gets 401 on each of them. They are registered on every frontend, and a frontend that is not distributed answers 503. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…eal broker and a real database The specs run frontend replicas the way the application builds them: a NATS set over a real server, a tunnel set over pgbus and the claim queue, the holders, the swapper, the window and a leader loop under the advisory lock. A worker answers on both carriers, and an agent worker consumes on each. From NATS to the tunnel, with an install, a held run, renewals of a load and chatter from both replicas in flight: the change settles in about half a second, every message is heard once by every replica, every job runs once on the carrier it was queued on, the install and the run end where they began, and no renewal fails. Back to NATS, the units that were queued on the tunnel run once on NATS and the LISTEN connection is closed. A replica that dies in prepare does not hold the change, and the one that starts again follows the row. A leader that dies between the commit and the settle is replaced by the other replica. An abort in prepare drops what was built for the change on every replica. A worker that cannot follow blocks the change, and a forced change leaves it reachable through the window and then unroutable, and reaps nothing. The claim loop and the handover of a tunnel set are built by NewTunnelSet from the options, so the specs and the application use the same code. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The guide says how the cluster row, the leader, the swapper and the window work together, what the preflight lists, how a run that the end of a drain cuts off is settled, why the NATS set guards loopback addresses of workers that hold a tunnel, and where the settings of the cluster live. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…itch The metrics of a replica show the epoch of the row, the state of a change, the carrier in use, whether the replica is ready, and how long the drain has to go. The status of the switch reports the row, the replicas, the workers and the time left of a drain without a target, and blocks nothing. A deployment with no NATS URL has no NATS bus, so the warning about a bus without credentials is no longer printed for it. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The window routes a stop of a model by node, and the holder must still forward each method to the method of the same name on the set. The specs of the holders and of the window said so and were out of date after the first version of the window. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…1-holders Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
| sort.Strings(names) | ||
| h := sha256.New() | ||
| for _, name := range names { | ||
| data, err := os.ReadFile(filepath.Join(dir, name)) |
The spec stored the length of a payload in an int32, which gosec reports as an integer overflow conversion. A 64-bit counter holds any length, and the spec reads it as an int as before. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
A worker that the frontend hands the address of NATS to has no file for the CA of that server. TLSFiles takes the CA as a PEM bundle as well as a path, validates that it holds a certificate, and trusts it, together with the file when both are set. The specs connect to a server of a private CA over TLS with the bundle, and show that the same server is refused without it. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
A worker that follows a change of carrier registers again to get the credential of the carrier it attaches to. The registration takes the carrier that the worker asks for. The frontend gives the credential of that carrier when it is the active one, the target of a change under way, or the one that drains, and the credential of the active carrier for anything else, so that a node cannot get a tunnel while the cluster has none. Only that credential is minted: a worker that holds the tunnel and registers for NATS keeps its tunnel credential and its session. The answer names the carrier, the epoch, the state of a change, its target and the drain, and the carrier that the credential is for. A credential for NATS comes with the address that workers are told to use, the CA of the server as PEM, and a flag when the server asks for a client certificate, which cannot be handed over. A worker needs no NATS setting of its own for this. The answer to a heartbeat carries the same news, because the heartbeat reaches the frontend whichever carrier is active. A heartbeat can carry the carriers that the worker is attached to, the epoch it saw, the carriers it can follow and why it cannot follow another. They are stored for the frontends to route by, and a worker that sends none is left as it was. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The switch guessed where a worker was attached from the tunnel presence it held, because nothing reported it. Workers now report the carriers they are attached to in every heartbeat, so the guess is gone. A worker that reports its capabilities is attached to what it reports, and to nothing when it reports none: it is between two carriers and is not said to hold one it may not. A worker that reports nothing predates carrier switching, has only NATS, and is on NATS. The reader no longer takes a presence source or a reconnect grace. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…the heartbeat answer A worker that follows a change of carrier waits for the credential of the carrier it attaches to, which is not the active one while the change is prepared. A credential manager can be made for one carrier. It refuses an answer that carries the credential of another carrier, as a frontend that does not allow the request would send, and tries again. The registration answer carries the state of a change, the address of NATS, its CA and the client certificate flag. The heartbeat client can return the answer of the frontend, which names the carrier, the epoch, and the target of a change. A frontend that predates carriers answers without them. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
A worker used to attach to one carrier at start and keep it. It now reads the answer of every heartbeat, which reaches the frontend whichever carrier is active, and follows the carrier the frontend names. While a change is prepared it attaches to the target as well: it registers again for the credential of that carrier, waits a random time of up to ten seconds so that a fleet does not connect at once, and connects with a wait that grows from one to thirty seconds when the carrier refuses. It keeps the previous carrier open through the drain, tells the frontends where it is attached in every heartbeat, and closes a carrier only when the cluster released it, the carrier it moved to is connected, and no request runs on it. A backend process is not touched. A change that is aborted closes the target again. A worker that cannot follow keeps the carrier it has, says why in its heartbeat and tries again. It cannot follow to NATS without an address that the frontends can dial, or with a server that asks for a client certificate it does not have. A worker that boots on NATS counts as routable, because its backends listen on every interface. A frontend that predates carriers leaves the worker as it was. The NATS side takes the address, the CA and the credential from the frontend. A local address or CA still wins, so the flag keeps its meaning. The tunnel is attached and released through the same path at start and later, and the control plane of the HTTP server is empty whenever the worker holds no tunnel. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
An agent worker serves its runs, its MCP requests and its jobs on the carrier the cluster uses, and now follows a change as a backend worker does. The handlers are the same on both carriers, and the worker attaches them to NATS or to the tunnel as the frontend names the carrier: on NATS it connects with the credential and the address the frontend hands over (a local address or credential still wins), consumes the queues, and serves the MCP subjects; on the tunnel it serves the same verbs behind its loopback server. A run keeps the carrier it began on, and the carrier is not closed while a run is in flight on it. The agent worker opens no port, so the only thing that can stop it from following is a NATS server that asks for a client certificate, which cannot be handed over. It reports that, and its heartbeat carries the carriers it is attached to like the heartbeat of a backend worker. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…of carrier The specs of the change no longer attach a fake worker by hand. A backend worker and an agent worker register through the real endpoints of the frontends, attach to the carrier the cluster names with the planes and the follower that Run uses, and follow each change by themselves from the answer of their heartbeat. The preflight reads what they report. From NATS to the tunnel and back with an install, a model that is loading on the backend, a held run and the renewals of a lease in flight: the workers attach to the target, keep NATS through the drain, release it when the cluster does, and nothing is lost or repeated. The backend answers over the new carrier without a restart. Two workers follow at once. A worker that restarts while a change is prepared follows it, and one that restarts after it starts on the carrier the cluster has, with no flag. An abort makes the worker let go of the target. A worker that predates carrier switching blocks a change and is left on NATS, unroutable and unreaped, by a forced one. A worker with no address is listed as unable to follow to NATS with the reason it reported, and stays on the tunnel after a forced change. The planes of a backend worker and its verbs are exported so that the specs can serve verbs of their own on the real carriers. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…he cluster settles After a forced change, a worker that could not follow heartbeats over HTTP, which exists on both carriers, and cannot be reached. The scheduler kept choosing it, found no route, marked it unhealthy, and the next tick promoted it again, so a request that landed on it failed until the loop ended. The health monitor now reads what each worker reports. A worker that reports its carriers and is not attached to the active one is demoted, and it stays demoted until it reports the active carrier. Nothing is deleted: the status changes, its models stay, and its backend keeps running. Nothing is demoted while a change is under way or a carrier drains, because the window routes the workers then, and a worker that reports nothing is left to the tunnel read as before. A replica runs the check when a carrier is released, so the demotion does not wait for the next tick. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The guide says how a worker learns the carrier from the answer of its heartbeat, how it gets the credential of the target and attaches with a random wait, what it reports, when it releases a carrier, why it can fail to follow, and how the health monitor treats a worker that does not follow. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
The job registration of an agent worker now receives the carrier it is called for, and the agent worker command accepts the new signature. A worker that fails to follow for one reason reports that reason alone, without the name of the carrier in front of it, because the preflight already names the carrier. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…A path as operator input The heartbeat client ignored the error of closing the body, and the read of the CA file for the handover of NATS takes a path that the operator configured, which the scanner reports as a file inclusion. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…fb79603737547a` (#12554) ⬆️ Update PrismML-Eng/llama.cpp Signed-off-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> Co-authored-by: mudler <2420543+mudler@users.noreply.github.com>
…2ae26adc08c1aa4e329805` (#12553) ⬆️ Update ServeurpersoCom/omnivoice.cpp Signed-off-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> Co-authored-by: mudler <2420543+mudler@users.noreply.github.com>
…07ba0c07` (#12416) ⬆️ Update mudler/vllm.cpp Signed-off-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> Co-authored-by: mudler <2420543+mudler@users.noreply.github.com>
…33e8953b728` (#12461) * ⬆️ Update ggml-org/llama.cpp Signed-off-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> * fix(llama-cpp): refresh decision patch contexts Preserve the upstream decision-order assignment when adding score logits. Refresh the TTS metadata context after upstream realigns its comments. Both patches apply to the new pin, and the gRPC source compiles. Assisted-by: Codex:gpt-6 --------- Signed-off-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> Co-authored-by: mudler <2420543+mudler@users.noreply.github.com> Co-authored-by: localai-org-maint-bot <306269227+localai-org-maint-bot@users.noreply.github.com>
…ff0fb212c8` (#12421) * ⬆️ Update 0xShug0/audio.cpp Signed-off-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> * fix(audio-cpp): mirror the turn detection task The upstream enum adds TurnDetection, so the exhaustive conversion fails to compile. Extend both conversions and preserve the RPC admission rules. Test that turn detection cannot route through VAD or another existing RPC. Assisted-by: Codex:gpt-6 --------- Signed-off-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> Co-authored-by: mudler <2420543+mudler@users.noreply.github.com> Co-authored-by: localai-org-maint-bot <306269227+localai-org-maint-bot@users.noreply.github.com>
* fix(mcp): connect to legacy SSE servers Remote model MCP connections only use Streamable HTTP, so legacy SSE servers fail initialization. Retry with SSE when the initial POST returns 400, 404, or 405. Share one discovery timeout across attempts and cancel failed connections without truncating successful sessions. Preserve HTTP policy and reject foreign SSE message endpoints before attaching credentials. Add SDK integration tests and document automatic transport selection. Assisted-by: Codex:gpt-6 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(mcp): avoid narrowing HTTP status codes Store the initialization status in a 64-bit atomic value to avoid the integer overflow conversion reported by gosec. Assisted-by: Codex:GPT-6 gosec Signed-off-by: Ettore Di Giacinto <mudler@localai.io> --------- Signed-off-by: Ettore Di Giacinto <mudler@localai.io> Co-authored-by: localai-org-maint-bot <306269227+localai-org-maint-bot@users.noreply.github.com>
The specs cover the missing address, the client certificate that cannot be handed over, a worker with no frontend URL for a tunnel, and the choice between the CA of the operator and the one the frontend hands over. Assisted-by: Claude Code:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
…1-holders Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
This is the first slice of a series that makes the transport of a distributed deployment switchable at runtime, so a cluster can later move between NATS and a database-backed carrier without a restart. This slice adds the indirection and the stored state that later slices use. NATS is the only carrier. It changes no behaviour for an existing deployment.
It targets the long-lived integration branch
feat/distributed-live-carrier, notmaster. Later slices merge into that branch, and it lands onmasteronce.What it adds
cluster_carrierhas one row: the active carrier, an epoch, the change state, the target and the carrier that is draining.cluster.CarrierStorereads it, inserts it when it is absent (Seed), and changes it with a compare-and-set on the epoch (Transition). A check constraint keeps the table at one row.core/services/carrier). Acarrier.Setis the implementation of every seam for one carrier. One forwarding holder per seam loads anatomic.Pointer[carrier.Set]once per call and calls the same method on that set: fan-out, work queue, control verbs, file staging, backend clients, worker dialer and agent RPC. The holders take no lock and add no allocation.Listen(next)attaches all of them to a second set, so each handler listens on both carriers.Release(old)drops the first set again. Publish always goes to the set the pointer names. The holder also forwards the reconnect hook thatsyncstateand the failover sync register.nodes.NodeControl. The managers, the reconciler and the model loader took the NATS sender by its concrete type.NodeControlis the interface of every control call they make, so they can take a holder. It names the three unload interfaces that the model loader reads by type assertion, because a sender that lacks one still compiles and silently loses that behaviour.initDistributedreads the cluster row, or inserts it, and builds the set it names.DistributedServicesno longer exposes the NATS client. Code that only publishes and subscribes gets amessaging.Broadcaster.github.com/nats-io/*.tests-e2e-distributed.ymlrunsmake test-e2e-distributed. No workflow ran it before. The job is advisory (continue-on-error).Behaviour
natsat epoch 1 and every frontend follows it. Subjects, payloads, queue groups, direct gRPC dials and registration are unchanged.--nats-urlandLOCALAI_NATS_URLmean what they meant before.messaging.Newretries a failed connect, so a frontend can start before its broker. A new spec pins this.Listen,Release,NotifyReconnectandTransitioncallers are in later slices. Until then only specs call them.Design choices to review
NodeControlis wider thanNodeCommandSender. The managers, the reconciler and the model loader also use the process lister, the unload interfaces,InstallTimeout,DeleteModelFilesand the install-with-force fallback. I exported that fallback asInstallBackendForceso a sender in another package can implement it.WorkConsumer. Only the agent worker consumes work, in its own process, and it never swaps. A holder would have no caller.Broadcaster. It only publishes and subscribes. The frontend's bridge goes through the holder; the agent worker still builds its own on its NATS client.ReplicaReconcilerOptionsgets aClientFactory, so the default probe dials the way the router does.Transition.changed_atis for operators and nothing decides on it, so replicas need no agreed clock.go listskips files whose build constraints do not match the host, and the repository does not run depguard.--nats-urlstays the only source of the URL here.Testing
NodeControland of the file stager is forwarded to the method of the same name on the set the pointer names; a call already running finishes on its set while a new call uses the next one; subscriptions attach to both sets during a swap and detach on release; a subscription made whileListenruns is neither lost nor doubled; no call blocks across a swap; the hot path allocates nothing (testing.AllocsPerRun).Local runs with
-race:core/services/...,core/application/...,pkg/mcp/...,core/http/....golangci-lintis clean on the touched packages,go vetis clean for linux,GOOS=darwinand fortests/e2e/distributed, andGOOS=windows go buildpasses.gosecfinds nothing in the new packages. Failures that are also onmaster: a data race inside thegalleryoptest fixture. One timing-sensitivenodesspec (a 20 ms stall window) failed once on a loaded machine and passes alone and in a full run.A real smoke run compared this branch with
master(PostgreSQL and NATS in containers, two frontends, one worker, a mock backend): load through both frontends, a failing load (503,Retry-After: 15), a load that outlasts the wait (503,Retry-After: 5),POST /api/models/{id}/load-cancel(200, process and replica row gone), worker restart with SIGTERM and with SIGKILL, and shutdown. Status codes, headers and bodies were identical. The only differences are the new log lineCluster carrier seededand the order of two startup log lines. A frontend frommasterand then one from this branch on the same database also worked, and a row set to a carrier this build does not run stopped the start with a clear message.Not tested: a worker from an older release against this branch, TLS and JWT-protected NATS in the smoke run, and more than two frontends.
🤖 Generated with Claude Code