From 3fbaa65eeaeab008e3f867e742539c22c5b5f0c4 Mon Sep 17 00:00:00 2001 From: Codex Date: Fri, 7 Aug 2026 20:50:26 +0000 Subject: [PATCH] feat(observability): add notification and audit services --- internal/audit/service.go | 178 +++++++++ internal/audit/service_test.go | 55 +++ internal/notification/service.go | 425 ++++++++++++++++++++++ internal/notification/service_test.go | 84 +++++ internal/persistence/sqlite/store_test.go | 10 +- migrations/0009_notifications_audit.sql | 40 ++ 6 files changed, 787 insertions(+), 5 deletions(-) create mode 100644 internal/audit/service.go create mode 100644 internal/audit/service_test.go create mode 100644 internal/notification/service.go create mode 100644 internal/notification/service_test.go create mode 100644 migrations/0009_notifications_audit.sql diff --git a/internal/audit/service.go b/internal/audit/service.go new file mode 100644 index 0000000..93b6f3e --- /dev/null +++ b/internal/audit/service.go @@ -0,0 +1,178 @@ +// Package audit stores the deliberately small, redacted security audit trail. +package audit + +import ( + "context" + "crypto/rand" + "database/sql" + "encoding/base64" + "encoding/json" + "errors" + "fmt" + "strings" + "time" +) + +type Event struct { + ID string `json:"id"` + ActorID string `json:"actor_id"` + ActorLabel string `json:"actor_label"` + InstanceID string `json:"instance_id"` + Action string `json:"action"` + Outcome string `json:"outcome"` + OccurredAt time.Time `json:"occurred_at"` + Summary map[string]string `json:"summary"` +} + +type Filter struct { + ActorID, InstanceID, Action, Outcome string + Since, Until time.Time + Limit int +} + +type Policy struct { + RetentionDays int `json:"retention_days"` + MaximumCount int `json:"maximum_count"` +} + +type Service struct { + db *sql.DB + now func() time.Time +} + +func New(db *sql.DB) *Service { return &Service{db: db, now: time.Now} } + +var allowedSummaryKeys = map[string]bool{"target_id": true, "target_name": true, "channel_type": true, "event_type": true, "reason_code": true, "deleted_count": true, "before": true, "after": true} + +func (s *Service) Record(ctx context.Context, event Event) error { + if event.Action == "" || (event.Outcome != "allowed" && event.Outcome != "denied" && event.Outcome != "failed") { + return errors.New("invalid audit event") + } + clean := map[string]string{} + for key, value := range event.Summary { + if allowedSummaryKeys[key] && len(value) <= 200 { + clean[key] = value + } + } + body, _ := json.Marshal(clean) + when := event.OccurredAt + if when.IsZero() { + when = s.now().UTC() + } + var actor, instance any + if event.ActorID != "" { + actor = event.ActorID + } + if event.InstanceID != "" { + instance = event.InstanceID + } + _, err := s.db.ExecContext(ctx, `INSERT INTO audit_events(id,occurred_at,actor_id,actor_label,instance_id,action,outcome,summary_json) VALUES(?,?,?,?,?,?,?,?)`, randomID(), when.Format(time.RFC3339Nano), actor, event.ActorLabel, instance, event.Action, event.Outcome, string(body)) + if err != nil { + return fmt.Errorf("record audit event: %w", err) + } + return nil +} + +func (s *Service) List(ctx context.Context, f Filter) ([]Event, error) { + limit := f.Limit + if limit <= 0 || limit > 200 { + limit = 100 + } + clauses, args := []string{"1=1"}, []any{} + for _, item := range []struct{ column, value string }{{"actor_id", f.ActorID}, {"instance_id", f.InstanceID}, {"action", f.Action}, {"outcome", f.Outcome}} { + if item.value != "" { + clauses = append(clauses, item.column+"=?") + args = append(args, item.value) + } + } + if !f.Since.IsZero() { + clauses = append(clauses, "occurred_at>=?") + args = append(args, f.Since.UTC().Format(time.RFC3339Nano)) + } + if !f.Until.IsZero() { + clauses = append(clauses, "occurred_at 3650 || p.MaximumCount < 0 || p.MaximumCount > 1000000 { + return errors.New("audit policy is out of bounds") + } + body, _ := json.Marshal(map[string]int{"retention_days": p.RetentionDays, "maximum_count": p.MaximumCount}) + _, err := s.db.ExecContext(ctx, `UPDATE system_settings SET value_json=?,revision=revision+1,updated_at=? WHERE key='audit_policy'`, body, s.now().UTC().Format(time.RFC3339Nano)) + return err +} +func (s *Service) Purge(ctx context.Context, before time.Time) (int64, error) { + if before.IsZero() || before.After(s.now().UTC()) { + return 0, errors.New("invalid audit purge boundary") + } + result, err := s.db.ExecContext(ctx, `DELETE FROM audit_events WHERE occurred_at < ?`, before.UTC().Format(time.RFC3339Nano)) + if err != nil { + return 0, err + } + return result.RowsAffected() +} +func (s *Service) RunRetention(ctx context.Context) (int64, error) { + p, err := s.Policy(ctx) + if err != nil { + return 0, err + } + var total int64 + if p.RetentionDays > 0 { + n, e := s.Purge(ctx, s.now().UTC().AddDate(0, 0, -p.RetentionDays)) + if e != nil { + return 0, e + } + total += n + } + if p.MaximumCount > 0 { + result, e := s.db.ExecContext(ctx, `DELETE FROM audit_events WHERE id IN (SELECT id FROM audit_events ORDER BY occurred_at DESC,id DESC LIMIT -1 OFFSET ?)`, p.MaximumCount) + if e != nil { + return total, e + } + n, _ := result.RowsAffected() + total += n + } + return total, nil +} +func randomID() string { + b := make([]byte, 18) + _, _ = rand.Read(b) + return base64.RawURLEncoding.EncodeToString(b) +} diff --git a/internal/audit/service_test.go b/internal/audit/service_test.go new file mode 100644 index 0000000..69c81b6 --- /dev/null +++ b/internal/audit/service_test.go @@ -0,0 +1,55 @@ +package audit_test + +import ( + "context" + "path/filepath" + "testing" + "time" + + "git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/audit" + "git.zaynet.fr/DoGaMa/DoGaMa-serv/internal/persistence/sqlite" +) + +func TestRecordFiltersSummaryAndRetention(t *testing.T) { + ctx := context.Background() + db, err := sqlite.Open(ctx, filepath.Join(t.TempDir(), "dogama.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + service := audit.New(db) + old := time.Now().UTC().AddDate(0, 0, -40) + if err := service.Record(ctx, audit.Event{OccurredAt: old, ActorLabel: "admin", Action: "instance.update", Outcome: "allowed", Summary: map[string]string{"target_name": "Palworld", "secret": "must-not-persist"}}); err != nil { + t.Fatal(err) + } + events, err := service.List(ctx, audit.Filter{Action: "instance.update"}) + if err != nil || len(events) != 1 { + t.Fatalf("events=%#v err=%v", events, err) + } + if events[0].Summary["target_name"] != "Palworld" || events[0].Summary["secret"] != "" { + t.Fatalf("summary was not allow-listed: %#v", events[0].Summary) + } + if err := service.SetPolicy(ctx, audit.Policy{RetentionDays: 30, MaximumCount: 100}); err != nil { + t.Fatal(err) + } + deleted, err := service.RunRetention(ctx) + if err != nil || deleted != 1 { + t.Fatalf("deleted=%d err=%v", deleted, err) + } +} + +func TestPolicyAndPurgeBounds(t *testing.T) { + ctx := context.Background() + db, err := sqlite.Open(ctx, filepath.Join(t.TempDir(), "dogama.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + service := audit.New(db) + if err := service.SetPolicy(ctx, audit.Policy{RetentionDays: -1}); err == nil { + t.Fatal("negative retention accepted") + } + if _, err := service.Purge(ctx, time.Now().Add(time.Hour)); err == nil { + t.Fatal("future purge accepted") + } +} diff --git a/internal/notification/service.go b/internal/notification/service.go new file mode 100644 index 0000000..eec5bb5 --- /dev/null +++ b/internal/notification/service.go @@ -0,0 +1,425 @@ +// Package notification manages encrypted channels and bounded asynchronous delivery. +package notification + +import ( + "bytes" + "context" + "crypto/aes" + "crypto/cipher" + "crypto/hmac" + "crypto/rand" + "crypto/sha256" + "crypto/tls" + "database/sql" + "encoding/base64" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "net" + "net/http" + "net/netip" + "net/smtp" + "net/url" + "strconv" + "strings" + "time" +) + +type Channel struct { + ID, Name, Type string + Enabled bool + Events []string + Configured bool +} +type Input struct { + Name, Type string + Enabled bool + Events []string + Config map[string]string +} +type Event struct{ Type, Title, Message, InstanceName, OperationID string } +type Service struct { + db *sql.DB + aead cipher.AEAD + client *http.Client + resolver *net.Resolver + now func() time.Time +} + +func New(db *sql.DB, key []byte) (*Service, error) { + if len(key) != 32 { + return nil, errors.New("notification encryption key must be exactly 32 bytes") + } + block, err := aes.NewCipher(key) + if err != nil { + return nil, err + } + aead, err := cipher.NewGCM(block) + if err != nil { + return nil, err + } + s := &Service{db: db, aead: aead, resolver: net.DefaultResolver, now: time.Now} + transport := http.DefaultTransport.(*http.Transport).Clone() + transport.DialContext = s.dialSafe + s.client = &http.Client{Transport: transport, Timeout: 10 * time.Second, CheckRedirect: func(req *http.Request, via []*http.Request) error { + if len(via) >= 3 { + return errors.New("too many redirects") + } + return s.validateURL(req.Context(), req.URL) + }} + return s, nil +} + +func (s *Service) Upsert(ctx context.Context, id string, in Input) (Channel, error) { + if strings.TrimSpace(in.Name) == "" || !validType(in.Type) { + return Channel{}, errors.New("invalid notification channel") + } + if err := validateConfigShape(in.Type, in.Config); err != nil { + return Channel{}, err + } + encrypted, err := s.seal(in.Config) + if err != nil { + return Channel{}, err + } + events, _ := json.Marshal(normalizeEvents(in.Events)) + now := s.now().UTC().Format(time.RFC3339Nano) + if id == "" { + id = randomID() + } + _, err = s.db.ExecContext(ctx, `INSERT INTO notification_channels(id,name,type,enabled,encrypted_config,event_filter_json,created_at,updated_at) VALUES(?,?,?,?,?,?,?,?) ON CONFLICT(id) DO UPDATE SET name=excluded.name,type=excluded.type,enabled=excluded.enabled,encrypted_config=excluded.encrypted_config,event_filter_json=excluded.event_filter_json,updated_at=excluded.updated_at`, id, strings.TrimSpace(in.Name), in.Type, in.Enabled, encrypted, string(events), now, now) + if err != nil { + return Channel{}, fmt.Errorf("save notification channel: %w", err) + } + return Channel{ID: id, Name: strings.TrimSpace(in.Name), Type: in.Type, Enabled: in.Enabled, Events: normalizeEvents(in.Events), Configured: true}, nil +} +func (s *Service) List(ctx context.Context) ([]Channel, error) { + rows, err := s.db.QueryContext(ctx, `SELECT id,name,type,enabled,event_filter_json,length(encrypted_config)>0 FROM notification_channels ORDER BY name COLLATE NOCASE`) + if err != nil { + return nil, err + } + defer rows.Close() + var out []Channel + for rows.Next() { + var c Channel + var body string + if err := rows.Scan(&c.ID, &c.Name, &c.Type, &c.Enabled, &body, &c.Configured); err != nil { + return nil, err + } + _ = json.Unmarshal([]byte(body), &c.Events) + out = append(out, c) + } + return out, rows.Err() +} +func (s *Service) Delete(ctx context.Context, id string) error { + _, err := s.db.ExecContext(ctx, `DELETE FROM notification_channels WHERE id=?`, id) + return err +} +func (s *Service) Queue(ctx context.Context, event Event) error { + if !validEvent(event.Type) || len(event.Message) > 1000 { + return errors.New("invalid notification event") + } + payload, _ := json.Marshal(event) + rows, err := s.db.QueryContext(ctx, `SELECT id,event_filter_json FROM notification_channels WHERE enabled=1`) + if err != nil { + return err + } + var channelIDs []string + for rows.Next() { + var id, filter string + if err := rows.Scan(&id, &filter); err != nil { + rows.Close() + return err + } + var events []string + _ = json.Unmarshal([]byte(filter), &events) + if matches(events, event.Type) { + channelIDs = append(channelIDs, id) + } + } + if err := rows.Err(); err != nil { + rows.Close() + return err + } + rows.Close() + now := s.now().UTC().Format(time.RFC3339Nano) + for _, id := range channelIDs { + if _, err := s.db.ExecContext(ctx, `INSERT INTO notification_deliveries(id,channel_id,event_type,payload_redacted,next_attempt_at,created_at) VALUES(?,?,?,?,?,?)`, randomID(), id, event.Type, string(payload), now, now); err != nil { + return err + } + } + return nil +} +func (s *Service) Test(ctx context.Context, id string) error { + return s.queueForChannel(ctx, id, Event{Type: "notification.test", Title: "DoGaMa test notification", Message: "This is a test notification from DoGaMa."}) +} +func (s *Service) queueForChannel(ctx context.Context, id string, event Event) error { + var enabled bool + if err := s.db.QueryRowContext(ctx, `SELECT enabled FROM notification_channels WHERE id=?`, id).Scan(&enabled); err != nil { + return err + } + payload, _ := json.Marshal(event) + now := s.now().UTC().Format(time.RFC3339Nano) + _, err := s.db.ExecContext(ctx, `INSERT INTO notification_deliveries(id,channel_id,event_type,payload_redacted,next_attempt_at,created_at) VALUES(?,?,?,?,?,?)`, randomID(), id, event.Type, string(payload), now, now) + return err +} + +func (s *Service) RunDue(ctx context.Context) error { + rows, err := s.db.QueryContext(ctx, `SELECT d.id,d.attempt,d.payload_redacted,c.type,c.encrypted_config FROM notification_deliveries d JOIN notification_channels c ON c.id=d.channel_id WHERE d.status IN ('queued','retrying') AND d.next_attempt_at<=? ORDER BY d.next_attempt_at LIMIT 20`, s.now().UTC().Format(time.RFC3339Nano)) + if err != nil { + return err + } + type job struct { + id, typ, payload string + attempt int + encrypted []byte + } + var jobs []job + for rows.Next() { + var j job + if err := rows.Scan(&j.id, &j.attempt, &j.payload, &j.typ, &j.encrypted); err != nil { + rows.Close() + return err + } + jobs = append(jobs, j) + } + rows.Close() + for _, j := range jobs { + config, e := s.open(j.encrypted) + if e == nil { + e = s.deliver(ctx, j.id, j.typ, config, []byte(j.payload)) + } + attempt := j.attempt + 1 + if e == nil { + if _, updateErr := s.db.ExecContext(ctx, `UPDATE notification_deliveries SET status='succeeded',attempt=?,completed_at=?,last_error_code='' WHERE id=?`, attempt, s.now().UTC().Format(time.RFC3339Nano), j.id); updateErr != nil { + return updateErr + } + } else { + status := "retrying" + if attempt >= 5 { + status = "failed" + } + delay := time.Duration(1<= 0), + next_attempt_at TEXT NOT NULL, + last_error_code TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL, + completed_at TEXT +); +CREATE INDEX notification_deliveries_due_idx ON notification_deliveries(status, next_attempt_at); + +CREATE TABLE audit_events ( + id TEXT PRIMARY KEY, + occurred_at TEXT NOT NULL, + actor_id TEXT REFERENCES users(id) ON DELETE SET NULL, + actor_label TEXT NOT NULL, + instance_id TEXT REFERENCES instances(id) ON DELETE SET NULL, + action TEXT NOT NULL, + outcome TEXT NOT NULL CHECK (outcome IN ('allowed', 'denied', 'failed')), + summary_json TEXT NOT NULL DEFAULT '{}' +); +CREATE INDEX audit_events_time_idx ON audit_events(occurred_at DESC, id DESC); +CREATE INDEX audit_events_filters_idx ON audit_events(actor_id, instance_id, action, outcome); + +INSERT INTO system_settings(key, value_json, revision, updated_at) +VALUES ('audit_policy', '{"retention_days":30,"maximum_count":10000}', 1, strftime('%Y-%m-%dT%H:%M:%fZ','now'));