Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468package xrpc
import ( "context" "database/sql" "encoding/json" "errors" "fmt" "net/http" "strconv" "strings" "time"
"github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/org_tangled" "tangled.org/core/hostutil" "tangled.org/core/rbac" "tangled.org/core/spindle/models" "tangled.org/core/spindle/quota" 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, int) { repoDid, xerr, ok := x.resolveKnownRepoDid(repoDidStr) if !ok { status := http.StatusBadRequest if xerr.Tag == xrpcerr.RepoNotFoundError.Tag { status = http.StatusNotFound } return "", xerr, status } allowed, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, repoDid.String()) if err != nil || !allowed { return "", xrpcerr.AccessControlError(actorDid.String()), http.StatusForbidden } return repoDid, xrpcerr.XrpcError{}, http.StatusOK}
func webhookActor(r *http.Request) (syntax.DID, bool) { did, ok := r.Context().Value(ActorDid).(syntax.DID) return did, ok}
func (x *Xrpc) webhookLimit(ctx context.Context, did string, fallback int64) int64 { if x.QuotaStore == nil { return fallback } limit, err := x.QuotaStore.GetLimit(ctx, did, quota.ResourceWebhooks) if err != nil { x.Logger.Warn("failed to read webhook limit", "did", did, "err", err) return fallback } if limit == nil { return fallback } return limit.Limit}
func (x *Xrpc) webhookLimits(ctx context.Context, repoDid syntax.DID) (repoLimit, ownerLimit int64, err error) { repo, err := x.Db.GetRepoByDid(repoDid) if err != nil { return 0, 0, err } return x.webhookLimit(ctx, repoDid.String(), x.Config.Quota.Repo.Webhooks), x.webhookLimit(ctx, repo.Owner.String(), x.Config.Quota.User.Webhooks), nil}
func validateWebhookEvents(events []string) ([]string, xrpcerr.XrpcError, bool) { parsed, err := models.ParseWebhookEvents(events) if err != nil { return nil, xrpcerr.XrpcError{Tag: "InvalidEvent", Message: err.Error()}, false } if len(parsed) == 0 { return nil, xrpcerr.XrpcError{Tag: "NoEventsSelected", Message: "at least one event must be specified"}, false } names := make([]string, 0, len(parsed)) for _, event := range parsed { names = append(names, string(event)) } return names, xrpcerr.XrpcError{}, true}
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 org_tangled.TempWebhookCreateWebhook_Input if err := json.NewDecoder(r.Body).Decode(&input); err != nil { writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) return }
repoDid, xerr, status := x.resolveSettingsRepo(actorDid, input.RepoDid) if status != http.StatusOK { writeError(w, xerr, status) 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 } events, xerr, ok := validateWebhookEvents(input.Events) if !ok { writeError(w, xerr, http.StatusBadRequest) return }
repoLimit, ownerLimit, err := x.webhookLimits(r.Context(), repoDid) if err != nil { l.Error("failed to resolve webhook limits", "err", err) writeError(w, xrpcerr.RepoNotFoundError, http.StatusNotFound) 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: events, } created, err := x.Db.AddWebhook(r.Context(), webhook, repoLimit, ownerLimit) if err != nil { l.Error("failed to add webhook", "err", err) writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) return } if !created { writeError(w, xrpcerr.WebhookLimitExceededError, http.StatusTooManyRequests) return }
writeJson(w, http.StatusOK, org_tangled.TempWebhookCreateWebhook_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 org_tangled.TempWebhookUpdateWebhook_Input if err := json.NewDecoder(r.Body).Decode(&input); err != nil { writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) return }
repoDid, xerr, status := x.resolveSettingsRepo(actorDid, input.RepoDid) if status != http.StatusOK { writeError(w, xerr, status) 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 { events, xerr, ok := validateWebhookEvents(input.Events) if !ok { writeError(w, xerr, http.StatusBadRequest) return } webhook.Events = 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 org_tangled.TempWebhookDeleteWebhook_Input if err := json.NewDecoder(r.Body).Decode(&input); err != nil { writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) return }
repoDid, xerr, status := x.resolveSettingsRepo(actorDid, input.RepoDid) if status != http.StatusOK { writeError(w, xerr, status) 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 org_tangled.TempWebhookToggleWebhook_Input if err := json.NewDecoder(r.Body).Decode(&input); err != nil { writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) return }
repoDid, xerr, status := x.resolveSettingsRepo(actorDid, input.RepoDid) if status != http.StatusOK { writeError(w, xerr, status) 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, org_tangled.TempWebhookToggleWebhook_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 }
q := r.URL.Query() repoDid, xerr, status := x.resolveSettingsRepo(actorDid, q.Get("repoDid")) if status != http.StatusOK { writeError(w, xerr, status) 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([]*org_tangled.TempWebhookListWebhooks_Webhook, 0, len(webhooks)) for i := range webhooks { items = append(items, webhookToXrpc(&webhooks[i])) }
writeJson(w, http.StatusOK, org_tangled.TempWebhookListWebhooks_Output{Webhooks: items})}
func (x *Xrpc) WebhookGetDeliveriesForWebhooks(w http.ResponseWriter, r *http.Request) { l := x.Logger.With("handler", "WebhookGetDeliveriesForWebhooks")
actorDid, ok := webhookActor(r) if !ok { writeError(w, xrpcerr.MissingActorDidError, http.StatusForbidden) return }
q := r.URL.Query() repoDid, xerr, status := x.resolveSettingsRepo(actorDid, q.Get("repoDid")) if status != http.StatusOK { writeError(w, xerr, status) return }
raw := q["ids"] if len(raw) == 0 { writeError(w, xrpcerr.XrpcError{Tag: "InvalidRequest", Message: "at least one webhook id is required"}, http.StatusBadRequest) return } if len(raw) > maxWebhookIds { writeError(w, xrpcerr.GenericError(fmt.Errorf("at most %d webhooks per request", maxWebhookIds)), http.StatusBadRequest) return }
ids := make([]int64, 0, len(raw)) for _, value := range raw { id, err := strconv.ParseInt(value, 10, 64) if err != nil { writeError(w, xrpcerr.XrpcError{Tag: "InvalidRequest", Message: "invalid webhook id"}, http.StatusBadRequest) return } ids = append(ids, id) }
deliveries, err := x.Db.GetWebhookDeliveriesForWebhooks(repoDid, ids, webhookDeliveryLimit(q.Get("limit"), 5)) if err != nil { l.Error("failed to get webhook deliveries", "err", err) writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) return }
writeJson(w, http.StatusOK, org_tangled.TempWebhookGetDeliveriesForWebhooks_Output{ Deliveries: deliveriesToXrpc(deliveries), })}
const maxWebhookIds = 50
func webhookDeliveryLimit(raw string, fallback int) int { if raw == "" { return fallback } n, err := strconv.Atoi(raw) if err != nil || n < 1 || n > 100 { return fallback } return n}
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 org_tangled.TempWebhookRetryDelivery_Input if err := json.NewDecoder(r.Body).Decode(&input); err != nil { writeError(w, xrpcerr.GenericError(err), http.StatusBadRequest) return }
repoDid, xerr, status := x.resolveSettingsRepo(actorDid, input.RepoDid) if status != http.StatusOK { writeError(w, xerr, status) 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 webhookToXrpc(wh *models.Webhook) *org_tangled.TempWebhookListWebhooks_Webhook { updated := wh.UpdatedAt.UTC().Format(time.RFC3339) return &org_tangled.TempWebhookListWebhooks_Webhook{ Id: wh.Id, Url: wh.Url, Active: wh.Active, Events: wh.Events, CreatedAt: wh.CreatedAt.UTC().Format(time.RFC3339), UpdatedAt: &updated, }}
func deliveryToXrpc(d *models.WebhookDelivery) *org_tangled.TempWebhookGetDeliveriesForWebhooks_Delivery { item := &org_tangled.TempWebhookGetDeliveriesForWebhooks_Delivery{ Id: d.Id, WebhookId: d.WebhookId, 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 } return item}
func deliveriesToXrpc(deliveries []models.WebhookDelivery) []*org_tangled.TempWebhookGetDeliveriesForWebhooks_Delivery { items := make([]*org_tangled.TempWebhookGetDeliveriesForWebhooks_Delivery, 0, len(deliveries)) for i := range deliveries { items = append(items, deliveryToXrpc(&deliveries[i])) } return items}
func webhookNotFound() xrpcerr.XrpcError { return xrpcerr.XrpcError{Tag: "WebhookNotFound", Message: "webhook not found"}}