package application import ( "context" "encoding/json" "fmt" "io" "net" "strconv" "strings" "sync" "time" "github.com/google/uuid" "github.com/mudler/LocalAI/core/config" "github.com/mudler/LocalAI/core/services/agents" "github.com/mudler/LocalAI/core/services/cluster" "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/nodes" "github.com/mudler/LocalAI/core/services/nodes/prefixcache" "github.com/mudler/LocalAI/core/services/pgbus" "github.com/mudler/LocalAI/core/services/storage" "github.com/mudler/LocalAI/internal" "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 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 // Bus is the deployment's fan-out carrier, riding the auth database's // PostgreSQL rather than a message broker. Nothing publishes on it and // nothing subscribes yet; it is built at boot because its DSN has exactly // one legitimate source and that has to be settled once, here, rather than // invented by whichever call site is migrated onto it first. Bus *pgbus.Bus // Cluster is the replica-membership registry: which frontend replicas are // alive, at which address, and which of them holds a given worker's tunnel. Cluster *cluster.Registry // Membership publishes this replica's row and reaps the dead. Nil when no // peer-reachable address could be determined, which leaves this replica // invisible to its peers but otherwise fully functional. Membership *cluster.Membership // PeerSessions owns the peer links other replicas dialled into this one, // and relays the streams that arrive on them onto the worker tunnels this // replica holds. PeerSessions *cluster.SessionStore // Peers owns the peer links this replica dialled OUT, the mirror of // PeerSessions. It is what the relaying dialer opens a stream on when a // request arrives here for a worker another replica holds. Peers *cluster.PeerPool // Tunnels holds the worker tunnels this replica has accepted and keeps the // node_connections table agreeing with them. It is handed to the membership // loop, which re-claims what it holds after this replica has been reaped, // and to the route that accepts a worker's dial. Tunnels *cluster.TunnelRegistry // WorkerDialer is how anything in this process reaches a worker: locally // when this replica holds the tunnel, and through the owning replica when // it does not. The HTTP layer takes its WebSocket log proxy from here. WorkerDialer *cluster.WorkerDialer // BackendClients builds the gRPC clients for worker backend processes, over // WorkerDialer. Exposed so the model store built in startup.go reaches // remote models the same way every other caller does. BackendClients nodes.BackendClientFactory // AgentControl carries the frontend's MCP verbs to whichever agent worker // holds a tunnel this deployment can reach. It is what the chat, responses, // messages and MCP endpoints reach an agent worker through; a nil one means // this frontend cannot run MCP at all, which is why initDistributed refuses // to come up without it rather than leaving the endpoints to discover it // one request at a time. AgentControl *nodes.AgentControlClient 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() { // Peer state first: a replica that is going away should stop claiming // to be alive before it stops answering, so peers re-home rather than // dial a process in teardown. if ds.Membership != nil { ds.Membership.Stop() } if ds.PeerSessions != nil { ds.PeerSessions.CloseAll() } // Both halves of the peer mesh go down together. A pool left open // holds a WebSocket and two yamux loop goroutines per peer for as long // as the process lives, and an Open after this reports ErrPoolClosed, // which is a fact about this process and never node absence. if ds.Peers != nil { ds.Peers.Close() } 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() } // The broadcast carrier holds a PostgreSQL session pinned for the life // of the process, plus the goroutine parked on it. A replica that // leaves one behind on every restart runs the server out of // connections, and the symptom lands on whatever connects next. if ds.Bus != nil { ds.Bus.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. func initDistributed(cfg *config.ApplicationConfig, authDB *gorm.DB, configLoader *config.ModelConfigLoader) (*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") } // The fan-out carrier, opened before anything that might want it. It is // built here and not by its first adopter because its DSN has one // legitimate source, and a setting that decides whether every broadcast in // the deployment is delivered should not be settled under the time pressure // of a migration. bus, err := newBroadcastBus(cfg.Context, cfg, authDB) if err != nil { return nil, err } defer func() { if !success { bus.Close() } }() registry, err := nodes.NewNodeRegistry(authDB) if err != nil { return nil, fmt.Errorf("initializing node registry: %w", err) } xlog.Info("Node registry initialized") // Replica membership. NewNodeRegistry has just migrated the tables this // reads, so it has to come after it. clusterRegistry := cluster.NewRegistry(authDB) var membership *cluster.Membership if advertised, err := advertisedPeerAddr(cfg); err != nil { // Not fatal, and the cost is worth stating exactly rather than as // "peers cannot reach it", because it is larger than that now. // // Without a row in the instances table this replica is not a live // owner as far as Registry.Owner is concerned: that read joins a // connection against a live instance, so a worker whose tunnel lands // HERE is answered as unroutable at every OTHER replica, for as long // as it stays here. This replica serves that worker perfectly well // itself; nobody else can. On N replicas behind round robin that is // (N-1)/N of the traffic for that worker. // // It does not refuse to START. Refusing would take out every existing // single-host deployment, whose route to a local database is loopback // and which has no peers to be unreachable by; the deployments this // hurts are multi-replica ones, and telling those two apart at startup // is a change with its own design and its own specs rather than a line // here. // // What it does not get to do is stay quiet. One startup line scrolls // away in seconds and the cost is paid for the whole life of the // process, on a symptom (workers that 5xx from most of the fleet) whose // obvious reading is "the worker is broken". So this is an ERROR, not a // warning, and nagUnadvertisedReplica below repeats it for as long as // the state lasts, naming the workers it is currently costing. xlog.Error("This replica is not registered in the cluster: no advertised address. Peers cannot reach it, and any worker whose tunnel lands here will be unroutable from every other replica", "error", err, "knob", "LOCALAI_DISTRIBUTED_ADVERTISE_ADDR") } else { membership = cluster.NewMembership(clusterRegistry, cfg.Distributed.InstanceID, advertised, internal.PrintableVersion()) // Before Start, so the first sweep already purges on the retention this // deployment's grace requires rather than on the floor. membership.SetReconnectGrace(cfg.Distributed.ReconnectGraceOrDefault()) if err := membership.Start(cfg.Context); err != nil { return nil, fmt.Errorf("registering this replica in the cluster: %w", err) } } // The worker tunnels this replica accepts. It claims as the SAME instance // ID membership registers under, because that is the ID a peer's Owner // lookup joins a claim against to decide the owner is alive; two IDs here // would make every claim this replica writes look like it belongs to a // replica that does not exist. tunnels := cluster.NewTunnelRegistry(clusterRegistry, cfg.Distributed.InstanceID) // Without this the re-claim in the heartbeat loop is dead code: a replica // stalled long enough to be swept loses the connection rows it owned, and // nothing would ever write them back, so every other replica would answer // "not connected" for workers that are connected right here. // // Nil when no peer-reachable address could be determined above. There is no // heartbeat loop to hand it to in that case, and no other replica can reach // this one anyway; the registry is still built, because it is what the // tunnel endpoint attaches to and what this replica opens its own streams // through. if membership != nil { membership.SetTunnels(tunnels) } else { // The runtime symptom the startup line cannot be. See // nagUnadvertisedReplica. go nagUnadvertisedReplica(cfg.Context, tunnels.Held, unadvertisedNagInterval, logUnroutableWorkers) } // The links peers dial IN, with the relay installed on them. This is what // makes more than one replica work: a worker holds one tunnel, it lands on // one replica, and every request that arrives anywhere else reaches the // worker through this handler. Passing nil here would leave every such // request refused, promptly and only at debug level, which presents as a // worker that is connected and unusable from most of the deployment. peerSessions := cluster.NewSessionStore(cluster.NewRelay(tunnels).Stream) // The links this replica dials OUT, the other half of the same mesh. It // authenticates with the registration token because that is the token the // peer route checks (see RegisterClusterRoutes); two different tokens here // would make every peer dial 401 with nothing naming the mismatch. peers := cluster.NewPeerPool(cfg.Distributed.InstanceID, cfg.Distributed.RegistrationToken, clusterRegistry) // The one door to every worker. Nothing in the frontend may dial a worker's // advertised address any more: a worker holds ONE tunnel, it lands on ONE // replica, and this resolves which replica that is and relays through it // when it is not this one. The three transports the frontend speaks to a // worker (gRPC to backend processes, HTTP for file staging and logs, a // WebSocket for live log streaming) are all pointed at it below. workerDialer := cluster.NewWorkerDialer(tunnels, peers) backendClients, err := nodes.NewTunnelClientFactory(cfg.Distributed.RegistrationToken, workerDialer.GRPCDialerFor) if err != nil { return nil, fmt.Errorf("wiring the worker backend client factory: %w", err) } // Bound to the http tag: the worker ignores the target for it and routes to // its own file-transfer and log server, wherever that bound. workerHTTPDialer := nodes.WorkerNetDialerFor(func(nodeID string) func(ctx context.Context, network, addr string) (net.Conn, error) { return workerDialer.DialerFor(nodeID, cluster.StreamTagHTTP) }) // 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) } // The health monitor is the SECOND reader of absence, and it reads it from // the same place and against the same window as the scheduler: a heartbeat // says the worker's supervisor is alive, presence says whether anything // here can still reach its backends, and a worker can be the first without // being the second indefinitely. healthMon := nodes.NewHealthMonitor(registry, authDB, cfg.Distributed.HealthCheckIntervalOrDefault(), cfg.Distributed.StaleNodeThresholdOrDefault(), routerAuthToken, !cfg.Distributed.DisablePerModelHealthCheck, clusterRegistry, cfg.Distributed.ReconnectGraceOrDefault(), backendClients, ) // 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") // Initialize job dispatcher dispatcher := jobs.NewDispatcher(jobStore, natsClient, authDB, cfg.Distributed.InstanceID, cfg.Distributed.JobWorkerConcurrency) // 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) // The frontend's control plane client. It reaches every worker over that // worker's own tunnel, on the same `http` stream tag the file stager below // uses, so a control RPC to a worker another replica holds is relayed the // way an inference request is. // // ONE of these for the whole frontend, and the S3 file stager takes this // one rather than minting a second. The client caches an http.Client per // node, which is what keeps a worker's tunnel stream warm between verbs; a // second client would open its own and the two would never share one. controlClient := nodes.NewControlClient(workerHTTPDialer, cfg.Distributed.RegistrationToken) // The caller the agent worker's control plane has been waiting for. MCP // execution and discovery used to be a NATS request onto a queue group, // where the bus chose the worker and neither side could say which one had // answered; they are now a query against the connection rows plus an // ordinary control RPC over the chosen worker's tunnel. agentControl, err := newAgentControl(cfg.Distributed, registry, clusterRegistry, controlClient) if err != nil { return nil, fmt.Errorf("wiring the agent control client: %w", err) } // Create FileStager for distributed file transfer var fileStager nodes.FileStager if cfg.Distributed.StorageURL != "" { fileStager = nodes.NewS3FileStager(fileMgr, controlClient) xlog.Info("File stager initialized (object store + worker tunnel)") } else { fileStager = nodes.NewHTTPFileStager(func(nodeID string) (string, error) { node, err := registry.Get(context.Background(), nodeID) if err != nil { return "", err } // An empty HTTPAddress is no longer a refusal. A tunnel-only worker // reports none and does not need one: the http stream tag ignores // the target and the worker routes to its own server. The host is // only ever the URL's host component here, and WorkerHTTPHost // supplies one that resolves nowhere so it cannot become a dial. return nodes.WorkerHTTPHost(nodeID, node.HTTPAddress), nil }, cfg.Distributed.RegistrationToken, workerHTTPDialer) xlog.Info("File stager initialized (HTTP direct transfer)") } // Create RemoteUnloaderAdapter — needed by SmartRouter and startup.go remoteUnloader := nodes.NewRemoteUnloaderAdapter( registry, controlClient, 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.NewPressure(prefixCfg.PressureWindow) 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) } // 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 if configLoader != nil { conflictResolver = configLoader } modelCleanup := nodes.NewModelCleanupService(registry, remoteUnloader) // Absence is stamped on by distributedSchedulerOptions rather than written // here. It is the only source of absence the scheduler has -- a fact read // from the database, so every replica answers it identically, where the bus // sentinel it replaces was one frontend's observation that nobody answered // IT within a budget -- and a field carrying that in a literal this size is // the easiest thing in this file to lose without a symptom. router := nodes.NewSmartRouter(registry, distributedSchedulerOptions(cfg.Distributed, clusterRegistry, nodes.SmartRouterOptions{ Unloader: remoteUnloader, ModelCleanup: modelCleanup, FileStager: fileStager, GalleriesJSON: routerGalleriesJSON, AuthToken: routerAuthToken, ClientFactory: backendClients, DB: authDB, ConflictResolver: conflictResolver, 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, ClientFactory: backendClients, DB: authDB, Interval: 30 * time.Second, ScaleDownDelay: 5 * time.Minute, ProbeStaleAfter: 2 * time.Minute, Pressure: pressure, PressureThreshold: prefixCfg.PressureScaleThreshold, }) // Both readers of absence, checked once, here. See requireAbsenceWiring for // why a missing assignment has no other symptom. if err := requireAbsenceWiring(router, healthMon); err != nil { return nil, err } // Create ModelRouterAdapter to wire into ModelLoader modelAdapter := nodes.NewModelRouterAdapter(router) success = true ds := &DistributedServices{ Nats: 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, Cluster: clusterRegistry, Membership: membership, PeerSessions: peerSessions, Peers: peers, Tunnels: tunnels, WorkerDialer: workerDialer, BackendClients: backendClients, AgentControl: agentControl, Bus: bus, } // Checked once, here, on the assembled struct. See requireBroadcastCarrier. if err := requireBroadcastCarrier(ds); err != nil { return nil, err } return ds, nil } // requireBroadcastCarrier refuses to hand back a distributed deployment whose // broadcast carrier is missing. // // The carrier reaches the deployment over two lines: the newBroadcastBus call // in initDistributed, and the Bus field in the twenty-three field literal // above. Deleting either one compiles and leaves every suite in this repository // green, and the two failures are different. Without the construction, nothing // can ever be published between replicas. Without the assignment the carrier is // opened and connected but Shutdown cannot see it, so every restart leaves a // pinned PostgreSQL session and its goroutines behind until the server runs out // of connections, and the operator sees the failure land on whatever connects // next rather than on LocalAI. // // Neither line can be reddened by a spec today: initDistributed opens NATS // before it reaches any of this, so it cannot be called from a unit test, and a // pointer field left out of a struct literal is not a compile error. What this // converts both omissions into is a deployment that refuses to start and names // what is missing, which is as far as they can be pinned until initDistributed // is testable. The guard itself is spec'd. func requireBroadcastCarrier(ds *DistributedServices) error { if ds == nil || ds.Bus == nil { return fmt.Errorf("distributed mode was initialized without a broadcast carrier: nothing could be published between replicas, and the PostgreSQL session it pins could not be closed on shutdown") } return nil } // newBroadcastBus opens the deployment's fan-out carrier on the auth database. // // The DSN is cfg.Auth.DatabaseURL and it may never be anything else. A second // source, a flag of its own or a value read from the environment, would let the // pinned LISTEN connection and the connection pool address two different // databases; that carrier publishes successfully, delivers nothing, on every // replica, and reports no error anywhere. isPostgresURL above has already // refused a value this carrier could not use. // // It is a function rather than four lines inside initDistributed so that the // equality can be pinned by a spec. initDistributed opens NATS before it // reaches this point and so cannot be called from a unit test, which would // leave the assignment as one line in a long function that compiles perfectly // well when it names the wrong field. func newBroadcastBus(ctx context.Context, cfg *config.ApplicationConfig, authDB *gorm.DB) (*pgbus.Bus, error) { if err := pgbus.Migrate(ctx, authDB); err != nil { return nil, fmt.Errorf("migrating the broadcast carrier: %w", err) } bus, err := pgbus.New(ctx, pgbus.Config{DSN: cfg.Auth.DatabaseURL, DB: authDB}) if err != nil { return nil, fmt.Errorf("opening the broadcast carrier: %w", err) } return bus, nil } // unadvertisedNagInterval is how often a replica that could not advertise // itself says so again. // // Five minutes is chosen against the log it lands in, not against the urgency: // the condition never clears on its own, so this line is either read once and // acted on or it is noise for the life of the process, and a noisy line gets // filtered rather than fixed. It is still frequent enough that the state is // visible in any window of logs an operator pulls while investigating the // symptom it causes. const unadvertisedNagInterval = 5 * time.Minute // nagUnadvertisedReplica repeats, for as long as the process runs, that this // replica is invisible to its peers, and names what that is currently costing. // // It exists because the deferral it accompanies changed cost between phases and // nothing about the deployment says so. Before workers held tunnels, a replica // with no advertised address was merely unreachable BY peers and could still // dial every worker directly, so a startup warning was proportionate. Now a // worker's tunnel lands on one replica and every other replica reaches it by // relaying to the owner, and the owner is resolved by joining the connection // row against a LIVE INSTANCES ROW - which this replica does not have. So every // worker that lands here is answered as unroutable everywhere else: on N // replicas behind round robin, (N-1)/N of that worker's traffic fails, while // this replica serves it perfectly and reports nothing. // // held is passed as a function rather than the registry so this can be driven // without one, and alarm is passed rather than logged inline so a spec can // observe the alarms instead of scraping a log. func nagUnadvertisedReplica(ctx context.Context, held func() []string, every time.Duration, alarm func([]string)) { ticker := time.NewTicker(every) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: alarm(held()) } } } // logUnroutableWorkers says what the state costs RIGHT NOW. // // The two cases are kept apart because they call for different urgency and an // operator can tell them apart at a glance. With no worker held this is a // misconfiguration that has not been paid for yet; with workers held, every one // of them is named, because "which worker is broken" is the question the // symptom sends an operator to ask and the answer is that none of them is. func logUnroutableWorkers(held []string) { if len(held) == 0 { xlog.Warn("This replica is still not registered in the cluster: no advertised address. No worker holds a tunnel here yet; the first that does will be unroutable from every other replica", "knob", "LOCALAI_DISTRIBUTED_ADVERTISE_ADDR") return } xlog.Error("This replica is not registered in the cluster and holds worker tunnels: those workers are unroutable from every OTHER replica, and requests for their models fail there with no route. The workers are healthy; this replica is invisible", "workers", held, "worker_count", len(held), "knob", "LOCALAI_DISTRIBUTED_ADVERTISE_ADDR") } // advertisedPeerAddr is the host:port peers dial to reach this replica. // // The operator's value wins outright. Otherwise it is derived from the port // this process serves on and the local address that routes to PostgreSQL, which // is only a peer-reachable answer when the database is on another host; // DiscoverAdvertisedAddr refuses rather than guessing when it is not. func advertisedPeerAddr(cfg *config.ApplicationConfig) (string, error) { if configured := cfg.Distributed.AdvertiseAddr; configured != "" { // A configured address skips discovery, so it also skips every check // discovery makes. Unusable is refused; merely questionable (a // loopback address, correct on one host and wrong on three) is said // once and honoured, because refusing it would refuse single-host // deployments that use it correctly. reason, err := cluster.CheckAdvertisedAddr(configured) if err != nil { return "", err } if reason != "" { xlog.Warn("Configured peer address is not one another host can dial", "address", configured, "reason", reason, "knob", "LOCALAI_DISTRIBUTED_ADVERTISE_ADDR") } return configured, nil } if cfg.APIAddress == "" { return "", fmt.Errorf("no API address to derive a peer port from") } _, port, err := net.SplitHostPort(cfg.APIAddress) if err != nil { return "", fmt.Errorf("reading the peer port out of API address %q: %w", cfg.APIAddress, err) } portNumber, err := strconv.Atoi(port) if err != nil { return "", fmt.Errorf("API address %q has a non-numeric port: %w", cfg.APIAddress, err) } return cluster.DiscoverAdvertisedAddr(cfg.Auth.DatabaseURL, portNumber) } func isPostgresURL(url string) bool { return strings.HasPrefix(url, "postgres://") || strings.HasPrefix(url, "postgresql://") }