package gateway import ( "context" "crypto/sha1" "encoding/base64" "encoding/xml" "fmt" "github.com/pkg/errors" "io" "sort" "strconv" "strings" "sync" "dev.narayana.im/narayana/telegabber/badger" "dev.narayana.im/narayana/telegabber/e2ee" "dev.narayana.im/narayana/telegabber/e2ee/omemo" "dev.narayana.im/narayana/telegabber/xmpp/extensions" "github.com/google/uuid" log "github.com/sirupsen/logrus" "github.com/soheilhy/args" "github.com/xdg-go/stringprep" "gosrc.io/xmpp" "gosrc.io/xmpp/stanza" ) type Reply struct { Author string Id string Start uint64 End uint64 } // Reactions is a XEP-0444 reaction set: the sender's complete current // list of emoji on the message identified by Id (always a full // replacement, never a diff). type Reactions struct { Id string Reactions []string } type MarkerType byte const ( MarkerTypeReceived MarkerType = iota MarkerTypeDisplayed ) type marker struct { Type MarkerType Id string } type MUCUserItem struct { Affiliation string Jid string Role string } const NSNick string = "http://jabber.org/protocol/nick" const NodeVCard4 string = "urn:xmpp:vcard4" const NodeAvatarMetadata string = "urn:xmpp:avatar:metadata" const NodeAvatarMetadataNotify string = NodeAvatarMetadata + "+notify" const NodeAvatarData string = "urn:xmpp:avatar:data" const NSCommand string = "http://jabber.org/protocol/commands" // NSCaps is the XEP-0115 entity capabilities namespace const NSCaps string = "http://jabber.org/protocol/caps" // NSReactions is the XEP-0444 message reactions namespace const NSReactions string = "urn:xmpp:reactions:0" // CapsNode is this software's advertised XEP-0115 node URI const CapsNode string = "https://dev.narayana.im/narayana/telegabber/" const NS_MAM2 = "urn:xmpp:mam:2" const NS_MAM1 = "urn:xmpp:mam:1" const NS_MAM0 = "urn:xmpp:mam:0" // Queue stores presences to send later var Queue = make(map[string]*stanza.Presence) var QueueLock = sync.Mutex{} // Jid stores the component's JID object var Jid *stanza.Jid // Version stores this software's version var Version string // IdsDB provides a disk-backed bidirectional dictionary of Telegram and XMPP ids var IdsDB badger.IdsDB // E2EE holds the process-wide selected end-to-end encryption backend (set // up by xmpp.NewComponent from config), reachable from both the xmpp and // telegram packages without either importing the other's e2ee-specific // glue - mirrors IdsDB's role as a shared, package-level handle. Its // Backend() reports ok=false whenever the feature is disabled or not yet // set up; every call site must treat that as "behave exactly as if this // feature didn't exist" rather than nil-checking E2EE itself. var E2EE *e2ee.Manager // ShutdownCtx is cancelled once (see xmpp.Close) when the component begins // tearing down. Long-running background operations that would otherwise // wait out their own timeout regardless of process lifecycle - e.g. e2ee's // PEP bundle fetches - should derive their per-call timeout from this // (context.WithTimeout(gateway.ShutdownCtx, ...)) rather than // context.Background(), so teardown cancels them immediately instead of // leaving them to run for up to their full timeout past Close(). var ShutdownCtx, CancelShutdown = context.WithCancel(context.Background()) // DirtySessions denotes that some Telegram session configurations // were changed and need to be re-flushed to the YamlDB var DirtySessions = false // MessageOutgoingPermissionVersion contains a XEP-0356 version to fake outgoing messages by foreign JIDs var MessageOutgoingPermissionVersion = 0 // MAMThreshold specifies a day limit behind which history should not be requested to avoid abuse detection and storage overload var MAMThreshold uint32 // HistoryBackfillLimit caps how many messages a single chat's // downtime-backfill sweep will pull, regardless of gap size. 0 (unset) // means the default of 1000. var HistoryBackfillLimit uint32 = 1000 // HistoryBackfillWorkers bounds how many chats backfill concurrently. // Chats are independent, so this is a concurrency cap, not a queue depth: // one huge or FLOOD_WAIT-throttled chat only occupies one worker slot, // leaving the rest free to make progress on other chats instead of // waiting behind it. 0 (unset) means the default of 4. var HistoryBackfillWorkers uint32 = 4 // MaxReactionsPerEmoji caps named reactors relayed per (message, emoji) // pair for group/channel reactions - anonymous ones are exempt, see // telegram.updateMessageInteractionInfo. 0 (unset) means the default of // 10; values above 100 are clamped (TDlib's own per-call ceiling). var MaxReactionsPerEmoji uint32 = 10 // DownloadSpaceMargin is extra free-disk-space headroom, beyond a file's // own declared size, required before a TDlib download proceeds. Download // requests are already serialized system-wide via StorageLock, so this // only has to cover concurrent writes by other, unrelated processes on // the same machine - not a race against telegabber's own downloads. 0 // (unset) means the default of 64 MiB. var DownloadSpaceMargin uint64 = 64 * 1024 * 1024 // CHATNODE converts numeric id to node part of 1-1 chat JID func CHATNODE(chatId int64) string { return strconv.FormatInt(chatId, 10) } // CHATJID converts numeric id to 1-1 chat JID func CHATJID(chatId int64, full bool) string { var suffix string if full { suffix = Jid.Full() } else { suffix = Jid.Bare() } return CHATNODE(chatId) + "@" + suffix } // MUCNODE converts numeric id to node part of MUC JID func MUCNODE(chatId int64) string { return "c" + CHATNODE(chatId) } // MUCJID converts numeric id to MUC JID func MUCJID(chatId int64) string { return "c" + CHATJID(chatId, false) } var resourcePrepProfile = stringprep.Profile{ Mappings: []stringprep.Mapping{ stringprep.TableB1, }, Normalize: true, Prohibits: []stringprep.Set{ stringprep.TableC1_2, stringprep.TableC2_1, stringprep.TableC2_2, stringprep.TableC3, stringprep.TableC4, stringprep.TableC5, stringprep.TableC6, stringprep.TableC7, stringprep.TableC8, stringprep.TableC9, }, CheckBiDi: true, } // ResourcePrep normalizes a resource according to RFC 6122 func ResourcePrep(resource string) (string, error) { return resourcePrepProfile.Prepare(resource) } // SendMessage creates and sends a message stanza. See the SM* option // constructors below (SMBody, SMReply, SMOOB, ...) for what can be set. func SendMessage(to, from string, component *xmpp.Component, args ...args.V) { sendMessageWrapper(to, from, component, args...) } // SendServiceMessage creates and sends a simple message stanza from transport func SendServiceMessage(to, body string, component *xmpp.Component) { var id string if uuid, err := uuid.NewRandom(); err == nil { id = uuid.String() } sendMessageWrapper(to, "", component, SMBody(body), SMId(id)) } // SendTextMessage creates and sends a simple message stanza func SendTextMessage(to, from, body string, component *xmpp.Component, isGroupchat bool) { var id string if uuid, err := uuid.NewRandom(); err == nil { id = uuid.String() } sendMessageWrapper(to, from, component, SMBody(body), SMId(id), SMIsGroupchat(isGroupchat)) } // SendMUCAnnouncement creates and sends a message by a temporary occupant func SendMUCAnnouncement(to, from, body, nickname, id string, component *xmpp.Component) { if nickname == "" { nickname = "announcement" } fullFrom := from + "/" + nickname SendPresence( component, to, SPFullFrom(fullFrom), SPMUCAffiliation("admin"), SPMUCRole("moderator"), SPMUCJid(from), ) if id == "" { if uuid, err := uuid.NewRandom(); err == nil { id = uuid.String() } } sendMessageWrapper(to, fullFrom, component, SMBody(body), SMId(id), SMIsGroupchat(true)) SendPresence( component, to, SPType("unavailable"), SPFullFrom(fullFrom), SPMUCAffiliation("none"), SPMUCRole("none"), SPMUCJid(from), ) } // SendErrorMessage creates and sends an error message stanza func SendErrorMessage(to, from, text string, code int, isGroupchat bool, component *xmpp.Component) { sendMessageWrapper(to, from, component, SMErrorText(text), SMIsGroupchat(isGroupchat), SMErrorCode(code)) } // SendErrorMessageWithBody creates and sends an error message stanza with body payload func SendErrorMessageWithBody(to, from, body, errorText, id string, code int, isGroupchat bool, component *xmpp.Component) { sendMessageWrapper(to, from, component, SMBody(body), SMErrorText(errorText), SMId(id), SMIsGroupchat(isGroupchat), SMErrorCode(code)) } // SendSubjectMessage creates and sends a MUC subject func SendSubjectMessage(to, from, subject, id string, component *xmpp.Component, timestamp int64) { sendMessageWrapper(to, from, component, SMSubject(subject), SMId(id), SMTimestamp(timestamp), SMIsGroupchat(true), SMForceSubject(true)) } // SendMessageMarker creates and sends a message stanza with a XEP-0333 marker func SendMessageMarker(to string, from string, component *xmpp.Component, markerType MarkerType, markerId string) { sendMessageWrapper(to, from, component, SMMarker(&marker{ Type: markerType, Id: markerId, })) } // SendReactionMessage creates and sends a message stanza with a XEP-0444 // reaction set (see Reactions); an empty reactionEmojis clears id's reactions. func SendReactionMessage(to string, from string, component *xmpp.Component, id string, reactionEmojis []string, isGroupchat bool) { sendMessageWrapper(to, from, component, SMReactions(&Reactions{ Id: id, Reactions: reactionEmojis, }), SMIsGroupchat(isGroupchat)) } // SendMAMReactionMessage is SendReactionMessage wrapped as a XEP-0313 MAM // result: entryId is this result's own id (distinct from id, the // reacted-to message's own id); timestamp stands in for a historical // reaction timestamp, which doesn't exist (only a reaction's current // state, not when it was added). func SendMAMReactionMessage(to string, from string, component *xmpp.Component, id string, reactionEmojis []string, entryId string, timestamp int64, mucJID string, mamQueryId string) { sendMessageWrapper(to, from, component, SMReactions(&Reactions{Id: id, Reactions: reactionEmojis}), SMId(entryId), SMStanzaId(entryId), SMTimestamp(timestamp), SMMamQueryId(mamQueryId), SMMucJID(mucJID), SMIsGroupchat(true), ) } // SendMUCInvite creates and send a MUC invitation message func SendMUCInvite(to string, from string, component *xmpp.Component, inviteFrom string) { sendMessageWrapper(to, from, component, SMInviteFrom(inviteFrom)) } // SendMUCStatusCode creates a groupchat message with a muc#user status code func SendMUCStatusCode(to string, from string, component *xmpp.Component, statusCode int64) { sendMessageWrapper(to, from, component, SMIsGroupchat(true), SMStatusCode(statusCode)) } // SMBody is the message body var SMBody = args.NewString() // SMSubject is a MUC subject var SMSubject = args.NewString() // SMErrorText is the human-readable error text var SMErrorText = args.NewString() // SMId is the stanza id var SMId = args.NewString() // SMReply is a XEP-0461 reply reference (*Reply) var SMReply = args.New() // SMMarker is a XEP-0333 marker (*marker) var SMMarker = args.New() // SMTimestamp is a XEP-0203 delay timestamp var SMTimestamp = args.NewInt64() // SMOOB is a XEP-0066 out of band URL var SMOOB = args.NewString() // SMReplaceId is a XEP-0308 replaced message id var SMReplaceId = args.NewString() // SMIsCarbon marks the message as a XEP-0280 carbon copy var SMIsCarbon = args.NewBool() // SMIsGroupchat marks the message as a groupchat message var SMIsGroupchat = args.NewBool() // SMForceSubject forces an empty MUC subject element to be included var SMForceSubject = args.NewBool() // SMRequestReceipt requests a XEP-0184 receipt var SMRequestReceipt = args.NewBool() // SMOriginalFrom is a XEP-0033 original sender address var SMOriginalFrom = args.NewString() // SMErrorCode is a legacy numeric error code var SMErrorCode = args.NewInt() // SMInviteFrom is a XEP-0045 MUC direct invitation sender var SMInviteFrom = args.NewString() // SMStanzaId is a XEP-0359 stanza id var SMStanzaId = args.NewString() // SMStatusCode is a XEP-0045 muc#user status code var SMStatusCode = args.NewInt64() // SMMamQueryId is a XEP-0313 MAM query id (namespace-prefixed, space-separated) var SMMamQueryId = args.NewString() // SMMucJID is the real MUC room JID to send to var SMMucJID = args.NewString() // SMMucUserItem is a XEP-0045 muc#user item (*MUCUserItem) var SMMucUserItem = args.New() // SMOMEMOEnvelope is an already-encrypted OMEMO envelope (*e2ee.Envelope, // from Backend.Encrypt) to attach as an extension plus a // XEP-0380 EME hint - see omemo.go's omemoStanzaExtension. var SMOMEMOEnvelope = args.New() // SMReactions is a XEP-0444 reaction set (*Reactions) var SMReactions = args.New() func sendMessageWrapper(to, from string, component *xmpp.Component, args ...args.V) { body := SMBody.Get(args) subject := SMSubject.Get(args) errorText := SMErrorText.Get(args) id := SMId.Get(args) reply, _ := SMReply.Get(args).(*Reply) marker, _ := SMMarker.Get(args).(*marker) timestamp := SMTimestamp.Get(args) oob := SMOOB.Get(args) replaceId := SMReplaceId.Get(args) isCarbon := SMIsCarbon.Get(args) isGroupchat := SMIsGroupchat.Get(args) forceSubject := SMForceSubject.Get(args) requestReceipt := SMRequestReceipt.Get(args) originalFrom := SMOriginalFrom.Get(args) errorCode := SMErrorCode.Get(args) inviteFrom := SMInviteFrom.Get(args) stanzaId := SMStanzaId.Get(args) statusCode := SMStatusCode.Get(args) mamQueryId := SMMamQueryId.Get(args) mucJID := SMMucJID.Get(args) mucUserItem, _ := SMMucUserItem.Get(args).(*MUCUserItem) envelope, _ := SMOMEMOEnvelope.Get(args).(*e2ee.Envelope) reactions, _ := SMReactions.Get(args).(*Reactions) toJid, err := stanza.NewJid(to) if err != nil { log.WithFields(log.Fields{ "to": to, }).Error(errors.Wrap(err, "Invalid to JID!")) return } bareTo := toJid.Bare() componentJid := Jid.Full() var logFrom string var messageFrom string var messageTo string var bareFrom string if isGroupchat { logFrom = from messageFrom = from bareFrom, _, _ = SplitJID(from) } else { if from == "" { logFrom = componentJid messageFrom = componentJid bareFrom = componentJid } else if inviteFrom != "" { logFrom = from messageFrom = from + "@" + Jid.Bare() bareFrom = messageFrom } else { logFrom = from messageFrom = from + "@" + componentJid bareFrom = from + "@" + Jid.Bare() } } if isCarbon { messageTo = messageFrom messageFrom = bareTo + "/" + Jid.Resource } else if mucJID != "" { messageTo = mucJID } else { messageTo = to } log.WithFields(log.Fields{ "from": logFrom, "to": to, }).Warn("Got message") var messageType stanza.StanzaType if errorCode != 0 { messageType = stanza.MessageTypeError } else if isGroupchat { messageType = stanza.MessageTypeGroupchat } else if inviteFrom != "" { messageType = stanza.MessageTypeNormal } else { messageType = stanza.MessageTypeChat } message := stanza.Message{ Attrs: stanza.Attrs{ From: messageFrom, To: messageTo, Type: messageType, Id: id, }, Subject: subject, Body: body, } if errorCode != 0 { message.Error = stanza.Err{ Code: errorCode, Text: errorText, } switch errorCode { case 400: message.Error.Type = stanza.ErrorTypeModify message.Error.Reason = "bad-request" case 403: message.Error.Type = stanza.ErrorTypeAuth message.Error.Reason = "forbidden" case 404: message.Error.Type = stanza.ErrorTypeCancel message.Error.Reason = "item-not-found" case 406: message.Error.Type = stanza.ErrorTypeModify message.Error.Reason = "not-acceptable" case 500: message.Error.Type = stanza.ErrorTypeWait message.Error.Reason = "internal-server-error" default: log.Error("Unknown error code, falling back with empty reason") message.Error.Type = stanza.ErrorTypeCancel message.Error.Reason = "undefined-condition" } } if envelope != nil { if envelope.Backend != omemo.Name { log.Errorf("Unsupported e2ee backend %q, dropping envelope", envelope.Backend) } else if ext, eme, err := omemoStanzaExtension(envelope.Raw); err != nil { log.Error(errors.Wrap(err, "Failed to encode OMEMO envelope")) } else { message.Extensions = append(message.Extensions, ext, eme) if message.Body == "" { message.Body = omemoFallbackBody } } } if oob != "" { message.Extensions = append(message.Extensions, stanza.OOB{ URL: oob, }) } if reply != nil { message.Extensions = append(message.Extensions, extensions.Reply{ To: reply.Author, Id: reply.Id, }) if reply.End > 0 { message.Extensions = append(message.Extensions, extensions.NewReplyFallback(reply.Start, reply.End)) } } if !isGroupchat && !isCarbon && toJid.Resource != "" && inviteFrom == "" { message.Extensions = append(message.Extensions, stanza.HintNoCopy{}) } if timestamp != 0 && mamQueryId == "" { var delayFrom string if isGroupchat { delayFrom = bareFrom } message.Extensions = append(message.Extensions, extensions.NewMessageDelay(timestamp, delayFrom)) message.Extensions = append(message.Extensions, extensions.NewMessageDelayLegacy(timestamp, delayFrom)) } if originalFrom != "" { message.Extensions = append(message.Extensions, extensions.MessageAddresses{ Addresses: []extensions.MessageAddress{ extensions.MessageAddress{ Type: "ofrom", Jid: originalFrom, }, }, }) } if subject == "" && forceSubject { message.Extensions = append(message.Extensions, extensions.EmptySubject{}) } if marker != nil { if marker.Type == MarkerTypeReceived { message.Extensions = append(message.Extensions, stanza.MarkReceived{ID: marker.Id}) } else if marker.Type == MarkerTypeDisplayed { message.Extensions = append(message.Extensions, stanza.MarkDisplayed{ID: marker.Id}) message.Extensions = append(message.Extensions, stanza.ReceiptReceived{ID: marker.Id}) } } if requestReceipt { message.Extensions = append(message.Extensions, stanza.Markable{}) } if replaceId != "" { message.Extensions = append(message.Extensions, extensions.Replace{Id: replaceId}) } if reactions != nil { reactionExts := make([]extensions.Reaction, len(reactions.Reactions)) for i, r := range reactions.Reactions { reactionExts[i] = extensions.Reaction{Text: r} } message.Extensions = append(message.Extensions, extensions.Reactions{Id: reactions.Id, Reactions: reactionExts}, // XEP-0334: body-less messages are commonly excluded from // MAM by default archiving policies, so hint explicitly. stanza.HintStore{}, ) } var userExt extensions.MessageXMucUserExtension if inviteFrom != "" { userExt.Invite = &extensions.MessageXMucUserInvite{ From: inviteFrom, } message.Extensions = append(message.Extensions, extensions.MessageXLegacyInviteExtension{ Jid: messageFrom, }) } if statusCode != 0 { userExt.Status = &extensions.MessageXMucUserStatus{ Code: strconv.FormatInt(statusCode, 10), } } if mucUserItem != nil { userExt.Item = extensions.PresenceXMucUserItem{ Affiliation: mucUserItem.Affiliation, Jid: &mucUserItem.Jid, Role: mucUserItem.Role, } } if inviteFrom != "" || statusCode != 0 || mucUserItem != nil { message.Extensions = append(message.Extensions, userExt) } if stanzaId != "" { message.Extensions = append(message.Extensions, extensions.MessageStanzaId{ Id: stanzaId, By: bareFrom, }) if stanzaId != id { message.Extensions = append(message.Extensions, extensions.MessageOriginId{ Id: id, }) } } if isCarbon { carbonMessage := extensions.ClientMessage{ Attrs: stanza.Attrs{ From: bareTo, To: to, Type: messageType, }, } carbonMessage.Extensions = append(carbonMessage.Extensions, extensions.CarbonSent{ Forwarded: stanza.Forwarded{ Stanza: extensions.ClientMessage(message), }, }) privilegeMessage := stanza.Message{ Attrs: stanza.Attrs{ From: Jid.Bare(), To: toJid.Domain, }, } if MessageOutgoingPermissionVersion == 2 { privilegeMessage.Extensions = append(privilegeMessage.Extensions, extensions.ComponentPrivilege2{ Forwarded: stanza.Forwarded{ Stanza: carbonMessage, }, }) } else { privilegeMessage.Extensions = append(privilegeMessage.Extensions, extensions.ComponentPrivilege1{ Forwarded: stanza.Forwarded{ Stanza: carbonMessage, }, }) } sendMessage(&privilegeMessage, component) } else if mamQueryId != "" { mneVpadluProbrasyvatJoshParametryPoraRefaktorit := strings.Split(mamQueryId, " ") ns := mneVpadluProbrasyvatJoshParametryPoraRefaktorit[0] mamQueryId = mneVpadluProbrasyvatJoshParametryPoraRefaktorit[1] delay := extensions.NewMessageDelay(timestamp, "") clientMessage := extensions.ClientMessage{ Attrs: message.Attrs, Subject: message.Subject, Body: message.Body, Thread: message.Thread, Error: message.Error, Extensions: message.Extensions, } forwarded := extensions.ForwardedMessage{ Delay: &delay, Message: &clientMessage, } var ext stanza.MsgExtension switch ns { case NS_MAM2: ext = extensions.MAM2MessageResult{ Id: stanzaId, QueryId: mamQueryId, Forwarded: &forwarded, } case NS_MAM1: ext = extensions.MAM1MessageResult{ Id: stanzaId, QueryId: mamQueryId, Forwarded: &forwarded, } case NS_MAM0: ext = extensions.MAM0MessageResult{ Id: stanzaId, QueryId: mamQueryId, Forwarded: &forwarded, } } mamMessage := stanza.Message{ Attrs: stanza.Attrs{ From: mucJID, To: to, Type: messageType, }, Extensions: []stanza.MsgExtension{ext}, } sendMessage(&mamMessage, component) } else { sendMessage(&message, component) } } // SetNickname sets a new nickname for a contact func SetNickname(to string, from string, nickname string, component *xmpp.Component) { componentJid := Jid.Bare() messageFrom := from + "@" + componentJid log.WithFields(log.Fields{ "from": from, "to": to, }).Warn("Set nickname") message := stanza.Message{ Attrs: stanza.Attrs{ From: messageFrom, To: to, Type: "headline", }, Extensions: []stanza.MsgExtension{ stanza.PubSubEvent{ EventElement: stanza.ItemsEvent{ Node: NSNick, Items: []stanza.ItemEvent{ stanza.ItemEvent{ Any: &stanza.Node{ XMLName: xml.Name{Space: NSNick, Local: "nick"}, Content: nickname, }, }, }, }, }, }, } sendMessage(&message, component) } func sendMessage(message *stanza.Message, component *xmpp.Component) { // explicit check, as marshalling is expensive if log.GetLevel() == log.DebugLevel { xmlMessage, err := xml.Marshal(message) if err == nil { log.Debug(string(xmlMessage)) } else { log.Debugf("%#v", message) } } _ = ResumableSend(component, message) } // LogBadPresence verbosely logs a presence func LogBadPresence(presence *stanza.Presence) { log.Errorf("Couldn't send presence: %#v", presence) } // SPFrom is a Telegram user id var SPFrom = args.NewString() // SPFullFrom is for specifying a full from when desired var SPFullFrom = args.NewString() // SPType is a presence type var SPType = args.NewString() // SPShow is a availability status var SPShow = args.NewString() // SPStatus is a verbose status var SPStatus = args.NewString() // SPNickname is a XEP-0172 nickname var SPNickname = args.NewString() // SPPhoto is a XEP-0153 hash of avatar in vCard var SPPhoto = args.NewString() // SPResource is an optional resource var SPResource = args.NewString() // SPImmed skips queueing var SPImmed = args.NewBool(args.Default(true)) // SPCaps is a XEP-0115 verification string var SPCaps = args.NewString() // SPMUCAffiliation is a XEP-0045 MUC affiliation var SPMUCAffiliation = args.NewString() // SPMUCRole is a XEP-0045 MUC role var SPMUCRole = args.NewString() // SPMUCNick is a XEP-0045 MUC user nick var SPMUCNick = args.NewString() // SPMUCJid is a real jid of a MUC member var SPMUCJid = args.NewString() // SPMUCStatusCodes is a set of XEP-0045 MUC status codes var SPMUCStatusCodes = args.New() // SPMUCDestroy is a XEP-0045 room destruction element var SPMUCDestroy = args.NewString() // SPToJids achieves to send the presence to certain full jids only var SPToJids = args.New() func newPresence(bareJid string, to string, args ...args.V) stanza.Presence { var presenceFrom string if SPFullFrom.IsSet(args) { presenceFrom = SPFullFrom.Get(args) } else if SPFrom.IsSet(args) { presenceFrom = SPFrom.Get(args) + "@" + bareJid if SPResource.IsSet(args) { resource := SPResource.Get(args) if resource != "" { presenceFrom += "/" + resource } } } else { presenceFrom = bareJid } presence := stanza.Presence{Attrs: stanza.Attrs{ From: presenceFrom, To: to, }} if SPType.IsSet(args) { t := SPType.Get(args) if t != "" { presence.Attrs.Type = stanza.StanzaType(t) } } if SPShow.IsSet(args) { show := SPShow.Get(args) if show != "" { presence.Show = stanza.PresenceShow(show) } } if SPStatus.IsSet(args) { status := SPStatus.Get(args) if status != "" { presence.Status = status } } if SPNickname.IsSet(args) { nickname := SPNickname.Get(args) if nickname != "" { presence.Extensions = append(presence.Extensions, extensions.PresenceNickExtension{ Text: nickname, }) } } if SPPhoto.IsSet(args) { photo := SPPhoto.Get(args) if photo != "" { presence.Extensions = append(presence.Extensions, extensions.PresenceXVCardUpdateExtension{ Photo: extensions.PresenceXVCardUpdatePhoto{ Text: photo, }, }) } } if SPCaps.IsSet(args) { ver := SPCaps.Get(args) if ver != "" { presence.Extensions = append(presence.Extensions, stanza.Caps{ Hash: "sha-1", Node: CapsNode, Ver: ver, }) } } if SPMUCAffiliation.IsSet(args) { affiliation := SPMUCAffiliation.Get(args) if affiliation != "" { var role string if SPMUCRole.IsSet(args) { role = SPMUCRole.Get(args) } else { role = affiliationToRole(affiliation) } userExt := extensions.PresenceXMucUserExtension{ Item: extensions.PresenceXMucUserItem{ Affiliation: affiliation, Role: role, }, } if SPMUCNick.IsSet(args) { userExt.Item.Nick = SPMUCNick.Get(args) } if SPMUCJid.IsSet(args) { mucJid := SPMUCJid.Get(args) userExt.Item.Jid = &mucJid } if SPMUCStatusCodes.IsSet(args) { statusCodes := SPMUCStatusCodes.Get(args).([]uint16) for _, statusCode := range statusCodes { userExt.Statuses = append(userExt.Statuses, extensions.PresenceXMucUserStatus{ Code: statusCode, }) } } if SPMUCDestroy.IsSet(args) { userExt.Destroy = &extensions.MucDestroy{ Jid: SPMUCDestroy.Get(args), Reason: "Group was deleted", } } presence.Extensions = append(presence.Extensions, userExt) } } return presence } // SendPresence creates and sends a presence stanza func SendPresence(component *xmpp.Component, to string, args ...args.V) error { var logFrom string bareJid := Jid.Bare() if SPFrom.IsSet(args) { logFrom = SPFrom.Get(args) } else { logFrom = bareJid } log.WithFields(log.Fields{ "type": SPType.Get(args), "from": logFrom, "to": to, }).Info("Got presence") var tos []string if SPToJids.IsSet(args) { tos = SPToJids.Get(args).([]string) } else { tos = []string{to} } for _, to := range tos { presence := newPresence(bareJid, to, args...) // explicit check, as marshalling is expensive if log.GetLevel() == log.DebugLevel { xmlPresence, err := xml.Marshal(presence) if err == nil { log.Debug(string(xmlPresence)) } else { log.Debugf("%#v", presence) } } immed := SPImmed.Get(args) if immed { err := ResumableSend(component, presence) if err != nil { LogBadPresence(&presence) return err } } else { QueueLock.Lock() Queue[presence.From+presence.To] = &presence QueueLock.Unlock() } } return nil } // SPAppendFrom appends numeric from and resource to varargs func SPAppendFrom(oldArgs []args.V, id int64) []args.V { newArgs := append(oldArgs, SPFrom(CHATNODE(id))) newArgs = append(newArgs, SPResource(Jid.Resource)) return newArgs } // isSubscriptionType reports whether a presence type manages the subscription // state. Per RFC 6121 ยง 3 such presences must be addressed from a bare JID; // attaching the component resource makes clients treat the full JID as a // separate entity and produces a duplicate subscription request. func isSubscriptionType(typ string) bool { switch typ { case "subscribe", "subscribed", "unsubscribe", "unsubscribed": return true } return false } // SimplePresence crafts simple presence varargs func SimplePresence(from int64, typ string) []args.V { args := []args.V{SPType(typ)} if isSubscriptionType(typ) { // subscription presences must come from the bare JID args = append(args, SPFrom(CHATNODE(from))) } else { args = SPAppendFrom(args, from) } return args } // ResumableSend tries to resume the connection once and sends the packet again func ResumableSend(component *xmpp.Component, packet stanza.Packet) error { err := component.Send(packet) if err != nil && strings.HasPrefix(err.Error(), "cannot send packet") { log.Warn("Packet send failed, trying to resume the connection...") err = component.Connect() if err == nil { err = component.Send(packet) } } if err != nil { log.Error(err.Error()) } return err } // SubscribeToTransport ensures a two-way subscription to the transport func SubscribeToTransport(component *xmpp.Component, jid string) { SendPresence(component, jid, SPType("subscribe")) SendPresence(component, jid, SPType("subscribed")) } // SplitJID tokenizes a JID string to bare JID and resource func SplitJID(from string) (string, string, bool) { fromJid, err := stanza.NewJid(from) if err != nil { log.WithFields(log.Fields{ "from": from, }).Error(errors.Wrap(err, "Invalid from JID!")) return "", "", false } return fromJid.Bare(), fromJid.Resource, true } func affiliationToRole(affilation string) string { switch affilation { case "owner", "admin": return "moderator" case "member": return "participant" } return "none" } // SendPubSubAvatarNotification encourages clients to fetch an avatar func SendPubSubAvatarNotification(component *xmpp.Component, jid string, chatJid string, sha1 string, size int64) { info := stanza.Node{ XMLName: xml.Name{Local: "info"}, Attrs: []xml.Attr{ xml.Attr{Name: xml.Name{Local: "bytes"}, Value: strconv.FormatInt(size, 10)}, xml.Attr{Name: xml.Name{Local: "height"}, Value: "160"}, xml.Attr{Name: xml.Name{Local: "id"}, Value: sha1}, xml.Attr{Name: xml.Name{Local: "type"}, Value: "image/jpeg"}, xml.Attr{Name: xml.Name{Local: "width"}, Value: "160"}, }, } log.WithFields(log.Fields{ "chatJid": chatJid, }).Debugf("%#v", info) event := &stanza.PubSubEvent{ EventElement: &stanza.ItemsEvent{ Node: NodeAvatarMetadata, Items: []stanza.ItemEvent{ stanza.ItemEvent{ Id: sha1, Any: &stanza.Node{ XMLName: xml.Name{Local: "metadata", Space: NodeAvatarMetadata}, Nodes: []stanza.Node{info}, }, }, }, }, } message := stanza.Message{ Attrs: stanza.Attrs{ From: chatJid, To: jid, Type: stanza.MessageTypeHeadline, }, Extensions: []stanza.MsgExtension{event}, } _ = ResumableSend(component, message) } // SendPubSubDeviceListNotification pings jid (typically the real user's // bare JID, reaching every connected resource per XEP-0163's personal // eventing model - no explicit subscription needed for a "+notify" // feature) that chatJid's e2ee device-list content published under node // may have changed, embedding raw (an e2ee.DeviceListDoc's already- // marshaled XML) as the event's item - mirrors SendPubSubAvatarNotification's // shape (a headline-type with a pubsub ). This lets any of // the real user's OTHER devices discover a freshly published identity // without needing to have initiated OMEMO themselves. func SendPubSubDeviceListNotification(component *xmpp.Component, jid string, chatJid string, node string, raw []byte) error { var content stanza.Node if err := xml.Unmarshal(raw, &content); err != nil { return err } event := &stanza.PubSubEvent{ EventElement: &stanza.ItemsEvent{ Node: node, Items: []stanza.ItemEvent{ {Id: "current", Any: &content}, }, }, } message := stanza.Message{ Attrs: stanza.Attrs{ From: chatJid, To: jid, Type: stanza.MessageTypeHeadline, }, Extensions: []stanza.MsgExtension{event}, } return ResumableSend(component, message) } func InviteToMUC(chatID int64, jid string, component *xmpp.Component) { SendMUCInvite(jid, MUCNODE(chatID), component, Jid.Full()) } // AddPMDiscoFeatures appends the disco#info identity/features for a 1:1 Telegram contact. func AddPMDiscoFeatures(disco *stanza.DiscoInfo) { disco.AddIdentity("", "account", "registered") if backend, ok := E2EE.Backend(); ok { disco.AddFeatures(backend.Namespaces().Disco...) } disco.AddFeatures(stanza.NSMsgChatMarkers) disco.AddFeatures(stanza.NSMsgReceipts) disco.AddFeatures(NSReactions) // Jingle / JMI features so clients (e.g. Dino) treat the contact as // call-capable. Without these, JMI never fires: the client // disco-info's the caller before engaging its call state machine and // silently drops the call when it sees no RTP support advertised. disco.AddFeatures( "urn:xmpp:jingle:1", "urn:xmpp:jingle-message:0", "urn:xmpp:jingle:apps:rtp:1", "urn:xmpp:jingle:apps:rtp:audio", "urn:xmpp:jingle:transports:ice-udp:1", "urn:xmpp:jingle:apps:dtls:0", ) } // AddTransportDiscoFeatures appends the disco#info identity/features for the transport root JID. func AddTransportDiscoFeatures(disco *stanza.DiscoInfo) { disco.AddIdentity("Telegram Gateway", "gateway", "telegram") disco.AddFeatures("jabber:iq:register") } // AddCommonDiscoFeatures appends features common to every root-node disco#info response. func AddCommonDiscoFeatures(disco *stanza.DiscoInfo) { disco.AddFeatures( NSCommand, "jabber:iq:version", "urn:xmpp:time", stanza.NSDiscoInfo, NSCaps, ) } // ReactionRestrictionsForm builds the XEP-0444 disco#info restricted- // reactions data form (urn:xmpp:reactions:0:restrictions). Shared by the // real disco#info response and PMCapsVer so they stay identical. func ReactionRestrictionsForm(allowlist []string, maxPerUser int) *stanza.Form { return stanza.NewForm([]*stanza.Field{ &stanza.Field{ Var: "FORM_TYPE", Type: "hidden", ValuesList: []string{"urn:xmpp:reactions:0:restrictions"}, }, &stanza.Field{ Var: "allowlist", ValuesList: allowlist, }, &stanza.Field{ Var: "max_reactions_per_user", ValuesList: []string{strconv.Itoa(maxPerUser)}, }, }, "result") } // PMCapsVer computes the XEP-0115 ver for a 1:1 contact's disco#info. // forms must match that contact's live disco#info response's own Forms. // Not otherwise cached: if AddPMDiscoFeatures ever starts varying // silently (no signature change), a cached ver would go stale with no // compiler warning. func PMCapsVer(forms ...*stanza.Form) string { disco := &stanza.DiscoInfo{} AddPMDiscoFeatures(disco) AddCommonDiscoFeatures(disco) disco.Forms = forms return DiscoVer(disco) } // DiscoVer computes a XEP-0115 ver string for a disco#info response. func DiscoVer(disco *stanza.DiscoInfo) string { hash := sha1.New() discoToCaps(disco, hash) return base64.StdEncoding.EncodeToString(hash.Sum(nil)) } // discoToCaps writes disco's XEP-0115 S5.1 capabilities identifier string to w. func discoToCaps(disco *stanza.DiscoInfo, w io.Writer) { const sep = "<" identities := make([]string, len(disco.Identity)) for i, identity := range disco.Identity { identities[i] = fmt.Sprintf("%s/%s//%s", identity.Category, identity.Type, identity.Name) } sort.Strings(identities) for _, identity := range identities { io.WriteString(w, identity) io.WriteString(w, sep) } vars := make([]string, len(disco.Features)) for i, feature := range disco.Features { vars[i] = feature.Var } sort.Strings(vars) for _, v := range vars { io.WriteString(w, v) io.WriteString(w, sep) } var forms []*stanza.Form for _, f := range disco.Forms { if f != nil { forms = append(forms, f) } } // XEP-0115 S5.1 step 6: with more than one extended form, sort by // FORM_TYPE value before processing each in turn. sort.Slice(forms, func(i, j int) bool { return formType(forms[i]) < formType(forms[j]) }) for _, form := range forms { fields := make([]*stanza.Field, len(form.Fields)) copy(fields, form.Fields) sort.Slice(fields, func(i, j int) bool { a, b := fields[i], fields[j] if a.Var == "FORM_TYPE" { return true } if b.Var == "FORM_TYPE" { return false } return a.Var < b.Var }) for _, field := range fields { if field.Var == "FORM_TYPE" { if len(field.ValuesList) > 0 { io.WriteString(w, field.ValuesList[0]) io.WriteString(w, sep) } continue } io.WriteString(w, field.Var) io.WriteString(w, sep) values := field.ValuesList if len(values) > 1 { values = append([]string(nil), values...) sort.Strings(values) } for _, value := range values { io.WriteString(w, value) io.WriteString(w, sep) } } } } // formType returns form's FORM_TYPE field value, or "" if it has none. func formType(form *stanza.Form) string { for _, field := range form.Fields { if field.Var == "FORM_TYPE" && len(field.ValuesList) > 0 { return field.ValuesList[0] } } return "" }