package api_impl import ( "context" "database/sql" "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/slack" "atlas9.dev/c/demo/lib/tasks" ) type SlackImpl struct { DB *sql.DB Guard access.Guard Store dbi.Factory[slack.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 *slack.Sender } func (s *SlackImpl) ServeMux(mux *http.ServeMux) { mux.HandleFunc(api.Path_Slack_CreateEndpoint, s.CreateEndpoint) mux.HandleFunc(api.Path_Slack_UpdateEndpoint, s.UpdateEndpoint) mux.HandleFunc(api.Path_Slack_GetEndpoint, s.GetEndpoint) mux.HandleFunc(api.Path_Slack_ListEndpoints, s.ListEndpoints) mux.HandleFunc(api.Path_Slack_DeleteEndpoint, s.DeleteEndpoint) mux.HandleFunc(api.Path_Slack_Deliver, s.Deliver) } func (s *SlackImpl) CreateEndpoint(w http.ResponseWriter, r *http.Request) { var req api.Slack_CreateEndpointReq if read(w, r, &req) { return } if check(w, r, s.Guard, slack.Cap_Slack_CreateEndpoint, req.Endpoint.Tenant, "") { return } ctx := r.Context() err := dbi.ReadWrite(ctx, s.DB, func(tx dbi.DBI) error { // Sealing needs the tenant DEK; grant Dek_Use scoped to this tenant for // that sub-operation (the caller holds the Slack cap, not Dek_Use). if err := s.sealURL(ctx, tx, &req.Endpoint, req.URL); err != nil { return err } 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: "Slack_CreateEndpoint", Resource: req.Endpoint.ID.String(), Detail: req.Endpoint.Name, }) }) write(ctx, w, err, api.Slack_CreateEndpointRes{Endpoint: req.Endpoint}) } func (s *SlackImpl) UpdateEndpoint(w http.ResponseWriter, r *http.Request) { var req api.Slack_UpdateEndpointReq if read(w, r, &req) { return } if check(w, r, s.Guard, slack.Cap_Slack_UpdateEndpoint, req.Endpoint.Tenant, "") { return } ctx := r.Context() err := dbi.ReadWrite(ctx, s.DB, func(tx dbi.DBI) error { if req.URL != "" { // Rotate the webhook URL: re-seal the new value. if err := s.sealURL(ctx, tx, &req.Endpoint, req.URL); err != nil { return err } } else { // Preserve the sealed URL: it is never sent by the client and is not // rotated unless a new one is supplied. var existing slack.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.URLEnc = existing.URLEnc } 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: "Slack_UpdateEndpoint", Resource: req.Endpoint.ID.String(), Detail: req.Endpoint.Name, }) }) write(ctx, w, err, api.Slack_UpdateEndpointRes{Endpoint: req.Endpoint}) } // sealURL seals url under the tenant DEK and stores the ciphertext on e. It // grants Dek_Use scoped to the tenant for the seal, since the caller holds the // Slack cap and not Dek_Use. func (s *SlackImpl) sealURL(ctx context.Context, tx dbi.DBI, e *slack.Endpoint, url string) error { sealCtx := access.PutScope(ctx, e.Tenant, envelope.Cap_Dek_Use) enc, err := s.Encryptors.For(sealCtx, s.Deks(tx), e.Tenant) if err != nil { return err } dekID, sealed, err := enc.Seal([]byte(url)) if err != nil { return err } e.DekID = dekID e.URLEnc = sealed return nil } func (s *SlackImpl) GetEndpoint(w http.ResponseWriter, r *http.Request) { var req api.Slack_GetEndpointReq if read(w, r, &req) { return } if check(w, r, s.Guard, slack.Cap_Slack_ReadEndpoint, req.Tenant, "") { return } ctx := r.Context() var res api.Slack_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 *SlackImpl) ListEndpoints(w http.ResponseWriter, r *http.Request) { var req api.Slack_ListEndpointsReq if read(w, r, &req) { return } if check(w, r, s.Guard, slack.Cap_Slack_ReadEndpoint, req.Tenant, "") { return } ctx := r.Context() var res api.Slack_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 *SlackImpl) DeleteEndpoint(w http.ResponseWriter, r *http.Request) { var req api.Slack_DeleteEndpointReq if read(w, r, &req) { return } if check(w, r, s.Guard, slack.Cap_Slack_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: "Slack_DeleteEndpoint", Resource: req.ID.String(), }) }) write(ctx, w, err, api.Slack_DeleteEndpointRes{}) } // Deliver is the internal task endpoint the delivery worker replays. It loads // the endpoint, unseals the webhook URL, and POSTs the message 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 *SlackImpl) Deliver(w http.ResponseWriter, r *http.Request) { var req api.Slack_DeliverReq if read(w, r, &req) { return } ctx := r.Context() var url string err := dbi.ReadOnly(ctx, s.DB, func(tx dbi.DBI) error { var ep slack.Endpoint 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 } plain, err := enc.Open(ep.URLEnc) if err != nil { return err } url = string(plain) return nil }) 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, "slack endpoint deleted before delivery", "endpoint", req.Endpoint) write(ctx, w, nil, api.Slack_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, url, req.Body) if err != nil { // Transport failure (connection refused, timeout): worth a retry. write(ctx, w, nil, api.Slack_DeliverRes{ Task: tasks.Result{Outcome: tasks.OutcomeRetry, Error: err.Error()}, }) return } res := api.Slack_DeliverRes{Task: deliveryOutcome(status)} if res.Task.Outcome != tasks.OutcomeCompleted { slog.WarnContext(ctx, "slack delivery unsuccessful", "endpoint", req.Endpoint, "status", status, "outcome", res.Task.Outcome) } write(ctx, w, nil, res) }