diff --git a/appview/migration/migrate_add_repo_did.go b/appview/migration/migrate_add_repo_did.go --- a/appview/migration/migrate_add_repo_did.go +++ b/appview/migration/migrate_add_repo_did.go @@ -2,6 +2,7 @@ import ( "context" + "encoding/json" "fmt" "strings" @@ -14,7 +15,10 @@ ) func (s *Migration) migrateAddRepoDid(ctx context.Context, client *atclient.APIClient, did syntax.DID, record syntax.ATURI) error { - // TODO: use agnostic.RepoGetRecord instead + if record.Collection().String() == tangled.FeedStarNSID { + return s.migrateAddRepoDidStar(ctx, client, did, record) + } + ex, err := comatproto.RepoGetRecord(ctx, client, "", record.Collection().String(), did.String(), record.RecordKey().String()) if err != nil { return fmt.Errorf("pds: %w", err) @@ -22,7 +26,7 @@ val := ex.Value.Val - switch record.Collection() { + switch record.Collection().String() { case tangled.RepoNSID: rec, ok := val.(*tangled.Repo) if !ok { @@ -39,35 +43,35 @@ if !ok { return fmt.Errorf("unexpected type for issue record") } - if rec.Repo != nil { - repoAt := *rec.Repo - repo, err := db.GetRepoByAtUri(s.db, repoAt) - if err != nil { - return fmt.Errorf("db: failed to query repo: %w", err) - } - rec.RepoDid = &repo.RepoDid + if strings.HasPrefix(rec.Repo, "did:") { + return nil } + repo, err := db.GetRepoByAtUri(s.db, rec.Repo) + if err != nil { + return fmt.Errorf("db: failed to query repo by at_uri %q: %w", rec.Repo, err) + } + rec.Repo = repo.RepoDid case tangled.RepoPullNSID: rec, ok := val.(*tangled.RepoPull) if !ok { return fmt.Errorf("unexpected type for pull record") } - if rec.Target != nil && rec.Target.Repo != nil { - repoAt := *rec.Target.Repo - repo, err := db.GetRepoByAtUri(s.db, repoAt) - if err != nil { - return fmt.Errorf("db: failed to query repo: %w", err) - } - rec.Target.RepoDid = &repo.RepoDid + if rec.Target == nil { + return fmt.Errorf("pull record has nil target") } - if rec.Source != nil && rec.Source.Repo != nil { - repoAt := *rec.Source.Repo - repo, err := db.GetRepoByAtUri(s.db, repoAt) + if !strings.HasPrefix(rec.Target.Repo, "did:") { + repo, err := db.GetRepoByAtUri(s.db, rec.Target.Repo) if err != nil { - return fmt.Errorf("db: failed to query repo: %w", err) + return fmt.Errorf("db: failed to query target repo by at_uri %q: %w", rec.Target.Repo, err) } - rec.Source.RepoDid = &repo.RepoDid + rec.Target.Repo = repo.RepoDid + } + if rec.Source != nil && rec.Source.Repo != nil && !strings.HasPrefix(*rec.Source.Repo, "did:") { + sourceRepo, srcErr := db.GetRepoByAtUri(s.db, *rec.Source.Repo) + if srcErr == nil && sourceRepo.RepoDid != "" { + rec.Source.Repo = &sourceRepo.RepoDid + } } case tangled.RepoCollaboratorNSID: @@ -75,14 +79,14 @@ if !ok { return fmt.Errorf("unexpected type for collaborator record") } - if rec.Repo != nil { - repoAt := *rec.Repo - repo, err := db.GetRepoByAtUri(s.db, repoAt) - if err != nil { - return fmt.Errorf("db: failed to query repo: %w", err) - } - rec.RepoDid = &repo.RepoDid + if strings.HasPrefix(rec.Repo, "did:") { + return nil } + repo, err := db.GetRepoByAtUri(s.db, rec.Repo) + if err != nil { + return fmt.Errorf("db: failed to query repo by at_uri %q: %w", rec.Repo, err) + } + rec.Repo = repo.RepoDid case tangled.RepoArtifactNSID: rec, ok := val.(*tangled.RepoArtifact) @@ -90,26 +94,11 @@ return fmt.Errorf("unexpected type for artifact record") } if rec.Repo != nil { - repoAt := *rec.Repo - repo, err := db.GetRepoByAtUri(s.db, repoAt) + repo, err := db.GetRepoByAtUri(s.db, *rec.Repo) if err != nil { - return fmt.Errorf("db: failed to query repo: %w", err) + return fmt.Errorf("db: failed to query repo by at_uri %q: %w", *rec.Repo, err) } rec.RepoDid = &repo.RepoDid - } - - case tangled.FeedStarNSID: - rec, ok := val.(*tangled.FeedStar) - if !ok { - return fmt.Errorf("unexpected type for star record") - } - if rec.Subject != nil { - repoAt := *rec.Subject - repo, err := db.GetRepoByAtUri(s.db, repoAt) - if err != nil { - return fmt.Errorf("db: failed to query repo: %w", err) - } - rec.SubjectDid = &repo.RepoDid } case tangled.ActorProfileNSID: @@ -147,5 +136,59 @@ return fmt.Errorf("put record: %w", err) } + return nil +} + +func (s *Migration) migrateAddRepoDidStar(ctx context.Context, client *atclient.APIClient, did syntax.DID, record syntax.ATURI) error { + var raw struct { + Cid *string `json:"cid,omitempty"` + Uri string `json:"uri"` + Value json.RawMessage `json:"value"` + } + params := map[string]any{ + "collection": record.Collection().String(), + "repo": did.String(), + "rkey": record.RecordKey().String(), + } + if err := client.LexDo(ctx, lexutil.Query, "", "com.atproto.repo.getRecord", params, nil, &raw); err != nil { + return fmt.Errorf("get record: %w", err) + } + + var legacy struct { + CreatedAt string `json:"createdAt"` + Subject *string `json:"subject,omitempty"` + } + if err := json.Unmarshal(raw.Value, &legacy); err != nil { + return fmt.Errorf("decode old star fields: %w", err) + } + if legacy.Subject == nil { + return fmt.Errorf("star record has no subject field") + } + + repo, err := db.GetRepoByAtUri(s.db, *legacy.Subject) + if err != nil { + return fmt.Errorf("db: failed to query repo by at_uri %q: %w", *legacy.Subject, err) + } + if repo.RepoDid == "" { + return fmt.Errorf("repo has no repoDid: %s", *legacy.Subject) + } + + newRecord := &tangled.FeedStar{ + CreatedAt: legacy.CreatedAt, + Subject: &tangled.FeedStar_Subject{ + FeedStar_Repo: &tangled.FeedStar_Repo{Did: repo.RepoDid}, + }, + } + + _, err = comatproto.RepoPutRecord(ctx, client, &comatproto.RepoPutRecord_Input{ + Repo: did.String(), + Collection: record.Collection().String(), + Rkey: record.RecordKey().String(), + SwapRecord: raw.Cid, + Record: &lexutil.LexiconTypeDecoder{Val: newRecord}, + }) + if err != nil { + return fmt.Errorf("put record: %w", err) + } return nil } diff --git a/appview/notify/merged_notifier.go b/appview/notify/merged_notifier.go --- a/appview/notify/merged_notifier.go +++ b/appview/notify/merged_notifier.go @@ -38,6 +38,10 @@ m.fanout(func(n Notifier) { n.DeleteRepo(ctx, repo) }) } +func (m *mergedNotifier) RenameRepo(ctx context.Context, actor syntax.DID, oldRepo, newRepo *models.Repo) { + m.fanout(func(n Notifier) { n.RenameRepo(ctx, actor, oldRepo, newRepo) }) +} + func (m *mergedNotifier) NewStar(ctx context.Context, star *models.Star) { m.fanout(func(n Notifier) { n.NewStar(ctx, star) }) } diff --git a/appview/notify/notifier.go b/appview/notify/notifier.go --- a/appview/notify/notifier.go +++ b/appview/notify/notifier.go @@ -10,6 +10,7 @@ type Notifier interface { NewRepo(ctx context.Context, repo *models.Repo) DeleteRepo(ctx context.Context, repo *models.Repo) + RenameRepo(ctx context.Context, actor syntax.DID, oldRepo, newRepo *models.Repo) NewStar(ctx context.Context, star *models.Star) DeleteStar(ctx context.Context, star *models.Star) @@ -47,6 +48,8 @@ func (m *BaseNotifier) NewRepo(ctx context.Context, repo *models.Repo) {} func (m *BaseNotifier) DeleteRepo(ctx context.Context, repo *models.Repo) {} +func (m *BaseNotifier) RenameRepo(ctx context.Context, actor syntax.DID, oldRepo, newRepo *models.Repo) { +} func (m *BaseNotifier) NewStar(ctx context.Context, star *models.Star) {} func (m *BaseNotifier) DeleteStar(ctx context.Context, star *models.Star) {} diff --git a/appview/notify/db/db.go b/appview/notify/db/db.go --- a/appview/notify/db/db.go +++ b/appview/notify/db/db.go @@ -5,7 +5,6 @@ "slices" "github.com/bluesky-social/indigo/atproto/syntax" - "tangled.org/core/api/tangled" "tangled.org/core/appview/db" "tangled.org/core/appview/models" "tangled.org/core/appview/notify" @@ -40,15 +39,17 @@ // no-op for now } +func (n *databaseNotifier) RenameRepo(ctx context.Context, actor syntax.DID, oldRepo, newRepo *models.Repo) { +} + func (n *databaseNotifier) NewStar(ctx context.Context, star *models.Star) { l := log.FromContext(ctx) - if star.RepoAt.Collection().String() != tangled.RepoNSID { - // skip string stars for now + if star.SubjectType != models.StarSubjectRepo { return } - var err error - repo, err := db.GetRepo(n.db, orm.FilterEq("at_uri", string(star.RepoAt))) + + repo, err := db.GetRepo(n.db, orm.FilterEq("repo_did", star.Subject)) if err != nil { l.Error("failed to get repos", "err", err) return @@ -58,7 +59,7 @@ recipients := sets.Singleton(syntax.DID(repo.Did)) eventType := models.NotificationTypeRepoStarred entityType := "repo" - entityId := star.RepoAt.String() + entityId := star.Subject repoId := &repo.Id var issueId *int64 var pullId *int64 @@ -83,7 +84,7 @@ func (n *databaseNotifier) NewIssue(ctx context.Context, issue *models.Issue, mentions []syntax.DID) { l := log.FromContext(ctx) - collaborators, err := db.GetCollaborators(n.db, orm.FilterEq("repo_at", issue.Repo.RepoAt())) + collaborators, err := db.GetCollaborators(n.db, orm.FilterEq("repo_did", string(issue.RepoDid))) if err != nil { l.Error("failed to fetch collaborators", "err", err) return @@ -240,12 +241,12 @@ func (n *databaseNotifier) NewPull(ctx context.Context, pull *models.Pull) { l := log.FromContext(ctx) - repo, err := db.GetRepo(n.db, orm.FilterEq("at_uri", string(pull.RepoAt))) + repo, err := db.GetRepo(n.db, orm.FilterEq("repo_did", string(pull.RepoDid))) if err != nil { l.Error("failed to get repos", "err", err) return } - collaborators, err := db.GetCollaborators(n.db, orm.FilterEq("repo_at", repo.RepoAt())) + collaborators, err := db.GetCollaborators(n.db, orm.FilterEq("repo_did", string(pull.RepoDid))) if err != nil { l.Error("failed to fetch collaborators", "err", err) return @@ -285,7 +286,7 @@ l := log.FromContext(ctx) pull, err := db.GetPull(n.db, - orm.FilterEq("repo_at", syntax.ATURI(comment.RepoAt)), + orm.FilterEq("repo_did", comment.RepoDid), orm.FilterEq("pull_id", comment.PullId), ) if err != nil { @@ -293,7 +294,7 @@ return } - repo, err := db.GetRepo(n.db, orm.FilterEq("at_uri", comment.RepoAt)) + repo, err := db.GetRepo(n.db, orm.FilterEq("repo_did", comment.RepoDid)) if err != nil { l.Error("failed to get repos", "err", err) return @@ -371,7 +372,7 @@ func (n *databaseNotifier) NewIssueState(ctx context.Context, actor syntax.DID, issue *models.Issue) { l := log.FromContext(ctx) - collaborators, err := db.GetCollaborators(n.db, orm.FilterEq("repo_at", issue.Repo.RepoAt())) + collaborators, err := db.GetCollaborators(n.db, orm.FilterEq("repo_did", string(issue.RepoDid))) if err != nil { l.Error("failed to fetch collaborators", "err", err) return @@ -419,13 +420,13 @@ l := log.FromContext(ctx) // Get repo details - repo, err := db.GetRepo(n.db, orm.FilterEq("at_uri", string(pull.RepoAt))) + repo, err := db.GetRepo(n.db, orm.FilterEq("repo_did", string(pull.RepoDid))) if err != nil { l.Error("failed to get repos", "err", err) return } - collaborators, err := db.GetCollaborators(n.db, orm.FilterEq("repo_at", repo.RepoAt())) + collaborators, err := db.GetCollaborators(n.db, orm.FilterEq("repo_did", string(pull.RepoDid))) if err != nil { l.Error("failed to fetch collaborators", "err", err) return diff --git a/appview/notify/logging/notifier.go b/appview/notify/logging/notifier.go --- a/appview/notify/logging/notifier.go +++ b/appview/notify/logging/notifier.go @@ -31,6 +31,11 @@ l.inner.DeleteRepo(ctx, repo) } +func (l *loggingNotifier) RenameRepo(ctx context.Context, actor syntax.DID, oldRepo, newRepo *models.Repo) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "RenameRepo")) + l.inner.RenameRepo(ctx, actor, oldRepo, newRepo) +} + func (l *loggingNotifier) NewStar(ctx context.Context, star *models.Star) { ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewStar")) l.inner.NewStar(ctx, star) diff --git a/appview/notify/posthog/notifier.go b/appview/notify/posthog/notifier.go --- a/appview/notify/posthog/notifier.go +++ b/appview/notify/posthog/notifier.go @@ -35,11 +35,30 @@ } } +func (n *posthogNotifier) RenameRepo(ctx context.Context, actor syntax.DID, oldRepo, newRepo *models.Repo) { + err := n.client.Enqueue(posthog.Capture{ + DistinctId: actor.String(), + Event: "repo_renamed", + Properties: posthog.Properties{ + "repo_at": newRepo.RepoAt(), + "owner": newRepo.Did, + "old_name": oldRepo.Name, + "new_name": newRepo.Name, + }, + }) + if err != nil { + log.Println("failed to enqueue posthog event:", err) + } +} + func (n *posthogNotifier) NewStar(ctx context.Context, star *models.Star) { err := n.client.Enqueue(posthog.Capture{ DistinctId: star.Did, Event: "star", - Properties: posthog.Properties{"repo_at": star.RepoAt.String()}, + Properties: posthog.Properties{ + "subject_type": string(star.SubjectType), + "subject": star.Subject, + }, }) if err != nil { log.Println("failed to enqueue posthog event:", err) @@ -50,7 +69,10 @@ err := n.client.Enqueue(posthog.Capture{ DistinctId: star.Did, Event: "unstar", - Properties: posthog.Properties{"repo_at": star.RepoAt.String()}, + Properties: posthog.Properties{ + "subject_type": string(star.SubjectType), + "subject": star.Subject, + }, }) if err != nil { log.Println("failed to enqueue posthog event:", err) @@ -62,7 +84,7 @@ DistinctId: issue.Did, Event: "new_issue", Properties: posthog.Properties{ - "repo_at": issue.RepoAt.String(), + "repo_did": string(issue.RepoDid), "issue_id": issue.IssueId, "mentions": mentions, }, @@ -77,8 +99,8 @@ DistinctId: pull.OwnerDid, Event: "new_pull", Properties: posthog.Properties{ - "repo_at": pull.RepoAt, - "pull_id": pull.PullId, + "repo_did": string(pull.RepoDid), + "pull_id": pull.PullId, }, }) if err != nil { @@ -91,7 +113,7 @@ DistinctId: comment.OwnerDid, Event: "new_pull_comment", Properties: posthog.Properties{ - "repo_at": comment.RepoAt, + "repo_did": comment.RepoDid, "pull_id": comment.PullId, "mentions": mentions, }, @@ -106,8 +128,8 @@ DistinctId: pull.OwnerDid, Event: "pull_closed", Properties: posthog.Properties{ - "repo_at": pull.RepoAt, - "pull_id": pull.PullId, + "repo_did": string(pull.RepoDid), + "pull_id": pull.PullId, }, }) if err != nil { @@ -216,7 +238,7 @@ DistinctId: issue.Did, Event: event, Properties: posthog.Properties{ - "repo_at": issue.RepoAt.String(), + "repo_did": string(issue.RepoDid), "actor": actor, "issue_id": issue.IssueId, }, @@ -243,9 +265,9 @@ DistinctId: pull.OwnerDid, Event: event, Properties: posthog.Properties{ - "repo_at": pull.RepoAt, - "pull_id": pull.PullId, - "actor": actor, + "repo_did": string(pull.RepoDid), + "pull_id": pull.PullId, + "actor": actor, }, }) if err != nil { diff --git a/appview/notify/webhook/notifier.go b/appview/notify/webhook/notifier.go --- a/appview/notify/webhook/notifier.go +++ b/appview/notify/webhook/notifier.go @@ -14,6 +14,7 @@ "time" "github.com/avast/retry-go/v4" + "github.com/bluesky-social/indigo/atproto/syntax" "github.com/google/uuid" "tangled.org/core/appview/db" "tangled.org/core/appview/models" @@ -41,57 +42,85 @@ var _ notify.Notifier = &Notifier{} func (w *Notifier) Push(ctx context.Context, repo *models.Repo, ref, oldSha, newSha, committerDid string) { - webhooks, err := db.GetActiveWebhooksForRepo(w.db, repo.RepoAt()) + webhooks, err := w.activeWebhooksForEvent(repo.RepoDid, models.WebhookEventPush) if err != nil { - w.logger.Error("failed to get webhooks for repo", "repo", repo.RepoAt(), "err", err) + w.logger.Error("failed to get webhooks for repo", "repo_did", repo.RepoDid, "err", err) + return + } + if len(webhooks) == 0 { return } - var pushWebhooks []models.Webhook + payload := w.buildPushPayload(repo, ref, oldSha, newSha, committerDid) + payloadBytes, err := json.Marshal(payload) + if err != nil { + w.logger.Error("failed to marshal push payload", "repo_did", repo.RepoDid, "err", err) + return + } + + userAgent := "Tangled-Hook/" + newSha[:7] for _, webhook := range webhooks { - if webhook.HasEvent(models.WebhookEventPush) { - pushWebhooks = append(pushWebhooks, webhook) - } - } - - if len(pushWebhooks) == 0 { - return - } - - payload, err := w.buildPushPayload(repo, ref, oldSha, newSha, committerDid) - if err != nil { - w.logger.Error("failed to build push payload", "repo", repo.RepoAt(), "err", err) - return - } - - for _, webhook := range pushWebhooks { - go w.sendWebhook(ctx, webhook, string(models.WebhookEventPush), payload) + go w.sendWebhook(ctx, webhook, string(models.WebhookEventPush), payload.Repository.FullName, userAgent, payloadBytes) } } -func (w *Notifier) buildPushPayload(repo *models.Repo, ref, oldSha, newSha, committerDid string) (*models.WebhookPayload, error) { - owner := repo.Did - - pusher := committerDid - if committerDid == "" { - pusher = owner +func (w *Notifier) RenameRepo(ctx context.Context, actor syntax.DID, oldRepo, newRepo *models.Repo) { + webhooks, err := w.activeWebhooksForEvent(newRepo.RepoDid, models.WebhookEventRepoRenamed) + if err != nil { + w.logger.Error("failed to get webhooks for repo", "repo_did", newRepo.RepoDid, "err", err) + return + } + if len(webhooks) == 0 { + return } + payload := &models.WebhookRenamePayload{ + OldName: oldRepo.Name, + NewName: newRepo.Name, + Repository: buildWebhookRepository(newRepo), + Sender: models.WebhookUser{Did: actor.String()}, + } + payloadBytes, err := json.Marshal(payload) + if err != nil { + w.logger.Error("failed to marshal rename payload", "repo_did", newRepo.RepoDid, "err", err) + return + } + + userAgent := "Tangled-Hook/rename" + for _, webhook := range webhooks { + go w.sendWebhook(ctx, webhook, string(models.WebhookEventRepoRenamed), payload.Repository.FullName, userAgent, payloadBytes) + } +} + +func (w *Notifier) activeWebhooksForEvent(repoDid string, event models.WebhookEvent) ([]models.Webhook, error) { + webhooks, err := db.GetActiveWebhooksForRepo(w.db, repoDid) + if err != nil { + return nil, err + } + var matching []models.Webhook + for _, webhook := range webhooks { + if webhook.HasEvent(event) { + matching = append(matching, webhook) + } + } + return matching, nil +} + +func buildWebhookRepository(repo *models.Repo) models.WebhookRepository { repository := models.WebhookRepository{ Name: repo.Name, - FullName: fmt.Sprintf("%s/%s", repo.Did, repo.Name), + FullName: fmt.Sprintf("%s/%s", repo.Did, repo.Rkey), Description: repo.Description, Fork: repo.Source != "", - HtmlUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, repo.Did, repo.Name), - CloneUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, repo.Did, repo.Name), - SshUrl: fmt.Sprintf("ssh://git@%s/%s/%s", repo.Knot, repo.Did, repo.Name), + HtmlUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, repo.Did, repo.Rkey), + CloneUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, repo.Did, repo.Rkey), + SshUrl: fmt.Sprintf("ssh://git@%s/%s/%s", repo.Knot, repo.Did, repo.Rkey), CreatedAt: repo.Created.Format(time.RFC3339), UpdatedAt: repo.Created.Format(time.RFC3339), Owner: models.WebhookUser{ - Did: owner, + Did: repo.Did, }, } - if repo.Website != "" { repository.Website = repo.Website } @@ -99,28 +128,27 @@ repository.StarsCount = repo.RepoStats.StarCount repository.OpenIssues = repo.RepoStats.IssueCount.Open } + return repository +} - payload := &models.WebhookPayload{ +func (w *Notifier) buildPushPayload(repo *models.Repo, ref, oldSha, newSha, committerDid string) *models.WebhookPayload { + pusher := committerDid + if committerDid == "" { + pusher = repo.Did + } + return &models.WebhookPayload{ Ref: ref, Before: oldSha, After: newSha, - Repository: repository, + Repository: buildWebhookRepository(repo), Pusher: models.WebhookUser{ Did: pusher, }, } - - return payload, nil } -func (w *Notifier) sendWebhook(ctx context.Context, webhook models.Webhook, event string, payload *models.WebhookPayload) { +func (w *Notifier) sendWebhook(ctx context.Context, webhook models.Webhook, event, repoFullName, userAgent string, payloadBytes []byte) { deliveryId := uuid.New().String() - - payloadBytes, err := json.Marshal(payload) - if err != nil { - w.logger.Error("failed to marshal webhook payload", "webhook_id", webhook.Id, "err", err) - return - } req, err := http.NewRequestWithContext(ctx, "POST", webhook.Url, bytes.NewReader(payloadBytes)) if err != nil { @@ -128,14 +156,12 @@ return } - shortSha := payload.After[:7] - req.Header.Set("Content-Type", "application/json") - req.Header.Set("User-Agent", "Tangled-Hook/"+shortSha) + req.Header.Set("User-Agent", userAgent) req.Header.Set("X-Tangled-Event", event) req.Header.Set("X-Tangled-Hook-ID", fmt.Sprintf("%d", webhook.Id)) req.Header.Set("X-Tangled-Delivery", deliveryId) - req.Header.Set("X-Tangled-Repo", payload.Repository.FullName) + req.Header.Set("X-Tangled-Repo", repoFullName) if webhook.Secret != "" { signature := w.computeSignature(payloadBytes, webhook.Secret)