Files
LocalAI/tests/e2e/distributed/gallery_distributed_test.go
T
Ettore Di Giacinto e4e5dbd227 chore(distributed): stop telling an operator to run a NATS cluster
Every carrier had already moved and no process opened a bus connection, but
the surface an operator reads still described a deployment with a broker in
it: a compose service, a 220-line credential-generation script, two CI steps
pulling a container nothing started, two flag tables offering --nats-url, an
architecture diagram with a NATS box wired to the workers, a join-command
generator in the Nodes page that emitted --nats-url for agent workers, and a
test suite that stood a NATS server up for specs that no longer used it.

That is the one way this programme could still fail invisibly. Every test
passes, every binary works, and every production deployment goes on running
and paying for infrastructure that carries nothing.

Nothing in this repository starts a NATS server any more. The compose file is
four services, the docs say to shut the broker down and what to keep, and the
e2e suite runs on one PostgreSQL container.

The three LOCALAI_NATS_*_TIMEOUT env vars are KEPT, and are now documented
twice as being kept. They were never broker settings: each names a control-RPC
budget the frontend applies to a worker, still read and still enforced. They
carry the prefix only because they arrived with the bus, and renaming them
would break every existing deployment for cosmetics.

The agent worker's join command was the last surface still emitting the flag,
two tasks after the agent worker stopped dialling. The Playwright spec that
covered it asserted the opposite of what is now true, so it is inverted rather
than deleted, and it reads the rendered command string rather than the
component's variables: the variables are what the fix removes, so a spec
reading them would have stopped compiling instead of failing, and a compile
error is not evidence about what an operator is shown.

nats_jwt_test.go and its helpers are deleted. They pinned a real server
ENFORCING the minted permissions. The CONTENT of those allow lists is still
pinned, untouched, by pkg/natsauth's own suites, including the spec that
refuses to let the agent lists go empty, since an empty allow list in NATS
means unrestricted. The enforcement half is retired rather than moved:
enforcement is a property of a connection, and nothing opens one.

The suite's own NATS container goes with them, which the brief left for the
next task. Removing the pre-pull while BeforeSuite still ran the image would
have defeated the step rather than cleaned it up, and this change removes the
last reader of TestInfra.NC. agent_native_executor_test.go and
mcp_ci_job_test.go are moved onto infra.Bus() instead of deleted: they were
the last two specs building a bridge and a dispatcher on a client nobody uses,
which is exactly the drift TestInfra.Bus's own comment warns about.

cluster.Options.NatsURL is now fed a deliberately dead address rather than a
live container's. Frontends and agent workers still receive LOCALAI_NATS_URL,
because that is the coverage for the promise that an existing command line
still starts; sourcing it from a running server would have let a regression
that actually dialled it pass. The control in cluster_control_test.go keeps
its assertion and loses its explanation, which claimed the deployment had a
bus and no longer could.

One latent spec race surfaced and is fixed: the background-run spec waited for
a COUNT of events and then read a snapshot for the terminal status, which is
the last event of a run and therefore always arrives after the count is met.
Its immediate twin had already been fixed this way. Nothing in production
changed.

pkg/natsauth keeps its files. It is reachable from production only through the
natsauth.Config parameter thread, and that thread is the next task's.

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

178 lines
6.1 KiB
Go

