diff --git a/deliberi/db/cache.go b/deliberi/db/cache.go index 2f19c65a..2e19c266 100644 --- a/deliberi/db/cache.go +++ b/deliberi/db/cache.go @@ -21,11 +21,14 @@ func GetRepoOwner(e Execer, repoDid string) string { return owner } -func PutEntityTitle(e Execer, atUri, title string) error { +func PutEntityTitle(e Execer, atUri, title, repoDid string) error { _, err := e.Exec( - `insert into entity_titles (at_uri, title) values (?, ?) - on conflict(at_uri) do update set title = excluded.title`, - atUri, title, + // a blank repo_did means unknown, so never overwrite a stored one. + `insert into entity_titles (at_uri, title, repo_did) values (?, ?, ?) + on conflict(at_uri) do update set + title = excluded.title, + repo_did = case when excluded.repo_did != '' then excluded.repo_did else entity_titles.repo_did end`, + atUri, title, repoDid, ) return err } @@ -35,3 +38,9 @@ func GetEntityTitle(e Execer, atUri string) string { _ = e.QueryRow(`select title from entity_titles where at_uri = ?`, atUri).Scan(&title) return title } + +func GetEntityRepo(e Execer, atUri string) string { + var repoDid string + _ = e.QueryRow(`select repo_did from entity_titles where at_uri = ?`, atUri).Scan(&repoDid) + return repoDid +} diff --git a/deliberi/db/db.go b/deliberi/db/db.go index 51acbc6f..12806332 100644 --- a/deliberi/db/db.go +++ b/deliberi/db/db.go @@ -106,9 +106,11 @@ create table if not exists repo_names ( name text not null, owner_did text not null default '' ); +-- repo_did lets a comment find its parent's repo. create table if not exists entity_titles ( at_uri text primary key, - title text not null + title text not null, + repo_did text not null default '' ); create table if not exists jetstream_cursor ( diff --git a/deliberi/deliberi.go b/deliberi/deliberi.go index 69e0c56b..c19764f0 100644 --- a/deliberi/deliberi.go +++ b/deliberi/deliberi.go @@ -30,7 +30,7 @@ func Run(ctx context.Context, cfg *config.Config) error { resolver := idresolver.DefaultResolver(cfg.PlcUrl) bobbin := newBobbinClient(cfg.BobbinApiUrl) - ingester, err := NewIngester(database, bobbin, cfg.JetstreamEndpoint, cfg.Hostname, log.SubLogger(logger, "ingest")) + ingester, err := NewIngester(database, bobbin, resolver, cfg.JetstreamEndpoint, cfg.Hostname, log.SubLogger(logger, "ingest")) if err != nil { return fmt.Errorf("creating ingester: %w", err) } diff --git a/deliberi/ingest.go b/deliberi/ingest.go index 7b662f4c..db87f815 100644 --- a/deliberi/ingest.go +++ b/deliberi/ingest.go @@ -5,12 +5,17 @@ import ( "encoding/json" "fmt" "log/slog" + "net/http" + "time" + comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" + indigoxrpc "github.com/bluesky-social/indigo/xrpc" jmodels "github.com/bluesky-social/jetstream/pkg/models" "tangled.org/core/api/tangled" deldb "tangled.org/core/deliberi/db" models "tangled.org/core/deliberi/models" + "tangled.org/core/idresolver" js "tangled.org/core/jetstream" ) @@ -18,6 +23,7 @@ type Ingester struct { db *deldb.DB recipients recipientResolver jc *js.JetstreamClient + idResolver *idresolver.Resolver logger *slog.Logger } @@ -30,7 +36,7 @@ var ingestCollections = []string{ tangled.GraphFollowNSID, } -func NewIngester(database *deldb.DB, recipients recipientResolver, endpoint, ident string, logger *slog.Logger) (*Ingester, error) { +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) @@ -39,6 +45,7 @@ func NewIngester(database *deldb.DB, recipients recipientResolver, endpoint, ide db: database, recipients: recipients, jc: jc, + idResolver: idRes, logger: logger, }, nil } @@ -84,7 +91,7 @@ func (i *Ingester) process(ctx context.Context, e *jmodels.Event) error { i.logger.Warn("decoding issue record", "err", err, "uri", entityAt) return nil } - if err := deldb.PutEntityTitle(i.db, entityAt, rec.Title); err != nil { + if err := deldb.PutEntityTitle(i.db, entityAt, rec.Title, rec.Repo); err != nil { i.logger.Warn("caching entity title", "err", err, "uri", entityAt) } i.notifyEntity(ctx, actorDid, entityAt, entityAt, rec.Repo, models.NotificationTypeIssueCreated, rec.Title, rec.Mentions, "sh.tangled.repo.issue") @@ -99,7 +106,7 @@ func (i *Ingester) process(ctx context.Context, e *jmodels.Event) error { if rec.Target != nil { repoDid = rec.Target.Repo } - if err := deldb.PutEntityTitle(i.db, entityAt, rec.Title); err != nil { + if err := deldb.PutEntityTitle(i.db, entityAt, rec.Title, repoDid); err != nil { i.logger.Warn("caching entity title", "err", err, "uri", entityAt) } i.notifyEntity(ctx, actorDid, entityAt, entityAt, repoDid, models.NotificationTypePullCreated, rec.Title, rec.Mentions, "sh.tangled.repo.pull") @@ -114,8 +121,9 @@ func (i *Ingester) process(ctx context.Context, e *jmodels.Event) error { return nil } subjectUri := rec.Subject.Uri + collection := syntax.ATURI(subjectUri).Collection().String() var t models.NotificationType - switch syntax.ATURI(subjectUri).Collection().String() { + switch collection { case tangled.RepoIssueNSID: t = models.NotificationTypeIssueCommented case tangled.RepoPullNSID: @@ -123,10 +131,9 @@ func (i *Ingester) process(ctx context.Context, e *jmodels.Event) error { default: return nil } - // comment carries no repo did and no mentions field; leave both empty. - title := deldb.GetEntityTitle(i.db, subjectUri) - collection := syntax.ATURI(subjectUri).Collection().String() - i.notifyEntity(ctx, actorDid, entityAt, subjectUri, "", t, title, nil, collection) + // comments carry no mentions field. + title, repoDid := i.hydrateEntity(ctx, subjectUri, collection) + i.notifyEntity(ctx, actorDid, entityAt, subjectUri, repoDid, t, title, nil, collection) case tangled.FeedStarNSID: var rec tangled.FeedStar @@ -164,9 +171,18 @@ func (i *Ingester) process(ctx context.Context, e *jmodels.Event) error { func (i *Ingester) notifyEntity(ctx context.Context, actorDid, sourceAt, entityAt, repoDid string, t models.NotificationType, title string, mentions []string, collection string) { seen := make(map[string]struct{}) - subscribers, err := i.recipients.ListRecipients(ctx, entityAt, collection) - if err != nil { - i.logger.Warn("listing recipients", "err", err, "entity", entityAt) + // bobbin matches subjects exactly, so ask at both levels. + var subscribers []string + for _, subject := range []string{entityAt, repoDid} { + if subject == "" { + continue + } + dids, err := i.recipients.ListRecipients(ctx, subject, collection) + if err != nil { + i.logger.Warn("listing recipients", "err", err, "subject", subject) + continue + } + subscribers = append(subscribers, dids...) } for _, dids := range [][]string{subscribers, mentions} { @@ -211,6 +227,73 @@ func (i *Ingester) hydrateRepoOwner(ctx context.Context, repoDid string) string return owner } +// hydrateEntity reads the entity cache, falling back to the parent's pds on a +// miss so comments on entities the ingester never saw still resolve. +func (i *Ingester) hydrateEntity(ctx context.Context, uri, collection string) (title, repoDid string) { + title = deldb.GetEntityTitle(i.db, uri) + repoDid = deldb.GetEntityRepo(i.db, uri) + if repoDid != "" || i.idResolver == nil { + return title, repoDid + } + + fetchedTitle, fetchedRepoDid, err := i.fetchEntity(ctx, uri, collection) + if err != nil { + i.logger.Warn("hydrating parent entity", "err", err, "uri", uri) + return title, repoDid + } + if title == "" { + title = fetchedTitle + } + repoDid = fetchedRepoDid + if err := deldb.PutEntityTitle(i.db, uri, title, repoDid); err != nil { + i.logger.Warn("caching hydrated entity", "err", err, "uri", uri) + } + return title, repoDid +} + +func (i *Ingester) fetchEntity(ctx context.Context, uri, collection string) (string, string, error) { + at := syntax.ATURI(uri) + ident, err := i.idResolver.ResolveIdent(ctx, at.Authority().String()) + if err != nil { + return "", "", fmt.Errorf("resolving %s: %w", at.Authority(), err) + } + + xc := &indigoxrpc.Client{ + Host: ident.PDSEndpoint(), + Client: &http.Client{Timeout: 10 * time.Second}, + } + out, err := comatproto.RepoGetRecord(ctx, xc, "", collection, ident.DID.String(), at.RecordKey().String()) + if err != nil { + return "", "", fmt.Errorf("getting record: %w", err) + } + if out == nil || out.Value == nil { + return "", "", fmt.Errorf("record has no value") + } + raw, err := out.Value.MarshalJSON() + if err != nil { + return "", "", fmt.Errorf("re-encoding record: %w", err) + } + + switch collection { + case tangled.RepoIssueNSID: + var rec tangled.RepoIssue + if err := json.Unmarshal(raw, &rec); err != nil { + return "", "", fmt.Errorf("decoding issue: %w", err) + } + return rec.Title, rec.Repo, nil + case tangled.RepoPullNSID: + var rec tangled.RepoPull + if err := json.Unmarshal(raw, &rec); err != nil { + return "", "", fmt.Errorf("decoding pull: %w", err) + } + if rec.Target == nil { + return rec.Title, "", nil + } + return rec.Title, rec.Target.Repo, nil + } + return "", "", fmt.Errorf("unsupported collection %s", collection) +} + func (i *Ingester) notifyOne(ctx context.Context, recipientDid, actorDid, sourceAt, entityAt, repoDid string, t models.NotificationType, title string) { i.deliver(recipientDid, actorDid, sourceAt, entityAt, repoDid, t, title) } diff --git a/deliberi/ingest_test.go b/deliberi/ingest_test.go index abc94df2..165d0bff 100644 --- a/deliberi/ingest_test.go +++ b/deliberi/ingest_test.go @@ -8,6 +8,7 @@ import ( "path/filepath" "testing" + comatprototypes "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" jmodels "github.com/bluesky-social/jetstream/pkg/models" "tangled.org/core/api/tangled" @@ -17,7 +18,9 @@ import ( type fakeResolver struct { dids []string - err error + // bySubject, when set, answers per subject instead of unconditionally. + bySubject map[string][]string + err error // repo owner lookups, keyed by repo did. owners map[string]string @@ -28,6 +31,9 @@ type fakeResolver struct { } func (f fakeResolver) ListRecipients(ctx context.Context, uri string, collection string) ([]string, error) { + if f.bySubject != nil { + return f.bySubject[uri], f.err + } return f.dids, f.err } @@ -261,6 +267,128 @@ func TestRepoHandlerSeedsOwner(t *testing.T) { } } +func TestCommentNotifiesRepoSubscriber(t *testing.T) { + const ( + repoDid = "did:plc:therepo" + repoName = "my-repo" + issueUri = "at://did:plc:bob/sh.tangled.repo.issue/issue1" + ) + + // subscribed to the repo, not the issue, so a row proves the repo lookup ran. + i := newTestIngester(t, fakeResolver{bySubject: map[string][]string{ + repoDid: {"did:sub"}, + }}) + + if err := deldb.PutRepoName(i.db, repoDid, "did:plc:bob", repoName); err != nil { + t.Fatalf("PutRepoName: %v", err) + } + if err := deldb.PutEntityTitle(i.db, issueUri, "the issue", repoDid); err != nil { + t.Fatalf("PutEntityTitle: %v", err) + } + + raw, err := json.Marshal(tangled.FeedComment{ + CreatedAt: "2026-01-01T00:00:00Z", + Subject: &comatprototypes.RepoStrongRef{Uri: issueUri, Cid: "bafyfake"}, + }) + if err != nil { + t.Fatalf("marshal comment: %v", err) + } + ev := &jmodels.Event{ + Did: "did:plc:alice", + Kind: jmodels.EventKindCommit, + Commit: &jmodels.Commit{ + Operation: jmodels.CommitOperationCreate, + Collection: tangled.FeedCommentNSID, + RKey: "comment1", + Record: raw, + }, + } + if err := i.process(context.Background(), ev); err != nil { + t.Fatalf("process: %v", err) + } + + if got := countFor(t, i, "did:sub"); got != 1 { + t.Fatalf("repo subscriber rows = %d, want 1", got) + } + + var gotRepoDid, gotTitle string + err = i.db.QueryRow( + `select repo_did, entity_title from notifications where recipient_did = ?`, + "did:sub", + ).Scan(&gotRepoDid, &gotTitle) + if err != nil { + t.Fatalf("scan notification: %v", err) + } + if gotRepoDid != repoDid { + t.Errorf("repo_did = %q, want %q", gotRepoDid, repoDid) + } + if gotTitle != "the issue" { + t.Errorf("entity_title = %q, want %q", gotTitle, "the issue") + } +} + +func TestCommentNotifiesEntitySubscriber(t *testing.T) { + const ( + repoDid = "did:plc:therepo" + issueUri = "at://did:plc:bob/sh.tangled.repo.issue/issue1" + ) + + i := newTestIngester(t, fakeResolver{bySubject: map[string][]string{ + issueUri: {"did:entitysub"}, + }}) + if err := deldb.PutEntityTitle(i.db, issueUri, "the issue", repoDid); err != nil { + t.Fatalf("PutEntityTitle: %v", err) + } + + raw, err := json.Marshal(tangled.FeedComment{ + CreatedAt: "2026-01-01T00:00:00Z", + Subject: &comatprototypes.RepoStrongRef{Uri: issueUri, Cid: "bafyfake"}, + }) + if err != nil { + t.Fatalf("marshal comment: %v", err) + } + ev := &jmodels.Event{ + Did: "did:plc:alice", + Kind: jmodels.EventKindCommit, + Commit: &jmodels.Commit{ + Operation: jmodels.CommitOperationCreate, + Collection: tangled.FeedCommentNSID, + RKey: "comment1", + Record: raw, + }, + } + if err := i.process(context.Background(), ev); err != nil { + t.Fatalf("process: %v", err) + } + + if got := countFor(t, i, "did:entitysub"); got != 1 { + t.Fatalf("entity subscriber rows = %d, want 1", got) + } +} + +func TestPutEntityTitleKeepsRepoDid(t *testing.T) { + database, err := deldb.Make(context.Background(), filepath.Join(t.TempDir(), "x.db")) + if err != nil { + t.Fatalf("make db: %v", err) + } + t.Cleanup(func() { database.Close() }) + + const uri = "at://did:plc:bob/sh.tangled.repo.issue/issue1" + if err := deldb.PutEntityTitle(database, uri, "t1", "did:plc:therepo"); err != nil { + t.Fatalf("first put: %v", err) + } + // a later write that does not know the repo must not erase it. + if err := deldb.PutEntityTitle(database, uri, "t2", ""); err != nil { + t.Fatalf("second put: %v", err) + } + if got := deldb.GetEntityRepo(database, uri); got != "did:plc:therepo" { + t.Errorf("repo did = %q, want it preserved", got) + } + if got := deldb.GetEntityTitle(database, uri); got != "t2" { + t.Errorf("title = %q, want %q", got, "t2") + } +} + func TestDisabledPrefSuppressesRow(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:sub"}}) prefs := models.DefaultNotificationPreferences(syntax.DID("did:sub"))