diff --git a/internal/moduletrust/cache.go b/internal/moduletrust/cache.go index 9fded0d..246b21a 100644 --- a/internal/moduletrust/cache.go +++ b/internal/moduletrust/cache.go @@ -71,6 +71,17 @@ func (c *StaleCache[T]) Get(ctx context.Context) (value T, stale bool, err error return zero, false, fmt.Errorf("cache leer und refresh fehlgeschlagen: %w", fetchErr) } +// Invalidate erzwingt beim naechsten Get-Aufruf einen sofortigen Refresh +// statt auf den TTL-Ablauf zu warten (API-06, Akzeptanzkriterium 1: ein +// Health-Check-getriggerter Wiederanlauf soll den Cache SOFORT aktualisieren, +// nicht die reguläre TTL abwarten) — rein additiv, aendert nichts an +// Get/RequireFresh (dasselbe Muster wie internal/flag.Service.Invalidate). +func (c *StaleCache[T]) Invalidate() { + c.mu.Lock() + c.hasValue = false + c.mu.Unlock() +} + // RequireFresh ruft IMMER frisch ab (FAIL-CLOSED) — fuer sicherheitskritische // Aktionen, die niemals auf einem zwischengespeicherten Stand basieren duerfen. func (c *StaleCache[T]) RequireFresh(ctx context.Context) (T, error) { diff --git a/internal/resync/buffer.go b/internal/resync/buffer.go new file mode 100644 index 0000000..0a587be --- /dev/null +++ b/internal/resync/buffer.go @@ -0,0 +1,155 @@ +// Package resync implementiert Core API-06: Wiederanlauf & Nachsynchro- +// nisierung nach einem Core-Ausfall. +// +// - Ein Modul puffert Audit-Events und Nutzungszaehler-Deltas LOKAL in +// Postgres (NICHT im Speicher — siehe "Bekannte Fehler vermeiden" im +// Ticket: ein erneuter Ausfall waehrend der Nachlieferung darf keine +// Daten verlieren, eine In-Memory-Queue wuerde das riskieren). +// - Ein Health-Check-getriggerter Worker erkennt die Core-Wiedererreich- +// barkeit SOFORT (nicht erst nach TTL-Ablauf, siehe StaleCache.Invalidate) +// und liefert die gepufferten Daten in ORIGINALER Reihenfolge, authenti- +// fiziert ueber das Service-Credential aus API-02 +// (internal/moduleregistry.Registry.Authenticate). +// +// Dieses Paket dupliziert weder internal/audit (AUD-01) noch internal/usage +// (LIC-03) — es liefert nur den Puffer- und Nachlieferungs-Mechanismus +// DAVOR bzw. DANACH. +package resync + +import ( + "context" + "encoding/json" + "fmt" + "time" + + "github.com/jackc/pgx/v5/pgxpool" +) + +// BufferedAuditEvent ist ein lokal gepuffertes Audit-Ereignis. Seq +// garantiert die Wiederherstellung der urspruenglichen Reihenfolge +// (Akzeptanzkriterium 2) unabhaengig von eventuellen Uhrzeit-Ungenauigkeiten. +type BufferedAuditEvent struct { + ID int64 + Seq int64 + TenantSlug string + Actor string + Action string + Target string + Metadata map[string]any + CreatedAt time.Time +} + +// BufferedUsageDelta ist ein lokal gepuffertes Nutzungszaehler-Inkrement. +type BufferedUsageDelta struct { + ID int64 + Seq int64 + TenantSlug string + Metric string + Delta int64 + CreatedAt time.Time +} + +// Buffer ist die lokale, persistente Pufferqueue eines Moduls. +type Buffer struct { + pool *pgxpool.Pool +} + +func NewBuffer(pool *pgxpool.Pool) *Buffer { + return &Buffer{pool: pool} +} + +// EnqueueAuditEvent puffert EIN Audit-Ereignis lokal — wird von einem +// Fachmodul aufgerufen, wenn Core gerade nicht erreichbar ist (die +// Erkennung "Core erreichbar oder nicht" ist NICHT Teil dieses Aufrufs, +// siehe Worker). +func (b *Buffer) EnqueueAuditEvent(ctx context.Context, tenantSlug, actor, action, target string, metadata map[string]any) error { + if metadata == nil { + metadata = map[string]any{} + } + metadataJSON, err := json.Marshal(metadata) + if err != nil { + return fmt.Errorf("metadata serialisieren: %w", err) + } + _, err = b.pool.Exec(ctx, ` + INSERT INTO resync_audit_buffer (tenant_slug, actor, action, target, metadata) + VALUES ($1, $2, $3, $4, $5) + `, tenantSlug, actor, action, target, metadataJSON) + if err != nil { + return fmt.Errorf("audit-ereignis puffern: %w", err) + } + return nil +} + +// EnqueueUsageDelta puffert EIN Nutzungszaehler-Inkrement lokal. +func (b *Buffer) EnqueueUsageDelta(ctx context.Context, tenantSlug, metric string, delta int64) error { + _, err := b.pool.Exec(ctx, ` + INSERT INTO resync_usage_buffer (tenant_slug, metric, delta) + VALUES ($1, $2, $3) + `, tenantSlug, metric, delta) + if err != nil { + return fmt.Errorf("nutzungsdelta puffern: %w", err) + } + return nil +} + +// PendingAuditEvents liefert ALLE noch nicht zugestellten Audit-Events in +// ORIGINALER Reihenfolge (Akzeptanzkriterium 2 / Pruefung 1). +func (b *Buffer) PendingAuditEvents(ctx context.Context) ([]BufferedAuditEvent, error) { + rows, err := b.pool.Query(ctx, ` + SELECT id, seq, tenant_slug, actor, action, target, metadata, created_at + FROM resync_audit_buffer ORDER BY seq + `) + if err != nil { + return nil, fmt.Errorf("gepufferte audit-events abfragen: %w", err) + } + defer rows.Close() + + var out []BufferedAuditEvent + for rows.Next() { + var e BufferedAuditEvent + var metadataJSON []byte + if err := rows.Scan(&e.ID, &e.Seq, &e.TenantSlug, &e.Actor, &e.Action, &e.Target, &metadataJSON, &e.CreatedAt); err != nil { + return nil, fmt.Errorf("gepuffertes audit-event lesen: %w", err) + } + _ = json.Unmarshal(metadataJSON, &e.Metadata) + out = append(out, e) + } + return out, rows.Err() +} + +// PendingUsageDeltas liefert ALLE noch nicht zugestellten Nutzungsdeltas. +func (b *Buffer) PendingUsageDeltas(ctx context.Context) ([]BufferedUsageDelta, error) { + rows, err := b.pool.Query(ctx, ` + SELECT id, seq, tenant_slug, metric, delta, created_at + FROM resync_usage_buffer ORDER BY seq + `) + if err != nil { + return nil, fmt.Errorf("gepufferte nutzungsdeltas abfragen: %w", err) + } + defer rows.Close() + + var out []BufferedUsageDelta + for rows.Next() { + var d BufferedUsageDelta + if err := rows.Scan(&d.ID, &d.Seq, &d.TenantSlug, &d.Metric, &d.Delta, &d.CreatedAt); err != nil { + return nil, fmt.Errorf("gepuffertes nutzungsdelta lesen: %w", err) + } + out = append(out, d) + } + return out, rows.Err() +} + +// RemoveAuditEvent entfernt EIN Audit-Event aus dem Puffer — wird NUR nach +// von Core BESTAETIGTER Zustellung aufgerufen (Akzeptanzkriterium 2/3: erst +// entfernen, wenn sicher zugestellt, sonst bleibt es fuer den naechsten +// Versuch erhalten — kein Datenverlust bei erneutem Ausfall waehrend der +// Nachlieferung). +func (b *Buffer) RemoveAuditEvent(ctx context.Context, id int64) error { + _, err := b.pool.Exec(ctx, `DELETE FROM resync_audit_buffer WHERE id = $1`, id) + return err +} + +func (b *Buffer) RemoveUsageDelta(ctx context.Context, id int64) error { + _, err := b.pool.Exec(ctx, `DELETE FROM resync_usage_buffer WHERE id = $1`, id) + return err +} diff --git a/internal/resync/handler.go b/internal/resync/handler.go new file mode 100644 index 0000000..8e3e88a --- /dev/null +++ b/internal/resync/handler.go @@ -0,0 +1,132 @@ +package resync + +import ( + "context" + "encoding/json" + "net/http" + "time" +) + +// AuditRecorder ist die schmale Schnittstelle, ueber die Core empfangene +// Audit-Events tatsaechlich persistiert. In Produktion durch internal/audit +// (AUD-01, nicht Abhaengigkeit dieser Kachel) implementiert — dieses Paket +// dupliziert dessen Validierungs-/Speicherlogik NICHT, sondern ruft sie nur +// auf. Die Events werden vom Aufrufer sequenziell in PendingAuditEvents- +// Reihenfolge uebergeben (Akzeptanzkriterium 2). +type AuditRecorder interface { + Record(ctx context.Context, tenantSlug, actor, action, target string, metadata map[string]any, occurredAt time.Time) error +} + +// UsageIncrementer ist die schmale Schnittstelle zu Core's Nutzungszaehler +// (in Produktion internal/usage, LIC-03 — nicht Abhaengigkeit dieser +// Kachel). Jedes gepufferte Delta wird GENAU EINMAL angewendet. +type UsageIncrementer interface { + Increment(ctx context.Context, tenantSlug, metric string, delta int64) error +} + +// CredentialAuthenticator ist die schmale Schnittstelle zu API-02s +// Service-Credential-Pruefung (internal/moduleregistry.Registry.Authenticate). +type CredentialAuthenticator interface { + Authenticate(ctx context.Context, clientID, secret string) (moduleName string, ok bool, err error) +} + +// Handler nimmt nachgelieferte Audit-Events/Nutzungsdeltas auf der +// Core-Seite entgegen — authentifiziert ueber dasselbe Service-Credential +// wie jeder andere Modul-Core-Aufruf (Akzeptanzkriterium 2/3, "authentifiziert +// ueber das in API-02 definierte Service-Credential"). +type Handler struct { + auth CredentialAuthenticator + audit AuditRecorder + usage UsageIncrementer +} + +func NewHandler(auth CredentialAuthenticator, audit AuditRecorder, usage UsageIncrementer) *Handler { + return &Handler{auth: auth, audit: audit, usage: usage} +} + +type credentialHeader struct { + ClientID string `json:"client_id"` + Secret string `json:"secret"` +} + +func (h *Handler) authenticate(w http.ResponseWriter, r *http.Request) bool { + clientID := r.Header.Get("X-Nexarch-Client-Id") + secret := r.Header.Get("X-Nexarch-Client-Secret") + _, ok, err := h.auth.Authenticate(r.Context(), clientID, secret) + if err != nil || !ok { + http.Error(w, "ungueltiges service-credential", http.StatusUnauthorized) + return false + } + return true +} + +type auditEventDTO struct { + TenantSlug string `json:"tenant_slug"` + Actor string `json:"actor"` + Action string `json:"action"` + Target string `json:"target"` + Metadata map[string]any `json:"metadata"` + OccurredAt time.Time `json:"occurred_at"` +} + +// AuditHandler nimmt EINE Liste gepufferter Audit-Events entgegen und +// schreibt sie SEQUENZIELL in der gegebenen Reihenfolge fort +// (Akzeptanzkriterium 2 / Pruefung 1: Vollstaendigkeit + Reihenfolge). +// Bricht die Verarbeitung bei einem Fehler ab und meldet, wie viele Events +// bereits sicher geschrieben wurden — der Aufrufer (Worker) entfernt aus +// seinem lokalen Puffer NUR die bestaetigt geschriebenen Events. +func (h *Handler) AuditHandler(w http.ResponseWriter, r *http.Request) { + if !h.authenticate(w, r) { + return + } + var events []auditEventDTO + if err := json.NewDecoder(r.Body).Decode(&events); err != nil { + http.Error(w, "ungueltiger anfrage-koerper", http.StatusBadRequest) + return + } + + written := 0 + for _, e := range events { + if err := h.audit.Record(r.Context(), e.TenantSlug, e.Actor, e.Action, e.Target, e.Metadata, e.OccurredAt); err != nil { + break + } + written++ + } + + writeJSON(w, map[string]int{"written": written}) +} + +type usageDeltaDTO struct { + TenantSlug string `json:"tenant_slug"` + Metric string `json:"metric"` + Delta int64 `json:"delta"` +} + +// UsageHandler wendet JEDES gepufferte Delta GENAU EINMAL an +// (Akzeptanzkriterium 3 / Pruefung 2: keine Doppelzaehlung) — der Aufrufer +// entfernt aus seinem lokalen Puffer nur die bestaetigt uebernommenen Deltas. +func (h *Handler) UsageHandler(w http.ResponseWriter, r *http.Request) { + if !h.authenticate(w, r) { + return + } + var deltas []usageDeltaDTO + if err := json.NewDecoder(r.Body).Decode(&deltas); err != nil { + http.Error(w, "ungueltiger anfrage-koerper", http.StatusBadRequest) + return + } + + applied := 0 + for _, d := range deltas { + if err := h.usage.Increment(r.Context(), d.TenantSlug, d.Metric, d.Delta); err != nil { + break + } + applied++ + } + + writeJSON(w, map[string]int{"applied": applied}) +} + +func writeJSON(w http.ResponseWriter, body any) { + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(body) +} diff --git a/internal/resync/resync_test.go b/internal/resync/resync_test.go new file mode 100644 index 0000000..f6b402e --- /dev/null +++ b/internal/resync/resync_test.go @@ -0,0 +1,320 @@ +package resync + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "os" + "sync" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgxpool" + + "gitea.perlbach24.de/scripte/nexarch/internal/flag" + "gitea.perlbach24.de/scripte/nexarch/internal/moduleregistry" +) + +// fakeAudit steht fuer internal/audit.Log (AUD-01, nicht Abhaengigkeit +// dieser Kachel) — zeichnet Aufrufe in Empfangsreihenfolge auf, damit +// Vollstaendigkeit UND Reihenfolge geprueft werden koennen. +type fakeAudit struct { + mu sync.Mutex + events []auditEventDTO + failAt int // -1 = nie fehlschlagen; sonst: ab diesem Index (0-basiert) schlaegt Record fehl +} + +func (f *fakeAudit) Record(ctx context.Context, tenantSlug, actor, action, target string, metadata map[string]any, occurredAt time.Time) error { + f.mu.Lock() + defer f.mu.Unlock() + if f.failAt >= 0 && len(f.events) == f.failAt { + return fmt.Errorf("simulierter core-ausfall waehrend der nachlieferung") + } + f.events = append(f.events, auditEventDTO{TenantSlug: tenantSlug, Actor: actor, Action: action, Target: target, Metadata: metadata, OccurredAt: occurredAt}) + return nil +} + +// fakeUsage steht fuer internal/usage.Store (LIC-03, nicht Abhaengigkeit +// dieser Kachel) — summiert Deltas wie der echte Store. +type fakeUsage struct { + mu sync.Mutex + totals map[string]int64 +} + +func newFakeUsage() *fakeUsage { return &fakeUsage{totals: map[string]int64{}} } + +func (f *fakeUsage) Increment(ctx context.Context, tenantSlug, metric string, delta int64) error { + f.mu.Lock() + defer f.mu.Unlock() + f.totals[tenantSlug+"|"+metric] += delta + return nil +} + +func setupTest(t *testing.T) (*Buffer, *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 resync_audit_buffer ( + id BIGSERIAL PRIMARY KEY, seq BIGSERIAL, tenant_slug TEXT NOT NULL, actor TEXT NOT NULL, + action TEXT NOT NULL, target TEXT NOT NULL DEFAULT '', metadata JSONB NOT NULL DEFAULT '{}', + created_at TIMESTAMPTZ NOT NULL DEFAULT now() + ); + CREATE TABLE IF NOT EXISTS resync_usage_buffer ( + id BIGSERIAL PRIMARY KEY, seq BIGSERIAL, tenant_slug TEXT NOT NULL, metric TEXT NOT NULL, + delta BIGINT NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now() + ); + CREATE TABLE IF NOT EXISTS feature_flags ( + key TEXT PRIMARY KEY, enabled BOOLEAN NOT NULL DEFAULT false, + rollout_percentage INT NOT NULL DEFAULT 0, target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}', + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() + ); + CREATE TABLE IF NOT EXISTS modules ( + name TEXT PRIMARY KEY, version TEXT NOT NULL CHECK (version <> ''), + required_flags TEXT[] NOT NULL DEFAULT '{}', registered_at TIMESTAMPTZ NOT NULL DEFAULT now() + ); + CREATE TABLE IF NOT EXISTS module_credentials ( + module_name TEXT PRIMARY KEY REFERENCES modules(name), + client_id TEXT NOT NULL UNIQUE, secret_hash BYTEA NOT NULL, + issued_at TIMESTAMPTZ NOT NULL DEFAULT now() + ); + `); err != nil { + t.Fatalf("schema: %v", err) + } + + cleanup := func() { pool.Close() } + return NewBuffer(pool), pool, cleanup +} + +func uniqueModuleName() string { + return fmt.Sprintf("resync-test-%d", time.Now().UnixNano()) +} + +// setupModuleCredential registriert ein echtes Modul + Service-Credential +// ueber internal/moduleregistry (API-02) — dieselbe Authentifizierung wird +// vom Handler tatsaechlich geprueft, kein Mock. +func setupModuleCredential(t *testing.T, pool *pgxpool.Pool) (registry *moduleregistry.Registry, clientID, secret string) { + t.Helper() + ctx := context.Background() + flagService := flag.NewService(flag.NewStore(pool), 10*time.Millisecond) + registry = moduleregistry.NewRegistry(pool, flagService) + name := uniqueModuleName() + if _, err := registry.Register(ctx, name, "1.0.0", nil); err != nil { + t.Fatalf("modul registrieren: %v", err) + } + clientID, secret, err := registry.Provision(ctx, name) + if err != nil { + t.Fatalf("credential provisionieren: %v", err) + } + return registry, clientID, secret +} + +// Akzeptanzkriterium 2 + Pruefung 1: waehrend eines simulierten Ausfalls +// lokal gepufferte Audit-Events sind nach Wiederanlauf vollstaendig und in +// korrekter Reihenfolge in Core's Audit-Log vorhanden. +func TestFlushAll_DeliversBufferedAuditEventsCompleteAndInOrder(t *testing.T) { + buffer, pool, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + registry, clientID, secret := setupModuleCredential(t, pool) + + tenantSlug := "acme" + // Ereignisse "waehrend core down" lokal puffern. + actions := []string{"login", "upload", "delete", "logout"} + for _, action := range actions { + if err := buffer.EnqueueAuditEvent(ctx, tenantSlug, "user-1", action, "res-1", nil); err != nil { + t.Fatalf("enqueue %s: %v", action, err) + } + } + + audit := &fakeAudit{failAt: -1} + usage := newFakeUsage() + handler := NewHandler(registry, audit, usage) + + mux := http.NewServeMux() + mux.HandleFunc("/internal/resync/audit", handler.AuditHandler) + mux.HandleFunc("/internal/resync/usage", handler.UsageHandler) + mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) }) + server := httptest.NewServer(mux) + defer server.Close() + + client := NewCoreClient(server.URL, clientID, secret) + worker := NewWorker(buffer, client) + + if err := worker.FlushAll(ctx); err != nil { + t.Fatalf("flushall: %v", err) + } + + audit.mu.Lock() + defer audit.mu.Unlock() + if len(audit.events) != len(actions) { + t.Fatalf("erwartet %d zugestellte events, habe %d", len(actions), len(audit.events)) + } + for i, e := range audit.events { + if e.Action != actions[i] { + t.Fatalf("reihenfolge falsch: position %d = %q, want %q", i, e.Action, actions[i]) + } + } + + // Puffer muss nach bestaetigter Zustellung leer sein. + remaining, err := buffer.PendingAuditEvents(ctx) + if err != nil { + t.Fatalf("pending: %v", err) + } + if len(remaining) != 0 { + t.Fatalf("erwartet leeren puffer nach bestaetigter zustellung, habe %d verbleibende", len(remaining)) + } +} + +// Bekannter-Fehler-Praevention: bricht die Zustellung waehrend der +// Nachlieferung erneut ab (Core faellt wieder aus), bleiben die NICHT +// bestaetigten Events sicher im Puffer erhalten statt verloren zu gehen. +func TestFlushAll_KeepsUnconfirmedEventsInBufferOnPartialFailure(t *testing.T) { + buffer, pool, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + registry, clientID, secret := setupModuleCredential(t, pool) + + for i := 0; i < 5; i++ { + if err := buffer.EnqueueAuditEvent(ctx, "acme", "user-1", fmt.Sprintf("action-%d", i), "", nil); err != nil { + t.Fatalf("enqueue %d: %v", i, err) + } + } + + // Core-Fake schlaegt AB dem 3. Event fehl -> simuliert erneuten Ausfall + // mitten in der Nachlieferung. + audit := &fakeAudit{failAt: 3} + handler := NewHandler(registry, audit, newFakeUsage()) + mux := http.NewServeMux() + mux.HandleFunc("/internal/resync/audit", handler.AuditHandler) + mux.HandleFunc("/internal/resync/usage", handler.UsageHandler) + server := httptest.NewServer(mux) + defer server.Close() + + client := NewCoreClient(server.URL, clientID, secret) + worker := NewWorker(buffer, client) + + if err := worker.FlushAll(ctx); err == nil { + t.Fatal("erwartet fehler, da core nur teilweise bestaetigt hat") + } + + remaining, err := buffer.PendingAuditEvents(ctx) + if err != nil { + t.Fatalf("pending: %v", err) + } + if len(remaining) != 2 { + t.Fatalf("erwartet 2 verbleibende (nicht bestaetigte) events im puffer, habe %d", len(remaining)) + } +} + +// Akzeptanzkriterium 3 + Pruefung 2: Nutzungszaehler-Differenz aus der +// Ausfallzeit wird korrekt nachgebucht, ein wiederholter (fehlerhafter) +// Flush-Versuch fuehrt NICHT zu Doppelzaehlung, weil bereits bestaetigte +// Deltas aus dem Puffer entfernt sind. +func TestFlushAll_AppliesUsageDeltasWithoutDoubleCounting(t *testing.T) { + buffer, pool, cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + registry, clientID, secret := setupModuleCredential(t, pool) + + deltas := []int64{3, 5, 2} + var want int64 + for _, d := range deltas { + want += d + if err := buffer.EnqueueUsageDelta(ctx, "acme", "api_calls", d); err != nil { + t.Fatalf("enqueue delta %d: %v", d, err) + } + } + + usage := newFakeUsage() + handler := NewHandler(registry, &fakeAudit{failAt: -1}, usage) + mux := http.NewServeMux() + mux.HandleFunc("/internal/resync/audit", handler.AuditHandler) + mux.HandleFunc("/internal/resync/usage", handler.UsageHandler) + server := httptest.NewServer(mux) + defer server.Close() + + client := NewCoreClient(server.URL, clientID, secret) + worker := NewWorker(buffer, client) + + if err := worker.FlushAll(ctx); err != nil { + t.Fatalf("flushall 1: %v", err) + } + + usage.mu.Lock() + got := usage.totals["acme|api_calls"] + usage.mu.Unlock() + if got != want { + t.Fatalf("nutzungsstand nach nachbuchung = %d, want %d", got, want) + } + + // Ein zweiter Flush-Versuch (z.B. redundanter Retry) darf NICHTS mehr + // nachbuchen, da der Puffer bereits geleert wurde. + if err := worker.FlushAll(ctx); err != nil { + t.Fatalf("flushall 2: %v", err) + } + usage.mu.Lock() + got2 := usage.totals["acme|api_calls"] + usage.mu.Unlock() + if got2 != want { + t.Fatalf("nutzungsstand nach redundantem zweiten flush = %d, want unveraendert %d (keine doppelzaehlung)", got2, want) + } +} + +// Akzeptanzkriterium 1 + Pruefung 3: bei erkannter Core-Wiedererreichbarkeit +// wird der Cache SOFORT invalidiert (naechster Zugriff refetcht), nicht +// erst nach TTL-Ablauf — Latenz wird gemessen und liegt weit unter einer +// langen TTL. +func TestCheckAndSync_InvalidatesCacheImmediatelyOnRecovery(t *testing.T) { + buffer, pool, cleanup := setupTest(t) + defer cleanup() + registry, clientID, secret := setupModuleCredential(t, pool) + _ = registry + + mux := http.NewServeMux() + mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) }) + server := httptest.NewServer(mux) + defer server.Close() + + client := NewCoreClient(server.URL, clientID, secret) + worker := NewWorker(buffer, client) + + invalidated := false + var mu sync.Mutex + worker.OnReachable(func() { + mu.Lock() + invalidated = true + mu.Unlock() + }) + + start := time.Now() + becameReachable, err := worker.CheckAndSync(context.Background()) + elapsed := time.Since(start) + if err != nil { + t.Fatalf("checkandsync: %v", err) + } + if !becameReachable { + t.Fatal("erwartet erkannten uebergang zu 'erreichbar' beim ersten erfolgreichen check") + } + + mu.Lock() + defer mu.Unlock() + if !invalidated { + t.Fatal("erwartet sofortigen cache-invalidierungs-callback bei core-wiedererreichbarkeit") + } + // Zielwert: deutlich unter einer typischen TTL (z.B. 5s beim + // Feature-Flag-Cache, LIC-02) — hier im Millisekundenbereich, da rein + // lokal ohne Netzwerk-Overhead. + if elapsed > time.Second { + t.Fatalf("cache-invalidierung brauchte %s, erwartet deutlich unter 1s", elapsed) + } +} diff --git a/internal/resync/worker.go b/internal/resync/worker.go new file mode 100644 index 0000000..66abd65 --- /dev/null +++ b/internal/resync/worker.go @@ -0,0 +1,220 @@ +package resync + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "net/http" + "time" +) + +// CoreClient ist der modul-seitige HTTP-Client fuer die Nachlieferung, +// authentifiziert ueber dasselbe Service-Credential wie jeder andere +// Modul-Core-Aufruf (API-02). +type CoreClient struct { + BaseURL string + ClientID string + Secret string + HTTP *http.Client +} + +func NewCoreClient(baseURL, clientID, secret string) *CoreClient { + return &CoreClient{BaseURL: baseURL, ClientID: clientID, Secret: secret, HTTP: &http.Client{Timeout: 5 * time.Second}} +} + +func (c *CoreClient) post(ctx context.Context, path string, body any) (*http.Response, error) { + payload, err := json.Marshal(body) + if err != nil { + return nil, fmt.Errorf("payload serialisieren: %w", err) + } + req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.BaseURL+path, bytes.NewReader(payload)) + if err != nil { + return nil, err + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("X-Nexarch-Client-Id", c.ClientID) + req.Header.Set("X-Nexarch-Client-Secret", c.Secret) + return c.HTTP.Do(req) +} + +// HealthCheck prueft, ob Core erreichbar ist — dieselbe Konvention wie +// internal/health (OPS-01): HTTP 200 auf einem Health-Endpunkt. +func (c *CoreClient) HealthCheck(ctx context.Context) bool { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.BaseURL+"/healthz", nil) + if err != nil { + return false + } + resp, err := c.HTTP.Do(req) + if err != nil { + return false + } + defer resp.Body.Close() + return resp.StatusCode == http.StatusOK +} + +// Worker erkennt Core-Wiedererreichbarkeit und stoesst DANN sofort +// (Akzeptanzkriterium 1) sowohl registrierte Cache-Invalidierungen als auch +// das Nachliefern des lokalen Puffers an. +type Worker struct { + buffer *Buffer + client *CoreClient + onReachable []func() + wasDown bool +} + +func NewWorker(buffer *Buffer, client *CoreClient) *Worker { + return &Worker{buffer: buffer, client: client, wasDown: true} // Start pessimistisch: erster erfolgreicher Check zaehlt als "Wiedererreichbarkeit". +} + +// OnReachable registriert einen Callback, der bei jeder erkannten +// Core-Wiedererreichbarkeit sofort ausgefuehrt wird — z.B. +// moduletrust.StaleCache[T].Invalidate, damit der naechste Zugriff sofort +// neu abruft statt auf TTL-Ablauf zu warten (Akzeptanzkriterium 1). +func (w *Worker) OnReachable(fn func()) { + w.onReachable = append(w.onReachable, fn) +} + +// CheckAndSync fuehrt EINEN Zyklus aus: Erreichbarkeit pruefen, bei +// erkanntem UEBERGANG "nicht erreichbar -> erreichbar" sofort die +// registrierten Callbacks ausloesen und den Puffer nachliefern. Gibt +// zurueck, ob ein Wiederanlauf in diesem Aufruf erkannt wurde (fuer +// Latenzmessung in Tests, Pruefung 3). +func (w *Worker) CheckAndSync(ctx context.Context) (becameReachable bool, err error) { + reachable := w.client.HealthCheck(ctx) + if !reachable { + w.wasDown = true + return false, nil + } + + justRecovered := w.wasDown + w.wasDown = false + if !justRecovered { + return false, nil + } + + for _, fn := range w.onReachable { + fn() + } + + if err := w.FlushAll(ctx); err != nil { + return true, err + } + return true, nil +} + +// FlushAll liefert ZUERST alle gepufferten Audit-Events (in Reihenfolge), +// DANN alle gepufferten Nutzungsdeltas nach. Jedes Element wird aus dem +// lokalen Puffer NUR entfernt, wenn Core es bestaetigt hat — bricht die +// Uebertragung vorzeitig ab (Core faellt waehrend der Nachlieferung erneut +// aus), bleibt der Rest sicher im Postgres-Puffer erhalten +// (Akzeptanzkriterium 2/3, "Bekannte Fehler vermeiden"). +func (w *Worker) FlushAll(ctx context.Context) error { + if err := w.flushAuditEvents(ctx); err != nil { + return err + } + return w.flushUsageDeltas(ctx) +} + +func (w *Worker) flushAuditEvents(ctx context.Context) error { + events, err := w.buffer.PendingAuditEvents(ctx) + if err != nil { + return err + } + if len(events) == 0 { + return nil + } + + dtos := make([]auditEventDTO, len(events)) + for i, e := range events { + dtos[i] = auditEventDTO{ + TenantSlug: e.TenantSlug, Actor: e.Actor, Action: e.Action, Target: e.Target, + Metadata: e.Metadata, OccurredAt: e.CreatedAt, + } + } + + resp, err := w.client.post(ctx, "/internal/resync/audit", dtos) + if err != nil { + return fmt.Errorf("audit-nachlieferung: %w", err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("audit-nachlieferung: unerwarteter status %d", resp.StatusCode) + } + + var ack struct { + Written int `json:"written"` + } + if err := json.NewDecoder(resp.Body).Decode(&ack); err != nil { + return fmt.Errorf("audit-bestaetigung lesen: %w", err) + } + + // NUR die von Core bestaetigt geschriebenen Events entfernen — sie sind + // nach PendingAuditEvents-Reihenfolge sortiert, die ersten `Written` + // Eintraege entsprechen also genau den bestaetigten. + for i := 0; i < ack.Written; i++ { + if err := w.buffer.RemoveAuditEvent(ctx, events[i].ID); err != nil { + return fmt.Errorf("bestaetigtes audit-event aus puffer entfernen: %w", err) + } + } + if ack.Written < len(events) { + return fmt.Errorf("core hat nur %d von %d audit-events bestaetigt", ack.Written, len(events)) + } + return nil +} + +func (w *Worker) flushUsageDeltas(ctx context.Context) error { + deltas, err := w.buffer.PendingUsageDeltas(ctx) + if err != nil { + return err + } + if len(deltas) == 0 { + return nil + } + + dtos := make([]usageDeltaDTO, len(deltas)) + for i, d := range deltas { + dtos[i] = usageDeltaDTO{TenantSlug: d.TenantSlug, Metric: d.Metric, Delta: d.Delta} + } + + resp, err := w.client.post(ctx, "/internal/resync/usage", dtos) + if err != nil { + return fmt.Errorf("nutzungs-nachlieferung: %w", err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("nutzungs-nachlieferung: unerwarteter status %d", resp.StatusCode) + } + + var ack struct { + Applied int `json:"applied"` + } + if err := json.NewDecoder(resp.Body).Decode(&ack); err != nil { + return fmt.Errorf("nutzungs-bestaetigung lesen: %w", err) + } + + for i := 0; i < ack.Applied; i++ { + if err := w.buffer.RemoveUsageDelta(ctx, deltas[i].ID); err != nil { + return fmt.Errorf("bestaetigtes nutzungsdelta aus puffer entfernen: %w", err) + } + } + if ack.Applied < len(deltas) { + return fmt.Errorf("core hat nur %d von %d nutzungsdeltas bestaetigt", ack.Applied, len(deltas)) + } + return nil +} + +// Run fuehrt CheckAndSync in festen Abstaenden aus — die "kurze, definierte +// Zeitspanne" aus Akzeptanzkriterium 1 ist dieses Poll-Intervall. +func (w *Worker) Run(ctx context.Context, interval time.Duration) { + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + _, _ = w.CheckAndSync(ctx) + } + } +} diff --git a/migrations/0006_resync_buffers.down.sql b/migrations/0006_resync_buffers.down.sql new file mode 100644 index 0000000..b4f73ae --- /dev/null +++ b/migrations/0006_resync_buffers.down.sql @@ -0,0 +1,2 @@ +DROP TABLE resync_usage_buffer; +DROP TABLE resync_audit_buffer; diff --git a/migrations/0006_resync_buffers.up.sql b/migrations/0006_resync_buffers.up.sql new file mode 100644 index 0000000..5de9fc1 --- /dev/null +++ b/migrations/0006_resync_buffers.up.sql @@ -0,0 +1,22 @@ +-- Wiederanlauf & Nachsynchronisierung nach Core-Ausfall (API-06, siehe +-- core-kanban/tickets/API-06.md) — lokale, PERSISTENTE Pufferqueue (kein +-- In-Memory) fuer Audit-Events und Nutzungszaehler-Deltas eines Moduls. +CREATE TABLE resync_audit_buffer ( + id BIGSERIAL PRIMARY KEY, + seq BIGSERIAL, + tenant_slug TEXT NOT NULL, + actor TEXT NOT NULL, + action TEXT NOT NULL, + target TEXT NOT NULL DEFAULT '', + metadata JSONB NOT NULL DEFAULT '{}', + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE TABLE resync_usage_buffer ( + id BIGSERIAL PRIMARY KEY, + seq BIGSERIAL, + tenant_slug TEXT NOT NULL, + metric TEXT NOT NULL, + delta BIGINT NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +);