mirror of
https://github.com/mudler/LocalAI.git
synced 2026-10-05 12:34:43 -04:00
* 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>
749 lines
32 KiB
Go
749 lines
32 KiB
Go
package application
|
|
|
|
import (
|
|
"crypto/rand"
|
|
"encoding/hex"
|
|
"fmt"
|
|
"os"
|
|
"path/filepath"
|
|
"time"
|
|
|
|
"github.com/mudler/LocalAI/core/backend"
|
|
"github.com/mudler/LocalAI/core/config"
|
|
"github.com/mudler/LocalAI/core/gallery"
|
|
"github.com/mudler/LocalAI/core/http/auth"
|
|
"github.com/mudler/LocalAI/core/services/failover"
|
|
"github.com/mudler/LocalAI/core/services/galleryop"
|
|
"github.com/mudler/LocalAI/core/services/jobs"
|
|
"github.com/mudler/LocalAI/core/services/messaging"
|
|
"github.com/mudler/LocalAI/core/services/modeladmin"
|
|
"github.com/mudler/LocalAI/core/services/monitoring"
|
|
"github.com/mudler/LocalAI/core/services/nodes"
|
|
"github.com/mudler/LocalAI/core/services/routing/admission"
|
|
"github.com/mudler/LocalAI/core/services/routing/billing"
|
|
"github.com/mudler/LocalAI/core/services/routing/pii"
|
|
"github.com/mudler/LocalAI/core/services/routing/router"
|
|
"github.com/mudler/LocalAI/core/services/storage"
|
|
coreStartup "github.com/mudler/LocalAI/core/startup"
|
|
"github.com/mudler/LocalAI/core/trace"
|
|
"github.com/mudler/LocalAI/internal"
|
|
"github.com/mudler/LocalAI/pkg/downloader"
|
|
"github.com/mudler/LocalAI/pkg/modelartifacts"
|
|
"github.com/mudler/LocalAI/pkg/signals"
|
|
"github.com/mudler/LocalAI/pkg/vram"
|
|
|
|
"github.com/mudler/LocalAI/pkg/model"
|
|
"github.com/mudler/LocalAI/pkg/sanitize"
|
|
"github.com/mudler/LocalAI/pkg/xsysinfo"
|
|
"github.com/mudler/xlog"
|
|
)
|
|
|
|
func New(opts ...config.AppOption) (*Application, error) {
|
|
options := config.NewApplicationConfig(opts...)
|
|
|
|
// Store a copy of the startup config (env/CLI only, before file
|
|
// loading): the settings endpoint uses it to tell env-provided API
|
|
// keys apart from runtime-managed ones.
|
|
startupConfigCopy := *options
|
|
|
|
// Merge persisted runtime settings BEFORE anything consumes options:
|
|
// model-config defaults (ToConfigLoaderOptions), gallery services, the
|
|
// watchdog, and the MITM listener are all configured from options
|
|
// further down. Loading late (the old call site, after model configs
|
|
// were read) meant boot-loaded models saw pre-file defaults for one
|
|
// full restart. Env/CLI values win over the file (see
|
|
// ApplyRuntimeSettingsAtStartup).
|
|
loadRuntimeSettingsFromFile(options)
|
|
|
|
// WithThreads no longer eagerly resolves 0 (so the settings merge can
|
|
// tell "unset" from "env-set"); resolve the physical-core default now
|
|
// that env, CLI, and file have all had their say.
|
|
if options.Threads == 0 {
|
|
options.Threads = xsysinfo.CPUPhysicalCores()
|
|
}
|
|
|
|
trace.ConfigureBackendTracePersistence(options.DataPath)
|
|
application := newApplication(options)
|
|
application.startupConfig = &startupConfigCopy
|
|
|
|
xlog.Info("Starting LocalAI", "threads", options.Threads, "modelsPath", options.SystemState.Model.ModelsPath)
|
|
xlog.Info("LocalAI version", "version", internal.PrintableVersion())
|
|
|
|
if err := application.start(); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
caps, err := xsysinfo.CPUCapabilities()
|
|
if err == nil {
|
|
xlog.Debug("CPU capabilities", "capabilities", caps)
|
|
}
|
|
gpus, err := xsysinfo.GPUs()
|
|
if err == nil {
|
|
xlog.Debug("GPU count", "count", len(gpus))
|
|
for _, gpu := range gpus {
|
|
xlog.Debug("GPU", "gpu", gpu.String())
|
|
}
|
|
}
|
|
|
|
// Make sure directories exists
|
|
if options.SystemState.Model.ModelsPath == "" {
|
|
return nil, fmt.Errorf("models path cannot be empty")
|
|
}
|
|
|
|
err = os.MkdirAll(options.SystemState.Model.ModelsPath, 0o750)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("unable to create ModelPath: %q", err)
|
|
}
|
|
|
|
// Reap *.partial downloads abandoned by a previous run (killed mid-transfer
|
|
// by an OOM/restart, or stalled before cleanup could run). The 24h window
|
|
// is well beyond any legitimate in-flight download, so this never trims an
|
|
// active transfer; it just stops dead partials accumulating on the volume.
|
|
if removed, cErr := downloader.CleanupStalePartialFiles(options.SystemState.Model.ModelsPath, 24*time.Hour); cErr != nil {
|
|
xlog.Warn("Failed to reap stale partial downloads", "error", cErr)
|
|
} else if removed > 0 {
|
|
xlog.Info("Reaped stale partial downloads", "count", removed)
|
|
}
|
|
// Managed artifacts stage into a per-writer tree, which a crashed writer's
|
|
// successor no longer overwrites for it, so the tree itself needs reaping
|
|
// too. Sweeping here as well as on the materialization path is what
|
|
// reclaims a volume whose abandoned artifact is never requested again.
|
|
if removed, cErr := modelartifacts.SweepStalePartialTrees(options.SystemState.Model.ModelsPath, modelartifacts.PartialOrphanTTL, ""); cErr != nil {
|
|
xlog.Warn("Failed to reap abandoned artifact partials", "error", cErr)
|
|
} else if removed > 0 {
|
|
xlog.Info("Reaped abandoned artifact partials", "count", removed)
|
|
}
|
|
if options.GeneratedContentDir != "" {
|
|
err := os.MkdirAll(options.GeneratedContentDir, 0o750)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("unable to create ImageDir: %q", err)
|
|
}
|
|
}
|
|
if options.UploadDir != "" {
|
|
err := os.MkdirAll(options.UploadDir, 0o750)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("unable to create UploadDir: %q", err)
|
|
}
|
|
}
|
|
|
|
// Create and migrate data directory
|
|
if options.DataPath != "" {
|
|
if err := os.MkdirAll(options.DataPath, 0o750); err != nil {
|
|
return nil, fmt.Errorf("unable to create DataPath: %q", err)
|
|
}
|
|
// Migrate data from DynamicConfigsDir to DataPath if needed
|
|
if options.DynamicConfigsDir != "" && options.DataPath != options.DynamicConfigsDir {
|
|
migrateDataFiles(options.DynamicConfigsDir, options.DataPath)
|
|
}
|
|
}
|
|
// Initialize auth database if auth is enabled
|
|
if options.Auth.Enabled {
|
|
// Auto-generate HMAC secret if not provided
|
|
if options.Auth.APIKeyHMACSecret == "" {
|
|
secretFile := filepath.Join(options.DataPath, ".hmac_secret")
|
|
secret, err := loadOrGenerateHMACSecret(secretFile)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to initialize HMAC secret: %w", err)
|
|
}
|
|
options.Auth.APIKeyHMACSecret = secret
|
|
}
|
|
|
|
authDB, err := auth.InitDB(options.Auth.DatabaseURL)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to initialize auth database: %w", err)
|
|
}
|
|
application.authDB = authDB
|
|
xlog.Info("Auth enabled", "database", sanitize.URL(options.Auth.DatabaseURL))
|
|
|
|
// Start session and expired API key cleanup goroutine
|
|
go func() {
|
|
ticker := time.NewTicker(1 * time.Hour)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-options.Context.Done():
|
|
return
|
|
case <-ticker.C:
|
|
if err := auth.CleanExpiredSessions(authDB); err != nil {
|
|
xlog.Error("failed to clean expired sessions", "error", err)
|
|
}
|
|
if err := auth.CleanExpiredAPIKeys(authDB); err != nil {
|
|
xlog.Error("failed to clean expired API keys", "error", err)
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
}
|
|
|
|
// Initialize the OTel + Prometheus metric pipeline before any
|
|
// counter is created. monitoring.NewLocalAIMetricsService calls
|
|
// otel.SetMeterProvider, so any subsequent otel.Meter() call —
|
|
// including billing.NewRecorder below — sees the real provider
|
|
// rather than the no-op global. Initialising metrics later (in
|
|
// core/http/app.go) leaves billing's counters bound to a no-op
|
|
// meter and never reaches /metrics. We deliberately ignore
|
|
// DisableMetrics here for ordering purposes; the HTTP middleware
|
|
// that records api_call histograms is still gated.
|
|
if !options.DisableMetrics {
|
|
ms, err := monitoring.NewLocalAIMetricsService()
|
|
if err != nil {
|
|
xlog.Error("failed to initialize metrics provider", "error", err)
|
|
} else {
|
|
application.metricsService = ms
|
|
// Bind the billing package's counters to the same meter the
|
|
// metrics service exports. Without this, billing's counters
|
|
// resolve via the OTel global and never reach /metrics.
|
|
billing.SetMeter(ms.Meter)
|
|
}
|
|
}
|
|
|
|
// Wire the routing-module billing recorder. The recorder runs in
|
|
// every mode (auth on/off, distributed/single-node) so that token
|
|
// tracking is not gated on auth — a no-auth single-user box still
|
|
// gets dashboards and `/api/usage` populated.
|
|
//
|
|
// fallbackUser is wired *unconditionally* when stats are enabled.
|
|
// UsageMiddleware uses it as the attribution source whenever
|
|
// auth.GetUser(c) is nil — that covers (a) no-auth deployments and
|
|
// (b) internal callers under auth-on (cron flushers, distributed
|
|
// worker callbacks) that hit a recordable endpoint without a user
|
|
// in context. The billing.user_id_present invariant still rejects
|
|
// empty IDs; LocalUser() returns a stable UUID per data path.
|
|
if !options.DisableStats {
|
|
var statsBackend billing.StatsBackend
|
|
switch {
|
|
case application.authDB != nil:
|
|
statsBackend = billing.NewGormBackend(application.authDB, 0, 0)
|
|
xlog.Info("stats: using auth DB for usage records")
|
|
default:
|
|
statsBackend = billing.NewMemoryBackend(0)
|
|
xlog.Info("stats: using in-memory ring buffer (no-auth single-user mode)")
|
|
}
|
|
application.fallbackUser = billing.LocalUser(options.DataPath)
|
|
application.statsRecorder = billing.NewRecorder(statsBackend)
|
|
// Drain pending records on SIGTERM. The GORM backend buffers up
|
|
// to maxPending (5k) records across a 5s flush tick, so without
|
|
// this the last few seconds of usage disappear on graceful exit.
|
|
signals.RegisterGracefulTerminationHandler(func() {
|
|
_ = application.statsRecorder.Close()
|
|
})
|
|
xlog.Info("stats: fallback user wired", "local_user_id", application.fallbackUser.ID)
|
|
} else {
|
|
xlog.Info("stats: disabled by --disable-stats")
|
|
}
|
|
|
|
// Wire the PII filter subsystem. The redactor is now a stateless
|
|
// handle — detection is driven by per-model NER detectors
|
|
// (pii.detectors → the detector model's pii_detection policy), run
|
|
// request-side by the chat middleware and the MITM input path. The
|
|
// regex tier was removed; redaction is opt-in per model via
|
|
// PIIIsEnabled(). The event store backs the /api/pii/events audit log.
|
|
application.piiRedactor = &pii.Redactor{}
|
|
application.piiEvents = pii.NewMemoryEventStore(0)
|
|
|
|
// Wire the routing decision log. Always-on when stats are enabled —
|
|
// the per-router admin page reads this as the live activity feed
|
|
// and as input to drift checks for subsystem 5.
|
|
if !options.DisableStats {
|
|
application.routerDecisions = router.NewMemoryDecisionStore(0)
|
|
}
|
|
// Process-wide classifier cache shared across all route middlewares so
|
|
// the embedding-cache stats endpoint sees a single source of truth.
|
|
application.routerRegistry = router.NewRegistry()
|
|
|
|
// Failover chains: probe targets and track which one is active per
|
|
// chain. WithOnWarmChanged pins and preloads warm local targets so a
|
|
// switch to them does not wait for a cold load.
|
|
application.failoverManager = failover.New(application.ModelConfigLoader(),
|
|
failover.WithProber(failover.NewProber(failoverLoadedBackend(application.ModelLoader()), options.ProxyAPIKeyEnvLookup)),
|
|
failover.WithOnWarmChanged(application.applyFailoverWarmTargets),
|
|
)
|
|
// The assistant client was built in start() (above), before this
|
|
// manager existed; wire it now so list_failover_chains /
|
|
// pin_failover_target / unpin_failover_target see real chains.
|
|
if application.assistantClient != nil {
|
|
application.assistantClient.Failover = application.failoverManager
|
|
}
|
|
|
|
// Subsystem 5: admission control. Limiter is always wired so a
|
|
// model that gains a limits: block via gallery install or YAML
|
|
// edit takes effect on the next restart without conditional plumbing.
|
|
application.admissionLimiter = admission.New()
|
|
|
|
// Wire JobStore for DB-backed task/job persistence whenever auth DB is available.
|
|
// This ensures tasks and jobs survive restarts in both single-node and distributed modes.
|
|
if application.authDB != nil && application.agentJobService != nil {
|
|
dbJobStore, err := jobs.NewJobStore(application.authDB)
|
|
if err != nil {
|
|
xlog.Error("Failed to create job store for auth DB", "error", err)
|
|
} else {
|
|
application.agentJobService.SetDistributedJobStore(dbJobStore)
|
|
}
|
|
}
|
|
|
|
// Initialize distributed mode services (NATS, object storage, node registry)
|
|
// revisionStore is built inside the distributed block below but used after
|
|
// the model configs are loaded, so it is declared out here.
|
|
var revisionStore modeladmin.RevisionStore
|
|
|
|
distSvc, err := initDistributed(options, application.authDB, application.ModelConfigLoader(),
|
|
&failoverPinnedResolver{base: application.ModelConfigLoader(), fm: application.failoverManager})
|
|
if err != nil {
|
|
return nil, fmt.Errorf("distributed mode initialization failed: %w", err)
|
|
}
|
|
if distSvc != nil {
|
|
application.distributed = distSvc
|
|
// Before failoverManager.Run starts below: the gate and sync must be
|
|
// in place for its first tick.
|
|
application.startFailoverDistributed(options.Context)
|
|
// Wire remote model unloader so ShutdownModel works for remote nodes
|
|
// Uses NATS to tell serve-backend nodes to Free + kill their backend process
|
|
application.modelLoader.SetRemoteUnloader(distSvc.Unloader)
|
|
// Wire ModelRouter so grpcModel() delegates to SmartRouter in distributed mode
|
|
application.modelLoader.SetModelRouter(distSvc.ModelAdapter.AsModelRouter())
|
|
// Wire DistributedModelStore so shutdown/list/watchdog can find remote models
|
|
distStore := nodes.NewDistributedModelStore(
|
|
model.NewInMemoryModelStore(),
|
|
distSvc.Registry,
|
|
)
|
|
application.modelLoader.SetModelStore(distStore)
|
|
// Drop the local stub when a model's last replica leaves the registry.
|
|
// The store reports local stubs UNION registry rows, and every removal
|
|
// path deletes the row only, so without this the frontend keeps
|
|
// reporting a model as loaded long after the replica is gone.
|
|
// Registered unconditionally: this is independent of the prefix cache.
|
|
distSvc.Registry.AddReplicaRemovedHook(
|
|
nodes.NewLocalStubInvalidator(distSvc.Registry, distStore),
|
|
)
|
|
// Start health monitor
|
|
distSvc.Health.Start(options.Context)
|
|
// Start replica reconciler for auto-scaling model replicas
|
|
if distSvc.Reconciler != nil {
|
|
go distSvc.Reconciler.Run(options.Context)
|
|
}
|
|
go distSvc.ModelCleanup.Run(options.Context)
|
|
// In distributed mode, MCP CI jobs are executed by agent workers (not the frontend)
|
|
// because the frontend can't create MCP sessions (e.g., stdio servers using docker).
|
|
// The dispatcher only enqueues jobs and persists the results and traces workers publish.
|
|
|
|
// Wire model config loader so job events include model config for agent workers
|
|
distSvc.Dispatcher.SetModelConfigLoader(application.backendLoader)
|
|
|
|
// Start job dispatcher — abort startup if it fails, as jobs would be accepted but never dispatched
|
|
if err := distSvc.Dispatcher.Start(options.Context); err != nil {
|
|
return nil, fmt.Errorf("starting job dispatcher: %w", err)
|
|
}
|
|
// Start ephemeral file cleanup
|
|
storage.StartEphemeralCleanup(options.Context, distSvc.FileMgr, 0, 0)
|
|
// Wire distributed backends into AgentJobService (before Start)
|
|
if application.agentJobService != nil {
|
|
application.agentJobService.SetDistributedBackends(distSvc.Dispatcher)
|
|
application.agentJobService.SetDistributedJobStore(distSvc.JobStore)
|
|
// Keep agent tasks consistent across replicas (jobs already sync via the
|
|
// dispatcher + DB read-through). Same NATS client the dispatcher uses.
|
|
application.agentJobService.SetTaskSyncNATS(distSvc.Nats)
|
|
}
|
|
// Wire skill store into AgentPoolService (wired at pool start time via closure)
|
|
// The actual wiring happens in StartAgentPool since the pool doesn't exist yet.
|
|
|
|
// Wire NATS and gallery store into GalleryService for cross-instance progress/cancel
|
|
if application.galleryService != nil {
|
|
application.galleryService.SetNATSClient(distSvc.Nats)
|
|
if distSvc.DistStores != nil && distSvc.DistStores.Gallery != nil {
|
|
// Clean up stale in-progress operations from previous crashed instances
|
|
if _, err := distSvc.DistStores.Gallery.CleanStale(30 * time.Minute); err != nil {
|
|
xlog.Warn("Failed to clean stale gallery operations", "error", err)
|
|
}
|
|
application.galleryService.SetGalleryStore(distSvc.DistStores.Gallery)
|
|
|
|
// Reap stale ops periodically, not just at boot: an op orphaned by
|
|
// a replica that died mid-install (its foreground handler goroutine
|
|
// gone) would otherwise linger "processing" in the UI until the next
|
|
// restart. 30m matches the install/upgrade ceiling so a genuinely
|
|
// slow op is never reaped out from under itself.
|
|
gsvc := application.galleryService
|
|
go func() {
|
|
ticker := time.NewTicker(15 * time.Minute)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-options.Context.Done():
|
|
return
|
|
case <-ticker.C:
|
|
if _, err := gsvc.ReapStaleOperations(30 * time.Minute); err != nil {
|
|
xlog.Warn("Failed to reap stale gallery operations", "error", err)
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
}
|
|
// Hydrate from the store first so the wildcard subscriber finds an
|
|
// already-populated statuses map for any operations still in flight
|
|
// on a peer replica.
|
|
if err := application.galleryService.Hydrate(); err != nil {
|
|
xlog.Warn("Gallery service hydrate failed", "error", err)
|
|
}
|
|
// Bind cache-invalidation handler before SubscribeBroadcasts so the
|
|
// first inbound event is already routed. Peer replicas install a
|
|
// model and broadcast on SubjectCacheInvalidateModels; this
|
|
// callback re-runs LoadModelConfigsFromPath so a subsequent chat
|
|
// completion that load-balances onto this replica finds the new
|
|
// config. The originating replica reloads inline in modelHandler
|
|
// and never enters this path.
|
|
gs := application.galleryService
|
|
sys := options.SystemState
|
|
cfgLoaderOpts := options.ToConfigLoaderOptions()
|
|
modelRevisionLifecycle := modeladmin.NewDistributedModelRevisionLifecycle(distSvc.Registry, distSvc.ModelCleanup)
|
|
gs.SetModelRevisionLifecycle(modelRevisionLifecycle)
|
|
// Captured here, used after the model configs are loaded below: the
|
|
// resync reads the loader, which is still empty at this point.
|
|
revisionStore = modeladmin.NewRevisionStore(distSvc.Registry, modelRevisionLifecycle)
|
|
gs.OnModelsChanged = func(evt messaging.CacheInvalidateEvent) {
|
|
// ApplyRemoteChange honors the op: a "delete" prunes the element
|
|
// (a reload-from-path is additive and cannot drop it), anything
|
|
// else reloads from disk; a named element's running instance is
|
|
// shut down so the new config takes effect. The originating
|
|
// replica reloads inline and never depends on this path.
|
|
if err := modeladmin.ApplyRemoteChange(options.Context, application.ModelConfigLoader(), sys.Model.ModelsPath, evt, modelRevisionLifecycle, cfgLoaderOpts...); err != nil {
|
|
xlog.Warn("Failed to apply peer model config change", "error", err)
|
|
}
|
|
}
|
|
if err := application.galleryService.SubscribeBroadcasts(); err != nil {
|
|
xlog.Warn("Gallery service subscribe failed", "error", err)
|
|
}
|
|
// Wire distributed model/backend managers so delete propagates to workers
|
|
application.galleryService.SetModelManager(
|
|
nodes.NewDistributedModelManager(options, application.modelLoader, distSvc.Unloader),
|
|
)
|
|
application.galleryService.SetBackendManager(
|
|
nodes.NewDistributedBackendManager(options, application.modelLoader, distSvc.Unloader, distSvc.Registry, application.galleryService),
|
|
)
|
|
}
|
|
}
|
|
|
|
// Start AgentJobService (after distributed wiring so it knows whether to use local or NATS)
|
|
if application.agentJobService != nil {
|
|
if err := application.agentJobService.Start(options.Context); err != nil {
|
|
return nil, fmt.Errorf("starting agent job service: %w", err)
|
|
}
|
|
}
|
|
|
|
if err := coreStartup.InstallModels(options.Context, application.GalleryService(), options.Galleries, options.BackendGalleries, options.SystemState, application.ModelLoader(), options.EnforcePredownloadScans, options.AutoloadBackendGalleries, options.RequireBackendIntegrity, nil, options.ModelsURL...); err != nil {
|
|
xlog.Error("error installing models", "error", err)
|
|
}
|
|
|
|
for _, backend := range options.ExternalBackends {
|
|
if err := galleryop.InstallExternalBackend(options.Context, options.BackendGalleries, options.SystemState, application.ModelLoader(), nil, backend, "", "", false, options.RequireBackendIntegrity); err != nil {
|
|
xlog.Error("error installing external backend", "error", err)
|
|
}
|
|
}
|
|
|
|
configLoaderOpts := options.ToConfigLoaderOptions()
|
|
|
|
if err := application.ModelConfigLoader().LoadModelConfigsFromPath(options.SystemState.Model.ModelsPath, configLoaderOpts...); err != nil {
|
|
xlog.Error("error loading config files", "error", err)
|
|
}
|
|
|
|
// Bring the controller's stored revisions back in line with the
|
|
// configuration just loaded. An inference request may only establish a
|
|
// revision, never replace one, so a model whose stored value has drifted
|
|
// stays unroutable until something republishes it. This has to run after
|
|
// the load above: the loader is empty until then, and a resync against an
|
|
// empty loader silently reconciles nothing.
|
|
if revisionStore != nil {
|
|
if err := modeladmin.ResyncModelConfigRevisions(options.Context, application.ModelConfigLoader(), options, revisionStore); err != nil {
|
|
xlog.Warn("Failed to resync model config revisions", "error", err)
|
|
}
|
|
}
|
|
|
|
if err := gallery.RegisterBackends(options.SystemState, application.ModelLoader()); err != nil {
|
|
xlog.Error("error registering external backends", "error", err)
|
|
}
|
|
|
|
// Start background upgrade checker for backends.
|
|
// In distributed mode, uses PostgreSQL advisory lock so only one frontend
|
|
// instance runs periodic checks (avoids duplicate upgrades across replicas).
|
|
if len(options.BackendGalleries) > 0 {
|
|
// Pass a lazy getter for the backend manager so the checker always
|
|
// uses the active one — DistributedBackendManager is swapped in above
|
|
// and asks workers for their installed backends, which is what
|
|
// upgrade detection needs in distributed mode.
|
|
bmFn := func() galleryop.BackendManager { return application.GalleryService().BackendManager() }
|
|
uc := NewUpgradeChecker(options, application.ModelLoader(), application.distributedDB(), bmFn)
|
|
application.upgradeChecker = uc
|
|
// Refresh the upgrade cache the moment a backend op finishes — otherwise
|
|
// the UI keeps showing a just-upgraded backend as upgradeable until the
|
|
// next 6-hour tick. TriggerCheck is non-blocking.
|
|
if gs := application.GalleryService(); gs != nil {
|
|
gs.OnBackendOpCompleted = uc.TriggerCheck
|
|
}
|
|
go uc.Run(options.Context)
|
|
}
|
|
|
|
// Wire gallery generation counter into VRAM caches so they invalidate
|
|
// when gallery data refreshes instead of using a fixed TTL.
|
|
vram.SetGalleryGenerationFunc(gallery.GalleryGeneration)
|
|
if options.AutoloadGalleries {
|
|
if options.VRAMPersistentCache {
|
|
// Remote GGUF probes can transfer substantial metadata. Keep successful
|
|
// results across restarts so the startup warmer does not repeat that work.
|
|
vram.ConfigurePersistentCache(filepath.Join(options.SystemState.Model.ModelsPath, "..", "cache", "vram"), 24*time.Hour)
|
|
}
|
|
|
|
// Fill those caches ahead of the first visitor. An estimate for an entry
|
|
// nobody has asked about yet costs a remote probe of its weight files, and
|
|
// the model gallery asks for one per row, so without this the first page
|
|
// spends seconds filling in its own sizes while somebody watches it.
|
|
// Non-blocking, and bounded: see DefaultEstimateWarmConfig.
|
|
gallery.WarmEstimateCache(options.Context, options.Galleries, options.SystemState, gallery.EstimateWarmConfigFromEnv())
|
|
}
|
|
|
|
if options.ConfigFile != "" {
|
|
if err := application.ModelConfigLoader().LoadMultipleModelConfigsSingleFile(options.ConfigFile, configLoaderOpts...); err != nil {
|
|
xlog.Error("error loading config file", "error", err)
|
|
}
|
|
}
|
|
|
|
if err := application.ModelConfigLoader().PreloadWithContext(options.Context, options.SystemState.Model.ModelsPath); err != nil {
|
|
xlog.Error("error downloading models", "error", err)
|
|
}
|
|
|
|
if options.PreloadJSONModels != "" {
|
|
if err := galleryop.ApplyGalleryFromString(options.SystemState, application.ModelLoader(), options.EnforcePredownloadScans, options.AutoloadBackendGalleries, options.Galleries, options.BackendGalleries, options.PreloadJSONModels, options.RequireBackendIntegrity, gallery.WithArtifactMaterializer(options.ModelArtifactMaterializer)); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
|
|
if options.PreloadModelsFromPath != "" {
|
|
if err := galleryop.ApplyGalleryFromFile(options.SystemState, application.ModelLoader(), options.EnforcePredownloadScans, options.AutoloadBackendGalleries, options.Galleries, options.BackendGalleries, options.PreloadModelsFromPath, options.RequireBackendIntegrity, gallery.WithArtifactMaterializer(options.ModelArtifactMaterializer)); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
|
|
if options.Debug {
|
|
for _, v := range application.ModelConfigLoader().GetAllModelsConfigs() {
|
|
xlog.Debug("Model", "name", v.Name, "config", v)
|
|
}
|
|
}
|
|
|
|
// Wire the cloudproxy MITM listener. Opt-in: empty MITMListen
|
|
// means "no MITM" — operators must explicitly choose to start
|
|
// it because clients have to install the generated CA cert.
|
|
// The handler reuses the global redactor + event store so an
|
|
// admin who's already configured PII filtering for direct API
|
|
// traffic doesn't need a parallel config for MITM traffic.
|
|
// Runs after loadRuntimeSettingsFromFile so a listener configured
|
|
// via /api/settings is brought back up across restarts.
|
|
startMITMIfConfigured(application, options)
|
|
|
|
application.ModelLoader().SetBackendLoggingEnabled(options.EnableBackendLogging)
|
|
|
|
// Safety-net cleanup if the application context is cancelled without
|
|
// the caller invoking Shutdown directly. This is fire-and-forget — it
|
|
// races binary exit and is unreliable in tests; the deterministic path
|
|
// is application.Shutdown(), which Shutdown's sync.Once dedupes with
|
|
// this goroutine.
|
|
go func() {
|
|
<-options.Context.Done()
|
|
xlog.Debug("Context canceled, shutting down")
|
|
if err := application.Shutdown(); err != nil {
|
|
xlog.Error("error while stopping all grpc backends", "error", err)
|
|
}
|
|
}()
|
|
|
|
// Initialize watchdog with current settings (after loading from file)
|
|
initializeWatchdog(application, options)
|
|
|
|
if options.LoadToMemory != nil && !options.SingleBackend {
|
|
for _, m := range options.LoadToMemory {
|
|
xlog.Debug("Auto loading model into memory from file", "model", m)
|
|
// Same path as POST /backend/load: a realtime pipeline model expands
|
|
// to its sub-models, and load failures are recorded as model_load
|
|
// traces.
|
|
if _, err := backend.PreloadModelByName(options.Context, application.ModelConfigLoader(), application.ModelLoader(), options, m); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
}
|
|
|
|
// Start the failover scheduler: it syncs chains from config, runs
|
|
// liveness/recovery probes and dwell-based fail-back. Run is the only
|
|
// caller of Sync in production so onWarm callbacks stay ordered.
|
|
failover.RegisterMetrics(application.failoverManager)
|
|
go application.failoverManager.Run(options.Context)
|
|
|
|
// Watch the configuration directory
|
|
startWatcher(options)
|
|
|
|
// Everything that must happen before this process can serve a request has
|
|
// happened. Flip readiness last, and only on the success path — the early
|
|
// `return nil, err` exits above abort startup, and an application that
|
|
// never finished starting must never report itself ready.
|
|
application.markStartupComplete()
|
|
|
|
xlog.Info("core/startup process completed!")
|
|
return application, nil
|
|
}
|
|
|
|
func startWatcher(options *config.ApplicationConfig) {
|
|
if options.DynamicConfigsDir == "" {
|
|
// No need to start the watcher if the directory is not set
|
|
return
|
|
}
|
|
|
|
if _, err := os.Stat(options.DynamicConfigsDir); err != nil {
|
|
if os.IsNotExist(err) {
|
|
// We try to create the directory if it does not exist and was specified
|
|
if err := os.MkdirAll(options.DynamicConfigsDir, 0o700); err != nil {
|
|
xlog.Error("failed creating DynamicConfigsDir", "error", err)
|
|
}
|
|
} else {
|
|
// something else happened, we log the error and don't start the watcher
|
|
xlog.Error("failed to read DynamicConfigsDir, watcher will not be started", "error", err)
|
|
return
|
|
}
|
|
}
|
|
|
|
configHandler := newConfigFileHandler(options)
|
|
if err := configHandler.Watch(); err != nil {
|
|
xlog.Error("failed creating watcher", "error", err)
|
|
}
|
|
}
|
|
|
|
// loadRuntimeSettingsFromFile merges runtime_settings.json into options
|
|
// with env-over-file precedence. Field coverage and the precedence rules
|
|
// live in the config registry (ApplyRuntimeSettingsAtStartup); this wrapper
|
|
// only handles the file read. No-op when DynamicConfigsDir is unset.
|
|
func loadRuntimeSettingsFromFile(options *config.ApplicationConfig) {
|
|
if options.DynamicConfigsDir == "" {
|
|
return
|
|
}
|
|
// ReadPersistedSettings treats a missing file as zero settings, so probe
|
|
// for it here only to keep the log honest: "loaded" must mean a file was
|
|
// actually read, not that we merged an all-nil struct.
|
|
if _, err := os.Stat(filepath.Join(options.DynamicConfigsDir, "runtime_settings.json")); os.IsNotExist(err) {
|
|
xlog.Debug("runtime_settings.json not found, using defaults")
|
|
return
|
|
}
|
|
settings, err := options.ReadPersistedSettings()
|
|
if err != nil {
|
|
xlog.Warn("failed to read runtime_settings.json", "error", err)
|
|
return
|
|
}
|
|
options.ApplyRuntimeSettingsAtStartup(&settings)
|
|
xlog.Debug("Runtime settings loaded from runtime_settings.json")
|
|
}
|
|
|
|
// initializeWatchdog initializes the watchdog with current ApplicationConfig settings
|
|
func initializeWatchdog(application *Application, options *config.ApplicationConfig) {
|
|
// Get effective max active backends (considers both MaxActiveBackends and deprecated SingleBackend)
|
|
lruLimit := options.GetEffectiveMaxActiveBackends()
|
|
|
|
// Create watchdog if enabled OR if LRU limit is set OR if memory reclaimer is enabled
|
|
if options.WatchDog || lruLimit > 0 || options.MemoryReclaimerEnabled {
|
|
wd := model.NewWatchDog(
|
|
model.WithProcessManager(application.ModelLoader()),
|
|
model.WithBusyTimeout(options.WatchDogBusyTimeout),
|
|
model.WithIdleTimeout(options.WatchDogIdleTimeout),
|
|
model.WithWatchdogInterval(options.WatchDogInterval),
|
|
model.WithBusyCheck(options.WatchDogBusy),
|
|
model.WithIdleCheck(options.WatchDogIdle),
|
|
model.WithLRULimit(lruLimit),
|
|
model.WithMemoryReclaimer(options.MemoryReclaimerEnabled, options.MemoryReclaimerThreshold),
|
|
model.WithForceEvictionWhenBusy(options.ForceEvictionWhenBusy),
|
|
model.WithSizeAwareEviction(options.SizeAwareEviction),
|
|
)
|
|
application.ModelLoader().SetWatchDog(wd)
|
|
|
|
// Initialize ModelLoader LRU eviction retry settings
|
|
application.ModelLoader().SetLRUEvictionRetrySettings(
|
|
options.LRUEvictionMaxRetries,
|
|
options.LRUEvictionRetryInterval,
|
|
)
|
|
|
|
// Sync per-model state from configs to the watchdog. Without this,
|
|
// `pinned: true` and `concurrency_groups:` are only honored after a
|
|
// settings-driven RestartWatchdog and never at boot.
|
|
application.SyncPinnedModelsToWatchdog()
|
|
application.SyncModelGroupsToWatchdog()
|
|
|
|
// Start watchdog goroutine if any periodic checks are enabled
|
|
// LRU eviction doesn't need the Run() loop - it's triggered on model load
|
|
// But memory reclaimer needs the Run() loop for periodic checking
|
|
if options.WatchDogBusy || options.WatchDogIdle || options.MemoryReclaimerEnabled {
|
|
go wd.Run()
|
|
}
|
|
|
|
go func() {
|
|
<-options.Context.Done()
|
|
xlog.Debug("Context canceled, shutting down")
|
|
wd.Shutdown()
|
|
}()
|
|
}
|
|
}
|
|
|
|
// loadOrGenerateHMACSecret loads an HMAC secret from the given file path,
|
|
// or generates a random 32-byte secret and persists it if the file doesn't exist.
|
|
func loadOrGenerateHMACSecret(path string) (string, error) {
|
|
data, err := os.ReadFile(path)
|
|
if err == nil {
|
|
secret := string(data)
|
|
if len(secret) >= 32 {
|
|
return secret, nil
|
|
}
|
|
}
|
|
|
|
b := make([]byte, 32)
|
|
if _, err := rand.Read(b); err != nil {
|
|
return "", fmt.Errorf("failed to generate HMAC secret: %w", err)
|
|
}
|
|
secret := hex.EncodeToString(b)
|
|
|
|
if err := os.WriteFile(path, []byte(secret), 0o600); err != nil {
|
|
return "", fmt.Errorf("failed to persist HMAC secret: %w", err)
|
|
}
|
|
|
|
xlog.Info("Generated new HMAC secret for API key hashing", "path", path)
|
|
return secret, nil
|
|
}
|
|
|
|
// migrateDataFiles moves persistent data files from the old config directory
|
|
// to the new data directory. Only moves files that exist in src but not in dst.
|
|
func migrateDataFiles(srcDir, dstDir string) {
|
|
// Files and directories to migrate
|
|
items := []string{
|
|
"agent_tasks.json",
|
|
"agent_jobs.json",
|
|
"collections",
|
|
"assets",
|
|
}
|
|
|
|
migrated := false
|
|
for _, item := range items {
|
|
srcPath := filepath.Join(srcDir, item)
|
|
dstPath := filepath.Join(dstDir, item)
|
|
|
|
// Only migrate if source exists and destination does not
|
|
if _, err := os.Stat(srcPath); os.IsNotExist(err) {
|
|
continue
|
|
}
|
|
if _, err := os.Stat(dstPath); err == nil {
|
|
continue // destination already exists, skip
|
|
}
|
|
|
|
if err := os.Rename(srcPath, dstPath); err != nil {
|
|
xlog.Warn("Failed to migrate data file, will copy instead", "src", srcPath, "dst", dstPath, "error", err)
|
|
// os.Rename fails across filesystems, fall back to leaving in place
|
|
// and log a warning for the user to manually move
|
|
xlog.Warn("Data file remains in old location, please move manually", "src", srcPath, "dst", dstPath)
|
|
continue
|
|
}
|
|
migrated = true
|
|
xlog.Info("Migrated data file to new data path", "src", srcPath, "dst", dstPath)
|
|
}
|
|
|
|
if migrated {
|
|
xlog.Info("Data migration complete", "from", srcDir, "to", dstDir)
|
|
}
|
|
}
|