package api_impl import ( "crypto/rand" "database/sql" "encoding/hex" "errors" "log/slog" "net/http" "atlas9.dev/c/core" "atlas9.dev/c/core/dbi" "atlas9.dev/c/core/envelope" "atlas9.dev/c/core/iam" "atlas9.dev/c/demo/api" "atlas9.dev/c/demo/lib/access" "atlas9.dev/c/demo/lib/tasks" "atlas9.dev/c/demo/lib/webhooks" ) type WebhooksImpl struct { DB *sql.DB Guard access.Guard Store dbi.Factory[webhooks.Store] Audit dbi.Factory[iam.AuditStore] Encryptors *envelope.EncryptorFactory Deks dbi.Factory[envelope.DekStore] // Sender is the guarded egress layer used by the internal Deliver endpoint. Sender *webhooks.Sender } func (s *WebhooksImpl) ServeMux(mux *http.ServeMux) { mux.HandleFunc(api.Path_Webhooks_CreateEndpoint, s.CreateEndpoint) mux.HandleFunc(api.Path_Webhooks_UpdateEndpoint, s.UpdateEndpoint) mux.HandleFunc(api.Path_Webhooks_GetEndpoint, s.GetEndpoint) mux.HandleFunc(api.Path_Webhooks_ListEndpoints, s.ListEndpoints) mux.HandleFunc(api.Path_Webhooks_DeleteEndpoint, s.DeleteEndpoint) mux.HandleFunc(api.Path_Webhooks_Deliver, s.Deliver) } func (s *WebhooksImpl) CreateEndpoint(w http.ResponseWriter, r *http.Request) { var req api.Webhooks_CreateEndpointReq if read(w, r, &req) { return } if check(w, r, s.Guard, webhooks.Cap_Webhooks_CreateEndpoint, req.Endpoint.Tenant, "") { return } ctx := r.Context() // The signing secret is generated server-side and sealed under the tenant's // DEK; only the plaintext returned here can reproduce the HMAC signature. secret, err := newSecret() if writeErr(ctx, w, err) { return } err = dbi.ReadWrite(ctx, s.DB, func(tx dbi.DBI) error { // TODO might be interesting to build these rules into the access system. // Cedar sort of did this. Like, "if route is create webhook and caller has that cap, then also allow use dek". // // Sealing needs the tenant DEK; grant Dek_Use scoped to this tenant for // that sub-operation (the caller holds the webhook cap, not Dek_Use). sealCtx := access.PutScope(ctx, req.Endpoint.Tenant, envelope.Cap_Dek_Use) // TODO there are too many names involved with secrets: envelope, dek, encryptor, seal, etc. enc, err := s.Encryptors.For(sealCtx, s.Deks(tx), req.Endpoint.Tenant) if err != nil { return err } dekID, sealed, err := enc.Seal([]byte(secret)) if err != nil { return err } req.Endpoint.DekID = dekID req.Endpoint.SecretEnc = sealed if err := s.Store(tx).CreateEndpoint(ctx, &req.Endpoint); err != nil { return err } return audit(ctx, s.Audit(tx), iam.AuditEntry{ Tenant: req.Endpoint.Tenant, Action: "Webhooks_CreateEndpoint", Resource: req.Endpoint.ID.String(), Detail: req.Endpoint.Name, }) }) write(ctx, w, err, api.Webhooks_CreateEndpointRes{ Endpoint: req.Endpoint, Secret: secret, }) } func (s *WebhooksImpl) UpdateEndpoint(w http.ResponseWriter, r *http.Request) { var req api.Webhooks_UpdateEndpointReq if read(w, r, &req) { return } if check(w, r, s.Guard, webhooks.Cap_Webhooks_UpdateEndpoint, req.Endpoint.Tenant, "") { return } ctx := r.Context() err := dbi.ReadWrite(ctx, s.DB, func(tx dbi.DBI) error { // Preserve the sealed secret: it is never sent by the client and is not // rotated on update. var existing webhooks.Endpoint if err := s.Store(tx).GetEndpoint(ctx, req.Endpoint.Tenant, req.Endpoint.ID, &existing); err != nil { return err } req.Endpoint.DekID = existing.DekID req.Endpoint.SecretEnc = existing.SecretEnc if err := s.Store(tx).UpdateEndpoint(ctx, &req.Endpoint); err != nil { return err } return audit(ctx, s.Audit(tx), iam.AuditEntry{ Tenant: req.Endpoint.Tenant, Action: "Webhooks_UpdateEndpoint", Resource: req.Endpoint.ID.String(), Detail: req.Endpoint.Name, }) }) write(ctx, w, err, api.Webhooks_UpdateEndpointRes{Endpoint: req.Endpoint}) } func (s *WebhooksImpl) GetEndpoint(w http.ResponseWriter, r *http.Request) { var req api.Webhooks_GetEndpointReq if read(w, r, &req) { return } if check(w, r, s.Guard, webhooks.Cap_Webhooks_ReadEndpoint, req.Tenant, "") { return } ctx := r.Context() var res api.Webhooks_GetEndpointRes err := dbi.ReadOnly(ctx, s.DB, func(tx dbi.DBI) error { return s.Store(tx).GetEndpoint(ctx, req.Tenant, req.ID, &res.Endpoint) }) write(ctx, w, err, res) } func (s *WebhooksImpl) ListEndpoints(w http.ResponseWriter, r *http.Request) { var req api.Webhooks_ListEndpointsReq if read(w, r, &req) { return } if check(w, r, s.Guard, webhooks.Cap_Webhooks_ReadEndpoint, req.Tenant, "") { return } ctx := r.Context() var res api.Webhooks_ListEndpointsRes err := dbi.ReadOnly(ctx, s.DB, func(tx dbi.DBI) error { return s.Store(tx).ListEndpoints(ctx, req.Tenant, req.Page, &res.Page) }) write(ctx, w, err, res) } func (s *WebhooksImpl) DeleteEndpoint(w http.ResponseWriter, r *http.Request) { var req api.Webhooks_DeleteEndpointReq if read(w, r, &req) { return } if check(w, r, s.Guard, webhooks.Cap_Webhooks_DeleteEndpoint, req.Tenant, "") { return } ctx := r.Context() err := dbi.ReadWrite(ctx, s.DB, func(tx dbi.DBI) error { if err := s.Store(tx).DeleteEndpoint(ctx, req.Tenant, req.ID); err != nil { return err } return audit(ctx, s.Audit(tx), iam.AuditEntry{ Tenant: req.Tenant, Action: "Webhooks_DeleteEndpoint", Resource: req.ID.String(), }) }) write(ctx, w, err, api.Webhooks_DeleteEndpointRes{}) } // Deliver is the internal task endpoint the delivery worker replays. It loads // the endpoint, unseals the signing secret, and POSTs the signed body through // the guarded sender. It bypasses HTTP status mapping by returning a tasks.Result // envelope: transport failures and 429/5xx retry, other 4xx fail permanently. func (s *WebhooksImpl) Deliver(w http.ResponseWriter, r *http.Request) { var req api.Webhooks_DeliverReq if read(w, r, &req) { return } ctx := r.Context() var ep webhooks.Endpoint var secret []byte err := dbi.ReadOnly(ctx, s.DB, func(tx dbi.DBI) error { if err := s.Store(tx).GetEndpoint(ctx, req.Tenant, req.Endpoint, &ep); err != nil { return err } enc, err := s.Encryptors.For(ctx, s.Deks(tx), req.Tenant) if err != nil { return err } secret, err = enc.Open(ep.SecretEnc) return err }) if errors.Is(err, core.ErrNotFound) { // The endpoint was deleted while a delivery was queued: retrying cannot // fix that, so the task fails permanently. slog.InfoContext(ctx, "webhook endpoint deleted before delivery", "endpoint", req.Endpoint) write(ctx, w, nil, api.Webhooks_DeliverRes{ Task: tasks.Result{Outcome: tasks.OutcomeFailed, Error: "endpoint deleted before delivery"}, }) return } if writeErr(ctx, w, err) { return } status, err := s.Sender.Deliver(ctx, ep, secret, req.Type, req.Body) if err != nil { // Transport failure (connection refused, timeout): worth a retry. write(ctx, w, nil, api.Webhooks_DeliverRes{ Task: tasks.Result{Outcome: tasks.OutcomeRetry, Error: err.Error()}, }) return } res := api.Webhooks_DeliverRes{Task: deliveryOutcome(status)} if res.Task.Outcome != tasks.OutcomeCompleted { slog.WarnContext(ctx, "webhook delivery unsuccessful", "endpoint", req.Endpoint, "status", status, "outcome", res.Task.Outcome) } write(ctx, w, nil, res) } // deliveryOutcome maps a receiver's HTTP status to a task outcome: 2xx/3xx // succeed, 429 and 5xx are transient (retry), other 4xx are permanent (fail). func deliveryOutcome(status int) tasks.Result { switch { case status < 400: return tasks.Result{Outcome: tasks.OutcomeCompleted} case status == http.StatusTooManyRequests || status >= 500: return tasks.Result{Outcome: tasks.OutcomeRetry, Error: http.StatusText(status)} default: return tasks.Result{Outcome: tasks.OutcomeFailed, Error: http.StatusText(status)} } } // newSecret returns a random hex signing secret. The bytes of the returned // string are the HMAC key, so signing and verification agree on it directly. func newSecret() (string, error) { var raw [32]byte if _, err := rand.Read(raw[:]); err != nil { return "", err } return hex.EncodeToString(raw[:]), nil }