From 217038d4fed813ce6ce6042238216dde23d158d3 Mon Sep 17 00:00:00 2001 From: localai-org-maint-bot Date: Tue, 6 Oct 2026 22:45:23 +0200 Subject: [PATCH] fix(agents): share the collection list across frontends (#12536) With the PostgreSQL vector engine, each frontend kept the collections it had opened in an in-memory map and answered "collection not found" for any other. A collection created through one frontend was unknown to the others until they restarted, and each frontend listed a different set. Wrap the in-process collections backend in distributed mode so that the database is the source of truth: - lists come from the registry, - a lookup miss checks the registry before it returns 404, and opens a collection that exists there (once per name, even under concurrency), - a cached collection that left the registry is dropped and closed, with a re-check at most every 5 seconds, - create and reset publish an event on the existing collection invalidation subject, so other frontends re-check at once. Without the postgres engine, or outside distributed mode, nothing changes. Assisted-by: Claude:claude-sonnet-5-5 go Signed-off-by: Ettore Di Giacinto Co-authored-by: Ettore Di Giacinto --- core/application/application.go | 3 + core/services/agentpool/agent_pool.go | 27 +- core/services/agentpool/collections_shared.go | 333 ++++++++++++++++++ .../agentpool/collections_shared_test.go | 295 ++++++++++++++++ core/services/messaging/subjects.go | 4 + docs/content/features/distributed-mode.md | 11 + go.mod | 2 +- go.sum | 2 + 8 files changed, 674 insertions(+), 3 deletions(-) create mode 100644 core/services/agentpool/collections_shared.go create mode 100644 core/services/agentpool/collections_shared_test.go diff --git a/core/application/application.go b/core/application/application.go index 187bb2ab4..ce333a5a0 100644 --- a/core/application/application.go +++ b/core/application/application.go @@ -667,6 +667,9 @@ func (a *Application) agentPoolOptions() agentpool.AgentPoolOptions { opts.WorkQueue = d.WorkQueue opts.EventBridge = d.AgentBridge opts.AgentStore = d.AgentStore + if d.Nats != nil { + opts.Bus = d.Nats + } } return opts } diff --git a/core/services/agentpool/agent_pool.go b/core/services/agentpool/agent_pool.go index 46d1da2f2..c139dda41 100644 --- a/core/services/agentpool/agent_pool.go +++ b/core/services/agentpool/agent_pool.go @@ -73,6 +73,7 @@ type distributedBridge struct { agentStore *agents.AgentStore // PostgreSQL agent config store eventBridge AgentEventBridge // Event bridge for SSE + persistence skillStore *distributed.SkillStore // PostgreSQL skill metadata (distributed mode) + bus messaging.Broadcaster // Fan-out bus for collection events (distributed mode) } // userManager handles per-user services, storage, and auth. @@ -125,6 +126,8 @@ type AgentPoolOptions struct { WorkQueue messaging.WorkQueue EventBridge AgentEventBridge AgentStore *agents.AgentStore + // Bus carries collection events between frontends (nil in standalone). + Bus messaging.Broadcaster } func NewAgentPoolService(appConfig *config.ApplicationConfig, opts ...AgentPoolOptions) (*AgentPoolService, error) { @@ -148,6 +151,9 @@ func NewAgentPoolService(appConfig *config.ApplicationConfig, opts ...AgentPoolO if o.AgentStore != nil { svc.distributed.agentStore = o.AgentStore } + if o.Bus != nil { + svc.distributed.bus = o.Bus + } } return svc, nil } @@ -233,8 +239,9 @@ func (s *AgentPoolService) startDistributed(ctx context.Context, apiURL, apiKey } fileAssets := filepath.Join(stateDir, "assets") - collectionsBackend, _ := collections.NewInProcessBackend(s.buildCollectionsConfig(apiURL, apiKey, collectionDBPath, fileAssets)) - s.collectionsBackend = collectionsBackend + collectionsCfg := s.buildCollectionsConfig(apiURL, apiKey, collectionDBPath, fileAssets) + collectionsBackend, collectionsState := collections.NewInProcessBackend(collectionsCfg) + s.collectionsBackend = s.shareCollections(collectionsBackend, collectionsState, collectionsCfg) // User-scoped storage dataDir := cmp.Or(s.appConfig.DataPath, s.appConfig.DynamicConfigsDir) @@ -1201,3 +1208,19 @@ func (s *AgentPoolService) buildSkillProvider() agents.SkillContentProvider { return s.loadSkillsForUser(userID) } } + +// shareCollections makes the frontends agree on the list of collections. +// +// The in-process backend keeps the collections a frontend opened in memory. +// With several frontends that cache is per replica, so a collection created +// through one frontend would be unknown to the others. With the PostgreSQL +// vector engine the database is the shared registry, and the backend is +// wrapped so that it asks the database. With any other engine there is no +// shared registry, and the backend is returned unchanged. +func (s *AgentPoolService) shareCollections(backend collections.Backend, state *collections.State, cfg *collections.Config) collections.Backend { + if cfg.VectorEngine != "postgres" || cfg.DatabaseURL == "" { + xlog.Warn("Distributed mode without the postgres vector engine: each frontend keeps its own list of collections", "vectorEngine", cfg.VectorEngine) + return backend + } + return newSharedCollections(backend, state, postgresCollectionRegistry{databaseURL: cfg.DatabaseURL}, s.distributed.bus) +} diff --git a/core/services/agentpool/collections_shared.go b/core/services/agentpool/collections_shared.go new file mode 100644 index 000000000..9ac232762 --- /dev/null +++ b/core/services/agentpool/collections_shared.go @@ -0,0 +1,333 @@ +package agentpool + +import ( + "context" + "fmt" + "io" + "sync" + "time" + + "github.com/mudler/LocalAGI/webui/collections" + "github.com/mudler/LocalAI/core/services/messaging" + "github.com/mudler/localrecall/rag/engine" + "github.com/mudler/xlog" + "golang.org/x/sync/singleflight" +) + +// Shared mode: with the PostgreSQL vector engine, the database is the source +// of truth for which collections exist. The in-process backend keeps the +// collections a frontend has opened in a map, which is only a cache. A +// collection created through one frontend is absent from that map on every +// other frontend, so without this layer a lookup on another frontend answers +// "collection not found" until that frontend restarts. +// +// sharedCollections wraps the in-process backend and: +// - answers lists from the shared registry, +// - re-checks the registry before it reports a collection as not found, and +// opens a collection that exists there but not in the local cache, +// - drops a cached collection once the registry no longer has it, +// - publishes create/reset events so that the other frontends refresh at once +// instead of waiting for the next re-check. + +const ( + // collectionRecheckInterval bounds how long a frontend keeps serving a + // cached collection that another frontend removed from the registry. + collectionRecheckInterval = 5 * time.Second + // collectionRegistryTimeout bounds each query to the shared registry. + collectionRegistryTimeout = 5 * time.Second + + collectionOpCreate = "create" + collectionOpReset = "reset" +) + +// collectionRegistry is the shared record of which collections exist. +type collectionRegistry interface { + Exists(ctx context.Context, name string) (bool, error) + List(ctx context.Context) ([]string, error) +} + +// postgresCollectionRegistry reads the collection registry of the PostgreSQL +// vector engine. +type postgresCollectionRegistry struct{ databaseURL string } + +func (r postgresCollectionRegistry) Exists(ctx context.Context, name string) (bool, error) { + return engine.PostgresCollectionExists(ctx, r.databaseURL, name) +} + +func (r postgresCollectionRegistry) List(ctx context.Context) ([]string, error) { + return engine.ListPostgresCollections(ctx, r.databaseURL) +} + +// collectionEvent is the payload published on +// messaging.SubjectCacheInvalidateCollection. +type collectionEvent struct { + Name string `json:"name"` + Op string `json:"op"` +} + +// sharedCollections implements collections.Backend on top of the in-process +// backend and the shared registry. +type sharedCollections struct { + inner collections.Backend + state *collections.State + registry collectionRegistry + bus messaging.Broadcaster // nil without a message bus + recheck time.Duration + + opens singleflight.Group + + mu sync.Mutex + checked map[string]time.Time // last time the registry confirmed a name + sub messaging.Subscription +} + +var _ collections.Backend = (*sharedCollections)(nil) + +// newSharedCollections wraps inner. bus may be nil: the frontends then +// converge through the registry alone, at the latest after the re-check +// interval. +func newSharedCollections(inner collections.Backend, state *collections.State, registry collectionRegistry, bus messaging.Broadcaster) *sharedCollections { + s := &sharedCollections{ + inner: inner, + state: state, + registry: registry, + bus: bus, + recheck: collectionRecheckInterval, + checked: map[string]time.Time{}, + } + if bus != nil { + sub, err := messaging.SubscribeJSON(bus, messaging.SubjectCacheInvalidateCollectionAll, s.onEvent) + if err != nil { + xlog.Warn("Failed to subscribe to collection events; relying on registry re-checks", "error", err) + } else { + s.sub = sub + } + } + return s +} + +// Close stops listening for collection events. +func (s *sharedCollections) Close() { + if s.sub != nil { + _ = s.sub.Unsubscribe() + } +} + +// onEvent handles an event from any frontend, including this one. It never +// publishes, so it cannot loop. +func (s *sharedCollections) onEvent(evt collectionEvent) { + if evt.Name == "" { + return + } + s.forget(evt.Name) + if evt.Op == collectionOpCreate { + // Nothing to drop. The next lookup opens it from the registry. + return + } + s.reconcile(evt.Name) +} + +func (s *sharedCollections) publish(name, op string) { + if s.bus == nil { + return + } + if err := s.bus.Publish(messaging.SubjectCacheInvalidateCollection(name), collectionEvent{Name: name, Op: op}); err != nil { + xlog.Warn("Failed to publish collection event", "collection", name, "op", op, "error", err) + } +} + +func (s *sharedCollections) fresh(name string) bool { + s.mu.Lock() + defer s.mu.Unlock() + t, ok := s.checked[name] + return ok && time.Since(t) < s.recheck +} + +func (s *sharedCollections) markChecked(name string) { + s.mu.Lock() + s.checked[name] = time.Now() + s.mu.Unlock() +} + +func (s *sharedCollections) forget(name string) { + s.mu.Lock() + delete(s.checked, name) + s.mu.Unlock() +} + +func (s *sharedCollections) cached(name string) bool { + s.state.Mu.RLock() + defer s.state.Mu.RUnlock() + kb, ok := s.state.Collections[name] + return ok && kb != nil +} + +// drop forgets a collection that the registry no longer has and releases its +// resources. +func (s *sharedCollections) drop(name string) { + s.forget(name) + s.state.Mu.Lock() + kb, ok := s.state.Collections[name] + delete(s.state.Collections, name) + s.state.Mu.Unlock() + if !ok { + return + } + s.state.SourceManager.UnregisterCollection(name) + if kb != nil { + kb.Close() + } + xlog.Info("Dropped a collection that no longer exists in the shared registry", "collection", name) +} + +// reconcile drops the cached entry when the registry no longer has the name. +// A registry error keeps the entry: a database blip must not become a 404. +func (s *sharedCollections) reconcile(name string) { + ctx, cancel := context.WithTimeout(context.Background(), collectionRegistryTimeout) + defer cancel() + exists, err := s.registry.Exists(ctx, name) + if err != nil { + xlog.Warn("Could not check the collection registry", "collection", name, "error", err) + return + } + if !exists { + s.drop(name) + } +} + +// ensure makes the local cache agree with the registry for one collection. It +// returns a "collection not found" error when the registry does not have it. +func (s *sharedCollections) ensure(name string) error { + if s.fresh(name) { + return nil + } + ctx, cancel := context.WithTimeout(context.Background(), collectionRegistryTimeout) + defer cancel() + exists, err := s.registry.Exists(ctx, name) + if err != nil { + // Keep serving what this frontend already has. + xlog.Warn("Could not check the collection registry", "collection", name, "error", err) + return nil + } + if !exists { + s.drop(name) + return fmt.Errorf("collection not found: %s", name) + } + if !s.cached(name) { + // Concurrent first lookups open the collection (and its connection + // pool) once. + _, err, _ := s.opens.Do(name, func() (any, error) { + if s.cached(name) { + return nil, nil + } + if _, ok := s.state.EnsureCollection(name); !ok { + return nil, fmt.Errorf("failed to open collection %s from the shared registry", name) + } + return nil, nil + }) + if err != nil { + return err + } + } + s.markChecked(name) + return nil +} + +func (s *sharedCollections) ListCollections() ([]string, error) { + ctx, cancel := context.WithTimeout(context.Background(), collectionRegistryTimeout) + defer cancel() + return s.registry.List(ctx) +} + +func (s *sharedCollections) CreateCollection(name string) error { + if err := s.inner.CreateCollection(name); err != nil { + return err + } + s.markChecked(name) + s.publish(name, collectionOpCreate) + return nil +} + +func (s *sharedCollections) Reset(collection string) error { + if err := s.ensure(collection); err != nil { + return err + } + if err := s.inner.Reset(collection); err != nil { + return err + } + // The inner backend dropped its entry. The registry decides whether the + // collection still exists, so the next lookup opens it again. + s.forget(collection) + s.publish(collection, collectionOpReset) + return nil +} + +func (s *sharedCollections) Upload(collection, filename string, fileBody io.Reader) (string, error) { + if err := s.ensure(collection); err != nil { + return "", err + } + return s.inner.Upload(collection, filename, fileBody) +} + +func (s *sharedCollections) ListEntries(collection string) ([]string, error) { + if err := s.ensure(collection); err != nil { + return nil, err + } + return s.inner.ListEntries(collection) +} + +func (s *sharedCollections) GetEntryContent(collection, entry string) (string, int, error) { + if err := s.ensure(collection); err != nil { + return "", 0, err + } + return s.inner.GetEntryContent(collection, entry) +} + +func (s *sharedCollections) Search(collection, query string, maxResults int) ([]collections.SearchResult, error) { + if err := s.ensure(collection); err != nil { + return nil, err + } + return s.inner.Search(collection, query, maxResults) +} + +func (s *sharedCollections) DeleteEntry(collection, entry string) ([]string, error) { + if err := s.ensure(collection); err != nil { + return nil, err + } + return s.inner.DeleteEntry(collection, entry) +} + +func (s *sharedCollections) AddSource(collection, url string, intervalMin int) error { + if err := s.ensure(collection); err != nil { + return err + } + return s.inner.AddSource(collection, url, intervalMin) +} + +func (s *sharedCollections) RemoveSource(collection, url string) error { + if err := s.ensure(collection); err != nil { + return err + } + return s.inner.RemoveSource(collection, url) +} + +func (s *sharedCollections) ListSources(collection string) ([]collections.SourceInfo, error) { + if err := s.ensure(collection); err != nil { + return nil, err + } + return s.inner.ListSources(collection) +} + +func (s *sharedCollections) EntryExists(collection, entry string) bool { + if err := s.ensure(collection); err != nil { + return false + } + return s.inner.EntryExists(collection, entry) +} + +func (s *sharedCollections) GetEntryFilePath(collection, entry string) (string, error) { + if err := s.ensure(collection); err != nil { + return "", err + } + return s.inner.GetEntryFilePath(collection, entry) +} diff --git a/core/services/agentpool/collections_shared_test.go b/core/services/agentpool/collections_shared_test.go new file mode 100644 index 000000000..0b3e3df7f --- /dev/null +++ b/core/services/agentpool/collections_shared_test.go @@ -0,0 +1,295 @@ +package agentpool + +// White-box tests: two frontends share one registry (the database) and one bus, +// but each has its own in-memory collections cache, like two replicas behind +// one Service. + +import ( + "context" + "errors" + "fmt" + "io" + "sort" + "sync" + "sync/atomic" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + + "github.com/mudler/LocalAGI/webui/collections" + "github.com/mudler/LocalAI/core/services/testutil" + "github.com/mudler/localrecall/rag" + "github.com/mudler/localrecall/rag/sources" +) + +// fakeRegistry stands in for the collection_config table. +type fakeRegistry struct { + mu sync.Mutex + names map[string]bool + err error +} + +func newFakeRegistry() *fakeRegistry { return &fakeRegistry{names: map[string]bool{}} } + +func (r *fakeRegistry) add(n string) { r.mu.Lock(); r.names[n] = true; r.mu.Unlock() } +func (r *fakeRegistry) remove(n string) { r.mu.Lock(); delete(r.names, n); r.mu.Unlock() } +func (r *fakeRegistry) setErr(e error) { r.mu.Lock(); r.err = e; r.mu.Unlock() } + +func (r *fakeRegistry) Exists(_ context.Context, n string) (bool, error) { + r.mu.Lock() + defer r.mu.Unlock() + return r.names[n], r.err +} + +func (r *fakeRegistry) List(_ context.Context) ([]string, error) { + r.mu.Lock() + defer r.mu.Unlock() + if r.err != nil { + return nil, r.err + } + out := []string{} + for n := range r.names { + out = append(out, n) + } + sort.Strings(out) + return out, nil +} + +// fakeInner models the in-process backend: it only knows the collections in +// its own cache and answers "collection not found" for the rest, which is the +// bug under test. +type fakeInner struct { + state *collections.State + registry *fakeRegistry + opens atomic.Int32 + uploads []string + // resetDeletes models a store where a reset removes the registry row. + resetDeletes bool +} + +func newFakeInner(reg *fakeRegistry) *fakeInner { + f := &fakeInner{ + registry: reg, + state: &collections.State{ + Collections: collections.CollectionList{}, + SourceManager: rag.NewSourceManager(&sources.Config{}), + }, + } + f.state.EnsureCollection = func(name string) (*rag.PersistentKB, bool) { + f.opens.Add(1) + f.state.Mu.Lock() + defer f.state.Mu.Unlock() + if kb := f.state.Collections[name]; kb != nil { + return kb, true + } + kb := &rag.PersistentKB{} + f.state.Collections[name] = kb + return kb, true + } + return f +} + +func (f *fakeInner) known(name string) error { + f.state.Mu.RLock() + defer f.state.Mu.RUnlock() + if kb := f.state.Collections[name]; kb == nil { + return fmt.Errorf("collection not found: %s", name) + } + return nil +} + +func (f *fakeInner) ListCollections() ([]string, error) { + f.state.Mu.RLock() + defer f.state.Mu.RUnlock() + out := []string{} + for n := range f.state.Collections { + out = append(out, n) + } + return out, nil +} + +func (f *fakeInner) CreateCollection(name string) error { + f.registry.add(name) + f.state.EnsureCollection(name) + return nil +} + +func (f *fakeInner) Upload(c, name string, _ io.Reader) (string, error) { + if err := f.known(c); err != nil { + return "", err + } + f.uploads = append(f.uploads, name) + return name, nil +} + +func (f *fakeInner) ListEntries(c string) ([]string, error) { + if err := f.known(c); err != nil { + return nil, err + } + return f.uploads, nil +} + +func (f *fakeInner) GetEntryContent(c, _ string) (string, int, error) { + return "", 0, f.known(c) +} + +func (f *fakeInner) Search(c, _ string, _ int) ([]collections.SearchResult, error) { + return nil, f.known(c) +} + +func (f *fakeInner) Reset(c string) error { + if err := f.known(c); err != nil { + return err + } + f.state.Mu.Lock() + delete(f.state.Collections, c) + f.state.Mu.Unlock() + if f.resetDeletes { + f.registry.remove(c) + } + return nil +} + +func (f *fakeInner) DeleteEntry(c, _ string) ([]string, error) { return nil, f.known(c) } +func (f *fakeInner) AddSource(c, _ string, _ int) error { return f.known(c) } +func (f *fakeInner) RemoveSource(c, _ string) error { return f.known(c) } +func (f *fakeInner) ListSources(c string) ([]collections.SourceInfo, error) { + return nil, f.known(c) +} +func (f *fakeInner) EntryExists(c, _ string) bool { return f.known(c) == nil } +func (f *fakeInner) GetEntryFilePath(c, _ string) (string, error) { + return "", f.known(c) +} + +var _ = Describe("sharedCollections across two frontends", func() { + var ( + reg *fakeRegistry + bus *testutil.FakeBus + innerA *fakeInner + innerB *fakeInner + a, b *sharedCollections + ) + + BeforeEach(func() { + reg = newFakeRegistry() + bus = testutil.NewFakeBus() + innerA, innerB = newFakeInner(reg), newFakeInner(reg) + a = newSharedCollections(innerA, innerA.state, reg, bus) + b = newSharedCollections(innerB, innerB.state, reg, bus) + }) + + AfterEach(func() { a.Close(); b.Close() }) + + It("fails without the shared layer: the bug", func() { + Expect(innerA.CreateCollection("col-1")).To(Succeed()) + _, err := innerB.Upload("col-1", "f.txt", nil) + Expect(err).To(MatchError(ContainSubstring("collection not found"))) + }) + + It("serves a collection created on A through B", func() { + Expect(a.CreateCollection("col-1")).To(Succeed()) + + _, err := b.Upload("col-1", "f.txt", nil) + Expect(err).NotTo(HaveOccurred()) + _, err = b.Search("col-1", "q", 3) + Expect(err).NotTo(HaveOccurred()) + _, err = b.ListEntries("col-1") + Expect(err).NotTo(HaveOccurred()) + Expect(b.EntryExists("col-1", "f.txt")).To(BeTrue()) + Expect(b.Reset("col-1")).To(Succeed()) + }) + + It("lists the same collections on both frontends", func() { + Expect(a.CreateCollection("col-a")).To(Succeed()) + Expect(b.CreateCollection("col-b")).To(Succeed()) + + la, err := a.ListCollections() + Expect(err).NotTo(HaveOccurred()) + lb, err := b.ListCollections() + Expect(err).NotTo(HaveOccurred()) + Expect(la).To(Equal([]string{"col-a", "col-b"})) + Expect(lb).To(Equal(la)) + }) + + It("answers not found for a collection that does not exist anywhere", func() { + _, err := b.Upload("nope", "f.txt", nil) + Expect(err).To(MatchError(ContainSubstring("collection not found: nope"))) + Expect(innerB.opens.Load()).To(BeZero(), "a miss must never create the collection") + }) + + It("answers not found on B after the collection is gone, without waiting for the re-check", func() { + Expect(a.CreateCollection("col-1")).To(Succeed()) + _, err := b.Search("col-1", "q", 1) + Expect(err).NotTo(HaveOccurred()) + + // The collection leaves the shared registry; A tells the other frontends. + innerA.resetDeletes = true + Expect(a.Reset("col-1")).To(Succeed()) + + _, err = b.Search("col-1", "q", 1) + Expect(err).To(MatchError(ContainSubstring("collection not found"))) + Expect(innerB.known("col-1")).To(HaveOccurred(), "the stale cache entry is dropped") + }) + + It("notices a removal at the next re-check when no event arrives", func() { + quiet := newSharedCollections(innerB, innerB.state, reg, nil) + quiet.recheck = 0 + Expect(a.CreateCollection("col-1")).To(Succeed()) + _, err := quiet.Search("col-1", "q", 1) + Expect(err).NotTo(HaveOccurred()) + + reg.remove("col-1") + _, err = quiet.Search("col-1", "q", 1) + Expect(err).To(MatchError(ContainSubstring("collection not found"))) + }) + + It("keeps serving the cached collection when the registry is unreachable", func() { + Expect(a.CreateCollection("col-1")).To(Succeed()) + _, err := b.Search("col-1", "q", 1) + Expect(err).NotTo(HaveOccurred()) + + b.forget("col-1") + reg.setErr(errors.New("connection refused")) + _, err = b.Search("col-1", "q", 1) + Expect(err).NotTo(HaveOccurred()) + + _, err = b.ListCollections() + Expect(err).To(HaveOccurred()) + }) + + It("re-checks at once when an event arrives", func() { + Expect(a.CreateCollection("col-1")).To(Succeed()) + _, err := b.Search("col-1", "q", 1) + Expect(err).NotTo(HaveOccurred()) + Expect(b.fresh("col-1")).To(BeTrue()) + + Expect(a.Reset("col-1")).To(Succeed()) + Expect(b.fresh("col-1")).To(BeFalse()) + Expect(bus.PublishCount("cache.invalidate.collections.col-1")).To(Equal(2), "create and reset") + }) + + It("opens a collection once for concurrent first lookups", func() { + Expect(a.CreateCollection("col-1")).To(Succeed()) + + var wg sync.WaitGroup + for range 20 { + wg.Add(1) + go func() { + defer wg.Done() + defer GinkgoRecover() + _, err := b.Search("col-1", "q", 1) + Expect(err).NotTo(HaveOccurred()) + }() + } + wg.Wait() + Expect(innerB.opens.Load()).To(Equal(int32(1))) + }) + + It("works without a bus", func() { + c := newSharedCollections(innerB, innerB.state, reg, nil) + Expect(a.CreateCollection("col-1")).To(Succeed()) + _, err := c.Upload("col-1", "f.txt", nil) + Expect(err).NotTo(HaveOccurred()) + c.Close() + }) +}) diff --git a/core/services/messaging/subjects.go b/core/services/messaging/subjects.go index e72222fc2..44e475229 100644 --- a/core/services/messaging/subjects.go +++ b/core/services/messaging/subjects.go @@ -287,6 +287,10 @@ type CacheInvalidateEvent struct { ConfigRevision string `json:"config_revision,omitempty"` } +// SubjectCacheInvalidateCollectionAll matches the collection cache +// invalidation subject of every collection. +const SubjectCacheInvalidateCollectionAll = "cache.invalidate.collections.*" + // SubjectCacheInvalidateCollection returns the NATS subject for collection cache invalidation. func SubjectCacheInvalidateCollection(name string) string { return "cache.invalidate.collections." + sanitizeSubjectToken(name) diff --git a/docs/content/features/distributed-mode.md b/docs/content/features/distributed-mode.md index 96f6fd08f..3d99cd63b 100644 --- a/docs/content/features/distributed-mode.md +++ b/docs/content/features/distributed-mode.md @@ -835,6 +835,17 @@ gauge and the query above see only LocalAI's own sessions -- and the transaction wedges the horizon is typically the co-located vector store connecting as a different role, which is exactly the case they exist to catch. +## Agent collections across frontends + +Collections (the `/api/agents/collections` routes) are **shared state**: every frontend shows the same list. With the `postgres` vector engine, the database is the source of truth for which collections exist. Each frontend keeps the collections it has opened in memory, but only as a cache: + +- Listing a collection set reads it from the database. +- A request for a collection that this frontend has not opened yet checks the database before it answers `404`, and opens the collection if it exists. A collection created through one frontend is therefore usable through any other one at once. +- A frontend re-checks a cached collection at most every 5 seconds, so a collection removed through another frontend stops being served within that time. +- Creating or resetting a collection publishes an event on the message bus, so the other frontends re-check immediately. + +This needs `LOCALAI_AGENT_POOL_VECTOR_ENGINE=postgres` and `LOCALAI_AGENT_POOL_DATABASE_URL`. With another vector engine there is no shared registry, and each frontend keeps its own list. Per-user collections (authentication enabled) are not yet coordinated this way. + ## Agent Workers Agent workers are dedicated processes for executing agent chats and MCP CI jobs. Unlike backend workers (which run gRPC model inference), agent workers use cogito to orchestrate multi-step conversations with tool calls. diff --git a/go.mod b/go.mod index c558895ab..e16d680f5 100644 --- a/go.mod +++ b/go.mod @@ -258,7 +258,7 @@ require ( github.com/labstack/gommon v0.4.2 // indirect github.com/mschoch/smat v0.2.0 // indirect github.com/mudler/LocalAGI v0.0.0-20260927202351-7e0947d7ebca - github.com/mudler/localrecall v0.6.6 // indirect + github.com/mudler/localrecall v0.6.7-0.20261006203612-d7c99211a3c7 github.com/mudler/skillserver v0.0.7-0.20260520220837-a7317cbf9145 github.com/olekukonko/tablewriter v0.0.5 // indirect github.com/oxffaa/gopher-parse-sitemap v0.0.0-20191021113419-005d2eb1def4 // indirect diff --git a/go.sum b/go.sum index 9edb8b051..244f166e0 100644 --- a/go.sum +++ b/go.sum @@ -1044,6 +1044,8 @@ github.com/mudler/go-processmanager v0.1.2-0.20260823202314-dfa0ed852db6 h1:/nFm github.com/mudler/go-processmanager v0.1.2-0.20260823202314-dfa0ed852db6/go.mod h1:h6kmHUZeafr+k5hRYpGLMzJFH4hItHffgpRo2QIkP+o= github.com/mudler/localrecall v0.6.6 h1:yN0vSeAnKjJ0BthkQuMpailWXHKs4W8yZcW8SVsfl/M= github.com/mudler/localrecall v0.6.6/go.mod h1:28k5n19raUrkuwXkacdNsBlj8yuSnGhpT16tu+2+4dU= +github.com/mudler/localrecall v0.6.7-0.20261006203612-d7c99211a3c7 h1:IKcbMzjpBD/aENARXz1SvLBeYqyvL1ddEVskXVBmAmA= +github.com/mudler/localrecall v0.6.7-0.20261006203612-d7c99211a3c7/go.mod h1:MRyhZMDZOBn1zmTmpVNB0F4Qjjx1tn9fK5CIWH6PyEA= github.com/mudler/memory v0.0.0-20260406210934-424c1ecf2cf8 h1:Ry8RiWy8fZ6Ff4E7dPmjRsBrnHOnPeOOj2LhCgyjQu0= github.com/mudler/memory v0.0.0-20260406210934-424c1ecf2cf8/go.mod h1:EA8Ashhd56o32qN7ouPKFSRUs/Z+LrRCF4v6R2Oarm8= github.com/mudler/nib v0.12.1 h1:7yKZeOWcIac62UoNb3J2R/oEtFYzebFwgap9HPjIJdI=