Skip to content

Audit Log & Redis Streams (Server)

[!NOTE] L’audit est une trace durable des actions métier (qui, quoi, quand, sur quelle ressource), distincte des logs techniques. Il est asynchrone pour ne pas allonger le temps de réponse des commandes ; PostgreSQL reste la source de vérité consultée par l’API.

Berth utilise déjà Redis pour le cache et le déploiement self-hosted vise un volume d’actions modéré. Redis Streams apporte une persistance suffisante pour ce cas, des consumer groups, des accusés de réception et une relecture des messages, sans ajouter un service à opérer. RabbitMQ, Pulsar et Redpanda deviennent pertinents lorsque le débit, le routage, la rétention longue ou l’écosystème événementiel dépassent ce besoin.

Critère Redis Streams RabbitMQ Apache Pulsar Redpanda
Intégration Berth Redis déjà présent (cache + stream) Service supplémentaire Cluster supplémentaire, plus lourd Cluster supplémentaire, Kafka-compatible
Modèle Log persistant + consumer groups Files, exchanges, routage riche Topics/partitions, stockage segmenté Topics/partitions, API Kafka
Opérations Faible complexité pour une installation Moyenne ; tuning DLQ/routage Élevée ; brokers + bookies selon déploiement Moyenne à élevée ; cluster à maintenir
Cas adapté Audit modéré, traitement batch, self-hosted Routage complexe, RPC asynchrone, workloads variés Très gros volume, multi-tenant, rétention longue Gros débit Kafka et compatibilité Kafka
Choix pour Berth Retenu Surdimensionné à ce stade Surdimensionné à ce stade Surdimensionné à ce stade

[!TIP] Le stream ne remplace pas une transaction outbox si la publication de l’audit doit être strictement atomique avec la mutation métier. Pour un audit best-effort à faible latence, le decorator publie après le succès de la commande et le worker assure la durabilité côté PostgreSQL.

Les valeurs payload et metadata sont conservées en jsonb afin de permettre l’évolution du schéma sans migration à chaque ajout de détail. user_id peut être nul pour une action système.

-- migrations/000X_create_audit_logs.up.sql
CREATE TABLE audit_logs (
id uuid PRIMARY KEY,
user_id uuid NULL,
action varchar(100) NOT NULL,
resource_type varchar(100) NOT NULL,
resource_id uuid NULL,
payload jsonb NOT NULL DEFAULT '{}'::jsonb,
metadata jsonb NOT NULL DEFAULT '{}'::jsonb,
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE INDEX idx_audit_logs_user_resource_created
ON audit_logs (user_id, resource_type, created_at DESC);
CREATE INDEX idx_audit_logs_resource_created
ON audit_logs (resource_type, resource_id, created_at DESC);
CREATE INDEX idx_audit_logs_action_created
ON audit_logs (action, created_at DESC);
CREATE INDEX idx_audit_logs_created_at ON audit_logs (created_at DESC);
Command Handler (métier)
│ succès uniquement
Audit Decorator (publication non bloquante)
│ XADD audit:events
Redis Stream + Consumer Group (audit-writers)
│ XREADGROUP, batch de 100
Worker : validation → Bulk insert PostgreSQL
│ succès du batch
XACK (chaque message traité)

La commande ne doit jamais attendre l’insertion PostgreSQL de l’audit. Une erreur XADD est journalisée et métriquée ; selon la criticité, l’installation peut activer une outbox transactionnelle dans une évolution ultérieure. Le worker est idempotent : id est la clé primaire et un doublon est ignoré ou traité comme déjà écrit avant XACK.

Contrat d’événement et decorator générique

Section titled “Contrat d’événement et decorator générique”

Le contexte d’authentification transporte l’utilisateur courant, tandis que le decorator ne dépend pas de Gin. Le CommandHandler générique suit le style Go 1.22+ et peut envelopper n’importe quel handler métier.

internal/audit/decorator.go
package audit
import (
"context"
"encoding/json"
"github.com/google/uuid"
)
type CommandHandler[C any, R any] interface {
Handle(context.Context, C) (R, error)
}
type Event struct {
ID uuid.UUID `json:"id"`
UserID *uuid.UUID `json:"user_id,omitempty"`
Action string `json:"action"`
ResourceType string `json:"resource_type"`
ResourceID *uuid.UUID `json:"resource_id,omitempty"`
Payload map[string]any `json:"payload,omitempty"`
Metadata map[string]any `json:"metadata,omitempty"`
}
type Publisher interface { Publish(context.Context, Event) error }
type Decorated[C any, R any] struct { next CommandHandler[C, R]; publisher Publisher; describe func(C, R) Event }
func NewDecorator[C any, R any](next CommandHandler[C, R], p Publisher, describe func(C, R) Event) *Decorated[C, R] {
return &Decorated[C, R]{next: next, publisher: p, describe: describe}
}
func (d *Decorated[C, R]) Handle(ctx context.Context, cmd C) (R, error) {
result, err := d.next.Handle(ctx, cmd)
if err != nil { return result, err }
event := d.describe(cmd, result)
if raw, ok := ctx.Value(userIDKey{}).(uuid.UUID); ok { event.UserID = &raw }
if event.ID == uuid.Nil { event.ID = uuid.New() }
if err := d.publisher.Publish(context.WithoutCancel(ctx), event); err != nil {
// La commande est déjà réussie : ne pas la faire échouer à cause de l'audit.
}
return result, nil
}

XREADGROUP laisse les messages non acquittés dans la Pending Entries List (PEL). Un worker de récupération inspecte les entrées anciennes avec XPENDING/XPendingExt, les réclame via XCLAIM et applique un backoff.

L’audit se lit directement dans PostgreSQL pour éviter qu’une commande récente reste masquée par un cache. Un cache TTL court (par exemple 5 secondes) est possible pour une page très consultée, mais il doit être invalidé ou accepté comme légèrement eventual-consistent.