// Package events: domain-event dispatch. One Emit call fans a mutation's event // to the audit log (always) and to any notification channels that reach it. // // The dispatcher is concrete — it names each channel fanout as a field and calls // it explicitly. It decides what an event supports by asserting the event // against the small interfaces below, rather than a case per concrete event // type, so it never imports the feature packages: a new event that implements // the interfaces is dispatched without any change here. package events import ( "context" "fmt" "log/slog" "atlas9.dev/c/core" "atlas9.dev/c/core/iam" "atlas9.dev/c/demo/lib" "atlas9.dev/c/demo/lib/slack" "atlas9.dev/c/demo/lib/webhooks" ) // auditable is the floor every domain event must satisfy: it names itself (so a // channel can match its subscriptions) and produces its audit entry. Emit // rejects an event that isn't auditable, so a mutation cannot be dispatched // without leaving an audit record. type auditable interface { EventType() string Audit() iam.AuditEntry } // slackRenderable is the opt-in for the Slack channel: an event implements it to // be delivered to subscribed Slack endpoints, rendered as the message it // returns. Events without it are simply not sent to Slack. type slackRenderable interface { SlackMessage() slack.Message } // webhookPayloader lets an event choose its webhook body. Without it the whole // event is marshaled; with it, WebhookPayload() is (e.g. the bare resource). type webhookPayloader interface { WebhookPayload() any } // Dispatcher fans one domain event to the audit log and the notification // channels. All three collaborators write in the caller's transaction, so build // a Dispatcher per transaction (see the factory wired in boot) and call Emit // inside the mutation's dbi.ReadWrite closure: the audit row and any queued // deliveries then commit atomically with the change that produced them. type Dispatcher struct { Audit iam.AuditStore Webhooks webhooks.Fanout Slack slack.Fanout } // Emit records the event's audit entry — the floor, which no event can skip — // then fans the event out to each channel it supports. Channels self-gate on // subscription (see the fanouts' ListActiveForEvent), so nothing is enqueued for // a tenant that hasn't enabled a channel for this event. Audit and deliveries // share the caller's transaction, so any error rolls the whole mutation back. func (d Dispatcher) Emit(ctx context.Context, tenant core.ID, event any) error { ev, ok := event.(auditable) if !ok { return fmt.Errorf("events: %T is not auditable", event) } if err := writeAudit(ctx, d.Audit, tenant, ev.Audit()); err != nil { return err } // Webhooks are universal: every event is offered to webhook subscribers. The // body is the whole event unless the event narrows it via WebhookPayload. var payload any = event if p, ok := event.(webhookPayloader); ok { payload = p.WebhookPayload() } if err := d.Webhooks.Emit(ctx, tenant, ev.EventType(), payload); err != nil { return err } // Slack is opt-in: only events that render a message are delivered there. if r, ok := event.(slackRenderable); ok { if err := d.Slack.Emit(ctx, tenant, ev.EventType(), r.SlackMessage()); err != nil { return err } } return nil } // writeAudit appends the entry in the caller's transaction, filling Tenant from // the Emit argument and Actor/Request from the request context. It mirrors the // audit helper in api_impl. func writeAudit(ctx context.Context, s iam.AuditStore, tenant core.ID, e iam.AuditEntry) error { e.Tenant = tenant if e.Actor == "" { e.Actor = iam.GetPrincipal(ctx).Subject } e.Request = lib.GetRequestID(ctx) if err := s.Append(ctx, &e); err != nil { return err } slog.DebugContext(ctx, "audit recorded", "action", e.Action, "resource", e.Resource) return nil }