mirror of
https://github.com/mudler/LocalAI.git
synced 2026-09-20 05:07:07 -04:00
agent.<name>.cancel was the last family on a message bus, and the only reason an agent worker dialled one. Its subscriber is the worker running the execution, and a worker has no database, so the family could not move to the PostgreSQL fan-out carrier: a cancel published there would reach no worker while reporting that it had been sent. It is a control verb now. An agent worker mounts workerctl.PathAgentCancel on the loopback control plane behind its tunnel and applies the cancel to the same registry the executor registers a run on. The frontend issues it through nodes.AgentControlClient.CancelAgentRun. That call is a FAN-OUT and not a pick, because nothing records which worker holds a given execution: the claim row names the claiming replica, and it is deleted when the run ends. Every agent worker a live replica can reach is asked over its own tunnel, relayed by the peer mesh when a peer holds it, and each worker answers only for itself. The answers stay apart, which is why this family was held back. A cancel a worker made is nil. A cancel some worker could not be asked is ErrAgentCancelUndelivered, which is neither a refusal nor a missing run. A cancel every reachable worker declined to own is ErrAgentRunNotOnAnyWorker. A deployment with no agent worker is ErrNoAgentWorker. Neither new sentinel wraps ErrWorkerUnroutable and neither is a worker answer, so nothing is reaped, demoted or evicted because of a cancel. A worker in the ABSENT CONNECTION condition, one whose tunnel was lost inside the reconnect grace, counts as undelivered. It is not retried in the call and not queued: a retry would spend a budget the caller did not choose, and a queue would need durable state whose only consumer is a run whose control stream went with the tunnel. A worker whose departure has outlived the grace is the one routing fact a caller may act on and is excluded, or a single retired agent node would make every cancel undelivered for ever. The fan-out reads a different node set from the pick. A draining worker takes no new work but is still finishing what it holds, so it is offered the cancel; a pending one is refused by the tunnel route on every dial and is not. With that, nothing in LocalAI connects to NATS. The agent worker's dial, its credential ladder and its refresh loop are gone, and so is the frontend's cancel carrier. LOCALAI_NATS_URL is accepted and ignored everywhere, and distributed mode no longer requires it. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
176 lines
7.1 KiB
Go
176 lines
7.1 KiB
Go
// SPDX-License-Identifier: MIT
|
|
|
|
package application
|
|
|
|
import (
|
|
"context"
|
|
|
|
. "github.com/onsi/ginkgo/v2"
|
|
. "github.com/onsi/gomega"
|
|
|
|
"github.com/mudler/LocalAI/core/config"
|
|
"github.com/mudler/LocalAI/core/services/messaging"
|
|
"github.com/mudler/LocalAI/core/services/pgbus"
|
|
"github.com/mudler/LocalAI/core/services/syncstate"
|
|
"github.com/mudler/LocalAI/core/services/testutil"
|
|
)
|
|
|
|
// The guard on the one setting that decides whether any broadcast in the
|
|
// deployment is ever delivered.
|
|
//
|
|
// The carrier holds a pinned LISTEN connection opened from a DSN, and publishes
|
|
// travel on a pooled handle opened from another. When those two name different
|
|
// databases every publish succeeds, every subscribe succeeds, and nothing
|
|
// arrives, on every replica, with no error anywhere. There is exactly one
|
|
// legitimate DSN, and these specs are what say so in a way that fails when it
|
|
// stops being true.
|
|
var _ = Describe("opening the deployment's broadcast carrier", func() {
|
|
It("listens on the same database URL the auth pool was built from", func() {
|
|
db, dsn := testutil.SetupTestDBWithDSN()
|
|
cfg := &config.ApplicationConfig{}
|
|
cfg.Auth.DatabaseURL = dsn
|
|
|
|
bus, err := newBroadcastBus(context.Background(), cfg, db)
|
|
Expect(err).ToNot(HaveOccurred())
|
|
DeferCleanup(bus.Close)
|
|
|
|
// Equality with the field, not "is a PostgreSQL URL": the failure being
|
|
// excluded is two databases, and any DSN passes a shape check.
|
|
Expect(bus.DSN()).To(Equal(cfg.Auth.DatabaseURL))
|
|
})
|
|
|
|
It("migrates the spill table, so an oversized broadcast has somewhere to go", func() {
|
|
db, dsn := testutil.SetupTestDBWithDSN()
|
|
cfg := &config.ApplicationConfig{}
|
|
cfg.Auth.DatabaseURL = dsn
|
|
|
|
bus, err := newBroadcastBus(context.Background(), cfg, db)
|
|
Expect(err).ToNot(HaveOccurred())
|
|
DeferCleanup(bus.Close)
|
|
|
|
Expect(db.Migrator().HasTable(&pgbus.BusMessage{})).To(BeTrue())
|
|
})
|
|
|
|
It("refuses to open a carrier whose DSN is not the pool's database", func() {
|
|
db, _ := testutil.SetupTestDBWithDSN()
|
|
_, otherDSN := testutil.SetupTestDBWithDSN()
|
|
cfg := &config.ApplicationConfig{}
|
|
cfg.Auth.DatabaseURL = otherDSN
|
|
|
|
_, err := newBroadcastBus(context.Background(), cfg, db)
|
|
|
|
Expect(err).To(HaveOccurred())
|
|
})
|
|
})
|
|
|
|
// The partial pin on two wiring lines that cannot be reddened by a spec: the
|
|
// newBroadcastBus call, and `Bus: bus` in the returned literal. Neither is a
|
|
// compile error when deleted and initDistributed cannot be unit tested while it
|
|
// opens NATS first, so what is available is a boot refusal, and this is what
|
|
// keeps that refusal honest.
|
|
var _ = Describe("refusing a deployment with no broadcast carrier", func() {
|
|
It("accepts services that carry one", func() {
|
|
db, dsn := testutil.SetupTestDBWithDSN()
|
|
cfg := &config.ApplicationConfig{}
|
|
cfg.Auth.DatabaseURL = dsn
|
|
bus, err := newBroadcastBus(context.Background(), cfg, db)
|
|
Expect(err).ToNot(HaveOccurred())
|
|
DeferCleanup(bus.Close)
|
|
|
|
Expect(requireBroadcastCarrier(&DistributedServices{Bus: bus})).To(Succeed())
|
|
})
|
|
|
|
It("refuses services whose carrier was never assigned, and says what it costs", func() {
|
|
err := requireBroadcastCarrier(&DistributedServices{})
|
|
|
|
Expect(err).To(HaveOccurred())
|
|
Expect(err.Error()).To(ContainSubstring("published between replicas"))
|
|
Expect(err.Error()).To(ContainSubstring("shutdown"))
|
|
})
|
|
|
|
It("refuses a nil deployment rather than dereferencing it", func() {
|
|
Expect(requireBroadcastCarrier(nil)).ToNot(Succeed())
|
|
})
|
|
})
|
|
|
|
var _ = Describe("shutting the distributed services down", func() {
|
|
It("closes the broadcast carrier", func() {
|
|
// A pinned PostgreSQL session and the goroutine parked on it, per
|
|
// replica restart. Nothing else in this process ever closes it, so the
|
|
// line in the shutdown closure is the whole lifecycle.
|
|
db, dsn := testutil.SetupTestDBWithDSN()
|
|
cfg := &config.ApplicationConfig{}
|
|
cfg.Auth.DatabaseURL = dsn
|
|
bus, err := newBroadcastBus(context.Background(), cfg, db)
|
|
Expect(err).ToNot(HaveOccurred())
|
|
Expect(bus.IsConnected()).To(BeTrue())
|
|
|
|
(&DistributedServices{Bus: bus}).Shutdown()
|
|
|
|
Expect(bus.IsConnected()).To(BeFalse())
|
|
})
|
|
})
|
|
|
|
// The one place the four state.*.delta families are told which carrier they
|
|
// travel on.
|
|
//
|
|
// It was five field reads before this: the fine-tune service, the quantization
|
|
// service, the agent-task setter on two startup paths, the per-user services
|
|
// manager and the Open Responses store. Every one of them takes a
|
|
// messaging.Broadcaster, which *messaging.Client satisfies too, so a site left
|
|
// holding the struct's NATS field compiled, started, published and was
|
|
// delivered onto a carrier only agent workers read, and nothing failed until
|
|
// NATS did. Collapsing the choice into one function is what makes it a fact
|
|
// these specs can hold.
|
|
var _ = Describe("handing the broadcast carrier to its adopters", func() {
|
|
It("returns the carrier the deployment opened", func() {
|
|
db, dsn := testutil.SetupTestDBWithDSN()
|
|
cfg := &config.ApplicationConfig{}
|
|
cfg.Auth.DatabaseURL = dsn
|
|
bus, err := newBroadcastBus(context.Background(), cfg, db)
|
|
Expect(err).ToNot(HaveOccurred())
|
|
DeferCleanup(bus.Close)
|
|
|
|
// Identity and not "is a Broadcaster". There is no second carrier on
|
|
// this struct any more: the family that needed one, agent.<name>.cancel,
|
|
// rides the agent worker's own tunnel now. The identity assertion stays
|
|
// because what it pins is that adopters get THIS bus rather than
|
|
// anything else that satisfies the interface.
|
|
ds := &DistributedServices{Bus: bus}
|
|
|
|
Expect(ds.Broadcast()).To(BeIdenticalTo(messaging.Broadcaster(bus)))
|
|
})
|
|
|
|
It("returns an interface that reads as absent, not a typed nil, when there is no carrier", func() {
|
|
// Every adopter branches on `bus == nil` to mean standalone. A nil
|
|
// *pgbus.Bus placed in an interface is NOT nil, so that branch would be
|
|
// skipped and the first Set would panic on a request rather than at
|
|
// boot.
|
|
//
|
|
// Compared with == and not with BeNil(). Gomega's BeNil reports a nil
|
|
// POINTER inside an interface as nil, so it passes on exactly the value
|
|
// this spec exists to reject; the first draft of this spec did, and the
|
|
// mutation that removed the guard stayed green.
|
|
var ds *DistributedServices
|
|
Expect(ds.Broadcast() == nil).To(BeTrue(), "a nil deployment must yield an interface that is itself nil")
|
|
Expect((&DistributedServices{}).Broadcast() == nil).To(BeTrue(),
|
|
"a deployment with no carrier must yield an interface that is itself nil, not one wrapping a nil *pgbus.Bus")
|
|
})
|
|
|
|
It("gives an adopter a carrier-less map rather than one that panics on the first write", func() {
|
|
// The consequence, driven through the component every adopter builds.
|
|
// A typed nil satisfies `!= nil`, so Start subscribes on it and Set
|
|
// publishes on it, and both dereference a nil *pgbus.Bus on a request
|
|
// path rather than at boot.
|
|
m := syncstate.New(syncstate.Config[string, string]{
|
|
Name: "test.jobs",
|
|
Key: func(v string) string { return v },
|
|
Bus: (&DistributedServices{}).Broadcast(),
|
|
})
|
|
Expect(m.Start(context.Background())).To(Succeed())
|
|
DeferCleanup(func() { Expect(m.Close()).To(Succeed()) })
|
|
|
|
Expect(func() { Expect(m.Set(context.Background(), "v")).To(Succeed()) }).ToNot(Panic())
|
|
})
|
|
})
|