package tenant import ( "container/list" "context" "errors" "fmt" "sync" "github.com/jackc/pgx/v5/pgxpool" ) // ErrMissingTenantContext wird geliefert, wenn keine Tenant-Kennung // uebergeben wurde — es gibt bewusst keinen impliziten Default-Tenant // (TEN-06 Akzeptanzkriterium 3). var ErrMissingTenantContext = errors.New("tenant: kein tenant-kontext angegeben") // ErrUnknownTenant wird geliefert, wenn die Tenant-Kennung in der Registry // nicht existiert. var ErrUnknownTenant = errors.New("tenant: unbekannter tenant") // Router loest den Tenant-Kontext (Slug, aus dem JWT-Claim von API-05) in // eine wiederverwendbare Verbindung zur richtigen Tenant-Datenbank auf. // Ein LRU-verwalteter Cache begrenzt die Zahl gleichzeitig offener // pgxpool.Pool-Instanzen, damit die Zahl offener Postgres-Verbindungen NICHT // linear mit der Mandantenzahl waechst (Akzeptanzkriterium 2). type Router struct { registry *Registry maxOpen int mu sync.Mutex order *list.List // vorne = zuletzt benutzt items map[string]*list.Element // slug -> element mit *routerEntry } type routerEntry struct { slug string pool *pgxpool.Pool } func NewRouter(registry *Registry, maxOpen int) *Router { if maxOpen < 1 { maxOpen = 1 } return &Router{ registry: registry, maxOpen: maxOpen, order: list.New(), items: make(map[string]*list.Element), } } // Resolve liefert einen wiederverwendeten Pool fuer den angegebenen Tenant. // Ist der Tenant bereits im Cache, wird KEINE neue Verbindung aufgebaut // (Akzeptanzkriterium 2 / Pruefung 3). func (r *Router) Resolve(ctx context.Context, tenantSlug string) (*pgxpool.Pool, error) { if tenantSlug == "" { return nil, ErrMissingTenantContext } r.mu.Lock() if el, ok := r.items[tenantSlug]; ok { r.order.MoveToFront(el) pool := el.Value.(*routerEntry).pool r.mu.Unlock() return pool, nil } r.mu.Unlock() // Registry-Lookup und Verbindungsaufbau bewusst ausserhalb des Locks, // damit ein langsamer Verbindungsaufbau nicht alle anderen Tenants blockiert. t, err := r.registry.GetBySlug(ctx, tenantSlug) if err != nil { return nil, fmt.Errorf("%w: %s", ErrUnknownTenant, tenantSlug) } pool, err := pgxpool.New(ctx, t.DBDSN) if err != nil { return nil, fmt.Errorf("verbindung zu tenant %q aufbauen: %w", tenantSlug, err) } r.mu.Lock() defer r.mu.Unlock() // Zwischen Unlock oben und hier koennte ein paralleler Aufruf denselben // Tenant bereits eingefuegt haben — dann die eigene, ueberzaehlige // Verbindung wieder schliessen und die vorhandene verwenden. if el, ok := r.items[tenantSlug]; ok { r.order.MoveToFront(el) existing := el.Value.(*routerEntry).pool pool.Close() return existing, nil } el := r.order.PushFront(&routerEntry{slug: tenantSlug, pool: pool}) r.items[tenantSlug] = el if r.order.Len() > r.maxOpen { r.evictOldest() } return pool, nil } // evictOldest schliesst den am laengsten nicht genutzten Pool. Muss mit // gehaltenem r.mu aufgerufen werden. func (r *Router) evictOldest() { oldest := r.order.Back() if oldest == nil { return } entry := oldest.Value.(*routerEntry) r.order.Remove(oldest) delete(r.items, entry.slug) entry.pool.Close() } // OpenCount liefert die aktuelle Zahl offen gehaltener Tenant-Pools — // dient Tests/Monitoring, um AC2 nachzuweisen. func (r *Router) OpenCount() int { r.mu.Lock() defer r.mu.Unlock() return r.order.Len() } // Close schliesst alle offen gehaltenen Tenant-Pools, z.B. beim // Herunterfahren des Core-Prozesses. func (r *Router) Close() { r.mu.Lock() defer r.mu.Unlock() for el := r.order.Front(); el != nil; el = el.Next() { el.Value.(*routerEntry).pool.Close() } r.order.Init() r.items = make(map[string]*list.Element) }