From f63fd4d72141b589291d029a8d05b209cdc1e9da Mon Sep 17 00:00:00 2001 From: mintaka Date: Tue, 6 Oct 2026 01:03:48 -0400 Subject: [PATCH 1/5] feat(store): owner-set visibility and derived TREE members on channel reads (RIG-4188) The three channel predicate copies grant agent-attached channels to the anchor owner set; projections carry parent_agent_id and membership_mode to the wire; message and topic reads gate on participation; and loadChannelMembers materializes each TREE channel subtree, with subscribers intersected with it. CreateChannel returns the post-commit getChannel read. The RPC response and event assertion lands with the RPC wiring. Co-authored-by: Matt Wilkinson --- .../comms/channel_attach_pgtest_test.go | 49 ++++++ go/internal/comms/channel_attach_wire_test.go | 11 ++ go/internal/comms/mapping.go | 8 + .../store/channel_participant_pgtest_test.go | 159 ++++++++++++++++++ go/internal/store/channel_tree_pgtest_test.go | 50 +++++- go/internal/store/channels.go | 52 +++--- go/internal/store/coordination_pgtest_test.go | 4 +- go/internal/store/db/channels.sql.go | 83 ++++++++- go/internal/store/db/messages.sql.go | 149 +++++++++++++++- go/internal/store/db/topics.sql.go | 32 +++- go/internal/store/queries/channels.sql | 72 +++++++- go/internal/store/queries/messages.sql | 152 ++++++++++++++++- go/internal/store/queries/topics.sql | 32 +++- go/internal/store/types.go | 4 + 14 files changed, 790 insertions(+), 67 deletions(-) create mode 100644 go/internal/comms/channel_attach_pgtest_test.go 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 8b227d2a6..c1cd4b120 100644 --- a/go/internal/comms/mapping.go +++ b/go/internal/comms/mapping.go @@ -77,9 +77,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_participant_pgtest_test.go b/go/internal/store/channel_participant_pgtest_test.go index 5b9175949..55888d416 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,162 @@ 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 TestChannelTreeReadPathsUseDerivedParticipants(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) + } + + 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) + } + + 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) + } + 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) + } + _, _, 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) + } + asks, err := s.q.FindAskMessage(t.Context(), db.FindAskMessageParams{AccountID: string(f.leaf.ID), Column2: askFilter}) + if err != nil || len(asks) != 1 { + t.Fatalf("FindAskMessage for subtree agent = %d rows, %v; want one ask", len(asks), err) + } + if _, _, err := s.AnswerAsk(t.Context(), f.sibling.ID, "tree-ask", []AskAnswer{{QuestionID: "q1", ChosenOptionIDs: []string{"opt-a"}}}); err == nil { + t.Fatal("owner-set sibling answered TREE history, want membership refusal") + } + 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..bf3d73ae4 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 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 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 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 788961801..0432864fa 100644 --- a/go/internal/store/channels.go +++ b/go/internal/store/channels.go @@ -139,8 +139,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 @@ -220,22 +219,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 @@ -427,7 +417,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 @@ -547,7 +537,7 @@ 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 @@ -996,7 +986,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 } @@ -1010,16 +1000,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) @@ -1030,15 +1019,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), diff --git a/go/internal/store/coordination_pgtest_test.go b/go/internal/store/coordination_pgtest_test.go index 863203c3d..eeb3a8f78 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 a6154449e..03ff5dd07 100644 --- a/go/internal/store/db/channels.sql.go +++ b/go/internal/store/db/channels.sql.go @@ -90,10 +90,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 +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 { @@ -177,6 +201,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 @@ -190,6 +219,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) + ) ) ) ` @@ -219,9 +253,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 ( @@ -233,6 +273,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 ` @@ -250,6 +295,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) { @@ -269,6 +316,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 } @@ -338,7 +387,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 ` @@ -350,6 +400,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) { @@ -363,6 +415,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 } @@ -536,9 +590,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 ( @@ -550,6 +610,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 ` @@ -562,6 +627,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) { @@ -581,6 +648,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 96cfead3d..c6a3986f9 100644 --- a/go/internal/store/db/messages.sql.go +++ b/go/internal/store/db/messages.sql.go @@ -10,13 +10,39 @@ 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 + AND EXISTS (SELECT 1 FROM channels WHERE membership_mode = 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 +) 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 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 WHERE m.blocks @> $2::jsonb + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) + ) FOR UPDATE OF m ` @@ -149,10 +175,36 @@ 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 + AND EXISTS (SELECT 1 FROM channels WHERE id = $3 AND membership_mode = 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 +) 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 WHERE m.id = $2 AND t.channel_id = $3 + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) + ) ` type GetPageCursorSeqParams struct { @@ -272,14 +324,40 @@ 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 + AND EXISTS (SELECT 1 FROM channels WHERE id = $2 AND membership_mode = 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 +) 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 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 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) + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) + ) ORDER BY m.seq DESC LIMIT $4 ` @@ -357,15 +435,44 @@ 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 + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.membership_mode = 1 AND ($3 = '' OR c.id = $3) + ) + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +) 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 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 WHERE m.search_tsv @@ websearch_to_tsquery('english', $2) AND ($3 = '' OR t.channel_id = $3) AND ($5 = 0 OR m.seq <= $5) + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) + ) ORDER BY ts_rank(m.search_tsv, websearch_to_tsquery('english', $2)) DESC, m.seq DESC LIMIT $4 ` @@ -441,15 +548,43 @@ 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 + AND EXISTS ( + SELECT 1 FROM messages m + JOIN topics t ON t.id = m.topic_id + JOIN channels c ON c.id = t.channel_id + WHERE m.id = $3 AND c.membership_mode = 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 +) 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 + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $4 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $4 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) ) 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, diff --git a/go/internal/store/db/topics.sql.go b/go/internal/store/db/topics.sql.go index 4112ce0a3..221a5ff5c 100644 --- a/go/internal/store/db/topics.sql.go +++ b/go/internal/store/db/topics.sql.go @@ -131,9 +131,39 @@ 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 + AND EXISTS ( + SELECT 1 FROM topics t + JOIN channels c ON c.id = t.channel_id + WHERE t.id = $2 AND c.membership_mode = 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 +) SELECT t.channel_id FROM topics t -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 WHERE t.id = $2 + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) + ) FOR UPDATE OF t ` diff --git a/go/internal/store/queries/channels.sql b/go/internal/store/queries/channels.sql index 42d80a04a..0b58518fc 100644 --- a/go/internal/store/queries/channels.sql +++ b/go/internal/store/queries/channels.sql @@ -104,15 +104,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 +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 @@ -176,9 +200,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 ( @@ -190,6 +220,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; @@ -206,6 +241,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 @@ -219,6 +259,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) + ) ) ); @@ -235,9 +280,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 ( @@ -249,6 +300,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 ec2dbe7ab..35f73728e 100644 --- a/go/internal/store/queries/messages.sql +++ b/go/internal/store/queries/messages.sql @@ -43,15 +43,43 @@ 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 + AND EXISTS ( + SELECT 1 FROM messages m + JOIN topics t ON t.id = m.topic_id + JOIN channels c ON c.id = t.channel_id + WHERE m.id = $3 AND c.membership_mode = 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 +) 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 + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $4 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $4 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) ) 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, @@ -61,44 +89,150 @@ RETURNING m.id, m.topic_id, m.author_account_id, SELECT blocks FROM messages WHERE 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 + AND EXISTS (SELECT 1 FROM channels WHERE id = $3 AND membership_mode = 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 +) 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 -WHERE m.id = $2 AND t.channel_id = $3; +WHERE m.id = $2 AND t.channel_id = $3 + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) + ); -- name: ListMessages :many +WITH RECURSIVE chain AS ( + SELECT aa.account_id, aa.parent_agent_id + FROM agent_accounts aa + WHERE aa.account_id = $1 + AND EXISTS (SELECT 1 FROM channels WHERE id = $2 AND membership_mode = 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 +) 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 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 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) + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) + ) 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 + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.membership_mode = 1 AND ($3 = '' OR c.id = $3) + ) + UNION + SELECT a.account_id, a.parent_agent_id + FROM agent_accounts a + JOIN chain ch ON a.account_id = ch.parent_agent_id +) 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 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 WHERE m.search_tsv @@ websearch_to_tsquery('english', $2) AND ($3 = '' OR t.channel_id = $3) AND ($5 = 0 OR m.seq <= $5) + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) + ) ORDER BY ts_rank(m.search_tsv, websearch_to_tsquery('english', $2)) DESC, m.seq DESC 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 + AND EXISTS (SELECT 1 FROM channels WHERE membership_mode = 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 +) 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 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 WHERE m.blocks @> $2::jsonb + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) + ) FOR UPDATE OF m; -- name: GetMessageByRequestID :many diff --git a/go/internal/store/queries/topics.sql b/go/internal/store/queries/topics.sql index 5e8b2125c..dea59cfe8 100644 --- a/go/internal/store/queries/topics.sql +++ b/go/internal/store/queries/topics.sql @@ -12,9 +12,39 @@ 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 + AND EXISTS ( + SELECT 1 FROM topics t + JOIN channels c ON c.id = t.channel_id + WHERE t.id = $2 AND c.membership_mode = 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 +) SELECT t.channel_id FROM topics t -JOIN channel_members cm ON cm.channel_id = t.channel_id AND cm.account_id = $1 WHERE t.id = $2 + AND ( + EXISTS ( + SELECT 1 FROM channel_members cm + WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 + ) + OR ( + EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) + AND EXISTS ( + SELECT 1 FROM channels c + WHERE c.id = t.channel_id AND ( + c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) + OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa + WHERE aa.account_id = c.parent_agent_id) + ) + ) + ) + ) FOR UPDATE OF t; -- name: SetTopicArchived :exec diff --git a/go/internal/store/types.go b/go/internal/store/types.go index 623b0f870..de05f0bc6 100644 --- a/go/internal/store/types.go +++ b/go/internal/store/types.go @@ -211,6 +211,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 From 7dbb8ce0349c3fd8e786ed2dd857f4abaf61696e Mon Sep 17 00:00:00 2001 From: mintaka Date: Tue, 6 Oct 2026 03:18:08 -0400 Subject: [PATCH 2/5] fix(store): join a participant channel set on reads; narrow name resolves to participants (RIG-4188) Message and topic reads join one participant-channel CTE instead of per-row OR subplans, restoring the driving relation. ChannelByNameForViewer narrows a multi-match to channels the viewer participates in, so a same-owner sibling channel of the same name no longer makes an agent own channel ambiguous. The product spec states that owner-set visibility does not grant history. Co-authored-by: Matt Wilkinson --- docs/specs/product/compass.md | 8 +- .../store/channel_by_name_pgtest_test.go | 30 ++++ .../store/channel_participant_pgtest_test.go | 72 +++++++-- go/internal/store/channel_tree_pgtest_test.go | 6 +- go/internal/store/channels.go | 23 ++- go/internal/store/db/channels.sql.go | 2 +- go/internal/store/db/messages.sql.go | 151 ++++++----------- go/internal/store/db/topics.sql.go | 32 ++-- go/internal/store/messages.go | 122 +++++--------- go/internal/store/queries/channels.sql | 2 +- go/internal/store/queries/messages.sql | 153 ++++++------------ go/internal/store/queries/topics.sql | 32 ++-- go/internal/store/topics.go | 39 ++--- 13 files changed, 293 insertions(+), 379 deletions(-) diff --git a/docs/specs/product/compass.md b/docs/specs/product/compass.md index 993786240..9790d4ee1 100644 --- a/docs/specs/product/compass.md +++ b/docs/specs/product/compass.md @@ -539,9 +539,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/store/channel_by_name_pgtest_test.go b/go/internal/store/channel_by_name_pgtest_test.go index e7075ed68..42256bc30 100644 --- a/go/internal/store/channel_by_name_pgtest_test.go +++ b/go/internal/store/channel_by_name_pgtest_test.go @@ -109,3 +109,33 @@ 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") +} diff --git a/go/internal/store/channel_participant_pgtest_test.go b/go/internal/store/channel_participant_pgtest_test.go index 55888d416..73e6ab3c5 100644 --- a/go/internal/store/channel_participant_pgtest_test.go +++ b/go/internal/store/channel_participant_pgtest_test.go @@ -137,7 +137,7 @@ func TestChannelVisibilityOwnerSetParity(t *testing.T) { } } -func TestChannelTreeReadPathsUseDerivedParticipants(t *testing.T) { +func TestChannelTreeParticipantReadsAndWrites(t *testing.T) { s := newTestStore(t) f := newTreeFixture(t, s) @@ -155,6 +155,10 @@ func TestChannelTreeReadPathsUseDerivedParticipants(t *testing.T) { 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) + } cursorSeq, err := s.q.GetPageCursorSeq(t.Context(), db.GetPageCursorSeqParams{ AccountID: string(f.leaf.ID), ID: string(message.ID), ChannelID: string(f.channel.ID), @@ -167,16 +171,61 @@ func TestChannelTreeReadPathsUseDerivedParticipants(t *testing.T) { t.Fatalf("ListMessages cursor for subtree agent = %v, %v; want empty page", page, 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) - } 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) + } + 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) @@ -185,16 +234,11 @@ func TestChannelTreeReadPathsUseDerivedParticipants(t *testing.T) { if err != nil { t.Fatalf("askIDContainmentFilter: %v", err) } - asks, err := s.q.FindAskMessage(t.Context(), db.FindAskMessageParams{AccountID: string(f.leaf.ID), Column2: askFilter}) - if err != nil || len(asks) != 1 { - t.Fatalf("FindAskMessage for subtree agent = %d rows, %v; want one ask", len(asks), err) - } - if _, _, err := s.AnswerAsk(t.Context(), f.sibling.ID, "tree-ask", []AskAnswer{{QuestionID: "q1", ChosenOptionIDs: []string{"opt-a"}}}); err == nil { - t.Fatal("owner-set sibling answered TREE history, want membership refusal") - } - 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) + 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") } func TestChannelTreeMembersAreAttributedAndSubscriptionsIntersect(t *testing.T) { diff --git a/go/internal/store/channel_tree_pgtest_test.go b/go/internal/store/channel_tree_pgtest_test.go index bf3d73ae4..9d68a7176 100644 --- a/go/internal/store/channel_tree_pgtest_test.go +++ b/go/internal/store/channel_tree_pgtest_test.go @@ -56,12 +56,12 @@ 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 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 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) } diff --git a/go/internal/store/channels.go b/go/internal/store/channels.go index 0432864fa..ed1bd4bae 100644 --- a/go/internal/store/channels.go +++ b/go/internal/store/channels.go @@ -542,13 +542,26 @@ func (s *Store) ChannelByNameForViewer(ctx context.Context, viewer AccountID, na if err := loadChannelMembers(ctx, s.scopedPool(), channels); err != nil { return Channel{}, err } + if len(channels) > 1 { + participants := channels[:0] + 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) + } + } + 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 the viewer participates in; address it by id", ErrInvalidArgument, name, len(channels)) } } @@ -1037,10 +1050,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/db/channels.sql.go b/go/internal/store/db/channels.sql.go index 03ff5dd07..a41d61272 100644 --- a/go/internal/store/db/channels.sql.go +++ b/go/internal/store/db/channels.sql.go @@ -112,7 +112,7 @@ WITH RECURSIVE stored AS ( WHERE c.id = ANY($1::text[]) AND c.membership_mode = 1 ) SELECT channel_id, account_id, subscribed FROM stored -UNION +UNION ALL SELECT p.channel_id, p.account_id, COALESCE(cs.subscribed, FALSE) AS subscribed FROM participants p LEFT JOIN channel_subscriptions cs diff --git a/go/internal/store/db/messages.sql.go b/go/internal/store/db/messages.sql.go index c6a3986f9..036651025 100644 --- a/go/internal/store/db/messages.sql.go +++ b/go/internal/store/db/messages.sql.go @@ -14,35 +14,27 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $1 - AND EXISTS (SELECT 1 FROM channels WHERE membership_mode = 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 +), visible 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 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 visible v ON v.channel_id = t.channel_id WHERE m.blocks @> $2::jsonb - AND ( - EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) FOR UPDATE OF m ` @@ -179,32 +171,24 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $1 - AND EXISTS (SELECT 1 FROM channels WHERE id = $3 AND membership_mode = 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 +), visible 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 visible v ON v.channel_id = t.channel_id WHERE m.id = $2 AND t.channel_id = $3 - AND ( - EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) ` type GetPageCursorSeqParams struct { @@ -328,36 +312,28 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $1 - AND EXISTS (SELECT 1 FROM channels WHERE id = $2 AND membership_mode = 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 +), visible 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 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 visible v ON v.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) - AND ( - EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) ORDER BY m.seq DESC LIMIT $4 ` @@ -439,40 +415,29 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $1 - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.membership_mode = 1 AND ($3 = '' OR c.id = $3) - ) UNION SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id +), visible 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 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 visible v ON v.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) - AND ( - EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) ORDER BY ts_rank(m.search_tsv, websearch_to_tsquery('english', $2)) DESC, m.seq DESC LIMIT $4 ` @@ -552,40 +517,24 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $4 - AND EXISTS ( - SELECT 1 FROM messages m - JOIN topics t ON t.id = m.topic_id - JOIN channels c ON c.id = t.channel_id - WHERE m.id = $3 AND c.membership_mode = 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 +), visible 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 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $4 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) +JOIN visible v ON v.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 diff --git a/go/internal/store/db/topics.sql.go b/go/internal/store/db/topics.sql.go index 221a5ff5c..a69aae239 100644 --- a/go/internal/store/db/topics.sql.go +++ b/go/internal/store/db/topics.sql.go @@ -135,35 +135,23 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $1 - AND EXISTS ( - SELECT 1 FROM topics t - JOIN channels c ON c.id = t.channel_id - WHERE t.id = $2 AND c.membership_mode = 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 +), visible 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 visible v ON v.channel_id = t.channel_id WHERE t.id = $2 - AND ( - EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) FOR UPDATE OF t ` diff --git a/go/internal/store/messages.go b/go/internal/store/messages.go index a9e529cd5..b12675579 100644 --- a/go/internal/store/messages.go +++ b/go/internal/store/messages.go @@ -12,22 +12,10 @@ import ( "github.com/RigelBuild/compass/go/internal/store/db" ) -// AppendMessage stores a new message under a topic in channelID, assigning the -// row id and timestamp (comms.proto:463-479). The topic is resolved inside the -// insert tx: a TopicRef.Name is get-or-created on (channel_id, lower(name)); a -// TopicRef.ID names an existing topic, validated to live under channelID. The -// blocks are serialized to JSONB and their text content extracted for the -// full-text index. When clientRequestID is non-empty, a retry with the same key -// returns the already-stored message rather than duplicating (idempotency, -// comms.proto:470-474). The returned bool reports whether a row was genuinely -// inserted: it is false on the idempotent-retry return, so the caller -// 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, -// 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. +// AppendMessage stores a message under a topic, assigning its id and timestamp. +// Topic resolution, participation validation, insert, and last_seq update run in +// one transaction, so a concurrent tree move cannot race the author check. +// Idempotent request ids suppress duplicate MessagePosted fan-out. func (s *Store) AppendMessage(ctx context.Context, m Message, channelID string, topic TopicRef, clientRequestID string) (Message, bool, error) { if channelID == "" { return Message{}, false, fmt.Errorf("%w: message channel is required", ErrInvalidArgument) @@ -50,10 +38,8 @@ 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. + // The author must participate through a member row or TREE derivation. + // Check in the insert transaction so a concurrent tree move cannot race. if err := requireChannelMember(ctx, tx, m.AuthorAccountID, ChannelID(channelID)); err != nil { return Message{}, false, err } @@ -294,22 +280,16 @@ 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 -// 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 -// 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. +// addresses a bare MessageID with no participation or authorship check. AnswerAsk +// resolves its target through the participant gate and locks it FOR UPDATE; it +// also allows a participant who is not the author to answer. These write paths +// need different predicates. // -// 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 -// 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 -// 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 -// AnswerAsk documents at :307-310). +// 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 old message after a tree move, +// and participation alone would let it rewrite another account's words. Either +// refusal and an unknown id return the same ErrNotFound to hide the message. // // The updated row is returned via RETURNING, so the caller fans out the // post-update state without a second read that could observe a later write. @@ -346,10 +326,9 @@ func (s *Store) UpdateMessageBlocksAsAuthor(ctx context.Context, actor AccountID return Message{}, err } - // 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. + // One statement keeps participation and authorship checks atomic with the + // write. A concurrent tree move or member removal lands before the UPDATE + // (which then matches no row) or after it, never between. row, err := s.q.UpdateMessageBlocksAsAuthor(ctx, db.UpdateMessageBlocksAsAuthorParams{ Blocks: blocksJSON, TextContent: textContent(blocks), @@ -358,7 +337,7 @@ func (s *Store) UpdateMessageBlocksAsAuthor(ctx context.Context, actor AccountID }) if err != nil { if noRows(err) { - // Unknown id, not the author, or no longer a member — one answer for + // Unknown id, not the author, or no longer a participant — one answer for // all three, so a refusal enumerates nothing. return Message{}, fmt.Errorf("%w: message %q", ErrNotFound, id) } @@ -414,12 +393,10 @@ func (s *Store) MessageAskIDs(ctx context.Context, id MessageID) ([]string, erro // every topic. // // 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 -// 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. +// channel_id: messages join topics on topic_id, filtered by topics.channel_id. +// SQL also requires channel participation: an explicit member row or TREE +// derivation. A non-participant sees no history, even when channel visibility +// grants access to the channel row. The store, not the RPC edge, applies this gate. func (s *Store) ListMessages(ctx context.Context, q ListMessagesQuery) ([]Message, error) { if q.ChannelID == "" { return nil, fmt.Errorf("%w: list channel is required", ErrInvalidArgument) @@ -428,10 +405,9 @@ 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 - // 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. + // Scope the cursor probe to channel participants too, so a non-participant + // naming a real message gets the same result as a fake id. The channel comes + // from the cursor message's topic. seq, err := s.q.GetPageCursorSeq(ctx, db.GetPageCursorSeqParams{ AccountID: string(q.Actor), ID: string(q.Page.BeforeMessageID), @@ -446,15 +422,9 @@ func (s *Store) ListMessages(ctx context.Context, q ListMessagesQuery) ([]Messag beforeSeq = seq } - // 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 - // 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 - // 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. + // The query snapshots participant access, while MessageUpdated carries later + // edits to the live tail. A client that deduplicates by id converges to the last + // write. Freezing content too would require a change sequence and schema change. 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), @@ -473,11 +443,9 @@ func (s *Store) ListMessages(ctx context.Context, q ListMessagesQuery) ([]Messag // SearchMessages runs a Postgres full-text search over message text, scoped to // the actor's visible set server-side (comms.proto:489-504, design.md:1137-1139). // scope optionally narrows to one channel; otherwise it searches every channel -// the actor is a member of. Results are best-match-first, clamped to the page -// 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. +// the actor participates in (member row or TREE derivation). Results are +// best-match-first and page-bounded. SQL enforces participation. Under FORCE RLS +// the search index cannot serve the @@ filter because it is not LEAKPROOF. 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) @@ -485,9 +453,9 @@ func (s *Store) SearchMessages(ctx context.Context, actor AccountID, scope Searc limit := clampLimit(page.Limit) // websearch_to_tsquery parses a human query safely — no injection, and an - // all-stopword query yields no rows. Visibility: the message's channel - // (via the topic join) must be one the actor is a member of; the optional - // scope narrows within that set. SnapshotSeq is int64-safe; see ListMessages. + // all-stopword query yields no rows. The message's channel must be in the + // actor's participant set; the optional scope narrows that set. SnapshotSeq is + // int64-safe; see ListMessages. snap := int64(page.SnapshotSeq) //nolint:gosec // G115: server-issued seq, int64 domain (see ListMessages) rows, err := s.q.SearchMessages(ctx, db.SearchMessagesParams{ AccountID: string(actor), @@ -504,15 +472,10 @@ 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 -// 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 -// single ask_answer block snapshotting the just-answered ask — inserted in the -// SAME tx so "ask answered" and "answer message exists" are one atomic fact. -// It returns both messages: the updated ask (the handler publishes -// MessageUpdated) and the answer (the handler publishes MessagePosted). +// locates the message whose blocks carry askID in the actor's participant set — +// an explicit member row or TREE derivation — before applying the answer. +// It applies the per-question answers, persists the immutable ask id, and posts +// the answer as a new message in the same transaction. // // Answering is atomic: answers must cover EXACTLY the ask's question_id set — // every question answered once, no unknown or repeated question_id. An answer @@ -541,10 +504,9 @@ func (s *Store) AnswerAsk(ctx context.Context, actor AccountID, askID string, an } defer func() { _ = tx.Rollback(ctx) }() // deferred cleanup; the Commit below is the real outcome. - // Visibility + existence in one gate: the message's channel must be one the - // actor is a member of, and its blocks must contain an ask with askID. Zero - // rows -> ErrNotFound, so ask existence cannot leak across a membership - // boundary. FOR UPDATE OF m locks the message row for the tx. + // One gate checks participation and ask existence. Zero rows map to + // ErrNotFound so ask existence cannot leak across the boundary. FOR UPDATE OF m + // locks the message row for this transaction. filter, err := askIDContainmentFilter(askID) if err != nil { return Message{}, Message{}, fmt.Errorf("store: marshal ask filter: %w", err) diff --git a/go/internal/store/queries/channels.sql b/go/internal/store/queries/channels.sql index 0b58518fc..0c7bf5f2b 100644 --- a/go/internal/store/queries/channels.sql +++ b/go/internal/store/queries/channels.sql @@ -131,7 +131,7 @@ WITH RECURSIVE stored AS ( WHERE c.id = ANY($1::text[]) AND c.membership_mode = 1 ) SELECT channel_id, account_id, subscribed FROM stored -UNION +UNION ALL SELECT p.channel_id, p.account_id, COALESCE(cs.subscribed, FALSE) AS subscribed FROM participants p LEFT JOIN channel_subscriptions cs diff --git a/go/internal/store/queries/messages.sql b/go/internal/store/queries/messages.sql index 35f73728e..22e6477be 100644 --- a/go/internal/store/queries/messages.sql +++ b/go/internal/store/queries/messages.sql @@ -47,40 +47,24 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $4 - AND EXISTS ( - SELECT 1 FROM messages m - JOIN topics t ON t.id = m.topic_id - JOIN channels c ON c.id = t.channel_id - WHERE m.id = $3 AND c.membership_mode = 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 +), visible 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 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $4 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) +JOIN visible v ON v.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; @@ -93,68 +77,52 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $1 - AND EXISTS (SELECT 1 FROM channels WHERE id = $3 AND membership_mode = 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 +), visible 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 -WHERE m.id = $2 AND t.channel_id = $3 - AND ( - EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ); +JOIN visible v ON v.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 - AND EXISTS (SELECT 1 FROM channels WHERE id = $2 AND membership_mode = 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 +), visible 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 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 visible v ON v.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) - AND ( - EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) ORDER BY m.seq DESC LIMIT $4; -- name: SearchMessages :many @@ -162,40 +130,29 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $1 - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.membership_mode = 1 AND ($3 = '' OR c.id = $3) - ) UNION SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id +), visible 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 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 visible v ON v.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) - AND ( - EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) ORDER BY ts_rank(m.search_tsv, websearch_to_tsquery('english', $2)) DESC, m.seq DESC LIMIT $4; @@ -204,35 +161,27 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $1 - AND EXISTS (SELECT 1 FROM channels WHERE membership_mode = 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 +), visible 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 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 visible v ON v.channel_id = t.channel_id WHERE m.blocks @> $2::jsonb - AND ( - EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) FOR UPDATE OF m; -- name: GetMessageByRequestID :many diff --git a/go/internal/store/queries/topics.sql b/go/internal/store/queries/topics.sql index dea59cfe8..e38b14c9a 100644 --- a/go/internal/store/queries/topics.sql +++ b/go/internal/store/queries/topics.sql @@ -16,35 +16,23 @@ WITH RECURSIVE chain AS ( SELECT aa.account_id, aa.parent_agent_id FROM agent_accounts aa WHERE aa.account_id = $1 - AND EXISTS ( - SELECT 1 FROM topics t - JOIN channels c ON c.id = t.channel_id - WHERE t.id = $2 AND c.membership_mode = 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 +), visible 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 visible v ON v.channel_id = t.channel_id WHERE t.id = $2 - AND ( - EXISTS ( - SELECT 1 FROM channel_members cm - WHERE cm.channel_id = t.channel_id AND cm.account_id = $1 - ) - OR ( - EXISTS (SELECT 1 FROM channels WHERE id = t.channel_id AND membership_mode = 1) - AND EXISTS ( - SELECT 1 FROM channels c - WHERE c.id = t.channel_id AND ( - c.parent_agent_id IN (SELECT ch.account_id FROM chain ch) - OR $1 = (SELECT aa.owner_user_id FROM agent_accounts aa - WHERE aa.account_id = c.parent_agent_id) - ) - ) - ) - ) FOR UPDATE OF t; -- name: SetTopicArchived :exec diff --git a/go/internal/store/topics.go b/go/internal/store/topics.go index 68e5adc5e..7b4f474bc 100644 --- a/go/internal/store/topics.go +++ b/go/internal/store/topics.go @@ -10,13 +10,10 @@ import ( "github.com/RigelBuild/compass/go/internal/store/db" ) -// 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 — -// 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". +// ListTopics returns a channel's topics newest-activity-first. It checks caller +// participation through a member row or TREE derivation. Archived topics are +// omitted unless requested. Unauthorized and unknown channels map to ErrNotFound +// so existence cannot leak through this read. func (s *Store) ListTopics(ctx context.Context, callerAccountID, channelID string, includeArchived bool) ([]Topic, error) { if channelID == "" { return nil, fmt.Errorf("%w: list topics channel is required", ErrInvalidArgument) @@ -38,21 +35,14 @@ func (s *Store) ListTopics(ctx context.Context, callerAccountID, channelID strin return topicsFromRows(rows), nil } -// UpdateTopic renames and/or archives a topic under an acting account, or — -// when a rename collides with an existing topic name in the same channel — -// MERGES the two: the source topic's messages are re-pointed at the target -// (message rows carry the target's topic_id), the target's last_seq absorbs the -// source's, and the emptied source row is deleted, all in one transaction. The -// surviving topic is returned. +// UpdateTopic renames and/or archives a topic under the caller's participation +// gate. A name collision moves source messages into the target and deletes the +// source in one transaction; the surviving topic is returned. // -// The caller must be a member of 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 -// are each optional (nil = leave unchanged); the archived flag is applied to -// the SURVIVING topic (the target on a merge, the source otherwise). -// -// A rename whose lowercased name matches the topic's own current name is a -// harmless in-place rename (no merge — the collision check excludes self). +// The caller must participate in the channel through a member row or TREE +// derivation. A non-participant or unknown topic maps to ErrNotFound. Name and +// archived are optional. Renaming to the current case-folded name updates the +// source in place. func (s *Store) UpdateTopic(ctx context.Context, callerAccountID, topicID string, name *string, archived *bool) (Topic, error) { if topicID == "" { return Topic{}, fmt.Errorf("%w: topic id is required", ErrInvalidArgument) @@ -64,10 +54,9 @@ 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) -> - // ErrNotFound. FOR UPDATE OF t locks the source topic row for the tx so a - // concurrent rename/merge serializes. + // Resolve the topic and channel-participant gate in one statement. An explicit + // member row or TREE derivation grants participation; no row (unknown topic or + // non-participant) maps to ErrNotFound. FOR UPDATE locks the topic for this tx. q := db.New(tx) channelID, err := q.ResolveTopicForUpdate(ctx, db.ResolveTopicForUpdateParams{ AccountID: callerAccountID, From 74a06376df629978869c3d2bfd1b9923a226226c Mon Sep 17 00:00:00 2001 From: mintaka Date: Tue, 6 Oct 2026 04:20:10 -0400 Subject: [PATCH 3/5] fix(store): keep name ambiguity for non-participants; restore message doc contracts (RIG-4188) A name that matches several visible channels narrows to the viewer's participant channels only when that set is non-empty, so a non-participant sibling gets the same ambiguity its channel list shows, not a not-found. Message and topic doc comments go back to their full contracts, reworded from membership to participation. Co-authored-by: Matt Wilkinson --- .../store/channel_by_name_pgtest_test.go | 5 + go/internal/store/channels.go | 13 +- go/internal/store/messages.go | 121 ++++++++++++------ go/internal/store/topics.go | 41 +++--- 4 files changed, 120 insertions(+), 60 deletions(-) diff --git a/go/internal/store/channel_by_name_pgtest_test.go b/go/internal/store/channel_by_name_pgtest_test.go index 42256bc30..7cd5c6fc6 100644 --- a/go/internal/store/channel_by_name_pgtest_test.go +++ b/go/internal/store/channel_by_name_pgtest_test.go @@ -138,4 +138,9 @@ func TestChannelByNameForViewerNarrowsOwnerSetDuplicatesToParticipants(t *testin } _, 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/channels.go b/go/internal/store/channels.go index ed1bd4bae..e98b612ef 100644 --- a/go/internal/store/channels.go +++ b/go/internal/store/channels.go @@ -521,7 +521,9 @@ func resolveGroupRef(groups []ChannelGroup, ref string) (ChannelGroup, error) { // 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). // @@ -543,7 +545,8 @@ func (s *Store) ChannelByNameForViewer(ctx context.Context, viewer AccountID, na return Channel{}, err } if len(channels) > 1 { - participants := channels[:0] + // 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 { @@ -553,7 +556,9 @@ func (s *Store) ChannelByNameForViewer(ctx context.Context, viewer AccountID, na participants = append(participants, channel) } } - channels = participants + if len(participants) > 0 { + channels = participants + } } switch len(channels) { case 0: @@ -561,7 +566,7 @@ func (s *Store) ChannelByNameForViewer(ctx context.Context, viewer AccountID, na case 1: return channels[0], nil default: - return Channel{}, fmt.Errorf("%w: channel name %q is ambiguous — it names %d channels the viewer participates in; 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)) } } diff --git a/go/internal/store/messages.go b/go/internal/store/messages.go index b12675579..1e69cf2b3 100644 --- a/go/internal/store/messages.go +++ b/go/internal/store/messages.go @@ -12,10 +12,22 @@ import ( "github.com/RigelBuild/compass/go/internal/store/db" ) -// AppendMessage stores a message under a topic, assigning its id and timestamp. -// Topic resolution, participation validation, insert, and last_seq update run in -// one transaction, so a concurrent tree move cannot race the author check. -// Idempotent request ids suppress duplicate MessagePosted fan-out. +// AppendMessage stores a new message under a topic in channelID, assigning the +// row id and timestamp (comms.proto:463-479). The topic is resolved inside the +// insert tx: a TopicRef.Name is get-or-created on (channel_id, lower(name)); a +// TopicRef.ID names an existing topic, validated to live under channelID. The +// blocks are serialized to JSONB and their text content extracted for the +// full-text index. When clientRequestID is non-empty, a retry with the same key +// returns the already-stored message rather than duplicating (idempotency, +// comms.proto:470-474). The returned bool reports whether a row was genuinely +// inserted: it is false on the idempotent-retry return, so the caller +// 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 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. func (s *Store) AppendMessage(ctx context.Context, m Message, channelID string, topic TopicRef, clientRequestID string) (Message, bool, error) { if channelID == "" { return Message{}, false, fmt.Errorf("%w: message channel is required", ErrInvalidArgument) @@ -38,8 +50,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. - // The author must participate through a member row or TREE derivation. - // Check in the insert transaction so a concurrent tree move cannot race. + // D9 write-authz: the author must participate in the target channel (member + // row or TREE derivation), so a non-participant cannot persist into a channel + // it can't read. A non-participant gets ErrNotFound, never a hint the channel exists. Checked in + // the insert tx so a concurrent removal cannot race the gate. if err := requireChannelMember(ctx, tx, m.AuthorAccountID, ChannelID(channelID)); err != nil { return Message{}, false, err } @@ -280,16 +294,22 @@ 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 a bare MessageID with no participation or authorship check. AnswerAsk -// resolves its target through the participant gate and locks it FOR UPDATE; it -// also allows a participant who is not the author to answer. These write paths -// need different predicates. +// 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 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 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 old message after a tree move, -// and participation alone would let it rewrite another account's words. Either -// refusal and an unknown id return the same ErrNotFound to hide the message. +// 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 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 +// AnswerAsk documents at :307-310). // // The updated row is returned via RETURNING, so the caller fans out the // post-update state without a second read that could observe a later write. @@ -326,9 +346,11 @@ func (s *Store) UpdateMessageBlocksAsAuthor(ctx context.Context, actor AccountID return Message{}, err } - // One statement keeps participation and authorship checks atomic with the - // write. A concurrent tree move or member removal lands before the UPDATE - // (which then matches no row) or after it, never between. + // One statement, so the authz predicate and the write cannot race: a + // concurrent membership revocation or tree move lands before the UPDATE + // (matches no row) or after it, never between. The visible join is the + // participation half and + // the author_account_id equality the authorship half; both must hold. row, err := s.q.UpdateMessageBlocksAsAuthor(ctx, db.UpdateMessageBlocksAsAuthorParams{ Blocks: blocksJSON, TextContent: textContent(blocks), @@ -337,7 +359,7 @@ func (s *Store) UpdateMessageBlocksAsAuthor(ctx context.Context, actor AccountID }) if err != nil { if noRows(err) { - // Unknown id, not the author, or no longer a participant — one answer for + // Unknown id, not the author, or no longer a member — one answer for // all three, so a refusal enumerates nothing. return Message{}, fmt.Errorf("%w: message %q", ErrNotFound, id) } @@ -393,10 +415,12 @@ func (s *Store) MessageAskIDs(ctx context.Context, id MessageID) ([]string, erro // every topic. // // 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. -// SQL also requires channel participation: an explicit member row or TREE -// derivation. A non-participant sees no history, even when channel visibility -// grants access to the channel row. The store, not the RPC edge, applies this gate. +// 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 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. func (s *Store) ListMessages(ctx context.Context, q ListMessagesQuery) ([]Message, error) { if q.ChannelID == "" { return nil, fmt.Errorf("%w: list channel is required", ErrInvalidArgument) @@ -405,9 +429,10 @@ func (s *Store) ListMessages(ctx context.Context, q ListMessagesQuery) ([]Messag var beforeSeq int64 if q.Page.BeforeMessageID != "" { - // Scope the cursor probe to channel participants too, so a non-participant - // naming a real message gets the same result as a fake id. The channel comes - // from the cursor message's topic. + // 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. seq, err := s.q.GetPageCursorSeq(ctx, db.GetPageCursorSeqParams{ AccountID: string(q.Actor), ID: string(q.Page.BeforeMessageID), @@ -422,9 +447,15 @@ func (s *Store) ListMessages(ctx context.Context, q ListMessagesQuery) ([]Messag beforeSeq = seq } - // The query snapshots participant access, while MessageUpdated carries later - // edits to the live tail. A client that deduplicates by id converges to the last - // write. Freezing content too would require a change sequence and schema change. + // A zero beforeSeq reads the newest page; a positive one pages strictly + // 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. + + // 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; 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), @@ -443,9 +474,11 @@ func (s *Store) ListMessages(ctx context.Context, q ListMessagesQuery) ([]Messag // SearchMessages runs a Postgres full-text search over message text, scoped to // the actor's visible set server-side (comms.proto:489-504, design.md:1137-1139). // scope optionally narrows to one channel; otherwise it searches every channel -// the actor participates in (member row or TREE derivation). Results are -// best-match-first and page-bounded. SQL enforces participation. Under FORCE RLS -// the search index cannot serve the @@ filter because it is not LEAKPROOF. +// the actor is a member of. Results are best-match-first, clamped to the page +// 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 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) @@ -453,9 +486,9 @@ func (s *Store) SearchMessages(ctx context.Context, actor AccountID, scope Searc limit := clampLimit(page.Limit) // websearch_to_tsquery parses a human query safely — no injection, and an - // all-stopword query yields no rows. The message's channel must be in the - // actor's participant set; the optional scope narrows that set. SnapshotSeq is - // int64-safe; see ListMessages. + // all-stopword query yields no rows. Visibility: the message's channel + // (via the topic join) must be one the actor is a member of; the optional + // scope narrows within that set. SnapshotSeq is int64-safe; see ListMessages. snap := int64(page.SnapshotSeq) //nolint:gosec // G115: server-issued seq, int64 domain (see ListMessages) rows, err := s.q.SearchMessages(ctx, db.SearchMessagesParams{ AccountID: string(actor), @@ -472,10 +505,15 @@ 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 askID in the actor's participant set — -// an explicit member row or TREE derivation — before applying the answer. -// It applies the per-question answers, persists the immutable ask id, and posts -// the answer as a new message in the same transaction. +// locates the message whose blocks carry an ask with askID within the actor's +// 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 +// single ask_answer block snapshotting the just-answered ask — inserted in the +// SAME tx so "ask answered" and "answer message exists" are one atomic fact. +// It returns both messages: the updated ask (the handler publishes +// MessageUpdated) and the answer (the handler publishes MessagePosted). // // Answering is atomic: answers must cover EXACTLY the ask's question_id set — // every question answered once, no unknown or repeated question_id. An answer @@ -504,9 +542,10 @@ func (s *Store) AnswerAsk(ctx context.Context, actor AccountID, askID string, an } defer func() { _ = tx.Rollback(ctx) }() // deferred cleanup; the Commit below is the real outcome. - // One gate checks participation and ask existence. Zero rows map to - // ErrNotFound so ask existence cannot leak across the boundary. FOR UPDATE OF m - // locks the message row for this transaction. + // Visibility + existence in one gate: the message's channel must be one the + // actor is a member of, and its blocks must contain an ask with askID. Zero + // rows -> ErrNotFound, so ask existence cannot leak across a membership + // boundary. FOR UPDATE OF m locks the message row for the tx. filter, err := askIDContainmentFilter(askID) if err != nil { return Message{}, Message{}, fmt.Errorf("store: marshal ask filter: %w", err) diff --git a/go/internal/store/topics.go b/go/internal/store/topics.go index 7b4f474bc..4d8b12bbb 100644 --- a/go/internal/store/topics.go +++ b/go/internal/store/topics.go @@ -10,10 +10,13 @@ import ( "github.com/RigelBuild/compass/go/internal/store/db" ) -// ListTopics returns a channel's topics newest-activity-first. It checks caller -// participation through a member row or TREE derivation. Archived topics are -// omitted unless requested. Unauthorized and unknown channels map to ErrNotFound -// so existence cannot leak through this read. +// 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 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". func (s *Store) ListTopics(ctx context.Context, callerAccountID, channelID string, includeArchived bool) ([]Topic, error) { if channelID == "" { return nil, fmt.Errorf("%w: list topics channel is required", ErrInvalidArgument) @@ -23,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) } @@ -35,14 +38,21 @@ func (s *Store) ListTopics(ctx context.Context, callerAccountID, channelID strin return topicsFromRows(rows), nil } -// UpdateTopic renames and/or archives a topic under the caller's participation -// gate. A name collision moves source messages into the target and deletes the -// source in one transaction; the surviving topic is returned. +// UpdateTopic renames and/or archives a topic under an acting account, or — +// when a rename collides with an existing topic name in the same channel — +// MERGES the two: the source topic's messages are re-pointed at the target +// (message rows carry the target's topic_id), the target's last_seq absorbs the +// source's, and the emptied source row is deleted, all in one transaction. The +// surviving topic is returned. // -// The caller must participate in the channel through a member row or TREE -// derivation. A non-participant or unknown topic maps to ErrNotFound. Name and -// archived are optional. Renaming to the current case-folded name updates the -// source in place. +// 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 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). +// +// A rename whose lowercased name matches the topic's own current name is a +// harmless in-place rename (no merge — the collision check excludes self). func (s *Store) UpdateTopic(ctx context.Context, callerAccountID, topicID string, name *string, archived *bool) (Topic, error) { if topicID == "" { return Topic{}, fmt.Errorf("%w: topic id is required", ErrInvalidArgument) @@ -54,9 +64,10 @@ 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 and channel-participant gate in one statement. An explicit - // member row or TREE derivation grants participation; no row (unknown topic or - // non-participant) maps to ErrNotFound. FOR UPDATE locks the topic for this tx. + // 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) channelID, err := q.ResolveTopicForUpdate(ctx, db.ResolveTopicForUpdateParams{ AccountID: callerAccountID, From e6c5539eed96d3abea74894bc144cc58f5e913e8 Mon Sep 17 00:00:00 2001 From: mintaka Date: Tue, 6 Oct 2026 04:38:50 -0400 Subject: [PATCH 4/5] test(store): pin TREE ask and unscoped-search participation; name the participant CTE (RIG-4188) Restores the subtree-agent FindAskMessage/AnswerAsk positives and adds unscoped search and owner ListMessages cases. The per-query participant CTE is renamed participating, so it is not mistaken for channel visibility, and both query files say the copies must stay equal to ChannelParticipant. Co-authored-by: Matt Wilkinson --- .../store/channel_participant_pgtest_test.go | 19 ++++++++++++++ go/internal/store/db/messages.sql.go | 25 +++++++++++-------- go/internal/store/db/querier.go | 8 ++++++ go/internal/store/db/topics.sql.go | 9 +++++-- go/internal/store/messages.go | 11 ++++---- go/internal/store/queries/messages.sql | 25 +++++++++++-------- go/internal/store/queries/topics.sql | 9 +++++-- 7 files changed, 76 insertions(+), 30 deletions(-) diff --git a/go/internal/store/channel_participant_pgtest_test.go b/go/internal/store/channel_participant_pgtest_test.go index 73e6ab3c5..161b8c1cc 100644 --- a/go/internal/store/channel_participant_pgtest_test.go +++ b/go/internal/store/channel_participant_pgtest_test.go @@ -159,6 +159,14 @@ func TestChannelTreeParticipantReadsAndWrites(t *testing.T) { 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), @@ -204,6 +212,10 @@ func TestChannelTreeOwnerSetSiblingReadWriteDenials(t *testing.T) { 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), }) @@ -239,6 +251,13 @@ func TestChannelTreeOwnerSetSiblingReadWriteDenials(t *testing.T) { } _, _, 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) { diff --git a/go/internal/store/db/messages.sql.go b/go/internal/store/db/messages.sql.go index 036651025..b604cec48 100644 --- a/go/internal/store/db/messages.sql.go +++ b/go/internal/store/db/messages.sql.go @@ -18,7 +18,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 UNION SELECT c.id FROM channels c @@ -33,7 +33,7 @@ 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 visible v ON v.channel_id = t.channel_id +JOIN participating p ON p.channel_id = t.channel_id WHERE m.blocks @> $2::jsonb FOR UPDATE OF m ` @@ -81,6 +81,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 ` @@ -99,6 +100,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 projection so the Go maps each row through messageFromParts. +// Participant-channel copies: every `chain` + `participating` CTE in this +// file and topics.sql MUST stay identical and equal to ChannelParticipant +// (authz.sql). It is participation, not channel visibility: never widen it to +// the owner-set visibility predicate. A future ACL conjunct goes in each copy. func (q *Queries) GetChannelPostPolicy(ctx context.Context, id string) (GetChannelPostPolicyRow, error) { row := q.db.QueryRow(ctx, getChannelPostPolicy, id) var i GetChannelPostPolicyRow @@ -175,7 +180,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 UNION SELECT c.id FROM channels c @@ -187,7 +192,7 @@ WITH RECURSIVE chain AS ( ) SELECT m.seq FROM messages m JOIN topics t ON t.id = m.topic_id -JOIN visible v ON v.channel_id = t.channel_id +JOIN participating p ON p.channel_id = t.channel_id WHERE m.id = $2 AND t.channel_id = $3 ` @@ -316,7 +321,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 UNION SELECT c.id FROM channels c @@ -331,7 +336,7 @@ 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 visible v ON v.channel_id = t.channel_id +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 @@ -419,7 +424,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 UNION SELECT c.id FROM channels c @@ -434,7 +439,7 @@ 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 visible v ON v.channel_id = t.channel_id +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) @@ -521,7 +526,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $4 UNION SELECT c.id FROM channels c @@ -533,7 +538,7 @@ WITH RECURSIVE chain AS ( UPDATE messages m SET blocks = $1, text_content = $2 FROM topics t -JOIN visible v ON v.channel_id = t.channel_id +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, diff --git a/go/internal/store/db/querier.go b/go/internal/store/db/querier.go index 97b05874c..8ae410823 100644 --- a/go/internal/store/db/querier.go +++ b/go/internal/store/db/querier.go @@ -189,6 +189,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 projection so the Go maps each row through messageFromParts. + // Participant-channel copies: every `chain` + `participating` CTE in this + // file and topics.sql MUST stay identical and equal to ChannelParticipant + // (authz.sql). It is participation, not channel visibility: never widen it to + // the owner-set visibility predicate. A future ACL conjunct goes in each copy. 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 @@ -364,6 +368,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 topics.sql MUST stay identical and equal to ChannelParticipant + // (authz.sql). It is participation, not channel visibility: never widen it to + // the owner-set visibility predicate. A future ACL conjunct goes in each copy. ListTopics(ctx context.Context, arg ListTopicsParams) ([]Topic, error) ListVisibleAccounts(ctx context.Context, id string) ([]ListVisibleAccountsRow, error) LoadDeliveryCursor(ctx context.Context, arg LoadDeliveryCursorParams) (LoadDeliveryCursorRow, error) diff --git a/go/internal/store/db/topics.sql.go b/go/internal/store/db/topics.sql.go index a69aae239..83dd4d4d9 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 topics.sql MUST stay identical and equal to ChannelParticipant +// (authz.sql). It is participation, not channel visibility: never widen it to +// the owner-set visibility predicate. A future ACL conjunct goes in each copy. 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 { @@ -139,7 +144,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 UNION SELECT c.id FROM channels c @@ -150,7 +155,7 @@ WITH RECURSIVE chain AS ( ) ) SELECT t.channel_id FROM topics t -JOIN visible v ON v.channel_id = t.channel_id +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 1e69cf2b3..3d7e9b894 100644 --- a/go/internal/store/messages.go +++ b/go/internal/store/messages.go @@ -51,9 +51,9 @@ 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 participate in the target channel (member - // row or TREE derivation), so a non-participant cannot persist into a channel - // it can't read. A non-participant gets ErrNotFound, never a hint the channel exists. Checked in - // the insert tx so a concurrent removal cannot race the gate. + // 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 } @@ -348,9 +348,8 @@ func (s *Store) UpdateMessageBlocksAsAuthor(ctx context.Context, actor AccountID // One statement, so the authz predicate and the write cannot race: a // concurrent membership revocation or tree move lands before the UPDATE - // (matches no row) or after it, never between. The visible join is the - // participation half and - // the author_account_id equality the authorship half; both must hold. + // (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), diff --git a/go/internal/store/queries/messages.sql b/go/internal/store/queries/messages.sql index 22e6477be..1e47ebe42 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 projection so the Go maps each row through messageFromParts. +-- Participant-channel copies: every `chain` + `participating` CTE in this +-- file and topics.sql MUST stay identical and equal to ChannelParticipant +-- (authz.sql). It is participation, not channel visibility: never widen it to +-- the owner-set visibility predicate. A future ACL conjunct goes in each copy. + -- name: GetChannelPostPolicy :one SELECT post_policy, COALESCE(owner_account_id, '') AS owner_account_id, name FROM channels WHERE id = $1; @@ -51,7 +56,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $4 UNION SELECT c.id FROM channels c @@ -63,7 +68,7 @@ WITH RECURSIVE chain AS ( UPDATE messages m SET blocks = $1, text_content = $2 FROM topics t -JOIN visible v ON v.channel_id = t.channel_id +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, @@ -81,7 +86,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 UNION SELECT c.id FROM channels c @@ -93,7 +98,7 @@ WITH RECURSIVE chain AS ( ) SELECT m.seq FROM messages m JOIN topics t ON t.id = m.topic_id -JOIN visible v ON v.channel_id = t.channel_id +JOIN participating p ON p.channel_id = t.channel_id WHERE m.id = $2 AND t.channel_id = $3; -- name: ListMessages :many @@ -105,7 +110,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 UNION SELECT c.id FROM channels c @@ -120,7 +125,7 @@ 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 visible v ON v.channel_id = t.channel_id +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 @@ -134,7 +139,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 UNION SELECT c.id FROM channels c @@ -149,7 +154,7 @@ 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 visible v ON v.channel_id = t.channel_id +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) @@ -165,7 +170,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 UNION SELECT c.id FROM channels c @@ -180,7 +185,7 @@ 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 visible v ON v.channel_id = t.channel_id +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 e38b14c9a..d62f55419 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 topics.sql MUST stay identical and equal to ChannelParticipant +-- (authz.sql). It is participation, not channel visibility: never widen it to +-- the owner-set visibility predicate. A future ACL conjunct goes in each copy. + -- name: ListTopics :many SELECT id, channel_id, name, created_by_account_id, created_at_unix_ms, archived, last_seq, tenant_id FROM topics @@ -20,7 +25,7 @@ WITH RECURSIVE chain AS ( SELECT a.account_id, a.parent_agent_id FROM agent_accounts a JOIN chain ch ON a.account_id = ch.parent_agent_id -), visible AS ( +), participating AS ( SELECT cm.channel_id FROM channel_members cm WHERE cm.account_id = $1 UNION SELECT c.id FROM channels c @@ -31,7 +36,7 @@ WITH RECURSIVE chain AS ( ) ) SELECT t.channel_id FROM topics t -JOIN visible v ON v.channel_id = t.channel_id +JOIN participating p ON p.channel_id = t.channel_id WHERE t.id = $2 FOR UPDATE OF t; From a265f919597ba5f65646a9b642777280ca8d9f3b Mon Sep 17 00:00:00 2001 From: mintaka Date: Tue, 6 Oct 2026 04:52:39 -0400 Subject: [PATCH 5/5] docs(store): correct the participant-copy header pointers (RIG-4188) Co-authored-by: Matt Wilkinson --- go/internal/store/db/messages.sql.go | 6 +++--- go/internal/store/db/querier.go | 12 ++++++------ go/internal/store/db/topics.sql.go | 6 +++--- go/internal/store/queries/messages.sql | 6 +++--- go/internal/store/queries/topics.sql | 6 +++--- 5 files changed, 18 insertions(+), 18 deletions(-) diff --git a/go/internal/store/db/messages.sql.go b/go/internal/store/db/messages.sql.go index b604cec48..4d1bb52e3 100644 --- a/go/internal/store/db/messages.sql.go +++ b/go/internal/store/db/messages.sql.go @@ -101,9 +101,9 @@ type GetChannelPostPolicyRow struct { // calls. Every message read shares the id/topic_id/author_account_id/author_handle/ // at_unix_ms/blocks projection so the Go maps each row through messageFromParts. // Participant-channel copies: every `chain` + `participating` CTE in this -// file and topics.sql MUST stay identical and equal to ChannelParticipant -// (authz.sql). It is participation, not channel visibility: never widen it to -// the owner-set visibility predicate. A future ACL conjunct goes in each copy. +// 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 diff --git a/go/internal/store/db/querier.go b/go/internal/store/db/querier.go index 8ae410823..72ffc5755 100644 --- a/go/internal/store/db/querier.go +++ b/go/internal/store/db/querier.go @@ -190,9 +190,9 @@ type Querier interface { // calls. Every message read shares the id/topic_id/author_account_id/author_handle/ // at_unix_ms/blocks projection so the Go maps each row through messageFromParts. // Participant-channel copies: every `chain` + `participating` CTE in this - // file and topics.sql MUST stay identical and equal to ChannelParticipant - // (authz.sql). It is participation, not channel visibility: never widen it to - // the owner-set visibility predicate. A future ACL conjunct goes in each copy. + // 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 @@ -369,9 +369,9 @@ type Querier interface { // (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 topics.sql MUST stay identical and equal to ChannelParticipant - // (authz.sql). It is participation, not channel visibility: never widen it to - // the owner-set visibility predicate. A future ACL conjunct goes in each copy. + // 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) ListVisibleAccounts(ctx context.Context, id string) ([]ListVisibleAccountsRow, error) LoadDeliveryCursor(ctx context.Context, arg LoadDeliveryCursorParams) (LoadDeliveryCursorRow, error) diff --git a/go/internal/store/db/topics.sql.go b/go/internal/store/db/topics.sql.go index 83dd4d4d9..832cd4995 100644 --- a/go/internal/store/db/topics.sql.go +++ b/go/internal/store/db/topics.sql.go @@ -60,9 +60,9 @@ type ListTopicsParams struct { // (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 topics.sql MUST stay identical and equal to ChannelParticipant -// (authz.sql). It is participation, not channel visibility: never widen it to -// the owner-set visibility predicate. A future ACL conjunct goes in each copy. +// 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 { diff --git a/go/internal/store/queries/messages.sql b/go/internal/store/queries/messages.sql index 1e47ebe42..387193e4b 100644 --- a/go/internal/store/queries/messages.sql +++ b/go/internal/store/queries/messages.sql @@ -8,9 +8,9 @@ -- at_unix_ms/blocks projection so the Go maps each row through messageFromParts. -- Participant-channel copies: every `chain` + `participating` CTE in this --- file and topics.sql MUST stay identical and equal to ChannelParticipant --- (authz.sql). It is participation, not channel visibility: never widen it to --- the owner-set visibility predicate. A future ACL conjunct goes in each copy. +-- 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 diff --git a/go/internal/store/queries/topics.sql b/go/internal/store/queries/topics.sql index d62f55419..f33484638 100644 --- a/go/internal/store/queries/topics.sql +++ b/go/internal/store/queries/topics.sql @@ -6,9 +6,9 @@ -- 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 topics.sql MUST stay identical and equal to ChannelParticipant --- (authz.sql). It is participation, not channel visibility: never widen it to --- the owner-set visibility predicate. A future ACL conjunct goes in each copy. +-- 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