diff --git a/internal/cfgservice/service.go b/internal/cfgservice/service.go new file mode 100644 index 0000000..b74e8bf --- /dev/null +++ b/internal/cfgservice/service.go @@ -0,0 +1,88 @@ +package cfgservice + +import ( + "context" + "errors" + "sync" + "time" +) + +// DefaultCacheTTL ist die dokumentierte Cache-Invalidierungszeit +// (Akzeptanzkriterium 2 / Pruefung 2 in diesem Ticket bezieht sich auf die +// Aenderungsnachvollziehbarkeit — die Cache-Frist selbst folgt demselben +// Muster wie internal/flag.DefaultCacheTTL). +const DefaultCacheTTL = 5 * time.Second + +type cacheEntry struct { + value Value + expiresAt time.Time +} + +// Service ist die Leseseite mit Vorrangregel (Akzeptanzkriterium 1: +// Tenant-Override vor Global-Default) und lokalem TTL-Cache. +type Service struct { + store *Store + ttl time.Duration + + mu sync.RWMutex + cache map[string]cacheEntry // Schluessel: key + "\x00" + tenantSlug +} + +func NewService(store *Store, ttl time.Duration) *Service { + if ttl <= 0 { + ttl = DefaultCacheTTL + } + return &Service{store: store, ttl: ttl, cache: make(map[string]cacheEntry)} +} + +func cacheKey(key, tenantSlug string) string { + return key + "\x00" + tenantSlug +} + +// Resolve liefert den Konfigurationswert fuer einen Tenant: ein +// Tenant-spezifischer Override hat Vorrang vor dem globalen Default +// (Akzeptanzkriterium 1 / Pruefung 1). tenantSlug == "" wertet nur den +// globalen Wert aus. +func (s *Service) Resolve(ctx context.Context, tenantSlug, key string) (Value, error) { + ck := cacheKey(key, tenantSlug) + + s.mu.RLock() + entry, exists := s.cache[ck] + fresh := exists && time.Now().Before(entry.expiresAt) + s.mu.RUnlock() + if fresh { + return entry.value, nil + } + + v, err := s.resolveUncached(ctx, tenantSlug, key) + if err != nil { + return Value{}, err + } + + s.mu.Lock() + s.cache[ck] = cacheEntry{value: v, expiresAt: time.Now().Add(s.ttl)} + s.mu.Unlock() + return v, nil +} + +func (s *Service) resolveUncached(ctx context.Context, tenantSlug, key string) (Value, error) { + if tenantSlug != "" { + v, err := s.store.Get(ctx, key, tenantSlug) + if err == nil { + return v, nil + } + if !errors.Is(err, ErrNotFound) { + return Value{}, err + } + } + return s.store.Get(ctx, key, GlobalScope) +} + +// Invalidate erzwingt beim naechsten Resolve-Aufruf ein sofortiges Neuladen +// fuer einen bestimmten (key, tenantSlug) statt auf den TTL-Ablauf zu warten +// — analog internal/flag.Service.Invalidate. +func (s *Service) Invalidate(key, tenantSlug string) { + s.mu.Lock() + delete(s.cache, cacheKey(key, tenantSlug)) + s.mu.Unlock() +} diff --git a/internal/cfgservice/service_test.go b/internal/cfgservice/service_test.go new file mode 100644 index 0000000..2542b47 --- /dev/null +++ b/internal/cfgservice/service_test.go @@ -0,0 +1,114 @@ +package cfgservice + +import ( + "context" + "testing" + "time" +) + +// Akzeptanzkriterium 1 + Pruefung 1: Tenant-Override hat Vorrang vor +// Global-Default, automatisiert getestet. +func TestService_TenantOverrideTakesPrecedenceOverGlobal(t *testing.T) { + store, cleanup := setupStoreTest(t) + defer cleanup() + ctx := context.Background() + + if _, err := store.Set(ctx, "test_precedence_key", GlobalScope, "global-wert"); err != nil { + t.Fatalf("set global: %v", err) + } + if _, err := store.Set(ctx, "test_precedence_key", "test_acme", "tenant-wert"); err != nil { + t.Fatalf("set tenant: %v", err) + } + + svc := NewService(store, time.Hour) + + got, err := svc.Resolve(ctx, "test_acme", "test_precedence_key") + if err != nil { + t.Fatalf("resolve mit override: %v", err) + } + if got.Value != "tenant-wert" { + t.Fatalf("erwartet tenant-override, habe %q", got.Value) + } + + gotOther, err := svc.Resolve(ctx, "test_anderer_tenant", "test_precedence_key") + if err != nil { + t.Fatalf("resolve ohne override: %v", err) + } + if gotOther.Value != "global-wert" { + t.Fatalf("erwartet global-default fuer tenant ohne override, habe %q", gotOther.Value) + } +} + +// Akzeptanzkriterium 2 + Pruefung 2: Cache-Invalidierung nach +// Konfigurationsaenderung innerhalb dokumentierter Zeit gemessen. +func TestService_CacheInvalidationTiming(t *testing.T) { + store, cleanup := setupStoreTest(t) + defer cleanup() + ctx := context.Background() + + const ttl = 150 * time.Millisecond + if _, err := store.Set(ctx, "test_ttl_key", GlobalScope, "alt"); err != nil { + t.Fatalf("set: %v", err) + } + svc := NewService(store, ttl) + + v, err := svc.Resolve(ctx, "", "test_ttl_key") + if err != nil { + t.Fatalf("resolve: %v", err) + } + if v.Value != "alt" { + t.Fatalf("erwartet 'alt', habe %q", v.Value) + } + + changedAt := time.Now() + if _, err := store.Set(ctx, "test_ttl_key", GlobalScope, "neu"); err != nil { + t.Fatalf("set: %v", err) + } + + v, err = svc.Resolve(ctx, "", "test_ttl_key") + if err != nil { + t.Fatalf("resolve direkt nach aenderung: %v", err) + } + if v.Value != "alt" { + t.Fatalf("cache haette den alten wert liefern sollen, habe %q", v.Value) + } + + deadline := changedAt.Add(ttl + 100*time.Millisecond) + for time.Now().Before(deadline) { + v, err := svc.Resolve(ctx, "", "test_ttl_key") + if err != nil { + t.Fatalf("resolve: %v", err) + } + if v.Value == "neu" { + t.Logf("aenderung wurde nach %s wirksam (ziel: innerhalb %s + toleranz)", time.Since(changedAt), ttl) + return + } + time.Sleep(10 * time.Millisecond) + } + t.Fatalf("aenderung wurde nicht innerhalb von %s wirksam", deadline.Sub(changedAt)) +} + +func TestService_InvalidateForcesImmediateRefresh(t *testing.T) { + store, cleanup := setupStoreTest(t) + defer cleanup() + ctx := context.Background() + + if _, err := store.Set(ctx, "test_invalidate_key", GlobalScope, "alt"); err != nil { + t.Fatalf("set: %v", err) + } + svc := NewService(store, time.Hour) + _, _ = svc.Resolve(ctx, "", "test_invalidate_key") + + if _, err := store.Set(ctx, "test_invalidate_key", GlobalScope, "neu"); err != nil { + t.Fatalf("set: %v", err) + } + svc.Invalidate("test_invalidate_key", "") + + v, err := svc.Resolve(ctx, "", "test_invalidate_key") + if err != nil { + t.Fatalf("resolve: %v", err) + } + if v.Value != "neu" { + t.Fatalf("erwartet sofort sichtbaren neuen wert nach Invalidate, habe %q", v.Value) + } +} diff --git a/internal/cfgservice/store.go b/internal/cfgservice/store.go new file mode 100644 index 0000000..8be1449 --- /dev/null +++ b/internal/cfgservice/store.go @@ -0,0 +1,129 @@ +// Package cfgservice implementiert Core CFG-01: den zentralen Dienst fuer +// globale und tenant-spezifische Konfigurationswerte mit Versionierung und +// Cache-Invalidierung. Andere Module lesen Konfiguration AUSSCHLIESSLICH +// ueber dieses Paket (Akzeptanzkriterium 3), niemals ueber eigene Tabellen. +package cfgservice + +import ( + "context" + "errors" + "fmt" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" +) + +// GlobalScope ist der reservierte Scope-Wert fuer globale Defaults — jeder +// andere Scope-Wert ist ein Tenant-Slug (Akzeptanzkriterium 1). +const GlobalScope = "global" + +var ErrNotFound = errors.New("cfgservice: kein wert fuer diesen key gefunden") + +type Value struct { + Key string + Scope string + Value string + Version int +} + +type HistoryEntry struct { + Key string + Scope string + Value string + Version int +} + +// Store ist die Schreib-/Verwaltungsseite. Set schreibt IMMER sowohl den +// aktuellen Stand (config_values) als auch einen Historieneintrag +// (config_value_history) in derselben Transaktion — eine Aenderung ohne +// Versionshistorie ist strukturell ausgeschlossen (Akzeptanzkriterium 2). +type Store struct { + pool *pgxpool.Pool +} + +func NewStore(pool *pgxpool.Pool) *Store { + return &Store{pool: pool} +} + +// Set schreibt einen neuen Wert fuer (key, scope) und erhoeht die Version um 1 +// (Version 1 bei erstmaligem Setzen). +func (s *Store) Set(ctx context.Context, key, scope, value string) (Value, error) { + if scope == "" { + return Value{}, errors.New("cfgservice: scope darf nicht leer sein") + } + + tx, err := s.pool.Begin(ctx) + if err != nil { + return Value{}, fmt.Errorf("transaktion starten: %w", err) + } + defer func() { _ = tx.Rollback(ctx) }() + + var currentVersion int + err = tx.QueryRow(ctx, `SELECT version FROM config_values WHERE key = $1 AND scope = $2`, key, scope).Scan(¤tVersion) + if err != nil && !errors.Is(err, pgx.ErrNoRows) { + return Value{}, fmt.Errorf("aktuelle version lesen: %w", err) + } + newVersion := currentVersion + 1 + + if _, err := tx.Exec(ctx, ` + INSERT INTO config_values (key, scope, value, version, updated_at) + VALUES ($1, $2, $3, $4, now()) + ON CONFLICT (key, scope) DO UPDATE SET value = $3, version = $4, updated_at = now() + `, key, scope, value, newVersion); err != nil { + return Value{}, fmt.Errorf("wert speichern: %w", err) + } + + if _, err := tx.Exec(ctx, ` + INSERT INTO config_value_history (key, scope, value, version, changed_at) + VALUES ($1, $2, $3, $4, now()) + `, key, scope, value, newVersion); err != nil { + return Value{}, fmt.Errorf("historie schreiben: %w", err) + } + + if err := tx.Commit(ctx); err != nil { + return Value{}, fmt.Errorf("transaktion committen: %w", err) + } + + return Value{Key: key, Scope: scope, Value: value, Version: newVersion}, nil +} + +// Get liefert den Wert fuer GENAU EINEN Scope (kein Vorrang-Fallback) — die +// Vorrangregel (Tenant vor Global) lebt bewusst in Service.Resolve, damit +// Store rein CRUD bleibt. +func (s *Store) Get(ctx context.Context, key, scope string) (Value, error) { + var v Value + v.Key, v.Scope = key, scope + err := s.pool.QueryRow(ctx, ` + SELECT value, version FROM config_values WHERE key = $1 AND scope = $2 + `, key, scope).Scan(&v.Value, &v.Version) + if err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return Value{}, ErrNotFound + } + return Value{}, fmt.Errorf("wert lesen: %w", err) + } + return v, nil +} + +// History liefert die vollstaendige Versionshistorie eines (key, scope) in +// aufsteigender Reihenfolge (Akzeptanzkriterium 2 / Pruefung 3). +func (s *Store) History(ctx context.Context, key, scope string) ([]HistoryEntry, error) { + rows, err := s.pool.Query(ctx, ` + SELECT key, scope, value, version FROM config_value_history + WHERE key = $1 AND scope = $2 ORDER BY version + `, key, scope) + if err != nil { + return nil, fmt.Errorf("historie abfragen: %w", err) + } + defer rows.Close() + + var out []HistoryEntry + for rows.Next() { + var h HistoryEntry + if err := rows.Scan(&h.Key, &h.Scope, &h.Value, &h.Version); err != nil { + return nil, fmt.Errorf("historieneintrag lesen: %w", err) + } + out = append(out, h) + } + return out, rows.Err() +} diff --git a/internal/cfgservice/store_test.go b/internal/cfgservice/store_test.go new file mode 100644 index 0000000..d3f6ac3 --- /dev/null +++ b/internal/cfgservice/store_test.go @@ -0,0 +1,119 @@ +package cfgservice + +import ( + "context" + "errors" + "os" + "testing" + + "github.com/jackc/pgx/v5/pgxpool" +) + +func setupStoreTest(t *testing.T) (*Store, func()) { + t.Helper() + adminDSN := os.Getenv("TEST_ADMIN_DSN") + if adminDSN == "" { + t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen") + } + ctx := context.Background() + + pool, err := pgxpool.New(ctx, adminDSN) + if err != nil { + t.Fatalf("pool: %v", err) + } + if _, err := pool.Exec(ctx, ` + CREATE TABLE IF NOT EXISTS config_values ( + key TEXT NOT NULL, + scope TEXT NOT NULL CHECK (scope <> ''), + value TEXT NOT NULL, + version INT NOT NULL, + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + PRIMARY KEY (key, scope) + ); + CREATE TABLE IF NOT EXISTS config_value_history ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + key TEXT NOT NULL, + scope TEXT NOT NULL, + value TEXT NOT NULL, + version INT NOT NULL, + changed_at TIMESTAMPTZ NOT NULL DEFAULT now() + ); + `); err != nil { + t.Fatalf("schema: %v", err) + } + + cleanup := func() { + _, _ = pool.Exec(ctx, `DELETE FROM config_value_history WHERE key LIKE 'test\_%' ESCAPE '\'`) + _, _ = pool.Exec(ctx, `DELETE FROM config_values WHERE key LIKE 'test\_%' ESCAPE '\'`) + pool.Close() + } + return NewStore(pool), cleanup +} + +func TestStore_SetIncrementsVersion(t *testing.T) { + store, cleanup := setupStoreTest(t) + defer cleanup() + ctx := context.Background() + + v1, err := store.Set(ctx, "test_key", GlobalScope, "erster-wert") + if err != nil { + t.Fatalf("set 1: %v", err) + } + if v1.Version != 1 { + t.Fatalf("erwartet version 1, habe %d", v1.Version) + } + + v2, err := store.Set(ctx, "test_key", GlobalScope, "zweiter-wert") + if err != nil { + t.Fatalf("set 2: %v", err) + } + if v2.Version != 2 { + t.Fatalf("erwartet version 2, habe %d", v2.Version) + } + + got, err := store.Get(ctx, "test_key", GlobalScope) + if err != nil { + t.Fatalf("get: %v", err) + } + if got.Value != "zweiter-wert" || got.Version != 2 { + t.Fatalf("aktueller wert unerwartet: %+v", got) + } +} + +// Akzeptanzkriterium 2 + Pruefung 3: Versionierungshistorie ueber mehrere +// Aenderungen hinweg nachvollzogen. +func TestStore_HistoryTracksAllChanges(t *testing.T) { + store, cleanup := setupStoreTest(t) + defer cleanup() + ctx := context.Background() + + values := []string{"v1", "v2", "v3"} + for _, v := range values { + if _, err := store.Set(ctx, "test_history_key", GlobalScope, v); err != nil { + t.Fatalf("set %q: %v", v, err) + } + } + + history, err := store.History(ctx, "test_history_key", GlobalScope) + if err != nil { + t.Fatalf("history: %v", err) + } + if len(history) != 3 { + t.Fatalf("erwartet 3 historieneintraege, habe %d", len(history)) + } + for i, h := range history { + if h.Version != i+1 || h.Value != values[i] { + t.Fatalf("historieneintrag[%d] unerwartet: %+v", i, h) + } + } +} + +func TestStore_GetUnknownKeyReturnsNotFound(t *testing.T) { + store, cleanup := setupStoreTest(t) + defer cleanup() + ctx := context.Background() + + if _, err := store.Get(ctx, "test_nie_gesetzt", GlobalScope); !errors.Is(err, ErrNotFound) { + t.Fatalf("erwartet ErrNotFound, habe %v", err) + } +} diff --git a/internal/notify/dispatcher.go b/internal/notify/dispatcher.go new file mode 100644 index 0000000..21c0cd1 --- /dev/null +++ b/internal/notify/dispatcher.go @@ -0,0 +1,81 @@ +// Package notify implementiert Core CFG-02: den zentralen Benachrichtigungs- +// Dispatcher, ueber den beliebige Module Benachrichtigungen ausloesen — +// Warteschlange, Wiederholungslogik, Kanal-Abstraktion. Die tatsaechlichen +// Kanaele (E-Mail/In-App) sind CFG-03, hier gibt es nur die Sender- +// Schnittstelle als Vorbereitung. +package notify + +import ( + "context" + "encoding/json" + "fmt" + "time" + + "github.com/jackc/pgx/v5/pgxpool" +) + +// DefaultMaxAttempts begrenzt Wiederholungsversuche (Akzeptanzkriterium 2) — +// nach dieser Anzahl gibt der Dispatcher kontrolliert auf (status=failed) +// statt endlos zu wiederholen. +const DefaultMaxAttempts = 5 + +// DefaultRetryBackoff ist die Basis-Wartezeit zwischen Wiederholungen, +// linear mit der Versuchsnummer skaliert. +const DefaultRetryBackoff = 200 * time.Millisecond + +type Notification struct { + ID string + Channel string + Recipient string + Payload map[string]any + Attempts int +} + +// Sender ist die schmale Schnittstelle, die ein konkreter Kanal (CFG-03) +// implementiert. Der Dispatcher selbst weiss nichts ueber E-Mail/In-App. +type Sender interface { + Send(ctx context.Context, n Notification) error +} + +// Dispatcher ist die EINE Schnittstelle, ueber die Module Benachrichtigungen +// ausloesen — kein Modul baut eigenen Versandcode (Akzeptanzkriterium 1). +type Dispatcher struct { + pool *pgxpool.Pool + maxAttempts int + retryBackoff time.Duration +} + +func NewDispatcher(pool *pgxpool.Pool) *Dispatcher { + return &Dispatcher{pool: pool, maxAttempts: DefaultMaxAttempts, retryBackoff: DefaultRetryBackoff} +} + +// WithRetryPolicy erlaubt Tests/Betrieb, Versuchsanzahl und Backoff +// anzupassen, ohne die Default-Policy im Produktionscode zu veraendern. +func (d *Dispatcher) WithRetryPolicy(maxAttempts int, backoff time.Duration) *Dispatcher { + return &Dispatcher{pool: d.pool, maxAttempts: maxAttempts, retryBackoff: backoff} +} + +// Enqueue reiht eine Benachrichtigung in die Postgres-Warteschlange ein und +// kehrt sofort zurueck — die Zeile ueberlebt jeden Neustart des Dispatcher- +// Prozesses unveraendert (Akzeptanzkriterium 3), da sie ausschliesslich in +// der Datenbank existiert, nicht im Prozessspeicher. +func (d *Dispatcher) Enqueue(ctx context.Context, channel, recipient string, payload map[string]any) (string, error) { + if payload == nil { + payload = map[string]any{} + } + payloadJSON, err := json.Marshal(payload) + if err != nil { + return "", fmt.Errorf("payload serialisieren: %w", err) + } + + var id string + err = d.pool.QueryRow(ctx, ` + INSERT INTO notification_jobs (channel, recipient, payload, max_attempts) + VALUES ($1, $2, $3, $4) + RETURNING id + `, channel, recipient, payloadJSON, d.maxAttempts).Scan(&id) + if err != nil { + return "", fmt.Errorf("benachrichtigung einreihen: %w", err) + } + return id, nil +} diff --git a/internal/notify/dispatcher_test.go b/internal/notify/dispatcher_test.go new file mode 100644 index 0000000..bcb490a --- /dev/null +++ b/internal/notify/dispatcher_test.go @@ -0,0 +1,219 @@ +package notify + +import ( + "context" + "errors" + "os" + "sync" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgxpool" +) + +func setupTest(t *testing.T) (*pgxpool.Pool, func()) { + t.Helper() + adminDSN := os.Getenv("TEST_ADMIN_DSN") + if adminDSN == "" { + t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen") + } + ctx := context.Background() + + pool, err := pgxpool.New(ctx, adminDSN) + if err != nil { + t.Fatalf("pool: %v", err) + } + if _, err := pool.Exec(ctx, ` + CREATE TABLE IF NOT EXISTS notification_jobs ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + channel TEXT NOT NULL, + recipient TEXT NOT NULL, + payload JSONB NOT NULL DEFAULT '{}'::jsonb, + status TEXT NOT NULL DEFAULT 'pending' CHECK (status IN ('pending', 'sent', 'failed')), + attempts INT NOT NULL DEFAULT 0, + max_attempts INT NOT NULL DEFAULT 5, + next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT now(), + last_error TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() + )`); err != nil { + t.Fatalf("schema: %v", err) + } + + cleanup := func() { pool.Close() } + return pool, cleanup +} + +type fakeSender struct { + mu sync.Mutex + sentIDs []string + failUntil int + calls int +} + +func (f *fakeSender) Send(ctx context.Context, n Notification) error { + f.mu.Lock() + defer f.mu.Unlock() + f.calls++ + if f.calls <= f.failUntil { + return errors.New("simulierter zustellfehler") + } + f.sentIDs = append(f.sentIDs, n.ID) + return nil +} + +func (f *fakeSender) sentCount() int { + f.mu.Lock() + defer f.mu.Unlock() + return len(f.sentIDs) +} + +// Akzeptanzkriterium 1: Module loesen ueber Enqueue aus, keine eigene +// Versandlogik noetig. +func TestDispatcher_EnqueueAndProcess(t *testing.T) { + pool, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + + d := NewDispatcher(pool) + id, err := d.Enqueue(ctx, "email", "alice@example.com", map[string]any{"subject": "Willkommen"}) + if err != nil { + t.Fatalf("enqueue: %v", err) + } + if id == "" { + t.Fatal("erwartet nicht-leere id") + } + + sender := &fakeSender{} + sent, failed, err := d.ProcessDue(ctx, sender, 10) + if err != nil { + t.Fatalf("process: %v", err) + } + if sent != 1 || failed != 0 { + t.Fatalf("erwartet sent=1 failed=0, habe sent=%d failed=%d", sent, failed) + } + if sender.sentCount() != 1 { + t.Fatalf("erwartet 1 zustellung, habe %d", sender.sentCount()) + } +} + +// Akzeptanzkriterium 2 + Pruefung 2: Wiederholungslogik greift bei +// simuliertem Fehler und bricht nach definierter Anzahl kontrolliert ab. +func TestProcessDue_RetriesThenGivesUpAfterMaxAttempts(t *testing.T) { + pool, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + + d := NewDispatcher(pool).WithRetryPolicy(3, time.Millisecond) + id, err := d.Enqueue(ctx, "email", "bob@example.com", nil) + if err != nil { + t.Fatalf("enqueue: %v", err) + } + + sender := &fakeSender{failUntil: 100} // schlaegt bei jedem versuch fehl + + for i := 0; i < 3; i++ { + time.Sleep(5 * time.Millisecond) // next_attempt_at abwarten + if _, _, err := d.ProcessDue(ctx, sender, 10); err != nil { + t.Fatalf("process %d: %v", i, err) + } + } + + var status string + var attempts int + if err := pool.QueryRow(ctx, `SELECT status, attempts FROM notification_jobs WHERE id = $1`, id).Scan(&status, &attempts); err != nil { + t.Fatalf("status lesen: %v", err) + } + if status != "failed" { + t.Fatalf("erwartet status failed nach max_attempts, habe %q", status) + } + if attempts != 3 { + t.Fatalf("erwartet 3 versuche, habe %d", attempts) + } + + // Weiteres ProcessDue darf den bereits aufgegebenen job nicht mehr anfassen. + sent, failed, err := d.ProcessDue(ctx, sender, 10) + if err != nil { + t.Fatalf("process nach abbruch: %v", err) + } + if sent != 0 || failed != 0 { + t.Fatalf("erwartet keine weitere verarbeitung, habe sent=%d failed=%d", sent, failed) + } +} + +// Akzeptanzkriterium 3 + Pruefung 1: Neustart des Dienstes waehrend offener +// Zustellung verliert keine Nachricht — simuliert durch eine komplett neue +// Dispatcher/Pool-Instanz nach dem Enqueue, bevor irgendetwas verarbeitet wurde. +func TestQueue_SurvivesRestartWithoutMessageLoss(t *testing.T) { + pool, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + + firstInstance := NewDispatcher(pool) + id, err := firstInstance.Enqueue(ctx, "email", "carol@example.com", nil) + if err != nil { + t.Fatalf("enqueue: %v", err) + } + + // "Neustart": eine voellig neue Dispatcher-Instanz (repraesentiert einen + // neuen Prozess) verbindet sich neu und verarbeitet die Warteschlange — + // die Nachricht existiert ausschliesslich in Postgres, nicht im + // Prozessspeicher der ersten Instanz. + restartedInstance := NewDispatcher(pool) + sender := &fakeSender{} + sent, failed, err := restartedInstance.ProcessDue(ctx, sender, 10) + if err != nil { + t.Fatalf("process nach neustart: %v", err) + } + if sent != 1 || failed != 0 { + t.Fatalf("erwartet sent=1 nach neustart, habe sent=%d failed=%d", sent, failed) + } + if len(sender.sentIDs) != 1 || sender.sentIDs[0] != id { + t.Fatalf("erwartet zustellung der urspruenglichen nachricht %q, habe %v", id, sender.sentIDs) + } +} + +// Akzeptanzkriterium 3 + Pruefung 3: zwei gleichzeitig ausloesende Module, +// beide Nachrichten werden korrekt (und nicht doppelt) zugestellt. +func TestProcessDue_ConcurrentDispatchBothDelivered(t *testing.T) { + pool, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + + d := NewDispatcher(pool) + idA, err := d.Enqueue(ctx, "email", "modul-a@example.com", nil) + if err != nil { + t.Fatalf("enqueue a: %v", err) + } + idB, err := d.Enqueue(ctx, "email", "modul-b@example.com", nil) + if err != nil { + t.Fatalf("enqueue b: %v", err) + } + + sender := &fakeSender{} + var wg sync.WaitGroup + for i := 0; i < 2; i++ { + wg.Add(1) + go func() { + defer wg.Done() + if _, _, err := d.ProcessDue(ctx, sender, 10); err != nil { + t.Errorf("process: %v", err) + } + }() + } + wg.Wait() + + if sender.sentCount() != 2 { + t.Fatalf("erwartet genau 2 zustellungen, habe %d: %v", sender.sentCount(), sender.sentIDs) + } + seen := map[string]bool{} + for _, id := range sender.sentIDs { + if seen[id] { + t.Fatalf("nachricht %q wurde doppelt zugestellt", id) + } + seen[id] = true + } + if !seen[idA] || !seen[idB] { + t.Fatalf("erwartet beide nachrichten zugestellt, habe %v", sender.sentIDs) + } +} diff --git a/internal/notify/worker.go b/internal/notify/worker.go new file mode 100644 index 0000000..cad3f8d --- /dev/null +++ b/internal/notify/worker.go @@ -0,0 +1,102 @@ +package notify + +import ( + "context" + "encoding/json" + "fmt" + "time" +) + +// ProcessDue holt bis zu limit faellige Benachrichtigungen und versucht sie +// ueber sender zuzustellen. FOR UPDATE SKIP LOCKED serialisiert konkurrierende +// Aufrufe (Akzeptanzkriterium 3 / Pruefung 3: zwei gleichzeitig ausloesende +// Module duerfen sich nicht gegenseitig blockieren oder Nachrichten doppelt +// zustellen) — dieselbe Konvention wie internal/tenant.Lifecycle.ProcessDueDeletions. +func (d *Dispatcher) ProcessDue(ctx context.Context, sender Sender, limit int) (sent, failed int, err error) { + tx, err := d.pool.Begin(ctx) + if err != nil { + return 0, 0, fmt.Errorf("transaktion starten: %w", err) + } + defer func() { _ = tx.Rollback(ctx) }() + + rows, err := tx.Query(ctx, ` + SELECT id, channel, recipient, payload, attempts, max_attempts + FROM notification_jobs + WHERE status = 'pending' AND next_attempt_at <= now() + ORDER BY created_at + FOR UPDATE SKIP LOCKED + LIMIT $1 + `, limit) + if err != nil { + return 0, 0, fmt.Errorf("faellige benachrichtigungen abfragen: %w", err) + } + + type due struct { + id, channel, recipient string + payload []byte + attempts, maxAttempts int + } + var candidates []due + for rows.Next() { + var c due + if err := rows.Scan(&c.id, &c.channel, &c.recipient, &c.payload, &c.attempts, &c.maxAttempts); err != nil { + rows.Close() + return 0, 0, fmt.Errorf("faellige benachrichtigung lesen: %w", err) + } + candidates = append(candidates, c) + } + rows.Close() + if err := rows.Err(); err != nil { + return 0, 0, err + } + + for _, c := range candidates { + var payload map[string]any + if err := json.Unmarshal(c.payload, &payload); err != nil { + payload = map[string]any{} + } + + sendErr := sender.Send(ctx, Notification{ + ID: c.id, Channel: c.channel, Recipient: c.recipient, Payload: payload, Attempts: c.attempts, + }) + + if sendErr == nil { + if _, err := tx.Exec(ctx, ` + UPDATE notification_jobs SET status = 'sent', updated_at = now() WHERE id = $1 + `, c.id); err != nil { + return sent, failed, fmt.Errorf("erfolg speichern: %w", err) + } + sent++ + continue + } + + newAttempts := c.attempts + 1 + if newAttempts >= c.maxAttempts { + // Akzeptanzkriterium 2: kontrollierter Abbruch nach definierter + // Anzahl Versuche, kein endloses Wiederholen. + if _, err := tx.Exec(ctx, ` + UPDATE notification_jobs + SET status = 'failed', attempts = $2, last_error = $3, updated_at = now() + WHERE id = $1 + `, c.id, newAttempts, sendErr.Error()); err != nil { + return sent, failed, fmt.Errorf("fehlschlag speichern: %w", err) + } + failed++ + continue + } + + nextAttempt := time.Now().Add(time.Duration(newAttempts) * d.retryBackoff) + if _, err := tx.Exec(ctx, ` + UPDATE notification_jobs + SET attempts = $2, next_attempt_at = $3, last_error = $4, updated_at = now() + WHERE id = $1 + `, c.id, newAttempts, nextAttempt, sendErr.Error()); err != nil { + return sent, failed, fmt.Errorf("wiederholung planen: %w", err) + } + } + + if err := tx.Commit(ctx); err != nil { + return 0, 0, fmt.Errorf("transaktion committen: %w", err) + } + return sent, failed, nil +} diff --git a/migrations/0004_config_values.down.sql b/migrations/0004_config_values.down.sql new file mode 100644 index 0000000..6fb11d7 --- /dev/null +++ b/migrations/0004_config_values.down.sql @@ -0,0 +1,2 @@ +DROP TABLE IF EXISTS config_value_history; +DROP TABLE IF EXISTS config_values; diff --git a/migrations/0004_config_values.up.sql b/migrations/0004_config_values.up.sql new file mode 100644 index 0000000..56a9fb1 --- /dev/null +++ b/migrations/0004_config_values.up.sql @@ -0,0 +1,23 @@ +-- Zentraler Konfigurationsdienst (CFG-01, siehe core-kanban/tickets/CFG-01.md). +-- scope = 'global' fuer globale Defaults, sonst der Tenant-Slug. config_values +-- haelt den AKTUELLEN Stand je (key, scope); config_value_history haelt JEDE +-- Aenderung fest (Akzeptanzkriterium 2: versioniert nachvollziehbar). +CREATE TABLE config_values ( + key TEXT NOT NULL, + scope TEXT NOT NULL CHECK (scope <> ''), + value TEXT NOT NULL, + version INT NOT NULL, + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + PRIMARY KEY (key, scope) +); + +CREATE TABLE config_value_history ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + key TEXT NOT NULL, + scope TEXT NOT NULL, + value TEXT NOT NULL, + version INT NOT NULL, + changed_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE INDEX config_value_history_key_scope_idx ON config_value_history (key, scope, version); diff --git a/migrations/0005_notification_jobs.down.sql b/migrations/0005_notification_jobs.down.sql new file mode 100644 index 0000000..1a73f4f --- /dev/null +++ b/migrations/0005_notification_jobs.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS notification_jobs; diff --git a/migrations/0005_notification_jobs.up.sql b/migrations/0005_notification_jobs.up.sql new file mode 100644 index 0000000..3bb2d5c --- /dev/null +++ b/migrations/0005_notification_jobs.up.sql @@ -0,0 +1,20 @@ +-- Benachrichtigungs-Dispatcher-Warteschlange (CFG-02, siehe +-- core-kanban/tickets/CFG-02.md). Postgres-basiert statt Redis/AMQP +-- (Projekt-Konvention, siehe nexarch-state.json techstack.job_queue) — +-- Zeilen ueberleben einen Neustart des Dispatcher-Prozesses unveraendert +-- (Akzeptanzkriterium 3). +CREATE TABLE notification_jobs ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + channel TEXT NOT NULL, + recipient TEXT NOT NULL, + payload JSONB NOT NULL DEFAULT '{}'::jsonb, + status TEXT NOT NULL DEFAULT 'pending' CHECK (status IN ('pending', 'sent', 'failed')), + attempts INT NOT NULL DEFAULT 0, + max_attempts INT NOT NULL DEFAULT 5, + next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT now(), + last_error TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE INDEX notification_jobs_due_idx ON notification_jobs (status, next_attempt_at); diff --git a/scripts/reset-test-env.sh b/scripts/reset-test-env.sh index fab5903..dc82497 100755 --- a/scripts/reset-test-env.sh +++ b/scripts/reset-test-env.sh @@ -1,11 +1,26 @@ #!/usr/bin/env bash +# Setzt die nexarch-Testumgebung zurueck: loescht die geteilte +# Registry-Tabelle "tenants" in der postgres-Wartungsdatenbank sowie alle +# tenant_*-Datenbanken. Noetig, weil verschiedene Feature-Branches +# unterschiedliche Registry-Schemata erwarten, aber dieselbe physische +# Postgres-Instanz auf dem Testhost teilen (siehe [[project-nexarch-test-infra]]). +# +# Aufruf: NEXARCH_TEST_DB_PASSWORD=... ./scripts/reset-test-env.sh set -euo pipefail + PASS="${NEXARCH_TEST_DB_PASSWORD:?Setze NEXARCH_TEST_DB_PASSWORD vor dem Aufruf}" ROLE="nexarch_test" + export PGPASSWORD="$PASS" + psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS tenants CASCADE;" +psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS audit_events CASCADE;" +psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS config_value_history CASCADE;" +psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS config_values CASCADE;" + dbs=$(psql -h localhost -U "$ROLE" -d postgres -tAc "SELECT datname FROM pg_database WHERE datname LIKE 'tenant\_%' ESCAPE '\'") for db in $dbs; do psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP DATABASE IF EXISTS \"${db}\";" done + echo "Testumgebung zurueckgesetzt: registry-tabelle + $(echo "$dbs" | grep -c . || true) tenant-datenbank(en) entfernt." diff --git a/scripts/run-checks.sh b/scripts/run-checks.sh index 1c28c3c..83f26b6 100755 --- a/scripts/run-checks.sh +++ b/scripts/run-checks.sh @@ -1,12 +1,24 @@ #!/usr/bin/env bash +# Ein-Kommando-Pruefung fuer den aktuellen Code-Stand auf dem Testhost: +# Registry+Tenant-DBs zuruecksetzen, dann build/vet/test in einem Rutsch. +# -p 1 ist Pflicht, da mehrere Pakete dieselbe physische Registry-Tabelle auf +# dem Testhost teilen (siehe [[project-nexarch-test-infra]]). +# +# Aufruf: NEXARCH_TEST_DB_PASSWORD=... ./scripts/run-checks.sh set -euo pipefail + PASS="${NEXARCH_TEST_DB_PASSWORD:?Setze NEXARCH_TEST_DB_PASSWORD vor dem Aufruf}" cd "$(dirname "$0")/.." + NEXARCH_TEST_DB_PASSWORD="$PASS" bash scripts/reset-test-env.sh + export TEST_ADMIN_DSN="postgresql://nexarch_test:${PASS}@localhost:5432/postgres?sslmode=disable" + echo "== go build ==" go build ./... + echo "== go vet ==" go vet ./... + echo "== go test (-p 1) ==" go test ./... -p 1 -count=1