2019-10-22 16:36:54 +00:00
|
|
|
package xmpp
|
|
|
|
|
|
|
|
import (
|
2019-11-12 15:50:25 +00:00
|
|
|
"github.com/pkg/errors"
|
2022-01-05 21:04:22 +00:00
|
|
|
"sync"
|
2019-11-16 18:44:13 +00:00
|
|
|
"time"
|
2019-11-12 15:50:25 +00:00
|
|
|
|
2019-10-22 16:36:54 +00:00
|
|
|
"dev.narayana.im/narayana/telegabber/config"
|
2019-11-10 23:50:50 +00:00
|
|
|
"dev.narayana.im/narayana/telegabber/persistence"
|
2019-11-07 21:09:53 +00:00
|
|
|
"dev.narayana.im/narayana/telegabber/telegram"
|
2019-11-24 17:10:29 +00:00
|
|
|
"dev.narayana.im/narayana/telegabber/xmpp/gateway"
|
2019-10-22 16:36:54 +00:00
|
|
|
|
2019-11-12 15:50:25 +00:00
|
|
|
log "github.com/sirupsen/logrus"
|
2019-10-22 16:36:54 +00:00
|
|
|
"gosrc.io/xmpp"
|
2021-12-18 16:04:24 +00:00
|
|
|
"gosrc.io/xmpp/stanza"
|
2019-10-22 16:36:54 +00:00
|
|
|
)
|
|
|
|
|
2019-11-03 22:15:43 +00:00
|
|
|
var tgConf config.TelegramConfig
|
2019-11-24 17:10:29 +00:00
|
|
|
var sessions map[string]*telegram.Client
|
2019-12-04 15:55:15 +00:00
|
|
|
var db *persistence.SessionsYamlDB
|
2022-01-05 21:04:22 +00:00
|
|
|
var sessionLock sync.Mutex
|
2019-11-03 22:15:43 +00:00
|
|
|
|
2019-10-29 01:23:57 +00:00
|
|
|
// NewComponent starts a new component and wraps it in
|
|
|
|
// a stream manager that you should start yourself
|
2019-11-19 20:25:14 +00:00
|
|
|
func NewComponent(conf config.XMPPConfig, tc config.TelegramConfig) (*xmpp.StreamManager, *xmpp.Component, error) {
|
2019-11-03 22:15:43 +00:00
|
|
|
var err error
|
2019-11-10 23:50:50 +00:00
|
|
|
|
2021-12-18 16:04:24 +00:00
|
|
|
gateway.Jid, err = stanza.NewJid(conf.Jid)
|
2019-11-03 22:15:43 +00:00
|
|
|
if err != nil {
|
2019-11-19 20:25:14 +00:00
|
|
|
return nil, nil, err
|
2019-11-03 22:15:43 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
tgConf = tc
|
|
|
|
|
2019-10-22 16:36:54 +00:00
|
|
|
options := xmpp.ComponentOptions{
|
2021-12-18 16:04:24 +00:00
|
|
|
TransportConfiguration: xmpp.TransportConfiguration{
|
|
|
|
Address: conf.Host + ":" + conf.Port,
|
|
|
|
Domain: conf.Jid,
|
|
|
|
},
|
2019-10-22 16:36:54 +00:00
|
|
|
Domain: conf.Jid,
|
|
|
|
Secret: conf.Password,
|
|
|
|
Name: "telegabber",
|
|
|
|
}
|
|
|
|
|
|
|
|
router := xmpp.NewRouter()
|
2019-11-03 22:15:43 +00:00
|
|
|
router.HandleFunc("iq", HandleIq)
|
|
|
|
router.HandleFunc("presence", HandlePresence)
|
2019-10-22 16:36:54 +00:00
|
|
|
router.HandleFunc("message", HandleMessage)
|
|
|
|
|
2021-12-18 16:04:24 +00:00
|
|
|
component, err := xmpp.NewComponent(options, router, func(err error) {
|
|
|
|
log.Error(err)
|
|
|
|
})
|
2019-10-22 16:36:54 +00:00
|
|
|
if err != nil {
|
2019-11-19 20:25:14 +00:00
|
|
|
return nil, nil, err
|
2019-10-22 16:36:54 +00:00
|
|
|
}
|
|
|
|
|
2019-12-22 01:04:45 +00:00
|
|
|
// probe all known sessions
|
2019-11-24 17:10:29 +00:00
|
|
|
err = loadSessions(conf.Db, component)
|
|
|
|
if err != nil {
|
|
|
|
return nil, nil, err
|
|
|
|
}
|
|
|
|
|
2019-12-04 15:55:15 +00:00
|
|
|
sm := xmpp.NewStreamManager(component, func(s xmpp.Sender) {
|
|
|
|
go heartbeat(component)
|
|
|
|
})
|
2019-11-12 15:50:25 +00:00
|
|
|
|
2019-11-19 20:25:14 +00:00
|
|
|
return sm, component, nil
|
2019-10-22 16:36:54 +00:00
|
|
|
}
|
2019-11-12 15:50:25 +00:00
|
|
|
|
2019-11-16 18:44:13 +00:00
|
|
|
func heartbeat(component *xmpp.Component) {
|
|
|
|
var err error
|
2019-11-24 22:20:07 +00:00
|
|
|
probeType := gateway.SPType("probe")
|
2019-11-16 18:44:13 +00:00
|
|
|
|
2022-01-05 21:04:22 +00:00
|
|
|
sessionLock.Lock()
|
2019-11-14 20:11:04 +00:00
|
|
|
for jid := range sessions {
|
2019-12-04 15:55:15 +00:00
|
|
|
err = gateway.SendPresence(component, jid, probeType)
|
|
|
|
if err != nil {
|
|
|
|
log.Error(err)
|
2019-11-18 19:01:45 +00:00
|
|
|
}
|
2019-11-14 20:11:04 +00:00
|
|
|
}
|
2022-01-05 21:04:22 +00:00
|
|
|
sessionLock.Unlock()
|
2019-11-16 18:44:13 +00:00
|
|
|
|
2019-11-18 19:01:45 +00:00
|
|
|
log.Info("Starting heartbeat queue")
|
|
|
|
|
2019-12-22 01:04:45 +00:00
|
|
|
// status updater thread
|
2019-11-16 18:44:13 +00:00
|
|
|
for {
|
2019-12-22 01:04:45 +00:00
|
|
|
time.Sleep(60e9)
|
2019-11-29 00:51:41 +00:00
|
|
|
for key, presence := range gateway.Queue {
|
2020-01-10 13:02:25 +00:00
|
|
|
err = gateway.ResumableSend(component, presence)
|
2019-11-16 18:44:13 +00:00
|
|
|
if err != nil {
|
2020-01-10 13:02:25 +00:00
|
|
|
gateway.LogBadPresence(presence)
|
2019-11-16 18:44:13 +00:00
|
|
|
} else {
|
2019-11-29 00:51:41 +00:00
|
|
|
delete(gateway.Queue, key)
|
2019-11-16 18:44:13 +00:00
|
|
|
}
|
|
|
|
}
|
2022-01-05 21:04:22 +00:00
|
|
|
|
|
|
|
if gateway.DirtySessions {
|
|
|
|
gateway.DirtySessions = false
|
|
|
|
// no problem if a dirty flag gets set again here,
|
|
|
|
// it would be resolved on the next iteration
|
|
|
|
SaveSessions()
|
|
|
|
}
|
2019-11-16 18:44:13 +00:00
|
|
|
}
|
2019-11-12 15:50:25 +00:00
|
|
|
}
|
|
|
|
|
2019-11-24 17:10:29 +00:00
|
|
|
func loadSessions(dbPath string, component *xmpp.Component) error {
|
2019-11-12 15:50:25 +00:00
|
|
|
var err error
|
|
|
|
|
2019-11-24 17:10:29 +00:00
|
|
|
sessions = make(map[string]*telegram.Client)
|
2019-11-12 15:50:25 +00:00
|
|
|
|
|
|
|
db, err = persistence.LoadSessions(dbPath)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
db.Transaction(func() bool {
|
|
|
|
for jid, session := range db.Data.Sessions {
|
2022-01-05 21:04:22 +00:00
|
|
|
// copy the session struct, otherwise all of them would reference
|
|
|
|
// the same temporary range variable
|
|
|
|
currentSession := session
|
|
|
|
getTelegramInstance(jid, ¤tSession, component)
|
2019-11-12 15:50:25 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
return false
|
|
|
|
}, persistence.SessionMarshaller)
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2019-11-24 17:10:29 +00:00
|
|
|
func getTelegramInstance(jid string, savedSession *persistence.Session, component *xmpp.Component) (*telegram.Client, bool) {
|
|
|
|
var err error
|
2019-11-12 15:50:25 +00:00
|
|
|
session, ok := sessions[jid]
|
|
|
|
if !ok {
|
2019-11-24 17:10:29 +00:00
|
|
|
session, err = telegram.NewClient(tgConf, jid, component, savedSession)
|
2019-11-12 15:50:25 +00:00
|
|
|
if err != nil {
|
|
|
|
log.Error(errors.Wrap(err, "TDlib initialization failure"))
|
|
|
|
return session, false
|
|
|
|
}
|
2022-01-05 21:04:22 +00:00
|
|
|
if savedSession.KeepOnline {
|
|
|
|
if err = session.Connect(""); err != nil {
|
|
|
|
log.Error(err)
|
|
|
|
return session, false
|
|
|
|
}
|
|
|
|
}
|
|
|
|
sessionLock.Lock()
|
2019-11-12 15:50:25 +00:00
|
|
|
sessions[jid] = session
|
2022-01-05 21:04:22 +00:00
|
|
|
sessionLock.Unlock()
|
2019-11-12 15:50:25 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
return session, true
|
|
|
|
}
|
2019-11-14 20:11:04 +00:00
|
|
|
|
2022-01-05 21:04:22 +00:00
|
|
|
// SaveSessions dumps current sessions to the file
|
|
|
|
func SaveSessions() {
|
|
|
|
sessionLock.Lock()
|
|
|
|
defer sessionLock.Unlock()
|
|
|
|
db.Transaction(func() bool {
|
|
|
|
for jid, session := range sessions {
|
|
|
|
db.Data.Sessions[jid] = *session.Session
|
|
|
|
}
|
|
|
|
|
|
|
|
return true
|
|
|
|
}, persistence.SessionMarshaller)
|
|
|
|
}
|
|
|
|
|
2019-11-19 20:25:14 +00:00
|
|
|
// Close gracefully terminates the component and saves active sessions
|
|
|
|
func Close(component *xmpp.Component) {
|
|
|
|
log.Error("Disconnecting...")
|
|
|
|
|
2022-01-05 21:04:22 +00:00
|
|
|
sessionLock.Lock()
|
2019-12-22 01:04:45 +00:00
|
|
|
// close all sessions
|
2019-11-19 20:25:14 +00:00
|
|
|
for _, session := range sessions {
|
2022-01-03 03:54:13 +00:00
|
|
|
session.Disconnect("", true)
|
2019-11-19 20:25:14 +00:00
|
|
|
}
|
2022-01-05 21:04:22 +00:00
|
|
|
sessionLock.Unlock()
|
2019-11-19 20:25:14 +00:00
|
|
|
|
2019-12-22 01:04:45 +00:00
|
|
|
// save sessions
|
2022-01-05 21:04:22 +00:00
|
|
|
SaveSessions()
|
2019-11-19 20:25:14 +00:00
|
|
|
|
2019-12-22 01:04:45 +00:00
|
|
|
// close stream
|
2019-11-19 20:25:14 +00:00
|
|
|
component.Disconnect()
|
|
|
|
}
|