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