package telegram import ( "github.com/pkg/errors" "path/filepath" "strconv" "sync" "time" "dev.narayana.im/narayana/telegabber/config" "dev.narayana.im/narayana/telegabber/persistence" "dev.narayana.im/narayana/telegabber/telegram/cache" "dev.narayana.im/narayana/telegabber/xmpp/gateway" "github.com/zelenin/go-tdlib/client" "gosrc.io/xmpp" ) // DelayedStatus describes an online status expiring on timeout type DelayedStatus struct { TimestampOnline int64 TimestampExpired int64 } // IntPair holds two int64 values type IntPair struct { ChatId int64 MessageId int64 } type barrier struct { mu sync.Mutex open bool released bool releaseCh chan struct{} } // Wait blocks until the barrier is released. func (b *barrier) Wait() { b.mu.Lock() if b.released { b.mu.Unlock() return } if !b.open { b.open = true b.releaseCh = make(chan struct{}) // Reinitialize the channel } b.mu.Unlock() // Wait for the barrier to be released <-b.releaseCh } // Done releases the barrier. func (b *barrier) Done() { b.mu.Lock() defer b.mu.Unlock() if b.open { close(b.releaseCh) // Close the channel to release waiting goroutines b.open = false // Mark the barrier as closed } else { b.released = true } } // IsPending checks if the barrier is currently being waited func (b *barrier) IsPending() bool { b.mu.Lock() defer b.mu.Unlock() return b.open } // NewId stores message ids and timestamps of their additions so old ones can be truncated to save memory type newId struct { Id int64 Ts int64 lock sync.Mutex ownLock sync.Mutex locked bool fired bool } func newNewId() *newId { return &newId{Ts: time.Now().Unix()} } func (i *newId) Lock() { i.ownLock.Lock() if i.fired { i.ownLock.Unlock() return } i.locked = true i.ownLock.Unlock() i.lock.Lock() } func (i *newId) Unlock() { i.ownLock.Lock() if i.locked { i.lock.Unlock() i.locked = false i.fired = true } i.ownLock.Unlock() } // Client stores the metadata for lazily invoked TDlib instance type Client struct { client *client.Client authorizer *clientAuthorizer parameters *client.SetTdlibParametersRequest options []client.Option me *client.User xmpp *xmpp.Component jid string Session *persistence.Session resources map[string]bool content *config.TelegramContentConfig cache *cache.Cache online bool loginWizard *loginWizardMetadata loginStage LoginStage lastAuthorizationStateType string outbox map[string]string editOutbox map[string]string pinOutbox map[IntPair]chan int64 DelayedStatuses map[int64]*DelayedStatus DelayedStatusesLock sync.Mutex lastMsgHashes map[int64]uint64 lastMsgIds map[int64]string mucCache map[int64]*MUCState uploadingFiles map[int32]string LastBotCmdString string XmppClientFeatures map[string]*[]string XmppClientFeaturesLock sync.Mutex avatarHashes map[int64]*gateway.HashedAvatar avatarHashesLock sync.Mutex MessageIdChanges map[int64]map[int64]*newId MessageIdChangesLock sync.Mutex locks clientLocks SendMessageLock sync.Mutex } type clientLocks struct { authorizationReady sync.Mutex chatMessageLocks map[int64]*sync.Mutex resourcesLock sync.Mutex outboxLock sync.Mutex mucCacheLock sync.Mutex editOutboxLock sync.Mutex pinOutboxLock sync.Mutex lastMsgHashesLock sync.Mutex lastMsgIdsLock sync.RWMutex loginFinish barrier uploadingFilesLock sync.Mutex authorizerReadLock sync.Mutex authorizerWriteLock sync.Mutex loginWizardReadLock sync.Mutex loginWizardWriteLock sync.Mutex } type loginWizardMetadata struct { nextStage chan LoginStage chanBusy bool commandSent bool } // NewClient instantiates a Telegram App func NewClient(conf config.TelegramConfig, jid string, component *xmpp.Component, session *persistence.Session) (*Client, error) { var options []client.Option if conf.Tdlib.Client.CatchTimeout != 0 { options = append(options, client.WithCatchTimeout( time.Duration(conf.Tdlib.Client.CatchTimeout)*time.Second, )) } apiID, err := strconv.ParseInt(conf.Tdlib.Client.APIID, 10, 32) if err != nil { return &Client{}, errors.Wrap(err, "Wrong api_id") } datadir := conf.Tdlib.Datadir if datadir == "" { datadir = "./sessions/" // ye olde defaute } parameters := client.SetTdlibParametersRequest{ UseTestDc: false, DatabaseDirectory: filepath.Join(datadir, jid), FilesDirectory: filepath.Join(datadir, jid, "/files/"), UseFileDatabase: true, UseChatInfoDatabase: conf.Tdlib.Client.UseChatInfoDatabase, UseMessageDatabase: true, UseSecretChats: conf.Tdlib.Client.UseSecretChats, ApiId: int32(apiID), ApiHash: conf.Tdlib.Client.APIHash, SystemLanguageCode: "en", DeviceModel: conf.Tdlib.Client.DeviceModel, SystemVersion: "1.0.0", ApplicationVersion: conf.Tdlib.Client.ApplicationVersion, EnableStorageOptimizer: true, IgnoreFileNames: false, } return &Client{ parameters: ¶meters, xmpp: component, jid: jid, Session: session, resources: make(map[string]bool), content: &conf.Content, cache: cache.NewCache(), outbox: make(map[string]string), editOutbox: make(map[string]string), pinOutbox: make(map[IntPair]chan int64), mucCache: make(map[int64]*MUCState), uploadingFiles: make(map[int32]string), options: options, DelayedStatuses: make(map[int64]*DelayedStatus), lastMsgHashes: make(map[int64]uint64), lastMsgIds: make(map[int64]string), XmppClientFeatures: make(map[string]*[]string), avatarHashes: make(map[int64]*gateway.HashedAvatar), MessageIdChanges: make(map[int64]map[int64]*newId), locks: clientLocks{ chatMessageLocks: make(map[int64]*sync.Mutex), }, }, nil } // GetPersistenceSession retrieves the internal session configuration func (c *Client) GetPersistenceSession() *persistence.Session { return c.Session }