From 57a579ddd9a54d7530fbf97c7f22dd2738f3d997 Mon Sep 17 00:00:00 2001 From: Bohdan Horbeshko Date: Mon, 17 Aug 2026 11:14:48 -0400 Subject: [PATCH] Improve resilience (downtime history backfilling, better free disk space control) --- config/config.go | 15 +++ config_schema.json | 9 ++ go.mod | 2 +- persistence/sessions.go | 7 ++ telegram/client.go | 20 +++- telegram/connect.go | 33 ++++++ telegram/handlers.go | 192 +++++++++++++++++++------------ telegram/utils.go | 245 ++++++++++++++++++++++++++++++++++++---- xmpp/component.go | 27 +++++ xmpp/gateway/gateway.go | 20 ++++ 10 files changed, 474 insertions(+), 96 deletions(-) diff --git a/config/config.go b/config/config.go index a397757..39d45df 100644 --- a/config/config.go +++ b/config/config.go @@ -49,6 +49,15 @@ type TelegramConfig struct { Verbosity uint8 `yaml:":tdlib_verbosity"` MAMThreshold uint32 `yaml:":mam_threshold"` Tdlib TelegramTdlibConfig `yaml:":tdlib"` + + // HistoryBackfillLimit caps how many messages a single chat's + // downtime-backfill sweep will pull on reconnect (see + // gateway.HistoryBackfillLimit). 0/unset falls back to a default of 1000. + HistoryBackfillLimit uint32 `yaml:":history_backfill_limit"` + + // HistoryBackfillWorkers bounds how many chats backfill concurrently + // (see gateway.HistoryBackfillWorkers). 0/unset falls back to 4. + HistoryBackfillWorkers uint32 `yaml:":history_backfill_workers"` } // TelegramContentConfig is for :content: subtree @@ -58,6 +67,12 @@ type TelegramContentConfig struct { Upload string `yaml:":upload"` User string `yaml:":user"` Quota string `yaml:":quota"` + + // MinFreeSpace is the minimum free disk space (e.g. "64M") required on + // the filesystem hosting TDlib's own files directory before a file + // download is allowed to proceed (see gateway.DownloadSpaceMargin). + // 0/unset falls back to a default of 64 MiB. + MinFreeSpace string `yaml:":min_free_space"` } // TelegramTdlibConfig is for :tdlib: subtree diff --git a/config_schema.json b/config_schema.json index 7e72cfa..d1df878 100644 --- a/config_schema.json +++ b/config_schema.json @@ -27,6 +27,9 @@ }, ":quota": { "type": "string" + }, + ":min_free_space": { + "type": "string" } } }, @@ -36,6 +39,12 @@ ":mam_threshold": { "type": "integer" }, + ":history_backfill_limit": { + "type": "integer" + }, + ":history_backfill_workers": { + "type": "integer" + }, ":tdlib": { "required": [":client"], "type": "object", diff --git a/go.mod b/go.mod index e232b4e..d170429 100644 --- a/go.mod +++ b/go.mod @@ -15,6 +15,7 @@ require ( github.com/xdg-go/stringprep v1.0.4 github.com/zelenin/go-tdlib v0.5.2 golang.org/x/crypto v0.48.0 + golang.org/x/sys v0.41.0 gopkg.in/hraban/opus.v2 v2.0.0-20230925203106-0188a62cb302 gopkg.in/yaml.v2 v2.2.4 gosrc.io/xmpp v0.5.2-0.20211214110136-5f99e1cd06e1 @@ -63,7 +64,6 @@ require ( github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect go.opencensus.io v0.22.5 // indirect golang.org/x/net v0.50.0 // indirect - golang.org/x/sys v0.41.0 // indirect golang.org/x/term v0.40.0 // indirect golang.org/x/text v0.34.0 // indirect golang.org/x/time v0.14.0 // indirect diff --git a/persistence/sessions.go b/persistence/sessions.go index 790b439..4af6a70 100644 --- a/persistence/sessions.go +++ b/persistence/sessions.go @@ -58,6 +58,13 @@ type Session struct { // (the zero value, e.g. for an account that predates this field) // means "auto", matching e2ee.ModeAuto's own default. OMEMO string `yaml:":omemo"` + + // LastOnline is a Unix timestamp, refreshed periodically while the + // TDlib session is online and on graceful shutdown. It's used only to + // size a downtime-backfill window on reconnect (telegram.Client.Connect) + // - not user-configurable, so it's not part of ConfigKeys. Zero means + // "never seen online yet", which skips backfill entirely. + LastOnline int64 `yaml:":lastonline"` } const ( diff --git a/telegram/client.go b/telegram/client.go index 14c8403..5245edc 100644 --- a/telegram/client.go +++ b/telegram/client.go @@ -219,8 +219,14 @@ func buildCallAdapters(tdClient *client.Client, jid string, deps CallDeps) ( type clientLocks struct { authorizationReady sync.Mutex - chatMessageLocks map[int64]*sync.Mutex - resourcesLock sync.Mutex + // chatMessageLocksLock guards chatMessageLocks itself (the map, not its + // values): every existing caller of getChatMessageLock runs serially on + // the single updateHandler/dispatchUpdate goroutine, so this was + // previously safe by construction, but downtime backfill now populates + // it from an independent goroutine too + chatMessageLocksLock sync.Mutex + chatMessageLocks map[int64]*sync.Mutex + resourcesLock sync.Mutex outboxLock sync.Mutex mucCacheLock sync.Mutex editOutboxLock sync.Mutex @@ -230,6 +236,13 @@ type clientLocks struct { loginFinish barrier uploadingFilesLock sync.Mutex + // deliveredMessageIds tracks, per chat, the highest message id already + // handed to ProcessIncomingMessage - guards against double delivery + // between live updateNewMessage processing and downtime backfill + // (message ids are guaranteed strictly increasing per chat in TDlib) + deliveredMessageIds map[int64]int64 + deliveredMessageIdsLock sync.Mutex + authorizerReadLock sync.Mutex authorizerWriteLock sync.Mutex @@ -307,7 +320,8 @@ func NewClient(conf config.TelegramConfig, jid string, component *xmpp.Component avatarHashes: make(map[int64]*HashedAvatar), MessageIdChanges: make(map[int64]map[int64]*newId), locks: clientLocks{ - chatMessageLocks: make(map[int64]*sync.Mutex), + chatMessageLocks: make(map[int64]*sync.Mutex), + deliveredMessageIds: make(map[int64]int64), }, }, nil } diff --git a/telegram/connect.go b/telegram/connect.go index 56b44af..3986680 100644 --- a/telegram/connect.go +++ b/telegram/connect.go @@ -2,6 +2,7 @@ package telegram import ( "github.com/pkg/errors" + "sync" "time" "dev.narayana.im/narayana/telegabber/xmpp/gateway" @@ -156,12 +157,44 @@ func (c *Client) Connect(resource string, wasSessionLoginEmpty bool) error { log.Debug("waiting for loginFinish") c.locks.loginFinish.Wait() + // captured before we touch LastOnline again, to size the downtime + // backfill window below + lastOnline := c.Session.LastOnline + + // Reserve chats needing a downtime backfill *before* live update + // processing starts, so a live message can never race ahead of that + // chat's older backfilled history (see reserveChatsForBackfill) - this + // means fetching the chat list here, synchronously, rather than in the + // background warm-up below. Only for KeepOnline sessions: KeepOnline=false + // already means "I'd rather lose a history gap than get flooded on + // reconnect" (relying on the XMPP server's own offline queue), and the + // same preference applies to a transport-wide downtime gap. lastOnline + // being zero means there's no prior watermark (first-ever connect), so + // there's nothing to catch up on either. + var backfillReserved map[int64]*sync.Mutex + if c.Session.KeepOnline && lastOnline != 0 { + chats, err := c.client.GetChats(&client.GetChatsRequest{ + Limit: chatsLimit, + }) + if err != nil { + log.Errorf("Could not retrieve chats for downtime backfill: %v", err) + } else { + backfillReserved = c.reserveChatsForBackfill(chats.ChatIds) + } + } + go c.updateHandler() log.Warn("Going online") c.online.Store(true) c.locks.authorizationReady.Unlock() c.addResource(resource) + if backfillReserved != nil { + // runs independently so a large backfill doesn't delay the "Logged + // in as" notification/subscribe below + go c.backfillDowntimeHistory(backfillReserved, lastOnline) + } + go func() { chats, err := c.client.GetChats(&client.GetChatsRequest{ Limit: chatsLimit, diff --git a/telegram/handlers.go b/telegram/handlers.go index b971779..0c5f164 100644 --- a/telegram/handlers.go +++ b/telegram/handlers.go @@ -35,6 +35,9 @@ func int64SliceToStringSlice(ints []int64) []string { } func (c *Client) getChatMessageLock(chatID int64) *sync.Mutex { + c.locks.chatMessageLocksLock.Lock() + defer c.locks.chatMessageLocksLock.Unlock() + lock, ok := c.locks.chatMessageLocks[chatID] if !ok { lock = &sync.Mutex{} @@ -44,6 +47,33 @@ func (c *Client) getChatMessageLock(chatID int64) *sync.Mutex { return lock } +// reserveChatsForBackfill locks (creating if needed) the per-chat message +// lock for every given chat id and returns them already held, keyed by +// chat id. Call this before starting live update processing for a +// downtime-backfill sweep: since updateNewMessage acquires the same lock +// before delivering a live message, any live message for a reserved chat +// will block until the backfill sweep delivers that chat's history and +// releases it - guaranteeing backfilled (older) messages are always +// delivered before a live (newer) one for the same chat, rather than +// racing with it. Chats not in chatIds are left untouched and unaffected. +func (c *Client) reserveChatsForBackfill(chatIds []int64) map[int64]*sync.Mutex { + c.locks.chatMessageLocksLock.Lock() + defer c.locks.chatMessageLocksLock.Unlock() + + reserved := make(map[int64]*sync.Mutex, len(chatIds)) + for _, chatId := range chatIds { + lock, ok := c.locks.chatMessageLocks[chatId] + if !ok { + lock = &sync.Mutex{} + c.locks.chatMessageLocks[chatId] = lock + } + lock.Lock() + reserved[chatId] = lock + } + + return reserved +} + func (c *Client) cleanTempFile(path string) { os.Remove(path) @@ -84,80 +114,93 @@ func (c *Client) updateHandler() { defer listener.Close() for update := range listener.Updates { - if update.GetClass() == client.ClassUpdate { - // debug probe: did UpdateCall reach the listener at all? - if t := update.GetType(); t == client.TypeUpdateCall || t == client.TypeUpdateNewCallSignalingData { - log.WithField("type", t).Debug("updateHandler: call-related update received") - } - switch update.GetType() { - case client.TypeUpdateUser: - typedUpdate, _ := update.(*client.UpdateUser) - c.updateUser(typedUpdate) - log.Debugf("%#v", typedUpdate.User) - case client.TypeUpdateUserStatus: - typedUpdate, _ := update.(*client.UpdateUserStatus) - c.updateUserStatus(typedUpdate) - log.Debugf("%#v", typedUpdate.Status) - case client.TypeUpdateNewChat: - typedUpdate, _ := update.(*client.UpdateNewChat) - c.updateNewChat(typedUpdate) - log.Debugf("%#v", typedUpdate.Chat) - case client.TypeUpdateChatPosition: - typedUpdate, _ := update.(*client.UpdateChatPosition) - c.updateChatPosition(typedUpdate) - log.Debugf("%#v", typedUpdate) - case client.TypeUpdateChatLastMessage: - typedUpdate, _ := update.(*client.UpdateChatLastMessage) - c.updateChatLastMessage(typedUpdate) - log.Debugf("%#v", typedUpdate) - case client.TypeUpdateNewMessage: - typedUpdate, _ := update.(*client.UpdateNewMessage) - c.updateNewMessage(typedUpdate) - log.Debugf("%#v", typedUpdate.Message) - case client.TypeUpdateMessageContent: - typedUpdate, _ := update.(*client.UpdateMessageContent) - c.updateMessageContent(typedUpdate) - log.Debugf("%#v", typedUpdate.NewContent) - case client.TypeUpdateDeleteMessages: - typedUpdate, _ := update.(*client.UpdateDeleteMessages) - c.updateDeleteMessages(typedUpdate) - case client.TypeUpdateAuthorizationState: - typedUpdate, _ := update.(*client.UpdateAuthorizationState) - c.updateAuthorizationState(typedUpdate) - case client.TypeUpdateMessageSendSucceeded: - typedUpdate, _ := update.(*client.UpdateMessageSendSucceeded) - c.updateMessageSendSucceeded(typedUpdate) - case client.TypeUpdateMessageSendFailed: - typedUpdate, _ := update.(*client.UpdateMessageSendFailed) - c.updateMessageSendFailed(typedUpdate) - case client.TypeUpdateChatTitle: - typedUpdate, _ := update.(*client.UpdateChatTitle) - c.updateChatTitle(typedUpdate) - case client.TypeUpdateChatReadOutbox: - typedUpdate, _ := update.(*client.UpdateChatReadOutbox) - c.updateChatReadOutbox(typedUpdate) - case client.TypeUpdateBasicGroupFullInfo: - typedUpdate, _ := update.(*client.UpdateBasicGroupFullInfo) - c.updateBasicGroupFullInfo(typedUpdate) - case client.TypeUpdateChatPermissions: - typedUpdate, _ := update.(*client.UpdateChatPermissions) - c.updateChatPermissions(typedUpdate) - case client.TypeUpdateFile: - typedUpdate, _ := update.(*client.UpdateFile) - c.updateFile(typedUpdate) - case client.TypeUpdateCall: - typedUpdate, _ := update.(*client.UpdateCall) - c.callDeps.TgManager.OnUpdateCall(c.jid, typedUpdate) - case client.TypeUpdateNewCallSignalingData: - typedUpdate, _ := update.(*client.UpdateNewCallSignalingData) - c.callDeps.TgManager.OnNewSignalingData(c.jid, typedUpdate) - default: - // log only handled types - continue - } + c.dispatchUpdate(update) + } +} - log.Debugf("%#v", update) +// dispatchUpdate handles a single TDlib update, recovering from any panic so +// that one malformed update (or a downstream failure like a disk-full file +// write) can't take down the whole process for every hosted session. +func (c *Client) dispatchUpdate(update client.Type) { + defer func() { + if r := recover(); r != nil { + log.Errorf("Recovered from panic while handling update %T: %v", update, r) } + }() + + if update.GetClass() == client.ClassUpdate { + // debug probe: did UpdateCall reach the listener at all? + if t := update.GetType(); t == client.TypeUpdateCall || t == client.TypeUpdateNewCallSignalingData { + log.WithField("type", t).Debug("updateHandler: call-related update received") + } + switch update.GetType() { + case client.TypeUpdateUser: + typedUpdate, _ := update.(*client.UpdateUser) + c.updateUser(typedUpdate) + log.Debugf("%#v", typedUpdate.User) + case client.TypeUpdateUserStatus: + typedUpdate, _ := update.(*client.UpdateUserStatus) + c.updateUserStatus(typedUpdate) + log.Debugf("%#v", typedUpdate.Status) + case client.TypeUpdateNewChat: + typedUpdate, _ := update.(*client.UpdateNewChat) + c.updateNewChat(typedUpdate) + log.Debugf("%#v", typedUpdate.Chat) + case client.TypeUpdateChatPosition: + typedUpdate, _ := update.(*client.UpdateChatPosition) + c.updateChatPosition(typedUpdate) + log.Debugf("%#v", typedUpdate) + case client.TypeUpdateChatLastMessage: + typedUpdate, _ := update.(*client.UpdateChatLastMessage) + c.updateChatLastMessage(typedUpdate) + log.Debugf("%#v", typedUpdate) + case client.TypeUpdateNewMessage: + typedUpdate, _ := update.(*client.UpdateNewMessage) + c.updateNewMessage(typedUpdate) + log.Debugf("%#v", typedUpdate.Message) + case client.TypeUpdateMessageContent: + typedUpdate, _ := update.(*client.UpdateMessageContent) + c.updateMessageContent(typedUpdate) + log.Debugf("%#v", typedUpdate.NewContent) + case client.TypeUpdateDeleteMessages: + typedUpdate, _ := update.(*client.UpdateDeleteMessages) + c.updateDeleteMessages(typedUpdate) + case client.TypeUpdateAuthorizationState: + typedUpdate, _ := update.(*client.UpdateAuthorizationState) + c.updateAuthorizationState(typedUpdate) + case client.TypeUpdateMessageSendSucceeded: + typedUpdate, _ := update.(*client.UpdateMessageSendSucceeded) + c.updateMessageSendSucceeded(typedUpdate) + case client.TypeUpdateMessageSendFailed: + typedUpdate, _ := update.(*client.UpdateMessageSendFailed) + c.updateMessageSendFailed(typedUpdate) + case client.TypeUpdateChatTitle: + typedUpdate, _ := update.(*client.UpdateChatTitle) + c.updateChatTitle(typedUpdate) + case client.TypeUpdateChatReadOutbox: + typedUpdate, _ := update.(*client.UpdateChatReadOutbox) + c.updateChatReadOutbox(typedUpdate) + case client.TypeUpdateBasicGroupFullInfo: + typedUpdate, _ := update.(*client.UpdateBasicGroupFullInfo) + c.updateBasicGroupFullInfo(typedUpdate) + case client.TypeUpdateChatPermissions: + typedUpdate, _ := update.(*client.UpdateChatPermissions) + c.updateChatPermissions(typedUpdate) + case client.TypeUpdateFile: + typedUpdate, _ := update.(*client.UpdateFile) + c.updateFile(typedUpdate) + case client.TypeUpdateCall: + typedUpdate, _ := update.(*client.UpdateCall) + c.callDeps.TgManager.OnUpdateCall(c.jid, typedUpdate) + case client.TypeUpdateNewCallSignalingData: + typedUpdate, _ := update.(*client.UpdateNewCallSignalingData) + c.callDeps.TgManager.OnNewSignalingData(c.jid, typedUpdate) + default: + // log only handled types + return + } + + log.Debugf("%#v", update) } } @@ -230,6 +273,11 @@ func (c *Client) updateNewMessage(update *client.UpdateNewMessage) { go func() { lock.Lock() defer lock.Unlock() + defer func() { + if r := recover(); r != nil { + log.Errorf("Recovered from panic while processing message %v in chat %v: %v", update.Message.Id, chatId, r) + } + }() var forceCmd bool if c.LastBotCmdString != "" && update.Message.IsOutgoing { diff --git a/telegram/utils.go b/telegram/utils.go index 957f63b..202255a 100644 --- a/telegram/utils.go +++ b/telegram/utils.go @@ -16,6 +16,8 @@ import ( "sort" "strconv" "strings" + "sync" + "syscall" "time" "unicode/utf8" @@ -27,6 +29,7 @@ import ( log "github.com/sirupsen/logrus" "github.com/soheilhy/args" "github.com/zelenin/go-tdlib/client" + "golang.org/x/sys/unix" ) type VCardInfo struct { @@ -136,6 +139,14 @@ func NewMessageLimitSince(since int64) *MessageLimit { return &limit } +// NewMessageLimitSinceCapped is like NewMessageLimitSince, but overrides the +// default 1000-message safety cap (see getNLastMessages) with maxMessages. +func NewMessageLimitSinceCapped(since int64, maxMessages int32) *MessageLimit { + limit := NewMessageLimitSince(since) + limit.Messages = maxMessages + return limit +} + const AVATAR_SIZE_LIMIT int64 = 128 * 1024 const ( @@ -1313,7 +1324,12 @@ func (c *Client) PermastoreFile(file *client.File, clone bool) (string, string) } size64 := uint64(file.Size) - c.prepareDiskSpace(size64) + if !c.prepareDiskSpace(size64) { + if !clone { + c.forgetTdlibFile(file.Id) + } + return "", "" + } // detect uploading files, there's no remote id for them yet c.locks.uploadingFilesLock.Lock() @@ -1329,6 +1345,8 @@ func (c *Client) PermastoreFile(file *client.File, clone bool) (string, string) dest := c.content.Path + "/" + basename // destination path link = c.content.Link + "/" + basename // download link + var alreadyStored bool + if clone { file, path, err := c.ForceOpenFile(file, 1) if err == nil { @@ -1344,20 +1362,27 @@ func (c *Client) PermastoreFile(file *client.File, clone bool) (string, string) // create destination tempFile, err := os.OpenFile(dest, os.O_CREATE|os.O_EXCL|os.O_WRONLY, mode) if err != nil { - pathErr := err.(*os.PathError) - if pathErr.Err.Error() == "file exists" { + if os.IsExist(err) { log.Warn(err.Error()) return src, link + } + if errors.Is(err, syscall.ENOSPC) { + log.Errorf("Not enough disk space to store %v, dropping it: %v", basename, err) } else { log.Errorf("File creation error: %v", err) - return "", "" } + return "", "" } defer tempFile.Close() // copy _, err = io.Copy(tempFile, file) if err != nil { - log.Errorf("File copying error: %v", err) + if errors.Is(err, syscall.ENOSPC) { + log.Errorf("Not enough disk space to store %v, dropping it: %v", basename, err) + } else { + log.Errorf("File copying error: %v", err) + } + os.Remove(dest) // don't leave a truncated file behind return "", "" } } else if path != "" { @@ -1371,9 +1396,19 @@ func (c *Client) PermastoreFile(file *client.File, clone bool) (string, string) // move err = os.Rename(src, dest) if err != nil { - linkErr := err.(*os.LinkError) - if linkErr.Err.Error() == "file exists" { + // whatever the reason, our local copy is either redundant + // (dest already exists) or one we're giving up on - either + // way tell TDlib it can reclaim the space, rather than + // leaving a file it moved out from under it stuck in its + // own storage forever + defer c.forgetTdlibFile(file.Id) + + if os.IsExist(err) { log.Warn(err.Error()) + alreadyStored = true + } else if errors.Is(err, syscall.ENOSPC) { + log.Errorf("Not enough disk space to move %v, dropping it: %v", basename, err) + return "", "" } else { log.Errorf("File moving error: %v", err) return "", "" @@ -1399,13 +1434,27 @@ func (c *Client) PermastoreFile(file *client.File, clone bool) (string, string) } } - // copy or move should have succeeded at this point - gateway.CachedStorageSize += size64 + // copy or move should have succeeded at this point, unless the + // destination was already there (e.g. a prior partial move) + if !alreadyStored { + gateway.CachedStorageSize += size64 + } } return src, link } +// forgetTdlibFile tells TDlib it no longer needs to keep fileId in its own +// local file cache. Used when PermastoreFile has decided a locally-moved +// file is redundant or unrecoverable, so that a file we've stopped tracking +// via gateway's own quota doesn't linger in TDlib's storage indefinitely. +func (c *Client) forgetTdlibFile(fileId int32) { + _, err := c.client.DeleteFile(&client.DeleteFileRequest{FileId: fileId}) + if err != nil { + log.Warnf("Could not tell TDlib to delete file %v from its own storage: %v", fileId, err) + } +} + func (c *Client) formatBantime(hours int64) int32 { var until int32 if hours > 0 { @@ -1799,7 +1848,12 @@ func (c *Client) ensureDownloadFile(file *client.File) *client.File { defer gateway.StorageLock.Unlock() if file != nil { - c.prepareDiskSpace(uint64(file.Size)) + if !c.prepareDiskSpace(uint64(file.Size)) { + return file + } + if !hasFreeDiskSpace(c.parameters.FilesDirectory, uint64(file.Size)) { + return file + } newFile, err := c.DownloadFile(file.Id, 1, true) if err == nil { @@ -1831,6 +1885,11 @@ func (c *Client) ProcessIncomingMessage(chatId int64, message *client.Message) { return } + if !c.markMessageDelivered(chatId, message.Id) { + log.Debugf("Skipping already-delivered message %v in chat %v", message.Id, chatId) + return + } + chat, _, _ := c.GetContactByID(chatId, nil, true) safeToSend := true groupChatFrom := "" @@ -2458,6 +2517,9 @@ func (c *Client) getNLastMessages(chatID int64, limit *MessageLimit) ([]*client. } case MessageLimitSince: safetyLimit = 1000 + if limit.Messages > 0 { + safetyLimit = limit.Messages + } } safetyLoop: @@ -2511,6 +2573,108 @@ func (c *Client) getNLastMessages(chatID int64, limit *MessageLimit) ([]*client. return messages, nil } +// backfillDowntimeHistory fetches, for every reserved chat, the messages +// that arrived while the whole transport process (not just one XMPP +// resource - see Client.Disconnect's KeepOnline handling for that case) +// was down, and feeds them through the same pipeline live messages use so +// attachments come through normally - unlike !history/!search, which never +// download attachments. sinceUnix is 0 on a session's first-ever connect +// (no LastOnline watermark persisted yet), which is deliberately left +// alone here: there's no prior "downtime" to speak of, so no backfill +// should run. +// +// reserved must be the result of reserveChatsForBackfill, called before +// live update processing started for this connection: each chat's lock is +// already held, and backfillChatHistory releases it once that chat's +// history has been delivered, so a live message for the same chat that +// arrives in the meantime blocks until backfill catches it up - rather +// than racing ahead and being delivered before older backfilled messages. +// Chats are processed concurrently (bounded by gateway.HistoryBackfillWorkers) +// so one huge or FLOOD_WAIT-throttled chat doesn't hold up the others. +func (c *Client) backfillDowntimeHistory(reserved map[int64]*sync.Mutex, sinceUnix int64) { + if sinceUnix == 0 { + // nothing to catch up on; still release every reservation + for _, lock := range reserved { + lock.Unlock() + } + return + } + + log.Warnf("Transport was down for ≈%v, backfilling history for %v chats", time.Since(time.Unix(sinceUnix, 0)).Round(time.Second), len(reserved)) + + workers := int(gateway.HistoryBackfillWorkers) + if workers <= 0 { + workers = 1 + } + + jobs := make(chan int64, len(reserved)) + for chatId := range reserved { + jobs <- chatId + } + close(jobs) + + var wg sync.WaitGroup + for i := 0; i < workers; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for chatId := range jobs { + c.backfillChatHistory(chatId, sinceUnix, reserved[chatId]) + } + }() + } + wg.Wait() +} + +// backfillChatHistory fetches and delivers one chat's history since the +// given watermark, oldest first, then releases lock (already held by the +// caller via reserveChatsForBackfill). If a chat accumulated more than +// gateway.HistoryBackfillLimit messages during the downtime, only the most +// recent ones are delivered - logged and skipped, rather than blocking +// reconnect indefinitely trying to fetch an unbounded gap. +func (c *Client) backfillChatHistory(chatId int64, since int64, lock *sync.Mutex) { + defer lock.Unlock() + defer func() { + if r := recover(); r != nil { + log.Errorf("Recovered from panic while backfilling chat %v: %v", chatId, r) + } + }() + + limit := int32(gateway.HistoryBackfillLimit) + messages, err := c.getNLastMessages(chatId, NewMessageLimitSinceCapped(since, limit)) + if err != nil { + log.Errorf("Failed to backfill history for chat %v, dropping this chat's gap: %v", chatId, err) + return + } + if len(messages) == 0 { + return + } + if limit > 0 && len(messages) >= int(limit) { + log.Warnf("Chat %v had ≥%v messages since downtime; delivering only the most recent ones", chatId, limit) + } + + // oldest first, matching live delivery order + for i := len(messages) - 1; i >= 0; i-- { + c.ProcessIncomingMessage(chatId, messages[i]) + } +} + +// markMessageDelivered records that messageId is about to be delivered for +// chatId, returning false if an equal-or-later message was already +// delivered. This guards against double delivery between live +// updateNewMessage processing and downtime backfill (message ids are +// guaranteed strictly increasing per chat in TDlib). +func (c *Client) markMessageDelivered(chatId, messageId int64) bool { + c.locks.deliveredMessageIdsLock.Lock() + defer c.locks.deliveredMessageIdsLock.Unlock() + + if c.locks.deliveredMessageIds[chatId] >= messageId { + return false + } + c.locks.deliveredMessageIds[chatId] = messageId + return true +} + // GetMessagesBetween lazily fetches message history between given ids (from exclusive, last inclusive), also calculating completeness flag; negative limit means messages from the end func (c *Client) GetMessagesBetween(chatID, fromMessageId, lastMessageId int64, limit int32, reverse bool) (messages []*client.Message, complete bool, err error) { log.WithFields(log.Fields{ @@ -2668,6 +2832,10 @@ func (c *Client) ForceOpenFile(tgFile *client.File, priority int32) (*os.File, s } else // obtain the photo right now if still not downloaded if !tgFile.Local.IsDownloadingCompleted { + if !hasFreeDiskSpace(c.parameters.FilesDirectory, uint64(tgFile.Size)) { + return nil, path, errors.New("Not enough disk space to download file") + } + tdFile, tdErr := c.DownloadFile(tgFile.Id, priority, true) if tdErr == nil { path = tdFile.Local.Path @@ -2818,17 +2986,54 @@ func (c *Client) sendPresence(args ...args.V) error { return gateway.SendPresence(c.xmpp, c.jid, args...) } -func (c *Client) prepareDiskSpace(size uint64) { - if gateway.StorageQuota > 0 && c.content.Path != "" { - var loweredQuota uint64 - if gateway.StorageQuota >= size { - loweredQuota = gateway.StorageQuota - size - } - if gateway.CachedStorageSize >= loweredQuota { - log.Warn("Storage is rapidly clogged") - gateway.CleanOldFiles(c.content.Path, loweredQuota) - } +// hasFreeDiskSpace reports whether the filesystem hosting path currently +// has at least size + gateway.DownloadSpaceMargin bytes free. A statfs +// failure is treated as "can't tell, proceed". +func hasFreeDiskSpace(path string, size uint64) bool { + var stat unix.Statfs_t + if err := unix.Statfs(path, &stat); err != nil { + log.Warnf("Could not check free disk space at %v, proceeding anyway: %v", path, err) + return true } + + free := stat.Bavail * uint64(stat.Bsize) + needed := size + gateway.DownloadSpaceMargin + + if free < needed { + log.Errorf("Not enough free disk space at %v: %v bytes free, need %v (%v file + %v margin)", path, free, needed, size, gateway.DownloadSpaceMargin) + return false + } + + return true +} + +// prepareDiskSpace makes room in c.content.Path for an incoming file of the +// given size by evicting old files down to a lowered quota. Returns false +// if the file should be discarded instead: either it exceeds the whole +// quota by itself (cleaning everything still wouldn't fit it), or real +// free disk space is still short after cleaning. +func (c *Client) prepareDiskSpace(size uint64) bool { + if gateway.StorageQuota == 0 || c.content.Path == "" { + return true + } + + if size > gateway.StorageQuota { + log.Errorf("File is %v bytes, exceeding the entire storage quota (%v bytes); discarding it rather than evicting everything else in a futile attempt to fit it", size, gateway.StorageQuota) + return false + } + + loweredQuota := gateway.StorageQuota - size + if gateway.CachedStorageSize >= loweredQuota { + log.Warn("Storage is rapidly clogged") + gateway.CleanOldFiles(c.content.Path, loweredQuota) + } + + if !hasFreeDiskSpace(c.content.Path, size) { + log.Errorf("Still not enough real disk space at %v after cleaning; discarding the file", c.content.Path) + return false + } + + return true } func (c *Client) GetVcardInfo(toID int64) (VCardInfo, error) { diff --git a/xmpp/component.go b/xmpp/component.go index 7dc3d5d..d5ef556 100644 --- a/xmpp/component.go +++ b/xmpp/component.go @@ -94,8 +94,24 @@ func NewComponent(conf config.XMPPConfig, tc config.TelegramConfig, idsPath stri } } + if tc.Content.MinFreeSpace != "" { + gateway.DownloadSpaceMargin, err = parseSize(tc.Content.MinFreeSpace) + if err != nil { + log.Warnf("Error parsing min_free_space: %v; using the default", err) + gateway.DownloadSpaceMargin = 64 * 1024 * 1024 + } + } + gateway.MAMThreshold = tc.MAMThreshold + if tc.HistoryBackfillLimit > 0 { + gateway.HistoryBackfillLimit = tc.HistoryBackfillLimit + } + + if tc.HistoryBackfillWorkers > 0 { + gateway.HistoryBackfillWorkers = tc.HistoryBackfillWorkers + } + options := xmpp.ComponentOptions{ TransportConfiguration: xmpp.TransportConfiguration{ Address: conf.Host + ":" + conf.Port, @@ -255,6 +271,11 @@ func heartbeat(component *xmpp.Component) { sessionLock.Lock() for _, session := range sessions { + if session.Online() { + session.Session.LastOnline = time.Now().Unix() + gateway.DirtySessions = true + } + session.DelayedStatusesLock.Lock() for chatID, delayedStatus := range session.DelayedStatuses { if delayedStatus.TimestampExpired <= now { @@ -392,6 +413,12 @@ func Close(component *xmpp.Component) { sessionLock.Lock() // close all sessions for _, session := range sessions { + if session.Online() { + // mark the exact moment we went offline, so a graceful + // shutdown doesn't cost the downtime-backfill window the + // up-to-60s staleness of the heartbeat-driven update above + session.Session.LastOnline = time.Now().Unix() + } session.Disconnect("", true) } sessionLock.Unlock() diff --git a/xmpp/gateway/gateway.go b/xmpp/gateway/gateway.go index f2a1e74..d5ab72f 100644 --- a/xmpp/gateway/gateway.go +++ b/xmpp/gateway/gateway.go @@ -98,6 +98,26 @@ 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 + +// 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)