diff --git a/docs/specs/product/compass.md b/docs/specs/product/compass.md index 3ad0a23a5..a5ba4be6c 100644 --- a/docs/specs/product/compass.md +++ b/docs/specs/product/compass.md @@ -540,9 +540,11 @@ carries a caller identity. The server SHALL scope every listing, read, and search to the channels, groups, and accounts the caller may see, enforced in the store (SQL), not at the RPC -edge. A message read for a channel the caller is not a member of SHALL return -nothing rather than the channel's contents — a non-member cannot read a private -channel's history by naming its id. +edge. An attached channel is visible to its members and to the anchor agent's +owner set. That visibility does not grant history access: message reads, search, +topic reads, and writes require channel participation. TREE participation is the +anchor agent, its owner, and the agents in the anchor's subtree. A caller who +does not participate SHALL receive no channel history by naming its id. #### Scenario: A non-member lists a private channel's messages diff --git a/go/internal/comms/channel_attach_pgtest_test.go b/go/internal/comms/channel_attach_pgtest_test.go new file mode 100644 index 000000000..f682c3d80 --- /dev/null +++ b/go/internal/comms/channel_attach_pgtest_test.go @@ -0,0 +1,49 @@ +//go:build pgtest + +package comms + +import ( + "testing" + + compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + "github.com/RigelBuild/compass/go/internal/store" +) + +func TestPublishChannelChangedCarriesTreeAttachment(t *testing.T) { + h := newStreamHarness(t) + owner := mustUser(t, h.store, "tree-event-owner") + anchor := mustAgent(t, h.store, owner.ID, "tree-event-anchor") + channel, err := h.store.CreateChannel(t.Context(), owner.ID, store.NewChannel{ + Name: "tree-event", Kind: store.ChannelKindChannel, + ParentAgentID: anchor.ID, MembershipMode: store.ChannelMembershipModeTree, + }) + if err != nil { + t.Fatalf("CreateChannel: %v", err) + } + + sub, err := h.bus.Subscribe(0, 0) + if err != nil { + t.Fatalf("subscribe event bus: %v", err) + } + defer sub.Cancel() + + h.svc.publishChannelChanged(channel, nil) + select { + case event := <-sub.Live: + changed := event.Payload.GetChannelChanged() + if changed == nil { + t.Fatalf("published payload = %T, want ChannelChanged", event.Payload.GetPayload()) + } + got := changed.GetChannel() + if got.GetParentAgentId() != string(anchor.ID) || got.GetMembershipMode() != compassv1.ChannelMembershipMode_CHANNEL_MEMBERSHIP_MODE_TREE { + t.Fatalf("ChannelChanged channel parent=%q mode=%v, want parent=%q TREE", got.GetParentAgentId(), got.GetMembershipMode(), anchor.ID) + } + default: + t.Fatal("publishChannelChanged emitted no event") + } + select { + case event := <-sub.Live: + t.Fatalf("publishChannelChanged emitted duplicate event %T", event.Payload.GetPayload()) + default: + } +} diff --git a/go/internal/comms/channel_attach_wire_test.go b/go/internal/comms/channel_attach_wire_test.go index 666b1a7d2..306bd6c54 100644 --- a/go/internal/comms/channel_attach_wire_test.go +++ b/go/internal/comms/channel_attach_wire_test.go @@ -7,6 +7,7 @@ import ( "google.golang.org/protobuf/reflect/protoreflect" compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + "github.com/RigelBuild/compass/go/internal/store" ) // The channel-attach fields are additive: their numbers are the wire contract @@ -69,3 +70,13 @@ func TestChannelAttachRoundTrip(t *testing.T) { t.Fatalf("round trip: got %v, want %v", &out, in) } } + +func TestChannelToWireCarriesTreeAttachment(t *testing.T) { + got := channelToWire(store.Channel{ + ParentAgentID: "agent-anchor", + MembershipMode: store.ChannelMembershipModeTree, + }) + if got.GetParentAgentId() != "agent-anchor" || got.GetMembershipMode() != compassv1.ChannelMembershipMode_CHANNEL_MEMBERSHIP_MODE_TREE { + t.Fatalf("channelToWire attachment = parent %q mode %v, want agent-anchor and TREE", got.GetParentAgentId(), got.GetMembershipMode()) + } +} diff --git a/go/internal/comms/mapping.go b/go/internal/comms/mapping.go index 275bf7095..3d9fee565 100644 --- a/go/internal/comms/mapping.go +++ b/go/internal/comms/mapping.go @@ -99,9 +99,17 @@ func channelToWire(c store.Channel) *compassv1.Channel { PostPolicy: channelPostPolicyToWire(c.Policy.PostPolicy), OwnerAccountId: string(c.Policy.OwnerAccountID), MandatorySubscription: c.Policy.MandatorySubscription, + ParentAgentId: string(c.ParentAgentID), + MembershipMode: channelMembershipModeToWire(c.MembershipMode), } } +func channelMembershipModeToWire(m store.ChannelMembershipMode) compassv1.ChannelMembershipMode { + if m == store.ChannelMembershipModeTree { + return compassv1.ChannelMembershipMode_CHANNEL_MEMBERSHIP_MODE_TREE + } + return compassv1.ChannelMembershipMode_CHANNEL_MEMBERSHIP_MODE_EXPLICIT +} func channelKindToWire(k store.ChannelKind) compassv1.ChannelKind { switch k { case store.ChannelKindDM: diff --git a/go/internal/store/channel_by_name_pgtest_test.go b/go/internal/store/channel_by_name_pgtest_test.go index e7075ed68..7cd5c6fc6 100644 --- a/go/internal/store/channel_by_name_pgtest_test.go +++ b/go/internal/store/channel_by_name_pgtest_test.go @@ -109,3 +109,38 @@ func TestChannelByNameForViewerAmbiguousIsInvalidArgument(t *testing.T) { t.Fatalf("ambiguous resolve returned ErrNotFound, want ErrInvalidArgument only") } } + +func TestChannelByNameForViewerNarrowsOwnerSetDuplicatesToParticipants(t *testing.T) { + s := newTestStore(t) + owner := mustUser(t, s, "standup-owner") + anchorA := mustAgent(t, s, owner.ID, "standup-a") + anchorB := mustAgent(t, s, owner.ID, "standup-b") + channelA, err := s.CreateChannel(t.Context(), owner.ID, NewChannel{ + Name: "standup", Kind: ChannelKindChannel, + ParentAgentID: anchorA.ID, MembershipMode: ChannelMembershipModeTree, + }) + if err != nil { + t.Fatalf("CreateChannel(A): %v", err) + } + if _, err := s.CreateChannel(t.Context(), owner.ID, NewChannel{ + Name: "standup", Kind: ChannelKindChannel, + ParentAgentID: anchorB.ID, MembershipMode: ChannelMembershipModeTree, + }); err != nil { + t.Fatalf("CreateChannel(B): %v", err) + } + + got, err := s.ChannelByNameForViewer(t.Context(), anchorA.ID, "standup") + if err != nil { + t.Fatalf("ChannelByNameForViewer(A1): %v", err) + } + if got.ID != channelA.ID { + t.Fatalf("A1 resolved channel %q, want its own %q", got.ID, channelA.ID) + } + _, err = s.ChannelByNameForViewer(t.Context(), owner.ID, "standup") + sentinelIs(t, err, ErrInvalidArgument, "owner resolves both standup channels") + // A sibling that sees both but participates in neither keeps the visible-set + // ambiguity rather than a not-found for a name its list shows. + sibling := mustAgent(t, s, owner.ID, "standup-c") + _, err = s.ChannelByNameForViewer(t.Context(), sibling.ID, "standup") + sentinelIs(t, err, ErrInvalidArgument, "non-participant sibling resolves both standup channels") +} diff --git a/go/internal/store/channel_participant_pgtest_test.go b/go/internal/store/channel_participant_pgtest_test.go index 5b9175949..161b8c1cc 100644 --- a/go/internal/store/channel_participant_pgtest_test.go +++ b/go/internal/store/channel_participant_pgtest_test.go @@ -4,8 +4,11 @@ package store import ( "context" + "reflect" "testing" "time" + + "github.com/RigelBuild/compass/go/internal/store/db" ) func appendAsParticipant(t *testing.T, s *Store, author AccountID, channel ChannelID, text string) (Message, error) { @@ -95,6 +98,225 @@ func TestChannelParticipantTreeArm(t *testing.T) { } } +func TestChannelVisibilityOwnerSetParity(t *testing.T) { + s := newTestStore(t) + owner := mustUser(t, s, "visibility-owner") + anchor := mustAgent(t, s, owner.ID, "visibility-anchor") + sibling := mustAgent(t, s, owner.ID, "visibility-sibling") + foreign := mustUser(t, s, "visibility-foreign") + channel := mustAttachedChannel(t, s, owner.ID, anchor.ID, "visibility-tree", ChannelMembershipModeTree) + + for _, tc := range []struct { + name string + who AccountID + want bool + }{ + {name: "owner", who: owner.ID, want: true}, + {name: "same-owner sibling", who: sibling.ID, want: true}, + {name: "foreign user", who: foreign.ID}, + } { + t.Run(tc.name, func(t *testing.T) { + channels, err := s.ListChannels(t.Context(), tc.who) + if err != nil { + t.Fatalf("ListChannels: %v", err) + } + listed := false + for _, got := range channels { + if got.ID == channel.ID { + listed = true + } + } + visible, err := s.ChannelVisibleTo(t.Context(), tc.who, channel.ID) + if err != nil { + t.Fatalf("ChannelVisibleTo: %v", err) + } + if listed != tc.want || visible != tc.want || visible != listed { + t.Fatalf("owner-set visibility: ListChannels=%v ChannelVisibleTo=%v, want %v", listed, visible, tc.want) + } + }) + } +} + +func TestChannelTreeParticipantReadsAndWrites(t *testing.T) { + s := newTestStore(t) + f := newTreeFixture(t, s) + + message, err := appendAsParticipant(t, s, f.leaf.ID, f.channel.ID, "subtree searchable history") + if err != nil { + t.Fatalf("AppendMessage by subtree agent: %v", err) + } + + listed, err := s.ListMessages(t.Context(), ListMessagesQuery{Actor: f.leaf.ID, ChannelID: f.channel.ID}) + if err != nil || len(listed) != 1 || listed[0].ID != message.ID { + t.Fatalf("ListMessages for subtree agent = %v, %v; want message %q", listed, err, message.ID) + } + + searched, err := s.SearchMessages(t.Context(), f.leaf.ID, SearchScope{ChannelID: f.channel.ID}, "subtree searchable", Page{}) + if err != nil || len(searched) != 1 || searched[0].ID != message.ID { + t.Fatalf("SearchMessages for subtree agent = %v, %v; want message %q", searched, err, message.ID) + } + ownerSearch, err := s.SearchMessages(t.Context(), f.owner.ID, SearchScope{ChannelID: f.channel.ID}, "subtree searchable", Page{}) + if err != nil || len(ownerSearch) != 1 || ownerSearch[0].ID != message.ID { + t.Fatalf("SearchMessages for anchor owner = %v, %v; want message %q", ownerSearch, err, message.ID) + } + unscoped, err := s.SearchMessages(t.Context(), f.leaf.ID, SearchScope{}, "subtree searchable", Page{}) + if err != nil || len(unscoped) != 1 || unscoped[0].ID != message.ID { + t.Fatalf("unscoped SearchMessages for subtree agent = %v, %v; want message %q", unscoped, err, message.ID) + } + ownerListed, err := s.ListMessages(t.Context(), ListMessagesQuery{Actor: f.owner.ID, ChannelID: f.channel.ID}) + if err != nil || len(ownerListed) != 1 || ownerListed[0].ID != message.ID { + t.Fatalf("ListMessages for anchor owner = %v, %v; want message %q", ownerListed, err, message.ID) + } + + cursorSeq, err := s.q.GetPageCursorSeq(t.Context(), db.GetPageCursorSeqParams{ + AccountID: string(f.leaf.ID), ID: string(message.ID), ChannelID: string(f.channel.ID), + }) + if err != nil || cursorSeq == 0 { + t.Fatalf("GetPageCursorSeq for subtree agent = %d, %v; want message sequence", cursorSeq, err) + } + page, err := s.ListMessages(t.Context(), ListMessagesQuery{Actor: f.leaf.ID, ChannelID: f.channel.ID, Page: Page{BeforeMessageID: message.ID}}) + if err != nil || len(page) != 0 { + t.Fatalf("ListMessages cursor for subtree agent = %v, %v; want empty page", page, err) + } + + if _, err := s.UpdateMessageBlocksAsAuthor(t.Context(), f.leaf.ID, message.ID, []MessageBlock{textBlock("edited by subtree author")}); err != nil { + t.Fatalf("UpdateMessageBlocksAsAuthor by subtree author: %v", err) + } + if _, err := s.UpdateTopic(t.Context(), string(f.leaf.ID), listed[0].TopicID, nil, nil); err != nil { + t.Fatalf("UpdateTopic by subtree agent: %v", err) + } + + if _, err := s.ReparentAgent(t.Context(), f.owner.ID, f.leaf.ID, ""); err != nil { + t.Fatalf("ReparentAgent(leaf out of channel subtree): %v", err) + } + if _, err := s.UpdateMessageBlocksAsAuthor(t.Context(), f.leaf.ID, message.ID, []MessageBlock{textBlock("reparented author edit")}); err == nil { + t.Fatal("UpdateMessageBlocksAsAuthor by reparented author succeeded") + } else { + sentinelIs(t, err, ErrNotFound, "reparented author UpdateMessageBlocksAsAuthor") + } +} + +func TestChannelTreeOwnerSetSiblingReadWriteDenials(t *testing.T) { + s := newTestStore(t) + f := newTreeFixture(t, s) + message, err := appendAsParticipant(t, s, f.leaf.ID, f.channel.ID, "subtree searchable history") + if err != nil { + t.Fatalf("AppendMessage by subtree agent: %v", err) + } + + ownerOnly, err := s.ListMessages(t.Context(), ListMessagesQuery{Actor: f.sibling.ID, ChannelID: f.channel.ID}) + if err != nil || len(ownerOnly) != 0 { + t.Fatalf("owner-set sibling history = %v, %v; want no messages", ownerOnly, err) + } + ownerSearch, err := s.SearchMessages(t.Context(), f.sibling.ID, SearchScope{ChannelID: f.channel.ID}, "subtree searchable", Page{}) + if err != nil || len(ownerSearch) != 0 { + t.Fatalf("owner-set sibling search = %v, %v; want no messages", ownerSearch, err) + } + unscopedSearch, err := s.SearchMessages(t.Context(), f.sibling.ID, SearchScope{}, "subtree searchable", Page{}) + if err != nil || len(unscopedSearch) != 0 { + t.Fatalf("owner-set sibling unscoped search = %v, %v; want no messages", unscopedSearch, err) + } + ownerCursorSeq, err := s.q.GetPageCursorSeq(t.Context(), db.GetPageCursorSeqParams{ + AccountID: string(f.sibling.ID), ID: string(message.ID), ChannelID: string(f.channel.ID), + }) + if !noRows(err) { + t.Fatalf("GetPageCursorSeq for owner-set sibling = %d, %v; want no rows", ownerCursorSeq, err) + } + _, err = s.UpdateTopic(t.Context(), string(f.sibling.ID), message.TopicID, nil, nil) + sentinelIs(t, err, ErrNotFound, "owner-set sibling UpdateTopic") + + outsiderMessageID := newID() + if _, err := s.scopedPool().Exec(t.Context(), ` + INSERT INTO messages (id, topic_id, author_account_id, at_unix_ms, blocks, text_content) + VALUES ($1, $2, $3, 1, '[{"kind":"text","text":"outsider authored"}]'::jsonb, 'outsider authored')`, + outsiderMessageID, message.TopicID, string(f.sibling.ID)); err != nil { + t.Fatalf("insert outsider-authored message: %v", err) + } + if _, err := s.UpdateMessageBlocksAsAuthor(t.Context(), f.sibling.ID, MessageID(outsiderMessageID), []MessageBlock{textBlock("unauthorized edit")}); err == nil { + t.Fatal("owner-set sibling edited a message without channel participation") + } else { + sentinelIs(t, err, ErrNotFound, "owner-set sibling UpdateMessageBlocksAsAuthor") + } + + _, _, err = s.AppendMessage(t.Context(), Message{AuthorAccountID: f.leaf.ID, Blocks: []MessageBlock{pendingAsk("tree-ask", false)}}, string(f.channel.ID), TopicRef{Name: "asks", Create: true}, "") + if err != nil { + t.Fatalf("AppendMessage(ask): %v", err) + } + askFilter, err := askIDContainmentFilter("tree-ask") + if err != nil { + t.Fatalf("askIDContainmentFilter: %v", err) + } + if rows, err := s.q.FindAskMessage(t.Context(), db.FindAskMessageParams{AccountID: string(f.sibling.ID), Column2: askFilter}); err != nil || len(rows) != 0 { + t.Fatalf("FindAskMessage for owner-set sibling = %d rows, %v; want none", len(rows), err) + } + _, _, err = s.AnswerAsk(t.Context(), f.sibling.ID, "tree-ask", []AskAnswer{{QuestionID: "q1", ChosenOptionIDs: []string{"opt-a"}}}) + sentinelIs(t, err, ErrNotFound, "owner-set sibling AnswerAsk") + + if rows, err := s.q.FindAskMessage(t.Context(), db.FindAskMessageParams{AccountID: string(f.leaf.ID), Column2: askFilter}); err != nil || len(rows) != 1 { + t.Fatalf("FindAskMessage for subtree agent = %d rows, %v; want one ask", len(rows), err) + } + if _, _, err := s.AnswerAsk(t.Context(), f.leaf.ID, "tree-ask", []AskAnswer{{QuestionID: "q1", ChosenOptionIDs: []string{"opt-a"}}}); err != nil { + t.Fatalf("AnswerAsk by subtree agent: %v", err) + } +} + +func TestChannelTreeMembersAreAttributedAndSubscriptionsIntersect(t *testing.T) { + s := newTestStore(t) + owner := mustUser(t, s, "materialize-owner") + rootA := mustAgent(t, s, owner.ID, "materialize-a") + childA := mustAgentWithParent(t, s, owner.ID, rootA.ID, "materialize-a-child") + rootB := mustAgent(t, s, owner.ID, "materialize-b") + childB := mustAgentWithParent(t, s, owner.ID, rootB.ID, "materialize-b-child") + otherOwner := mustUser(t, s, "materialize-other") + outside := mustAgent(t, s, otherOwner.ID, "materialize-outside") + channelA := mustAttachedChannel(t, s, owner.ID, rootA.ID, "materialize-tree-a", ChannelMembershipModeTree) + channelB := mustAttachedChannel(t, s, owner.ID, rootB.ID, "materialize-tree-b", ChannelMembershipModeTree) + + if _, err := s.ReparentAgent(t.Context(), owner.ID, childA.ID, rootB.ID); err != nil { + t.Fatalf("ReparentAgent(child A out): %v", err) + } + if _, err := s.scopedPool().Exec(t.Context(), + `INSERT INTO channel_subscriptions (channel_id, account_id, subscribed) VALUES ($1, $2, TRUE)`, + string(channelA.ID), string(childA.ID)); err != nil { + t.Fatalf("insert stale subscription: %v", err) + } + if _, err := s.scopedPool().Exec(t.Context(), + `INSERT INTO channel_subscriptions (channel_id, account_id, subscribed) VALUES ($1, $2, TRUE)`, + string(channelA.ID), string(rootA.ID)); err != nil { + t.Fatalf("insert live subscription: %v", err) + } + if _, err := s.scopedPool().Exec(t.Context(), + `INSERT INTO channel_subscriptions (channel_id, account_id, subscribed) VALUES ($1, $2, TRUE)`, + string(channelA.ID), string(outside.ID)); err != nil { + t.Fatalf("insert outside stale subscription: %v", err) + } + + channels, err := s.ListChannels(t.Context(), owner.ID) + if err != nil { + t.Fatalf("ListChannels: %v", err) + } + got := make(map[ChannelID]map[AccountID]bool) + for _, channel := range channels { + if channel.ID == channelA.ID || channel.ID == channelB.ID { + got[channel.ID] = memberSet(channel) + if channel.ID == channelA.ID && !reflect.DeepEqual(channel.SubscriberAccountIDs, []AccountID{rootA.ID}) { + t.Fatalf("channel A subscribers = %v, want only %s", channel.SubscriberAccountIDs, rootA.ID) + } + if channel.ID == channelB.ID && len(channel.SubscriberAccountIDs) != 0 { + t.Fatalf("channel B subscribers = %v, want none", channel.SubscriberAccountIDs) + } + } + } + wantA := map[AccountID]bool{owner.ID: true, rootA.ID: true} + wantB := map[AccountID]bool{owner.ID: true, rootB.ID: true, childA.ID: true, childB.ID: true} + if !reflect.DeepEqual(got[channelA.ID], wantA) || !reflect.DeepEqual(got[channelB.ID], wantB) { + t.Fatalf("derived participants are misattributed: A=%v want %v; B=%v want %v", got[channelA.ID], wantA, got[channelB.ID], wantB) + } + if n := countChannelMembers(t, s, channelA.ID); n != 0 { + t.Fatalf("TREE channel has %d stored member rows, want 0", n) + } +} + func TestChannelParticipantFollowsReparentAgent(t *testing.T) { s := newTestStore(t) f := newTreeFixture(t, s) diff --git a/go/internal/store/channel_tree_pgtest_test.go b/go/internal/store/channel_tree_pgtest_test.go index c7d059043..9d68a7176 100644 --- a/go/internal/store/channel_tree_pgtest_test.go +++ b/go/internal/store/channel_tree_pgtest_test.go @@ -5,6 +5,7 @@ package store import ( "context" "errors" + "reflect" "strings" "testing" "time" @@ -55,9 +56,54 @@ func TestChannelTreeCreateUnderAgent(t *testing.T) { if parent != string(anchor.ID) || mode != int16(ChannelMembershipModeTree) || members != 0 { t.Fatalf("tree channel state = (%q, %d, %d), want (%q, 1, 0)", parent, mode, members, anchor.ID) } - if tree.MemberAccountIDs != nil { - t.Fatalf("tree returned members = %v, want nil before derived reads", tree.MemberAccountIDs) + if tree.ParentAgentID != anchor.ID || tree.MembershipMode != ChannelMembershipModeTree { + t.Fatalf("tree create returned parent=%q mode=%d, want parent=%q mode=%d", tree.ParentAgentID, tree.MembershipMode, anchor.ID, ChannelMembershipModeTree) } + if got, want := memberSet(tree), map[AccountID]bool{owner.ID: true, anchor.ID: true}; !reflect.DeepEqual(got, want) { + t.Fatalf("tree returned members = %v, want owner and anchor %v", got, want) + } + if len(tree.SubscriberAccountIDs) != 0 { + t.Fatalf("tree returned subscribers = %v, want none", tree.SubscriberAccountIDs) + } +} + +func TestChannelTreeProjectionCarriesAnchorAndMode(t *testing.T) { + s := newTestStore(t) + owner := mustUser(t, s, "projection-owner") + anchor := mustAgent(t, s, owner.ID, "projection-anchor") + created := mustAttachedChannel(t, s, owner.ID, anchor.ID, "projection-tree", ChannelMembershipModeTree) + + assertProjection := func(name string, got Channel) { + t.Helper() + if got.ID != created.ID || got.ParentAgentID != anchor.ID || got.MembershipMode != ChannelMembershipModeTree { + t.Fatalf("%s projection = id %q parent %q mode %d, want id %q parent %q mode %d", name, got.ID, got.ParentAgentID, got.MembershipMode, created.ID, anchor.ID, ChannelMembershipModeTree) + } + } + + listed, err := s.ListChannels(t.Context(), owner.ID) + if err != nil { + t.Fatalf("ListChannels: %v", err) + } + found := false + for _, channel := range listed { + if channel.ID == created.ID { + assertProjection("ListChannels", channel) + found = true + } + } + if !found { + t.Fatalf("ListChannels omitted attached channel %q", created.ID) + } + byName, err := s.ChannelByNameForViewer(t.Context(), owner.ID, "projection-tree") + if err != nil { + t.Fatalf("ChannelByNameForViewer: %v", err) + } + assertProjection("ChannelByNameForViewer", byName) + got, err := s.GetChannel(t.Context(), created.ID) + if err != nil { + t.Fatalf("GetChannel: %v", err) + } + assertProjection("GetChannel", got) } func TestChannelTreeCreateAuthorizationAndInputRefusals(t *testing.T) { diff --git a/go/internal/store/channels.go b/go/internal/store/channels.go index a60306d41..787837987 100644 --- a/go/internal/store/channels.go +++ b/go/internal/store/channels.go @@ -161,8 +161,7 @@ func validateNewChannel(c NewChannel) error { // // A ParentAgentID attaches the channel under an agent the actor's owner set // owns (else ErrNotFound). A TREE channel writes no member rows: its -// participants derive from the anchor's subtree, which need not include the -// actor, and it returns no members until reads derive them. +// participants derive from the anchor's subtree. func (s *Store) CreateChannel(ctx context.Context, actor AccountID, c NewChannel) (Channel, error) { if err := validateNewChannel(c); err != nil { return Channel{}, err @@ -242,22 +241,13 @@ func (s *Store) CreateChannel(ctx context.Context, actor AccountID, c NewChannel if err := s.writeExplicitMembers(ctx, tx, ChannelID(id), c.Policy, members); err != nil { return Channel{}, err } - } else { - // No member rows exist, so report none: a later ListChannels reads the same. - members = nil } + // The post-commit read projects TREE participants instead of stored rows. if err := tx.Commit(ctx); err != nil { return Channel{}, fmt.Errorf("store: commit create channel: %w", err) } - return Channel{ - ID: ChannelID(id), - Name: c.Name, - GroupID: c.GroupID, - Kind: c.Kind, - MemberAccountIDs: members, - Policy: c.Policy, - }, nil + return s.getChannel(ctx, ChannelID(id)) } // writeExplicitMembers stores an EXPLICIT channel's member rows inside the @@ -450,7 +440,7 @@ func (s *Store) ListChannels(ctx context.Context, visibleTo AccountID) ([]Channe } var channels []Channel for _, row := range rows { - channels = append(channels, channelFromRow(row.ID, row.Name, row.GroupID, row.Kind, row.PostPolicy, row.OwnerAccountID, row.MandatorySubscription)) + channels = append(channels, channelFromRow(row.ID, row.Name, row.GroupID, row.Kind, row.PostPolicy, row.OwnerAccountID, row.MandatorySubscription, row.ParentAgentID, row.MembershipMode)) } if err := loadChannelMembers(ctx, s.scopedPool(), channels); err != nil { return nil, err @@ -639,7 +629,9 @@ func groupRefHints(groups []ChannelGroup, handles map[AccountID]string, matches // ErrNotFound, the two indistinguishable so a probe cannot enumerate names it // lacks visibility for (the D9 not-found/forbidden merge); // - exactly one → that channel; -// - two or more visible channels sharing the name → ErrInvalidArgument naming +// - two or more → narrowed to the ones the viewer participates in, when it +// participates in any, so an owner-set sibling's same-named channel does +// not shadow the viewer's own; still two or more → ErrInvalidArgument naming // the collision, so the caller disambiguates rather than the server guessing // (there is no ErrAmbiguous sentinel — invalid_argument is the R1 rule). // @@ -655,18 +647,34 @@ func (s *Store) ChannelByNameForViewer(ctx context.Context, viewer AccountID, na } var channels []Channel for _, row := range rows { - channels = append(channels, channelFromRow(row.ID, row.Name, row.GroupID, row.Kind, row.PostPolicy, row.OwnerAccountID, row.MandatorySubscription)) + channels = append(channels, channelFromRow(row.ID, row.Name, row.GroupID, row.Kind, row.PostPolicy, row.OwnerAccountID, row.MandatorySubscription, row.ParentAgentID, row.MembershipMode)) } if err := loadChannelMembers(ctx, s.scopedPool(), channels); err != nil { return Channel{}, err } + if len(channels) > 1 { + // The ascent probe is the participation authority, not the materialized list. + var participants []Channel + for _, channel := range channels { + member, err := isChannelMember(ctx, s.scopedPool(), viewer, channel.ID) + if err != nil { + return Channel{}, fmt.Errorf("store: check channel-name participant: %w", err) + } + if member { + participants = append(participants, channel) + } + } + if len(participants) > 0 { + channels = participants + } + } switch len(channels) { case 0: return Channel{}, fmt.Errorf("%w: channel %q", ErrNotFound, name) case 1: return channels[0], nil default: - return Channel{}, fmt.Errorf("%w: channel name %q is ambiguous — it names %d visible channels; address it by id", ErrInvalidArgument, name, len(channels)) + return Channel{}, fmt.Errorf("%w: channel name %q is ambiguous — it names %d channels; address it by id", ErrInvalidArgument, name, len(channels)) } } @@ -1104,7 +1112,7 @@ func (s *Store) getChannel(ctx context.Context, id ChannelID) (Channel, error) { } return Channel{}, fmt.Errorf("store: get channel: %w", err) } - channels := []Channel{channelFromRow(row.ID, row.Name, row.GroupID, row.Kind, row.PostPolicy, row.OwnerAccountID, row.MandatorySubscription)} + channels := []Channel{channelFromRow(row.ID, row.Name, row.GroupID, row.Kind, row.PostPolicy, row.OwnerAccountID, row.MandatorySubscription, row.ParentAgentID, row.MembershipMode)} if err := loadChannelMembers(ctx, s.scopedPool(), channels); err != nil { return Channel{}, err } @@ -1118,16 +1126,15 @@ func scanChannels(ctx context.Context, q db.DBTX, rows pgx.Rows) ([]Channel, err var channels []Channel for rows.Next() { var ( - id, name, groupID string - kind int16 - postPolicy int16 - ownerAccountID string - mandatorySubscription bool + id, name, groupID, parentAgentID string + kind, postPolicy, membershipMode int16 + ownerAccountID string + mandatorySubscription bool ) - if err := rows.Scan(&id, &name, &groupID, &kind, &postPolicy, &ownerAccountID, &mandatorySubscription); err != nil { + if err := rows.Scan(&id, &name, &groupID, &kind, &postPolicy, &ownerAccountID, &mandatorySubscription, &parentAgentID, &membershipMode); err != nil { return nil, fmt.Errorf("store: scan channel: %w", err) } - channels = append(channels, channelFromRow(id, name, groupID, kind, postPolicy, ownerAccountID, mandatorySubscription)) + channels = append(channels, channelFromRow(id, name, groupID, kind, postPolicy, ownerAccountID, mandatorySubscription, parentAgentID, membershipMode)) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("store: iterate channels: %w", err) @@ -1138,15 +1145,16 @@ func scanChannels(ctx context.Context, q db.DBTX, rows pgx.Rows) ([]Channel, err return channels, nil } -// channelFromRow builds the base Channel (id, name, group, kind, policy) from -// the shared seven-column channel projection every channel read selects; the -// caller populates the member/subscriber sets with loadChannelMembers. -func channelFromRow(id, name, groupID string, kind, postPolicy int16, ownerAccountID string, mandatorySubscription bool) Channel { +// channelFromRow builds the base Channel from the shared nine-column channel +// projection; the caller populates member/subscriber sets with loadChannelMembers. +func channelFromRow(id, name, groupID string, kind, postPolicy int16, ownerAccountID string, mandatorySubscription bool, parentAgentID string, membershipMode int16) Channel { return Channel{ - ID: ChannelID(id), - Name: name, - GroupID: ChannelGroupID(groupID), - Kind: ChannelKind(kind), + ID: ChannelID(id), + Name: name, + GroupID: ChannelGroupID(groupID), + Kind: ChannelKind(kind), + ParentAgentID: AccountID(parentAgentID), + MembershipMode: ChannelMembershipMode(membershipMode), Policy: ChannelPolicy{ PostPolicy: ChannelPostPolicy(postPolicy), OwnerAccountID: AccountID(ownerAccountID), @@ -1155,10 +1163,10 @@ func channelFromRow(id, name, groupID string, kind, postPolicy int16, ownerAccou } } -// loadChannelMembers populates each channel's member and subscriber sets with -// one follow-up query over the whole id set, so member loading is O(1) -// round-trips rather than one per channel. Runs against the pool or a tx (any -// db.DBTX), mirroring the former scanChannels member follow-up. +// loadChannelMembers uses one batch query to load EXPLICIT member rows and +// derived TREE participants for the whole channel id set. Subscription overrides +// apply only to derived participants, keeping subscribers a subset of members. +// The shared read works with the pool or a tx (any db.DBTX). func loadChannelMembers(ctx context.Context, q db.DBTX, channels []Channel) error { if len(channels) == 0 { return nil diff --git a/go/internal/store/coordination_pgtest_test.go b/go/internal/store/coordination_pgtest_test.go index d70c85a18..499681efa 100644 --- a/go/internal/store/coordination_pgtest_test.go +++ b/go/internal/store/coordination_pgtest_test.go @@ -70,7 +70,9 @@ func coordChannels(t *testing.T, s *Store, owner AccountID) []Channel { t.Helper() ctx := context.Background() rows, err := s.pool.Query(ctx, - `SELECT c.id, c.name, COALESCE(c.group_id,''), c.kind, c.post_policy, COALESCE(c.owner_account_id,''), c.mandatory_subscription + `SELECT c.id, c.name, COALESCE(c.group_id,''), c.kind, c.post_policy, + COALESCE(c.owner_account_id,''), c.mandatory_subscription, + COALESCE(c.parent_agent_id,''), c.membership_mode FROM channels c JOIN channel_groups g ON g.id = c.group_id WHERE g.owner_user_id = $1 AND g.name = '__coordination__' diff --git a/go/internal/store/db/channels.sql.go b/go/internal/store/db/channels.sql.go index c9bba0199..ccc0e8a0e 100644 --- a/go/internal/store/db/channels.sql.go +++ b/go/internal/store/db/channels.sql.go @@ -91,10 +91,34 @@ func (q *Queries) ChannelMemberExists(ctx context.Context, arg ChannelMemberExis } const channelMembersByChannelIDs = `-- name: ChannelMembersByChannelIDs :many -SELECT channel_id, account_id, subscribed -FROM channel_members -WHERE channel_id = ANY($1::text[]) -ORDER BY account_id +WITH RECURSIVE stored AS ( + SELECT cm.channel_id, cm.account_id, cm.subscribed + FROM channel_members cm + JOIN channels c ON c.id = cm.channel_id + WHERE cm.channel_id = ANY($1::text[]) AND c.membership_mode = 0 +), subtree AS ( + SELECT c.id AS channel_id, c.parent_agent_id AS account_id + FROM channels c + WHERE c.id = ANY($1::text[]) AND c.membership_mode = 1 + UNION + SELECT s.channel_id, a.account_id + FROM agent_accounts a + JOIN subtree s ON a.parent_agent_id = s.account_id +), participants AS ( + SELECT channel_id, account_id FROM subtree + UNION + SELECT c.id AS channel_id, aa.owner_user_id AS account_id + FROM channels c + JOIN agent_accounts aa ON aa.account_id = c.parent_agent_id + WHERE c.id = ANY($1::text[]) AND c.membership_mode = 1 +) +SELECT channel_id, account_id, subscribed FROM stored +UNION ALL +SELECT p.channel_id, p.account_id, COALESCE(cs.subscribed, FALSE) AS subscribed +FROM participants p +LEFT JOIN channel_subscriptions cs + ON cs.channel_id = p.channel_id AND cs.account_id = p.account_id +ORDER BY channel_id, account_id ` type ChannelMembersByChannelIDsRow struct { @@ -178,6 +202,11 @@ effective AS ( SELECT id, MIN(min_vis) AS eff_vis FROM ancestry GROUP BY id +), +viewer AS ( + SELECT owner_user_id AS uid FROM agent_accounts WHERE account_id = $1 + UNION ALL + SELECT $1 AS uid ) SELECT EXISTS ( SELECT 1 FROM channels c @@ -191,6 +220,11 @@ SELECT EXISTS ( SELECT 1 FROM effective e WHERE e.id = c.group_id AND e.eff_vis = 1 ) ) + OR ( + c.parent_agent_id IS NOT NULL + AND (SELECT aa.owner_user_id FROM agent_accounts aa WHERE aa.account_id = c.parent_agent_id) + IN (SELECT uid FROM viewer) + ) ) ) ` @@ -220,9 +254,15 @@ effective AS ( SELECT id, MIN(min_vis) AS eff_vis FROM ancestry GROUP BY id +), +viewer AS ( + SELECT owner_user_id AS uid FROM agent_accounts WHERE account_id = $1 + UNION ALL + SELECT $1 AS uid ) SELECT c.id, c.name, COALESCE(c.group_id, '') AS group_id, c.kind, c.post_policy, - COALESCE(c.owner_account_id, '') AS owner_account_id, c.mandatory_subscription + COALESCE(c.owner_account_id, '') AS owner_account_id, c.mandatory_subscription, + COALESCE(c.parent_agent_id, '') AS parent_agent_id, c.membership_mode FROM channels c WHERE c.name = $2 AND ( EXISTS ( @@ -234,6 +274,11 @@ WHERE c.name = $2 AND ( SELECT 1 FROM effective e WHERE e.id = c.group_id AND e.eff_vis = 1 ) ) + OR ( + c.parent_agent_id IS NOT NULL + AND (SELECT aa.owner_user_id FROM agent_accounts aa WHERE aa.account_id = c.parent_agent_id) + IN (SELECT uid FROM viewer) + ) ) ORDER BY c.id ` @@ -251,6 +296,8 @@ type ChannelsByNameForViewerRow struct { PostPolicy int16 OwnerAccountID string MandatorySubscription bool + ParentAgentID string + MembershipMode int16 } func (q *Queries) ChannelsByNameForViewer(ctx context.Context, arg ChannelsByNameForViewerParams) ([]ChannelsByNameForViewerRow, error) { @@ -270,6 +317,8 @@ func (q *Queries) ChannelsByNameForViewer(ctx context.Context, arg ChannelsByNam &i.PostPolicy, &i.OwnerAccountID, &i.MandatorySubscription, + &i.ParentAgentID, + &i.MembershipMode, ); err != nil { return nil, err } @@ -339,7 +388,8 @@ func (q *Queries) GetAgentWorkspaceID(ctx context.Context, agentAccountID string const getChannel = `-- name: GetChannel :one SELECT id, name, COALESCE(group_id, '') AS group_id, kind, post_policy, - COALESCE(owner_account_id, '') AS owner_account_id, mandatory_subscription + COALESCE(owner_account_id, '') AS owner_account_id, mandatory_subscription, + COALESCE(parent_agent_id, '') AS parent_agent_id, membership_mode FROM channels WHERE id = $1 ` @@ -351,6 +401,8 @@ type GetChannelRow struct { PostPolicy int16 OwnerAccountID string MandatorySubscription bool + ParentAgentID string + MembershipMode int16 } func (q *Queries) GetChannel(ctx context.Context, id string) (GetChannelRow, error) { @@ -364,6 +416,8 @@ func (q *Queries) GetChannel(ctx context.Context, id string) (GetChannelRow, err &i.PostPolicy, &i.OwnerAccountID, &i.MandatorySubscription, + &i.ParentAgentID, + &i.MembershipMode, ) return i, err } @@ -575,9 +629,15 @@ effective AS ( SELECT id, MIN(min_vis) AS eff_vis FROM ancestry GROUP BY id +), +viewer AS ( + SELECT owner_user_id AS uid FROM agent_accounts WHERE account_id = $1 + UNION ALL + SELECT $1 AS uid ) SELECT c.id, c.name, COALESCE(c.group_id, '') AS group_id, c.kind, c.post_policy, - COALESCE(c.owner_account_id, '') AS owner_account_id, c.mandatory_subscription + COALESCE(c.owner_account_id, '') AS owner_account_id, c.mandatory_subscription, + COALESCE(c.parent_agent_id, '') AS parent_agent_id, c.membership_mode FROM channels c WHERE ( EXISTS ( @@ -589,6 +649,11 @@ WHERE ( SELECT 1 FROM effective e WHERE e.id = c.group_id AND e.eff_vis = 1 ) ) + OR ( + c.parent_agent_id IS NOT NULL + AND (SELECT aa.owner_user_id FROM agent_accounts aa WHERE aa.account_id = c.parent_agent_id) + IN (SELECT uid FROM viewer) + ) ) ORDER BY c.name ` @@ -601,6 +666,8 @@ type ListChannelsRow struct { PostPolicy int16 OwnerAccountID string MandatorySubscription bool + ParentAgentID string + MembershipMode int16 } func (q *Queries) ListChannels(ctx context.Context, accountID string) ([]ListChannelsRow, error) { @@ -620,6 +687,8 @@ func (q *Queries) ListChannels(ctx context.Context, accountID string) ([]ListCha &i.PostPolicy, &i.OwnerAccountID, &i.MandatorySubscription, + &i.ParentAgentID, + &i.MembershipMode, ); err != nil { return nil, err } diff --git a/go/internal/store/db/messages.sql.go b/go/internal/store/db/messages.sql.go index 75054611e..388e00ec9 100644 --- a/go/internal/store/db/messages.sql.go +++ b/go/internal/store/db/messages.sql.go @@ -10,12 +10,30 @@ import ( ) const findAskMessage = `-- name: FindAskMessage :many +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 + AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $1) + ) +) SELECT m.id, m.topic_id, m.author_account_id, (CASE WHEN ah.owner_user_id IS NULL THEN COALESCE(ah.handle, '') WHEN oh.handle IS NULL THEN '' ELSE oh.handle || '/' || ah.handle END)::text AS author_handle, m.at_unix_ms, m.blocks, m.turn_sequence FROM messages m LEFT JOIN account_handles ah ON ah.account_id = m.author_account_id LEFT JOIN account_handles oh ON oh.account_id = ah.owner_user_id JOIN topics t ON t.id = m.topic_id -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 +JOIN participating p ON p.channel_id = t.channel_id WHERE m.blocks @> $2::jsonb FOR UPDATE OF m ` @@ -65,6 +83,7 @@ func (q *Queries) FindAskMessage(ctx context.Context, arg FindAskMessageParams) const getChannelPostPolicy = `-- name: GetChannelPostPolicy :one + SELECT post_policy, COALESCE(owner_account_id, '') AS owner_account_id, name FROM channels WHERE id = $1 ` @@ -83,6 +102,10 @@ type GetChannelPostPolicyRow struct { // not-found/forbidden error mapping — all hand-written around these generated // calls. Every message read shares the id/topic_id/author_account_id/author_handle/ // at_unix_ms/blocks/turn_sequence projection so Go maps each row through messageFromParts. +// Participant-channel copies: every `chain` + `participating` CTE in this +// file and topics.sql MUST stay identical (bar the actor's parameter number) and +// equal to ChannelParticipant (channels.sql). It is participation, not channel +// visibility: never widen it to the owner-set predicate. ACL conjuncts go in each. func (q *Queries) GetChannelPostPolicy(ctx context.Context, id string) (GetChannelPostPolicyRow, error) { row := q.db.QueryRow(ctx, getChannelPostPolicy, id) var i GetChannelPostPolicyRow @@ -165,9 +188,27 @@ func (q *Queries) GetMessageByRequestID(ctx context.Context, arg GetMessageByReq } const getPageCursorSeq = `-- name: GetPageCursorSeq :one +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 + AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $1) + ) +) SELECT m.seq FROM messages m JOIN topics t ON t.id = m.topic_id -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 +JOIN participating p ON p.channel_id = t.channel_id WHERE m.id = $2 AND t.channel_id = $3 ` @@ -290,12 +331,30 @@ func (q *Queries) InsertTopicIgnore(ctx context.Context, arg InsertTopicIgnorePa } const listMessages = `-- name: ListMessages :many +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 + AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $1) + ) +) SELECT m.id, m.topic_id, m.author_account_id, (CASE WHEN ah.owner_user_id IS NULL THEN COALESCE(ah.handle, '') WHEN oh.handle IS NULL THEN '' ELSE oh.handle || '/' || ah.handle END)::text AS author_handle, m.at_unix_ms, m.blocks, m.turn_sequence FROM messages m LEFT JOIN account_handles ah ON ah.account_id = m.author_account_id LEFT JOIN account_handles oh ON oh.account_id = ah.owner_user_id JOIN topics t ON t.id = m.topic_id -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 +JOIN participating p ON p.channel_id = t.channel_id WHERE t.channel_id = $2 AND ($3 = 0 OR m.seq < $3) AND ($5 = 0 OR m.seq <= $5) AND ($6 = '' OR m.topic_id = $6) ORDER BY m.seq DESC @@ -377,12 +436,30 @@ func (q *Queries) ReviveTopic(ctx context.Context, id string) error { } const searchMessages = `-- name: SearchMessages :many +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 + AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $1) + ) +) SELECT m.id, m.topic_id, m.author_account_id, (CASE WHEN ah.owner_user_id IS NULL THEN COALESCE(ah.handle, '') WHEN oh.handle IS NULL THEN '' ELSE oh.handle || '/' || ah.handle END)::text AS author_handle, m.at_unix_ms, m.blocks, t.channel_id, m.turn_sequence FROM messages m LEFT JOIN account_handles ah ON ah.account_id = m.author_account_id LEFT JOIN account_handles oh ON oh.account_id = ah.owner_user_id JOIN topics t ON t.id = m.topic_id -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 +JOIN participating p ON p.channel_id = t.channel_id WHERE m.search_tsv @@ websearch_to_tsquery('english', $2) AND ($3 = '' OR t.channel_id = $3) AND ($5 = 0 OR m.seq <= $5) @@ -463,16 +540,28 @@ func (q *Queries) UpdateMessageBlocks(ctx context.Context, arg UpdateMessageBloc } const updateMessageBlocksAsAuthor = `-- name: UpdateMessageBlocksAsAuthor :one +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $4 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $4 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $4) + ) +) UPDATE messages m SET blocks = $1, text_content = $2 FROM topics t -WHERE m.id = $3 - AND t.id = m.topic_id - AND m.author_account_id = $4 - AND EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $4 - ) +JOIN participating p ON p.channel_id = t.channel_id +WHERE m.id = $3 AND t.id = m.topic_id AND m.author_account_id = $4 RETURNING m.id, m.topic_id, m.author_account_id, COALESCE((SELECT (CASE WHEN author_handles.owner_user_id IS NULL THEN author_handles.handle WHEN owner_handles.handle IS NULL THEN '' ELSE owner_handles.handle || '/' || author_handles.handle END)::text FROM account_handles AS author_handles LEFT JOIN account_handles AS owner_handles ON owner_handles.account_id = author_handles.owner_user_id WHERE author_handles.account_id = $4), '')::text AS author_handle, m.at_unix_ms, m.blocks, m.turn_sequence diff --git a/go/internal/store/db/querier.go b/go/internal/store/db/querier.go index 15f7f2943..8e5dd9dfd 100644 --- a/go/internal/store/db/querier.go +++ b/go/internal/store/db/querier.go @@ -208,6 +208,10 @@ type Querier interface { // not-found/forbidden error mapping — all hand-written around these generated // calls. Every message read shares the id/topic_id/author_account_id/author_handle/ // at_unix_ms/blocks/turn_sequence projection so Go maps each row through messageFromParts. + // Participant-channel copies: every `chain` + `participating` CTE in this + // file and topics.sql MUST stay identical (bar the actor's parameter number) and + // equal to ChannelParticipant (channels.sql). It is participation, not channel + // visibility: never widen it to the owner-set predicate. ACL conjuncts go in each. GetChannelPostPolicy(ctx context.Context, id string) (GetChannelPostPolicyRow, error) GetCoordinationChannelByName(ctx context.Context, arg GetCoordinationChannelByNameParams) (GetCoordinationChannelByNameRow, error) // Coordination-store queries (sqlc adoption T3, RIG-3034). These replace the @@ -393,6 +397,10 @@ type Querier interface { // loop, and the D9 not-found/forbidden error mapping. The topic projection // (id, channel_id, name, created_by_account_id, created_at_unix_ms, archived, // last_seq) matches the former scanTopics order so the Go maps each row to Topic. + // Participant-channel copies: every `chain` + `participating` CTE in this + // file and messages.sql MUST stay identical (bar the actor's parameter number) and + // equal to ChannelParticipant (channels.sql). It is participation, not channel + // visibility: never widen it to the owner-set predicate. ACL conjuncts go in each. ListTopics(ctx context.Context, arg ListTopicsParams) ([]Topic, error) ListUserPeerings(ctx context.Context, userID string) ([]ListUserPeeringsRow, error) ListVisibleAccounts(ctx context.Context, id string) ([]ListVisibleAccountsRow, error) diff --git a/go/internal/store/db/topics.sql.go b/go/internal/store/db/topics.sql.go index 4112ce0a3..832cd4995 100644 --- a/go/internal/store/db/topics.sql.go +++ b/go/internal/store/db/topics.sql.go @@ -41,6 +41,7 @@ func (q *Queries) GetTopic(ctx context.Context, id string) (Topic, error) { const listTopics = `-- name: ListTopics :many + SELECT id, channel_id, name, created_by_account_id, created_at_unix_ms, archived, last_seq, tenant_id FROM topics WHERE channel_id = $1 AND ($2 OR NOT archived) @@ -58,6 +59,10 @@ type ListTopicsParams struct { // loop, and the D9 not-found/forbidden error mapping. The topic projection // (id, channel_id, name, created_by_account_id, created_at_unix_ms, archived, // last_seq) matches the former scanTopics order so the Go maps each row to Topic. +// Participant-channel copies: every `chain` + `participating` CTE in this +// file and messages.sql MUST stay identical (bar the actor's parameter number) and +// equal to ChannelParticipant (channels.sql). It is participation, not channel +// visibility: never widen it to the owner-set predicate. ACL conjuncts go in each. func (q *Queries) ListTopics(ctx context.Context, arg ListTopicsParams) ([]Topic, error) { rows, err := q.db.Query(ctx, listTopics, arg.ChannelID, arg.Column2) if err != nil { @@ -131,8 +136,26 @@ func (q *Queries) RenameTopic(ctx context.Context, arg RenameTopicParams) error } const resolveTopicForUpdate = `-- name: ResolveTopicForUpdate :one +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 + AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $1) + ) +) SELECT t.channel_id FROM topics t -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 +JOIN participating p ON p.channel_id = t.channel_id WHERE t.id = $2 FOR UPDATE OF t ` diff --git a/go/internal/store/messages.go b/go/internal/store/messages.go index 43616d645..7b87ae795 100644 --- a/go/internal/store/messages.go +++ b/go/internal/store/messages.go @@ -25,7 +25,7 @@ import ( // suppresses a duplicate MessagePosted fan-out for a row that did not change. A // message with no blocks, a TopicRef that is neither exactly-id nor // exactly-name, or a TopicRef.ID naming a topic in another channel (or no -// topic) is ErrInvalidArgument. The membership check, topic resolution, insert, +// topic) is ErrInvalidArgument. The participation check, topic resolution, insert, // and last_seq denormalization all run in one transaction, so a membership // revoked between them cannot slip a message into a channel the author can no // longer read, and a get-or-created topic never outlives a rolled-back insert. @@ -54,10 +54,10 @@ func (s *Store) AppendMessage(ctx context.Context, m Message, channelID string, } defer func() { _ = tx.Rollback(ctx) }() // deferred cleanup; the Commit below is the real outcome. - // D9 write-authz: the author must be a member of the target channel, so a - // non-member cannot persist into a private channel it can't see. A - // non-member gets ErrNotFound, never a hint the channel exists. Checked in - // the insert tx so a concurrent removal cannot race the gate. + // D9 write-authz: the author must participate in the target channel (member + // row or TREE derivation); a non-participant gets ErrNotFound, never a hint + // the channel exists. Checked in the insert tx so a concurrent removal or + // tree move cannot race the gate. if err := requireChannelMember(ctx, tx, m.AuthorAccountID, ChannelID(channelID)); err != nil { return Message{}, false, err } @@ -299,18 +299,18 @@ func updateMessageBlocksExec(ctx context.Context, dbtx db.DBTX, id MessageID, bl // frame, RIG-1364 T3). // // Why it is a fork and not a flag on the shared core. updateMessageBlocksExec -// addresses the row by a bare MessageID with NO membership and NO authorship +// addresses the row by a bare MessageID with NO participation and NO authorship // check. That is correct where it is used — AnswerAsk has already resolved the -// target through a membership JOIN and locked it FOR UPDATE, so re-checking -// would be redundant, and AnswerAsk deliberately permits a MEMBER who is not +// target through the participant join and locked it FOR UPDATE, so re-checking +// would be redundant, and AnswerAsk deliberately permits a PARTICIPANT who is not // the author to answer. Folding this path onto that core would either strip the // authz a relayed id requires or break every ask answered by anyone but the // asker. The two predicates genuinely differ, so they are two statements. // -// The predicate is membership AND authorship: the actor must be a member of the -// message's channel and be its author. Both halves are load-bearing — authorship +// The predicate is participation AND authorship: the actor must participate in +// the message's channel and be its author. Both halves are load-bearing — authorship // alone would let an account edit its own past message in a channel it has since -// been removed from, and membership alone would let any member rewrite another +// been removed from, and participation alone would let anyone rewrite another // account's words. A failure of EITHER half, and an id that names no row at all, // return the same ErrNotFound: the D9 not-found/forbidden merge, so an actor // cannot learn that a message it may not touch exists (the same collapse @@ -352,9 +352,9 @@ func (s *Store) UpdateMessageBlocksAsAuthor(ctx context.Context, actor AccountID } // One statement, so the authz predicate and the write cannot race: a - // concurrent membership revocation lands before the UPDATE (matches no row) - // or after it, never between. The EXISTS subquery is the membership half and - // the author_account_id equality the authorship half; both must hold. + // concurrent membership revocation or tree move lands before the UPDATE + // (matches no row) or after it, never between. The participating join is the + // participation half, author_account_id equality the authorship half. row, err := s.q.UpdateMessageBlocksAsAuthor(ctx, db.UpdateMessageBlocksAsAuthorParams{ Blocks: blocksJSON, TextContent: textContent(blocks), @@ -413,8 +413,8 @@ func (s *Store) MessageAskIDs(ctx context.Context, actor AccountID, id MessageID // // The channel is resolved THROUGH the topic join now that a message carries no // channel_id: messages JOIN topics ON topic_id, filtered by topics.channel_id. -// Visibility is enforced in SQL — the channel must be one the actor is a member -// of (JOIN channel_members on the topic's channel), so a non-member — or a +// Visibility is enforced in SQL — the channel must be one the actor participates +// in (member row or TREE derivation), so a non-participant — or a // caller naming a channel it cannot see — reads nothing rather than leaking a // private channel's history by id (the D9 not-found/forbidden merge, matching // SearchMessages). The visibility gate is the store's, not the RPC edge's. @@ -426,7 +426,7 @@ func (s *Store) ListMessages(ctx context.Context, q ListMessagesQuery) ([]Messag var beforeSeq int64 if q.Page.BeforeMessageID != "" { - // Scope the cursor probe to the actor's membership too, so a non-member + // Scope the cursor probe to the actor's participation too, so a non-participant // naming a real message in a channel it cannot see gets the same "not in // channel" result as a fake id — no existence oracle across the // visibility boundary. The channel is the cursor message's topic's channel. @@ -445,14 +445,14 @@ func (s *Store) ListMessages(ctx context.Context, q ListMessagesQuery) ([]Messag } // A zero beforeSeq reads the newest page; a positive one pages strictly - // older. The membership JOIN scopes the read to the actor's visible set; a + // older. The participant join scopes the read to the actor's channels; a // non-zero SnapshotSeq bounds it to the subscribe-time snapshot on set // MEMBERSHIP, not content. - // Membership-only is sufficient, not a lost update: the matching + // Set-membership-only is sufficient, not a lost update: the matching // MessageUpdated rides the live tail, so an id-deduping client converges // last-write-wins. Freezing content too would need a change-seq and a - // schema change; membership-only is the ratified scope. + // schema change; set-membership-only is the ratified scope. snap := int64(q.Page.SnapshotSeq) //nolint:gosec // G115: server-issued seq, int64 domain rows, err := s.q.ListMessages(ctx, db.ListMessagesParams{ AccountID: string(q.Actor), @@ -475,7 +475,7 @@ func (s *Store) ListMessages(ctx context.Context, q ListMessagesQuery) ([]Messag // bounds. Visibility is enforced in SQL — the actor sees a message only in a // channel it belongs to — so a scope pointing at a channel the actor cannot see // yields nothing rather than leaking. Under FORCE RLS messages_search_idx cannot -// serve the `@@` filter (it is not LEAKPROOF); the membership join may narrow rows. +// serve the `@@` filter (it is not LEAKPROOF); the participant join narrows rows. func (s *Store) SearchMessages(ctx context.Context, actor AccountID, scope SearchScope, query string, page Page) ([]Message, error) { if strings.TrimSpace(query) == "" { return nil, fmt.Errorf("%w: search query is required", ErrInvalidArgument) @@ -503,7 +503,7 @@ func (s *Store) SearchMessages(ctx context.Context, actor AccountID, scope Searc // AnswerAsk records a participant's atomic answer to a pending structured ask // (RespondToAsk; see docs/designs/agent/compass-ask-typed-derivation.md). It // locates the message whose blocks carry an ask with askID within the actor's -// visible set — the membership JOIN makes "the message exists" and "the actor +// visible set — the participant join makes "the message exists" and "the actor // participates" one gate — records the per-question answers on that ask block, // and persists via the immutable-ask_id update path. It ALSO posts the answer // as a new message — authored by actor, in the ask's channel/topic, carrying a diff --git a/go/internal/store/queries/channels.sql b/go/internal/store/queries/channels.sql index 5f998bf2f..68811c27d 100644 --- a/go/internal/store/queries/channels.sql +++ b/go/internal/store/queries/channels.sql @@ -111,15 +111,39 @@ UPDATE channels SET post_policy = $2, owner_account_id = NULLIF($3, ''), mandato -- name: GetChannel :one SELECT id, name, COALESCE(group_id, '') AS group_id, kind, post_policy, - COALESCE(owner_account_id, '') AS owner_account_id, mandatory_subscription + COALESCE(owner_account_id, '') AS owner_account_id, mandatory_subscription, + COALESCE(parent_agent_id, '') AS parent_agent_id, membership_mode FROM channels WHERE id = $1; -- name: ChannelMembersByChannelIDs :many -SELECT channel_id, account_id, subscribed -FROM channel_members -WHERE channel_id = ANY($1::text[]) -ORDER BY account_id; - +WITH RECURSIVE stored AS ( + SELECT cm.channel_id, cm.account_id, cm.subscribed + FROM channel_members cm + JOIN channels c ON c.id = cm.channel_id + WHERE cm.channel_id = ANY($1::text[]) AND c.membership_mode = 0 +), subtree AS ( + SELECT c.id AS channel_id, c.parent_agent_id AS account_id + FROM channels c + WHERE c.id = ANY($1::text[]) AND c.membership_mode = 1 + UNION + SELECT s.channel_id, a.account_id + FROM agent_accounts a + JOIN subtree s ON a.parent_agent_id = s.account_id +), participants AS ( + SELECT channel_id, account_id FROM subtree + UNION + SELECT c.id AS channel_id, aa.owner_user_id AS account_id + FROM channels c + JOIN agent_accounts aa ON aa.account_id = c.parent_agent_id + WHERE c.id = ANY($1::text[]) AND c.membership_mode = 1 +) +SELECT channel_id, account_id, subscribed FROM stored +UNION ALL +SELECT p.channel_id, p.account_id, COALESCE(cs.subscribed, FALSE) AS subscribed +FROM participants p +LEFT JOIN channel_subscriptions cs + ON cs.channel_id = p.channel_id AND cs.account_id = p.account_id +ORDER BY channel_id, account_id; -- name: ListChannelGroups :many WITH RECURSIVE ancestry AS ( SELECT id, parent_group_id, visibility AS min_vis @@ -185,9 +209,15 @@ effective AS ( SELECT id, MIN(min_vis) AS eff_vis FROM ancestry GROUP BY id +), +viewer AS ( + SELECT owner_user_id AS uid FROM agent_accounts WHERE account_id = $1 + UNION ALL + SELECT $1 AS uid ) SELECT c.id, c.name, COALESCE(c.group_id, '') AS group_id, c.kind, c.post_policy, - COALESCE(c.owner_account_id, '') AS owner_account_id, c.mandatory_subscription + COALESCE(c.owner_account_id, '') AS owner_account_id, c.mandatory_subscription, + COALESCE(c.parent_agent_id, '') AS parent_agent_id, c.membership_mode FROM channels c WHERE ( EXISTS ( @@ -199,6 +229,11 @@ WHERE ( SELECT 1 FROM effective e WHERE e.id = c.group_id AND e.eff_vis = 1 ) ) + OR ( + c.parent_agent_id IS NOT NULL + AND (SELECT aa.owner_user_id FROM agent_accounts aa WHERE aa.account_id = c.parent_agent_id) + IN (SELECT uid FROM viewer) + ) ) ORDER BY c.name; @@ -215,6 +250,11 @@ effective AS ( SELECT id, MIN(min_vis) AS eff_vis FROM ancestry GROUP BY id +), +viewer AS ( + SELECT owner_user_id AS uid FROM agent_accounts WHERE account_id = $1 + UNION ALL + SELECT $1 AS uid ) SELECT EXISTS ( SELECT 1 FROM channels c @@ -228,6 +268,11 @@ SELECT EXISTS ( SELECT 1 FROM effective e WHERE e.id = c.group_id AND e.eff_vis = 1 ) ) + OR ( + c.parent_agent_id IS NOT NULL + AND (SELECT aa.owner_user_id FROM agent_accounts aa WHERE aa.account_id = c.parent_agent_id) + IN (SELECT uid FROM viewer) + ) ) ); @@ -244,9 +289,15 @@ effective AS ( SELECT id, MIN(min_vis) AS eff_vis FROM ancestry GROUP BY id +), +viewer AS ( + SELECT owner_user_id AS uid FROM agent_accounts WHERE account_id = $1 + UNION ALL + SELECT $1 AS uid ) SELECT c.id, c.name, COALESCE(c.group_id, '') AS group_id, c.kind, c.post_policy, - COALESCE(c.owner_account_id, '') AS owner_account_id, c.mandatory_subscription + COALESCE(c.owner_account_id, '') AS owner_account_id, c.mandatory_subscription, + COALESCE(c.parent_agent_id, '') AS parent_agent_id, c.membership_mode FROM channels c WHERE c.name = $2 AND ( EXISTS ( @@ -258,6 +309,11 @@ WHERE c.name = $2 AND ( SELECT 1 FROM effective e WHERE e.id = c.group_id AND e.eff_vis = 1 ) ) + OR ( + c.parent_agent_id IS NOT NULL + AND (SELECT aa.owner_user_id FROM agent_accounts aa WHERE aa.account_id = c.parent_agent_id) + IN (SELECT uid FROM viewer) + ) ) ORDER BY c.id; diff --git a/go/internal/store/queries/messages.sql b/go/internal/store/queries/messages.sql index 9bff5d624..30c596d86 100644 --- a/go/internal/store/queries/messages.sql +++ b/go/internal/store/queries/messages.sql @@ -7,6 +7,11 @@ -- calls. Every message read shares the id/topic_id/author_account_id/author_handle/ -- at_unix_ms/blocks/turn_sequence projection so Go maps each row through messageFromParts. +-- Participant-channel copies: every `chain` + `participating` CTE in this +-- file and topics.sql MUST stay identical (bar the actor's parameter number) and +-- equal to ChannelParticipant (channels.sql). It is participation, not channel +-- visibility: never widen it to the owner-set predicate. ACL conjuncts go in each. + -- name: GetChannelPostPolicy :one SELECT post_policy, COALESCE(owner_account_id, '') AS owner_account_id, name FROM channels WHERE id = $1; @@ -43,16 +48,28 @@ SELECT COALESCE(MAX(seq), 0)::BIGINT AS head FROM messages; UPDATE messages SET blocks = $1, text_content = $2 WHERE id = $3; -- name: UpdateMessageBlocksAsAuthor :one +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $4 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $4 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $4) + ) +) UPDATE messages m SET blocks = $1, text_content = $2 FROM topics t -WHERE m.id = $3 - AND t.id = m.topic_id - AND m.author_account_id = $4 - AND EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $4 - ) +JOIN participating p ON p.channel_id = t.channel_id +WHERE m.id = $3 AND t.id = m.topic_id AND m.author_account_id = $4 RETURNING m.id, m.topic_id, m.author_account_id, COALESCE((SELECT (CASE WHEN author_handles.owner_user_id IS NULL THEN author_handles.handle WHEN owner_handles.handle IS NULL THEN '' ELSE owner_handles.handle || '/' || author_handles.handle END)::text FROM account_handles AS author_handles LEFT JOIN account_handles AS owner_handles ON owner_handles.account_id = author_handles.owner_user_id WHERE author_handles.account_id = $4), '')::text AS author_handle, m.at_unix_ms, m.blocks, m.turn_sequence; @@ -68,30 +85,83 @@ WHERE m.id = $1 ); -- name: GetPageCursorSeq :one +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 + AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $1) + ) +) SELECT m.seq FROM messages m JOIN topics t ON t.id = m.topic_id -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 +JOIN participating p ON p.channel_id = t.channel_id WHERE m.id = $2 AND t.channel_id = $3; -- name: ListMessages :many +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 + AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $1) + ) +) SELECT m.id, m.topic_id, m.author_account_id, (CASE WHEN ah.owner_user_id IS NULL THEN COALESCE(ah.handle, '') WHEN oh.handle IS NULL THEN '' ELSE oh.handle || '/' || ah.handle END)::text AS author_handle, m.at_unix_ms, m.blocks, m.turn_sequence FROM messages m LEFT JOIN account_handles ah ON ah.account_id = m.author_account_id LEFT JOIN account_handles oh ON oh.account_id = ah.owner_user_id JOIN topics t ON t.id = m.topic_id -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 +JOIN participating p ON p.channel_id = t.channel_id WHERE t.channel_id = $2 AND ($3 = 0 OR m.seq < $3) AND ($5 = 0 OR m.seq <= $5) AND ($6 = '' OR m.topic_id = $6) ORDER BY m.seq DESC LIMIT $4; - -- name: SearchMessages :many +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 + AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $1) + ) +) SELECT m.id, m.topic_id, m.author_account_id, (CASE WHEN ah.owner_user_id IS NULL THEN COALESCE(ah.handle, '') WHEN oh.handle IS NULL THEN '' ELSE oh.handle || '/' || ah.handle END)::text AS author_handle, m.at_unix_ms, m.blocks, t.channel_id, m.turn_sequence FROM messages m LEFT JOIN account_handles ah ON ah.account_id = m.author_account_id LEFT JOIN account_handles oh ON oh.account_id = ah.owner_user_id JOIN topics t ON t.id = m.topic_id -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 +JOIN participating p ON p.channel_id = t.channel_id WHERE m.search_tsv @@ websearch_to_tsquery('english', $2) AND ($3 = '' OR t.channel_id = $3) AND ($5 = 0 OR m.seq <= $5) @@ -99,12 +169,30 @@ ORDER BY ts_rank(m.search_tsv, websearch_to_tsquery('english', $2)) DESC, m.seq LIMIT $4; -- name: FindAskMessage :many +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 + AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $1) + ) +) SELECT m.id, m.topic_id, m.author_account_id, (CASE WHEN ah.owner_user_id IS NULL THEN COALESCE(ah.handle, '') WHEN oh.handle IS NULL THEN '' ELSE oh.handle || '/' || ah.handle END)::text AS author_handle, m.at_unix_ms, m.blocks, m.turn_sequence FROM messages m LEFT JOIN account_handles ah ON ah.account_id = m.author_account_id LEFT JOIN account_handles oh ON oh.account_id = ah.owner_user_id JOIN topics t ON t.id = m.topic_id -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 +JOIN participating p ON p.channel_id = t.channel_id WHERE m.blocks @> $2::jsonb FOR UPDATE OF m; diff --git a/go/internal/store/queries/topics.sql b/go/internal/store/queries/topics.sql index 5e8b2125c..f33484638 100644 --- a/go/internal/store/queries/topics.sql +++ b/go/internal/store/queries/topics.sql @@ -5,6 +5,11 @@ -- (id, channel_id, name, created_by_account_id, created_at_unix_ms, archived, -- last_seq) matches the former scanTopics order so the Go maps each row to Topic. +-- Participant-channel copies: every `chain` + `participating` CTE in this +-- file and messages.sql MUST stay identical (bar the actor's parameter number) and +-- equal to ChannelParticipant (channels.sql). It is participation, not channel +-- visibility: never widen it to the owner-set predicate. ACL conjuncts go in each. + -- name: ListTopics :many SELECT id, channel_id, name, created_by_account_id, created_at_unix_ms, archived, last_seq, tenant_id FROM topics @@ -12,8 +17,26 @@ WHERE channel_id = $1 AND ($2 OR NOT archived) ORDER BY last_seq DESC, created_at_unix_ms DESC, id; -- name: ResolveTopicForUpdate :one +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +), participating AS ( + SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 + UNION + SELECT c.id FROM channels c + WHERE c.membership_mode = 1 + AND ( + c.parent_agent_id IN (SELECT account_id FROM chain) + OR c.parent_agent_id IN (SELECT aa.account_id FROM agent_accounts aa WHERE aa.owner_user_id = $1) + ) +) SELECT t.channel_id FROM topics t -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 +JOIN participating p ON p.channel_id = t.channel_id WHERE t.id = $2 FOR UPDATE OF t; diff --git a/go/internal/store/topics.go b/go/internal/store/topics.go index 68e5adc5e..4d8b12bbb 100644 --- a/go/internal/store/topics.go +++ b/go/internal/store/topics.go @@ -13,7 +13,7 @@ import ( // ListTopics returns the topics in channelID, newest-activity-first (last_seq // descending, then birth time), scoped to the caller's visible set. Archived // topics are omitted unless includeArchived is set. Visibility is the same D9 -// gate the message reads apply: a caller who is not a member of the channel — +// gate the message reads apply: a caller who does not participate in the channel — // or names a channel it cannot see — gets ErrNotFound (the not-found/forbidden // merge), never a hint the channel exists or an empty list it could mistake for // "no topics". @@ -26,7 +26,7 @@ func (s *Store) ListTopics(ctx context.Context, callerAccountID, channelID strin return nil, err } if !member { - // D9 merge: a non-member cannot tell an unauthorized channel from a + // D9 merge: a non-participant cannot tell an unauthorized channel from a // nonexistent one, so the refusal enumerates nothing. return nil, fmt.Errorf("%w: channel %q", ErrNotFound, channelID) } @@ -45,9 +45,9 @@ func (s *Store) ListTopics(ctx context.Context, callerAccountID, channelID strin // source's, and the emptied source row is deleted, all in one transaction. The // surviving topic is returned. // -// The caller must be a member of the topic's channel; a topic it cannot see — +// The caller must participate in the topic's channel; a topic it cannot see — // or an unknown topicID — is ErrNotFound (the D9 not-found/forbidden merge, so -// topic existence cannot leak across a membership boundary). name and archived +// topic existence cannot leak across a participation boundary). name and archived // are each optional (nil = leave unchanged); the archived flag is applied to // the SURVIVING topic (the target on a merge, the source otherwise). // @@ -64,8 +64,8 @@ func (s *Store) UpdateTopic(ctx context.Context, callerAccountID, topicID string } defer func() { _ = tx.Rollback(ctx) }() // deferred cleanup; the Commit below is the real outcome. - // Resolve the topic + gate visibility in one statement: the caller must be a - // member of the topic's channel. Zero rows (unknown topic OR non-member) -> + // Resolve the topic + gate visibility in one statement: the caller must + // participate in the topic's channel. Zero rows (unknown topic OR non-participant) -> // ErrNotFound. FOR UPDATE OF t locks the source topic row for the tx so a // concurrent rename/merge serializes. q := db.New(tx) diff --git a/go/internal/store/types.go b/go/internal/store/types.go index eb7cc1f0a..7a1732d1e 100644 --- a/go/internal/store/types.go +++ b/go/internal/store/types.go @@ -252,6 +252,10 @@ type Channel struct { // owner-scoped to its creating caller (the OWNER default), not global. GroupID ChannelGroupID Kind ChannelKind + // ParentAgentID is the attached agent; empty means the channel is at the tree root. + ParentAgentID AccountID + // MembershipMode selects stored membership or agent-tree-derived membership. + MembershipMode ChannelMembershipMode // MemberAccountIDs are the accounts party to the channel. For DM/GROUP_DM // this set governs visibility directly (design.md:235-243). MemberAccountIDs []AccountID