diff --git a/go.mod b/go.mod index 75a79b7..05c10d0 100644 --- a/go.mod +++ b/go.mod @@ -3,3 +3,12 @@ module gitea.perlbach24.de/scripte/nexarch go 1.22 require github.com/jackc/pgx/v5 v5.6.0 + +require ( + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a // indirect + github.com/jackc/puddle/v2 v2.2.1 // indirect + golang.org/x/crypto v0.17.0 // indirect + golang.org/x/sync v0.1.0 // indirect + golang.org/x/text v0.14.0 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..5c39671 --- /dev/null +++ b/go.sum @@ -0,0 +1,28 @@ +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a h1:bbPeKD0xmW/Y25WS6cokEszi5g+S0QxI/d45PkRi7Nk= +github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.6.0 h1:SWJzexBzPL5jb0GEsrPMLIsi/3jOo7RHlzTjcAeDrPY= +github.com/jackc/pgx/v5 v5.6.0/go.mod h1:DNZ/vlrUnhWCoFGxHAG8U2ljioxukquj7utPDgtQdTw= +github.com/jackc/puddle/v2 v2.2.1 h1:RhxXJtFG022u4ibrCSMSiu5aOq1i77R3OHKNJj77OAk= +github.com/jackc/puddle/v2 v2.2.1/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk= +github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= +golang.org/x/crypto v0.17.0 h1:r8bRNjWL3GshPW3gkd+RpvzWrZAwPS49OmTGZ/uhM4k= +golang.org/x/crypto v0.17.0/go.mod h1:gCAAfMLgwOJRpTjQ2zCCt2OcSfYMTeZVSRtQlPC7Nq4= +golang.org/x/sync v0.1.0 h1:wsuoTGHzEhffawBOhz5CYhcrV4IdKZbEyZjBMuTp12o= +golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/text v0.14.0 h1:ScX5w1eTa3QqT8oi6+ziP7dTV1S2+ALU0bI+0zXKWiQ= +golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= 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) + } +}