Files
LocalAI/core/services/cluster/instance_test.go
T
Ettore Di Giacinto dca8ce7149 feat(cluster): give phase 1 a call site, and prove it against real replicas
Tasks 1 to 5 built an instances table, a splice, both halves of a peer link
and an epoch fence, and nothing in the tree called any of it: no replica
registered, no route was mounted, no sweeper ran. Proving phase 1 end to
end therefore had to start by wiring it.

A frontend in distributed mode now publishes the address its peers dial,
heartbeats it, and sweeps replicas that stopped answering along with the
connection rows they owned, in one pass so the two can never disagree about
who is alive. It serves the peer link and owns the sessions peers dial in,
refusing streams on them until phase 2 installs a relay: a session nobody
accepts on does not fail a peer's Open, it hangs it.

The address is the one peers use, not the one the process binds, and it is
derived from the route to PostgreSQL. That derivation only holds while the
database is remote, so LOCALAI_DISTRIBUTED_ADVERTISE_ADDR sets it
explicitly and a replica that can determine neither warns and keeps
serving rather than failing to start.

Three e2e scenarios run against real local-ai processes, real PostgreSQL
and real dials: replicas publish addresses that can actually be connected
to; a sibling opens a stream over the peer link and is refused without the
cluster token; and a killed replica is reported unreachable, never absent,
loses the claim it held, and takes no worker with it. Each was verified by
mutation: eight injected defects, each failing the scenario that claims to
catch it.

Also moves RegisterClusterRoutes to core/http/routes beside every other
registrar, folds AutoMigrate and the epoch sequence into one
cluster.Migrate, and turns the peer route's auth-coverage spec into a real
assertion: it drives the request through the actual auth middleware
instead of comparing two string constants, which the old spec would have
passed even with the exemption deleted.

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

116 lines
3.9 KiB
Go

package cluster_test
import (
"context"
"net"
"time"
"github.com/mudler/LocalAI/core/services/cluster"
"github.com/mudler/LocalAI/core/services/testutil"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"gorm.io/gorm"
)
var _ = Describe("Instance registry", func() {
var (
db *gorm.DB
reg *cluster.Registry
ctx context.Context
)
BeforeEach(func() {
db = testutil.SetupTestDB()
ctx = context.Background()
Expect(cluster.Migrate(ctx, db)).To(Succeed())
reg = cluster.NewRegistry(db)
})
It("registers an instance and reads it back", func() {
Expect(reg.Register(ctx, "inst-a", "10.0.0.1:8080", "v1")).To(Succeed())
got, err := reg.Get(ctx, "inst-a")
Expect(err).ToNot(HaveOccurred())
Expect(got.AdvertisedAddr).To(Equal("10.0.0.1:8080"))
Expect(got.Version).To(Equal("v1"))
})
It("re-registering the same id updates the address instead of duplicating", func() {
Expect(reg.Register(ctx, "inst-a", "10.0.0.1:8080", "v1")).To(Succeed())
Expect(reg.Register(ctx, "inst-a", "10.0.0.9:9090", "v2")).To(Succeed())
live, err := reg.Live(ctx, time.Hour)
Expect(err).ToNot(HaveOccurred())
Expect(live).To(HaveLen(1))
Expect(live[0].AdvertisedAddr).To(Equal("10.0.0.9:9090"))
})
It("reports a missing instance distinguishably", func() {
_, err := reg.Get(ctx, "nope")
Expect(err).To(MatchError(cluster.ErrInstanceNotFound))
})
It("excludes instances whose heartbeat has aged out", func() {
Expect(reg.Register(ctx, "stale", "10.0.0.1:8080", "v1")).To(Succeed())
// Age the row directly; sleeping in a spec is forbidden.
Expect(db.Model(&cluster.Instance{}).Where("id = ?", "stale").
Update("last_seen", time.Now().Add(-10*time.Minute)).Error).To(Succeed())
live, err := reg.Live(ctx, time.Minute)
Expect(err).ToNot(HaveOccurred())
Expect(live).To(BeEmpty())
})
It("brings a stale instance back with a heartbeat", func() {
Expect(reg.Register(ctx, "revive", "10.0.0.1:8080", "v1")).To(Succeed())
Expect(db.Model(&cluster.Instance{}).Where("id = ?", "revive").
Update("last_seen", time.Now().Add(-10*time.Minute)).Error).To(Succeed())
Expect(reg.Heartbeat(ctx, "revive")).To(Succeed())
live, err := reg.Live(ctx, time.Minute)
Expect(err).ToNot(HaveOccurred())
Expect(live).To(HaveLen(1))
})
It("heartbeating an unknown instance is an error, not a silent insert", func() {
Expect(reg.Heartbeat(ctx, "ghost")).To(MatchError(cluster.ErrInstanceNotFound))
})
})
var _ = Describe("Advertised address discovery", func() {
// The address itself depends on host networking and is deliberately not
// asserted. What is portable is the shape: whatever interface routes to the
// database, the port must be the one the caller asked for, not the
// database's.
It("combines a local interface with the caller's port", func() {
addr, err := cluster.DiscoverAdvertisedAddr("postgres://198.51.100.1:5432/testdb", 8080)
if err != nil {
Skip("no route to a database host on this machine: " + err.Error())
}
host, port, splitErr := net.SplitHostPort(addr)
Expect(splitErr).ToNot(HaveOccurred())
Expect(port).To(Equal("8080"))
Expect(net.ParseIP(host)).ToNot(BeNil())
})
It("refuses a DSN it cannot derive an address from", func() {
_, err := cluster.DiscoverAdvertisedAddr("", 8080)
Expect(err).To(HaveOccurred())
})
// A database on this same host routes over loopback on every platform, so
// this is deterministic rather than host-dependent. Returning 127.0.0.1
// would make a peer dialling this replica reach itself.
It("refuses a loopback route instead of advertising an address peers cannot use", func() {
addr, err := cluster.DiscoverAdvertisedAddr("postgres://user@127.0.0.1:5432/testdb", 8080)
Expect(addr).To(BeEmpty())
Expect(err).To(MatchError(ContainSubstring("loopback")))
})
It("refuses a port that cannot be dialled", func() {
_, err := cluster.DiscoverAdvertisedAddr("postgres://198.51.100.1:5432/testdb", 0)
Expect(err).To(MatchError(ContainSubstring("out of range")))
})
})