221 lines
6.5 KiB
Go
221 lines
6.5 KiB
Go
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)
|
|
}
|
|
}
|
|
}
|