package webhooks import ( "context" "encoding/json" "atlas9.dev/c/core" "atlas9.dev/c/demo/lib/access" ) // Delivery is one queued webhook delivery: which endpoint to POST to and the // event body to sign and send. It is the task payload the delivery worker // replays; api.Webhooks_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 webhook deliveries. type DeliveryQueue interface { Create(ctx context.Context, d Delivery) error } // Fanout enqueues webhook 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 — a crash can't leave a // notification half-sent. 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. event // is the payload; Emit marshals it to JSON once and each delivery signs and // sends those bytes verbatim. Callers pass the domain object directly — there is // no per-feature marshal-and-emit boilerplate. Marshaling is skipped when the // tenant has no subscribers for the event. func (f Fanout) Emit(ctx context.Context, tenant core.ID, eventType string, event any) error { // The subscription lookup is on the app's behalf, not the caller's: the // emitter is mid-mutation and needn't hold webhook 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_Webhooks_ReadEndpoint) eps, err := f.Endpoints.ListActiveForEvent(readCtx, tenant, eventType) if err != nil { return err } if len(eps) == 0 { return nil } body, err := json.Marshal(event) 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 }