Files
LocalAI/core/application/distributed.go
T
mudler-agentandEttore Di Giacinto e4fa051ee3 refactor(distributed): put the NATS-only paths behind interfaces (#12395)
* feat(messaging): add shared subject rules

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* test(messaging): cover BroadcastRoots, ControlRoots and SubjectRoot

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* feat(messaging): add Broadcaster and enforce subject rules in every carrier

Broadcaster is the fan-out half of MessagingClient. The NATS client and
the in-memory FakeBus now refuse a subject outside the served roots and any
wildcard other than a whole single token, and FakeBus shares MatchSubject
instead of its own copy. FakeBus Unsubscribe now removes its own
subscription instead of the first one with the same subject.

A shared conformance suite in messagingtest runs against both carriers.
The distributed e2e specs that used invented test.* subjects, and the one
that subscribed with a > filter, now use subjects from subjects.go.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: depend on Broadcaster where only publish and subscribe are used

Narrowed to messaging.Broadcaster: nodes/staging_progress.go,
nodes/install_progress_publisher.go, galleryop/operation.go,
galleryop/service.go, agentpool/user_services.go, agentpool/agent_jobs.go,
openresponses/store.go, openresponses/sync.go, syncstate/syncstate.go,
finetune/service.go, quantization/service.go and
failover/distsync/distsync.go. SubscribeJSON now takes a Broadcaster
because it only calls Subscribe, which lets the narrowed consumers use it.

Stayed wide: worker/supervisor.go, because its client field also serves
the SubscribeReply handlers in worker/lifecycle.go. The request/reply,
queue and wiring files (nodes/unloader.go, nodes/file_stager_s3.go,
jobs/dispatcher.go, agents/dispatcher.go, agents/events.go,
worker/file_staging.go, cli/agent_worker.go, http/app.go) are unchanged by
design.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): name the no-route condition and confine the carrier error

Consumers matched nats.ErrNoResponders, which names an absence, to demote a node. They now match ErrNoRoute, the control path maps the carrier's failure onto it, and timeouts and worker refusals are pinned as not being no-route.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs(nodes): state which FileStager implementations return ErrNoRoute

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): build backend clients through one node-aware seam

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: describe the distributed transport seams

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: correct comments that overclaim after the seams refactor

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* fix(agent-worker): refuse an unserved LOCALAI_AGENT_SUBJECT at startup

The messaging client now refuses a subject whose root no carrier serves.
An agent worker started with a custom LOCALAI_AGENT_SUBJECT such as
tenant-a.agent.execute used to start and then wait on a subject the
frontend never publishes to. After the subject rules landed it exited
at subscribe time with an error that did not name the setting.

Behaviour change: the worker now checks LOCALAI_AGENT_SUBJECT before it
registers or connects, and exits with an error that names the variable
and says to use a served subject under the agent root, for example
agent.execute. The served roots are not widened: a custom root was
never delivered by the frontend, and a wider set would reopen the
drift the subject rules exist to close. The flag help and the agent
worker docs state the constraint.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* test(nodes): pin the reactions to ErrNoRoute

Three callers react to ErrNoRoute and had no spec: the reconciler's
upgrade drain falls back to the legacy forced install, the reconciler
marks the node unhealthy when a pending op has no route, and the
backend-op fan-out marks the node unhealthy. Each spec drives the real
caller with a scripted no-responders reply and reads the result from
the registry or the recorded requests.

A fourth spec pins the other side: a pending op that times out leaves
the node healthy and only counts the attempt, so mapping timeouts onto
ErrNoRoute would fail here.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* test(messaging): pin client subject checks and fail the carrier suite in CI

Add specs that call Publish, Request, Subscribe, QueueSubscribe,
SubscribeReply and QueueSubscribeReply on a client with no connection.
Each call must return ErrUnservedSubject for bogus.thing and
ErrUnsupportedWildcard for jobs.>. This proves that the subject check
runs before the connection is used, and needs no server.

The NATS conformance suite is the only check that runs the subject rules
against a real carrier. Before this change it skipped without output
when Docker was missing. Now it fails when CI is set, so a Linux runner
without Docker cannot hide it. It still skips on local runs and on macOS
CI, which has no Docker.

Add SubjectNodeBackendInstallProgress to the list of constructors that
must build served subjects, and ask contributors to extend the list.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: state what ErrNoRoute may change, and group the distributed guides

