package 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 needed func (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 resolve func (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 yet func (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) } }