Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328package deliberi
import ( "context" "encoding/json" "fmt" "log/slog" "net/http" "time"
"github.com/bluesky-social/indigo/atproto/syntax" indigoxrpc "github.com/bluesky-social/indigo/xrpc" jmodels "github.com/bluesky-social/jetstream/pkg/models" deldb "tangled.org/core/deliberi/db" models "tangled.org/core/deliberi/models" "tangled.org/core/idresolver" js "tangled.org/core/jetstream")
type Ingester struct { db *deldb.DB recipients recipientResolver jc *js.JetstreamClient idResolver *idresolver.Resolver logger *slog.Logger}
var ingestCollections = []string{ repoManifestNSID, ticketNSID, commentNSID, starNSID, followNSID,}
func NewIngester(database *deldb.DB, recipients recipientResolver, idRes *idresolver.Resolver, endpoint, ident string, logger *slog.Logger) (*Ingester, error) { jc, err := js.NewJetstreamClient(endpoint, ident, ingestCollections, nil, logger, database, false, false) if err != nil { return nil, fmt.Errorf("creating jetstream client: %w", err) } return &Ingester{ db: database, recipients: recipients, jc: jc, idResolver: idRes, logger: logger, }, nil}
func (i *Ingester) Run(ctx context.Context) error { return i.jc.StartJetstream(ctx, i.process)}
func (i *Ingester) process(ctx context.Context, e *jmodels.Event) error { if e.Kind != jmodels.EventKindCommit || e.Commit == nil { return nil } if e.Commit.Operation != jmodels.CommitOperationCreate && e.Commit.Operation != jmodels.CommitOperationUpdate { return nil } // updates only refresh caches: ticket state changes arrive as ticket // updates, and edits must not notify again created := e.Commit.Operation == jmodels.CommitOperationCreate sourceAt := fmt.Sprintf("at://%s/%s/%s", e.Did, e.Commit.Collection, e.Commit.RKey)
switch e.Commit.Collection { case repoManifestNSID: var rec repoManifest if err := json.Unmarshal(e.Commit.Record, &rec); err != nil { i.logger.Warn("decoding manifest record", "err", err, "uri", sourceAt) return nil } if rec.Declaration == nil || rec.Declaration.Owner == "" { return nil } if err := deldb.PutRepoName(i.db, e.Did, rec.Declaration.Owner, rec.Declaration.Slug); err != nil { i.logger.Warn("caching repo name", "err", err, "repoDid", e.Did) }
case ticketNSID: var rec ticketRecord if err := json.Unmarshal(e.Commit.Record, &rec); err != nil { i.logger.Warn("decoding ticket record", "err", err, "uri", sourceAt) return nil } kind := rec.kind() if err := deldb.PutEntity(i.db, sourceAt, deldb.Entity{Title: rec.title(), RepoDid: e.Did, Kind: kind}); err != nil { i.logger.Warn("caching ticket", "err", err, "uri", sourceAt) } if !created { return nil } actor := editorOr(rec.Editor, e.Did) t := models.NotificationTypeIssueCreated if kind == models.EntityKindPull { t = models.NotificationTypePullCreated } i.notifyEntity(ctx, models.Notification{ AtUri: sourceAt, Type: t, ActorDid: actor, RepoDid: e.Did, EntityAt: sourceAt, EntityKind: kind, EntityTitle: rec.title(), }, i.mentionDids(ctx, actor, bodyText(rec.Body)))
case commentNSID: if !created { return nil } var rec commentRecord if err := json.Unmarshal(e.Commit.Record, &rec); err != nil { i.logger.Warn("decoding comment record", "err", err, "uri", sourceAt) return nil } if rec.Subject == nil { return nil } subject, err := syntax.ParseATURI(rec.Subject.Uri) if err != nil || subject.Collection().String() != ticketNSID || subject.Authority().String() != e.Did { return nil } ticket := i.hydrateTicket(ctx, subject) var t models.NotificationType switch ticket.Kind { case models.EntityKindIssue: t = models.NotificationTypeIssueCommented case models.EntityKindPull: t = models.NotificationTypePullCommented default: return nil } actor := editorOr(rec.Editor, e.Did) i.notifyEntity(ctx, models.Notification{ AtUri: sourceAt, Type: t, ActorDid: actor, RepoDid: e.Did, EntityAt: subject.String(), EntityKind: ticket.Kind, EntityTitle: ticket.Title, }, i.mentionDids(ctx, actor, bodyText(rec.Body)))
case starNSID: if !created { return nil } var rec starRecord if err := json.Unmarshal(e.Commit.Record, &rec); err != nil { i.logger.Warn("decoding star record", "err", err, "uri", sourceAt) return nil } if rec.Subject == nil || rec.Subject.Type != starRepoType || rec.Subject.Identity == "" { return nil } repoDid := rec.Subject.Identity ownerDid := i.hydrateRepoOwner(ctx, repoDid) if ownerDid == "" { i.logger.Warn("star: could not resolve repo owner, skipping notification", "repoDid", repoDid) return nil } i.deliver(ownerDid, models.Notification{ AtUri: sourceAt, Type: models.NotificationTypeRepoStarred, ActorDid: e.Did, RepoDid: repoDid, })
case followNSID: if !created { return nil } subject, err := syntax.ParseDID(e.Commit.RKey) if err != nil { return nil } i.deliver(subject.String(), models.Notification{ AtUri: sourceAt, Type: models.NotificationTypeFollowed, ActorDid: e.Did, }) }
return nil}
func (i *Ingester) notifyEntity(ctx context.Context, n models.Notification, mentions []string) { seen := make(map[string]struct{})
// mentions go first and claim the recipient: a mentioned subscriber gets // the "mentioned you" notification instead of a duplicate generic one mention := n mention.Type = models.NotificationTypeUserMentioned for _, did := range mentions { if _, ok := seen[did]; ok { continue } seen[did] = struct{}{} i.deliver(did, mention) }
// bobbin matches subjects exactly, so ask at both levels var subscribers []string for _, subject := range []string{n.EntityAt, n.RepoDid} { if subject == "" { continue } dids, err := i.recipients.ListRecipients(ctx, subject, ticketNSID) if err != nil { i.logger.Warn("listing recipients", "err", err, "subject", subject) continue } subscribers = append(subscribers, dids...) }
for _, did := range subscribers { if _, ok := seen[did]; ok { continue } seen[did] = struct{}{} i.deliver(did, n) }}
// falls back to bobbin on a cache miss so stars on repos the ingester never// saw still notify their owner; the answer is cached, so no backfill is neededfunc (i *Ingester) hydrateRepoOwner(ctx context.Context, repoDid string) string { if owner := deldb.GetRepoOwner(i.db, repoDid); owner != "" { return owner } if i.recipients == nil { return "" }
owner, name, err := i.recipients.RepoOwner(ctx, repoDid) if err != nil { i.logger.Warn("resolving repo owner", "err", err, "repoDid", repoDid) return "" } if owner == "" { return "" }
// a nameless record must not blank out a name we already cached if name == "" { name = deldb.GetRepoName(i.db, repoDid) } if err := deldb.PutRepoName(i.db, repoDid, owner, name); err != nil { i.logger.Warn("caching repo owner", "err", err, "repoDid", repoDid) } return owner}
// falls back to the repo's pds on a cache miss so comments on tickets the// ingester never saw still resolvefunc (i *Ingester) hydrateTicket(ctx context.Context, uri syntax.ATURI) deldb.Entity { cached := deldb.GetEntity(i.db, uri.String()) if cached.Kind != "" || i.idResolver == nil { return cached }
rec, err := i.fetchTicket(ctx, uri) if err != nil { i.logger.Warn("hydrating ticket", "err", err, "uri", uri) return cached } ent := deldb.Entity{Title: rec.title(), RepoDid: uri.Authority().String(), Kind: rec.kind()} if err := deldb.PutEntity(i.db, uri.String(), ent); err != nil { i.logger.Warn("caching hydrated ticket", "err", err, "uri", uri) } return ent}
// indigo's typed getRecord only decodes registered $types, and org.tangled// records have no go bindings yetfunc (i *Ingester) fetchTicket(ctx context.Context, uri syntax.ATURI) (*ticketRecord, error) { ident, err := i.idResolver.ResolveIdent(ctx, uri.Authority().String()) if err != nil { return nil, fmt.Errorf("resolving %s: %w", uri.Authority(), err) }
xc := &indigoxrpc.Client{ Host: ident.PDSEndpoint(), Client: &http.Client{Timeout: 10 * time.Second}, } var out struct { Value json.RawMessage `json:"value"` } params := map[string]any{ "repo": ident.DID.String(), "collection": ticketNSID, "rkey": uri.RecordKey().String(), } if err := xc.Do(ctx, indigoxrpc.Query, "", "com.atproto.repo.getRecord", params, nil, &out); err != nil { return nil, fmt.Errorf("getting record: %w", err) } if len(out.Value) == 0 { return nil, fmt.Errorf("record has no value") }
var rec ticketRecord if err := json.Unmarshal(out.Value, &rec); err != nil { return nil, fmt.Errorf("decoding ticket: %w", err) } return &rec, nil}
func (i *Ingester) deliver(recipientDid string, n models.Notification) { if recipientDid == "" || recipientDid == n.ActorDid { return }
prefs, err := deldb.GetNotificationPreference(i.db, recipientDid) if err != nil { i.logger.Warn("loading prefs", "err", err, "recipient", recipientDid) return } if !prefs.ShouldNotify(n.Type) { return }
n.RecipientDid = recipientDid if err := deldb.CreateNotification(i.db, &n); err != nil { i.logger.Warn("creating notification", "err", err, "recipient", recipientDid, "uri", n.AtUri) }}