The seams note said MarkUnhealthy was the only state change allowed
on ErrNoRoute. A pending backend op still records the failed attempt,
counts toward the reconciler's retry limit and is dead-lettered after
the maximum attempts. The note now says that MarkUnhealthy is the only
change to the node's own state, and that the per-op accounting is not
a verdict about the node.

The note also documents that the NATS conformance run fails under CI
when Docker is missing. The distributed-seams row moves next to the
distributed-state row in the topics table. The liveness ping spec
header now says no route is a reason to skip the worker, not proof
that the worker is gone.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): give the backend client factory the node id

Mechanical: the method gains a nodeID parameter and the eight test fakes are updated. No behaviour change.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): drop the optional node-aware factory

The node id is now in the main method, so the optional interface and its helper had no behaviour of their own.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): dial backend probes through the client factory

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): dial workers' file servers through a per-node dialer

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(http): proxy backend logs through the per-node worker dialer

The admin backend-logs proxy (list, lines and the WebSocket stream) now reaches a worker through the same per-node dialer as the HTTP file stager, so every frontend-to-worker dial goes through one seam. The shared direct dialer keeps alive for 15s where the proxy used 30s. Harmless for requests bounded at 15s.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* fix(http): keep the backend-logs proxy independent of the admin connection

The proxy request had no context before the dialer change and is bounded only by its 15s timeout. Keep it that way.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: move the worker control payloads to workerctl

Mechanical move of the request and reply structs, the install progress event and the file payloads out of messaging. The verbs no longer belong to one carrier. No alias is left behind.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(worker): serve the lifecycle verbs through a controlServer

The worker registers one handler per verb and a NATS server maps each verb to its subject. Registration errors now name the verb. node.stop is served with SubscribeReply, which is identical on the wire because the handler never replies.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(worker): report install progress through the control sink

Install and upgrade now emit download progress through the sink the control server hands them. The debounce and the terminal flush stay in the handler path, built over that sink by the new nodes.NewDebouncedInstallProgressSink, which replaces NewDebouncedInstallProgressPublisher. The subject and payload on the wire are unchanged. The supervisor no longer holds the bus, and installFn and upgradeFn let specs drive both verbs without a gallery.

The malformed-request log lines are restored for install, upgrade, backend.delete, model.unload, model.stop and model.delete, with the reply bytes unchanged. The signal adapter is renamed noReply, which also lets worker.go import os/signal without an alias again.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(worker): serve the file-staging verbs through a controlServer

An empty list-dir answer is now {} rather than {"files":null}, because the typed reply omits an empty Files slice. The frontend decodes both to a nil slice in nodes/file_stager_s3.go.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* feat(messaging): add WorkQueue and the NATS producer

This is the producer side of the competing-consumer seam. The work kinds map one to one to today's subjects and queue groups: task to jobs.new and mcp-ci to jobs.mcp-ci.new (both in group workers), agent-run to agent.execute (group agent-workers). Enqueue publishes the payload as Publish does today, with one JSON marshal. FakeBus now records queue groups and keeps reply handlers so later specs can pin and drive them.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* feat(messaging): add the NATS WorkConsumer

An in-flight limit of one runs the handler inline on the delivery goroutine, as the MCP CI consumer does today. Any other limit spawns per delivery, as the agent consumer does. Queue groups are unchanged.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: publish queued work through WorkQueue

The job dispatcher, the agent pool and the agent scheduler enqueue through messaging.WorkQueue; the NATS implementation publishes to the same subjects as before. DistributedServices builds the queue next to the NATS client and hands it to the dispatcher and the agent pool, whose distributed mode switch now reads a non-nil WorkQueue. The unused AgentPoolService.SetNATSClient is removed.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: consume queued work through WorkConsumer

The agent dispatcher and the MCP CI consumer register through messaging.WorkConsumer. The NATS implementation keeps the inline one-at-a-time model for MCP CI and the per-delivery model for agent runs. handleMCPCIJob reports on the events publisher the carrier hands it instead of a captured client.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: delete the consumers nothing in production reached

jobs.new has a producer and no production consumer, and the agent dispatcher's Dispatch was only called from tests. Publishing jobs.new is unchanged.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(mcp): send MCP requests to agent workers through AgentControl

Timeouts still honour only the deadline, not cancellation, exactly as today. The NATS no-responders error maps to ErrNoRoute and a timeout does not.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(agent-worker): serve MCP requests and backend.stop through agentRPCServer

