From 23a2b2aff773c3721b6adaab612af5ee412848d1 Mon Sep 17 00:00:00 2001 From: A Date: Sun, 7 Jun 2026 20:28:37 +0800 Subject: [PATCH] fix: restore sticker placeholders and media history (cherry picked from commit 488e409a1898e9c739cc0bd24cb9791636dfd6b3) --- docs/compatibility-matrix.md | 1 + docs/performance-audit.md | 2 +- internal/app/files/seed.go | 57 +----- internal/app/files/seed_test.go | 51 +---- internal/rpc/messages.go | 25 ++- internal/rpc/send_media.go | 4 + internal/store/postgres/channel.go | 177 +++++++++++++----- internal/store/postgres/dialog.go | 10 + internal/store/postgres/queries/dialog.sql | 10 +- internal/store/postgres/sqlcgen/dialog.sql.go | 16 +- internal/store/postgres/sqlcgen/models.go | 37 ++++ 11 files changed, 225 insertions(+), 165 deletions(-) diff --git a/docs/compatibility-matrix.md b/docs/compatibility-matrix.md index 91c4f1cd..343de35c 100644 --- a/docs/compatibility-matrix.md +++ b/docs/compatibility-matrix.md @@ -41,6 +41,7 @@ status 取值:done(真实实现) / stub(兼容响应) / todo(已发现未实 > 2026-06-03 sticker installed-cache 语义修复:继续对照 TDesktop `Storage::Account::writeInstalledStickers` / `Data::ParseStickersSetFlags` 后确认,`installed_date` 会让 set 进入 Installed;普通 installed stickers 集若仍处于 `NotLoaded`,客户端会中止本地 installed stickers 写入。telesrv seed 过去把 `InputStickerSetAnimatedEmoji` / dice / generic animations 等 system set 也标成 installed,导致系统资源混入普通 installed stickers 缓存路径。现 migration `0066_system_sticker_sets_not_installed` 修复既有库,并改 seed:`set_kind=system` 仅作为按 input 系统 key 解析的内置资源,不再声明为 installed。 > 2026-06-03 sticker 静态缩略图元数据修复:继续读 TDesktop `history_view_sticker.cpp` / `data_document_media.cpp` / `data_cloud_file.cpp` 后确认,历史页若看到 document `PhotoPathSize` 会把它当成 vector placeholder,`dataMediaCreated()` 因 `thumbnailPath()` 非空跳过 `thumbnailWanted()`,从而挡住同时存在的真实 raster thumb。seed 现将 ≤32KB 可下载 document thumb 转为 `PhotoCachedSize` 并写入 image cache 所需 bytes;当 document 已有 raster/default/cached/progressive thumb 时丢弃 `PhotoPathSize`,仅没有 raster 的 set cover 继续保留 path 占位;thumb blob MIME 由字节魔数写入,`upload.getFile` 也优先按魔数返回 `storage.fileWebp`,兼容旧库误标 `image/jpeg`。验证:`go test ./internal/app/files -run "TestSeedMediaFromRealExport|TestSeedPreferRasterDocumentThumbsDropsPathWhenRasterExists|TestDocumentsNeedInlineCachedThumbsDetectsPathWithRaster|TestDocumentsNeedInlineCachedThumbsDetectsStaleMime" -count=1`、`go test ./internal/rpc -run "TestStorageFileType" -count=1`、`go test ./... -count=1` 通过;新 server 启动 repair 后 PG `documents.thumbs` kind 分布为 `cached=2313`,两条样例 `5415908822013185960/5381935901284774004` 均仅 `cached/m`,`file_blobs` document thumb MIME 为 `image/webp=2313`;服务进程 PID 53892 监听 2398。 > 2026-06-03 sticker 文档身份修复:继续对比 TDesktop `data_session.cpp` / `data_document.cpp` / `history_view_sticker.cpp` 后确认,TDesktop 以 `document_id` 复用 `DocumentData`,且 `DocumentData::updateThumbnails()` 不会清掉旧的 inline/path thumbnail 状态;一旦服务端把外部导出 document id 当成本服资源主键,旧 Debug tdata 中同 id 的污染对象会持续影响历史页渲染。telesrv 现把导出 JSON/文件名中的 document id 只作为 seed source id,在 `internal/app/files` 导入阶段归一为 telesrv-owned storage id;RPC、`InputDocument`、`inputDocumentFileLocation`、`getCustomEmojiDocuments`、channel custom emoji status/reaction/color 均直接使用同一个服务端 document id,不再做边界映射。migration `0067_seed_document_id_namespace` 同步修复既有开发库中的 documents/file_blobs/sticker_sets/available_reactions/message media/channel appearance 引用。document thumb 在 store 中可保留 cached bytes,但 RPC 对 document 统一暴露 downloadable `photoSize m`,避免 `PhotoCachedSize` 与本地旧 cache 组合出不可替换状态。验证:`go test ./internal/app/files -run TestSeedMediaFromRealExport -count=1 -v`、`go test ./... -count=1` 通过;server `bin\\telesrv.exe` 监听 2398,migration 状态 `67|false`,PG `documents/file_blobs/message media/available_reactions/sticker_sets` 均无 `>4e18` 外部 document id;Computer Use 重启 Debug/Alice 与 DebugBob/Bob 后分别打开 Bob B/Alice A,250ms 截图已显示 sticker,server 与两端客户端日志无 `NOT_IMPLEMENTED` / `Unhandled RPC` / `bad_msg` / `LOCATION_INVALID` / `API Error` / sticker hash 异常。 +> 2026-06-07 sticker 静态图慢复盘 / 撤销 a.2 的 path 丢弃:再读 TDesktop `history_view_sticker.cpp`(`draw→paintPixmap→paintPath` 占位优先级,`thumbnailInline` 内联分支被源码注释禁用)与 `data_document_media.cpp`(`checkStickerLarge→automaticLoad` 下载完整 `.tgs`,与 path/thumbnail 无关)后确认:上方 2026-06-03「sticker 静态缩略图元数据修复」丢弃 `PhotoPathSize` 是误判(真根因是同日 document id 污染,已由 `0067` 修复);id 修复后再丢 path 反而使 animated sticker 失去唯一即时占位 ⇒ 打开会话先空白,并因 `thumbnailPath()` 变空多触发一次缩略图 `getFile`,与完整 `.tgs` 争用下载通道 ⇒「静态图特别慢、不像本地加载」。现 `internal/app/files/seed.go` 恢复 `PhotoPathSize` 占位(删 `seedPreferRasterDocumentThumbs` 及 path+raster stale 判定),document.thumbs = path + downloadable `photoSize m`。已实测排除磁盘 IO/缓存容量(blob 4330 文件 71.3MB,几乎全 ≤64KB)与 `document.Size` 不匹配(28/28 一致)。验证:`go test ./internal/app/files ./internal/rpc -count=1`、`TestSeedMediaFromRealExport -count=1`(真实导出 74 reactions/11 sets/1355 docs/2710 blobs,断言 seed 后 document 保留 path)、`go vet` 通过。**现有库需重新 seed 一次**:`DELETE FROM sticker_sets; DELETE FROM available_reactions;` 重启触发全量 reseed(`PutDocument`/`PutFileBlob` upsert,不删用户上传媒体)。待双 TDesktop 人工验证渲染体感。 ## Transport / MTProto 服务消息 diff --git a/docs/performance-audit.md b/docs/performance-audit.md index 046c1cad..84917a38 100644 --- a/docs/performance-audit.md +++ b/docs/performance-audit.md @@ -142,7 +142,7 @@ - **P1-媒体-a `upload.getFile` 整文件入内存 + 每 chunk 查 PG / sticker 小资源首开冷路径** — ✅ **已修(2026-06-03)**:原 `blobs.Get` 一次 `os.ReadFile` 读整个 blob 再切片,客户端按 ≤512KB/1MB 分块多次请求 ⇒ 同一大文件被重复整读 N 次(O(N²))+ 每 chunk 一次 `GetFileBlob` PG 往返;sticker/reaction/thumb 虽小,但 TDesktop 重启或打开历史首次渲染时仍可能在 `messages.getStickerSet` + `upload.getFile` + 本地文件冷读上形成可感知卡顿。现 `BlobBackend.GetRange`/`LocalFS`(`internal/app/files/blobfs.go`)用 `ReadAt` 只读 offset+limit 段(`n` 受文件大小约束,超大 limit 不会按客户端值分配内存);并加 `location_key→FileBlob` 进程内 LRU(容量 65536,元数据不可变故只读填充无需失效)消除每 chunk 的 PG 查;再加 `object_key→bytes` 小 blob LRU(单项 ≤256KB,总 64MB)、完整 sticker set cache 与启动 `WarmCaches`,从已 seed 元数据预热 sticker/reaction document 和可下载缩略图。参考实现 用 MinIO ranged read(`object.ReadAt`)+ SSDB 两级,但其下载命中冷存储不回填、无 LRU/阈值、photo/头像直连 MinIO;telesrv 当前先做进程内元数据/小字节 LRU + 段读,**多实例共享缓存(Redis)仍待二阶段换对象存储时评估**。实测启动预热 `49 sets / 2313 docs / 2298 blobs`,TDesktop `messages.getStickerSet` 从约 18-24ms 降至 0-1.7ms;seed document id 归一后,首次新服务端 id 会触发必要的 `upload.getFile` 主体/缩略图请求,但 server 端大多 0.5-6ms、样例最长 59ms,重复打开 Bob 会话 0.58s 截图已可见 sticker。 - **P1-媒体-a.1 system sticker set 污染 installed stickers 本地索引** — ✅ **已修(2026-06-03)**:TDesktop 贴纸正文缓存 key 为 `dc_id + document_id`,但 installed stickers 本地索引写入依赖 set flags;`installed_date` 会使 set 进入 Installed,普通 stickers 类型的 `Installed + NotLoaded` set 会让 `writeInstalledStickers()` 中止。此前 seed 把 animated emoji/dice/generic animations 这类 system set 也标成 installed,可能导致重启后 installed sticker 索引反复失效与 set 元数据重拉。现 migration `0066_system_sticker_sets_not_installed` 修复既有库,seed 后续也不再把 `set_kind=system` 声明为 installed。 -- **P1-媒体-a.2 sticker document 同时暴露 `PhotoPathSize` 与 raster thumb,历史页静态图被占位短路** — ✅ **已修(2026-06-03)**:TDesktop `history_view_sticker.cpp` 在创建 sticker media view 时若 `thumbnailPath()` 非空就不调用 `thumbnailWanted()`,而 `PhotoPathSize` 只是 vector placeholder;此前 seed 给同一 document 同时返回 path + raster thumb,导致历史页持续画 path/等完整 TGS 解码,用户反复重启 Debug 仍看不到静态图。现 document 有 default/cached/progressive raster thumb 时过滤 `PhotoPathSize`,thumb blob MIME 按魔数写入;domain 中可用 cached bytes 做修复依据,但 RPC 对 document 统一暴露 downloadable `photoSize m`,避免 `PhotoCachedSize` 与旧本地 cache 叠加出不可替换状态。启动 repair 后 PG 中 document thumbs kind 仅剩 `cached=2313`,document thumb `file_blobs.mime_type` 全部为 `image/webp`。 +- **P1-媒体-a.2 历史页 sticker 静态图慢 / 打开会话先空白** — ✅ **已修正(2026-06-07,撤销 2026-06-03 的过度修复)**:根因复盘 — TDesktop `history_view_sticker.cpp` 对 animated sticker 的占位渲染优先级为 `getStickerLarge()`(完整 `.tgs` 首帧,需下载)→ ~~`thumbnailInline()`(stripped 内联,源码显式注释禁用:sticker 需 alpha 通道)~~ → `goodThumbnail()` → 下载的 `thumbnail()` → **`paintPath()`(`PhotoPathSize` 矢量轮廓,内联即时,唯一不依赖下载的占位)** → 空白;且完整 `.tgs` 由 `DocumentMedia::checkStickerLarge()`→`automaticLoad()` 下载,**与 path/thumbnail 无关**。2026-06-03 的 a.2 误把"卡在 path 等完整 TGS"当作 path 短路(真根因是 a.3 的 document id 污染让完整文件下载/缓存失败),改为有 raster 时过滤 `PhotoPathSize`。a.3 修掉 id 后,这个过滤变成净负面:document 失去唯一即时占位 ⇒ 打开会话 sticker 先空白;且 `thumbnailPath()` 变空触发 `thumbnailWanted()` ⇒ 每个 sticker 多下载一次缩略图,与完整 `.tgs` 争用下载通道,几十个 sticker 排队 ⇒ "静态图特别慢、不像本地加载"。现 seed 恢复 `PhotoPathSize` 占位(与官方 sticker 一致,`internal/app/files/seed.go` 不再过滤 path,删 `seedPreferRasterDocumentThumbs`/`documentThumbsHave{Raster,Path}`),document.thumbs = path + downloadable `photoSize m`:打开即时画轮廓占位、有 path 后不再主动下载缩略图、下载通道全给完整 `.tgs`、第二次本地缓存秒显。已实测排除磁盘 IO/缓存容量(blob 4330 文件 71.3MB,几乎全 ≤64KB)与 `document.Size` 不匹配(28/28 一致)。**现有库需重新 seed 一次恢复 path**:`DELETE FROM sticker_sets; DELETE FROM available_reactions;` 后重启触发全量 reseed(`PutDocument`/`PutFileBlob` upsert,不删用户上传媒体)。`TestSeedMediaFromRealExport` 用真实导出验证 seed 后 document 保留 path。**待双 TDesktop 人工验证渲染体感。** - **P1-媒体-a.3 复用导出 document id 命中 TDesktop 旧 `DocumentData` 缓存** — ✅ **已修(2026-06-03)**:TDesktop `Data::Session::document(id)` 按 `document_id` 单例化,`DocumentData::updateThumbnails()` 不会清空旧 inline/path thumbnail;因此 server 修掉坏 thumb 字段后,旧 Debug tdata 仍可能用同 id 的污染对象,打开历史先空白再等完整 TGS。根因是 seed 曾把外部导出 document id 直接作为本服资源主键。现 `internal/app/files` 在导入阶段把 source id 归一为 telesrv-owned storage id,RPC/download/custom emoji/channel appearance 全部直接使用该服务端 id;migration `0067_seed_document_id_namespace` 修复既有库的 documents/file_blobs/sticker_sets/available_reactions/message media/channel appearance 引用。这样不需要改官方客户端,也不在 RPC 层保留客户端特判。复测:migration 后 PG 相关引用均无 `>4e18` 外部 document id;双 TDesktop 重启后打开 Bob B/Alice A 首屏 250ms 已显示 sticker,`messages.getHistory` 约 6-12ms,`upload.getFile` 命中新 id 成功,server/client 日志无 location/hash/API 错误。 - **P1-媒体-b `upload_parts` 无 GC / 无每用户配额** — `internal/app/files/service.go:30` + `deploy/migrations/0057_media.up.sql`:分片直接进 PG `bytea`,仅 `assembleUpload` 成功才删;**未 assemble 的分片永久滞留**,且不同 `file_id` 无上限 ⇒ 任意登录用户可用海量 fileID 各传几片撑爆 PG(容量 DoS)。建议每用户 in-flight 上传字节/分片配额 + 后台按 `created_at` 过期清理 worker。 - **P1-媒体-c `media` JSONB 内联放大 fan-out 写** — `deploy/migrations/0057_media.up.sql:129`:`media` 快照内联在 `message_boxes`(owner 双份)/`channel_messages`,含 stripped thumb/attributes 使单行变大;叠加主线 A 的 channel 全员写扇出时,大群每条媒体消息按成员数复制整个 media JSONB。建议随主线 A 改懒算/单副本时一并评估 media 是否只存引用、下沉 `documents`/`photos`。 diff --git a/internal/app/files/seed.go b/internal/app/files/seed.go index 58165dff..8401a893 100644 --- a/internal/app/files/seed.go +++ b/internal/app/files/seed.go @@ -338,8 +338,13 @@ func (s *Service) importDocument(ctx context.Context, dj seedDocumentJSON, binDi stats.Blobs++ } - // 缩略图:PhotoPathSize 内联;小的 PhotoSize 静态图同时写 blob 并作为 - // PhotoCachedSize 返回,让 TDesktop 处理 document 元数据时即可填本地 image cache。 + // 缩略图保留两类并存(与官方 sticker 一致): + // - PhotoPathSize 矢量轮廓:随 document 元数据内联下发,是 TDesktop 对 animated + // sticker 在完整 .tgs 下载完成前唯一可即时渲染的占位(history_view_sticker 显式 + // 禁用了 stripped 内联占位,cached 字节又经 RPC 出口转成 downloadable);丢掉它会让 + // 打开会话时 sticker 先空白、并多触发一次缩略图 getFile。 + // - 小 PhotoSize 静态图:写 blob 并暂存为 PhotoCachedSize(RPC 出口再转 downloadable + // photoSize m),供 sticker 面板等需要小缩略图的场景下载。 thumbs := make([]domain.PhotoSize, 0, len(dj.Thumbs)) for _, tj := range dj.Thumbs { ps, downloadable := seedPhotoSize(tj) @@ -375,7 +380,7 @@ func (s *Service) importDocument(ctx context.Context, dj seedDocumentJSON, binDi } thumbs = append(thumbs, ps) } - doc.Thumbs = seedPreferRasterDocumentThumbs(thumbs) + doc.Thumbs = thumbs if err := s.media.PutDocument(ctx, doc); err != nil { return domain.Document{}, err @@ -607,20 +612,6 @@ func seedInlineCachedDocumentThumb(ps domain.PhotoSize, data []byte) domain.Phot return ps } -func seedPreferRasterDocumentThumbs(sizes []domain.PhotoSize) []domain.PhotoSize { - if !documentThumbsHaveRaster(sizes) { - return sizes - } - out := sizes[:0] - for _, size := range sizes { - if size.Kind == domain.PhotoSizeKindPath { - continue - } - out = append(out, size) - } - return out -} - func seedThumbMimeType(data []byte) string { switch { case len(data) >= 12 && data[0] == 'R' && data[1] == 'I' && data[2] == 'F' && data[3] == 'F' && @@ -677,9 +668,6 @@ func (s *Service) documentsNeedInlineCachedThumbs(ctx context.Context, ids []int return false, err } for _, doc := range docs { - if documentThumbsHaveRaster(doc.Thumbs) && documentThumbsHavePath(doc.Thumbs) { - return true, nil - } for _, thumb := range doc.Thumbs { if thumb.Kind == domain.PhotoSizeKindDefault && thumb.Size > 0 && thumb.Size <= seedInlineCachedDocumentThumbMaxBytes { return true, nil @@ -701,35 +689,6 @@ func (s *Service) documentsNeedInlineCachedThumbs(ctx context.Context, ids []int return false, nil } -func documentThumbsHaveRaster(sizes []domain.PhotoSize) bool { - for _, size := range sizes { - switch size.Kind { - case domain.PhotoSizeKindCached: - if len(size.Bytes) > 0 { - return true - } - case domain.PhotoSizeKindDefault: - if size.Type != "" && size.Size > 0 { - return true - } - case domain.PhotoSizeKindProgressive: - if size.Type != "" && len(size.Sizes) > 0 { - return true - } - } - } - return false -} - -func documentThumbsHavePath(sizes []domain.PhotoSize) bool { - for _, size := range sizes { - if size.Kind == domain.PhotoSizeKindPath && len(size.Bytes) > 0 { - return true - } - } - return false -} - func seedStickerPacks(setPacks, resultPacks []seedStickerPackJSON, docIDBySource map[int64]int64) []domain.StickerPack { packs := setPacks if len(packs) == 0 { diff --git a/internal/app/files/seed_test.go b/internal/app/files/seed_test.go index 236d27bc..8c408f61 100644 --- a/internal/app/files/seed_test.go +++ b/internal/app/files/seed_test.go @@ -334,8 +334,8 @@ func TestSeedMediaFromRealExport(t *testing.T) { if want := seedThumbMimeType(thumb.Bytes); blob.MimeType != want { t.Fatalf("sample sticker thumb mime = %q, want %q", blob.MimeType, want) } - if hasPathThumb(doc.Thumbs) { - t.Fatalf("sample sticker document still exposes path thumb together with raster: %+v", doc.Thumbs) + if !hasPathThumb(doc.Thumbs) { + t.Fatalf("sample sticker document dropped its PhotoPathSize placeholder: %+v", doc.Thumbs) } } } @@ -397,25 +397,6 @@ func TestSeedThumbMimeType(t *testing.T) { } } -func TestSeedPreferRasterDocumentThumbsDropsPathWhenRasterExists(t *testing.T) { - sizes := []domain.PhotoSize{ - {Kind: domain.PhotoSizeKindPath, Type: "j", Bytes: []byte("path")}, - {Kind: domain.PhotoSizeKindCached, Type: "m", Bytes: []byte("webp")}, - } - got := seedPreferRasterDocumentThumbs(sizes) - if hasPathThumb(got) { - t.Fatalf("path thumb should be dropped when raster exists: %+v", got) - } - if !hasCachedThumb(got) { - t.Fatalf("cached thumb should be kept: %+v", got) - } - - onlyPath := []domain.PhotoSize{{Kind: domain.PhotoSizeKindPath, Type: "j", Bytes: []byte("path")}} - if got := seedPreferRasterDocumentThumbs(onlyPath); !hasPathThumb(got) { - t.Fatalf("path-only thumbs should be kept: %+v", got) - } -} - func TestDocumentsNeedInlineCachedThumbsDetectsStaleMime(t *testing.T) { ctx := context.Background() media := newFakeMediaStore() @@ -453,34 +434,6 @@ func TestDocumentsNeedInlineCachedThumbsDetectsStaleMime(t *testing.T) { } } -func TestDocumentsNeedInlineCachedThumbsDetectsPathWithRaster(t *testing.T) { - ctx := context.Background() - media := newFakeMediaStore() - doc := domain.Document{ - ID: 100, - Thumbs: []domain.PhotoSize{ - {Kind: domain.PhotoSizeKindPath, Type: "j", Bytes: []byte("path")}, - {Kind: domain.PhotoSizeKindCached, Type: "m", Bytes: []byte("webp")}, - }, - } - if err := media.PutDocument(ctx, doc); err != nil { - t.Fatalf("put doc: %v", err) - } - svc := NewService(media, nil, 2) - stale, err := svc.documentsNeedInlineCachedThumbs(ctx, []int64{doc.ID}) - if err != nil { - t.Fatalf("documentsNeedInlineCachedThumbs: %v", err) - } - if !stale { - t.Fatal("path thumb with raster should require repair") - } -} - -func hasCachedThumb(sizes []domain.PhotoSize) bool { - _, ok := findCachedThumb(sizes) - return ok -} - func findCachedThumb(sizes []domain.PhotoSize) (domain.PhotoSize, bool) { for _, size := range sizes { if size.Kind == domain.PhotoSizeKindCached && len(size.Bytes) > 0 { diff --git a/internal/rpc/messages.go b/internal/rpc/messages.go index b490a956..42bbe3e4 100644 --- a/internal/rpc/messages.go +++ b/internal/rpc/messages.go @@ -4348,6 +4348,7 @@ func (r *Router) onMessagesForwardMessages(ctx context.Context, req *tg.Messages RandomID: req.RandomID[i], Message: source.body, Entities: source.entities, + Media: source.media, Silent: req.Silent, NoForwards: req.Noforwards, ReplyTo: replyTo, @@ -4395,6 +4396,7 @@ func (r *Router) onMessagesForwardMessages(ctx context.Context, req *tg.Messages RandomID: req.RandomID[i], Message: source.body, Entities: source.entities, + Media: source.media, Silent: req.Silent, NoForwards: req.Noforwards, ReplyTo: replyTo, @@ -4472,6 +4474,7 @@ func mergeForwardTopMsgID(toPeer domain.Peer, replyTo *domain.MessageReply, topM type forwardSource struct { body string entities []domain.MessageEntity + media *domain.MessageMedia forward *domain.MessageForward from domain.Peer date int @@ -4496,17 +4499,14 @@ func (r *Router) forwardSources(ctx context.Context, userID int64, fromPeer doma if r.deps.Messages == nil { return nil, domain.ErrMessageIDInvalid } - list, err := r.deps.Messages.GetHistory(ctx, userID, domain.MessageFilter{ - HasPeer: true, - Peer: fromPeer, - Limit: 1, - MaxID: id, - MinID: id - 1, - }) + list, err := r.deps.Messages.GetMessages(ctx, userID, []int{id}) if err != nil || len(list.Messages) != 1 || list.Messages[0].ID != id { return nil, domain.ErrMessageIDInvalid } msg := list.Messages[0] + if msg.Peer != fromPeer { + return nil, domain.ErrMessageIDInvalid + } if msg.NoForwards { return nil, domain.ErrChatForwardsRestricted } @@ -4518,6 +4518,7 @@ func (r *Router) forwardSources(ctx context.Context, userID int64, fromPeer doma body: msg.Body, entities: append([]domain.MessageEntity(nil), msg.Entities...), + media: msg.Media, forward: forward, from: msg.From, date: msg.Date, @@ -4526,12 +4527,7 @@ func (r *Router) forwardSources(ctx context.Context, userID int64, fromPeer doma if r.deps.Channels == nil { return nil, domain.ErrMessageIDInvalid } - history, err := r.deps.Channels.GetHistory(ctx, userID, domain.ChannelHistoryFilter{ - ChannelID: fromPeer.ID, - Limit: 1, - MaxID: id, - MinID: id - 1, - }) + history, err := r.deps.Channels.GetMessages(ctx, userID, fromPeer.ID, []int{id}) if err != nil || len(history.Messages) != 1 || history.Messages[0].ID != id { return nil, domain.ErrMessageIDInvalid } @@ -4539,7 +4535,7 @@ func (r *Router) forwardSources(ctx context.Context, userID int64, fromPeer doma if msg.NoForwards || history.Channel.NoForwards { return nil, domain.ErrChatForwardsRestricted } - if msg.Body == "" || msg.Action != nil { + if msg.Action != nil || (msg.Body == "" && msg.Media.IsZero()) { return nil, domain.ErrMessageIDInvalid } forward := cloneDomainMessageForward(msg.Forward) @@ -4560,6 +4556,7 @@ func (r *Router) forwardSources(ctx context.Context, userID int64, fromPeer doma body: msg.Body, entities: append([]domain.MessageEntity(nil), msg.Entities...), + media: msg.Media, forward: forward, from: from, date: msg.Date, diff --git a/internal/rpc/send_media.go b/internal/rpc/send_media.go index 9457a5c6..634f8805 100644 --- a/internal/rpc/send_media.go +++ b/internal/rpc/send_media.go @@ -3,10 +3,12 @@ package rpc import ( "context" "errors" + "fmt" "strconv" "unicode/utf8" "github.com/gotd/td/tg" + "go.uber.org/zap" "telesrv/internal/domain" ) @@ -361,6 +363,7 @@ func (r *Router) resolveInputMedia(ctx context.Context, userID int64, input tg.I case *tg.InputMediaDocument: docIDs, ok := inputDocumentCandidateIDs(in.ID) if !ok { + r.log.Warn("sendMedia InputMediaDocument unresolvable id", zap.String("id_type", fmt.Sprintf("%T", in.ID))) return nil, mediaInvalidErr() } var doc domain.Document @@ -376,6 +379,7 @@ func (r *Router) resolveInputMedia(ctx context.Context, userID int64, input tg.I } } if !found { + r.log.Warn("sendMedia references unknown document", zap.Int64s("doc_ids", docIDs), zap.Int64("user_id", userID)) return nil, mediaInvalidErr() } return messageMediaFromDocument(doc, in.Spoiler, in.TTLSeconds), nil diff --git a/internal/store/postgres/channel.go b/internal/store/postgres/channel.go index 8eabd576..c33104ec 100644 --- a/internal/store/postgres/channel.go +++ b/internal/store/postgres/channel.go @@ -4438,71 +4438,158 @@ func (s *ChannelStore) ListChannelHistory(ctx context.Context, viewerUserID int6 if limit <= 0 || limit > 100 { limit = 100 } - args := []any{filter.ChannelID} - where := "channel_id = $1 AND NOT deleted" + // 公共过滤条件(不含 offset 锚点的方向条件,供 add_offset 各模式复用) + baseArgs := []any{filter.ChannelID} + base := "channel_id = $1 AND NOT deleted" if member.AvailableMinID > 0 { - args = append(args, member.AvailableMinID) - where += fmt.Sprintf(" AND id > $%d", len(args)) + baseArgs = append(baseArgs, member.AvailableMinID) + base += fmt.Sprintf(" AND id > $%d", len(baseArgs)) } if filter.Query != "" { - args = append(args, filter.Query) - where += fmt.Sprintf(" AND body ILIKE '%%' || $%d || '%%'", len(args)) + baseArgs = append(baseArgs, filter.Query) + base += fmt.Sprintf(" AND body ILIKE '%%' || $%d || '%%'", len(baseArgs)) } if filter.SenderUserID != 0 { - args = append(args, filter.SenderUserID) - where += fmt.Sprintf(" AND sender_user_id = $%d", len(args)) + baseArgs = append(baseArgs, filter.SenderUserID) + base += fmt.Sprintf(" AND sender_user_id = $%d", len(baseArgs)) } if filter.MinDate > 0 { - args = append(args, filter.MinDate) - where += fmt.Sprintf(" AND message_date > $%d", len(args)) + baseArgs = append(baseArgs, filter.MinDate) + base += fmt.Sprintf(" AND message_date > $%d", len(baseArgs)) } if filter.MaxDate > 0 { - args = append(args, filter.MaxDate) - where += fmt.Sprintf(" AND message_date < $%d", len(args)) - } - if filter.OffsetID > 0 { - args = append(args, filter.OffsetID) - where += fmt.Sprintf(" AND id < $%d", len(args)) - } else if filter.OffsetDate > 0 { - args = append(args, filter.OffsetDate) - where += fmt.Sprintf(" AND message_date < $%d", len(args)) + baseArgs = append(baseArgs, filter.MaxDate) + base += fmt.Sprintf(" AND message_date < $%d", len(baseArgs)) } if filter.MaxID > 0 { - args = append(args, filter.MaxID) - where += fmt.Sprintf(" AND id <= $%d", len(args)) + baseArgs = append(baseArgs, filter.MaxID) + base += fmt.Sprintf(" AND id <= $%d", len(baseArgs)) } if filter.MinID > 0 { - args = append(args, filter.MinID) - where += fmt.Sprintf(" AND id > $%d", len(args)) + baseArgs = append(baseArgs, filter.MinID) + base += fmt.Sprintf(" AND id > $%d", len(baseArgs)) } - queryLimit := limit + 1 - args = append(args, queryLimit) - rows, err := s.db.Query(ctx, ` -SELECT `+channelMessageColumns+` -FROM channel_messages -WHERE `+where+` -ORDER BY id DESC -LIMIT $`+fmt.Sprint(len(args)), args...) - if err != nil { - return domain.ChannelHistory{}, fmt.Errorf("list channel history: %w", err) + scanList := func(sql string, queryArgs []any) ([]domain.ChannelMessage, error) { + rows, err := s.db.Query(ctx, sql, queryArgs...) + if err != nil { + return nil, fmt.Errorf("list channel history: %w", err) + } + defer rows.Close() + var list []domain.ChannelMessage + for rows.Next() { + msg, scanErr := scanChannelMessage(rows) + if scanErr != nil { + return nil, scanErr + } + list = append(list, msg) + } + return list, rows.Err() } - defer rows.Close() + // add_offset 决定加载方向(对齐私聊 ListMessagesByUser): + // >= 0 backward:锚点更旧方向,先跳过 add_offset 条 + // < 0 且 +limit>0 around:以锚点为中心,向更新取 -add_offset 条 + 向更旧取 limit+add_offset 条 + // 否则 forward:仅锚点更新方向(拉未读消息) + addOffset := filter.AddOffset out := domain.ChannelHistory{Channel: channel, Self: member} - for rows.Next() { - msg, err := scanChannelMessage(rows) + hasMoreOlder := false + // 锚点条件:offset_date 优先按日期、否则按消息 id(对齐私聊/orange); + // 二者皆空时向更新方向退化为空、向更旧方向退化为全部(取最新)。 + forwardCond := func(args *[]any) string { + if filter.OffsetDate > 0 { + *args = append(*args, filter.OffsetDate) + return fmt.Sprintf("message_date >= $%d", len(*args)) + } + if filter.OffsetID > 0 { + *args = append(*args, filter.OffsetID) + return fmt.Sprintf("id > $%d", len(*args)) + } + return "false" + } + aroundOlderCond := func(args *[]any) string { + if filter.OffsetDate > 0 { + *args = append(*args, filter.OffsetDate) + return fmt.Sprintf("message_date < $%d", len(*args)) + } + if filter.OffsetID > 0 { + *args = append(*args, filter.OffsetID) + return fmt.Sprintf("id <= $%d", len(*args)) + } + return "true" + } + switch { + case addOffset < 0 && addOffset+limit > 0: + // around:以锚点为中心,向更新取 -add_offset 条 + 向更旧(含锚点)取 limit+add_offset 条 + fwdLimit := minInt(-addOffset, limit) + bwdLimit := maxInt(limit+addOffset, 0) + fwdArgs := append([]any{}, baseArgs...) + fwdWhere := forwardCond(&fwdArgs) + fwdArgs = append(fwdArgs, fwdLimit) + newer, err := scanList(fmt.Sprintf("SELECT "+channelMessageColumns+" FROM channel_messages WHERE %s AND %s ORDER BY id ASC LIMIT $%d", base, fwdWhere, len(fwdArgs)), fwdArgs) if err != nil { return domain.ChannelHistory{}, err } - out.Messages = append(out.Messages, msg) + bwdArgs := append([]any{}, baseArgs...) + bwdWhere := aroundOlderCond(&bwdArgs) + bwdArgs = append(bwdArgs, bwdLimit+1) + older, err := scanList(fmt.Sprintf("SELECT "+channelMessageColumns+" FROM channel_messages WHERE %s AND %s ORDER BY id DESC LIMIT $%d", base, bwdWhere, len(bwdArgs)), bwdArgs) + if err != nil { + return domain.ChannelHistory{}, err + } + if len(older) > bwdLimit { + older = older[:bwdLimit] + hasMoreOlder = true + } + for i := len(newer) - 1; i >= 0; i-- { + out.Messages = append(out.Messages, newer[i]) + } + out.Messages = append(out.Messages, older...) + case addOffset < 0: + // forward:仅锚点更新方向(拉未读/更新消息) + fwdArgs := append([]any{}, baseArgs...) + fwdWhere := forwardCond(&fwdArgs) + fwdArgs = append(fwdArgs, limit+1) + newer, err := scanList(fmt.Sprintf("SELECT "+channelMessageColumns+" FROM channel_messages WHERE %s AND %s ORDER BY id ASC LIMIT $%d", base, fwdWhere, len(fwdArgs)), fwdArgs) + if err != nil { + return domain.ChannelHistory{}, err + } + if len(newer) > limit { + newer = newer[:limit] + } + for i := len(newer) - 1; i >= 0; i-- { + out.Messages = append(out.Messages, newer[i]) + } + default: + // backward:锚点更旧方向(不含锚点),先跳过 add_offset 条 + where := base + args := append([]any{}, baseArgs...) + if filter.OffsetDate > 0 { + args = append(args, filter.OffsetDate) + where += fmt.Sprintf(" AND message_date < $%d", len(args)) + } else if filter.OffsetID > 0 { + args = append(args, filter.OffsetID) + where += fmt.Sprintf(" AND id < $%d", len(args)) + } + args = append(args, limit+1) + limIdx := len(args) + sql := "SELECT " + channelMessageColumns + " FROM channel_messages WHERE " + where + " ORDER BY id DESC" + if addOffset > 0 { + args = append(args, addOffset) + sql += fmt.Sprintf(" OFFSET $%d", len(args)) + } + sql += fmt.Sprintf(" LIMIT $%d", limIdx) + older, err := scanList(sql, args) + if err != nil { + return domain.ChannelHistory{}, err + } + if len(older) > limit { + older = older[:limit] + hasMoreOlder = true + } + out.Messages = older } - if err := rows.Err(); err != nil { - return domain.ChannelHistory{}, err - } - if len(out.Messages) > limit { - out.Messages = out.Messages[:limit] - out.Count = limit + 1 - } else { - out.Count = len(out.Messages) + out.Count = len(out.Messages) + if hasMoreOlder { + out.Count = len(out.Messages) + 1 } if err := s.populateChannelMessageReplies(ctx, s.db, viewerUserID, channel, out.Messages); err != nil { return domain.ChannelHistory{}, err diff --git a/internal/store/postgres/dialog.go b/internal/store/postgres/dialog.go index 8046040e..fc6e0387 100644 --- a/internal/store/postgres/dialog.go +++ b/internal/store/postgres/dialog.go @@ -155,6 +155,10 @@ func (s *DialogStore) ListByUser(ctx context.Context, userID int64, filter domai if err != nil { return domain.DialogList{}, fmt.Errorf("decode message entities: %w", err) } + media, err := decodeMessageMedia(row.MessageMediaJson) + if err != nil { + return domain.DialogList{}, fmt.Errorf("decode message media: %w", err) + } out.Messages = append(out.Messages, domain.Message{ ID: int(row.MessageID), OwnerUserID: row.UserID, @@ -164,6 +168,7 @@ func (s *DialogStore) ListByUser(ctx context.Context, userID int64, filter domai Out: row.MessageOutgoing, Body: row.MessageBody, Entities: entities, + Media: media, }) } } @@ -241,6 +246,10 @@ func (s *DialogStore) ListByPeers(ctx context.Context, userID int64, peers []dom if err != nil { return domain.DialogList{}, fmt.Errorf("decode message entities: %w", err) } + media, err := decodeMessageMedia(row.MessageMediaJson) + if err != nil { + return domain.DialogList{}, fmt.Errorf("decode message media: %w", err) + } out.Messages = append(out.Messages, domain.Message{ ID: int(row.MessageID), OwnerUserID: row.UserID, @@ -250,6 +259,7 @@ func (s *DialogStore) ListByPeers(ctx context.Context, userID int64, peers []dom Out: row.MessageOutgoing, Body: row.MessageBody, Entities: entities, + Media: media, }) } } diff --git a/internal/store/postgres/queries/dialog.sql b/internal/store/postgres/queries/dialog.sql index ba7cf24e..1af12f5a 100644 --- a/internal/store/postgres/queries/dialog.sql +++ b/internal/store/postgres/queries/dialog.sql @@ -33,7 +33,8 @@ WITH base AS ( COALESCE(m.message_date, 0)::int AS message_date, COALESCE(m.outgoing, false)::boolean AS message_outgoing, COALESCE(m.body, '')::text AS message_body, - COALESCE(m.entities::text, '[]')::text AS message_entities_json + COALESCE(m.entities::text, '[]')::text AS message_entities_json, + COALESCE(m.media::text, '{}')::text AS message_media_json FROM dialogs d LEFT JOIN users u ON d.peer_type = 'user' AND u.id = d.peer_id LEFT JOIN contacts c ON d.peer_type = 'user' AND c.user_id = d.user_id AND c.contact_user_id = d.peer_id @@ -155,7 +156,8 @@ SELECT message_date, message_outgoing, message_body, - message_entities_json + message_entities_json, + message_media_json FROM paged ORDER BY pinned DESC, @@ -288,6 +290,7 @@ base AS ( COALESCE(m.outgoing, false)::boolean AS message_outgoing, COALESCE(m.body, '')::text AS message_body, COALESCE(m.entities::text, '[]')::text AS message_entities_json, + COALESCE(m.media::text, '{}')::text AS message_media_json, r.ord FROM deduped r LEFT JOIN dialogs d @@ -331,7 +334,8 @@ SELECT message_date, message_outgoing, message_body, - message_entities_json + message_entities_json, + message_media_json FROM base ORDER BY ord; diff --git a/internal/store/postgres/sqlcgen/dialog.sql.go b/internal/store/postgres/sqlcgen/dialog.sql.go index a557f0c2..bdec78a9 100644 --- a/internal/store/postgres/sqlcgen/dialog.sql.go +++ b/internal/store/postgres/sqlcgen/dialog.sql.go @@ -602,6 +602,7 @@ base AS ( COALESCE(m.outgoing, false)::boolean AS message_outgoing, COALESCE(m.body, '')::text AS message_body, COALESCE(m.entities::text, '[]')::text AS message_entities_json, + COALESCE(m.media::text, '{}')::text AS message_media_json, r.ord FROM deduped r LEFT JOIN dialogs d @@ -645,7 +646,8 @@ SELECT message_date, message_outgoing, message_body, - message_entities_json + message_entities_json, + message_media_json FROM base ORDER BY ord ` @@ -690,6 +692,7 @@ type ListDialogsByPeersRow struct { MessageOutgoing bool MessageBody string MessageEntitiesJson string + MessageMediaJson string } func (q *Queries) ListDialogsByPeers(ctx context.Context, arg ListDialogsByPeersParams) ([]ListDialogsByPeersRow, error) { @@ -735,6 +738,7 @@ func (q *Queries) ListDialogsByPeers(ctx context.Context, arg ListDialogsByPeers &i.MessageOutgoing, &i.MessageBody, &i.MessageEntitiesJson, + &i.MessageMediaJson, ); err != nil { return nil, err } @@ -781,7 +785,8 @@ WITH base AS ( COALESCE(m.message_date, 0)::int AS message_date, COALESCE(m.outgoing, false)::boolean AS message_outgoing, COALESCE(m.body, '')::text AS message_body, - COALESCE(m.entities::text, '[]')::text AS message_entities_json + COALESCE(m.entities::text, '[]')::text AS message_entities_json, + COALESCE(m.media::text, '{}')::text AS message_media_json FROM dialogs d LEFT JOIN users u ON d.peer_type = 'user' AND u.id = d.peer_id LEFT JOIN contacts c ON d.peer_type = 'user' AND c.user_id = d.user_id AND c.contact_user_id = d.peer_id @@ -834,7 +839,7 @@ WITH base AS ( AND (NOT $16::boolean OR NOT d.pinned) ), paged AS ( - SELECT user_id, peer_type, peer_id, folder_id, top_message_id, top_message_date, read_inbox_max_id, read_outbox_max_id, unread_count, unread_mentions_count, unread_reactions_count, pinned, pinned_order, unread_mark, hidden_peer_settings_bar, peer_user_id, peer_access_hash, peer_phone, peer_first_name, peer_last_name, peer_username, peer_country_code, peer_verified, peer_support, peer_last_seen_at, peer_contact, peer_mutual, message_id, message_from_user_id, message_date, message_outgoing, message_body, message_entities_json + SELECT user_id, peer_type, peer_id, folder_id, top_message_id, top_message_date, read_inbox_max_id, read_outbox_max_id, unread_count, unread_mentions_count, unread_reactions_count, pinned, pinned_order, unread_mark, hidden_peer_settings_bar, peer_user_id, peer_access_hash, peer_phone, peer_first_name, peer_last_name, peer_username, peer_country_code, peer_verified, peer_support, peer_last_seen_at, peer_contact, peer_mutual, message_id, message_from_user_id, message_date, message_outgoing, message_body, message_entities_json, message_media_json FROM base WHERE ( ($17::int <= 0 AND $18::int <= 0) @@ -903,7 +908,8 @@ SELECT message_date, message_outgoing, message_body, - message_entities_json + message_entities_json, + message_media_json FROM paged ORDER BY pinned DESC, @@ -971,6 +977,7 @@ type ListDialogsByUserRow struct { MessageOutgoing bool MessageBody string MessageEntitiesJson string + MessageMediaJson string } func (q *Queries) ListDialogsByUser(ctx context.Context, arg ListDialogsByUserParams) ([]ListDialogsByUserRow, error) { @@ -1037,6 +1044,7 @@ func (q *Queries) ListDialogsByUser(ctx context.Context, arg ListDialogsByUserPa &i.MessageOutgoing, &i.MessageBody, &i.MessageEntitiesJson, + &i.MessageMediaJson, ); err != nil { return nil, err } diff --git a/internal/store/postgres/sqlcgen/models.go b/internal/store/postgres/sqlcgen/models.go index 36655ae2..c261b2e2 100644 --- a/internal/store/postgres/sqlcgen/models.go +++ b/internal/store/postgres/sqlcgen/models.go @@ -20,6 +20,20 @@ type AccountPassword struct { UpdatedAt pgtype.Timestamptz } +type AccountReactionSetting struct { + UserID int64 + MessagesNotifyFrom string + StoriesNotifyFrom string + PollVotesNotifyFrom string + ShowPreviews bool + DefaultReactionType string + DefaultReactionValue string + PaidPrivacyKind string + PaidPrivacyPeerType *string + PaidPrivacyPeerID *int64 + UpdatedAt pgtype.Timestamptz +} + type AppConfig struct { Client string Hash int32 @@ -539,6 +553,19 @@ type PrivateMessage struct { Media []byte } +type PrivateMessageReaction struct { + MessageSenderID int64 + PrivateMessageID int64 + UserID int64 + ReactionType string + ReactionValue string + Big bool + ReactionDate int32 + ChosenOrder int32 + CreatedAt pgtype.Timestamptz + UpdatedAt pgtype.Timestamptz +} + type ProfilePhoto struct { OwnerPeerType string OwnerPeerID int64 @@ -622,6 +649,16 @@ type User struct { LastSeenAt int64 } +type UserChannelMemberIndex struct { + UserID int64 + ChannelID int64 + Status string + Megagroup bool + Broadcast bool + Deleted bool + UpdatedAt pgtype.Timestamptz +} + type UserRecentReaction struct { UserID int64 ReactionType string