This commit is contained in:
Bohdan Horbeshko 2026-08-26 20:52:58 -04:00
parent c9700e931e
commit 58ba0198a3
5 changed files with 518 additions and 69 deletions

View file

@ -260,6 +260,11 @@ type clientLocks struct {
deliveredMessageIds map[int64]int64
deliveredMessageIdsLock sync.Mutex
// albumBuffers collects the messages of an in-progress Telegram
// album (media group) until its debounce timer fires
albumBuffersLock sync.Mutex
albumBuffers map[albumKey]*albumBuffer
authorizerReadLock sync.Mutex
authorizerWriteLock sync.Mutex
@ -342,6 +347,7 @@ func NewClient(conf config.TelegramConfig, jid string, component *xmpp.Component
locks: clientLocks{
chatMessageLocks: make(map[int64]*sync.Mutex),
deliveredMessageIds: make(map[int64]int64),
albumBuffers: make(map[albumKey]*albumBuffer),
},
}, nil
}

View file

@ -2072,6 +2072,15 @@ func (c *Client) ProcessIncomingMessage(chatId int64, message *client.Message) {
return
}
// Albums (Telegram media groups) arrive as several independent
// messages sharing MediaAlbumId; only audios, documents, photos and
// videos can be grouped this way, so none of the group-management
// content types handled below ever apply to them.
if message.MediaAlbumId != 0 {
c.bufferAlbumMessage(chatId, message)
return
}
chat, _, _ := c.GetContactByID(chatId, nil, true)
safeToSend := true
groupChatFrom := ""
@ -2130,6 +2139,98 @@ func (c *Client) ProcessIncomingMessage(chatId int64, message *client.Message) {
}
}
type albumKey struct {
chatId int64
albumId client.JsonInt64
}
type albumBuffer struct {
messages []*client.Message
timer *time.Timer
}
// albumDebounce is how long to wait after an album's last item arrives
// before flushing the group - Telegram clients send every item of an
// album back-to-back in a tight burst, well under a second apart.
const albumDebounce = 1500 * time.Millisecond
// bufferAlbumMessage accumulates one message of a Telegram album (media
// group) and (re)arms its flush timer, so that all of the album's items
// - which TDlib always delivers as independent messages - are sent to
// XMPP together as a single XEP-0447/XEP-0367 group instead of N
// unrelated messages.
func (c *Client) bufferAlbumMessage(chatId int64, message *client.Message) {
key := albumKey{chatId: chatId, albumId: message.MediaAlbumId}
c.locks.albumBuffersLock.Lock()
defer c.locks.albumBuffersLock.Unlock()
buf, ok := c.locks.albumBuffers[key]
if !ok {
buf = &albumBuffer{}
c.locks.albumBuffers[key] = buf
}
buf.messages = append(buf.messages, message)
if buf.timer != nil {
buf.timer.Stop()
}
buf.timer = time.AfterFunc(albumDebounce, func() {
c.flushAlbum(key)
})
}
// flushAlbum sends a complete, debounced album to XMPP. It re-derives
// MUC routing itself, rather than reusing state from any single
// buffered message, since the group as a whole is only now being
// dispatched; the group-management content-type switch a single
// message would otherwise go through never applies to album items, so
// it's skipped here.
func (c *Client) flushAlbum(key albumKey) {
c.locks.albumBuffersLock.Lock()
buf, ok := c.locks.albumBuffers[key]
if ok {
delete(c.locks.albumBuffers, key)
}
c.locks.albumBuffersLock.Unlock()
if !ok || len(buf.messages) == 0 {
return
}
messages := buf.messages
sort.Slice(messages, func(i, j int) bool { return messages[i].Id < messages[j].Id })
lock := c.getChatMessageLock(key.chatId)
lock.Lock()
defer lock.Unlock()
representative := messages[0]
chat, _, _ := c.GetContactByID(key.chatId, nil, true)
groupChatFrom := ""
groupChatTos := []string{}
safeToSend := true
if c.Session.MUC && c.IsGroup(chat) {
senderId := c.getMessageSenderId(representative)
if senderId == 0 {
log.Errorf("Invalid sender id for album message %#v", representative)
return
}
safeToSend = c.assureMUCOccupant(key.chatId, senderId, representative.SenderId, chat)
groupChatFrom = gateway.MUCJID(key.chatId) + "/" + c.GetMUCNickname(key.chatId, senderId)
var ok bool
ok, groupChatTos = c.getMUCJoinedJIDs(key.chatId, nil, true)
if !ok {
safeToSend = false
}
}
if !safeToSend {
mucJID := gateway.MUCJID(key.chatId)
gateway.SendErrorMessage(c.jid, mucJID, "Cannot show a message", 500, true, c.xmpp)
return
}
c.SendAlbumToGateway(key.chatId, messages, false, groupChatFrom, groupChatTos, "")
}
// SendMessageToGateway transfers a message to XMPP side and marks it as read on Telegram side
func (c *Client) SendMessageToGateway(chatId int64, message *client.Message, id string, delay bool, groupChatFrom string, groupChatTos []string, mamQueryId string) {
var isCarbon bool
@ -2153,6 +2254,7 @@ func (c *Client) SendMessageToGateway(chatId int64, message *client.Message, id
var rawCaption string
var fileMetaName, fileMetaType, fileMetaHash string
var fileMetaSize int64
var fileMetaDate string
var reply *gateway.Reply
var replyObtained bool
@ -2263,6 +2365,7 @@ func (c *Client) SendMessageToGateway(chatId int64, message *client.Message, id
// to attach file-sharing metadata to in the first
// place.
fileMetaName, fileMetaType, fileMetaSize, fileMetaHash = c.fileMetadata(content, file, link)
fileMetaDate = time.Unix(int64(message.Date), 0).UTC().Format(time.RFC3339)
}
}
}
@ -2423,6 +2526,7 @@ func (c *Client) SendMessageToGateway(chatId int64, message *client.Message, id
Name: fileMetaName,
MediaType: fileMetaType,
Size: fileMetaSize,
Date: fileMetaDate,
Desc: sfsDesc,
HashAlgo: "sha-256",
HashValue: fileMetaHash,
@ -2442,6 +2546,229 @@ func (c *Client) SendMessageToGateway(chatId int64, message *client.Message, id
c.UpdateLastChatMessageId(chatId, sId)
}
type albumSendItem struct {
id string
oob string
name string
mediaType string
size int64
date string
hash string
}
// SendAlbumToGateway sends a Telegram album (media group) to XMPP as one
// XEP-0447 metadata-only anchor message (one <file-sharing> per item),
// one OOB/source message per item attaching to that anchor, and (if any
// item carries a caption) a single shared caption message also
// attaching to the anchor - all correlated via XEP-0367, so a compliant
// client renders the whole group as one message. delay/mamQueryId carry
// the same meaning as SendMessageToGateway's own. There's no id-reuse
// parameter, unlike SendMessageToGateway: album items can never be
// XMPP-originated (the bridge never sends albums outbound to Telegram),
// so there's never a pre-existing id to reuse or stay stable with - the
// deterministic ids derived here are already reproduced identically
// whether a message is delivered live or replayed later.
func (c *Client) SendAlbumToGateway(chatId int64, messages []*client.Message, delay bool, groupChatFrom string, groupChatTos []string, mamQueryId string) {
if len(messages) == 0 {
return
}
representative := messages[0]
var isCarbon bool
var jids []string
var isGroupchat bool
var originalFrom string
senderId := c.getMessageSenderId(representative)
if len(groupChatTos) == 0 {
isCarbon = c.isCarbonsEnabled() && representative.IsOutgoing
jids = c.GetCarbonFullJids(isCarbon, "", true)
} else {
isGroupchat = true
jids = groupChatTos
if senderId != 0 {
originalFrom = gateway.CHATJID(senderId, true)
}
}
var items []albumSendItem
var rawCaption string
for _, message := range messages {
content := message.Content
if content == nil {
continue
}
if rawCaption == "" {
if capText := c.messageToText(message, false); capText != "" {
rawCaption = capText
}
}
file, _ := c.contentToFile(content)
file = c.ensureDownloadFile(file)
link, _ := c.formatFile(file, true)
if link == "" {
log.Errorf("Album item %v in chat %v produced no OOB link, skipping", message.Id, chatId)
continue
}
name, mediaType, size, hash := c.fileMetadata(content, file, link)
items = append(items, albumSendItem{
id: strconv.FormatInt(message.Id, 10),
oob: link,
name: name,
mediaType: mediaType,
size: size,
date: time.Unix(int64(message.Date), 0).UTC().Format(time.RFC3339),
hash: hash,
})
}
if len(items) == 0 {
log.Errorf("Album %v in chat %v produced no usable items", representative.MediaAlbumId, chatId)
return
}
reply, _ := c.getMessageReply(representative, false, true)
if !c.Session.Receipts {
for _, message := range messages {
c.MarkAsRead(chatId, message.Id)
}
}
// Derived from the first item's own message id, not MediaAlbumId: a
// reply to the anchor or caption message needs to resolve back to a
// real Telegram message id (parseMessageId strips the "a"/"c" prefix
// and parses the rest as one) - MediaAlbumId is a separate id space
// that doesn't identify any message at all.
ancId := "a" + items[0].id
capId := "c" + items[0].id
var from string
if groupChatFrom == "" {
from = gateway.CHATNODE(chatId)
} else {
from = groupChatFrom
}
var timestamp int64
if delay {
timestamp = int64(representative.Date)
}
var mucUserItem *gateway.MUCUserItem
var mucJID string
if mamQueryId != "" {
chatMember, err := c.client.GetChatMember(&client.GetChatMemberRequest{
ChatId: chatId,
MemberId: representative.SenderId,
})
var status client.ChatMemberStatus
if err == nil {
status = chatMember.Status
}
chat, err := c.GetChatByID(chatId, nil, true)
if err == nil {
affiliation, role := c.memberStatusToAffiliationAndRole(status, chat)
mucUserItem = &gateway.MUCUserItem{
Affiliation: affiliation,
Jid: gateway.CHATJID(chatId, false),
Role: role,
}
}
mucJID = gateway.MUCJID(chatId)
}
// Same OMEMO policy as SendMessageToGateway (see its comment),
// restricted to the shared caption - OOB file links are always left
// in the clear.
var captionEnvelope *e2ee.Envelope
omemoFailed := false
if !isGroupchat && rawCaption != "" {
if backend, ok := gateway.E2EE.Backend(); ok {
owner := e2ee.OwnedPeer(c.Session.Login, gateway.CHATJID(chatId, false))
mode, modeErr := e2ee.ParseMode(c.Session.OMEMO)
if modeErr != nil {
log.Error(errors.Wrap(modeErr, "Invalid omemo mode"))
}
if active, err := e2ee.ShouldEncrypt(backend, owner, mode); err != nil {
log.Error(errors.Wrap(err, "Failed to evaluate OMEMO policy"))
} else if active {
peer := e2ee.PeerID(c.jid)
if err := c.ensureOMEMOSession(backend, owner, peer); err != nil {
log.Error(errors.Wrap(err, "Failed to establish OMEMO session"))
omemoFailed = true
} else if env, err := backend.Encrypt(owner, []e2ee.PeerID{peer}, []byte(rawCaption)); err != nil {
log.Error(errors.Wrap(err, "Failed to encrypt OMEMO album caption"))
omemoFailed = true
} else {
captionEnvelope = &env
}
}
}
}
var occupantId string
if isGroupchat && senderId != 0 {
occupantId = gateway.CHATNODE(senderId)
}
sfsDesc := ""
if rawCaption != "" && captionEnvelope == nil {
sfsDesc = rawCaption
}
fileSharings := make([]*gateway.FileSharing, len(items))
for i, it := range items {
fileSharings[i] = &gateway.FileSharing{
Id: it.id,
Name: it.name,
MediaType: it.mediaType,
Size: it.size,
Date: it.date,
Desc: sfsDesc,
HashAlgo: "sha-256",
HashValue: it.hash,
}
}
for _, jid := range jids {
commonArgs := []args.V{
gateway.SMReply(reply), gateway.SMTimestamp(timestamp), gateway.SMIsCarbon(isCarbon),
gateway.SMIsGroupchat(isGroupchat), gateway.SMRequestReceipt(c.Session.Receipts),
gateway.SMOriginalFrom(originalFrom), gateway.SMMamQueryId(mamQueryId),
gateway.SMMucJID(mucJID), gateway.SMMucUserItem(mucUserItem),
gateway.SMOccupantId(occupantId),
}
if omemoFailed {
gateway.SendMessage(jid, from, c.xmpp, append(commonArgs,
gateway.SMBody(gateway.OMEMOSendFailedBody), gateway.SMId(ancId), gateway.SMStanzaId(ancId))...)
continue
}
gateway.SendMessage(jid, from, c.xmpp, append(commonArgs,
gateway.SMId(ancId), gateway.SMStanzaId(ancId), gateway.SMFileSharings(fileSharings))...)
for _, it := range items {
gateway.SendMessage(jid, from, c.xmpp, append(commonArgs,
gateway.SMBody(it.oob), gateway.SMId(it.id), gateway.SMStanzaId(it.id), gateway.SMOOB(it.oob),
gateway.SMAttachToId(ancId), gateway.SMSourcesId(it.id))...)
}
if rawCaption != "" {
capArgs := append(commonArgs,
gateway.SMBody(rawCaption), gateway.SMId(capId), gateway.SMStanzaId(capId),
gateway.SMOMEMOEnvelope(captionEnvelope), gateway.SMAttachToId(ancId))
if sfsDesc != "" {
capArgs = append(capArgs, gateway.SMSFSFallback(true))
}
gateway.SendMessage(jid, from, c.xmpp, capArgs...)
}
}
c.UpdateLastChatMessageId(chatId, items[len(items)-1].id)
}
// ensureOMEMOSession makes sure backend has at least one established
// session with peer under owner's identity, lazily fetching and ingesting
// peer's device list/bundles over PEP if not (the first message to a peer
@ -2472,6 +2799,45 @@ func (c *Client) SendDelayedMUCMessage(chatId int64, message *client.Message, to
)
}
// SendDelayedAlbumToGateway is SendDelayedMUCMessage's album-aware
// counterpart, used for MUC history/MAM replay of a Telegram album.
func (c *Client) SendDelayedAlbumToGateway(chatId int64, messages []*client.Message, toJid string, mamQueryId string) {
if len(messages) == 0 {
return
}
senderId := c.getMessageSenderId(messages[0])
c.SendAlbumToGateway(
chatId,
messages,
true,
gateway.MUCJID(chatId)+"/"+c.GetMUCNickname(chatId, senderId),
[]string{toJid},
mamQueryId,
)
}
// GroupMessagesByAlbum partitions messages into runs sharing the same
// non-zero MediaAlbumId (order-preserving by first occurrence; a
// message with MediaAlbumId == 0 always starts its own singleton
// group).
func GroupMessagesByAlbum(messages []*client.Message) [][]*client.Message {
var groups [][]*client.Message
albumIndex := make(map[client.JsonInt64]int)
for _, message := range messages {
if message.MediaAlbumId == 0 {
groups = append(groups, []*client.Message{message})
continue
}
if idx, ok := albumIndex[message.MediaAlbumId]; ok {
groups[idx] = append(groups[idx], message)
continue
}
albumIndex[message.MediaAlbumId] = len(groups)
groups = append(groups, []*client.Message{message})
}
return groups
}
// MarkAsRead marks a message as read
func (c *Client) MarkAsRead(chatId, messageId int64) {
c.client.ViewMessages(&client.ViewMessagesRequest{
@ -3532,10 +3898,24 @@ func (c *Client) sendMessagesReverse(chatID int64, messages []*client.Message, p
plainTos = []string{c.jid}
}
if !plain {
reversed := make([]*client.Message, len(messages))
for i, message := range messages {
reversed[len(messages)-1-i] = message
}
for _, group := range GroupMessagesByAlbum(reversed) {
if len(group) > 1 {
c.SendDelayedAlbumToGateway(chatID, group, toJid, "")
} else {
c.SendDelayedMUCMessage(chatID, group[0], toJid, "")
}
}
return
}
for i := len(messages) - 1; i >= 0; i-- {
message := messages[i]
if plain {
reply, _ := c.getMessageReply(message, false, true)
sId := strconv.FormatInt(message.Id, 10)
@ -3596,9 +3976,6 @@ func (c *Client) sendMessagesReverse(chatID int64, messages []*client.Message, p
gateway.SMOMEMOEnvelope(envelope), gateway.SMOccupantId(occupantId),
)
}
} else {
c.SendDelayedMUCMessage(chatID, message, toJid, "")
}
}
}

View file

@ -521,10 +521,14 @@ type AttachTo struct {
Id string `xml:"id,attr"`
}
// FileSharing is a XEP-0447 file-sharing element, with its source (a
// XEP-0066 OOB URL, typically) nested inline
// FileSharing is a XEP-0447 file-sharing element. Id disambiguates
// multiple FileSharing entries in the same message from each other,
// for a standalone Sources element to name which one it supplies a
// source for; left empty when Sources is nested inline instead, since
// then there's no ambiguity to resolve.
type FileSharing struct {
XMLName xml.Name `xml:"urn:xmpp:sfs:0 file-sharing"`
Id string `xml:"id,attr,omitempty"`
File FileMetadata `xml:"urn:xmpp:file:metadata:0 file"`
Sources *Sources `xml:"urn:xmpp:sfs:0 sources,omitempty"`
}
@ -535,8 +539,9 @@ type FileMetadata struct {
Name string `xml:"name,omitempty"`
MediaType string `xml:"media-type,omitempty"`
Size int64 `xml:"size,omitempty"`
Hashes []Hash `xml:"urn:xmpp:hashes:2 hash,omitempty"`
Date string `xml:"date,omitempty"`
Desc string `xml:"desc,omitempty"`
Hashes []Hash `xml:"urn:xmpp:hashes:2 hash,omitempty"`
}
// Hash is a XEP-0300 cryptographic hash, used inline within XEP-0446
@ -547,9 +552,12 @@ type Hash struct {
Value string `xml:",chardata"`
}
// Sources is a XEP-0447 sources element
// Sources is a XEP-0447 sources element. Id names which FileSharing
// entry (by its own Id) this source is for; only meaningful standalone,
// in a separate message from the FileSharing it supplies.
type Sources struct {
XMLName xml.Name `xml:"urn:xmpp:sfs:0 sources"`
Id string `xml:"id,attr,omitempty"`
UrlData UrlData `xml:"http://jabber.org/protocol/url-data url-data"`
}

View file

@ -53,11 +53,15 @@ type marker struct {
Id string
}
// FileSharing is a XEP-0447/XEP-0446 minimal file-sharing metadata set
// FileSharing is a XEP-0447/XEP-0446 minimal file-sharing metadata set.
// Id is only needed when several entries ride in the same message and a
// source sent standalone needs to name which one it's for.
type FileSharing struct {
Id string
Name string
MediaType string
Size int64
Date string
Desc string
HashAlgo string
HashValue string
@ -345,6 +349,20 @@ var SMAttachToId = args.NewString()
// nesting the message's OOB URL as its source
var SMFileSharing = args.New()
// SMFileSharings is a set of XEP-0447/XEP-0446 file-sharing entries
// ([]*FileSharing) sent without an inline source - each entry's Id is
// meant to be referenced by a source sent standalone via SMSourcesId
var SMFileSharings = args.New()
// SMSourcesId marks this message as standing in for a source (the OOB
// URL, via SMOOB) of the file-sharing entry with this id, sent
// standalone rather than nested inline
var SMSourcesId = args.NewString()
// SMSFSFallback marks the body as XEP-0428 fallback content for
// urn:xmpp:sfs:0, standing in for a file-sharing <desc> sent elsewhere
var SMSFSFallback = args.NewBool()
// SMIsCarbon marks the message as a XEP-0280 carbon copy
var SMIsCarbon = args.NewBool()
@ -404,6 +422,9 @@ func sendMessageWrapper(to, from string, component *xmpp.Component, args ...args
replaceId := SMReplaceId.Get(args)
attachToId := SMAttachToId.Get(args)
fileSharing, _ := SMFileSharing.Get(args).(*FileSharing)
fileSharings, _ := SMFileSharings.Get(args).([]*FileSharing)
sourcesId := SMSourcesId.Get(args)
sfsFallback := SMSFSFallback.Get(args)
isCarbon := SMIsCarbon.Get(args)
isGroupchat := SMIsGroupchat.Get(args)
forceSubject := SMForceSubject.Get(args)
@ -589,6 +610,7 @@ func sendMessageWrapper(to, from string, component *xmpp.Component, args ...args
Name: fileSharing.Name,
MediaType: fileSharing.MediaType,
Size: fileSharing.Size,
Date: fileSharing.Date,
Desc: fileSharing.Desc,
}
if fileSharing.HashValue != "" {
@ -608,6 +630,36 @@ func sendMessageWrapper(to, from string, component *xmpp.Component, args ...args
})
}
}
for _, entry := range fileSharings {
file := extensions.FileMetadata{
Name: entry.Name,
MediaType: entry.MediaType,
Size: entry.Size,
Date: entry.Date,
Desc: entry.Desc,
}
if entry.HashValue != "" {
file.Hashes = []extensions.Hash{{Algo: entry.HashAlgo, Value: entry.HashValue}}
}
message.Extensions = append(message.Extensions, extensions.FileSharing{Id: entry.Id, File: file})
}
if len(fileSharings) > 0 && body == "" {
// XEP-0334: body-less messages are commonly excluded from MAM by
// default archiving policies, so hint explicitly.
message.Extensions = append(message.Extensions, stanza.HintStore{})
}
if sourcesId != "" && oob != "" {
message.Extensions = append(message.Extensions, extensions.Sources{
Id: sourcesId,
UrlData: extensions.UrlData{Target: oob},
})
}
if sfsFallback {
message.Extensions = append(message.Extensions, extensions.Fallback{
For: "urn:xmpp:sfs:0",
Body: []extensions.FallbackBody{{}},
})
}
if reactions != nil {
reactionExts := make([]extensions.Reaction, len(reactions.Reactions))
for i, r := range reactions.Reactions {

View file

@ -2879,10 +2879,16 @@ func handleSetQueryMAM(s xmpp.Sender, iq *stanza.IQ, query extensions.MAMQuery)
log.Debugf("obtained %v messages", len(messages))
// ty zhe lopnesh, detochka
queryId := ns + " " + query.GetQueryId()
for _, message := range messages {
session.SendDelayedMUCMessage(toID, message, iq.From, queryId)
for _, group := range telegram.GroupMessagesByAlbum(messages) {
if len(group) > 1 {
session.SendDelayedAlbumToGateway(toID, group, iq.From, queryId)
} else {
session.SendDelayedMUCMessage(toID, group[0], iq.From, queryId)
}
for _, message := range group {
session.SendDelayedMUCReactions(toID, message, iq.From, queryId)
}
}
rs := stanza.ResultSet{}
switch ns {
@ -3211,7 +3217,7 @@ func parseMessageId(sId string) (int64, bool) {
if len(idParts) >= 1 {
sId = idParts[0]
}
} else if sId[0] == 'c' {
} else if sId[0] == 'c' || sId[0] == 'a' {
sId = sId[1:]
}
id, err := strconv.ParseInt(sId, 10, 64)