API-06: wiederanlauf-nachsynchronisierung-nach-core-ausfall (postgres-puffer, service-credential, sofort-invalidate)

This commit is contained in:
sysops
2026-08-28 09:22:08 +02:00
parent b23cd1961f
commit 0fd9856b10
11 changed files with 1391 additions and 0 deletions
+155
View File
@@ -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
}
+132
View File
@@ -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)
}
+320
View File
@@ -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)
}
}
+220
View File
@@ -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)
}
}
}