mirror of
https://dev.narayana.im/narayana/telegabber.git
synced 2026-10-07 16:51:47 +00:00
Improve resilience (downtime history backfilling, better free disk space control)
This commit is contained in:
parent
b315eeff1e
commit
57a579ddd9
10 changed files with 474 additions and 96 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
2
go.mod
2
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
|
||||
|
|
|
|||
|
|
@ -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 (
|
||||
|
|
|
|||
|
|
@ -219,6 +219,12 @@ func buildCallAdapters(tdClient *client.Client, jid string, deps CallDeps) (
|
|||
|
||||
type clientLocks struct {
|
||||
authorizationReady 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
|
||||
|
|
@ -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
|
||||
|
||||
|
|
@ -308,6 +321,7 @@ func NewClient(conf config.TelegramConfig, jid string, component *xmpp.Component
|
|||
MessageIdChanges: make(map[int64]map[int64]*newId),
|
||||
locks: clientLocks{
|
||||
chatMessageLocks: make(map[int64]*sync.Mutex),
|
||||
deliveredMessageIds: make(map[int64]int64),
|
||||
},
|
||||
}, nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,6 +114,20 @@ func (c *Client) updateHandler() {
|
|||
defer listener.Close()
|
||||
|
||||
for update := range listener.Updates {
|
||||
c.dispatchUpdate(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 {
|
||||
|
|
@ -153,13 +197,12 @@ func (c *Client) updateHandler() {
|
|||
c.callDeps.TgManager.OnNewSignalingData(c.jid, typedUpdate)
|
||||
default:
|
||||
// log only handled types
|
||||
continue
|
||||
return
|
||||
}
|
||||
|
||||
log.Debugf("%#v", update)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// new user discovered
|
||||
func (c *Client) updateUser(update *client.UpdateUser) {
|
||||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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 "<ERROR>", ""
|
||||
}
|
||||
|
||||
// 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 "<ERROR>", ""
|
||||
}
|
||||
return "<ERROR>", ""
|
||||
}
|
||||
defer tempFile.Close()
|
||||
// copy
|
||||
_, err = io.Copy(tempFile, file)
|
||||
if err != nil {
|
||||
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 "<ERROR>", ""
|
||||
}
|
||||
} 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 "<ERROR>", ""
|
||||
} else {
|
||||
log.Errorf("File moving error: %v", err)
|
||||
return "<ERROR>", ""
|
||||
|
|
@ -1399,13 +1434,27 @@ func (c *Client) PermastoreFile(file *client.File, clone bool) (string, string)
|
|||
}
|
||||
}
|
||||
|
||||
// copy or move should have succeeded at this point
|
||||
// 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
|
||||
// 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) {
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Reference in a new issue