diff --git a/spindle/xrpc/webhooks.go b/spindle/xrpc/webhooks.go new file mode 100644 index 000000000..46e942f42 --- /dev/null +++ b/spindle/xrpc/webhooks.go @@ -0,0 +1,377 @@ +package xrpc + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "net/http" + "strconv" + "strings" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/hostutil" + "tangled.org/core/rbac" + "tangled.org/core/spindle/models" + xrpcerr "tangled.org/core/xrpc/errors" +) + +// resolveSettingsRepo resolves a repository DID and checks the actor may manage +// its settings (the same permission gate used for secrets). +func (x *Xrpc) resolveSettingsRepo(actorDid syntax.DID, repoDidStr string) (syntax.DID, xrpcerr.XrpcError, bool) { + repoDid, xerr, ok := x.resolveKnownRepoDid(repoDidStr) + if !ok { + return "", xerr, false + } + allowed, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, repoDid.String()) + if err != nil || !allowed { + return "", xrpcerr.AccessControlError(actorDid.String()), false + } + return repoDid, xrpcerr.XrpcError{}, true +} + +func webhookActor(r *http.Request) (syntax.DID, bool) { + did, ok := r.Context().Value(ActorDid).(syntax.DID) + return did, ok +} + +func (x *Xrpc) WebhookCreate(w http.ResponseWriter, r *http.Request) { + l := x.Logger.With("handler", "WebhookCreate") + + actorDid, ok := webhookActor(r) + if !ok { + writeError(w, xrpcerr.MissingActorDidError, http.StatusForbidden) + return + } + + var input tangled.WebhookCreateWebhook_Input + if err := json.NewDecoder(r.Body).Decode(&input); err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + repoDid, xerr, ok := x.resolveSettingsRepo(actorDid, input.RepoDid) + if !ok { + writeError(w, xerr, http.StatusForbidden) + return + } + + url := strings.TrimSpace(input.Url) + if err := hostutil.ValidateExternalURL(url, x.Config.Server.Dev); err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + if len(input.Events) == 0 { + writeError(w, xrpcerr.XrpcError{Tag: "NoEventsSelected", Message: "at least one event must be specified"}, http.StatusBadRequest) + return + } + + active := true + if input.Active != nil { + active = *input.Active + } + secret := "" + if input.Secret != nil { + secret = strings.TrimSpace(*input.Secret) + } + + webhook := &models.Webhook{ + RepoDid: repoDid, + Url: url, + Secret: secret, + Active: active, + Events: input.Events, + } + if err := x.Db.AddWebhook(webhook); err != nil { + l.Error("failed to add webhook", "err", err) + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + writeJson(w, http.StatusOK, tangled.WebhookCreateWebhook_Output{Id: webhook.Id}) +} + +func (x *Xrpc) WebhookUpdate(w http.ResponseWriter, r *http.Request) { + l := x.Logger.With("handler", "WebhookUpdate") + + actorDid, ok := webhookActor(r) + if !ok { + writeError(w, xrpcerr.MissingActorDidError, http.StatusForbidden) + return + } + + var input tangled.WebhookUpdateWebhook_Input + if err := json.NewDecoder(r.Body).Decode(&input); err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + repoDid, xerr, ok := x.resolveSettingsRepo(actorDid, input.RepoDid) + if !ok { + writeError(w, xerr, http.StatusForbidden) + return + } + + webhook, err := x.Db.GetWebhook(input.Id) + if err != nil || webhook.RepoDid != repoDid { + writeError(w, webhookNotFound(), http.StatusNotFound) + return + } + + if input.Url != nil { + url := strings.TrimSpace(*input.Url) + if url != "" { + if err := hostutil.ValidateExternalURL(url, x.Config.Server.Dev); err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + webhook.Url = url + } + } + if input.Secret != nil { + webhook.Secret = strings.TrimSpace(*input.Secret) + } + if input.Active != nil { + webhook.Active = *input.Active + } + if len(input.Events) > 0 { + webhook.Events = input.Events + } + + if err := x.Db.UpdateWebhook(webhook); err != nil { + l.Error("failed to update webhook", "err", err) + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + w.WriteHeader(http.StatusOK) +} + +func (x *Xrpc) WebhookDelete(w http.ResponseWriter, r *http.Request) { + l := x.Logger.With("handler", "WebhookDelete") + + actorDid, ok := webhookActor(r) + if !ok { + writeError(w, xrpcerr.MissingActorDidError, http.StatusForbidden) + return + } + + var input tangled.WebhookDeleteWebhook_Input + if err := json.NewDecoder(r.Body).Decode(&input); err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + repoDid, xerr, ok := x.resolveSettingsRepo(actorDid, input.RepoDid) + if !ok { + writeError(w, xerr, http.StatusForbidden) + return + } + + webhook, err := x.Db.GetWebhook(input.Id) + if err != nil || webhook.RepoDid != repoDid { + writeError(w, webhookNotFound(), http.StatusNotFound) + return + } + + if err := x.Db.DeleteWebhook(input.Id); err != nil { + l.Error("failed to delete webhook", "err", err) + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + w.WriteHeader(http.StatusOK) +} + +func (x *Xrpc) WebhookToggle(w http.ResponseWriter, r *http.Request) { + l := x.Logger.With("handler", "WebhookToggle") + + actorDid, ok := webhookActor(r) + if !ok { + writeError(w, xrpcerr.MissingActorDidError, http.StatusForbidden) + return + } + + var input tangled.WebhookToggleWebhook_Input + if err := json.NewDecoder(r.Body).Decode(&input); err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + repoDid, xerr, ok := x.resolveSettingsRepo(actorDid, input.RepoDid) + if !ok { + writeError(w, xerr, http.StatusForbidden) + return + } + + webhook, err := x.Db.GetWebhook(input.Id) + if err != nil || webhook.RepoDid != repoDid { + writeError(w, webhookNotFound(), http.StatusNotFound) + return + } + + webhook.Active = !webhook.Active + if err := x.Db.UpdateWebhook(webhook); err != nil { + l.Error("failed to toggle webhook", "err", err) + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + writeJson(w, http.StatusOK, tangled.WebhookToggleWebhook_Output{Active: webhook.Active}) +} + +func (x *Xrpc) WebhookList(w http.ResponseWriter, r *http.Request) { + l := x.Logger.With("handler", "WebhookList") + + actorDid, ok := webhookActor(r) + if !ok { + writeError(w, xrpcerr.MissingActorDidError, http.StatusForbidden) + return + } + + repoDid, xerr, ok := x.resolveSettingsRepo(actorDid, r.URL.Query().Get("repoDid")) + if !ok { + writeError(w, xerr, http.StatusForbidden) + return + } + + webhooks, err := x.Db.GetWebhooksForRepo(repoDid) + if err != nil { + l.Error("failed to get webhooks", "err", err) + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + items := make([]*tangled.WebhookListWebhooks_Webhook, 0, len(webhooks)) + for i := range webhooks { + wh := &webhooks[i] + updated := wh.UpdatedAt.UTC().Format(time.RFC3339) + items = append(items, &tangled.WebhookListWebhooks_Webhook{ + Id: wh.Id, + Url: wh.Url, + Active: wh.Active, + Events: wh.Events, + CreatedAt: wh.CreatedAt.UTC().Format(time.RFC3339), + UpdatedAt: &updated, + }) + } + + writeJson(w, http.StatusOK, tangled.WebhookListWebhooks_Output{Webhooks: items}) +} + +func (x *Xrpc) WebhookListDeliveries(w http.ResponseWriter, r *http.Request) { + l := x.Logger.With("handler", "WebhookListDeliveries") + + actorDid, ok := webhookActor(r) + if !ok { + writeError(w, xrpcerr.MissingActorDidError, http.StatusForbidden) + return + } + + q := r.URL.Query() + repoDid, xerr, ok := x.resolveSettingsRepo(actorDid, q.Get("repoDid")) + if !ok { + writeError(w, xerr, http.StatusForbidden) + return + } + + id, err := strconv.ParseInt(q.Get("id"), 10, 64) + if err != nil { + writeError(w, xrpcerr.XrpcError{Tag: "InvalidRequest", Message: "invalid webhook id"}, http.StatusBadRequest) + return + } + + webhook, err := x.Db.GetWebhook(id) + if err != nil || webhook.RepoDid != repoDid { + writeError(w, webhookNotFound(), http.StatusNotFound) + return + } + + limit := 100 + if s := q.Get("limit"); s != "" { + if n, err := strconv.Atoi(s); err == nil && n > 0 && n <= 100 { + limit = n + } + } + + deliveries, err := x.Db.GetWebhookDeliveries(webhook.Id, limit) + if err != nil { + l.Error("failed to get webhook deliveries", "err", err) + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + items := make([]*tangled.WebhookListDeliveries_Delivery, 0, len(deliveries)) + for i := range deliveries { + d := &deliveries[i] + item := &tangled.WebhookListDeliveries_Delivery{ + Id: d.Id, + DeliveryId: d.DeliveryId, + Event: d.Event, + Url: d.Url, + Success: d.Success, + CreatedAt: d.CreatedAt.UTC().Format(time.RFC3339), + } + if d.RequestBody != "" { + rb := d.RequestBody + item.RequestBody = &rb + } + if d.ResponseBody != "" { + rb := d.ResponseBody + item.ResponseBody = &rb + } + if d.ResponseCode != 0 { + rc := int64(d.ResponseCode) + item.ResponseCode = &rc + } + items = append(items, item) + } + + writeJson(w, http.StatusOK, tangled.WebhookListDeliveries_Output{Deliveries: items}) +} + +func (x *Xrpc) WebhookRetryDelivery(w http.ResponseWriter, r *http.Request) { + actorDid, ok := webhookActor(r) + if !ok { + writeError(w, xrpcerr.MissingActorDidError, http.StatusForbidden) + return + } + + var input tangled.WebhookRetryDelivery_Input + if err := json.NewDecoder(r.Body).Decode(&input); err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + repoDid, xerr, ok := x.resolveSettingsRepo(actorDid, input.RepoDid) + if !ok { + writeError(w, xerr, http.StatusForbidden) + return + } + + webhook, err := x.Db.GetWebhook(input.WebhookId) + if err != nil || webhook.RepoDid != repoDid { + writeError(w, webhookNotFound(), http.StatusNotFound) + return + } + + delivery, err := x.Db.GetWebhookDelivery(input.DeliveryId) + if err != nil || delivery.WebhookId != webhook.Id { + if err != nil && !errors.Is(err, sql.ErrNoRows) { + x.Logger.Error("failed to load delivery", "err", err) + } + writeError(w, xrpcerr.XrpcError{Tag: "DeliveryNotFound", Message: "delivery not found"}, http.StatusNotFound) + return + } + + // re-dispatch async; the new attempt is recorded as its own delivery + go x.Webhooks.Redeliver(context.Background(), *webhook, *delivery) + + w.WriteHeader(http.StatusOK) +} + +func webhookNotFound() xrpcerr.XrpcError { + return xrpcerr.XrpcError{Tag: "WebhookNotFound", Message: "webhook not found"} +} diff --git a/spindle/xrpc/xrpc.go b/spindle/xrpc/xrpc.go index d8232be4b..74783de19 100644 --- a/spindle/xrpc/xrpc.go +++ b/spindle/xrpc/xrpc.go @@ -19,6 +19,7 @@ import ( "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" "tangled.org/core/spindle/secrets" + "tangled.org/core/spindle/webhook" xrpcerr "tangled.org/core/xrpc/errors" "tangled.org/core/xrpc/serviceauth" ) @@ -49,6 +50,7 @@ type Xrpc struct { Resolver *idresolver.Resolver Vault secrets.Manager Notifier *notifier.Notifier + Webhooks *webhook.Service ServiceAuth *serviceauth.ServiceAuth Trigger PipelineTrigger } @@ -64,6 +66,14 @@ func (x *Xrpc) Router() http.Handler { r.Get("/"+tangled.RepoListSecretsNSID, x.ListSecrets) r.Post("/"+tangled.CiCancelPipelineNSID, x.CancelPipeline) r.Post("/"+tangled.CiTriggerPipelineNSID, x.TriggerPipeline) + + r.Post("/"+tangled.WebhookCreateWebhookNSID, x.WebhookCreate) + r.Post("/"+tangled.WebhookUpdateWebhookNSID, x.WebhookUpdate) + r.Post("/"+tangled.WebhookDeleteWebhookNSID, x.WebhookDelete) + r.Post("/"+tangled.WebhookToggleWebhookNSID, x.WebhookToggle) + r.Post("/"+tangled.WebhookRetryDeliveryNSID, x.WebhookRetryDelivery) + r.Get("/"+tangled.WebhookListWebhooksNSID, x.WebhookList) + r.Get("/"+tangled.WebhookListDeliveriesNSID, x.WebhookListDeliveries) }) // service query endpoints (no auth required)