diff --git a/docs/channel-module.md b/docs/channel-module.md index 4c3e3dc6..2f6ce2ea 100644 --- a/docs/channel-module.md +++ b/docs/channel-module.md @@ -378,7 +378,7 @@ Redis miss 恢复来源: - 阶段外高级发送/转发 flags 不再返回 `NOT_IMPLEMENTED`:quick reply、effect、paid、suggested post、monoforum/todo/poll/story reply 分别映射到 `SHORTCUT_INVALID`、`EFFECT_ID_INVALID`、`PAYMENT_UNSUPPORTED`/`STARS_AMOUNT_INVALID`、`SUGGESTED_POST_PEER_INVALID`、`REPLY_TO_MONOFORUM_PEER_INVALID`、`REPLY_MESSAGE_ID_INVALID`、`POLL_OPTION_INVALID`、`STORY_ID_INVALID`;后续接入对应模型前不伪造成功 update。 - 已实现 linked discussion/comment 基础闭环:broadcast post 写入单份 source message 时同步在 linked megagroup 写一条 forwarded root message;source message 保存 discussion ref 并在 `messageReplies` 填 `comments/channel_id/replies_pts/max_id/read_max_id`;`messages.getDiscussionMessage` 返回 root、`messages.getReplies` 读 linked group thread,`messages.readDiscussion` 推进 linked group read watermark。 - 已实现文本转发的 channel 路径:channel→channel、channel→user、user→channel;目的 channel 仍按单份消息写入并生成 channel pts;请求携带的目标会话 `reply_to` 会正常校验和持久化,TDesktop 转发到 forum topic 时发送的 `top_msg_id` 会映射为 topic-only reply 并更新 topic top message,但源消息自身 reply 不继承;RPC 响应、账号级 difference/outbox 与 channel difference 都会为 `fwd_from/reply_to` 中可解析的 user/channel peer 补齐 users/chats,避免 TDesktop 因 header peer 未加载而延迟 apply。 -- 已实现 `messages.search(InputPeerChannel)`:支持 channel/supergroup 单份消息文本搜索、`from_id` 用户过滤、`min_date/max_date`、`offset_id/max_id/min_id` 与 limit cap,结果附带当前 channel 与消息 sender/forward/reply 所需 users/chats。 +- 已实现 `messages.search(InputPeerChannel)`:支持 channel/supergroup 单份消息文本搜索、`from_id` 用户过滤、`min_date/max_date`、`offset_id/max_id/min_id` 与 limit cap,结果附带当前 channel 与消息 sender/forward/reply 所需 users/chats;`inputMessagesFilterPinned` 按当前单置顶模型只返回 `channels.pinned_message_id` 对应消息,并只给该消息设置 TL `message.pinned`。 - 已实现 `messages.searchGlobal` 的 channel/supergroup 分支:TDesktop 主搜索的 Channels/Groups tab 会按当前账号 active membership 搜索单份 channel message,支持 `broadcasts_only/groups_only/users_only`、`offset_rate+offset_peer+offset_id` seek、folder_id=0/1 下推与 limit cap=50;未加入/left/kicked/view_messages banned 的频道不暴露,media filter 在 media store 接入前返回空结果。 - 已实现 `messages.getMessageReadParticipants(InputPeerChannel)`:基于小 megagroup 成员读水位返回 `readParticipantDate`,含 50 人阈值、7 天过期窗口、`available_min_id` 可见性过滤与 PG 索引。 - 已实现 `messages.getMessageEditData`:私聊与 channel peer 都做 message/peer/作者或管理员编辑权限校验;当前文本-only 编辑没有媒体 caption 状态,返回 `caption=false`。 diff --git a/docs/compatibility-matrix.md b/docs/compatibility-matrix.md index 13f316ff..cde777bf 100644 --- a/docs/compatibility-matrix.md +++ b/docs/compatibility-matrix.md @@ -215,7 +215,7 @@ outbox 多 worker 并发 + 发送事务乱序提交 → **主动推送可能乱 | messages.editMessage | done | real-partial | 当前阶段支持私聊文本消息编辑:仅原发送者可编辑自己的 outgoing message,更新 shared private_messages 与所有可见 owner message_boxes,生成 updateEditMessage(pts_count=1) 并可靠投递;`inputMediaWebPage/inputMediaEmpty` 降级为文本编辑,真实 media 返回 `MEDIA_INVALID`,reply_markup 返回 `REPLY_MARKUP_INVALID`,quick replies 返回 `MESSAGE_ID_INVALID`,scheduled 返回 `SCHEDULE_DATE_INVALID` | | messages.deleteMessages | done | real-partial | 当前 owner 视角软删除指定私聊消息;`revoke` 会删除同一 private_message 在其它 owner 视角的 message_box;按 owner 生成 `updateDeleteMessages`,`pts_count=len(message_ids)`,并重算或删除 dialog | | messages.deleteHistory | done | real-partial | 当前 owner 视角按 peer/max_id 清空历史;默认清空后无可见消息则删除 dialog,后续新消息可重建;`just_clear` 保留空 dialog;`revoke` 同步清理对端 owner 视角;min_date/max_date 第一阶段兼容 no-op | -| messages.search | done | real | 当前账号消息搜索;user peer 走 `message_boxes`,channel/supergroup peer 走单份 `channel_messages`;支持 q/from_id/min_date/max_date/offset_id/add_offset/limit/max_id/min_id/hash,limit cap,文本命中有 pg_trgm GIN 索引兜底;channel 结果附带 sender/fwd/reply 所需 users/chats | +| messages.search | done | real | 当前账号消息搜索;user peer 走 `message_boxes`,channel/supergroup peer 走单份 `channel_messages`;支持 q/from_id/min_date/max_date/offset_id/add_offset/limit/max_id/min_id/hash,limit cap,文本命中有 pg_trgm GIN 索引兜底;`inputMessagesFilterPinned` 只返回当前 `channels.pinned_message_id` 对应消息并设置单条 `message.pinned`,避免 Android pinned list 把普通历史全标置顶;channel 结果附带 sender/fwd/reply 所需 users/chats | | messages.searchGlobal | done | real-partial | TDesktop 搜索框全局消息分支;当前查当前 owner 私聊 message_boxes,并合并当前账号 active membership 的频道/超级群单份文本消息;支持 users_only/groups_only/broadcasts_only、filterEmpty、q、limit cap=50、offset_rate+offset_peer+offset_id channel seek、min/max_date 与 folder_id=0/1 下推,media filters 显式空结果;参考实现 的 joined channel list 语义与 参考实现 的 bounded search,后续再接外部全文搜索/媒体索引 | | messages.toggleDialogPin | done | real | 当前 owner 的 dialog pinned/pinned_order 持久化;user peer 写 dialogs,channel peer 写 channel_dialogs;写 durable updateDialogPinned + dispatch_outbox,可靠投递给其它在线 session | | messages.reorderPinnedDialogs | done | real | 当前 owner 的 pinned dialog 顺序持久化;支持混合 user/channel peer 与 force 清理未在 order 中的 pinned dialog,写 durable updatePinnedDialogs(order) + dispatch_outbox,可靠投递给其它在线 session | diff --git a/internal/domain/channel.go b/internal/domain/channel.go index 80d74a6d..de32e173 100644 --- a/internal/domain/channel.go +++ b/internal/domain/channel.go @@ -347,6 +347,7 @@ type ChannelMessage struct { Reactions *ChannelMessageReactions Action *ChannelMessageAction Media *MessageMedia + Pinned bool Mentioned bool MediaUnread bool Pts int @@ -1377,6 +1378,7 @@ type ChannelHistoryFilter struct { ChannelID int64 Query string SenderUserID int64 + PinnedOnly bool OffsetID int OffsetDate int AddOffset int diff --git a/internal/rpc/convert.go b/internal/rpc/convert.go index 61107d19..79268dfe 100644 --- a/internal/rpc/convert.go +++ b/internal/rpc/convert.go @@ -1042,6 +1042,9 @@ func tgChannelMessage(viewerUserID int64, m domain.ChannelMessage) tg.MessageCla Message: m.Body, Entities: tgMessageEntities(m.Entities), } + if m.Pinned { + msg.SetPinned(true) + } if m.EditDate != 0 { msg.SetEditDate(m.EditDate) } diff --git a/internal/rpc/messages.go b/internal/rpc/messages.go index 28730317..1c7a18a3 100644 --- a/internal/rpc/messages.go +++ b/internal/rpc/messages.go @@ -453,6 +453,12 @@ func (r *Router) registerMessages(d *tg.ServerDispatcher) { } return tgChannelHistoryMessages(userID, history), nil } + if _, ok := req.Filter.(*tg.InputMessagesFilterPinned); ok { + if _, err := r.checkedDomainPeerFromInputPeer(ctx, userID, req.Peer); err != nil { + return nil, err + } + return &tg.MessagesMessages{}, nil + } if r.deps.Messages == nil { return messagesNotModifiedOrEmpty(req.Hash), nil } @@ -5746,16 +5752,17 @@ func (r *Router) channelHistoryFilterFromSearchRequest(userID int64, req *tg.Mes limit = 100 } filter := domain.ChannelHistoryFilter{ - ChannelID: channelID, - Query: req.Q, - OffsetID: req.OffsetID, - AddOffset: domain.ClampMessageHistoryAddOffset(req.AddOffset), - Limit: limit, - MinDate: req.MinDate, - MaxDate: req.MaxDate, - MaxID: req.MaxID, - MinID: req.MinID, - Hash: req.Hash, + ChannelID: channelID, + Query: req.Q, + PinnedOnly: messagesSearchFilterPinned(req.Filter), + OffsetID: req.OffsetID, + AddOffset: domain.ClampMessageHistoryAddOffset(req.AddOffset), + Limit: limit, + MinDate: req.MinDate, + MaxDate: req.MaxDate, + MaxID: req.MaxID, + MinID: req.MinID, + Hash: req.Hash, } if req.FromID != nil { from, ok := r.domainPeerFromInputPeer(userID, req.FromID) @@ -5767,6 +5774,11 @@ func (r *Router) channelHistoryFilterFromSearchRequest(userID int64, req *tg.Mes return filter, true } +func messagesSearchFilterPinned(filter tg.MessagesFilterClass) bool { + _, ok := filter.(*tg.InputMessagesFilterPinned) + return ok +} + func searchFilterNeedsMediaStore(filter tg.MessagesFilterClass) bool { switch filter.(type) { case nil, *tg.InputMessagesFilterEmpty: diff --git a/internal/rpc/router_test.go b/internal/rpc/router_test.go index cb99b0ec..29b0ee75 100644 --- a/internal/rpc/router_test.go +++ b/internal/rpc/router_test.go @@ -2723,6 +2723,7 @@ func TestMessagesSearchChannelPeerReturnsSingleCopyMessages(t *testing.T) { t.Fatalf("create chat: %v", err) } channel := created.Updates.(*tg.Updates).Chats[0].(*tg.Channel) + pinnedMsgID := 0 for _, item := range []struct { userID int64 text string @@ -2732,13 +2733,21 @@ func TestMessagesSearchChannelPeerReturnsSingleCopyMessages(t *testing.T) { {friend.ID, "not this one", 5002}, {friend.ID, "needle from friend", 5003}, } { - if _, err := r.onMessagesSendMessage(WithUserID(ctx, item.userID), &tg.MessagesSendMessageRequest{ + sent, err := r.onMessagesSendMessage(WithUserID(ctx, item.userID), &tg.MessagesSendMessageRequest{ Peer: &tg.InputPeerChannel{ChannelID: channel.ID, AccessHash: channel.AccessHash}, Message: item.text, RandomID: item.random, - }); err != nil { + }) + if err != nil { t.Fatalf("send %q: %v", item.text, err) } + if item.text == "not this one" { + channelUpdates := sent.(*tg.Updates) + if len(channelUpdates.Updates) == 0 { + t.Fatalf("send %q updates = %+v, want updateMessageID", item.text, channelUpdates.Updates) + } + pinnedMsgID = channelUpdates.Updates[0].(*tg.UpdateMessageID).ID + } } req := &tg.MessagesSearchRequest{ @@ -2804,6 +2813,34 @@ func TestMessagesSearchChannelPeerReturnsSingleCopyMessages(t *testing.T) { if channelMessages.Count != 0 || len(channelMessages.Messages) != 0 { t.Fatalf("shared media count search = count %d messages %d, want empty without media store", channelMessages.Count, len(channelMessages.Messages)) } + + if _, err := r.onMessagesUpdatePinnedMessage(WithUserID(ctx, owner.ID), &tg.MessagesUpdatePinnedMessageRequest{ + Peer: &tg.InputPeerChannel{ChannelID: channel.ID, AccessHash: channel.AccessHash}, + ID: pinnedMsgID, + }); err != nil { + t.Fatalf("pin channel message: %v", err) + } + pinnedReq := &tg.MessagesSearchRequest{ + Peer: &tg.InputPeerChannel{ChannelID: channel.ID, AccessHash: channel.AccessHash}, + Filter: &tg.InputMessagesFilterPinned{}, + Limit: 40, + } + in.Reset() + if err := pinnedReq.Encode(&in); err != nil { + t.Fatalf("encode pinned search: %v", err) + } + enc, err = r.Dispatch(WithUserID(ctx, friend.ID), [8]byte{}, 0, &in) + if err != nil { + t.Fatalf("dispatch pinned search: %v", err) + } + pinnedMessages, _, _ := searchMessagesPayload(t, enc) + if len(pinnedMessages) != 1 { + t.Fatalf("pinned search returned %d messages, want 1", len(pinnedMessages)) + } + pinnedMessage, ok := pinnedMessages[0].(*tg.Message) + if !ok || pinnedMessage.ID != pinnedMsgID || !pinnedMessage.GetPinned() { + t.Fatalf("pinned search message = %#v, want pinned message id=%d", pinnedMessages[0], pinnedMsgID) + } } func searchMessagesPayload(t *testing.T, enc bin.Encoder) ([]tg.MessageClass, []tg.ChatClass, []tg.UserClass) { diff --git a/internal/store/memory/channel.go b/internal/store/memory/channel.go index 2337c3ce..d59d027a 100644 --- a/internal/store/memory/channel.go +++ b/internal/store/memory/channel.go @@ -3133,6 +3133,9 @@ func (s *ChannelStore) ListChannelHistory(_ context.Context, viewerUserID int64, if msg.ID <= member.AvailableMinID { continue } + if filter.PinnedOnly && msg.ID != channel.PinnedMessageID { + continue + } if query != "" && !strings.Contains(strings.ToLower(msg.Body), query) { continue } @@ -3159,7 +3162,9 @@ func (s *ChannelStore) ListChannelHistory(_ context.Context, viewerUserID int64, } matched++ if len(out) < limit { - out = append(out, cloneChannelMessage(msg)) + item := cloneChannelMessage(msg) + item.Pinned = channel.PinnedMessageID != 0 && item.ID == channel.PinnedMessageID + out = append(out, item) } } s.populateChannelMessageRepliesLocked(viewerUserID, filter.ChannelID, out) @@ -3402,7 +3407,9 @@ func (s *ChannelStore) GetChannelMessages(_ context.Context, viewerUserID, chann if msg.Deleted || msg.ID <= member.AvailableMinID { continue } - messages = append(messages, cloneChannelMessage(msg)) + item := cloneChannelMessage(msg) + item.Pinned = channel.PinnedMessageID != 0 && item.ID == channel.PinnedMessageID + messages = append(messages, item) } sort.Slice(messages, func(i, j int) bool { return messages[i].ID > messages[j].ID }) s.populateChannelMessageRepliesLocked(viewerUserID, channelID, messages) diff --git a/internal/store/postgres/channel.go b/internal/store/postgres/channel.go index ac218622..8b16a253 100644 --- a/internal/store/postgres/channel.go +++ b/internal/store/postgres/channel.go @@ -4435,6 +4435,13 @@ func (s *ChannelStore) ListChannelHistory(ctx context.Context, viewerUserID int6 baseArgs = append(baseArgs, member.AvailableMinID) base += fmt.Sprintf(" AND id > $%d", len(baseArgs)) } + if filter.PinnedOnly { + if channel.PinnedMessageID <= 0 { + return domain.ChannelHistory{Channel: channel, Self: member}, nil + } + baseArgs = append(baseArgs, channel.PinnedMessageID) + base += fmt.Sprintf(" AND id = $%d", len(baseArgs)) + } if filter.Query != "" { baseArgs = append(baseArgs, filter.Query) base += fmt.Sprintf(" AND body ILIKE '%%' || $%d || '%%'", len(baseArgs)) @@ -4577,6 +4584,7 @@ func (s *ChannelStore) ListChannelHistory(ctx context.Context, viewerUserID int6 } out.Messages = older } + markPinnedChannelMessages(channel, out.Messages) out.Count = len(out.Messages) if hasMoreOlder { out.Count = len(out.Messages) + 1 @@ -4830,6 +4838,7 @@ ORDER BY id DESC`, args...) if err := rows.Err(); err != nil { return domain.ChannelHistory{}, err } + markPinnedChannelMessages(channel, out.Messages) out.Count = len(out.Messages) if err := s.populateChannelMessageReplies(ctx, s.db, viewerUserID, channel, out.Messages); err != nil { return domain.ChannelHistory{}, err @@ -4840,6 +4849,17 @@ ORDER BY id DESC`, args...) return out, nil } +func markPinnedChannelMessages(channel domain.Channel, messages []domain.ChannelMessage) { + if channel.PinnedMessageID <= 0 { + return + } + for i := range messages { + if messages[i].ChannelID == channel.ID && messages[i].ID == channel.PinnedMessageID { + messages[i].Pinned = true + } + } +} + func (s *ChannelStore) ReadChannelMessageContents(ctx context.Context, req domain.ReadChannelMessageContentsRequest) (domain.ReadChannelMessageContentsResult, error) { if req.UserID == 0 || req.ChannelID == 0 { return domain.ReadChannelMessageContentsResult{}, domain.ErrChannelInvalid