Implement OMEMO IQ queries

This commit is contained in:
Bohdan Horbeshko 2026-07-29 01:32:16 -04:00
parent 44ff56742b
commit 2b349c9f4c
6 changed files with 223 additions and 7 deletions

View file

@ -18,6 +18,9 @@ func (b *Backend) Encrypt(from e2ee.PeerID, to []e2ee.PeerID, plaintext []byte)
defer b.mu.Unlock()
store := b.db.Store(b.libctx, login, owner)
if err := b.ensureIdentity(store); err != nil {
return e2ee.Envelope{}, fmt.Errorf("omemo: ensureIdentity: %w", err)
}
// Which variant to actually encrypt as is negotiated per chat, not
// fixed on the backend - see negotiatedVersion's doc comment.
@ -110,6 +113,9 @@ func (b *Backend) Decrypt(from, to e2ee.PeerID, env e2ee.Envelope) ([]byte, erro
defer b.mu.Unlock()
store := b.db.Store(b.libctx, login, owner)
if err := b.ensureIdentity(store); err != nil {
return nil, fmt.Errorf("omemo: ensureIdentity: %w", err)
}
regID, err := store.GetLocalRegistrationID()
if err != nil {
return nil, fmt.Errorf("omemo: GetLocalRegistrationID: %w", err)

View file

@ -66,6 +66,9 @@ func (b *Backend) Name() string { return Name }
// Close implements e2ee.Backend.
func (b *Backend) Close() error { return b.db.Close() }
// OwnDevice implements e2ee.Backend.
func (b *Backend) OwnDevice() e2ee.DeviceID { return OwnDeviceID }
func (b *Backend) Namespaces() e2ee.Namespaces {
disco := []string{omemo0NS + ".devicelist+notify"}
if libsignal.ProtocolV4Supported {
@ -193,8 +196,21 @@ func (b *Backend) EnsureIdentity(peer e2ee.PeerID) error {
b.mu.Lock()
defer b.mu.Unlock()
store := b.db.Store(b.libctx, login, owner)
return b.ensureIdentity(b.db.Store(b.libctx, login, owner))
}
// ensureIdentity is EnsureIdentity's store-scoped counterpart, called
// defensively at the top of every other method that touches identity/
// prekey material (PublishedIdentity, PublishedBundle, Encrypt, Decrypt) -
// nothing in this package used to call the public EnsureIdentity at all,
// so every one of those would have failed with "EnsureIdentity was not
// called" the first time a chat went OMEMO-active, whichever entry point
// hit first (an inbound decrypt, an outbound encrypt, or a peer's PEP GET
// arriving before either). Bootstrapping here instead of trusting callers
// to remember makes that invariant impossible to violate. store must
// already be scoped to the right (login, owner) and b.mu must already be
// held.
func (b *Backend) ensureIdentity(store *badgerstore.Store) error {
has, err := store.HasIdentityKeyPair()
if err != nil {
return fmt.Errorf("omemo: HasIdentityKeyPair: %w", err)
@ -258,10 +274,8 @@ func (b *Backend) PublishedIdentity(peer e2ee.PeerID, variant e2ee.Variant) (e2e
defer b.mu.Unlock()
store := b.db.Store(b.libctx, login, owner)
if has, err := store.HasIdentityKeyPair(); err != nil {
if err := b.ensureIdentity(store); err != nil {
return e2ee.DeviceListDoc{}, err
} else if !has {
return e2ee.DeviceListDoc{}, errors.New("omemo: PublishedIdentity: EnsureIdentity was not called")
}
raw, err := encodeDeviceList(version, ownDeviceIDNum)
@ -288,20 +302,23 @@ func (b *Backend) PublishedBundle(peer e2ee.PeerID, device e2ee.DeviceID, varian
defer b.mu.Unlock()
store := b.db.Store(b.libctx, login, owner)
if err := b.ensureIdentity(store); err != nil {
return e2ee.BundleDoc{}, err
}
identityPublic, _, err := store.GetIdentityKeyPair()
if err != nil {
return e2ee.BundleDoc{}, fmt.Errorf("omemo: GetIdentityKeyPair: %w", err)
}
// EnsureIdentity always creates signed prekey id 1 and never rotates it
// ensureIdentity always creates signed prekey id 1 and never rotates it
// yet (a future enhancement), so id 1 is the only one that will exist.
spkRecord, found, err := store.LoadSignedPreKey(1)
if err != nil {
return e2ee.BundleDoc{}, fmt.Errorf("omemo: LoadSignedPreKey: %w", err)
}
if !found {
return e2ee.BundleDoc{}, errors.New("omemo: PublishedBundle: EnsureIdentity was not called")
return e2ee.BundleDoc{}, errors.New("omemo: PublishedBundle: signed prekey 1 missing despite ensureIdentity succeeding")
}
spkInfo, err := libsignal.DecodeSignedPreKey(b.libctx, spkRecord)
if err != nil {

View file

@ -101,9 +101,21 @@ type Backend interface {
// OwnedPeer, since a chat pseudo-JID alone isn't globally unique across
// telegabber's multiple bridged Telegram logins), generating and
// persisting whatever key material the backend needs. One identity
// serves every Variant simultaneously.
// serves every Variant simultaneously. Every other method that touches
// identity/prekey material calls this internally too (it's cheap and
// idempotent), so callers never strictly need to call it themselves -
// it's exposed mainly for a future eager-bootstrap use case (e.g.
// pre-creating identities for every personal chat when an account's
// OMEMO toggle is turned on, instead of lazily on first use).
EnsureIdentity(peer PeerID) error
// OwnDevice is this backend's own (and only) published device id for
// any owned peer - telegabber only ever needs one device per bridged
// chat identity (see EnsureIdentity's doc comment), so unlike Devices
// (which lists a remote peer's possibly-many devices) this takes no
// arguments and returns no error.
OwnDevice() DeviceID
// PublishedIdentity returns the material this gateway must SERVE to
// pubsub GET requests for its own PeerID (an OwnedPeer - see
// EnsureIdentity) device list under the given variant, pre-marshaled to

117
xmpp/e2ee.go Normal file
View file

@ -0,0 +1,117 @@
package xmpp
import (
"encoding/xml"
"github.com/pkg/errors"
"dev.narayana.im/narayana/telegabber/e2ee"
"dev.narayana.im/narayana/telegabber/xmpp/gateway"
log "github.com/sirupsen/logrus"
"gosrc.io/xmpp"
"gosrc.io/xmpp/stanza"
)
// handleGetE2EEPubSubIq answers a PEP GET request for one of the e2ee
// backend's own device-list/bundle nodes, in whichever of its Variants the
// request actually names - available whenever the E2EE feature is
// deploy-enabled at all (config's :xmpp: :e2ee: :enabled), regardless of
// whether this specific chat has actually gone OMEMO-active yet. This is
// deliberate, not an oversight: a real client can only ever decide to
// start speaking OMEMO to us once it's already able to fetch our bundle,
// so serving must not be gated behind the same per-chat/account flag that
// gates outbound encryption (see the plan's enablement model).
func handleGetE2EEPubSubIq(s xmpp.Sender, iq *stanza.IQ, pubsub *stanza.PubSubGeneric) {
component, answer, ok := iqResultStub(s, iq)
if !ok {
return
}
defer gateway.ResumableSend(component, answer)
backend, ok := gateway.E2EE.Backend()
if !ok {
iqAnswerSetError(answer, 404)
return
}
_, toOk, toIsGroup := toToID(iq.To)
if !toOk || toIsGroup {
// e2ee identities only ever exist for personal-chat pseudo-JIDs -
// true MUC/XEP-0045 groupchat OMEMO is out of scope entirely (see
// the plan doc).
iqAnswerSetError(answer, 404)
return
}
toJid, err := stanza.NewJid(iq.To)
if err != nil {
iqAnswerSetError(answer, 400)
return
}
bare, _, fromOk := gateway.SplitJID(iq.From)
if !fromOk {
iqAnswerSetError(answer, 400)
return
}
session, sessionOk := sessions[bare]
if !sessionOk {
log.Errorf("E2EE PEP IQ from stranger %v", iq.From)
iqAnswerSetError(answer, 403)
return
}
owner := e2ee.OwnedPeer(session.Session.Login, toJid.Bare())
if err := backend.EnsureIdentity(owner); err != nil {
log.Error(errors.Wrap(err, "Failed to ensure e2ee identity"))
iqAnswerSetError(answer, 500)
return
}
for _, variant := range backend.Variants() {
if pubsub.Items.Node == backend.DeviceListNode(variant) {
doc, err := backend.PublishedIdentity(owner, variant)
if err != nil {
log.Error(errors.Wrap(err, "Failed to publish e2ee device list"))
iqAnswerSetError(answer, 500)
return
}
setE2EEItemsAnswer(answer, pubsub.Items.Node, "current", doc.Raw)
return
}
if pubsub.Items.Node == backend.BundleNode(variant, backend.OwnDevice()) {
doc, err := backend.PublishedBundle(owner, backend.OwnDevice(), variant)
if err != nil {
log.Error(errors.Wrap(err, "Failed to publish e2ee bundle"))
iqAnswerSetError(answer, 500)
return
}
setE2EEItemsAnswer(answer, pubsub.Items.Node, string(backend.OwnDevice()), doc.Raw)
return
}
}
// Node didn't match any variant this backend currently supports.
iqAnswerSetError(answer, 404)
}
// setE2EEItemsAnswer embeds raw (an e2ee.DeviceListDoc/BundleDoc's already-
// marshaled XML) as the single <item> of a pubsub#items result - raw is
// re-parsed into a generic stanza.Node rather than injected as a string,
// since stanza.Item.Any expects a structured node, not raw bytes.
func setE2EEItemsAnswer(answer *stanza.IQ, node, itemID string, raw []byte) {
var content stanza.Node
if err := xml.Unmarshal(raw, &content); err != nil {
log.Error(errors.Wrap(err, "Failed to unmarshal e2ee document"))
iqAnswerSetError(answer, 500)
return
}
answer.Payload = &stanza.PubSubGeneric{
Items: &stanza.Items{
Node: node,
List: []stanza.Item{
{Id: itemID, Any: &content},
},
},
}
}

View file

@ -994,6 +994,42 @@ func SendPubSubAvatarNotification(component *xmpp.Component, jid string, chatJid
_ = 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 <message> with a pubsub <event>). 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())
}

View file

@ -63,6 +63,14 @@ func HandleIq(s xmpp.Sender, p stanza.Packet) {
go handleGetAvatarDataIq(s, iq, pubsub)
return
}
if backend, ok := gateway.E2EE.Backend(); ok {
for _, variant := range backend.Variants() {
if pubsub.Items.Node == backend.DeviceListNode(variant) || pubsub.Items.Node == backend.BundleNode(variant, backend.OwnDevice()) {
go handleGetE2EEPubSubIq(s, iq, pubsub)
return
}
}
}
}
discoInfo, ok := iq.Payload.(*stanza.DiscoInfo)
if ok {
@ -272,6 +280,23 @@ func HandleMessage(s xmpp.Sender, p stanza.Packet) {
if err := backend.SetEnabled(owner, true); err != nil {
log.Error(errors.Wrap(err, "Failed to mark chat as OMEMO-active"))
}
// This chat's OMEMO involvement just became real
// for the first time - ping the real user's
// OTHER devices (if any) that our identity is
// available, in case they haven't independently
// discovered it yet (see
// gateway.SendPubSubDeviceListNotification's doc
// comment).
for _, variant := range backend.Variants() {
doc, docErr := backend.PublishedIdentity(owner, variant)
if docErr != nil {
log.Error(errors.Wrap(docErr, "Failed to publish e2ee device list for notification"))
continue
}
if err := gateway.SendPubSubDeviceListNotification(component, bare, toJid.Bare(), backend.DeviceListNode(variant), doc.Raw); err != nil {
log.Error(errors.Wrap(err, "Failed to send e2ee device-list notification"))
}
}
}
}
}
@ -1120,6 +1145,9 @@ func handleGetDiscoInfo(s xmpp.Sender, iq *stanza.IQ, di *stanza.DiscoInfo) {
}
} else {
disco.AddIdentity("", "account", "registered")
if backend, ok := gateway.E2EE.Backend(); ok {
disco.AddFeatures(backend.Namespaces().Disco...)
}
}
disco.AddFeatures(stanza.NSMsgChatMarkers)
disco.AddFeatures(stanza.NSMsgReceipts)