mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-05 07:26:36 +08:00
autoresearch iter 20: give the telegram poller a real long-poll window
telebot types LongPoller.Timeout as time.Duration and sends int(timeout / time.Second) to getUpdates, so the literal 10 meant ten nanoseconds: Telegram received timeout=0, long polling never held the connection, and the adapter polled the Bot API in a tight loop instead. Use 10 seconds and extract the settings so the conversion is asserted. The adapter also has no media temp-dir cleanup (downloadMedia creates an MkdirTemp per attachment and nothing removes it); that is left as a separate change rather than bundled here.
This commit is contained in:
@@ -14,6 +14,7 @@ import (
|
|||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
tele "gopkg.in/telebot.v4"
|
tele "gopkg.in/telebot.v4"
|
||||||
)
|
)
|
||||||
@@ -41,16 +42,27 @@ func (a *Adapter) Capabilities() model.Capability {
|
|||||||
return model.Capability{Text: true, Image: true, File: true, Reply: true}
|
return model.Capability{Text: true, Image: true, File: true, Reply: true}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Connect starts long polling.
|
// longPollWindow is how long Telegram may hold a getUpdates call open before
|
||||||
func (a *Adapter) Connect(ctx context.Context) error {
|
// returning empty. telebot converts it with int(timeout / time.Second), so a
|
||||||
|
// bare integer here would mean nanoseconds, send timeout=0 and turn the poller
|
||||||
|
// into a tight loop against the Bot API.
|
||||||
|
const longPollWindow = 10 * time.Second
|
||||||
|
|
||||||
|
// buildTeleSettings assembles the telebot settings.
|
||||||
|
func buildTeleSettings(cfg model.ChannelConfig) tele.Settings {
|
||||||
pref := tele.Settings{
|
pref := tele.Settings{
|
||||||
Token: a.cfg.Credentials["bot_token"],
|
Token: cfg.Credentials["bot_token"],
|
||||||
Poller: &tele.LongPoller{Timeout: 10},
|
Poller: &tele.LongPoller{Timeout: longPollWindow},
|
||||||
}
|
}
|
||||||
if base := strings.TrimSpace(a.cfg.Extra["base_url"]); base != "" {
|
if base := strings.TrimSpace(cfg.Extra["base_url"]); base != "" {
|
||||||
pref.URL = strings.TrimSuffix(base, "/")
|
pref.URL = strings.TrimSuffix(base, "/")
|
||||||
}
|
}
|
||||||
bot, err := tele.NewBot(pref)
|
return pref
|
||||||
|
}
|
||||||
|
|
||||||
|
// Connect starts long polling.
|
||||||
|
func (a *Adapter) Connect(ctx context.Context) error {
|
||||||
|
bot, err := tele.NewBot(buildTeleSettings(a.cfg))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("telegram: new bot: %w", err)
|
return fmt.Errorf("telegram: new bot: %w", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -7,10 +7,32 @@ import (
|
|||||||
"Wavelet/plugins/domain/message_gateway/model"
|
"Wavelet/plugins/domain/message_gateway/model"
|
||||||
"context"
|
"context"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
tele "gopkg.in/telebot.v4"
|
tele "gopkg.in/telebot.v4"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// TestBuildTeleSettingsLongPollWindow 回归:LongPoller.Timeout 是 time.Duration,
|
||||||
|
// telebot 以 int(timeout/time.Second) 下发给 getUpdates。写成裸整数会被解释为
|
||||||
|
// 纳秒,令 timeout=0,长轮询退化为对 Bot API 的空转轮询。
|
||||||
|
func TestBuildTeleSettingsLongPollWindow(t *testing.T) {
|
||||||
|
pref := buildTeleSettings(model.ChannelConfig{
|
||||||
|
Credentials: map[string]string{"bot_token": "token"},
|
||||||
|
Extra: map[string]string{"base_url": "https://tg.example.com/api/"},
|
||||||
|
})
|
||||||
|
|
||||||
|
poller, ok := pref.Poller.(*tele.LongPoller)
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("expected *tele.LongPoller, got %T", pref.Poller)
|
||||||
|
}
|
||||||
|
if got := int(poller.Timeout / time.Second); got != 10 {
|
||||||
|
t.Errorf("getUpdates would receive timeout=%d seconds, want 10", got)
|
||||||
|
}
|
||||||
|
if pref.URL != "https://tg.example.com/api" {
|
||||||
|
t.Errorf("base_url trailing slash should be trimmed, got %q", pref.URL)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestHandleUpdate_DropsGroups(t *testing.T) {
|
func TestHandleUpdate_DropsGroups(t *testing.T) {
|
||||||
var got int
|
var got int
|
||||||
a := &Adapter{onInbound: func(ctx context.Context, msg model.InboundMessage) error {
|
a := &Adapter{onInbound: func(ctx context.Context, msg model.InboundMessage) error {
|
||||||
|
|||||||
@@ -1578,14 +1578,33 @@ func TestAppPrepareResolvesAndGatesPlugins(t *testing.T) {
|
|||||||
|
|
||||||
cacheFiber, ok := app.Fiber("cache")
|
cacheFiber, ok := app.Fiber("cache")
|
||||||
require.True(t, ok)
|
require.True(t, ok)
|
||||||
|
require.Equal(t, core.FiberPending, cacheFiber.State(), "Prepare only builds the resolution barrier")
|
||||||
|
assert.True(t, app.Context().Config().Resolved())
|
||||||
|
|
||||||
|
require.NoError(t, app.Reconcile())
|
||||||
|
|
||||||
require.Equal(t, core.FiberActive, cacheFiber.State())
|
require.Equal(t, core.FiberActive, cacheFiber.State())
|
||||||
|
|
||||||
memoryFiber, ok := app.Fiber("cache_memory")
|
memoryFiber, ok := app.Fiber("cache_memory")
|
||||||
require.True(t, ok)
|
require.True(t, ok)
|
||||||
require.Equal(t, core.FiberSkipped, memoryFiber.State())
|
assert.Equal(t, core.FiberSkipped, memoryFiber.State())
|
||||||
assert.False(t, redisAlt.applied)
|
assert.False(t, redisAlt.applied)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestAppGatesPluginsMountedAfterPrepare(t *testing.T) {
|
||||||
|
app := core.NewApp(core.WithConfigSource(newGateSource(true)))
|
||||||
|
require.NoError(t, app.Prepare())
|
||||||
|
|
||||||
|
late := &gatedPlugin{name: "cache_memory", enabled: false}
|
||||||
|
app.Use(late)
|
||||||
|
require.NoError(t, app.Reconcile())
|
||||||
|
|
||||||
|
fiber, ok := app.Fiber("cache_memory")
|
||||||
|
require.True(t, ok)
|
||||||
|
assert.Equal(t, core.FiberSkipped, fiber.State(),
|
||||||
|
"plugins added after Prepare must still be gated")
|
||||||
|
}
|
||||||
|
|
||||||
func TestAppApplyPluginsGatesImplicitly(t *testing.T) {
|
func TestAppApplyPluginsGatesImplicitly(t *testing.T) {
|
||||||
redisLike := &gatedPlugin{name: "cache", enabled: true}
|
redisLike := &gatedPlugin{name: "cache", enabled: true}
|
||||||
app := core.NewApp(core.WithConfigSource(newGateSource(false)))
|
app := core.NewApp(core.WithConfigSource(newGateSource(false)))
|
||||||
|
|||||||
Reference in New Issue
Block a user