diff --git a/spindle/db/db.go b/spindle/db/db.go index eb7484cc..42fd3e14 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -132,6 +132,32 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { foreign key (pipeline_id) references pipelines(id) on delete cascade ); + create table if not exists webhooks ( + id integer primary key autoincrement, + repo_did text not null, + url text not null, + secret text, + active integer not null default 1, + events text not null, -- comma-separated event types + 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')) + ); + create index if not exists idx_webhooks_repo_did on webhooks(repo_did); + + create table if not exists webhook_deliveries ( + id integer primary key autoincrement, + webhook_id integer not null references webhooks(id) on delete cascade, + event text not null, + delivery_id text not null, + url text not null, + request_body text, + 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')) + ); + create index if not exists idx_webhook_deliveries_webhook_id on webhook_deliveries(webhook_id); + create table if not exists migrations ( id integer primary key autoincrement, name text unique @@ -270,6 +296,26 @@ func runMigrations(_ context.Context, conn *sql.Conn, logger *slog.Logger) error return err } + // re-introduce a repo display name, needed to build repository:renamed + // webhook payloads. the earlier repos-to-repo-did migration dropped the + // legacy name column; this adds it back as nullable metadata. + if err := orm.RunMigration(conn, logger, "repos-add-name-column", func(tx *sql.Tx) error { + var hasName int + if err := tx.QueryRow( + `select count(*) from pragma_table_info('repos') where name = 'name'`, + ).Scan(&hasName); err != nil { + return err + } + if hasName == 0 { + if _, err := tx.Exec(`alter table repos add column name text`); err != nil { + return err + } + } + return nil + }); err != nil { + return err + } + return nil } diff --git a/spindle/db/repos.go b/spindle/db/repos.go index fa8bc0de..e65eda04 100644 --- a/spindle/db/repos.go +++ b/spindle/db/repos.go @@ -11,6 +11,7 @@ type Repo struct { Owner syntax.DID Rkey syntax.RecordKey RepoDid syntax.DID + Name string CreatedAt string } @@ -19,14 +20,19 @@ func (d *DB) AddRepo(repo Repo) error { if repo.CreatedAt != "" { createdAt = sql.NullString{String: repo.CreatedAt, Valid: true} } + var name sql.NullString + if repo.Name != "" { + name = sql.NullString{String: repo.Name, Valid: true} + } _, err := d.Exec( - `insert into repos (knot, owner, rkey, repo_did, created_at) - values (?, ?, ?, ?, ?) + `insert into repos (knot, owner, rkey, repo_did, name, created_at) + values (?, ?, ?, ?, ?, ?) on conflict(owner, rkey) do update set knot = excluded.knot, repo_did = excluded.repo_did, + name = coalesce(excluded.name, repos.name), created_at = coalesce(excluded.created_at, repos.created_at)`, - repo.Knot, repo.Owner.String(), repo.Rkey.String(), repo.RepoDid.String(), createdAt, + repo.Knot, repo.Owner.String(), repo.Rkey.String(), repo.RepoDid.String(), name, createdAt, ) return err } @@ -81,8 +87,8 @@ func (d *DB) Knots() ([]string, error) { } func scanRepo(row interface{ Scan(...any) error }) (*Repo, error) { - var knot, owner, rkey, repoDid string - if err := row.Scan(&knot, &owner, &rkey, &repoDid); err != nil { + var knot, owner, rkey, repoDid, name string + if err := row.Scan(&knot, &owner, &rkey, &repoDid, &name); err != nil { return nil, err } return &Repo{ @@ -90,6 +96,7 @@ func scanRepo(row interface{ Scan(...any) error }) (*Repo, error) { Owner: syntax.DID(owner), Rkey: syntax.RecordKey(rkey), RepoDid: syntax.DID(repoDid), + Name: name, }, nil } @@ -122,20 +129,20 @@ func (d *DB) SiblingRkeysForRepoDid(owner, repoDid syntax.DID, excludeRkey synta func (d *DB) GetRepoByDid(repoDid syntax.DID) (*Repo, error) { return scanRepo(d.QueryRow( - `select knot, owner, rkey, coalesce(repo_did, '') from repos where repo_did = ?`, + `select knot, owner, rkey, coalesce(repo_did, ''), coalesce(name, '') from repos where repo_did = ?`, repoDid.String(), )) } func (d *DB) GetRepoByOwnerRkey(owner syntax.DID, rkey syntax.RecordKey) (*Repo, error) { return scanRepo(d.QueryRow( - `select knot, owner, rkey, coalesce(repo_did, '') from repos where owner = ? and rkey = ?`, + `select knot, owner, rkey, coalesce(repo_did, ''), coalesce(name, '') from repos where owner = ? and rkey = ?`, owner.String(), rkey.String(), )) } func (d *DB) AllRepos() ([]Repo, error) { - rows, err := d.Query(`select knot, owner, rkey, coalesce(repo_did, '') from repos`) + rows, err := d.Query(`select knot, owner, rkey, coalesce(repo_did, ''), coalesce(name, '') from repos`) if err != nil { return nil, err } diff --git a/spindle/db/webhooks.go b/spindle/db/webhooks.go new file mode 100644 index 00000000..1ac0e670 --- /dev/null +++ b/spindle/db/webhooks.go @@ -0,0 +1,252 @@ +package db + +import ( + "database/sql" + "fmt" + "strings" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/spindle/models" +) + +const webhookColumns = `id, repo_did, url, secret, active, events, created_at, updated_at` + +func scanWebhook(row interface{ Scan(...any) error }) (*models.Webhook, error) { + var wh models.Webhook + var repoDid, createdAt, updatedAt, eventsStr string + var secret sql.NullString + var active int + + if err := row.Scan( + &wh.Id, + &repoDid, + &wh.Url, + &secret, + &active, + &eventsStr, + &createdAt, + &updatedAt, + ); err != nil { + return nil, err + } + + wh.RepoDid = syntax.DID(repoDid) + 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 + } + return &wh, nil +} + +// GetWebhooksForRepo returns all webhooks configured for a repository, newest first. +func (d *DB) GetWebhooksForRepo(repoDid syntax.DID) ([]models.Webhook, error) { + rows, err := d.Query( + `select `+webhookColumns+` from webhooks where repo_did = ? order by created_at desc`, + repoDid.String(), + ) + if err != nil { + return nil, fmt.Errorf("failed to query webhooks: %w", err) + } + defer rows.Close() + + var webhooks []models.Webhook + for rows.Next() { + wh, err := scanWebhook(rows) + if err != nil { + return nil, fmt.Errorf("failed to scan webhook: %w", err) + } + webhooks = append(webhooks, *wh) + } + return webhooks, rows.Err() +} + +// GetActiveWebhooksForRepo returns only active webhooks for a repository. +func (d *DB) GetActiveWebhooksForRepo(repoDid syntax.DID) ([]models.Webhook, error) { + rows, err := d.Query( + `select `+webhookColumns+` from webhooks where repo_did = ? and active = 1 order by created_at desc`, + repoDid.String(), + ) + if err != nil { + return nil, fmt.Errorf("failed to query active webhooks: %w", err) + } + defer rows.Close() + + var webhooks []models.Webhook + for rows.Next() { + wh, err := scanWebhook(rows) + if err != nil { + return nil, fmt.Errorf("failed to scan webhook: %w", err) + } + webhooks = append(webhooks, *wh) + } + return webhooks, rows.Err() +} + +// GetWebhook returns a single webhook by ID. +func (d *DB) GetWebhook(id int64) (*models.Webhook, error) { + return scanWebhook(d.QueryRow( + `select `+webhookColumns+` from webhooks where id = ?`, id, + )) +} + +// AddWebhook creates a new webhook, setting webhook.Id on success. +func (d *DB) AddWebhook(webhook *models.Webhook) error { + eventsStr := strings.Join(webhook.Events, ",") + active := 0 + if webhook.Active { + active = 1 + } + + result, err := d.Exec( + `insert into webhooks (repo_did, url, secret, active, events) + values (?, ?, ?, ?, ?)`, + webhook.RepoDid.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's mutable fields. +func (d *DB) UpdateWebhook(webhook *models.Webhook) error { + eventsStr := strings.Join(webhook.Events, ",") + active := 0 + if webhook.Active { + active = 1 + } + + _, err := d.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 and (via cascade) its deliveries. +func (d *DB) DeleteWebhook(id int64) error { + if _, err := d.Exec(`delete from webhooks where id = ?`, id); err != nil { + return fmt.Errorf("failed to delete webhook: %w", err) + } + return nil +} + +// AddWebhookDelivery records a webhook delivery attempt, setting delivery.Id. +func (d *DB) AddWebhookDelivery(delivery *models.WebhookDelivery) error { + success := 0 + if delivery.Success { + success = 1 + } + + result, err := d.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 +} + +const deliveryColumns = `id, webhook_id, event, delivery_id, url, request_body, response_code, response_body, success, created_at` + +func scanDelivery(row interface{ Scan(...any) error }) (*models.WebhookDelivery, error) { + var dv models.WebhookDelivery + var createdAt string + var success int + var responseCode sql.NullInt64 + var responseBody sql.NullString + + if err := row.Scan( + &dv.Id, + &dv.WebhookId, + &dv.Event, + &dv.DeliveryId, + &dv.Url, + &dv.RequestBody, + &responseCode, + &responseBody, + &success, + &createdAt, + ); err != nil { + return nil, err + } + + dv.Success = success == 1 + if responseCode.Valid { + dv.ResponseCode = int(responseCode.Int64) + } + if responseBody.Valid { + dv.ResponseBody = responseBody.String + } + if t, err := time.Parse(time.RFC3339, createdAt); err == nil { + dv.CreatedAt = t + } + return &dv, nil +} + +// GetWebhookDeliveries returns recent deliveries for a webhook, newest first. +func (d *DB) GetWebhookDeliveries(webhookId int64, limit int) ([]models.WebhookDelivery, error) { + if limit <= 0 { + limit = 20 + } + rows, err := d.Query( + `select `+deliveryColumns+` from webhook_deliveries where webhook_id = ? order by created_at desc limit ?`, + 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() { + dv, err := scanDelivery(rows) + if err != nil { + return nil, fmt.Errorf("failed to scan delivery: %w", err) + } + deliveries = append(deliveries, *dv) + } + return deliveries, rows.Err() +} + +// GetWebhookDelivery returns a single delivery by its delivery_id (a uuid). +func (d *DB) GetWebhookDelivery(deliveryId string) (*models.WebhookDelivery, error) { + return scanDelivery(d.QueryRow( + `select `+deliveryColumns+` from webhook_deliveries where delivery_id = ?`, deliveryId, + )) +} diff --git a/spindle/db/webhooks_smoke_test.go b/spindle/db/webhooks_smoke_test.go new file mode 100644 index 00000000..1807deb9 --- /dev/null +++ b/spindle/db/webhooks_smoke_test.go @@ -0,0 +1,59 @@ +package db + +import ( + "context" + "path/filepath" + "testing" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/spindle/models" +) + +func TestWebhookSmoke(t *testing.T) { + d, err := Make(context.Background(), filepath.Join(t.TempDir(), "test.db")) + if err != nil { + t.Fatalf("make db: %v", err) + } + + repoDid := syntax.DID("did:plc:repo123") + if err := d.AddRepo(Repo{ + Knot: "knot.example", Owner: "did:plc:owner", Rkey: "r1", + RepoDid: repoDid, Name: "myrepo", CreatedAt: "2024-01-01T00:00:00Z", + }); err != nil { + t.Fatalf("add repo: %v", err) + } + got, err := d.GetRepoByDid(repoDid) + if err != nil || got.Name != "myrepo" { + t.Fatalf("get repo: %+v err=%v", got, err) + } + + wh := &models.Webhook{RepoDid: repoDid, Url: "https://example.com/hook", Active: true, Events: []string{"push", "pull_request:created"}} + if err := d.AddWebhook(wh); err != nil || wh.Id == 0 { + t.Fatalf("add webhook: id=%d err=%v", wh.Id, err) + } + active, err := d.GetActiveWebhooksForRepo(repoDid) + if err != nil || len(active) != 1 || !active[0].HasEvent(models.WebhookEventPush) { + t.Fatalf("active webhooks: %+v err=%v", active, err) + } + + dv := &models.WebhookDelivery{WebhookId: wh.Id, Event: "push", DeliveryId: "uuid-1", Url: wh.Url, RequestBody: "{}", Success: true, ResponseCode: 200} + if err := d.AddWebhookDelivery(dv); err != nil { + t.Fatalf("add delivery: %v", err) + } + list, err := d.GetWebhookDeliveries(wh.Id, 10) + if err != nil || len(list) != 1 { + t.Fatalf("deliveries: %d err=%v", len(list), err) + } + byId, err := d.GetWebhookDelivery("uuid-1") + if err != nil || !byId.Success || byId.ResponseCode != 200 { + t.Fatalf("delivery by id: %+v err=%v", byId, err) + } + + // cascade: deleting the webhook removes its deliveries + if err := d.DeleteWebhook(wh.Id); err != nil { + t.Fatalf("delete webhook: %v", err) + } + if left, _ := d.GetWebhookDeliveries(wh.Id, 10); len(left) != 0 { + t.Fatalf("expected deliveries cascade-deleted, got %d", len(left)) + } +} diff --git a/spindle/models/webhook.go b/spindle/models/webhook.go new file mode 100644 index 00000000..58d649cc --- /dev/null +++ b/spindle/models/webhook.go @@ -0,0 +1,128 @@ +package models + +import ( + "slices" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +type WebhookEvent string + +const ( + WebhookEventPush WebhookEvent = "push" + WebhookEventRepoRenamed WebhookEvent = "repository:renamed" + WebhookEventPullRequestCreated WebhookEvent = "pull_request:created" + WebhookEventPullRequestResubmitted WebhookEvent = "pull_request:resubmitted" + WebhookEventPullRequestMerged WebhookEvent = "pull_request:merged" + WebhookEventPullRequestClosed WebhookEvent = "pull_request:closed" + WebhookEventPullRequestReopened WebhookEvent = "pull_request:reopened" +) + +type Webhook struct { + Id int64 + RepoDid syntax.DID + 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. +// +// Note: spindle stores less repo metadata than the appview did, so +// Description, Website, StarsCount, OpenIssues and Fork are best-effort and +// may be empty/false when spindle has no record of them. +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"` +} + +// WebhookRenamePayload represents the payload for a repository:renamed event +type WebhookRenamePayload struct { + OldName string `json:"old_name"` + NewName string `json:"new_name"` + Repository WebhookRepository `json:"repository"` + Sender WebhookUser `json:"sender"` +} + +// WebhookPullRequestPayload represents the payload for pull_request:* events +type WebhookPullRequestPayload struct { + Action string `json:"action"` + PullRequest WebhookPullRequest `json:"pull_request"` + Repository WebhookRepository `json:"repository"` + Sender WebhookUser `json:"sender"` +} + +// WebhookPullRequest represents pull request information in webhook payload. +// +// Note: the sh.tangled.repo.pull record carries no appview-assigned pull +// number or state, so Number is 0 and State is derived from the lifecycle +// action. HtmlUrl/PatchUrl are omitted because spindle does not know the +// appview base URL. +type WebhookPullRequest struct { + Number int `json:"number"` + Title string `json:"title"` + Body string `json:"body"` + State string `json:"state"` + TargetBranch string `json:"target_branch"` + Source *WebhookPullRequestSource `json:"source,omitempty"` + RoundNumber int `json:"round_number"` + Owner WebhookUser `json:"owner"` + HtmlUrl string `json:"html_url,omitempty"` + PatchUrl string `json:"patch_url,omitempty"` + CreatedAt string `json:"created_at"` +} + +// WebhookPullRequestSource represents the source of a branch- or fork-based +// pull request; absent for patch-based pull requests +type WebhookPullRequestSource struct { + Branch string `json:"branch"` + Repo string `json:"repo,omitempty"` + Sha string `json:"sha,omitempty"` +}