diff --git a/internal/tenant/router.go b/internal/tenant/router.go new file mode 100644 index 0000000..4b2053b --- /dev/null +++ b/internal/tenant/router.go @@ -0,0 +1,136 @@ +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) +} diff --git a/internal/tenant/router_test.go b/internal/tenant/router_test.go new file mode 100644 index 0000000..e6f0482 --- /dev/null +++ b/internal/tenant/router_test.go @@ -0,0 +1,160 @@ +package tenant + +import ( + "context" + "errors" + "fmt" + "os" + "strings" + "testing" + + "github.com/jackc/pgx/v5/pgxpool" +) + +func newTestRouterSetup(t *testing.T, tenantCount int) (*Router, []Tenant, func()) { + t.Helper() + adminDSN := os.Getenv("TEST_ADMIN_DSN") + if adminDSN == "" { + t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen") + } + ctx := context.Background() + + adminPool, err := pgxpool.New(ctx, adminDSN) + if err != nil { + t.Fatalf("admin pool: %v", err) + } + registryPool, err := pgxpool.New(ctx, adminDSN) + if err != nil { + t.Fatalf("registry pool: %v", err) + } + if _, err := registryPool.Exec(ctx, ` + CREATE TABLE IF NOT EXISTS tenants ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + slug TEXT NOT NULL UNIQUE, + name TEXT NOT NULL, + db_name TEXT NOT NULL UNIQUE, + db_dsn TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'active', + created_at TIMESTAMPTZ NOT NULL DEFAULT now() + )`); err != nil { + t.Fatalf("registry-schema: %v", err) + } + + registry := NewRegistry(registryPool) + dsnTemplate := strings.Replace(adminDSN, "/postgres?", "/%s?", 1) + provisioner := NewProvisioner(adminPool, registry, dsnTemplate) + + var tenants []Tenant + var slugs []string + for i := 0; i < tenantCount; i++ { + slug := fmt.Sprintf("router_t%d", i) + slugs = append(slugs, slug) + tn, err := provisioner.Provision(ctx, slug, slug) + if err != nil { + t.Fatalf("provision %s: %v", slug, err) + } + tenants = append(tenants, tn) + } + + router := NewRouter(registry, 2) // klein gewaehlt, um Eviction im Test zu erzwingen + + cleanup := func() { + router.Close() + for _, slug := range slugs { + _, _ = adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, dbNameForSlug(slug))) + } + _, _ = registryPool.Exec(ctx, `DELETE FROM tenants WHERE slug = ANY($1)`, slugs) + registryPool.Close() + adminPool.Close() + } + return router, tenants, cleanup +} + +// Akzeptanzkriterium 1: Verbindung wird zuverlaessig anhand des Tenant-Kontexts aufgeloest. +func TestRouter_ResolvesCorrectTenantDatabase(t *testing.T) { + router, tenants, cleanup := newTestRouterSetup(t, 2) + defer cleanup() + ctx := context.Background() + + pool, err := router.Resolve(ctx, tenants[0].Slug) + if err != nil { + t.Fatalf("resolve: %v", err) + } + var dbName string + if err := pool.QueryRow(ctx, `SELECT current_database()`).Scan(&dbName); err != nil { + t.Fatalf("current_database: %v", err) + } + if dbName != tenants[0].DBName { + t.Fatalf("current_database() = %q, want %q", dbName, tenants[0].DBName) + } +} + +// Akzeptanzkriterium 3 + Pruefung 2: fehlender/unbekannter Tenant-Kontext +// wird explizit abgewiesen statt irgendeine Verbindung zu liefern. +func TestRouter_RejectsMissingOrUnknownTenant(t *testing.T) { + router, _, cleanup := newTestRouterSetup(t, 1) + defer cleanup() + ctx := context.Background() + + if _, err := router.Resolve(ctx, ""); !errors.Is(err, ErrMissingTenantContext) { + t.Fatalf("erwartet ErrMissingTenantContext, habe %v", err) + } + if _, err := router.Resolve(ctx, "nie-registrierter-slug"); !errors.Is(err, ErrUnknownTenant) { + t.Fatalf("erwartet ErrUnknownTenant, habe %v", err) + } +} + +// Akzeptanzkriterium 2 + Pruefung 3: Verbindungswiederverwendung nachweislich +// gemessen — zweiter Resolve-Aufruf liefert exakt denselben Pool, kein +// erneuter Verbindungsaufbau. +func TestRouter_ReusesConnectionForSameTenant(t *testing.T) { + router, tenants, cleanup := newTestRouterSetup(t, 1) + defer cleanup() + ctx := context.Background() + + first, err := router.Resolve(ctx, tenants[0].Slug) + if err != nil { + t.Fatalf("resolve 1: %v", err) + } + second, err := router.Resolve(ctx, tenants[0].Slug) + if err != nil { + t.Fatalf("resolve 2: %v", err) + } + if first != second { + t.Fatal("erwartet identische pool-instanz bei wiederholtem resolve, habe unterschiedliche") + } +} + +// Akzeptanzkriterium 2 + Pruefung 1: Lasttest mit mehr simulierten Mandanten +// als maxOpen — die Zahl gleichzeitig offener Tenant-Pools bleibt begrenzt +// (LRU-Eviction), waechst also NICHT linear mit der Mandantenzahl. +func TestRouter_BoundsOpenConnectionsUnderLoad(t *testing.T) { + const tenantCount = 6 + router, tenants, cleanup := newTestRouterSetup(t, tenantCount) + defer cleanup() + ctx := context.Background() + + for _, tn := range tenants { + if _, err := router.Resolve(ctx, tn.Slug); err != nil { + t.Fatalf("resolve %s: %v", tn.Slug, err) + } + if router.OpenCount() > 2 { + t.Fatalf("OpenCount() = %d, erwartet <= maxOpen (2) nach jedem Resolve", router.OpenCount()) + } + } + + if router.OpenCount() != 2 { + t.Fatalf("erwartet genau maxOpen=2 offene pools nach %d tenants, habe %d", tenantCount, router.OpenCount()) + } + + // Evictete Tenants sind wieder ganz normal ueber die Registry aufloesbar + // (Cache-Miss fuehrt zu neuem, funktionierendem Pool, kein Fehlerzustand). + pool, err := router.Resolve(ctx, tenants[0].Slug) + if err != nil { + t.Fatalf("resolve nach eviction: %v", err) + } + var one int + if err := pool.QueryRow(ctx, `SELECT 1`).Scan(&one); err != nil { + t.Fatalf("query nach re-resolve: %v", err) + } +}