From c710987f22bb25c4bdf713ffb2006bbc4767bf9a Mon Sep 17 00:00:00 2001 From: Anirudh Oppiliappan Date: Wed, 11 Feb 2026 20:02:31 +0000 Subject: [PATCH] appview/{db,models}: webhook tables and crud ops Signed-off-by: Anirudh Oppiliappan --- appview/db/db.go | 30 ++++++++++++++++++++++++++++++ appview/db/webhooks.go | 298 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ appview/models/webhook.go | 74 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 3 file(s) changed, 402 insertion(s)(+), 0 deletion(s)(-) diff --git a/appview/db/db.go b/appview/db/db.go --- a/appview/db/db.go +++ b/appview/db/db.go @@ -568,6 +568,34 @@ to_at text not null, unique (from_at, to_at) ); + create table if not exists webhooks ( + id integer primary key autoincrement, + repo_at text not null, + url text not null, + secret text, + active integer not null default 1, + events text not null, -- comma-separated list of events + created_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + updated_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + + foreign key (repo_at) references repos(at_uri) on delete cascade + ); + + create table if not exists webhook_deliveries ( + id integer primary key autoincrement, + webhook_id integer not null, + event text not null, + delivery_id text not null, + url text not null, + request_body text not null, + response_code integer, + response_body text, + success integer not null default 0, + created_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + + foreign key (webhook_id) references webhooks(id) on delete cascade + ); + create table if not exists migrations ( id integer primary key autoincrement, name text unique @@ -578,6 +606,8 @@ create index if not exists idx_notifications_recipient_created on notifications(recipient_did, created desc); create index if not exists idx_notifications_recipient_read on notifications(recipient_did, read); create index if not exists idx_references_from_at on reference_links(from_at); create index if not exists idx_references_to_at on reference_links(to_at); + create index if not exists idx_webhooks_repo_at on webhooks(repo_at); + create index if not exists idx_webhook_deliveries_webhook_id on webhook_deliveries(webhook_id); `) if err != nil { return nil, err diff --git a/appview/db/webhooks.go b/appview/db/webhooks.go new file mode 100644 --- /dev/null +++ b/appview/db/webhooks.go @@ -0,0 +1,298 @@ +package db + +import ( + "database/sql" + "fmt" + "strings" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/appview/models" + "tangled.org/core/orm" +) + +// GetWebhooks returns all webhooks for a repository +func GetWebhooks(e Execer, filters ...orm.Filter) ([]models.Webhook, error) { + var conditions []string + var args []any + for _, filter := range filters { + conditions = append(conditions, filter.Condition()) + args = append(args, filter.Arg()...) + } + + whereClause := "" + if conditions != nil { + whereClause = " where " + strings.Join(conditions, " and ") + } + + query := fmt.Sprintf(` + select + id, + repo_at, + url, + secret, + active, + events, + created_at, + updated_at + from webhooks + %s + order by created_at desc + `, whereClause) + + rows, err := e.Query(query, args...) + if err != nil { + return nil, fmt.Errorf("failed to query webhooks: %w", err) + } + defer rows.Close() + + var webhooks []models.Webhook + for rows.Next() { + var wh models.Webhook + var createdAt, updatedAt, eventsStr string + var secret sql.NullString + var active int + + err := rows.Scan( + &wh.Id, + &wh.RepoAt, + &wh.Url, + &secret, + &active, + &eventsStr, + &createdAt, + &updatedAt, + ) + if err != nil { + return nil, fmt.Errorf("failed to scan webhook: %w", err) + } + + if secret.Valid { + wh.Secret = secret.String + } + wh.Active = active == 1 + if eventsStr != "" { + wh.Events = strings.Split(eventsStr, ",") + } + + if t, err := time.Parse(time.RFC3339, createdAt); err == nil { + wh.CreatedAt = t + } + if t, err := time.Parse(time.RFC3339, updatedAt); err == nil { + wh.UpdatedAt = t + } + + webhooks = append(webhooks, wh) + } + + if err = rows.Err(); err != nil { + return nil, fmt.Errorf("failed to iterate webhooks: %w", err) + } + + return webhooks, nil +} + +// GetWebhook returns a single webhook by ID +func GetWebhook(e Execer, id int64) (*models.Webhook, error) { + webhooks, err := GetWebhooks(e, orm.FilterEq("id", id)) + if err != nil { + return nil, err + } + + if len(webhooks) == 0 { + return nil, sql.ErrNoRows + } + + if len(webhooks) != 1 { + return nil, fmt.Errorf("expected 1 webhook, got %d", len(webhooks)) + } + + return &webhooks[0], nil +} + +// AddWebhook creates a new webhook +func AddWebhook(e Execer, webhook *models.Webhook) error { + eventsStr := strings.Join(webhook.Events, ",") + active := 0 + if webhook.Active { + active = 1 + } + + result, err := e.Exec(` + insert into webhooks (repo_at, url, secret, active, events) + values (?, ?, ?, ?, ?) + `, webhook.RepoAt.String(), webhook.Url, webhook.Secret, active, eventsStr) + + if err != nil { + return fmt.Errorf("failed to insert webhook: %w", err) + } + + id, err := result.LastInsertId() + if err != nil { + return fmt.Errorf("failed to get webhook id: %w", err) + } + + webhook.Id = id + return nil +} + +// UpdateWebhook updates an existing webhook +func UpdateWebhook(e Execer, webhook *models.Webhook) error { + eventsStr := strings.Join(webhook.Events, ",") + active := 0 + if webhook.Active { + active = 1 + } + + _, err := e.Exec(` + update webhooks + set url = ?, secret = ?, active = ?, events = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') + where id = ? + `, webhook.Url, webhook.Secret, active, eventsStr, webhook.Id) + + if err != nil { + return fmt.Errorf("failed to update webhook: %w", err) + } + + return nil +} + +// DeleteWebhook deletes a webhook +func DeleteWebhook(e Execer, id int64) error { + _, err := e.Exec(`delete from webhooks where id = ?`, id) + if err != nil { + return fmt.Errorf("failed to delete webhook: %w", err) + } + return nil +} + +// AddWebhookDelivery records a webhook delivery attempt +func AddWebhookDelivery(e Execer, delivery *models.WebhookDelivery) error { + success := 0 + if delivery.Success { + success = 1 + } + + result, err := e.Exec(` + insert into webhook_deliveries ( + webhook_id, + event, + delivery_id, + url, + request_body, + response_code, + response_body, + success + ) values (?, ?, ?, ?, ?, ?, ?, ?) + `, + delivery.WebhookId, + delivery.Event, + delivery.DeliveryId, + delivery.Url, + delivery.RequestBody, + delivery.ResponseCode, + delivery.ResponseBody, + success, + ) + + if err != nil { + return fmt.Errorf("failed to insert webhook delivery: %w", err) + } + + id, err := result.LastInsertId() + if err != nil { + return fmt.Errorf("failed to get delivery id: %w", err) + } + + delivery.Id = id + return nil +} + +// GetWebhookDeliveries returns recent deliveries for a webhook +func GetWebhookDeliveries(e Execer, webhookId int64, limit int) ([]models.WebhookDelivery, error) { + if limit <= 0 { + limit = 20 + } + + query := ` + select + id, + webhook_id, + event, + delivery_id, + url, + request_body, + response_code, + response_body, + success, + created_at + from webhook_deliveries + where webhook_id = ? + order by created_at desc + limit ? + ` + + rows, err := e.Query(query, webhookId, limit) + if err != nil { + return nil, fmt.Errorf("failed to query webhook deliveries: %w", err) + } + defer rows.Close() + + var deliveries []models.WebhookDelivery + for rows.Next() { + var d models.WebhookDelivery + var createdAt string + var success int + var responseCode sql.NullInt64 + var responseBody sql.NullString + + err := rows.Scan( + &d.Id, + &d.WebhookId, + &d.Event, + &d.DeliveryId, + &d.Url, + &d.RequestBody, + &responseCode, + &responseBody, + &success, + &createdAt, + ) + if err != nil { + return nil, fmt.Errorf("failed to scan delivery: %w", err) + } + + d.Success = success == 1 + if responseCode.Valid { + d.ResponseCode = int(responseCode.Int64) + } + if responseBody.Valid { + d.ResponseBody = responseBody.String + } + + if t, err := time.Parse(time.RFC3339, createdAt); err == nil { + d.CreatedAt = t + } + + deliveries = append(deliveries, d) + } + + if err = rows.Err(); err != nil { + return nil, fmt.Errorf("failed to iterate deliveries: %w", err) + } + + return deliveries, nil +} + +// GetWebhooksForRepo is a convenience function to get all webhooks for a repository +func GetWebhooksForRepo(e Execer, repoAt syntax.ATURI) ([]models.Webhook, error) { + return GetWebhooks(e, orm.FilterEq("repo_at", repoAt.String())) +} + +// GetActiveWebhooksForRepo returns only active webhooks for a repository +func GetActiveWebhooksForRepo(e Execer, repoAt syntax.ATURI) ([]models.Webhook, error) { + return GetWebhooks(e, + orm.FilterEq("repo_at", repoAt.String()), + orm.FilterEq("active", 1), + ) +} diff --git a/appview/models/webhook.go b/appview/models/webhook.go new file mode 100644 --- /dev/null +++ b/appview/models/webhook.go @@ -0,0 +1,74 @@ +package models + +import ( + "slices" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +type WebhookEvent string + +const ( + WebhookEventPush WebhookEvent = "push" +) + +type Webhook struct { + Id int64 + RepoAt syntax.ATURI + Url string + Secret string + Active bool + Events []string // comma-separated event types + CreatedAt time.Time + UpdatedAt time.Time +} + +// HasEvent checks if the webhook is subscribed to a specific event +func (w *Webhook) HasEvent(event WebhookEvent) bool { + return slices.Contains(w.Events, string(event)) +} + +type WebhookDelivery struct { + Id int64 + WebhookId int64 + Event string + DeliveryId string // UUID for tracking + Url string + RequestBody string + ResponseCode int + ResponseBody string + Success bool + CreatedAt time.Time +} + +// WebhookPayload represents the webhook payload structure +type WebhookPayload struct { + Ref string `json:"ref"` + Before string `json:"before"` + After string `json:"after"` + Repository WebhookRepository `json:"repository"` + Pusher WebhookUser `json:"pusher"` +} + +// WebhookRepository represents repository information in webhook payload +type WebhookRepository struct { + Name string `json:"name"` + FullName string `json:"full_name"` + Description string `json:"description"` + Fork bool `json:"fork"` + HtmlUrl string `json:"html_url"` + CloneUrl string `json:"clone_url"` + SshUrl string `json:"ssh_url"` + Website string `json:"website,omitempty"` + StarsCount int `json:"stars_count,omitempty"` + OpenIssues int `json:"open_issues_count,omitempty"` + CreatedAt string `json:"created_at"` + UpdatedAt string `json:"updated_at"` + Owner WebhookUser `json:"owner"` +} + +// WebhookUser represents user information in webhook payload +type WebhookUser struct { + Did string `json:"did"` +} -- tangled.sh