Monorepo for Tangled forked from tangled.org/core
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495package db
import ( "context" "log" "maps" "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" "tangled.org/core/idresolver")
const ( maxMentions = 5)
type databaseNotifier struct { db *db.DB res *idresolver.Resolver}
func NewDatabaseNotifier(database *db.DB, resolver *idresolver.Resolver) notify.Notifier { return &databaseNotifier{ db: database, res: resolver, }}
var _ notify.Notifier = &databaseNotifier{}
func (n *databaseNotifier) NewRepo(ctx context.Context, repo *models.Repo) { // no-op for now}
func (n *databaseNotifier) NewStar(ctx context.Context, star *models.Star) { if star.RepoAt.Collection().String() != tangled.RepoNSID { // skip string stars for now return } var err error repo, err := db.GetRepo(n.db, db.FilterEq("at_uri", string(star.RepoAt))) if err != nil { log.Printf("NewStar: failed to get repos: %v", err) return }
actorDid := syntax.DID(star.Did) recipients := []syntax.DID{syntax.DID(repo.Did)} eventType := models.NotificationTypeRepoStarred entityType := "repo" entityId := star.RepoAt.String() repoId := &repo.Id var issueId *int64 var pullId *int64
n.notifyEvent( actorDid, recipients, eventType, entityType, entityId, repoId, issueId, pullId, )}
func (n *databaseNotifier) DeleteStar(ctx context.Context, star *models.Star) { // no-op}
func (n *databaseNotifier) NewIssue(ctx context.Context, issue *models.Issue, mentions []syntax.DID) {
// build the recipients list // - owner of the repo // - collaborators in the repo var recipients []syntax.DID recipients = append(recipients, syntax.DID(issue.Repo.Did)) collaborators, err := db.GetCollaborators(n.db, db.FilterEq("repo_at", issue.Repo.RepoAt())) if err != nil { log.Printf("failed to fetch collaborators: %v", err) return } for _, c := range collaborators { recipients = append(recipients, c.SubjectDid) }
actorDid := syntax.DID(issue.Did) entityType := "issue" entityId := issue.AtUri().String() repoId := &issue.Repo.Id issueId := &issue.Id var pullId *int64
n.notifyEvent( actorDid, recipients, models.NotificationTypeIssueCreated, entityType, entityId, repoId, issueId, pullId, ) n.notifyEvent( actorDid, mentions, models.NotificationTypeUserMentioned, entityType, entityId, repoId, issueId, pullId, )}
func (n *databaseNotifier) NewIssueComment(ctx context.Context, comment *models.IssueComment, mentions []syntax.DID) { issues, err := db.GetIssues(n.db, db.FilterEq("at_uri", comment.IssueAt)) if err != nil { log.Printf("NewIssueComment: failed to get issues: %v", err) return } if len(issues) == 0 { log.Printf("NewIssueComment: no issue found for %s", comment.IssueAt) return } issue := issues[0]
var recipients []syntax.DID recipients = append(recipients, syntax.DID(issue.Repo.Did))
if comment.IsReply() { // if this comment is a reply, then notify everybody in that thread parentAtUri := *comment.ReplyTo allThreads := issue.CommentList()
// find the parent thread, and add all DIDs from here to the recipient list for _, t := range allThreads { if t.Self.AtUri().String() == parentAtUri { recipients = append(recipients, t.Participants()...) } } } else { // not a reply, notify just the issue author recipients = append(recipients, syntax.DID(issue.Did)) }
actorDid := syntax.DID(comment.Did) entityType := "issue" entityId := issue.AtUri().String() repoId := &issue.Repo.Id issueId := &issue.Id var pullId *int64
n.notifyEvent( actorDid, recipients, models.NotificationTypeIssueCommented, entityType, entityId, repoId, issueId, pullId, ) n.notifyEvent( actorDid, mentions, models.NotificationTypeUserMentioned, entityType, entityId, repoId, issueId, pullId, )}
func (n *databaseNotifier) DeleteIssue(ctx context.Context, issue *models.Issue) { // no-op for now}
func (n *databaseNotifier) NewFollow(ctx context.Context, follow *models.Follow) { actorDid := syntax.DID(follow.UserDid) recipients := []syntax.DID{syntax.DID(follow.SubjectDid)} eventType := models.NotificationTypeFollowed entityType := "follow" entityId := follow.UserDid var repoId, issueId, pullId *int64
n.notifyEvent( actorDid, recipients, eventType, entityType, entityId, repoId, issueId, pullId, )}
func (n *databaseNotifier) DeleteFollow(ctx context.Context, follow *models.Follow) { // no-op}
func (n *databaseNotifier) NewPull(ctx context.Context, pull *models.Pull) { repo, err := db.GetRepo(n.db, db.FilterEq("at_uri", string(pull.RepoAt))) if err != nil { log.Printf("NewPull: failed to get repos: %v", err) return }
// build the recipients list // - owner of the repo // - collaborators in the repo var recipients []syntax.DID recipients = append(recipients, syntax.DID(repo.Did)) collaborators, err := db.GetCollaborators(n.db, db.FilterEq("repo_at", repo.RepoAt())) if err != nil { log.Printf("failed to fetch collaborators: %v", err) return } for _, c := range collaborators { recipients = append(recipients, c.SubjectDid) }
actorDid := syntax.DID(pull.OwnerDid) eventType := models.NotificationTypePullCreated entityType := "pull" entityId := pull.AtUri().String() repoId := &repo.Id var issueId *int64 p := int64(pull.ID) pullId := &p
n.notifyEvent( actorDid, recipients, eventType, entityType, entityId, repoId, issueId, pullId, )}
func (n *databaseNotifier) NewPullComment(ctx context.Context, comment *models.PullComment, mentions []syntax.DID) { pull, err := db.GetPull(n.db, syntax.ATURI(comment.RepoAt), comment.PullId, ) if err != nil { log.Printf("NewPullComment: failed to get pulls: %v", err) return }
repo, err := db.GetRepo(n.db, db.FilterEq("at_uri", comment.RepoAt)) if err != nil { log.Printf("NewPullComment: failed to get repos: %v", err) return }
// build up the recipients list: // - repo owner // - all pull participants var recipients []syntax.DID recipients = append(recipients, syntax.DID(repo.Did)) for _, p := range pull.Participants() { recipients = append(recipients, syntax.DID(p)) }
actorDid := syntax.DID(comment.OwnerDid) eventType := models.NotificationTypePullCommented entityType := "pull" entityId := pull.AtUri().String() repoId := &repo.Id var issueId *int64 p := int64(pull.ID) pullId := &p
n.notifyEvent( actorDid, recipients, eventType, entityType, entityId, repoId, issueId, pullId, ) n.notifyEvent( actorDid, mentions, models.NotificationTypeUserMentioned, entityType, entityId, repoId, issueId, pullId, )}
func (n *databaseNotifier) UpdateProfile(ctx context.Context, profile *models.Profile) { // no-op}
func (n *databaseNotifier) DeleteString(ctx context.Context, did, rkey string) { // no-op}
func (n *databaseNotifier) EditString(ctx context.Context, string *models.String) { // no-op}
func (n *databaseNotifier) NewString(ctx context.Context, string *models.String) { // no-op}
func (n *databaseNotifier) NewIssueState(ctx context.Context, actor syntax.DID, issue *models.Issue) { // build up the recipients list: // - repo owner // - repo collaborators // - all issue participants var recipients []syntax.DID recipients = append(recipients, syntax.DID(issue.Repo.Did)) collaborators, err := db.GetCollaborators(n.db, db.FilterEq("repo_at", issue.Repo.RepoAt())) if err != nil { log.Printf("failed to fetch collaborators: %v", err) return } for _, c := range collaborators { recipients = append(recipients, c.SubjectDid) } for _, p := range issue.Participants() { recipients = append(recipients, syntax.DID(p)) }
entityType := "pull" entityId := issue.AtUri().String() repoId := &issue.Repo.Id issueId := &issue.Id var pullId *int64 var eventType models.NotificationType
if issue.Open { eventType = models.NotificationTypeIssueReopen } else { eventType = models.NotificationTypeIssueClosed }
n.notifyEvent( actor, recipients, eventType, entityType, entityId, repoId, issueId, pullId, )}
func (n *databaseNotifier) NewPullState(ctx context.Context, actor syntax.DID, pull *models.Pull) { // Get repo details repo, err := db.GetRepo(n.db, db.FilterEq("at_uri", string(pull.RepoAt))) if err != nil { log.Printf("NewPullState: failed to get repos: %v", err) return }
// build up the recipients list: // - repo owner // - all pull participants var recipients []syntax.DID recipients = append(recipients, syntax.DID(repo.Did)) collaborators, err := db.GetCollaborators(n.db, db.FilterEq("repo_at", repo.RepoAt())) if err != nil { log.Printf("failed to fetch collaborators: %v", err) return } for _, c := range collaborators { recipients = append(recipients, c.SubjectDid) } for _, p := range pull.Participants() { recipients = append(recipients, syntax.DID(p)) }
entityType := "pull" entityId := pull.AtUri().String() repoId := &repo.Id var issueId *int64 var eventType models.NotificationType switch pull.State { case models.PullClosed: eventType = models.NotificationTypePullClosed case models.PullOpen: eventType = models.NotificationTypePullReopen case models.PullMerged: eventType = models.NotificationTypePullMerged default: log.Println("NewPullState: unexpected new PR state:", pull.State) return } p := int64(pull.ID) pullId := &p
n.notifyEvent( actor, recipients, eventType, entityType, entityId, repoId, issueId, pullId, )}
func (n *databaseNotifier) notifyEvent( actorDid syntax.DID, recipients []syntax.DID, eventType models.NotificationType, entityType string, entityId string, repoId *int64, issueId *int64, pullId *int64,) { if eventType == models.NotificationTypeUserMentioned && len(recipients) > maxMentions { recipients = recipients[:maxMentions] } recipientSet := make(map[syntax.DID]struct{}) for _, did := range recipients { // everybody except actor themselves if did != actorDid { recipientSet[did] = struct{}{} } }
prefMap, err := db.GetNotificationPreferences( n.db, db.FilterIn("user_did", slices.Collect(maps.Keys(recipientSet))), ) if err != nil { // failed to get prefs for users return }
// create a transaction for bulk notification storage tx, err := n.db.Begin() if err != nil { // failed to start tx return } defer tx.Rollback()
// filter based on preferences for recipientDid := range recipientSet { prefs, ok := prefMap[recipientDid] if !ok { prefs = models.DefaultNotificationPreferences(recipientDid) }
// skip users who don’t want this type if !prefs.ShouldNotify(eventType) { continue }
// create notification notif := &models.Notification{ RecipientDid: recipientDid.String(), ActorDid: actorDid.String(), Type: eventType, EntityType: entityType, EntityId: entityId, RepoId: repoId, IssueId: issueId, PullId: pullId, }
if err := db.CreateNotification(tx, notif); err != nil { log.Printf("notifyEvent: failed to create notification for %s: %v", recipientDid, err) } }
if err := tx.Commit(); err != nil { // failed to commit return }}