// Package retentionnotify implementiert RET-07: Benachrichtigung einer // konfigurierten zuständigen Rolle (Tenant-Admin — RET-01 führt bewusst // keine Objekt-Owner-Beziehung) vor Ablauf einer Aufbewahrungsfrist, // konfigurierbarer Vorlauf je Aufbewahrungsklasse. Erzeugt NUR das // Ereignis über Core CFG-05 (archive/internal/notifyclient) — versendet // selbst keine E-Mail. package retentionnotify import ( "context" "fmt" "time" "github.com/jackc/pgx/v5/pgxpool" "gitea.perlbach24.de/scripte/nexarch/archive/internal/notifyclient" "gitea.perlbach24.de/scripte/nexarch/archive/internal/retentionengine" ) // EventType ist der an CFG-05 übergebene Ereignistyp — muss mit dem im // Frontend/Core bekannten Namen übereinstimmen (siehe CFG-05-Tests). const EventType = "retention_due_soon" // Recipient benennt Tenant-Slug, User-ID und E-Mail-Adresse der // konfigurierten zuständigen Rolle (Tenant-Admin), an die alle // Fristablauf-Benachrichtigungen dieses Tenants gehen. type Recipient struct { TenantSlug string UserID string Email string } // Result ist das Ergebnis EINES benachrichtigten (oder fehlgeschlagenen) // Objekts — der Aufrufer (cmd/retention-notify-job) protokolliert Err // explizit, kein stillschweigendes Verwerfen (Pflichtprüfung 3). type Result struct { RetentionObjectID string RetentionClass string JobID string Skipped bool Err error } type classRuleLead struct { leadDays int enabled bool } // Run führt EINEN Durchlauf des Benachrichtigungs-Jobs aus: ermittelt je // aktiver, benachrichtigungs-aktivierter Aufbewahrungsklasse die Objekte, // deren Stichtag innerhalb des konfigurierten Vorlaufs liegt, überspringt // bereits benachrichtigte Objekte (Akzeptanzkriterium 2, Postgres- // persistent — übersteht einen Job-Neustart) und löst für den Rest je ein // Ereignis über CFG-05 aus. func Run(ctx context.Context, pool *pgxpool.Pool, client *notifyclient.Client, now time.Time, recipient Recipient) ([]Result, error) { // notify_lead_days/notify_enabled sind nicht Teil von // retentionengine.ClassRule (RET-02/RET-06 kennen sie nicht) — direkt // gelesen, um retentionengine nicht um RET-07-eigene Felder zu // erweitern (kein Umbau angrenzender Bereiche). leadByClass := make(map[string]classRuleLead) maxLeadDays := 0 leadRows, err := pool.Query(ctx, `SELECT retention_class, notify_lead_days, notify_enabled FROM retention_class_rules WHERE active`) if err != nil { return nil, fmt.Errorf("retentionnotify: benachrichtigungs-konfiguration laden: %w", err) } for leadRows.Next() { var class string var lead classRuleLead if err := leadRows.Scan(&class, &lead.leadDays, &lead.enabled); err != nil { leadRows.Close() return nil, fmt.Errorf("retentionnotify: konfigurationszeile lesen: %w", err) } leadByClass[class] = lead if lead.leadDays > maxLeadDays { maxLeadDays = lead.leadDays } } if err := leadRows.Err(); err != nil { return nil, fmt.Errorf("retentionnotify: benachrichtigungs-konfiguration lesen: %w", err) } leadRows.Close() if maxLeadDays == 0 { return nil, nil } // Nutzt DIESELBE Funktion wie der RET-02-Job/RET-06-API-Preview // (kein zweiter Ermittlungspfad) — asOf auf den größten konfigurierten // Vorlauf gesetzt, je Klasse wird unten mit deren EIGENEM Vorlauf // gefiltert. candidates, err := retentionengine.ListExpiringObjects(ctx, pool, now.AddDate(0, 0, maxLeadDays)) if err != nil { return nil, fmt.Errorf("retentionnotify: ablaufende objekte ermitteln: %w", err) } alreadyNotified, err := loadAlreadyNotified(ctx, pool) if err != nil { return nil, err } var results []Result for _, obj := range candidates { lead, known := leadByClass[obj.RetentionClass] if !known || !lead.enabled { continue } if !obj.DueDate.Before(now.AddDate(0, 0, lead.leadDays+1)) { // Ausserhalb des klassen-eigenen Vorlaufs (nur mit dem // globalen maxLeadDays vorselektiert). continue } if alreadyNotified[obj.RetentionObjectID] { continue } res := Result{RetentionObjectID: obj.RetentionObjectID, RetentionClass: obj.RetentionClass} enq, err := client.Enqueue(ctx, recipient.TenantSlug, recipient.UserID, EventType, "email", recipient.Email, map[string]any{ "object_type": obj.ObjectType, "object_reference": obj.ObjectReference, "retention_class": obj.RetentionClass, "due_date": obj.DueDate.Format(time.RFC3339), }) if err != nil { res.Err = err results = append(results, res) // Kein INSERT in retention_notifications bei Fehler — das // Objekt wird beim naechsten Durchlauf erneut versucht, // statt stillschweigend als erledigt zu gelten. continue } res.JobID = enq.JobID res.Skipped = enq.Skipped if _, err := pool.Exec(ctx, `INSERT INTO retention_notifications (retention_object_id) VALUES ($1)`, obj.RetentionObjectID); err != nil { res.Err = fmt.Errorf("retentionnotify: benachrichtigung als versendet markieren: %w", err) } results = append(results, res) } return results, nil } func loadAlreadyNotified(ctx context.Context, pool *pgxpool.Pool) (map[string]bool, error) { rows, err := pool.Query(ctx, `SELECT retention_object_id FROM retention_notifications`) if err != nil { return nil, fmt.Errorf("retentionnotify: bereits benachrichtigte objekte laden: %w", err) } defer rows.Close() out := make(map[string]bool) for rows.Next() { var id string if err := rows.Scan(&id); err != nil { return nil, fmt.Errorf("retentionnotify: zeile lesen: %w", err) } out[id] = true } return out, rows.Err() }