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/monitoring" "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 { 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. Every cross-replica family the // frontend broadcasts is on it now, with the single exception named above. // Adopters reach it through Broadcast() rather than through this field, so // that "which carrier does this family travel on" is decided once instead // of at every adopter; its DSN likewise has exactly one legitimate source, // settled in newBroadcastBus rather than invented per call site. 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 // JobDispatch takes queued work off the job store and drives it on an agent // worker over that worker's tunnel. It is what replaces the three NATS // queue groups: dispatch is a claim, and a claim is a row and a lock. JobDispatch *jobs.DispatchLoop 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.JobDispatch != nil { ds.JobDispatch.Stop() } if ds.Dispatcher != nil { ds.Dispatcher.Stop() } if closer, ok := ds.Store.(io.Closer); ok { closer.Close() } // AgentBridge has no Close method and needs none: it holds no // process-lifetime subscription of its own beyond the observable // persister, whose carrier is closed below. // 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. // // galleryProgress is the gallery service, narrowed to the one method a node // departure needs. It is a PARAMETER and not a later setter because the // registration of every per-node cache a departure evicts happens here, in one // place, and a cache registered somewhere else is a cache a reader cannot find // by reading this function. func initDistributed(cfg *config.ApplicationConfig, authDB *gorm.DB, configLoader *config.ModelConfigLoader, galleryProgress nodeProgressDropper) (*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) // No message bus is dialled here, and there is none left to dial. The last // family that needed one was agent..cancel, whose subscriber is an // agent worker that has no database and so could not join the broadcast // carrier below; it is now a control RPC on the tunnel that worker holds. // A distributed deployment needs PostgreSQL and the frontends' own HTTP // listener, and nothing else. // // success guards the carriers opened below, which must not be left pinned // when a later initialization step fails. success := false // 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") // 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) } // Replica membership. NewNodeRegistry has just migrated the tables this // reads, so it has to come after it. clusterRegistry := cluster.NewRegistry(authDB) // This replica's proof of which replica it is, minted ONCE per process and // handed to both halves of the peer mesh: the membership loop publishes its // hash in the instances row, and the peer pool presents its plaintext on // every outbound dial. Two mints would leave a replica whose published hash // and presented secret disagree, which every peer refuses and which reads // from the logs like a peer running an older release. // // The plaintext never leaves this process except in a peer dial's header. // There is nothing to configure and nothing to rotate: a restart mints a // new one, and the same registration that republishes this replica's // address republishes the hash beside it. peerCredential := cluster.NewPeerCredential() 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. // // It now costs the OTHER direction too, which is why the line says so. // A peer link is authenticated by the dialling replica's own credential, // published in the instances row this replica never writes, so this // replica cannot dial a peer either: its own relayed requests are // refused as an unproven identity rather than merely arriving nowhere. xlog.Error("This replica is not registered in the cluster: no advertised address. Peers cannot reach it, any worker whose tunnel lands here will be unroutable from every other replica, and this replica cannot relay OUT either, because a peer link is authenticated by the credential an instance row publishes and this replica has no row", "error", err, "knob", "LOCALAI_DISTRIBUTED_ADVERTISE_ADDR") } else { membership = cluster.NewMembership(clusterRegistry, cfg.Distributed.InstanceID, advertised, internal.PrintableVersion(), peerCredential) // 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. // // And with this replica's own credential, which is the half that says WHICH // replica is dialling. It is the same value membership published the hash // of, by construction: there is one mint above and both call sites read it. peers := cluster.NewPeerPool(cfg.Distributed.InstanceID, cfg.Distributed.RegistrationToken, peerCredential, 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. // // The departure notifier is built HERE, before its only caller, and its // subscribers are registered further down once the caches they drop exist. // One object, one caller, so "what does a departure evict" is answered by // reading registerDepartureEvictions and nothing else. departures := nodes.NewDepartureNotifier() healthMon := nodes.NewHealthMonitor(registry, authDB, cfg.Distributed.HealthCheckIntervalOrDefault(), cfg.Distributed.StaleNodeThresholdOrDefault(), routerAuthToken, !cfg.Distributed.DisablePerModelHealthCheck, clusterRegistry, cfg.Distributed.ReconnectGraceOrDefault(), departures, 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 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") // 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) } // The job dispatcher and the agent event bridge, both on the broadcast // carrier. See newFanoutBridges for why the two constructors are reached // through one function that names *pgbus.Bus. // // The bridge takes agentControl and not a carrier: a cancel is the one // family whose far end is an agent worker, and it now rides that worker's // tunnel as a control RPC. It is built ABOVE for that reason, rather than // with the rest of the control plane below. dispatcher, agentBridge, rebroadcast, err := newFanoutBridges(bus, agentControl, jobStore, agentStore, authDB, cfg.Distributed.InstanceID) if err != nil { return nil, err } // 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 consumer side of the claim queue, built and started in one act: see // startJobDispatchLoop for why those are not two lines. jobDispatch, err := startJobDispatchLoop(cfg.Context, cfg.Distributed, authDB, jobStore, registry, clusterRegistry, controlClient, rebroadcast) if err != nil { return nil, fmt.Errorf("wiring the job dispatch loop: %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 Sync on the broadcast carrier (which both publishes our // observations and applies peers'), 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 // The CONCRETE Sync, declared out here so the departure wiring below can be // handed it. Concrete and not prefixProvider, because a nil interface value // carrying a nil *Sync is not nil, and the disabled deployment would then // register an eviction that dereferences it. var prefixDrop *prefixcache.Sync 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) // S4. One call puts this replica's observations and its peers' on the // same carrier, and it takes the CONCRETE carrier so that no other // thing satisfying messaging.Broadcaster can be handed to it by // accident. There is no second carrier in this scope to hand over any // more; the type stays narrow so there is still none on the day one is // added. See cache_fanout_wiring.go for the whole argument. prefixSync, err := wirePrefixCacheBroadcasts(bus, prefixCfg, idx) if err != nil { return nil, err } pressure = prefixcache.NewPressure(prefixCfg.PressureWindow) prefixProvider = prefixSync prefixDrop = 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) } // 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, on the same carrier. // 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") } // S3, and it is ONE call rather than a publisher and a subscriber: see // StagingTracker.SetBroadcaster for why a tracker that could name two // carriers is a progress bar that only the originating replica shows. if _, err := wireStagingBroadcasts(bus, router.StagingTracker()); err != nil { xlog.Warn("Failed to subscribe to staging progress broadcasts", "error", err) } // Every per-node cache a departure leaves stale, onto the one notification // point, after the router that owns two of them exists. if err := registerDepartureEvictions(departures, prefixDrop, router, galleryProgress, controlClient, fileStager); err != nil { return nil, 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{ 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, JobDispatch: jobDispatch, Bus: bus, } // Checked once, here, on the assembled struct. See requireBroadcastCarrier. if err := requireBroadcastCarrier(ds); err != nil { return nil, err } return ds, nil } // nodeProgressDropper is the gallery service narrowed to the one method a node // departure needs. An interface so the wiring below cannot reach for anything // else on the gallery service from inside an eviction hook. type nodeProgressDropper interface { DropNodeProgress(nodeID string) } // The names each per-node cache is registered under. They are constants because // they are asserted: a wiring spec that spelled them itself would keep passing // after a subscriber was renamed and lost. const ( departurePrefixCache = "prefix-cache" departureProbeCache = "probe-cache" departureStagingTracker = "staging-tracker" departureGalleryNodes = "gallery-node-progress" departureControlClients = "control-http-clients" departureStagerClients = "file-stager-http-clients" ) // registerDepartureEvictions registers every per-node cache that a node's // departure leaves stale on the deployment's one departure notification point. // // One function, and every subscriber in it, because the reason the notifier // exists is that a reader cannot otherwise enumerate what a demotion // invalidates: before it, one node type could depart and each stale cache was // dropped from wherever its owner happened to notice. Adding a per-node cache // without adding a line here is the failure this shape exists to make visible, // which is why the subscribers are NAMED and the names are asserted. // // It REFUSES rather than skipping when the router or the gallery service is // missing. A deployment whose departed nodes keep their probe entries, staging // rows and per-node operation progress does not fail, log or slow down: it // answers with state for a node that left, indefinitely. // // prefix may be nil, and only prefix. That is --distributed-prefix-cache=false, // where there is no index to drop from, and it stays a true no-op: nothing is // registered rather than a hook registered onto nothing. It is the CONCRETE // *prefixcache.Sync for that decision to be safe, since a nil provider inside // an interface would compare non-nil here and dereference on the first // departure. func registerDepartureEvictions(departures *nodes.DepartureNotifier, prefix *prefixcache.Sync, router *nodes.SmartRouter, gallery nodeProgressDropper, control *nodes.ControlClient, stager nodes.FileStager) error { if departures == nil { return fmt.Errorf("wiring departure evictions: no departure notifier, so a departed node would keep every per-node cache entry it has for the life of the process") } if router == nil { return fmt.Errorf("wiring departure evictions: no router, so a departed node would keep its probe-freshness entries and its staging operations") } if gallery == nil { return fmt.Errorf("wiring departure evictions: no gallery service, so a departed node would stay in every open operation's per-node breakdown") } if control == nil { return fmt.Errorf("wiring departure evictions: no control client, so a departed node would keep its cached HTTP client and that client's idle streams on a tunnel that is gone") } if stager == nil { return fmt.Errorf("wiring departure evictions: no file stager, so a departed node would keep the cached HTTP client its transfers ran on") } // S1. Inside a nil check and not inside the prefix-cache-enabled block, so // that "the disabled deployment registers nothing" is a fact a spec can // hold rather than a property of where a line was written. if prefix != nil { departures.OnDeparture(departurePrefixCache, func(node nodes.DepartedNode) { prefix.DropNode(node.ID) }) } // S2 and S3 are two registrations and not one, because they are two rules: // a probe entry is keyed by node ID and a staging op by node NAME, and a // single hook doing both would hide which of them was lost. departures.OnDeparture(departureProbeCache, func(node nodes.DepartedNode) { router.InvalidateNodeProbes(node.ID) }) departures.OnDeparture(departureStagingTracker, func(node nodes.DepartedNode) { router.StagingTracker().DropNode(node.Name) }) // S4. departures.OnDeparture(departureGalleryNodes, func(node nodes.DepartedNode) { gallery.DropNodeProgress(node.ID) }) // S5 and S6, the two per-node http.Client caches. Two registrations again, // because they are two caches with two owners: the control client's entry // is built on the first verb issued to a node and the stager's on the first // file staged to it, so a node can be in either without being in the other, // and one hook doing both would say only that some client was kept. // // Both are keyed by node ID and both are DROPPED rather than emptied. A // worker that comes back builds a fresh client on its next verb, over // whatever tunnel it has by then; keeping the old one would keep a // transport whose idle streams belong to a session that has ended. departures.OnDeparture(departureControlClients, func(node nodes.DepartedNode) { control.ForgetNode(node.ID) }) departures.OnDeparture(departureStagerClients, func(node nodes.DepartedNode) { stager.ForgetNode(node.ID) }) return 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 } // Broadcast is the ONE place a wiring site gets the deployment's fan-out // carrier, and it exists so that "this family travels on the broadcast carrier // and not on NATS" is decided once instead of at every adopter. // // It was five field reads before this: the fine-tune service, the quantization // service, the agent-task setter (twice, on two startup paths), the per-user // services manager and the Open Responses store. Every one of them takes a // messaging.Broadcaster, and the struct used to carry a second field that // satisfied it, so a site that reached for the wrong one compiled, started, // published and was delivered onto a carrier almost nothing read. That second // field is gone with the last family that needed a bus. Collapsing the choice // to one function is what keeps it a fact a spec can pin. // // The return is the interface and not *pgbus.Bus on purpose: handing a nil // *pgbus.Bus to an adopter would produce a non-nil interface wrapping a nil // pointer, and every adopter reads a nil carrier as "standalone, do not // broadcast". A typed nil would instead panic on the first Set. initDistributed // already refuses to return a deployment with no carrier (see // requireBroadcastCarrier), so the nil branch here is belt and braces for a // zero-valued struct in a test. func (ds *DistributedServices) Broadcast() messaging.Broadcaster { if ds == nil || ds.Bus == nil { return nil } return ds.Bus } // 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://") }