mirror of
https://github.com/mudler/LocalAI.git
synced 2026-10-09 22:54:42 -04:00
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 <mudler@localai.io> Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
This commit is contained in:
1 parent
e2faec2bf2
commit
217038d4fe
8 files changed
+674
-3
No files matched your search
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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()
|
||||
})
|
||||
})
|
||||
@@ -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)
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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=
|
||||
|
||||
Reference in new issue
Block a user