The agent worker's MCP tool and discovery reply subscriptions and its backend stop listener move behind an unexported agentRPCServer interface, served on NATS by nodes.NATSAgentRPCServer. The handlers become typed mcp.ToolHandler and mcp.DiscoveryHandler values that answer every failure with a reply carrying Error.

Queue group (agent-workers), inline execution on the delivery goroutine, the background handler context, the unmarshal error reply texts and the reply-less backend stop subscription are unchanged. The backend stop handler takes the decoded backend name, so it can still close that backend's MCP sessions.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(messaging): remove helpers that only tests used

BroadcastRoots, ControlRoots and SubjectRoot had no production caller. The roots spec now asserts every served root through ValidateSubject instead. MatchSubject moves back into the test support package, the only place that used it, with its table. NATSAgentRPCServer drops the subscription list it stored and never read, and NewNATSAgentRPCServer gets a doc comment.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* test(mcp): round trip the agent RPC server over a real NATS server

One spec sends a tool request and a discovery request through NATSAgentControl to NATSAgentRPCServer and checks that the handlers see the decoded requests and the replies come back. It also puts an undecodable body on the tool subject and checks the server answers with an unmarshal error instead of leaving the requester to time out.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: describe the distributed transport seams

The developer note now lists the final seams: fan-out, queues, both halves of the control verbs and of agent RPC, and the dial. It records the open items a second carrier has to handle.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* test: pin the in-flight limit each queue consumer asks for

The work queue specs pin what Consume does for a given limit, but nothing
pinned which limit each production consumer passes. Changing the agent
worker's MCP CI limit from 1 to 0 would have let MCP CI jobs run
concurrently on each worker with every test green.

Move the MCP CI Consume call into startMCPCIConsumer with the same wiring
and pin that it asks for (WorkMCPCI, 1). Pin that NATSDispatcher.Start asks
for (WorkAgentRun, maxConcurrent) for several limits.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: remove helpers the branch left without a caller

SubjectJobCancelWildcard lost its last subscriber when the frontend
stopped listening on jobs.*.cancel; the NATS permissions and conformance
suite spell the subject out, so nothing reads the constant.

decodeBackendStopRequest returned a stopAll flag that production dropped
and only a test read. decodeBackendStop is now the single decoder with the
same semantics: an empty body is stop-all, an empty Backend is stop-all,
malformed JSON is an error. stopBackends still derives stop-all from
Backend, so no reply changes.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* fix(messaging): keep an explicitly empty agent queue a plain subscription

Before the work queue seam the agent worker passed LOCALAI_AGENT_QUEUE
straight to QueueSubscribe, so an explicitly empty value made a plain
subscription and every agent worker ran every agent run.
WithAgentRunRoute replaced an empty queue with agent-workers, which
silently changed that.

Keep the queue as given once the option is applied. An empty subject still
falls back to agent.execute, since it never had a meaning of its own. The
flag default stays agent-workers, so only an explicitly empty value
reaches this.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: correct comments and record the PR B notes

Fix the recordingFactory comment (it also records the parallel flag),
document that a negative maxInFlight is unbounded and that Unsubscribe from
a handler deadlocks, and say a permanently undecodable payload returns nil.

Record controlHandler's undecodable return as a kept exception, and add the
second carrier notes to the developer note: the reconciler has no
ClientFactory option, the logs proxy honours HTTP_PROXY, verbs one carrier
serves need an opt-out, terminal replies come from the result event, and
agent runs publish through the NATS-bound EventBridge, which is not an
additive change.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

---------

Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
2026-10-02 23:58:41 +02:00

512 lines
20 KiB
Go

