package deliberi import ( "context" "io" "log/slog" "path/filepath" "testing" "github.com/bluesky-social/indigo/atproto/syntax" jmodels "github.com/bluesky-social/jetstream/pkg/models" deldb "tangled.org/core/deliberi/db" models "tangled.org/core/deliberi/models" ) type fakeResolver struct { dids []string // bySubject, when set, answers per subject instead of unconditionally bySubject map[string][]string err error owners map[string]string repoNames map[string]string ownerErr error // proves cache hits skip bobbin ownerCalls *int } 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 } func (f fakeResolver) RepoOwner(ctx context.Context, repoDid string) (string, string, error) { if f.ownerCalls != nil { *f.ownerCalls++ } if f.ownerErr != nil { return "", "", f.ownerErr } return f.owners[repoDid], f.repoNames[repoDid], nil } func commitEvent(did string, op, collection, rkey, record string) *jmodels.Event { return &jmodels.Event{ Did: did, Kind: jmodels.EventKindCommit, Commit: &jmodels.Commit{ Operation: op, Collection: collection, RKey: rkey, Record: []byte(record), }, } } func starEvent(actorDid, repoDid, rkey string) *jmodels.Event { return commitEvent(actorDid, jmodels.CommitOperationCreate, starNSID, rkey, `{"$type":"org.tangled.feed.star","createdAt":"2026-01-01T00:00:00Z","subject":{"$type":"org.tangled.feed.star#repo","identity":"`+repoDid+`"}}`) } func newTestIngester(t *testing.T, r recipientResolver) *Ingester { t.Helper() 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() }) return &Ingester{ db: database, recipients: r, logger: slog.New(slog.NewTextHandler(io.Discard, nil)), } } func mustProcess(t *testing.T, i *Ingester, e *jmodels.Event) { t.Helper() if err := i.process(context.Background(), e); err != nil { t.Fatalf("process: %v", err) } } func countFor(t *testing.T, i *Ingester, did string) int64 { t.Helper() n, err := deldb.CountNotifications(i.db, did) if err != nil { t.Fatalf("count: %v", err) } return n } func onlyFor(t *testing.T, i *Ingester, did string) *models.Notification { t.Helper() ns, err := deldb.GetNotifications(i.db, did, 0) if err != nil { t.Fatalf("get notifications: %v", err) } if len(ns) != 1 { t.Fatalf("%s has %d notifications, want 1", did, len(ns)) } return ns[0] } func issueCreated() models.Notification { return models.Notification{ AtUri: "at://src", Type: models.NotificationTypeIssueCreated, ActorDid: "did:actor", RepoDid: "did:repo", EntityAt: "at://entity", EntityKind: models.EntityKindIssue, } } func TestNotifyEntitySubscriberGetsRow(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:sub"}}) i.notifyEntity(context.Background(), issueCreated(), nil) if got := countFor(t, i, "did:sub"); got != 1 { t.Fatalf("subscriber rows = %d, want 1", got) } } func TestActorNeverNotified(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:actor"}}) i.notifyEntity(context.Background(), issueCreated(), []string{"did:actor"}) if got := countFor(t, i, "did:actor"); got != 0 { t.Fatalf("actor rows = %d, want 0", got) } } func TestMentionDeliveredOnResolverError(t *testing.T) { i := newTestIngester(t, fakeResolver{err: io.ErrUnexpectedEOF}) i.notifyEntity(context.Background(), issueCreated(), []string{"did:mention"}) if got := countFor(t, i, "did:mention"); got != 1 { t.Fatalf("mention rows = %d, want 1", got) } } func TestMentionTakesPriorityOverSubscription(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:both"}}) i.notifyEntity(context.Background(), issueCreated(), []string{"did:both"}) n := onlyFor(t, i, "did:both") if n.Type != models.NotificationTypeUserMentioned { t.Errorf("type = %q, want %q", n.Type, models.NotificationTypeUserMentioned) } if n.EntityKind != models.EntityKindIssue { t.Errorf("kind = %q, want the mention to keep the entity's kind", n.EntityKind) } } func TestCreateNotificationDedupe(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:sub"}}) i.notifyEntity(context.Background(), issueCreated(), nil) i.notifyEntity(context.Background(), issueCreated(), nil) if got := countFor(t, i, "did:sub"); got != 1 { t.Fatalf("deduped rows = %d, want 1", got) } } func TestDisabledPrefSuppressesRow(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:sub"}}) prefs := models.DefaultNotificationPreferences(syntax.DID("did:sub")) prefs.IssueCreated = false if err := deldb.UpsertNotificationPreferences(i.db, prefs); err != nil { t.Fatalf("upsert prefs: %v", err) } i.notifyEntity(context.Background(), issueCreated(), nil) if got := countFor(t, i, "did:sub"); got != 0 { t.Fatalf("disabled-pref rows = %d, want 0", got) } } // repoDid comes from invites_test.go const ( ownerDid = "did:plc:bob" editorDid = "did:plc:alice" ticketUri = "at://" + repoDid + "/org.tangled.track.ticket/1" ) func ticketJSON(state, extra string) string { return `{"$type":"org.tangled.track.ticket","createdAt":"2026-01-01T00:00:00Z","state":"` + state + `",` + `"title":{"$type":"org.tangled.markup.markdown#inline","text":"the ticket"},` + extra + `"x-tngl-editor":"` + editorDid + `"}` } func TestTicketCreateNotifiesAsEditor(t *testing.T) { i := newTestIngester(t, fakeResolver{bySubject: map[string][]string{repoDid: {"did:sub", editorDid}}}) mustProcess(t, i, commitEvent(repoDid, jmodels.CommitOperationCreate, ticketNSID, "1", ticketJSON("open", ""))) n := onlyFor(t, i, "did:sub") if n.Type != models.NotificationTypeIssueCreated || n.EntityKind != models.EntityKindIssue { t.Errorf("type/kind = %q/%q, want issue_created/issue", n.Type, n.EntityKind) } if n.ActorDid != editorDid { t.Errorf("actor = %q, want the editor %q, not the repo account", n.ActorDid, editorDid) } if n.RepoDid != repoDid || n.EntityAt != ticketUri || n.EntityTitle != "the ticket" { t.Errorf("repo/entity/title = %q/%q/%q", n.RepoDid, n.EntityAt, n.EntityTitle) } if got := countFor(t, i, editorDid); got != 0 { t.Errorf("editor rows = %d, want 0 (author must not self-notify)", got) } } func TestTicketWithTargetBranchIsPull(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:sub"}}) mustProcess(t, i, commitEvent(repoDid, jmodels.CommitOperationCreate, ticketNSID, "1", ticketJSON("open", `"targetBranch":"main",`))) n := onlyFor(t, i, "did:sub") if n.Type != models.NotificationTypePullCreated || n.EntityKind != models.EntityKindPull { t.Errorf("type/kind = %q/%q, want pull_created/pull", n.Type, n.EntityKind) } } func TestTicketStateUpdateDoesNotNotify(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:sub"}}) for _, state := range []string{"closed", "completed", "open"} { mustProcess(t, i, commitEvent(repoDid, jmodels.CommitOperationUpdate, ticketNSID, "1", ticketJSON(state, `"targetBranch":"main",`))) } if got := countFor(t, i, "did:sub"); got != 0 { t.Fatalf("subscriber rows = %d, want 0 (state changes must not notify)", got) } if got := deldb.GetEntity(i.db, ticketUri); got.Kind != models.EntityKindPull || got.RepoDid != repoDid { t.Errorf("cached ticket = %+v, want the update to still refresh the cache", got) } } func commentJSON(subjectUri string) string { return `{"$type":"org.tangled.feed.comment","createdAt":"2026-01-01T00:00:00Z",` + `"body":{"$type":"org.tangled.markup.markdown","text":"looks good"},` + `"subject":{"uri":"` + subjectUri + `","cid":"bafyfake"},"x-tngl-editor":"` + editorDid + `"}` } func TestCommentNotifiesRepoSubscriber(t *testing.T) { // subscribed to the repo, not the ticket, so a row proves the repo lookup ran i := newTestIngester(t, fakeResolver{bySubject: map[string][]string{repoDid: {"did:sub"}}}) if err := deldb.PutEntity(i.db, ticketUri, deldb.Entity{Title: "the pull", RepoDid: repoDid, Kind: models.EntityKindPull}); err != nil { t.Fatalf("PutEntity: %v", err) } mustProcess(t, i, commitEvent(repoDid, jmodels.CommitOperationCreate, commentNSID, "c1", commentJSON(ticketUri))) n := onlyFor(t, i, "did:sub") if n.Type != models.NotificationTypePullCommented || n.EntityKind != models.EntityKindPull { t.Errorf("type/kind = %q/%q, want pull_commented/pull", n.Type, n.EntityKind) } if n.ActorDid != editorDid || n.EntityAt != ticketUri || n.EntityTitle != "the pull" { t.Errorf("actor/entity/title = %q/%q/%q", n.ActorDid, n.EntityAt, n.EntityTitle) } } func TestCommentNotifiesEntitySubscriber(t *testing.T) { i := newTestIngester(t, fakeResolver{bySubject: map[string][]string{ticketUri: {"did:entitysub"}}}) if err := deldb.PutEntity(i.db, ticketUri, deldb.Entity{Title: "the issue", RepoDid: repoDid, Kind: models.EntityKindIssue}); err != nil { t.Fatalf("PutEntity: %v", err) } mustProcess(t, i, commitEvent(repoDid, jmodels.CommitOperationCreate, commentNSID, "c1", commentJSON(ticketUri))) if n := onlyFor(t, i, "did:entitysub"); n.Type != models.NotificationTypeIssueCommented { t.Errorf("type = %q, want issue_commented", n.Type) } } func TestCommentOnAnotherReposTicketIgnored(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:sub"}}) const foreign = "at://did:plc:elsewhere/org.tangled.track.ticket/1" if err := deldb.PutEntity(i.db, foreign, deldb.Entity{Title: "t", RepoDid: "did:plc:elsewhere", Kind: models.EntityKindIssue}); err != nil { t.Fatalf("PutEntity: %v", err) } mustProcess(t, i, commitEvent(repoDid, jmodels.CommitOperationCreate, commentNSID, "c1", commentJSON(foreign))) if got := countFor(t, i, "did:sub"); got != 0 { t.Fatalf("rows = %d, want 0 (a repo can't comment into another repo's ticket)", got) } } func TestCommentEditDoesNotNotify(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:sub"}}) if err := deldb.PutEntity(i.db, ticketUri, deldb.Entity{Title: "t", RepoDid: repoDid, Kind: models.EntityKindIssue}); err != nil { t.Fatalf("PutEntity: %v", err) } mustProcess(t, i, commitEvent(repoDid, jmodels.CommitOperationUpdate, commentNSID, "c1", commentJSON(ticketUri))) if got := countFor(t, i, "did:sub"); got != 0 { t.Fatalf("rows = %d, want 0", got) } } func TestManifestSeedsOwnerAndSlug(t *testing.T) { calls := 0 i := newTestIngester(t, fakeResolver{ownerCalls: &calls}) mustProcess(t, i, commitEvent(repoDid, jmodels.CommitOperationCreate, repoManifestNSID, "self", `{"$type":"org.tangled.repo.manifest","createdAt":"2026-01-01T00:00:00Z","declaration":{"owner":"`+ownerDid+`","slug":"my-repo"}}`)) mustProcess(t, i, starEvent(editorDid, repoDid, "star1")) if got := deldb.GetRepoName(i.db, repoDid); got != "my-repo" { t.Errorf("cached name = %q, want my-repo", got) } if got := countFor(t, i, ownerDid); got != 1 { t.Fatalf("owner rows = %d, want 1", got) } if calls != 0 { t.Errorf("RepoOwner calls = %d, want 0 (manifest cache hit must not hit bobbin)", calls) } } func TestStarResolvesOwnerFromBobbinOnCacheMiss(t *testing.T) { i := newTestIngester(t, fakeResolver{ owners: map[string]string{repoDid: ownerDid}, repoNames: map[string]string{repoDid: "my-repo"}, }) mustProcess(t, i, starEvent(editorDid, repoDid, "star1")) n := onlyFor(t, i, ownerDid) if n.Type != models.NotificationTypeRepoStarred || n.ActorDid != editorDid || n.RepoDid != repoDid { t.Errorf("notification = %+v", n) } if got := countFor(t, i, repoDid); got != 0 { t.Errorf("repo DID rows = %d, want 0 (stars go to the owner)", got) } if got := deldb.GetRepoOwner(i.db, repoDid); got != ownerDid { t.Errorf("cached owner = %q, want %q", got, ownerDid) } } func TestStarKeepsCachedNameWhenSlugMissing(t *testing.T) { i := newTestIngester(t, fakeResolver{owners: map[string]string{repoDid: ownerDid}}) if err := deldb.PutRepoName(i.db, repoDid, "", "my-repo"); err != nil { t.Fatalf("PutRepoName: %v", err) } mustProcess(t, i, starEvent(editorDid, repoDid, "star1")) if got := deldb.GetRepoName(i.db, repoDid); got != "my-repo" { t.Errorf("cached name = %q, want my-repo", got) } } func TestStarSkipsWhenOwnerUnresolvable(t *testing.T) { i := newTestIngester(t, fakeResolver{ownerErr: io.ErrUnexpectedEOF}) mustProcess(t, i, starEvent(editorDid, repoDid, "star1")) if got := countFor(t, i, ownerDid); got != 0 { t.Fatalf("owner rows = %d, want 0", got) } } func TestFollowNotifiesRkeySubject(t *testing.T) { i := newTestIngester(t, fakeResolver{}) mustProcess(t, i, commitEvent(editorDid, jmodels.CommitOperationCreate, followNSID, ownerDid, `{"$type":"org.tangled.graph.follow","createdAt":"2026-01-01T00:00:00Z"}`)) if n := onlyFor(t, i, ownerDid); n.Type != models.NotificationTypeFollowed || n.ActorDid != editorDid { t.Errorf("notification = %+v", n) } } func TestPutEntityKeepsKnownRepoAndKind(t *testing.T) { i := newTestIngester(t, fakeResolver{}) if err := deldb.PutEntity(i.db, ticketUri, deldb.Entity{Title: "t1", RepoDid: repoDid, Kind: models.EntityKindPull}); err != nil { t.Fatalf("first put: %v", err) } if err := deldb.PutEntity(i.db, ticketUri, deldb.Entity{Title: "t2"}); err != nil { t.Fatalf("second put: %v", err) } want := deldb.Entity{Title: "t2", RepoDid: repoDid, Kind: models.EntityKindPull} if got := deldb.GetEntity(i.db, ticketUri); got != want { t.Errorf("entity = %+v, want %+v", got, want) } } func TestEditorlessCommentAttributedToRepo(t *testing.T) { i := newTestIngester(t, fakeResolver{dids: []string{"did:sub"}}) if err := deldb.PutEntity(i.db, ticketUri, deldb.Entity{Title: "t", RepoDid: repoDid, Kind: models.EntityKindIssue}); err != nil { t.Fatalf("PutEntity: %v", err) } mustProcess(t, i, commitEvent(repoDid, jmodels.CommitOperationCreate, commentNSID, "c1", `{"$type":"org.tangled.feed.comment","createdAt":"2026-01-01T00:00:00Z",`+ `"body":{"$type":"org.tangled.markup.markdown","text":"hi"},"subject":{"uri":"`+ticketUri+`","cid":"bafyfake"}}`)) if n := onlyFor(t, i, "did:sub"); n.ActorDid != repoDid { t.Errorf("actor = %q, want the repo did %q", n.ActorDid, repoDid) } }