From e22fac1c3f809ccfd39cd8a5816526e1e9d1ff07 Mon Sep 17 00:00:00 2001 From: iamxvbaba <28732408+iamxvbaba@users.noreply.github.com> Date: Wed, 29 Jul 2026 13:51:02 +0800 Subject: [PATCH] fix(messages): sync preserve private dialogs after clearing history --- internal/admin/service.go | 5 +- internal/domain/media.go | 5 + internal/domain/message.go | 52 ++- internal/rpc/convert_history_clear_test.go | 30 ++ internal/rpc/convert_messages.go | 2 + internal/rpc/messages_delete.go | 12 +- internal/rpc/messages_delete_rpc_test.go | 6 +- internal/store/memory/message_delete.go | 185 +++++++--- internal/store/memory/message_history.go | 40 ++- internal/store/memory/message_test.go | 186 +++++++++- internal/store/memory/saved_dialog.go | 2 +- internal/store/postgres/message_delete.go | 332 ++++++++++++++---- .../message_delete_integration_test.go | 169 +++++++-- internal/store/postgres/message_history.go | 44 ++- internal/store/postgres/queries/dialog.sql | 24 -- internal/store/postgres/queries/message.sql | 2 + internal/store/postgres/saved_dialog.go | 2 +- internal/store/postgres/sqlcgen/dialog.sql.go | 36 -- .../store/postgres/sqlcgen/message.sql.go | 20 +- 19 files changed, 923 insertions(+), 231 deletions(-) create mode 100644 internal/rpc/convert_history_clear_test.go diff --git a/internal/admin/service.go b/internal/admin/service.go index 4807e241..2ce8c2ed 100644 --- a/internal/admin/service.go +++ b/internal/admin/service.go @@ -3137,11 +3137,12 @@ func summarizeDeleteResult(res domain.DeleteMessagesResult) []any { for _, item := range res.Deleted { ids := append([]int(nil), item.MessageIDs...) sort.Ints(ids) + pts, ptsCount := item.AffectedPts() out = append(out, map[string]any{ "user_id": item.UserID, "message_ids": ids, - "pts": item.Event.Pts, - "pts_count": item.Event.PtsCount, + "pts": pts, + "pts_count": ptsCount, }) } return out diff --git a/internal/domain/media.go b/internal/domain/media.go index 77120d05..38d5531f 100644 --- a/internal/domain/media.go +++ b/internal/domain/media.go @@ -537,6 +537,11 @@ type MessageServiceActionKind string const ( MessageServiceActionSuggestProfilePhoto MessageServiceActionKind = "suggest_profile_photo" + // MessageServiceActionHistoryClear 映射 messageActionHistoryClear。私聊 + // messages.deleteHistory(just_clear) 复用清理开始时的 top box id,把它 + // 原位转换成 owner-local 服务消息,使 getDialogs/getHistory 在冷启动时 + // 仍能从真实 top message 重建会话。 + MessageServiceActionHistoryClear MessageServiceActionKind = "history_clear" // MessageServiceActionPinMessage 映射 messageActionPinMessage:非 // pm_oneside 私聊置顶生成的服务消息,被置顶消息经 reply_to 指向。 MessageServiceActionPinMessage MessageServiceActionKind = "pin_message" diff --git a/internal/domain/message.go b/internal/domain/message.go index f868c608..4be181e0 100644 --- a/internal/domain/message.go +++ b/internal/domain/message.go @@ -165,6 +165,37 @@ type Message struct { SavedPeer Peer } +// NewHistoryClearMessage 把一个 owner 视角的现有 box 原位投影为 +// messageActionHistoryClear 服务消息。调用方只传递仍需保留的身份字段; +// 其余正文、媒体、reply、reaction、TTL、pin 等载荷全部按不变量清空。 +func NewHistoryClearMessage(ownerUserID int64, peer Peer, boxID int, uid int64, date, pts int) Message { + self := Peer{Type: PeerTypeUser, ID: ownerUserID} + return Message{ + ID: boxID, + UID: uid, + OwnerUserID: ownerUserID, + Peer: peer, + From: self, + Date: date, + Out: true, + Pts: pts, + Media: &MessageMedia{ + Kind: MessageMediaKindService, + ServiceAction: &MessageServiceAction{ + Kind: MessageServiceActionHistoryClear, + }, + }, + } +} + +// IsHistoryClearServiceMessage 报告该 box 是否已经是清空历史锚点。 +func IsHistoryClearServiceMessage(msg Message) bool { + return msg.Media != nil && + msg.Media.Kind == MessageMediaKindService && + msg.Media.ServiceAction != nil && + msg.Media.ServiceAction.Kind == MessageServiceActionHistoryClear +} + // MessageRichMessage 是 Layer 228 富文本消息(richMessage)的协议中立快照:一组 IV // PageBlock(Blocks)+ 内嵌已解析的 Photos/Documents。 // @@ -593,7 +624,23 @@ type DeleteHistoryRequest struct { type DeletedMessagesForUser struct { UserID int64 MessageIDs []int - Event UpdateEvent + // Event 保留单一 delete_messages 事件,供普通 deleteMessages 路径与 + // replay receipt 使用;just_clear 还会在 Events 中携带 read/edit 事件。 + Event UpdateEvent + Events []UpdateEvent + // Pts/PtsCount 是本次 method 对该 owner 的 affected watermark 汇总。 + // 普通删除等于 Event;just_clear 等于最后一条真实 update 的 pts 与 + // delete/read/edit 三段 pts_count 之和。 + Pts int + PtsCount int +} + +// AffectedPts 返回 messages.Affected* 应使用的最终 PTS 与本次总增量。 +func (d DeletedMessagesForUser) AffectedPts() (int, int) { + if d.Pts != 0 { + return d.Pts, d.PtsCount + } + return d.Event.Pts, d.Event.PtsCount } // DeleteMessagesResult 描述消息删除后的 owner 维度结果。 @@ -616,7 +663,8 @@ func (r DeleteMessagesResult) Self() DeletedMessagesForUser { // Changed 表示本次删除是否实际影响了任何 owner 视角。 func (r DeleteMessagesResult) Changed() bool { for _, item := range r.Deleted { - if len(item.MessageIDs) > 0 { + pts, _ := item.AffectedPts() + if len(item.MessageIDs) > 0 || pts > 0 { return true } } diff --git a/internal/rpc/convert_history_clear_test.go b/internal/rpc/convert_history_clear_test.go new file mode 100644 index 00000000..38434b74 --- /dev/null +++ b/internal/rpc/convert_history_clear_test.go @@ -0,0 +1,30 @@ +package rpc + +import ( + "testing" + + "github.com/iamxvbaba/td/tg" + + "telesrv/internal/domain" +) + +func TestTGMessageProjectsHistoryClearServiceAction(t *testing.T) { + message := domain.NewHistoryClearMessage( + 1001, + domain.Peer{Type: domain.PeerTypeUser, ID: 1002}, + 77, + 88, + 1700000000, + 9, + ) + got, ok := tgMessage(message).(*tg.MessageService) + if !ok { + t.Fatalf("message = %T, want *tg.MessageService", tgMessage(message)) + } + if got.ID != 77 || !got.Out || got.PeerID == nil || got.FromID == nil { + t.Fatalf("service message = %+v, want owner-local id/peer/from", got) + } + if _, ok := got.Action.(*tg.MessageActionHistoryClear); !ok { + t.Fatalf("action = %T, want *tg.MessageActionHistoryClear", got.Action) + } +} diff --git a/internal/rpc/convert_messages.go b/internal/rpc/convert_messages.go index a33fdaff..0c5485a9 100644 --- a/internal/rpc/convert_messages.go +++ b/internal/rpc/convert_messages.go @@ -137,6 +137,8 @@ func tgMessageServiceAction(msg domain.Message) tg.MessageActionClass { return nil } switch m.ServiceAction.Kind { + case domain.MessageServiceActionHistoryClear: + return &tg.MessageActionHistoryClear{} case domain.MessageServiceActionSuggestProfilePhoto: if m.ServiceAction.Photo == nil || m.ServiceAction.Photo.ID == 0 { return &tg.MessageActionEmpty{} diff --git a/internal/rpc/messages_delete.go b/internal/rpc/messages_delete.go index b66fd3a7..2228b730 100644 --- a/internal/rpc/messages_delete.go +++ b/internal/rpc/messages_delete.go @@ -34,10 +34,11 @@ func (r *Router) onMessagesDeleteMessages(ctx context.Context, req *tg.MessagesD return nil, internalErr() } self := res.Self() - if len(self.MessageIDs) == 0 || self.Event.Pts == 0 { + pts, ptsCount := self.AffectedPts() + if len(self.MessageIDs) == 0 || pts == 0 { return r.affectedMessages(ctx, authKeyID, userID) } - return &tg.MessagesAffectedMessages{Pts: self.Event.Pts, PtsCount: self.Event.PtsCount}, nil + return &tg.MessagesAffectedMessages{Pts: pts, PtsCount: ptsCount}, nil } func (r *Router) onMessagesDeleteHistory(ctx context.Context, req *tg.MessagesDeleteHistoryRequest) (*tg.MessagesAffectedHistory, error) { @@ -109,12 +110,13 @@ func (r *Router) onMessagesDeleteHistory(ctx context.Context, req *tg.MessagesDe return nil, internalErr() } self := res.Self() - if len(self.MessageIDs) == 0 || self.Event.Pts == 0 { + pts, ptsCount := self.AffectedPts() + if pts == 0 { return r.affectedHistory(ctx, authKeyID, userID, 0) } return &tg.MessagesAffectedHistory{ - Pts: self.Event.Pts, - PtsCount: self.Event.PtsCount, + Pts: pts, + PtsCount: ptsCount, Offset: res.Offset, }, nil } diff --git a/internal/rpc/messages_delete_rpc_test.go b/internal/rpc/messages_delete_rpc_test.go index ced63a3e..f76f40d7 100644 --- a/internal/rpc/messages_delete_rpc_test.go +++ b/internal/rpc/messages_delete_rpc_test.go @@ -149,6 +149,8 @@ func TestMessagesDeleteHistoryPassesJustClearContext(t *testing.T) { UserID: userID, MessageIDs: []int{1, 2, 3}, Event: domain.UpdateEvent{Pts: 12, PtsCount: 3}, + Pts: 14, + PtsCount: 5, }}, }} r := New(Config{}, Deps{Messages: messages}, zaptest.NewLogger(t), clock.System) @@ -171,8 +173,8 @@ func TestMessagesDeleteHistoryPassesJustClearContext(t *testing.T) { if !ok { t.Fatalf("response = %T, want *tg.MessagesAffectedHistory", enc) } - if got.Pts != 12 || got.PtsCount != 3 { - t.Fatalf("affected = %+v, want pts=12 pts_count=3", got) + if got.Pts != 14 || got.PtsCount != 5 { + t.Fatalf("affected = %+v, want aggregate pts=14 pts_count=5", got) } reqGot := messages.deleteHistoryReq if reqGot.OwnerUserID != userID || reqGot.Peer.ID != peerID || reqGot.MaxID != 15 || !reqGot.JustClear || !reqGot.Revoke || reqGot.OriginSessionID != 88 || reqGot.OriginAuthKeyID != authKeyID { diff --git a/internal/store/memory/message_delete.go b/internal/store/memory/message_delete.go index 076204ae..dfb72736 100644 --- a/internal/store/memory/message_delete.go +++ b/internal/store/memory/message_delete.go @@ -33,7 +33,7 @@ func (s *MessageStore) DeleteMessages(_ context.Context, req domain.DeleteMessag if req.Revoke && len(revokeUIDs) > 0 { deleted = append(deleted, s.deleteMemoryMessagesByUIDLocked(revokeUIDs, req.OwnerUserID)...) } - return s.finishMemoryDeleteLocked(res, deleted, req.Date, false), nil + return s.finishMemoryDeleteLocked(res, deleted, req.Date, nil), nil } type deletedMemoryMessage struct { @@ -45,8 +45,13 @@ type deletedMemoryMessage struct { randomID int64 } -func (s *MessageStore) finishMemoryDeleteLocked(res domain.DeleteMessagesResult, deleted []deletedMemoryMessage, date int, preserveEmptyDialogs bool) domain.DeleteMessagesResult { - if len(deleted) == 0 { +type memoryHistoryClearAnchor struct { + message domain.Message + materialized bool +} + +func (s *MessageStore) finishMemoryDeleteLocked(res domain.DeleteMessagesResult, deleted []deletedMemoryMessage, date int, anchors map[int64]memoryHistoryClearAnchor) domain.DeleteMessagesResult { + if len(deleted) == 0 && len(anchors) == 0 { return res } idsByOwner := make(map[int64][]int) @@ -64,57 +69,153 @@ func (s *MessageStore) finishMemoryDeleteLocked(res domain.DeleteMessagesResult, } peersByOwner[row.userID][row.peer] = struct{}{} } - if s.dialogs != nil { - s.dialogs.mu.Lock() - for userID, peers := range peersByOwner { - for peer := range peers { - s.rebuildMemoryDialogLocked(userID, peer, preserveEmptyDialogs) - } + for userID, anchor := range anchors { + if peersByOwner[userID] == nil { + peersByOwner[userID] = make(map[domain.Peer]struct{}) } - s.dialogs.mu.Unlock() + peersByOwner[userID][anchor.message.Peer] = struct{}{} } - ownerIDs := make([]int64, 0, len(idsByOwner)) + ownerSet := make(map[int64]struct{}, len(idsByOwner)+len(anchors)) for userID := range idsByOwner { + ownerSet[userID] = struct{}{} + } + for userID, anchor := range anchors { + if !anchor.materialized { + ownerSet[userID] = struct{}{} + } + } + ownerIDs := make([]int64, 0, len(ownerSet)) + for userID := range ownerSet { ownerIDs = append(ownerIDs, userID) } sort.Slice(ownerIDs, func(i, j int) bool { return ownerIDs[i] < ownerIDs[j] }) for _, userID := range ownerIDs { ids := normalizeMemoryMessageIDs(idsByOwner[userID]) - if len(ids) == 0 { + anchor, hasAnchor := anchors[userID] + materializeAnchor := hasAnchor && !anchor.materialized + totalPtsCount := len(ids) + if materializeAnchor { + totalPtsCount += 2 + } + if totalPtsCount == 0 { continue } - pts := s.nextPtsNLocked(userID, len(ids)) - event := domain.UpdateEvent{ + pts := s.nextPtsNLocked(userID, totalPtsCount) + cursor := pts - totalPtsCount + item := domain.DeletedMessagesForUser{ UserID: userID, - Type: domain.UpdateEventDeleteMessages, + MessageIDs: ids, Pts: pts, - PtsCount: len(ids), - Date: date, - MessageIDs: ids, + PtsCount: totalPtsCount, + Events: make([]domain.UpdateEvent, 0, 3), } - for _, row := range deleted { - if row.userID != userID || row.messageSenderID != userID || row.randomID == 0 || row.privateMessageID == 0 { - continue + if len(ids) > 0 { + cursor += len(ids) + event := domain.UpdateEvent{ + UserID: userID, + Type: domain.UpdateEventDeleteMessages, + Pts: cursor, + PtsCount: len(ids), + Date: date, + MessageIDs: ids, } - key := privateSendDedupKey{senderUserID: userID, randomID: row.randomID} - record, ok := s.privateSendDedup[key] - if !ok { - continue + for _, row := range deleted { + if row.userID != userID || row.messageSenderID != userID || row.randomID == 0 || row.privateMessageID == 0 { + continue + } + key := privateSendDedupKey{senderUserID: userID, randomID: row.randomID} + record, ok := s.privateSendDedup[key] + if !ok { + continue + } + cloned := cloneUpdateEvent(event) + record.senderDeleteEvent = &cloned + s.privateSendDedup[key] = record } - cloned := cloneUpdateEvent(event) - record.senderDeleteEvent = &cloned - s.privateSendDedup[key] = record + item.Event = event + item.Events = append(item.Events, event) } - res.Deleted = append(res.Deleted, domain.DeletedMessagesForUser{ - UserID: userID, - MessageIDs: ids, - Event: event, - }) + if materializeAnchor { + readPts := cursor + 1 + editPts := readPts + 1 + msg := domain.NewHistoryClearMessage( + userID, + anchor.message.Peer, + anchor.message.ID, + anchor.message.UID, + anchor.message.Date, + editPts, + ) + for i := range s.m[userID] { + if s.m[userID][i].ID == anchor.message.ID && s.m[userID][i].Peer == anchor.message.Peer { + s.m[userID][i] = msg + break + } + } + if byMessage := s.savedMessageTags[userID]; byMessage != nil { + delete(byMessage, anchor.message.ID) + if len(byMessage) == 0 { + delete(s.savedMessageTags, userID) + } + } + readEvent := domain.UpdateEvent{ + UserID: userID, + Type: domain.UpdateEventReadHistoryInbox, + Pts: readPts, + PtsCount: 1, + Date: date, + Peer: anchor.message.Peer, + MaxID: anchor.message.ID, + StillUnreadCount: 0, + } + editEvent := domain.UpdateEvent{ + UserID: userID, + Type: domain.UpdateEventEditMessage, + Pts: editPts, + PtsCount: 1, + Date: date, + Message: cloneMessage(msg), + } + item.Events = append(item.Events, readEvent, editEvent) + cursor = editPts + } + if s.dialogs != nil { + s.dialogs.mu.Lock() + for peer := range peersByOwner[userID] { + s.rebuildMemoryDialogLocked(userID, peer) + } + if materializeAnchor { + s.advanceMemoryHistoryClearDialogLocked(userID, anchor.message.Peer, anchor.message.ID) + } + s.dialogs.mu.Unlock() + } + if cursor != pts { + panic(fmt.Sprintf("memory delete history pts cursor %d does not reach reserved pts %d", cursor, pts)) + } + res.Deleted = append(res.Deleted, item) } return res } -func (s *MessageStore) rebuildMemoryDialogLocked(userID int64, peer domain.Peer, preserveEmpty bool) { +func (s *MessageStore) advanceMemoryHistoryClearDialogLocked(userID int64, peer domain.Peer, maxID int) { + list := s.dialogs.m[userID] + for i := range list.Dialogs { + if list.Dialogs[i].Peer != peer { + continue + } + if list.Dialogs[i].ReadInboxMaxID < maxID { + list.Dialogs[i].ReadInboxMaxID = maxID + } + list.Dialogs[i].UnreadCount = 0 + list.Dialogs[i].UnreadMark = false + list.Dialogs[i].UnreadMentions = 0 + list.Dialogs[i].UnreadReactions = 0 + break + } + s.dialogs.m[userID] = list +} + +func (s *MessageStore) rebuildMemoryDialogLocked(userID int64, peer domain.Peer) { list := s.dialogs.m[userID] topID := 0 topDate := 0 @@ -135,22 +236,6 @@ func (s *MessageStore) rebuildMemoryDialogLocked(userID int64, peer domain.Peer, continue } if topID == 0 { - if preserveEmpty { - oldTop := dialog.TopMessage - dialog.TopMessage = 0 - dialog.TopMessageDate = 0 - if dialog.ReadInboxMaxID < oldTop { - dialog.ReadInboxMaxID = oldTop - } - if dialog.ReadOutboxMaxID < oldTop { - dialog.ReadOutboxMaxID = oldTop - } - dialog.UnreadCount = 0 - dialog.UnreadMark = false - dialog.UnreadMentions = 0 - dialog.UnreadReactions = 0 - dialogs = append(dialogs, dialog) - } continue } for _, msg := range s.m[userID] { diff --git a/internal/store/memory/message_history.go b/internal/store/memory/message_history.go index 5a7d21fa..1413799c 100644 --- a/internal/store/memory/message_history.go +++ b/internal/store/memory/message_history.go @@ -232,30 +232,66 @@ func (s *MessageStore) DeleteHistory(_ context.Context, req domain.DeleteHistory } return true } + var anchors map[int64]memoryHistoryClearAnchor + fullJustClear := req.JustClear && req.MaxID <= 0 && req.MinDate <= 0 && req.MaxDate <= 0 + if fullJustClear { + anchors = make(map[int64]memoryHistoryClearAnchor, 2) + if anchor, found := s.memoryHistoryClearAnchorLocked(req.OwnerUserID, req.Peer); found { + anchors[req.OwnerUserID] = anchor + } + if req.Revoke && req.Peer.ID != req.OwnerUserID { + peer := domain.Peer{Type: domain.PeerTypeUser, ID: req.OwnerUserID} + if anchor, found := s.memoryHistoryClearAnchorLocked(req.Peer.ID, peer); found { + anchors[req.Peer.ID] = anchor + } + } + } deleted, revokeUIDs, more := s.deleteMemoryMessagesLocked(req.OwnerUserID, domain.MaxDeleteHistoryBatch, func(msg domain.Message) bool { + if anchor, ok := anchors[req.OwnerUserID]; ok && msg.ID == anchor.message.ID { + return false + } return msg.Peer == req.Peer && (req.MaxID <= 0 || msg.ID <= req.MaxID) && inDateRange(msg) }) if req.Revoke { - if len(revokeUIDs) > 0 { + if req.MaxID > 0 && len(revokeUIDs) > 0 { deleted = append(deleted, s.deleteMemoryMessagesByUIDLocked(revokeUIDs, req.OwnerUserID)...) } // 与 PG 同语义:全量/按日期的双向清史直扫对端残余,我方早已 // 单向删除的消息不能在对端残留。 if req.MaxID <= 0 && req.Peer.ID != req.OwnerUserID { peerDeleted, _, peerMore := s.deleteMemoryMessagesLocked(req.Peer.ID, domain.MaxDeleteHistoryBatch, func(msg domain.Message) bool { + if anchor, ok := anchors[req.Peer.ID]; ok && msg.ID == anchor.message.ID { + return false + } return msg.Peer == (domain.Peer{Type: domain.PeerTypeUser, ID: req.OwnerUserID}) && inDateRange(msg) }) deleted = append(deleted, peerDeleted...) more = more || peerMore } } - res = s.finishMemoryDeleteLocked(res, deleted, req.Date, req.JustClear) + res = s.finishMemoryDeleteLocked(res, deleted, req.Date, anchors) if more { res.Offset = 1 } return res, nil } +func (s *MessageStore) memoryHistoryClearAnchorLocked(userID int64, peer domain.Peer) (memoryHistoryClearAnchor, bool) { + var top domain.Message + for _, msg := range s.m[userID] { + if msg.Peer == peer && msg.ID > top.ID { + top = msg + } + } + if top.ID == 0 { + return memoryHistoryClearAnchor{}, false + } + return memoryHistoryClearAnchor{ + message: cloneMessage(top), + materialized: domain.IsHistoryClearServiceMessage(top), + }, true +} + func filterMessageList(messages []domain.Message, filter domain.MessageFilter) domain.MessageList { filter.AddOffset = domain.ClampMessageHistoryAddOffset(filter.AddOffset) sort.SliceStable(messages, func(i, j int) bool { diff --git a/internal/store/memory/message_test.go b/internal/store/memory/message_test.go index 7c2ce4af..c6e43e79 100644 --- a/internal/store/memory/message_test.go +++ b/internal/store/memory/message_test.go @@ -1218,29 +1218,191 @@ func TestMessageStoreDeleteHistoryDeletesOrPreservesDialogAndRebuilds(t *testing preservedOwner := int64(1000000003) preservedPeerID := int64(1000000004) preservedPeer := domain.Peer{Type: domain.PeerTypeUser, ID: preservedPeerID} - if _, err := messages.SendPrivateText(ctx, domain.SendPrivateTextRequest{ - SenderUserID: preservedOwner, - RecipientUserID: preservedPeerID, - RandomID: 300, - Message: "clear but keep dialog", - Date: 1700000500, - }); err != nil { - t.Fatalf("seed preserved send: %v", err) + var preservedTop domain.Message + for i := 0; i < 2; i++ { + sent, err := messages.SendPrivateText(ctx, domain.SendPrivateTextRequest{ + SenderUserID: preservedOwner, + RecipientUserID: preservedPeerID, + RandomID: int64(300 + i), + Message: "clear but keep dialog", + Date: 1700000500 + i, + }) + if err != nil { + t.Fatalf("seed preserved send %d: %v", i, err) + } + preservedTop = sent.SenderMessage } - if _, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ + clearResult, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ OwnerUserID: preservedOwner, Peer: preservedPeer, JustClear: true, Date: 1700000600, - }); err != nil { + }) + if err != nil { t.Fatalf("DeleteHistory just_clear: %v", err) } + clearSelf := clearResult.Self() + if clearSelf.Pts != 5 || clearSelf.PtsCount != 3 || len(clearSelf.MessageIDs) != 1 || len(clearSelf.Events) != 3 { + t.Fatalf("just_clear result = %+v, want delete+read+edit ending pts=5 count=3", clearSelf) + } + if clearSelf.Events[0].Type != domain.UpdateEventDeleteMessages || + clearSelf.Events[1].Type != domain.UpdateEventReadHistoryInbox || + clearSelf.Events[2].Type != domain.UpdateEventEditMessage { + t.Fatalf("just_clear events = %+v, want delete/read/edit order", clearSelf.Events) + } preservedDialogs, err := dialogs.ListByUser(ctx, preservedOwner, domain.DialogFilter{Limit: 10}) if err != nil { t.Fatalf("preserved dialogs: %v", err) } - if len(preservedDialogs.Dialogs) != 1 || preservedDialogs.Dialogs[0].Peer != preservedPeer || preservedDialogs.Dialogs[0].TopMessage != 0 || len(preservedDialogs.Messages) != 0 { - t.Fatalf("preserved dialogs = %+v messages=%+v, want empty dialog kept after just_clear", preservedDialogs.Dialogs, preservedDialogs.Messages) + if len(preservedDialogs.Dialogs) != 1 || preservedDialogs.Dialogs[0].Peer != preservedPeer || + preservedDialogs.Dialogs[0].TopMessage != preservedTop.ID || len(preservedDialogs.Messages) != 1 { + t.Fatalf("preserved dialogs = %+v messages=%+v, want history-clear top %d", preservedDialogs.Dialogs, preservedDialogs.Messages, preservedTop.ID) + } + clearMessage := preservedDialogs.Messages[0] + if !domain.IsHistoryClearServiceMessage(clearMessage) || clearMessage.ID != preservedTop.ID || + !clearMessage.Out || clearMessage.From.ID != preservedOwner || clearMessage.Body != "" || + clearMessage.ReplyTo != nil || clearMessage.Forward != nil || clearMessage.MediaUnread || + clearMessage.ReactionUnread || clearMessage.Pinned { + t.Fatalf("history clear anchor = %+v, want clean owner-local service message", clearMessage) + } + repeated, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ + OwnerUserID: preservedOwner, + Peer: preservedPeer, + JustClear: true, + Date: 1700000601, + }) + if err != nil { + t.Fatalf("repeat DeleteHistory just_clear: %v", err) + } + if repeated.Changed() || len(repeated.Deleted) != 0 || messages.nextPts[preservedOwner] != 5 { + t.Fatalf("repeat just_clear = %+v pts=%d, want idempotent no-op", repeated, messages.nextPts[preservedOwner]) + } +} + +func TestMessageStoreDeleteHistoryJustClearRevokeKeepsPerOwnerAnchors(t *testing.T) { + ctx := context.Background() + dialogs := NewDialogStore() + messages := NewMessageStore(dialogs) + const alice, bob = int64(1101), int64(1102) + var sent domain.SendPrivateTextResult + for i := 0; i < 2; i++ { + var err error + sent, err = messages.SendPrivateText(ctx, domain.SendPrivateTextRequest{ + SenderUserID: alice, RecipientUserID: bob, RandomID: int64(800 + i), + Message: "revoke clear", Date: 1700000700 + i, + }) + if err != nil { + t.Fatalf("send %d: %v", i, err) + } + } + res, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ + OwnerUserID: alice, + Peer: domain.Peer{Type: domain.PeerTypeUser, ID: bob}, + JustClear: true, + Revoke: true, + Date: 1700000800, + }) + if err != nil { + t.Fatalf("revoke just_clear: %v", err) + } + if len(res.Deleted) != 2 { + t.Fatalf("deleted owners = %+v, want alice and bob", res.Deleted) + } + for _, tc := range []struct { + userID int64 + peerID int64 + topID int + }{ + {alice, bob, sent.SenderMessage.ID}, + {bob, alice, sent.RecipientMessage.ID}, + } { + history, err := messages.ListByUser(ctx, tc.userID, domain.MessageFilter{ + HasPeer: true, Peer: domain.Peer{Type: domain.PeerTypeUser, ID: tc.peerID}, Limit: 10, + }) + if err != nil { + t.Fatalf("history user %d: %v", tc.userID, err) + } + if len(history.Messages) != 1 || history.Messages[0].ID != tc.topID || + !domain.IsHistoryClearServiceMessage(history.Messages[0]) || + history.Messages[0].From.ID != tc.userID || !history.Messages[0].Out { + t.Fatalf("history user %d = %+v, want owner-local anchor %d", tc.userID, history.Messages, tc.topID) + } + } +} + +func TestMessageStoreDeleteHistoryDateRangeDoesNotCreateHistoryClearAnchor(t *testing.T) { + ctx := context.Background() + dialogs := NewDialogStore() + messages := NewMessageStore(dialogs) + const owner, peerID = int64(1201), int64(1202) + peer := domain.Peer{Type: domain.PeerTypeUser, ID: peerID} + for i, date := range []int{100, 200} { + if _, err := messages.SendPrivateText(ctx, domain.SendPrivateTextRequest{ + SenderUserID: owner, RecipientUserID: peerID, RandomID: int64(900 + i), + Message: "dated", Date: date, + }); err != nil { + t.Fatalf("send %d: %v", i, err) + } + } + if _, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ + OwnerUserID: owner, Peer: peer, JustClear: true, MinDate: 150, MaxDate: 250, Date: 300, + }); err != nil { + t.Fatalf("date delete: %v", err) + } + history, err := messages.ListByUser(ctx, owner, domain.MessageFilter{HasPeer: true, Peer: peer, Limit: 10}) + if err != nil { + t.Fatalf("history: %v", err) + } + if len(history.Messages) != 1 || history.Messages[0].Date != 100 || domain.IsHistoryClearServiceMessage(history.Messages[0]) { + t.Fatalf("date history = %+v, want surviving ordinary message only", history.Messages) + } +} + +func TestMessageStoreDeleteHistoryJustClearKeepsAnchorAcrossBatches(t *testing.T) { + ctx := context.Background() + dialogs := NewDialogStore() + messages := NewMessageStore(dialogs) + const owner, peerID = int64(1301), int64(1302) + peer := domain.Peer{Type: domain.PeerTypeUser, ID: peerID} + total := domain.MaxDeleteHistoryBatch + 2 + var topID int + for i := 0; i < total; i++ { + sent, err := messages.SendPrivateText(ctx, domain.SendPrivateTextRequest{ + SenderUserID: owner, RecipientUserID: peerID, RandomID: int64(10000 + i), + Message: "batch clear", Date: 1700010000 + i, + }) + if err != nil { + t.Fatalf("send %d: %v", i, err) + } + topID = sent.SenderMessage.ID + } + first, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ + OwnerUserID: owner, Peer: peer, JustClear: true, Date: 1700020000, + }) + if err != nil { + t.Fatalf("first clear: %v", err) + } + if first.Offset == 0 || len(first.Self().MessageIDs) != domain.MaxDeleteHistoryBatch || + first.Self().PtsCount != domain.MaxDeleteHistoryBatch+2 { + t.Fatalf("first clear = %+v, want full batch plus one read/edit", first.Self()) + } + second, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ + OwnerUserID: owner, Peer: peer, JustClear: true, Date: 1700020001, + }) + if err != nil { + t.Fatalf("second clear: %v", err) + } + if second.Offset != 0 || len(second.Self().MessageIDs) != 1 || second.Self().PtsCount != 1 || + len(second.Self().Events) != 1 || second.Self().Events[0].Type != domain.UpdateEventDeleteMessages { + t.Fatalf("second clear = %+v, want remaining delete only", second.Self()) + } + history, err := messages.ListByUser(ctx, owner, domain.MessageFilter{HasPeer: true, Peer: peer, Limit: 10}) + if err != nil { + t.Fatalf("history: %v", err) + } + if len(history.Messages) != 1 || history.Messages[0].ID != topID || + !domain.IsHistoryClearServiceMessage(history.Messages[0]) { + t.Fatalf("history = %+v, want stable top anchor %d", history.Messages, topID) } } diff --git a/internal/store/memory/saved_dialog.go b/internal/store/memory/saved_dialog.go index 952a68d8..b7662d4c 100644 --- a/internal/store/memory/saved_dialog.go +++ b/internal/store/memory/saved_dialog.go @@ -243,7 +243,7 @@ func (s *MessageStore) DeleteSavedHistory(_ context.Context, req domain.DeleteSa return true } deleted, _, more := s.deleteMemoryMessagesLocked(req.OwnerUserID, domain.MaxDeleteHistoryBatch, match) - delRes := s.finishMemoryDeleteLocked(domain.DeleteMessagesResult{OwnerUserID: req.OwnerUserID}, deleted, req.Date, false) + delRes := s.finishMemoryDeleteLocked(domain.DeleteMessagesResult{OwnerUserID: req.OwnerUserID}, deleted, req.Date, nil) res.More = more for _, d := range delRes.Deleted { if d.UserID == req.OwnerUserID { diff --git a/internal/store/postgres/message_delete.go b/internal/store/postgres/message_delete.go index 3d4406ef..4cb48bb2 100644 --- a/internal/store/postgres/message_delete.go +++ b/internal/store/postgres/message_delete.go @@ -71,7 +71,7 @@ func (s *MessageStore) DeleteMessages(ctx context.Context, req domain.DeleteMess } deleted = append(deleted, deletedRowsFromPrivateRows(peerRows)...) } - res, err = s.finishDeleteMessagesTx(ctx, tx, qtx, req.OwnerUserID, req.OriginAuthKeyID, req.OriginSessionID, req.Date, deleted, false) + res, err = s.finishDeleteMessagesTx(ctx, tx, qtx, req.OwnerUserID, req.OriginAuthKeyID, req.OriginSessionID, req.Date, deleted, nil) if err != nil { return res, err } @@ -118,9 +118,53 @@ type deletedOwnerPeerKey struct { peer domain.Peer } -func (s *MessageStore) finishDeleteMessagesTx(ctx context.Context, db sqlcgen.DBTX, q *sqlcgen.Queries, ownerUserID int64, excludeAuthKeyID [8]byte, excludeSessionID int64, date int, rows []deletedBox, preserveEmptyDialogs bool) (domain.DeleteMessagesResult, error) { +type historyClearAnchor struct { + userID int64 + peer domain.Peer + boxID int + uid int64 + messageDate int + materialized bool +} + +func (s *MessageStore) loadHistoryClearAnchor(ctx context.Context, q *sqlcgen.Queries, userID int64, peer domain.Peer) (historyClearAnchor, bool, error) { + top, err := q.TopVisibleMessageBoxByPeer(ctx, sqlcgen.TopVisibleMessageBoxByPeerParams{ + OwnerUserID: userID, + PeerType: string(peer.Type), + PeerID: peer.ID, + }) + if errors.Is(err, pgx.ErrNoRows) { + return historyClearAnchor{}, false, nil + } + if err != nil { + return historyClearAnchor{}, false, fmt.Errorf("load history clear top: %w", err) + } + row, err := q.GetMessageBoxForEdit(ctx, sqlcgen.GetMessageBoxForEditParams{ + OwnerUserID: userID, + BoxID: top.BoxID, + PeerType: string(peer.Type), + PeerID: peer.ID, + }) + if err != nil { + return historyClearAnchor{}, false, fmt.Errorf("lock history clear top: %w", err) + } + media, err := decodeMessageMedia(row.MediaJson) + if err != nil { + return historyClearAnchor{}, false, fmt.Errorf("decode history clear top media: %w", err) + } + return historyClearAnchor{ + userID: userID, + peer: peer, + boxID: int(row.BoxID), + uid: row.PrivateMessageID, + messageDate: int(row.MessageDate), + materialized: domain.IsHistoryClearServiceMessage(domain.Message{Media: media}), + }, true, nil +} + +func (s *MessageStore) finishDeleteMessagesTx(ctx context.Context, db sqlcgen.DBTX, q *sqlcgen.Queries, ownerUserID int64, excludeAuthKeyID [8]byte, excludeSessionID int64, date int, rows []deletedBox, anchors map[int64]historyClearAnchor) (domain.DeleteMessagesResult, error) { res := domain.DeleteMessagesResult{OwnerUserID: ownerUserID} - if len(rows) == 0 { + if len(rows) == 0 && len(anchors) == 0 { return res, nil } peersByOwner := make(map[int64]map[domain.Peer]struct{}) @@ -145,6 +189,20 @@ func (s *MessageStore) finishDeleteMessagesTx(ctx context.Context, db sqlcgen.DB incomingDeletedByPeer[key][row.boxID] = struct{}{} } } + for userID, anchor := range anchors { + if anchor.boxID <= 0 || anchor.peer.ID == 0 { + continue + } + if peersByOwner[userID] == nil { + peersByOwner[userID] = make(map[domain.Peer]struct{}) + } + peersByOwner[userID][anchor.peer] = struct{}{} + if !anchor.materialized { + // 首次物化锚点会用一条 max_id=anchor 的真实 read update 覆盖该 + // peer 的全部已读校正;不能再为本批删除的 incoming prefix 重复推进。 + delete(incomingDeletedByPeer, deletedOwnerPeerKey{userID: userID, peer: anchor.peer}) + } + } // 按 owner 升序重建 dialog,使两个反向 delete(X 删与 Y 的会话 / Y 删与 X 的会话)以一致顺序 // 获取 dialog 行锁,配合下方 watermark 的升序推进,彻底避免 delete-delete 之间的 AB-BA 死锁。 rebuildOwners := make([]int64, 0, len(peersByOwner)) @@ -154,11 +212,12 @@ func (s *MessageStore) finishDeleteMessagesTx(ctx context.Context, db sqlcgen.DB sort.Slice(rebuildOwners, func(i, j int) bool { return rebuildOwners[i] < rebuildOwners[j] }) for _, userID := range rebuildOwners { for peer := range peersByOwner[userID] { - // just_clear(preserveEmptyDialogs)是请求者的本端语义"清空但保留我这侧 - // 空会话"。revoke 反查出的对端并未选择 just_clear,其空会话应按普通删除 - // 处理(无存活消息则移除 dialog),不能也被保留成空会话。 - preserve := preserveEmptyDialogs && userID == ownerUserID - if err := rebuildDialogAfterMessageDelete(ctx, q, userID, peer, preserve); err != nil { + if anchor, ok := anchors[userID]; ok && anchor.peer == peer && !anchor.materialized { + // 先分配 edit PTS 并原位转换锚点,再按转换后的 outgoing + // 状态重算 dialog;否则会短暂把 incoming anchor 计为未读。 + continue + } + if err := rebuildDialogAfterMessageDelete(ctx, q, userID, peer); err != nil { return res, err } } @@ -168,8 +227,17 @@ func (s *MessageStore) finishDeleteMessagesTx(ctx context.Context, db sqlcgen.DB return res, err } - ownerIDs := make([]int64, 0, len(idsByOwner)) + ownerSet := make(map[int64]struct{}, len(idsByOwner)+len(anchors)) for userID := range idsByOwner { + ownerSet[userID] = struct{}{} + } + for userID, anchor := range anchors { + if !anchor.materialized { + ownerSet[userID] = struct{}{} + } + } + ownerIDs := make([]int64, 0, len(ownerSet)) + for userID := range ownerSet { ownerIDs = append(ownerIDs, userID) } sort.Slice(ownerIDs, func(i, j int) bool { return ownerIDs[i] < ownerIDs[j] }) @@ -177,36 +245,56 @@ func (s *MessageStore) finishDeleteMessagesTx(ctx context.Context, db sqlcgen.DB res.Deleted = make([]domain.DeletedMessagesForUser, 0, len(ownerIDs)) for _, userID := range ownerIDs { ids := normalizeMessageIDs(idsByOwner[userID]) - if len(ids) == 0 { + corrections := readCorrectionsByOwner[userID] + anchor, hasAnchor := anchors[userID] + materializeAnchor := hasAnchor && !anchor.materialized + totalPtsCount := len(ids) + len(corrections) + if materializeAnchor { + totalPtsCount += 2 // updateReadHistoryInbox + updateEditMessage + } + if totalPtsCount == 0 { continue } - corrections := readCorrectionsByOwner[userID] - totalPtsCount := len(ids) + len(corrections) pts, err := s.reservePtsN(ctx, db, userID, totalPtsCount) if err != nil { return res, fmt.Errorf("allocate delete messages pts: %w", err) } - deletePts := pts - len(corrections) - event := domain.UpdateEvent{ + cursor := pts - totalPtsCount + item := domain.DeletedMessagesForUser{ UserID: userID, - Type: domain.UpdateEventDeleteMessages, - Pts: deletePts, - PtsCount: len(ids), - Date: date, MessageIDs: ids, + Pts: pts, + PtsCount: totalPtsCount, + Events: make([]domain.UpdateEvent, 0, 1+len(corrections)+2), } - deleteIDsJSON, err := encodeEventMessageIDs(event.MessageIDs) - if err != nil { - return res, fmt.Errorf("encode sender delete receipt ids: %w", err) + dispatchAuthKeyID := [8]byte{} + dispatchSessionID := int64(0) + if userID == ownerUserID { + dispatchAuthKeyID = excludeAuthKeyID + dispatchSessionID = excludeSessionID } - senderPrivateIDs := make([]int64, 0, len(rows)) - for _, row := range rows { - if row.ownerUserID == userID && row.messageSenderID == userID && row.privateMessageID != 0 { - senderPrivateIDs = append(senderPrivateIDs, row.privateMessageID) + if len(ids) > 0 { + cursor += len(ids) + event := domain.UpdateEvent{ + UserID: userID, + Type: domain.UpdateEventDeleteMessages, + Pts: cursor, + PtsCount: len(ids), + Date: date, + MessageIDs: ids, } - } - if len(senderPrivateIDs) > 0 { - if _, err := db.Exec(ctx, ` + deleteIDsJSON, err := encodeEventMessageIDs(event.MessageIDs) + if err != nil { + return res, fmt.Errorf("encode sender delete receipt ids: %w", err) + } + senderPrivateIDs := make([]int64, 0, len(rows)) + for _, row := range rows { + if row.ownerUserID == userID && row.messageSenderID == userID && row.privateMessageID != 0 { + senderPrivateIDs = append(senderPrivateIDs, row.privateMessageID) + } + } + if len(senderPrivateIDs) > 0 { + if _, err := db.Exec(ctx, ` UPDATE private_messages SET sender_delete_pts = $3, sender_delete_pts_count = $4, @@ -215,34 +303,27 @@ SET sender_delete_pts = $3, WHERE sender_user_id = $1 AND id = ANY($2::bigint[]) AND sender_box_id > 0`, userID, senderPrivateIDs, event.Pts, event.PtsCount, event.Date, deleteIDsJSON); err != nil { - return res, fmt.Errorf("save sender delete replay receipt: %w", err) + return res, fmt.Errorf("save sender delete replay receipt: %w", err) + } } + if err := appendDeleteMessagesEvent(ctx, q, event); err != nil { + return res, err + } + if err := enqueueDispatch(ctx, q, sqlcgen.EnqueueDispatchParams{ + TargetUserID: userID, + Pts: int32(event.Pts), + EventType: string(domain.UpdateEventDeleteMessages), + ExcludeAuthKeyID: authKeyIDToInt64(dispatchAuthKeyID), + ExcludeSessionID: dispatchSessionID, + }); err != nil { + return res, fmt.Errorf("enqueue delete messages dispatch: %w", err) + } + item.Event = event + item.Events = append(item.Events, event) } - if err := appendDeleteMessagesEvent(ctx, q, event); err != nil { - return res, err - } - dispatchAuthKeyID := [8]byte{} - dispatchSessionID := int64(0) - if userID == ownerUserID { - dispatchAuthKeyID = excludeAuthKeyID - dispatchSessionID = excludeSessionID - } - if err := enqueueDispatch(ctx, q, sqlcgen.EnqueueDispatchParams{ - TargetUserID: userID, - Pts: int32(deletePts), - EventType: string(domain.UpdateEventDeleteMessages), - ExcludeAuthKeyID: authKeyIDToInt64(dispatchAuthKeyID), - ExcludeSessionID: dispatchSessionID, - }); err != nil { - return res, fmt.Errorf("enqueue delete messages dispatch: %w", err) - } - res.Deleted = append(res.Deleted, domain.DeletedMessagesForUser{ - UserID: userID, - MessageIDs: ids, - Event: event, - }) - for i, correction := range corrections { - correction.Pts = deletePts + i + 1 + for _, correction := range corrections { + cursor++ + correction.Pts = cursor if err := appendUserUpdateEvent(ctx, db, q, userID, correction); err != nil { return res, fmt.Errorf("append delete unread correction event: %w", err) } @@ -263,11 +344,142 @@ WHERE sender_user_id = $1 }); err != nil { return res, fmt.Errorf("enqueue delete unread correction dispatch: %w", err) } + item.Events = append(item.Events, correction) } + if materializeAnchor { + readPts := cursor + 1 + editPts := readPts + 1 + msg, err := materializeHistoryClearAnchorTx(ctx, db, anchor, editPts) + if err != nil { + return res, err + } + if err := rebuildDialogAfterMessageDelete(ctx, q, userID, anchor.peer); err != nil { + return res, err + } + readEvent := domain.UpdateEvent{ + UserID: userID, + Type: domain.UpdateEventReadHistoryInbox, + Pts: readPts, + PtsCount: 1, + Date: date, + Peer: anchor.peer, + MaxID: anchor.boxID, + StillUnreadCount: 0, + } + if err := appendUserUpdateEvent(ctx, db, q, userID, readEvent); err != nil { + return res, fmt.Errorf("append history clear read event: %w", err) + } + if err := q.AdvanceDialogReadInboxFloor(ctx, sqlcgen.AdvanceDialogReadInboxFloorParams{ + UserID: userID, + PeerType: string(anchor.peer.Type), + PeerID: anchor.peer.ID, + ReadInboxMaxID: int32(anchor.boxID), + }); err != nil { + return res, fmt.Errorf("advance history clear read inbox: %w", err) + } + if err := enqueueDispatch(ctx, q, sqlcgen.EnqueueDispatchParams{ + TargetUserID: userID, + Pts: int32(readPts), + EventType: string(domain.UpdateEventReadHistoryInbox), + ExcludeAuthKeyID: authKeyIDToInt64(dispatchAuthKeyID), + ExcludeSessionID: dispatchSessionID, + }); err != nil { + return res, fmt.Errorf("enqueue history clear read dispatch: %w", err) + } + editEvent := domain.UpdateEvent{ + UserID: userID, + Type: domain.UpdateEventEditMessage, + Pts: editPts, + PtsCount: 1, + Date: date, + Message: msg, + } + if err := appendUserUpdateEvent(ctx, db, q, userID, editEvent); err != nil { + return res, fmt.Errorf("append history clear edit event: %w", err) + } + if err := enqueueDispatch(ctx, q, sqlcgen.EnqueueDispatchParams{ + TargetUserID: userID, + Pts: int32(editPts), + EventType: string(domain.UpdateEventEditMessage), + ExcludeAuthKeyID: authKeyIDToInt64(dispatchAuthKeyID), + ExcludeSessionID: dispatchSessionID, + }); err != nil { + return res, fmt.Errorf("enqueue history clear edit dispatch: %w", err) + } + item.Events = append(item.Events, readEvent, editEvent) + cursor = editPts + } + if cursor != pts { + return res, fmt.Errorf("delete history pts cursor %d does not reach reserved pts %d", cursor, pts) + } + res.Deleted = append(res.Deleted, item) } return res, nil } +func materializeHistoryClearAnchorTx(ctx context.Context, db sqlcgen.DBTX, anchor historyClearAnchor, pts int) (domain.Message, error) { + msg := domain.NewHistoryClearMessage(anchor.userID, anchor.peer, anchor.boxID, anchor.uid, anchor.messageDate, pts) + mediaJSON, err := encodeMessageMedia(msg.Media) + if err != nil { + return domain.Message{}, fmt.Errorf("encode history clear media: %w", err) + } + tag, err := db.Exec(ctx, ` +UPDATE message_boxes +SET from_user_id = $3, + ttl_period = 0, + expires_at = 0, + edit_date = 0, + hide_edited = false, + outgoing = true, + body = '', + entities = '[]'::jsonb, + silent = false, + noforwards = false, + reply_to_msg_id = 0, + reply_to_peer_type = '', + reply_to_peer_id = 0, + reply_to_top_id = 0, + reply_to_story_id = 0, + quote_text = '', + quote_entities = '[]'::jsonb, + quote_offset = 0, + fwd_from_peer_type = '', + fwd_from_peer_id = 0, + fwd_from_name = '', + fwd_date = 0, + fwd_saved_from_peer_type = '', + fwd_saved_from_peer_id = 0, + fwd_saved_from_msg_id = 0, + saved_peer_type = '', + saved_peer_id = 0, + pts = $4, + media = $5::jsonb, + media_unread = false, + reaction_unread = false, + pinned = false, + via_bot_id = 0, + grouped_id = 0, + effect = 0, + reply_markup = '{}'::jsonb, + rich_message = '{}'::jsonb +WHERE owner_user_id = $1 + AND box_id = $2 + AND NOT deleted`, anchor.userID, int32(anchor.boxID), anchor.userID, int32(pts), mediaJSON) + if err != nil { + return domain.Message{}, fmt.Errorf("materialize history clear anchor: %w", err) + } + if tag.RowsAffected() != 1 { + return domain.Message{}, fmt.Errorf("materialize history clear anchor: box %d disappeared", anchor.boxID) + } + if _, err := db.Exec(ctx, `DELETE FROM message_box_media WHERE owner_user_id = $1 AND box_id = $2`, anchor.userID, int32(anchor.boxID)); err != nil { + return domain.Message{}, fmt.Errorf("delete history clear media index: %w", err) + } + if _, err := db.Exec(ctx, `DELETE FROM saved_message_reaction_tags WHERE user_id = $1 AND message_box_id = $2`, anchor.userID, int32(anchor.boxID)); err != nil { + return domain.Message{}, fmt.Errorf("delete history clear saved tags: %w", err) + } + return msg, nil +} + func maxDeletedMessageID(ids map[int]struct{}) int { maxID := 0 for id := range ids { @@ -278,23 +490,13 @@ func maxDeletedMessageID(ids map[int]struct{}) int { return maxID } -func rebuildDialogAfterMessageDelete(ctx context.Context, q *sqlcgen.Queries, userID int64, peer domain.Peer, preserveEmpty bool) error { +func rebuildDialogAfterMessageDelete(ctx context.Context, q *sqlcgen.Queries, userID int64, peer domain.Peer) error { top, err := q.TopVisibleMessageBoxByPeer(ctx, sqlcgen.TopVisibleMessageBoxByPeerParams{ OwnerUserID: userID, PeerType: string(peer.Type), PeerID: peer.ID, }) if errors.Is(err, pgx.ErrNoRows) { - if preserveEmpty { - if err := q.ClearDialogAfterHistoryDelete(ctx, sqlcgen.ClearDialogAfterHistoryDeleteParams{ - UserID: userID, - PeerType: string(peer.Type), - PeerID: peer.ID, - }); err != nil { - return fmt.Errorf("clear empty dialog after history delete: %w", err) - } - return nil - } if err := q.DeleteDialogByPeer(ctx, sqlcgen.DeleteDialogByPeerParams{ UserID: userID, PeerType: string(peer.Type), diff --git a/internal/store/postgres/message_delete_integration_test.go b/internal/store/postgres/message_delete_integration_test.go index ec47945b..ca1a2026 100644 --- a/internal/store/postgres/message_delete_integration_test.go +++ b/internal/store/postgres/message_delete_integration_test.go @@ -421,7 +421,7 @@ func TestMessageStoreDeleteHistoryRebuildsDialogAndEmitsDeleteUpdates(t *testing } } -func TestMessageStoreDeleteHistoryJustClearPreservesEmptyDialog(t *testing.T) { +func TestMessageStoreDeleteHistoryJustClearPreservesHistoryClearMessage(t *testing.T) { pool := testPool(t) ctx := context.Background() suffix := randomSuffix(t) @@ -444,31 +444,147 @@ func TestMessageStoreDeleteHistoryJustClearPreservesEmptyDialog(t *testing.T) { t.Fatalf("seed send: %v", err) } peer := domain.Peer{Type: domain.PeerTypeUser, ID: peerUser.ID} - if _, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ - OwnerUserID: owner.ID, - Peer: peer, - JustClear: true, - Date: 1700001200, - }); err != nil { + clearResult, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ + OwnerUserID: owner.ID, + Peer: peer, + JustClear: true, + Date: 1700001200, + OriginAuthKeyID: [8]byte{7}, + OriginSessionID: 99, + }) + if err != nil { t.Fatalf("DeleteHistory just_clear: %v", err) } + self := clearResult.Self() + if self.PtsCount != 2 || len(self.MessageIDs) != 0 || len(self.Events) != 2 || + self.Events[0].Type != domain.UpdateEventReadHistoryInbox || + self.Events[1].Type != domain.UpdateEventEditMessage { + t.Fatalf("clear result = %+v, want read+edit", self) + } dialogs, err := NewDialogStore(pool).ListByUser(ctx, owner.ID, domain.DialogFilter{Limit: 10}) if err != nil { t.Fatalf("dialogs after just_clear: %v", err) } - if len(dialogs.Dialogs) != 1 || dialogs.Dialogs[0].Peer != peer || dialogs.Dialogs[0].TopMessage != 0 || len(dialogs.Messages) != 0 { - t.Fatalf("dialogs = %+v messages=%+v, want empty dialog preserved after just_clear", dialogs.Dialogs, dialogs.Messages) + if len(dialogs.Dialogs) != 1 || dialogs.Dialogs[0].Peer != peer || + dialogs.Dialogs[0].TopMessage == 0 || len(dialogs.Messages) != 1 || + !domain.IsHistoryClearServiceMessage(dialogs.Messages[0]) { + t.Fatalf("dialogs = %+v messages=%+v, want real history-clear top", dialogs.Dialogs, dialogs.Messages) } history, err := messages.ListByUser(ctx, owner.ID, domain.MessageFilter{HasPeer: true, Peer: peer, Limit: 10, NeedTotalCount: true}) if err != nil { t.Fatalf("history after just_clear: %v", err) } - if len(history.Messages) != 0 { - t.Fatalf("history = %+v, want cleared", history.Messages) + if len(history.Messages) != 1 || history.Messages[0].ID != dialogs.Dialogs[0].TopMessage || + !domain.IsHistoryClearServiceMessage(history.Messages[0]) || + !history.Messages[0].Out || history.Messages[0].From.ID != owner.ID || + history.Messages[0].Body != "" || history.Messages[0].MediaUnread || + history.Messages[0].ReactionUnread || history.Messages[0].Pinned { + t.Fatalf("history = %+v, want clean owner-local history-clear anchor", history.Messages) + } + events, err := NewUpdateEventStore(pool).ListAfter(ctx, owner.ID, self.Pts-2, 10) + if err != nil { + t.Fatalf("list clear events: %v", err) + } + if len(events) != 2 || events[0].Type != domain.UpdateEventReadHistoryInbox || + events[1].Type != domain.UpdateEventEditMessage || + !domain.IsHistoryClearServiceMessage(events[1].Message) { + t.Fatalf("durable clear events = %+v, want read/edit with service message", events) + } + var outboxRows int + if err := pool.QueryRow(ctx, ` +SELECT count(*)::int +FROM dispatch_outbox +WHERE target_user_id = $1 + AND pts = ANY($2::int[]) + AND exclude_auth_key_id = $3 + AND exclude_session_id = $4`, + owner.ID, []int32{int32(self.Pts - 1), int32(self.Pts)}, int64(7), int64(99)).Scan(&outboxRows); err != nil { + t.Fatalf("count clear outbox: %v", err) + } + if outboxRows != 2 { + t.Fatalf("clear outbox rows = %d, want read/edit excluding origin session", outboxRows) + } + repeated, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ + OwnerUserID: owner.ID, Peer: peer, JustClear: true, Date: 1700001201, + }) + if err != nil { + t.Fatalf("repeat just_clear: %v", err) + } + if repeated.Changed() || len(repeated.Deleted) != 0 { + t.Fatalf("repeat just_clear = %+v, want idempotent no-op", repeated) } } -func TestMessageStoreDeleteHistoryBatchesHugeMaxID(t *testing.T) { +func TestMessageStoreDeleteHistoryJustClearRevokeKeepsPerOwnerAnchors(t *testing.T) { + pool := testPool(t) + ctx := context.Background() + suffix := randomSuffix(t) + users := NewUserStore(pool) + alice := createTestUser(t, ctx, users, "+1994"+suffix+"01", "ClearAlice", "") + bob := createTestUser(t, ctx, users, "+1994"+suffix+"02", "ClearBob", "") + t.Cleanup(func() { + _, _ = pool.Exec(ctx, "DELETE FROM users WHERE id = ANY($1::bigint[])", []int64{alice.ID, bob.ID}) + }) + + messages := NewMessageStore(pool) + var last domain.SendPrivateTextResult + for i := 0; i < 2; i++ { + var err error + last, err = messages.SendPrivateText(ctx, domain.SendPrivateTextRequest{ + SenderUserID: alice.ID, RecipientUserID: bob.ID, RandomID: int64(9200 + i), + Message: "clear for both", Date: 1700001300 + i, + }) + if err != nil { + t.Fatalf("seed send %d: %v", i, err) + } + } + res, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ + OwnerUserID: alice.ID, + Peer: domain.Peer{Type: domain.PeerTypeUser, ID: bob.ID}, + JustClear: true, + Revoke: true, + Date: 1700001400, + }) + if err != nil { + t.Fatalf("revoke just_clear: %v", err) + } + if len(res.Deleted) != 2 { + t.Fatalf("deleted owners = %+v, want alice and bob", res.Deleted) + } + for _, tc := range []struct { + userID int64 + peerID int64 + topID int + }{ + {alice.ID, bob.ID, last.SenderMessage.ID}, + {bob.ID, alice.ID, last.RecipientMessage.ID}, + } { + history, err := messages.ListByUser(ctx, tc.userID, domain.MessageFilter{ + HasPeer: true, Peer: domain.Peer{Type: domain.PeerTypeUser, ID: tc.peerID}, Limit: 10, + }) + if err != nil { + t.Fatalf("history user %d: %v", tc.userID, err) + } + if len(history.Messages) != 1 || history.Messages[0].ID != tc.topID || + !domain.IsHistoryClearServiceMessage(history.Messages[0]) || + !history.Messages[0].Out || history.Messages[0].From.ID != tc.userID { + t.Fatalf("history user %d = %+v, want owner-local anchor %d", tc.userID, history.Messages, tc.topID) + } + var ownerResult domain.DeletedMessagesForUser + for _, item := range res.Deleted { + if item.UserID == tc.userID { + ownerResult = item + break + } + } + if ownerResult.PtsCount != 3 || len(ownerResult.MessageIDs) != 1 || + len(ownerResult.Events) != 3 { + t.Fatalf("owner %d result = %+v, want delete/read/edit", tc.userID, ownerResult) + } + } +} + +func TestMessageStoreDeleteHistoryJustClearKeepsAnchorAcrossBatches(t *testing.T) { pool := testPool(t) ctx := context.Background() suffix := randomSuffix(t) @@ -566,34 +682,47 @@ func TestMessageStoreDeleteHistoryBatchesHugeMaxID(t *testing.T) { first, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ OwnerUserID: owner.ID, Peer: peer, - MaxID: domain.MaxMessageBoxID, + JustClear: true, Date: 1700003000, }) if err != nil { t.Fatalf("DeleteHistory first batch: %v", err) } self := first.Self() - if first.Offset != 1 || self.Event.Pts != domain.MaxDeleteHistoryBatch || self.Event.PtsCount != domain.MaxDeleteHistoryBatch || len(self.MessageIDs) != domain.MaxDeleteHistoryBatch { - t.Fatalf("first batch = %+v self=%+v, want offset=1 and exactly %d deleted ids", first, self, domain.MaxDeleteHistoryBatch) + if first.Offset != 1 || self.Event.Pts != domain.MaxDeleteHistoryBatch || + self.Event.PtsCount != domain.MaxDeleteHistoryBatch || + self.Pts != domain.MaxDeleteHistoryBatch+2 || self.PtsCount != domain.MaxDeleteHistoryBatch+2 || + len(self.MessageIDs) != domain.MaxDeleteHistoryBatch || len(self.Events) != 3 { + t.Fatalf("first batch = %+v self=%+v, want delete batch plus read/edit", first, self) } history, err := messages.ListByUser(ctx, owner.ID, domain.MessageFilter{HasPeer: true, Peer: peer, Limit: 10, NeedTotalCount: true}) if err != nil { t.Fatalf("history after first batch: %v", err) } - if history.Count != 2 || len(history.Messages) != 2 || history.Messages[0].ID != 2 { - t.Fatalf("history after first batch = %+v, want only two oldest messages left", history) + if history.Count != 2 || len(history.Messages) != 2 || history.Messages[0].ID != total || + !domain.IsHistoryClearServiceMessage(history.Messages[0]) || history.Messages[1].ID != 1 { + t.Fatalf("history after first batch = %+v, want stable top anchor plus oldest message", history) } second, err := messages.DeleteHistory(ctx, domain.DeleteHistoryRequest{ OwnerUserID: owner.ID, Peer: peer, - MaxID: domain.MaxMessageBoxID, + JustClear: true, Date: 1700003001, }) if err != nil { t.Fatalf("DeleteHistory second batch: %v", err) } - if second.Offset != 0 || second.Self().Event.PtsCount != 2 { - t.Fatalf("second batch = %+v, want final offset=0 pts_count=2", second) + if second.Offset != 0 || second.Self().Event.PtsCount != 1 || + second.Self().PtsCount != 1 || len(second.Self().Events) != 1 { + t.Fatalf("second batch = %+v, want final old-message delete only", second) + } + history, err = messages.ListByUser(ctx, owner.ID, domain.MessageFilter{HasPeer: true, Peer: peer, Limit: 10, NeedTotalCount: true}) + if err != nil { + t.Fatalf("history after second batch: %v", err) + } + if history.Count != 1 || len(history.Messages) != 1 || history.Messages[0].ID != total || + !domain.IsHistoryClearServiceMessage(history.Messages[0]) { + t.Fatalf("history after second batch = %+v, want only stable history-clear anchor", history) } } diff --git a/internal/store/postgres/message_history.go b/internal/store/postgres/message_history.go index 92372430..c5eadb48 100644 --- a/internal/store/postgres/message_history.go +++ b/internal/store/postgres/message_history.go @@ -520,10 +520,36 @@ func (s *MessageStore) DeleteHistory(ctx context.Context, req domain.DeleteHisto maxID := pgInt32NonNegative(req.MaxID) minDate := pgInt32NonNegative(req.MinDate) maxDate := pgInt32NonNegative(req.MaxDate) + var anchors map[int64]historyClearAnchor + // TDesktop 的完整 Clear history 形态固定为 just_clear + max_id=0 且 + // 不带日期范围;按日期删除只销毁选中区间,不在本地生成 history-clear + // 服务消息。只有完整清空才跨批保留 top 锚点。 + fullJustClear := req.JustClear && req.MaxID <= 0 && req.MinDate <= 0 && req.MaxDate <= 0 + if fullJustClear { + anchors = make(map[int64]historyClearAnchor, 2) + if anchor, found, err := s.loadHistoryClearAnchor(ctx, qtx, req.OwnerUserID, req.Peer); err != nil { + return res, err + } else if found { + anchors[req.OwnerUserID] = anchor + } + if req.Revoke && req.Peer.ID != req.OwnerUserID { + peer := domain.Peer{Type: domain.PeerTypeUser, ID: req.OwnerUserID} + if anchor, found, err := s.loadHistoryClearAnchor(ctx, qtx, req.Peer.ID, peer); err != nil { + return res, err + } else if found { + anchors[req.Peer.ID] = anchor + } + } + } + ownerKeepBoxID := int32(0) + if anchor, ok := anchors[req.OwnerUserID]; ok { + ownerKeepBoxID = int32(anchor.boxID) + } rows, err := qtx.DeleteMessageBoxesByPeerBatch(ctx, sqlcgen.DeleteMessageBoxesByPeerBatchParams{ OwnerUserID: req.OwnerUserID, PeerType: string(req.Peer.Type), PeerID: req.Peer.ID, + KeepBoxID: ownerKeepBoxID, MaxID: maxID, MinDate: minDate, MaxDate: maxDate, @@ -534,7 +560,10 @@ func (s *MessageStore) DeleteHistory(ctx context.Context, req domain.DeleteHisto } deleted := deletedRowsFromPeerBatchRows(rows) if req.Revoke { - if len(deleted) > 0 { + // max_id>0 是 owner-local box 边界,只能通过逻辑 private-message + // 映射删除对端副本。max_id=0 则双方按同一日期范围各扫一批,避免 + // linked delete 与双方保留锚点互相删除。 + if req.MaxID > 0 && len(deleted) > 0 { peerRows, err := qtx.DeleteMessageBoxesByPrivateMessages(ctx, privateMessageDeleteParams(deleted)) if err != nil { return res, fmt.Errorf("delete revoked private history boxes: %w", err) @@ -547,10 +576,15 @@ func (s *MessageStore) DeleteHistory(ctx context.Context, req domain.DeleteHisto // 适用同一区间,box_id 上限则是 owner 私有序无法映射,部分 // max_id 清史保持反查模型(官方 UI 无此入口)。 if req.MaxID <= 0 && req.Peer.ID != req.OwnerUserID { + peerKeepBoxID := int32(0) + if anchor, ok := anchors[req.Peer.ID]; ok { + peerKeepBoxID = int32(anchor.boxID) + } peerSideRows, err := qtx.DeleteMessageBoxesByPeerBatch(ctx, sqlcgen.DeleteMessageBoxesByPeerBatchParams{ OwnerUserID: req.Peer.ID, PeerType: string(domain.PeerTypeUser), PeerID: req.OwnerUserID, + KeepBoxID: peerKeepBoxID, MaxID: 0, MinDate: minDate, MaxDate: maxDate, @@ -562,7 +596,7 @@ func (s *MessageStore) DeleteHistory(ctx context.Context, req domain.DeleteHisto deleted = append(deleted, deletedRowsFromPeerBatchRows(peerSideRows)...) } } - res, err = s.finishDeleteMessagesTx(ctx, tx, qtx, req.OwnerUserID, req.OriginAuthKeyID, req.OriginSessionID, req.Date, deleted, req.JustClear) + res, err = s.finishDeleteMessagesTx(ctx, tx, qtx, req.OwnerUserID, req.OriginAuthKeyID, req.OriginSessionID, req.Date, deleted, anchors) if err != nil { return res, err } @@ -570,6 +604,7 @@ func (s *MessageStore) DeleteHistory(ctx context.Context, req domain.DeleteHisto OwnerUserID: req.OwnerUserID, PeerType: string(req.Peer.Type), PeerID: req.Peer.ID, + KeepBoxID: ownerKeepBoxID, MaxID: maxID, MinDate: minDate, MaxDate: maxDate, @@ -578,10 +613,15 @@ func (s *MessageStore) DeleteHistory(ctx context.Context, req domain.DeleteHisto return res, fmt.Errorf("check remaining history after delete: %w", err) } if !more && req.Revoke && req.MaxID <= 0 && req.Peer.ID != req.OwnerUserID { + peerKeepBoxID := int32(0) + if anchor, ok := anchors[req.Peer.ID]; ok { + peerKeepBoxID = int32(anchor.boxID) + } more, err = qtx.HasDeletableMessageBoxByPeer(ctx, sqlcgen.HasDeletableMessageBoxByPeerParams{ OwnerUserID: req.Peer.ID, PeerType: string(domain.PeerTypeUser), PeerID: req.OwnerUserID, + KeepBoxID: peerKeepBoxID, MaxID: 0, MinDate: minDate, MaxDate: maxDate, diff --git a/internal/store/postgres/queries/dialog.sql b/internal/store/postgres/queries/dialog.sql index 142e444f..5510d79b 100644 --- a/internal/store/postgres/queries/dialog.sql +++ b/internal/store/postgres/queries/dialog.sql @@ -880,30 +880,6 @@ WHERE d.user_id = sqlc.arg(user_id)::bigint AND d.peer_type = sqlc.arg(peer_type)::text AND d.peer_id = sqlc.arg(peer_id)::bigint; --- name: ClearDialogAfterHistoryDelete :exec -UPDATE dialogs d -SET - top_message_id = 0, - top_message_date = 0, - read_inbox_max_id = GREATEST(d.read_inbox_max_id, d.top_message_id), - read_outbox_max_id = GREATEST(d.read_outbox_max_id, d.top_message_id), - unread_count = 0, - unread_mark = false, - unread_mentions_count = 0, - unread_reactions_count = ( - SELECT COUNT(*)::int - FROM message_boxes m2 - WHERE m2.owner_user_id = d.user_id - AND m2.peer_type = d.peer_type - AND m2.peer_id = d.peer_id - AND NOT m2.deleted - AND m2.reaction_unread - ), - updated_at = now() -WHERE d.user_id = sqlc.arg(user_id)::bigint - AND d.peer_type = sqlc.arg(peer_type)::text - AND d.peer_id = sqlc.arg(peer_id)::bigint; - -- name: DeleteDialogByPeer :exec WITH dropped_drafts AS ( -- 删除会话同时丢弃该 peer 的云草稿,避免对端重建会话后旧草稿复活。 diff --git a/internal/store/postgres/queries/message.sql b/internal/store/postgres/queries/message.sql index 920d1af6..c29a8029 100644 --- a/internal/store/postgres/queries/message.sql +++ b/internal/store/postgres/queries/message.sql @@ -1343,6 +1343,7 @@ WITH target AS ( WHERE m.owner_user_id = sqlc.arg(owner_user_id)::bigint AND m.peer_type = sqlc.arg(peer_type)::text AND m.peer_id = sqlc.arg(peer_id)::bigint + AND (sqlc.arg(keep_box_id)::int <= 0 OR m.box_id <> sqlc.arg(keep_box_id)::int) AND (sqlc.arg(max_id)::int <= 0 OR m.box_id <= sqlc.arg(max_id)::int) AND (sqlc.arg(min_date)::int <= 0 OR m.message_date >= sqlc.arg(min_date)::int) AND (sqlc.arg(max_date)::int <= 0 OR m.message_date <= sqlc.arg(max_date)::int) @@ -1382,6 +1383,7 @@ SELECT EXISTS ( WHERE m.owner_user_id = sqlc.arg(owner_user_id)::bigint AND m.peer_type = sqlc.arg(peer_type)::text AND m.peer_id = sqlc.arg(peer_id)::bigint + AND (sqlc.arg(keep_box_id)::int <= 0 OR m.box_id <> sqlc.arg(keep_box_id)::int) AND (sqlc.arg(max_id)::int <= 0 OR m.box_id <= sqlc.arg(max_id)::int) AND (sqlc.arg(min_date)::int <= 0 OR m.message_date >= sqlc.arg(min_date)::int) AND (sqlc.arg(max_date)::int <= 0 OR m.message_date <= sqlc.arg(max_date)::int) diff --git a/internal/store/postgres/saved_dialog.go b/internal/store/postgres/saved_dialog.go index 87d42f5b..07a9fddd 100644 --- a/internal/store/postgres/saved_dialog.go +++ b/internal/store/postgres/saved_dialog.go @@ -299,7 +299,7 @@ func (s *MessageStore) DeleteSavedHistory(ctx context.Context, req domain.Delete peer: domain.Peer{Type: domain.PeerType(row.PeerType), ID: row.PeerID}, }) } - delRes, err := s.finishDeleteMessagesTx(ctx, tx, qtx, req.OwnerUserID, req.OriginAuthKeyID, req.OriginSessionID, req.Date, deleted, false) + delRes, err := s.finishDeleteMessagesTx(ctx, tx, qtx, req.OwnerUserID, req.OriginAuthKeyID, req.OriginSessionID, req.Date, deleted, nil) if err != nil { return res, err } diff --git a/internal/store/postgres/sqlcgen/dialog.sql.go b/internal/store/postgres/sqlcgen/dialog.sql.go index 794ddc95..56ba2477 100644 --- a/internal/store/postgres/sqlcgen/dialog.sql.go +++ b/internal/store/postgres/sqlcgen/dialog.sql.go @@ -35,42 +35,6 @@ func (q *Queries) AdvanceDialogReadInboxFloor(ctx context.Context, arg AdvanceDi return err } -const clearDialogAfterHistoryDelete = `-- name: ClearDialogAfterHistoryDelete :exec -UPDATE dialogs d -SET - top_message_id = 0, - top_message_date = 0, - read_inbox_max_id = GREATEST(d.read_inbox_max_id, d.top_message_id), - read_outbox_max_id = GREATEST(d.read_outbox_max_id, d.top_message_id), - unread_count = 0, - unread_mark = false, - unread_mentions_count = 0, - unread_reactions_count = ( - SELECT COUNT(*)::int - FROM message_boxes m2 - WHERE m2.owner_user_id = d.user_id - AND m2.peer_type = d.peer_type - AND m2.peer_id = d.peer_id - AND NOT m2.deleted - AND m2.reaction_unread - ), - updated_at = now() -WHERE d.user_id = $1::bigint - AND d.peer_type = $2::text - AND d.peer_id = $3::bigint -` - -type ClearDialogAfterHistoryDeleteParams struct { - UserID int64 - PeerType string - PeerID int64 -} - -func (q *Queries) ClearDialogAfterHistoryDelete(ctx context.Context, arg ClearDialogAfterHistoryDeleteParams) error { - _, err := q.db.Exec(ctx, clearDialogAfterHistoryDelete, arg.UserID, arg.PeerType, arg.PeerID) - return err -} - const clearDialogDrafts = `-- name: ClearDialogDrafts :many WITH doomed AS ( SELECT d.user_id, d.peer_type, d.peer_id, d.top_message_id diff --git a/internal/store/postgres/sqlcgen/message.sql.go b/internal/store/postgres/sqlcgen/message.sql.go index b86bfec9..d14d7f8b 100644 --- a/internal/store/postgres/sqlcgen/message.sql.go +++ b/internal/store/postgres/sqlcgen/message.sql.go @@ -886,12 +886,13 @@ WITH target AS ( WHERE m.owner_user_id = $1::bigint AND m.peer_type = $2::text AND m.peer_id = $3::bigint - AND ($4::int <= 0 OR m.box_id <= $4::int) - AND ($5::int <= 0 OR m.message_date >= $5::int) - AND ($6::int <= 0 OR m.message_date <= $6::int) + AND ($4::int <= 0 OR m.box_id <> $4::int) + AND ($5::int <= 0 OR m.box_id <= $5::int) + AND ($6::int <= 0 OR m.message_date >= $6::int) + AND ($7::int <= 0 OR m.message_date <= $7::int) AND NOT m.deleted ORDER BY m.box_id DESC - LIMIT $7::int + LIMIT $8::int FOR UPDATE SKIP LOCKED ), updated AS ( @@ -923,6 +924,7 @@ type DeleteMessageBoxesByPeerBatchParams struct { OwnerUserID int64 PeerType string PeerID int64 + KeepBoxID int32 MaxID int32 MinDate int32 MaxDate int32 @@ -943,6 +945,7 @@ func (q *Queries) DeleteMessageBoxesByPeerBatch(ctx context.Context, arg DeleteM arg.OwnerUserID, arg.PeerType, arg.PeerID, + arg.KeepBoxID, arg.MaxID, arg.MinDate, arg.MaxDate, @@ -2108,9 +2111,10 @@ SELECT EXISTS ( WHERE m.owner_user_id = $1::bigint AND m.peer_type = $2::text AND m.peer_id = $3::bigint - AND ($4::int <= 0 OR m.box_id <= $4::int) - AND ($5::int <= 0 OR m.message_date >= $5::int) - AND ($6::int <= 0 OR m.message_date <= $6::int) + AND ($4::int <= 0 OR m.box_id <> $4::int) + AND ($5::int <= 0 OR m.box_id <= $5::int) + AND ($6::int <= 0 OR m.message_date >= $6::int) + AND ($7::int <= 0 OR m.message_date <= $7::int) AND NOT m.deleted LIMIT 1 )::boolean AS more @@ -2120,6 +2124,7 @@ type HasDeletableMessageBoxByPeerParams struct { OwnerUserID int64 PeerType string PeerID int64 + KeepBoxID int32 MaxID int32 MinDate int32 MaxDate int32 @@ -2130,6 +2135,7 @@ func (q *Queries) HasDeletableMessageBoxByPeer(ctx context.Context, arg HasDelet arg.OwnerUserID, arg.PeerType, arg.PeerID, + arg.KeepBoxID, arg.MaxID, arg.MinDate, arg.MaxDate,