Files
LocalAI/tests/e2e/distributed/model_routing_test.go
Ettore Di Giacinto 016686a3db chore(distributed): take the nats-io modules out of the build
Distributed mode has not dialled a message broker since the control plane
moved onto the workers' own outward tunnels and every fan-out family moved
onto PostgreSQL LISTEN/NOTIFY. What was left was the dependency itself, and
the code that existed only to feed it.

Dropped from go.mod: nats-io/jwt/v2, nats-io/nats.go, nats-io/nkeys,
nats-io/nuid and testcontainers-go/modules/nats, along with the fourteen
indirect requires that only the NATS testcontainer pulled in. go.sum carries
no nats line either, so the removal is not the partial kind where the require
goes and the checksum stays.

Deleted with them: pkg/natsauth in full, the broker client's remaining
options and TLS files, the per-node JWT minting on both the register and the
approve path, and the natsauth.Config parameter threaded through the node
routes. The credential manager is renamed and stripped rather than deleted,
because it still holds the tunnel token that every re-registration rotates.

The bus flags stay accepted and ignored, and are now hidden, on every command
that had them, so an existing unit file, compose file or Helm values file
still starts on the day of the upgrade. What is not kept is the validation
that REQUIRED one: a distributed frontend started with no bus URL is no
longer fatal. The TLS paths lose type:"existingfile" deliberately, so a
certificate deleted along with the broker cannot fail a startup.

One operator-visible behaviour change: --nats-require-auth no longer makes an
agent worker wait through admin approval. Ask for that wait with
--distributed-require-auth, which already implied it. It is documented in the
migration section and pinned from both sides.

A deployment now needs PostgreSQL and the frontends' own HTTP listener, and
nothing else.

coverage-baseline.txt moves from 54.2 to 62.0.

Assisted-by: Claude Opus 5 [claude-code]
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
2026-09-20 03:05:35 +00:00

126 lines
4.5 KiB
Go

package distributed_test
import (
"context"
"github.com/mudler/LocalAI/core/config"
"github.com/mudler/LocalAI/core/services/nodes"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
pgdriver "gorm.io/driver/postgres"
"gorm.io/gorm"
"gorm.io/gorm/logger"
)
var _ = Describe("Model Routing", Label("Distributed"), func() {
var (
infra *TestInfra
db *gorm.DB
registry *nodes.NodeRegistry
)
BeforeEach(func() {
infra = SetupInfra("localai_routing_test")
var err error
db, err = gorm.Open(pgdriver.Open(infra.PGURL), &gorm.Config{
Logger: logger.Default.LogMode(logger.Silent),
})
Expect(err).ToNot(HaveOccurred())
registry, err = nodes.NewNodeRegistry(db)
Expect(err).ToNot(HaveOccurred())
})
Context("ModelRouterAdapter from SmartRouter", func() {
It("should create ModelRouterAdapter from SmartRouter", func() {
router := nodes.NewSmartRouter(registry, nodes.SmartRouterOptions{})
Expect(router).ToNot(BeNil())
adapter := nodes.NewModelRouterAdapter(router)
Expect(adapter).ToNot(BeNil())
// The adapter should provide a ModelRouter callback
routerFunc := adapter.AsModelRouter()
Expect(routerFunc).ToNot(BeNil())
})
It("should release in-flight counter on model unload", func() {
// Register a node with a loaded model
node := &nodes.BackendNode{
Name: "gpu-1", Address: "h1:50051",
}
Expect(registry.Register(context.Background(), node, true)).To(Succeed())
Expect(registry.SetNodeModel(context.Background(), node.ID, "llama3", 0, "loaded", "", 0)).To(Succeed())
Expect(registry.IncrementInFlight(context.Background(), node.ID, "llama3", 0)).To(Succeed())
Expect(registry.IncrementInFlight(context.Background(), node.ID, "llama3", 0)).To(Succeed())
// Verify in-flight count
models, err := registry.GetNodeModels(context.Background(), node.ID)
Expect(err).ToNot(HaveOccurred())
Expect(models[0].InFlight).To(Equal(2))
// FindAndLockNodeWithModel should return this node and atomically increment in-flight
foundNode, foundModel, err := registry.FindAndLockNodeWithModel(context.Background(), "llama3", nil, nil)
Expect(err).ToNot(HaveOccurred())
Expect(foundNode.ID).To(Equal(node.ID))
Expect(foundModel.ModelName).To(Equal("llama3"))
Expect(foundModel.InFlight).To(Equal(2), "InFlight returned is the pre-increment snapshot from the query")
// Verify the DB now has in_flight = 3 (2 manual + 1 from FindAndLock)
models, err = registry.GetNodeModels(context.Background(), node.ID)
Expect(err).ToNot(HaveOccurred())
Expect(models[0].InFlight).To(Equal(3))
// Simulate decrement (what Release does)
Expect(registry.DecrementInFlight(context.Background(), node.ID, "llama3", 0)).To(Succeed())
models, _ = registry.GetNodeModels(context.Background(), node.ID)
Expect(models[0].InFlight).To(Equal(2))
// The ModelRouterAdapter.ReleaseModel calls the stored Release function
router := nodes.NewSmartRouter(registry, nodes.SmartRouterOptions{})
adapter := nodes.NewModelRouterAdapter(router)
// ReleaseModel on an unknown model should be a no-op (no panic)
Expect(func() { adapter.ReleaseModel("nonexistent-model") }).ToNot(Panic())
})
It("should use SmartRouter to find nodes with a model", func() {
// Register multiple nodes
node1 := &nodes.BackendNode{
Name: "node-a", Address: "h1:50051",
}
node2 := &nodes.BackendNode{
Name: "node-b", Address: "h2:50051",
}
Expect(registry.Register(context.Background(), node1, true)).To(Succeed())
Expect(registry.Register(context.Background(), node2, true)).To(Succeed())
// Load model on node1
Expect(registry.SetNodeModel(context.Background(), node1.ID, "llama3", 0, "loaded", "", 0)).To(Succeed())
// Verify routing can find the model
nodesWithModel, err := registry.FindNodesWithModel(context.Background(), "llama3")
Expect(err).ToNot(HaveOccurred())
Expect(nodesWithModel).To(HaveLen(1))
Expect(nodesWithModel[0].ID).To(Equal(node1.ID))
})
})
Context("Without --distributed", func() {
It("should fall through to local loading without --distributed", func() {
appCfg := config.NewApplicationConfig()
Expect(appCfg.Distributed.Enabled).To(BeFalse())
// Without distributed mode, no SmartRouter is created.
// The ModelLoader uses its local process management.
// This test documents the design decision.
//
// The bus-URL half of this assertion went with the field it read;
// core/config's "broker surface" spec pins its absence.
})
})
})