package application
import (
"context"
"encoding/json"
"fmt"
"io"
"strings"
"sync"
"time"
"github.com/google/uuid"
"github.com/mudler/LocalAI/core/config"
mcpTools "github.com/mudler/LocalAI/core/http/endpoints/mcp"
"github.com/mudler/LocalAI/core/services/agents"
"github.com/mudler/LocalAI/core/services/distributed"
"github.com/mudler/LocalAI/core/services/jobs"
"github.com/mudler/LocalAI/core/services/messaging"
"github.com/mudler/LocalAI/core/services/monitoring"
"github.com/mudler/LocalAI/core/services/nodes"
"github.com/mudler/LocalAI/core/services/nodes/prefixcache"
"github.com/mudler/LocalAI/core/services/storage"
"github.com/mudler/LocalAI/pkg/distributedhdr"
"github.com/mudler/LocalAI/pkg/sanitize"
"github.com/mudler/xlog"
"gorm.io/gorm"
)
// DistributedServices holds all services initialized for distributed mode.
type DistributedServices struct {
Nats *messaging.Client
WorkQueue messaging.WorkQueue
AgentControl mcpTools.AgentControl
Store storage.ObjectStore
Registry *nodes.NodeRegistry
Router *nodes.SmartRouter
Health *nodes.HealthMonitor
Reconciler *nodes.ReplicaReconciler
JobStore *jobs.JobStore
Dispatcher *jobs.Dispatcher
AgentStore *agents.AgentStore
AgentBridge *agents.EventBridge
DistStores *distributed.Stores
FileMgr *storage.FileManager
FileStager nodes.FileStager
ModelAdapter *nodes.ModelRouterAdapter
Unloader *nodes.RemoteUnloaderAdapter
ModelCleanup *nodes.ModelCleanupService
// WorkerHTTPDial reaches a worker's own HTTP server for the admin
// backend-logs proxy, the same way the HTTP file stager does.
WorkerHTTPDial nodes.WorkerNetDialerFor
shutdownOnce sync.Once
}
// Shutdown stops all distributed services in reverse initialization order.
// It is safe to call on a nil receiver and is idempotent (uses sync.Once).
func (ds *DistributedServices) Shutdown() {
if ds == nil {
return
}
ds.shutdownOnce.Do(func() {
if ds.Health != nil {
ds.Health.Stop()
}
if ds.Dispatcher != nil {
ds.Dispatcher.Stop()
}
if closer, ok := ds.Store.(io.Closer); ok {
closer.Close()
}
// AgentBridge has no Close method — its NATS subscriptions are cleaned up
// when the NATS client is closed below.
if ds.Nats != nil {
ds.Nats.Close()
}
xlog.Info("Distributed services shut down")
})
}
// initDistributed validates distributed mode prerequisites and initializes
// NATS, object storage, node registry, and instance identity.
// Returns nil if distributed mode is not enabled.
// configLoader is used by the SmartRouter to compute concurrency-group
// anti-affinity at placement time (#9659); it may be nil in tests.
// pinned, when set, replaces configLoader as the source of models the router
// and reconciler must keep loaded (it adds warm failover targets).
func initDistributed(cfg *config.ApplicationConfig, authDB *gorm.DB, configLoader *config.ModelConfigLoader, pinned nodes.PinnedModelResolver) (*DistributedServices, error) {
if !cfg.Distributed.Enabled {
return nil, nil
}
xlog.Info("Distributed mode enabled — validating prerequisites")
// Validate distributed config (NATS URL, S3 credential pairing, durations, etc.)
if err := cfg.Distributed.Validate(); err != nil {
return nil, err
}
// Validate PostgreSQL is configured (auth DB must be PostgreSQL for distributed mode)
if !cfg.Auth.Enabled {
return nil, fmt.Errorf("distributed mode requires authentication to be enabled (--auth / LOCALAI_AUTH=true)")
}
if !isPostgresURL(cfg.Auth.DatabaseURL) {
return nil, fmt.Errorf("distributed mode requires PostgreSQL for auth database (got %q)", sanitize.URL(cfg.Auth.DatabaseURL))
}
// Generate instance ID if not set
if cfg.Distributed.InstanceID == "" {
cfg.Distributed.InstanceID = uuid.New().String()
}
xlog.Info("Distributed instance", "id", cfg.Distributed.InstanceID)
// Connect to NATS
natsAuth := cfg.Distributed.NatsAuthConfig()
if natsAuth.RequireAuth && (natsAuth.ServiceUserJWT == "" || natsAuth.ServiceUserSeed == "") {
return nil, fmt.Errorf("LOCALAI_NATS_REQUIRE_AUTH requires LOCALAI_NATS_SERVICE_JWT and LOCALAI_NATS_SERVICE_SEED")
}
natsOpts := cfg.Distributed.NatsMessagingOptions("", "")
natsClient, err := messaging.New(cfg.Distributed.NatsURL, natsOpts...)
if err != nil {
return nil, fmt.Errorf("connecting to NATS: %w", err)
}
xlog.Info("Connected to NATS", "url", sanitize.URL(cfg.Distributed.NatsURL))
// Ensure NATS is closed if any subsequent initialization step fails.
success := false
defer func() {
if !success {
natsClient.Close()
}
}()
// Initialize object storage
var store storage.ObjectStore
if cfg.Distributed.StorageURL != "" {
if cfg.Distributed.StorageBucket == "" {
return nil, fmt.Errorf("distributed storage bucket must be set when storage URL is configured")
}
s3Store, err := storage.NewS3Store(context.Background(), storage.S3Config{
Endpoint: cfg.Distributed.StorageURL,
Region: cfg.Distributed.StorageRegion,
Bucket: cfg.Distributed.StorageBucket,
AccessKeyID: cfg.Distributed.StorageAccessKey,
SecretAccessKey: cfg.Distributed.StorageSecretKey,
ForcePathStyle: true, // required for MinIO
})
if err != nil {
return nil, fmt.Errorf("initializing S3 storage: %w", err)
}
xlog.Info("Object storage initialized (S3)", "endpoint", cfg.Distributed.StorageURL, "bucket", cfg.Distributed.StorageBucket)
store = s3Store
} else {
// Fallback to filesystem storage in distributed mode (useful for single-node testing)
fsStore, err := storage.NewFilesystemStore(cfg.DataPath + "/objectstore")
if err != nil {
return nil, fmt.Errorf("initializing filesystem storage: %w", err)
}
xlog.Info("Object storage initialized (filesystem fallback)", "path", cfg.DataPath+"/objectstore")
store = fsStore
}
// Initialize node registry (requires the auth DB which is PostgreSQL)
if authDB == nil {
return nil, fmt.Errorf("distributed mode requires auth database to be initialized first")
}
registry, err := nodes.NewNodeRegistry(authDB)
if err != nil {
return nil, fmt.Errorf("initializing node registry: %w", err)
}
xlog.Info("Node registry initialized")
// Bound durable heartbeat writes: a beat that only carries a fresher
// timestamp is what turned backend_nodes into a 460 MB six-row table.
registry.SetHeartbeatCheckpoint(cfg.Distributed.NodeHeartbeatCheckpointOrDefault())
// Measure the vacuum horizon. The 42 days it stayed open went unnoticed
// because no gauge reported it until models started failing to load.
if err := monitoring.RegisterControlPlaneDBMetrics(authDB, 30*time.Second); err != nil {
// Metrics are diagnostic; a failure here must not stop the frontend.
xlog.Warn("Control-plane database metrics unavailable", "error", err)
}
// Let scheduling rules be keyed by a model alias. The registry resolves a
// rule's name through the config loader to find the model it governs, so an
// operator can pin placement to a stable name like "production" and have it
// follow the alias when the alias is repointed. Wired before the seed below
// and before the reconciler starts, so the first tick already resolves.
if configLoader != nil {
registry.SetAliasResolver(configLoader)
}
// Seed declarative per-model scheduling config (LOCALAI_MODEL_SCHEDULING /
// LOCALAI_MODEL_SCHEDULING_CONFIG). Authoritative: overwrites matching models
// on every boot. Runs before the reconciler starts so the first tick already
// sees the desired state. Models not listed are left untouched.
if cfg.Distributed.ModelSchedulingJSON != "" || cfg.Distributed.ModelSchedulingConfigPath != "" {
schedConfigs, err := nodes.ParseSchedulingSeed(cfg.Distributed.ModelSchedulingJSON, cfg.Distributed.ModelSchedulingConfigPath)
if err != nil {
return nil, fmt.Errorf("parsing declarative model scheduling config: %w", err)
}
if err := registry.SeedModelScheduling(context.Background(), schedConfigs); err != nil {
return nil, fmt.Errorf("seeding declarative model scheduling config: %w", err)
}
xlog.Info("Applied declarative model scheduling config", "models", len(schedConfigs))
}
// Collect SmartRouter option values; the router itself is created after all
// dependencies (including FileStager and Unloader) are ready.
var routerAuthToken string
if cfg.Distributed.RegistrationToken != "" {
routerAuthToken = cfg.Distributed.RegistrationToken
}
var routerGalleriesJSON string
if galleriesJSON, err := json.Marshal(cfg.BackendGalleries); err == nil {
routerGalleriesJSON = string(galleriesJSON)
}
healthMon := nodes.NewHealthMonitor(registry, authDB,
cfg.Distributed.HealthCheckIntervalOrDefault(),
cfg.Distributed.StaleNodeThresholdOrDefault(),
routerAuthToken,
!cfg.Distributed.DisablePerModelHealthCheck,
)
// Initialize job store
jobStore, err := jobs.NewJobStore(authDB)
if err != nil {
return nil, fmt.Errorf("initializing job store: %w", err)
}
xlog.Info("Distributed job store initialized")
workQueue := messaging.NewNATSWorkQueue(natsClient)
// Initialize job dispatcher
dispatcher := jobs.NewDispatcher(jobStore, workQueue, natsClient, authDB, cfg.Distributed.InstanceID)
// Initialize agent store
agentStore, err := agents.NewAgentStore(authDB)
if err != nil {
return nil, fmt.Errorf("initializing agent store: %w", err)
}
xlog.Info("Distributed agent store initialized")
// Initialize agent event bridge
agentBridge := agents.NewEventBridge(natsClient, agentStore, cfg.Distributed.InstanceID)
// Start observable persister — captures observable_update events from workers
// (which have no DB access) and persists them to PostgreSQL.
if err := agentBridge.StartObservablePersister(); err != nil {
xlog.Warn("Failed to start observable persister", "error", err)
} else {
xlog.Info("Observable persister started")
}
// Initialize Phase 4 stores (MCP, Gallery, FineTune, Skills)
distStores, err := distributed.InitStores(authDB)
if err != nil {
return nil, fmt.Errorf("initializing distributed stores: %w", err)
}
// Initialize file manager with local cache
cacheDir := cfg.DataPath + "/cache"
fileMgr, err := storage.NewFileManager(store, cacheDir)
if err != nil {
return nil, fmt.Errorf("initializing file manager: %w", err)
}
xlog.Info("File manager initialized", "cacheDir", cacheDir)
// Create FileStager for distributed file transfer
workerHTTPDial := nodes.DirectWorkerNetDialer()
var fileStager nodes.FileStager
if cfg.Distributed.StorageURL != "" {
fileStager = nodes.NewS3NATSFileStager(fileMgr, natsClient)
xlog.Info("File stager initialized (S3+NATS)")
} else {
fileStager = nodes.NewHTTPFileStager(func(nodeID string) (string, error) {
node, err := registry.Get(context.Background(), nodeID)
if err != nil {
return "", err
}
if node.HTTPAddress == "" {
return "", fmt.Errorf("node %s has no HTTP address for file transfer", nodeID)
}
return node.HTTPAddress, nil
}, cfg.Distributed.RegistrationToken, workerHTTPDial)
xlog.Info("File stager initialized (HTTP direct transfer)")
}
// Create RemoteUnloaderAdapter — needed by SmartRouter and startup.go
remoteUnloader := nodes.NewRemoteUnloaderAdapter(
registry,
natsClient,
cfg.Distributed.BackendInstallTimeoutOrDefault(),
cfg.Distributed.BackendUpgradeTimeoutOrDefault(),
)
// Prefix-cache-aware routing. Enabled by default; an operator can opt out
// with --distributed-prefix-cache=false, which leaves prefixProvider and
// pressure nil so the SmartRouter and reconciler behave exactly as the
// round-robin floor (true no-op). When enabled we build the local index,
// wrap it in a NATS-backed Sync (publishes our observations, applies peers'
// via the subscriptions below), install the extraction hook used by
// core/backend/llm.go, and run a background eviction ticker on the app ctx.
var prefixProvider prefixcache.Provider
var pressure *prefixcache.Pressure
var prefixCfg prefixcache.Config
if !cfg.Distributed.PrefixCacheDisabled {
prefixCfg = prefixcache.DefaultConfig()
if cfg.Distributed.PrefixCacheTTL > 0 {
prefixCfg.TTL = cfg.Distributed.PrefixCacheTTL
}
if err := prefixCfg.Validate(); err != nil {
return nil, fmt.Errorf("invalid prefix-cache configuration: %w", err)
}
idx := prefixcache.NewIndex(prefixCfg)
prefixSync := prefixcache.NewSync(idx, natsClient)
pressure = prefixcache.NewSyncedPressure(prefixCfg.PressureWindow, natsClient)
prefixProvider = prefixSync
// Invalidate the prefix-cache index whenever a replica row is removed.
// AddReplicaRemovedHook fires from the single chokepoint all removal paths
// funnel through (RemoveNodeModel / RemoveAllNodeModelReplicas), so this
// one hook covers every path: reconciler scale-down, probe reaper,
// health-monitor reap, RemoteUnloaderAdapter, and the router. Registering
// it only inside this enabled block keeps the disabled path a true no-op
// for the prefix cache; other subsystems register their own hooks
// independently and are unaffected either way.
registry.AddReplicaRemovedHook(func(model, node string, replica int) {
if replica < 0 {
prefixSync.InvalidateNode(model, node)
} else {
prefixSync.Invalidate(model, prefixcache.ReplicaKey{NodeID: node, Replica: replica})
}
})
distributedhdr.PrefixChainHook = func(model, prompt string) []uint64 {
return prefixcache.ExtractChain(model, prompt, prefixCfg)
}
// Apply peers' observations/invalidations to the same Sync. ApplyObserve
// and ApplyInvalidate update only the local index and do not re-publish,
// so there is no broadcast loop.
if _, err := messaging.SubscribeJSON(natsClient, messaging.SubjectPrefixCacheObserve, func(ev messaging.PrefixCacheObserveEvent) {
prefixSync.ApplyObserve(ev, time.Now())
}); err != nil {
return nil, fmt.Errorf("subscribing to %s: %w", messaging.SubjectPrefixCacheObserve, err)
}
if _, err := messaging.SubscribeJSON(natsClient, messaging.SubjectPrefixCacheInvalidate, func(ev messaging.PrefixCacheInvalidateEvent) {
prefixSync.ApplyInvalidate(ev)
}); err != nil {
return nil, fmt.Errorf("subscribing to %s: %w", messaging.SubjectPrefixCacheInvalidate, err)
}
if _, err := messaging.SubscribeJSON(natsClient, messaging.SubjectPrefixCachePressure, func(ev messaging.PrefixCachePressureEvent) {
pressure.ApplyPressure(ev, time.Now())
}); err != nil {
return nil, fmt.Errorf("subscribing to %s: %w", messaging.SubjectPrefixCachePressure, err)
}
// Keep an exact-residency index current so backend producers can report
// their real KV state without coupling to router internals. Routing stays
// on the guessed provider until a backend producer is available.
reportedIndex := prefixcache.NewReportedIndex()
if _, err := messaging.SubscribeJSON(natsClient, messaging.SubjectPrefixCacheResidency, reportedIndex.Apply); err != nil {
return nil, fmt.Errorf("subscribing to %s: %w", messaging.SubjectPrefixCacheResidency, err)
}
// Background eviction: sweep idle entries on the app context. Stopped
// when the app context is cancelled (mirrors the reconciler loop which
// also runs on options.Context). TTL/2 keeps stale entries from
// outliving their idle window by more than half a TTL.
evictInterval := prefixCfg.TTL / 2
go func() {
ticker := time.NewTicker(evictInterval)
defer ticker.Stop()
for {
select {
case <-cfg.Context.Done():
return
case <-ticker.C:
prefixSync.Evict(time.Now())
}
}
}()
xlog.Info("Prefix-cache-aware routing enabled", "ttl", prefixCfg.TTL, "evictInterval", evictInterval)
} else {
xlog.Info("Prefix-cache-aware routing disabled: using round-robin routing")
}
// All dependencies ready — build SmartRouter with all options at once
var conflictResolver nodes.ConcurrencyConflictResolver
var pinnedResolver nodes.PinnedModelResolver
var modelFiles func(string) []string
if configLoader != nil {
conflictResolver = configLoader
pinnedResolver = configLoader
if cfg.SystemState != nil {
modelFiles = declaredModelFiles(configLoader, cfg.SystemState.Model.ModelsPath)
}
}
if pinned != nil {
pinnedResolver = pinned
}
modelCleanup := nodes.NewModelCleanupService(registry, remoteUnloader)
router := nodes.NewSmartRouter(registry, nodes.SmartRouterOptions{
Unloader: remoteUnloader,
ModelCleanup: modelCleanup,
FileStager: fileStager,
GalleriesJSON: routerGalleriesJSON,
AuthToken: routerAuthToken,
DB: authDB,
ConflictResolver: conflictResolver,
PinnedResolver: pinnedResolver,
ModelFiles: modelFiles,
PrefixProvider: prefixProvider,
PrefixConfig: prefixCfg,
Pressure: pressure,
SharedModels: cfg.Distributed.SharedModels,
// A closure over the live ApplicationConfig, NOT a snapshot: the
// runtime setting (distributed_disk_headroom_check) mutates this exact
// member, so a snapshot here would make the toggle a no-op until
// restart. env/CLI sets the boot value, POST /api/settings overrides it
// live, and this is the single member both write.
DiskHeadroomEnabled: func() bool { return !cfg.Distributed.DiskHeadroomDisabled },
// RAW, not OrDefault: zero means "derive the budget per model from the
// checkpoint size" (config.ModelLoadTimeoutForSize), which is what makes
// a 70 GB video checkpoint work without the operator first hitting a
// DeadlineExceeded and going looking for a knob. A non-zero value here is
// an explicit override and is used verbatim.
ModelLoadTimeout: cfg.Distributed.ModelLoadTimeout,
// Cap how long a cold load may hold the per-model advisory lock. Derived
// from BOTH configured budgets it has to cover, so raising either the
// install timeout (slow links pulling multi-GB images) or the model load
// timeout (very large checkpoints) widens the ceiling too, instead of
// letting a stale bound cut a legitimately slow load short.
ModelLoadCeiling: nodes.ModelLoadCeilingFor(
cfg.Distributed.BackendInstallTimeoutOrDefault(),
cfg.Distributed.ModelLoadTimeoutOrDefault(),
),
// Bounds the REQUEST, not the load: a caller out of budget gets 503 with
// live staging progress while the job keeps running underneath.
ModelLoadWait: cfg.Distributed.ModelLoadWait,
})
// Wire staging-progress broadcasting so file-staging shows up on every
// replica, not just the one performing the transfer. Without this, a
// /api/operations poll that round-robins onto a peer sees no staging row and
// the progress flickers. The origin publishes; peers mirror via the wildcard.
// A silently disabled safety check is how the original incident stayed
// invisible for sixteen minutes. Say so once, loudly, at startup.
if cfg.Distributed.DiskHeadroomDisabled {
xlog.Info("Disk-headroom admission check is DISABLED: node selection will ignore whether a worker can store the model, and staging may fail with ENOSPC partway through a transfer",
"knob", config.FlagDiskHeadroomCheck, "env", "LOCALAI_DISTRIBUTED_DISK_HEADROOM_CHECK")
}
router.StagingTracker().SetPublisher(natsClient)
if _, err := router.StagingTracker().SubscribeBroadcasts(natsClient); err != nil {
xlog.Warn("Failed to subscribe to staging progress broadcasts", "error", err)
}
// Create ReplicaReconciler for auto-scaling model replicas. Adapter +
// RegistrationToken feed the state-reconciliation passes: pending op
// drain uses the adapter, and model health probes use the token to auth
// against workers' gRPC HealthCheck.
reconciler := nodes.NewReplicaReconciler(nodes.ReplicaReconcilerOptions{
Registry: registry,
Scheduler: router,
Unloader: remoteUnloader,
Adapter: remoteUnloader,
RegistrationToken: cfg.Distributed.RegistrationToken,
DB: authDB,
Interval: 30 * time.Second,
ScaleDownDelay: 5 * time.Minute,
ProbeStaleAfter: 2 * time.Minute,
Pressure: pressure,
PressureThreshold: prefixCfg.PressureScaleThreshold,
PinnedResolver: pinnedResolver,
})
// Create ModelRouterAdapter to wire into ModelLoader
modelAdapter := nodes.NewModelRouterAdapter(router)
success = true
return &DistributedServices{
Nats: natsClient,
WorkQueue: workQueue,
AgentControl: nodes.NewNATSAgentControl(natsClient),
Store: store,
Registry: registry,
Router: router,
Health: healthMon,
Reconciler: reconciler,
JobStore: jobStore,
Dispatcher: dispatcher,
AgentStore: agentStore,
AgentBridge: agentBridge,
DistStores: distStores,
FileMgr: fileMgr,
FileStager: fileStager,
ModelAdapter: modelAdapter,
Unloader: remoteUnloader,
ModelCleanup: modelCleanup,
WorkerHTTPDial: workerHTTPDial,
}, nil
}
func isPostgresURL(url string) bool {
return strings.HasPrefix(url, "postgres://") || strings.HasPrefix(url, "postgresql://")
}