telegabber/calls/signaling/xmppsig/adapter.go
2026-05-25 06:27:20 -07:00

344 lines
11 KiB
Go

// xmpp side of signaling.Bridge; one Adapter per session,
// sharing a *jingle.Manager for stanza dispatch
package xmppsig
import (
"fmt"
"strconv"
"strings"
"sync"
"sync/atomic"
"dev.narayana.im/narayana/telegabber/calls/signaling"
"dev.narayana.im/narayana/telegabber/xmpp/jingle"
"github.com/pion/webrtc/v4"
log "github.com/sirupsen/logrus"
)
type AdapterConfig struct {
Sender jingle.Sender // *xmpp.Component satisfies this
LocalJID string // gateway's bare JID
Manager *jingle.Manager
// fresh PC per call, with the transceivers/tracks SDP negotiation needs
PCFactory func() (*webrtc.PeerConnection, error)
}
type Adapter struct {
cfg AdapterConfig
}
func NewAdapter(cfg AdapterConfig) *Adapter { return &Adapter{cfg: cfg} }
func (a *Adapter) Manager() *jingle.Manager { return a.cfg.Manager }
func (a *Adapter) LocalJID() string { return a.cfg.LocalJID }
// xmpp side for a TG-originated call; PC and session built upfront so a
// failure unwinds both. Bind must run before bridge.Start.
func (a *Adapter) NewCaller(remoteJID string, callerUserID int64) (*Caller, error) {
pc, err := a.cfg.PCFactory()
if err != nil {
return nil, fmt.Errorf("PCFactory: %w", err)
}
// XEP-0353 requires a full JID; /telegabber matches the gateway's other
// per-contact outbound stanzas
localJID := fmt.Sprintf("%d@%s/telegabber", callerUserID, a.cfg.LocalJID)
c := &Caller{xmppBase: xmppBase{
a: a, remote: remoteJID, pc: pc, localJID: localJID,
side: signaling.CallerSide,
}}
// must install OnTrack forwarder before SetRemoteDescription, otherwise
// pion fires OnTrack during SDP processing and a late handler misses it
c.hookTrack(pc)
sid := newSID()
sess, err := jingle.New(jingle.SessionOpts{
PC: pc, Sender: a.cfg.Sender,
LocalJID: localJID, RemoteJID: remoteJID,
SID: sid, Role: jingle.RoleInitiator, Media: []string{"audio"},
Observer: c,
})
if err != nil {
_ = pc.Close()
return nil, fmt.Errorf("jingle.New: %w", err)
}
c.set(sid, sess)
return c, nil
}
// xmpp side for an XMPP-originated call; Bind must run before bridge.Start
func (a *Adapter) NewCallee(p jingle.IncomingProposal) (*Callee, *jingle.Session, error) {
pc, err := a.cfg.PCFactory()
if err != nil {
return nil, nil, fmt.Errorf("PCFactory: %w", err)
}
c := &Callee{xmppBase: xmppBase{a: a, remote: p.From, pc: pc, side: signaling.CalleeSide}}
c.hookTrack(pc)
// LocalJID is the propose's `to`, not the bare component domain;
// anotherim's libdino filters proceed with from.equals_bare(peer_state.jid)
// so a proceed from the bare component domain silently fails to match
// and anotherim never sends session-initiate
localJID := p.To
if localJID == "" {
localJID = a.cfg.LocalJID
}
// Conversations' RtpSessionActivity crashes on bare-JID `with` once ICE
// completes (IllegalStateException "No RTP connection found"); real
// callees always send <proceed> from a full JID. Stick the SID on as
// resource so concurrent calls don't collide.
if !strings.ContainsRune(localJID, '/') && p.SID != "" {
localJID = localJID + "/" + p.SID
}
sess, err := jingle.New(jingle.SessionOpts{
PC: pc, Sender: a.cfg.Sender,
LocalJID: localJID, RemoteJID: p.From,
SID: p.SID, Role: jingle.RoleResponder, Media: p.Media,
Observer: c,
})
if err != nil {
_ = pc.Close()
return nil, nil, err
}
c.set(p.SID, sess)
return c, sess, nil
}
// wires pion's connection-state callback to bridge media transitions.
// Disconnected can be transient (ICE restart) so only Failed reacts.
// side dedups MediaConnected per source. Called from xmppBase.Bind.
func hookPCState(pc *webrtc.PeerConnection, br *signaling.Bridge, side signaling.Side) {
if pc == nil || br == nil {
return
}
// DTLS logger lazily on first state change - the transceiver's
// DTLSTransport doesn't exist until SetLocalDescription runs
var dtlsHookOnce sync.Once
pc.OnConnectionStateChange(func(s webrtc.PeerConnectionState) {
dtlsHookOnce.Do(func() { hookDTLSState(pc) })
entry := log.WithField("state", s.String())
// log the selected ICE pair on Connected - flaky calls have happened
// when pion picked a docker-bridge or Tailscale pair whose STUN pings
// succeed but real RTP doesn't transit. SCTP isn't negotiated for
// audio-only, so reach the ICE transport via any RTP sender's
// DTLSTransport (BUNDLE means they all share one).
if s == webrtc.PeerConnectionStateConnected {
var dtls *webrtc.DTLSTransport
for _, snd := range pc.GetSenders() {
if snd != nil && snd.Transport() != nil {
dtls = snd.Transport()
break
}
}
if dtls == nil {
for _, r := range pc.GetReceivers() {
if r != nil && r.Transport() != nil {
dtls = r.Transport()
break
}
}
}
switch {
case dtls == nil:
entry = entry.WithField("ice_pair", "no-dtls-transport")
case dtls.ICETransport() == nil:
entry = entry.WithField("ice_pair", "no-ice-transport")
default:
pair, err := dtls.ICETransport().GetSelectedCandidatePair()
switch {
case err != nil:
entry = entry.WithError(err).WithField("ice_pair", "get-selected-failed")
case pair == nil || pair.Local == nil || pair.Remote == nil:
entry = entry.WithField("ice_pair", "nil")
default:
entry = entry.WithFields(log.Fields{
"local": fmt.Sprintf("%s:%d (%s)", pair.Local.Address, pair.Local.Port, pair.Local.Typ.String()),
"remote": fmt.Sprintf("%s:%d (%s)", pair.Remote.Address, pair.Remote.Port, pair.Remote.Typ.String()),
"local_related": pair.Local.RelatedAddress,
})
}
}
}
entry.Debug("pion: PC connection state")
switch s {
case webrtc.PeerConnectionStateConnected:
br.MediaConnected(side)
case webrtc.PeerConnectionStateFailed:
br.MediaFailed(signaling.ReasonMediaFailed)
}
})
pc.OnICEConnectionStateChange(func(s webrtc.ICEConnectionState) {
log.WithField("state", s.String()).Debug("pion: ICE connection state")
})
pc.OnICEGatheringStateChange(func(s webrtc.ICEGatheringState) {
log.WithField("state", s.String()).Debug("pion: ICE gathering state")
})
// do NOT register pc.OnICECandidate here - jingle.New attaches the
// session's trickle handler there, pion's API is replace-only, and
// overwriting it strands srflx/relay candidates -> breaks TURN-only paths
}
// logs DTLS state - separates "handshake never completed" from "completed
// but no audio" (SRTP/SSRC/codec). Called from the first state change so
// the transport actually exists.
func hookDTLSState(pc *webrtc.PeerConnection) {
var dtls *webrtc.DTLSTransport
for _, s := range pc.GetSenders() {
if s != nil && s.Transport() != nil {
dtls = s.Transport()
break
}
}
if dtls == nil {
for _, r := range pc.GetReceivers() {
if r != nil && r.Transport() != nil {
dtls = r.Transport()
break
}
}
}
if dtls == nil {
log.Debug("pion: no DTLSTransport available at hook time; DTLS state changes will go unlogged")
return
}
dtls.OnStateChange(func(s webrtc.DTLSTransportState) {
log.WithField("state", s.String()).Info("pion: DTLS transport state")
})
}
func mapReason(r signaling.TerminationReason) string {
switch r {
case signaling.ReasonHangup:
return jingle.ReasonSuccess
case signaling.ReasonDecline:
return jingle.ReasonDecline
case signaling.ReasonBusy:
return jingle.ReasonBusy
case signaling.ReasonRingTimeout:
return jingle.ReasonTimeout
case signaling.ReasonExchangeTimeout, signaling.ReasonConnectTimeout:
return jingle.ReasonConnectivityError
case signaling.ReasonMediaFailed:
return jingle.ReasonFailedTransport
case signaling.ReasonPeerGone:
return jingle.ReasonGone
default:
return jingle.ReasonGeneralError
}
}
func reasonFromCondition(cond string) signaling.TerminationReason {
switch cond {
case "rejected", jingle.ReasonDecline:
return signaling.ReasonDecline
case "retracted", jingle.ReasonGone, jingle.ReasonCancel:
return signaling.ReasonPeerGone
case jingle.ReasonBusy:
return signaling.ReasonBusy
case jingle.ReasonTimeout:
return signaling.ReasonRingTimeout
case jingle.ReasonConnectivityError:
return signaling.ReasonConnectTimeout
case jingle.ReasonFailedTransport, jingle.ReasonFailedApplication:
return signaling.ReasonMediaFailed
default:
return signaling.ReasonHangup
}
}
var sidCounter atomic.Uint64
func newSID() string {
return fmt.Sprintf("tg-call-%d", sidCounter.Add(1))
}
// telegram user-id from a JID localpart: "12345@gw.example" -> (12345, true)
func ParseTargetJID(jid string) (int64, bool) {
parts := strings.SplitN(jid, "@", 2)
if len(parts) == 0 || parts[0] == "" {
return 0, false
}
id, err := strconv.ParseInt(parts[0], 10, 64)
if err != nil {
return 0, false
}
return id, true
}
// PC with no ICE servers - for tests; production uses MakePCFactory
func DefaultPCFactory() (*webrtc.PeerConnection, error) {
return buildPC(nil)
}
// PCFactory pulling ICE servers from provider per call, so short-lived TURN
// creds refresh without restarting the component
func MakePCFactory(provider func() []webrtc.ICEServer) func() (*webrtc.PeerConnection, error) {
return func() (*webrtc.PeerConnection, error) {
var servers []webrtc.ICEServer
if provider != nil {
servers = provider()
}
return buildPC(servers)
}
}
func buildPC(iceServers []webrtc.ICEServer) (*webrtc.PeerConnection, error) {
// SetHandleUndeclaredSSRCWithoutAnswer: anotherim's libnice starts
// sending RTP the moment ICE is selected, which often beats the local
// SetRemoteDescription(answer) - so pion drops the inbound track with
// "Incoming unhandled RTP ssrc(N), OnTrack will not be fired" and the
// peer's audio never reaches TG. With this flag pion buffers early RTP.
s := webrtc.SettingEngine{}
s.SetHandleUndeclaredSSRCWithoutAnswer(true)
m := &webrtc.MediaEngine{}
if err := m.RegisterDefaultCodecs(); err != nil {
return nil, fmt.Errorf("RegisterDefaultCodecs: %w", err)
}
api := webrtc.NewAPI(webrtc.WithSettingEngine(s), webrtc.WithMediaEngine(m))
pc, err := api.NewPeerConnection(webrtc.Configuration{ICEServers: iceServers})
if err != nil {
return nil, err
}
// placeholder opus track BEFORE the transceiver (not just a kind).
// Otherwise pion's RTPSender.Send leaves hasSent=false, the later
// ReplaceTrack short-circuits without Bind, the real track has no
// packetizer, and WriteSample silently drops every sample. Caller
// path always hit this because setupAudio attaches the real track
// after SetLocalDescription.
placeholder, err := webrtc.NewTrackLocalStaticSample(
webrtc.RTPCodecCapability{
MimeType: webrtc.MimeTypeOpus,
ClockRate: 48000,
Channels: 2,
SDPFmtpLine: "minptime=10;useinbandfec=1",
},
"audio-placeholder",
"telegabber-placeholder",
)
if err != nil {
_ = pc.Close()
return nil, fmt.Errorf("placeholder track: %w", err)
}
tr, err := pc.AddTransceiverFromTrack(placeholder, webrtc.RTPTransceiverInit{
Direction: webrtc.RTPTransceiverDirectionSendrecv,
})
if err != nil {
_ = pc.Close()
return nil, err
}
// drain RTCP off the outgoing sender, otherwise its queue can stall
// under load. exits when pc.Close makes Read error out.
go drainSenderRTCP(tr.Sender())
return pc, nil
}
func drainSenderRTCP(sender *webrtc.RTPSender) {
if sender == nil {
return
}
buf := make([]byte, 1500)
for {
if _, _, err := sender.Read(buf); err != nil {
return
}
}
}