package slack import ( "context" "encoding/json" "atlas9.dev/c/core" "atlas9.dev/c/demo/lib/access" ) // Delivery is one queued Slack delivery: which endpoint to POST to and the // already-rendered Slack message body to send. It is the task payload the // delivery worker replays; api.Slack_DeliverReq mirrors these fields so the // worker decodes the same JSON. It lives here, not in api, so fanout can stay in // the library (lib cannot import api). type Delivery struct { Tenant core.ID Endpoint core.ID Type string Body []byte } // DeliveryQueue enqueues Slack deliveries. type DeliveryQueue interface { Create(ctx context.Context, d Delivery) error } // Fanout enqueues Slack deliveries for domain events. Call Emit from a mutation // handler, in the same transaction, so the deliveries commit atomically with the // change that triggered them. Both collaborators are transaction-bound, so build // a Fanout per transaction (see the factory wired in boot). type Fanout struct { Endpoints Store Deliveries DeliveryQueue } // Emit enqueues one delivery per active endpoint subscribed to eventType. msg is // the Slack message the caller authors for this event — plain text or full Block // Kit — the per-event template lives at the call site, and Emit only marshals it // to Slack's contract once and sends those bytes to each endpoint verbatim. The // marshal and enqueue are skipped when the tenant has no subscribers. func (f Fanout) Emit(ctx context.Context, tenant core.ID, eventType string, msg Message) error { // TODO rethink this kind of delegation or escalation or whatever its called. // The subscription lookup is on the app's behalf, not the caller's: the // emitter is mid-mutation and needn't hold Slack caps. Grant the read scoped // to this tenant for the lookup, the way the create path grants Dek_Use for // sealing. readCtx := access.PutScope(ctx, tenant, Cap_Slack_ReadEndpoint) eps, err := f.Endpoints.ListActiveForEvent(readCtx, tenant, eventType) if err != nil { return err } if len(eps) == 0 { return nil } // Message.MarshalJSON emits Slack's incoming-webhook contract (lowercase // "text"/"blocks"). body, err := json.Marshal(msg) if err != nil { return err } for _, ep := range eps { err := f.Deliveries.Create(ctx, Delivery{ Tenant: tenant, Endpoint: ep.ID, Type: eventType, Body: body, }) if err != nil { return err } } return nil }