package distributed_test
import (
"sync/atomic"
"github.com/mudler/LocalAI/core/config"
"github.com/mudler/LocalAI/core/services/distributed"
"github.com/mudler/LocalAI/core/services/messaging"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
pgdriver "gorm.io/driver/postgres"
"gorm.io/gorm"
"gorm.io/gorm/logger"
)
var _ = Describe("Gallery Distributed", Label("Distributed"), func() {
var (
infra *TestInfra
db *gorm.DB
galleryStore *distributed.GalleryStore
)
BeforeEach(func() {
infra = SetupInfra("localai_gallery_dist_test")
var err error
db, err = gorm.Open(pgdriver.Open(infra.PGURL), &gorm.Config{
Logger: logger.Default.LogMode(logger.Silent),
})
Expect(err).ToNot(HaveOccurred())
galleryStore, err = distributed.NewGalleryStore(db)
Expect(err).ToNot(HaveOccurred())
})
Context("PostgreSQL gallery operations", func() {
It("should write gallery operation status to PostgreSQL", func() {
op := &distributed.GalleryOperationRecord{
GalleryElementName: "llama3-8b",
OpType: "model_install",
Status: "downloading",
Cancellable: true,
FrontendID: "f1",
}
Expect(galleryStore.Create(op)).To(Succeed())
Expect(op.ID).ToNot(BeEmpty())
retrieved, err := galleryStore.Get(op.ID)
Expect(err).ToNot(HaveOccurred())
Expect(retrieved.GalleryElementName).To(Equal("llama3-8b"))
Expect(retrieved.Status).To(Equal("downloading"))
Expect(retrieved.FrontendID).To(Equal("f1"))
// Update progress (cancellable: a downloading install can be cancelled)
Expect(galleryStore.UpdateProgress(op.ID, 0.75, "75% complete", "6GB", true)).To(Succeed())
updated, _ := galleryStore.Get(op.ID)
Expect(updated.Progress).To(BeNumerically("~", 0.75, 0.01))
Expect(updated.Message).To(Equal("75% complete"))
Expect(updated.Cancellable).To(BeTrue())
// Complete
Expect(galleryStore.UpdateStatus(op.ID, "completed", "")).To(Succeed())
completed, _ := galleryStore.Get(op.ID)
Expect(completed.Status).To(Equal("completed"))
})
})
// The gallery families ride the broadcast carrier, not NATS.
//
// These used to publish and subscribe on a message-bus client, which
// asserted that the bus delivers to itself and nothing about this
// deployment: they would have
// stayed green through the whole migration while the gallery service had
// already moved. Two carriers on the deployment's own database is the shape
// a fleet has, and it is the shape that fails when one end moves and the
// other does not.
//
// No flush, unlike the NATS version: pgbus.Subscribe has already issued its
// LISTEN by the time it returns, and Subscribers() counts only live
// handlers, so there is no window to wait out.
Context("gallery progress on the broadcast carrier", func() {
It("delivers a peer replica's progress updates", func() {
op := &distributed.GalleryOperationRecord{
GalleryElementName: "whisper-large",
OpType: "model_install",
Status: "downloading",
}
Expect(galleryStore.Create(op)).To(Succeed())
publisher, subscriber := infra.Bus(), infra.Bus()
var received atomic.Int32
sub, err := subscriber.Subscribe(messaging.SubjectGalleryProgress(op.ID), func([]byte) {
received.Add(1)
})
Expect(err).ToNot(HaveOccurred())
defer func() { Expect(sub.Unsubscribe()).To(Succeed()) }()
Expect(publisher.Publish(messaging.SubjectGalleryProgress(op.ID), map[string]any{
"op_id": op.ID, "progress": 0.25, "message": "25%",
})).To(Succeed())
Expect(publisher.Publish(messaging.SubjectGalleryProgress(op.ID), map[string]any{
"op_id": op.ID, "progress": 0.50, "message": "50%",
})).To(Succeed())
Eventually(func() int32 { return received.Load() }, "20s").Should(Equal(int32(2)))
})
})
Context("gallery cancel on the broadcast carrier", func() {
It("delivers a cancel to the replica holding the operation", func() {
op := &distributed.GalleryOperationRecord{
GalleryElementName: "cancel-model",
OpType: "model_install",
Status: "downloading",
Cancellable: true,
}
Expect(galleryStore.Create(op)).To(Succeed())
publisher, subscriber := infra.Bus(), infra.Bus()
var cancelReceived atomic.Bool
sub, err := subscriber.Subscribe(messaging.SubjectGalleryCancel(op.ID), func([]byte) {
cancelReceived.Store(true)
})
Expect(err).ToNot(HaveOccurred())
defer func() { Expect(sub.Unsubscribe()).To(Succeed()) }()
Expect(publisher.Publish(messaging.SubjectGalleryCancel(op.ID), map[string]string{
"op_id": op.ID,
})).To(Succeed())
Eventually(func() bool { return cancelReceived.Load() }, "20s").Should(BeTrue())
// The row is what survives a replica that was not listening. The
// broadcast is the hint to go and look at it.
Expect(galleryStore.Cancel(op.ID)).To(Succeed())
updated, _ := galleryStore.Get(op.ID)
Expect(updated.Status).To(Equal("cancelled"))
})
})
Context("Deduplication", func() {
It("should deduplicate concurrent downloads of same model", func() {
op := &distributed.GalleryOperationRecord{
GalleryElementName: "same-model-v2",
OpType: "model_install",
Status: "downloading",
}
Expect(galleryStore.Create(op)).To(Succeed())
// Another instance tries to download the same model
dup, err := galleryStore.FindDuplicate("same-model-v2")
Expect(err).ToNot(HaveOccurred())
Expect(dup.ID).To(Equal(op.ID))
// Completed operations should not be considered duplicates
Expect(galleryStore.UpdateStatus(op.ID, "completed", "")).To(Succeed())
_, err = galleryStore.FindDuplicate("same-model-v2")
Expect(err).To(HaveOccurred()) // no active duplicate
})
})
Context("Without --distributed", func() {
It("should use in-memory map without --distributed", func() {
appCfg := config.NewApplicationConfig()
Expect(appCfg.Distributed.Enabled).To(BeFalse())
// Without distributed mode, gallery operations use the existing
// in-memory galleryApplier map. No PostgreSQL needed.
Expect(appCfg.Distributed.NatsURL).To(BeEmpty())
})
})
})