// 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 }