fix(alias): answer inline queries from one store read under a deadline

Telegram expires an inline query and then rejects the answer with "query
is too old and response timeout expired or query ID is invalid". The
picker invited that: it listed the names and then read the store once per
name — up to 50 round trips per keystroke — and it was the only handler
in the module with no deadline of its own. Updates are dispatched one at
a time, so a single slow answer also held up the queries queued behind
it, each ageing while it waited, and one slow read expired a whole burst
of typing.

Add DocStore.Scan, which reads a key prefix with its values in one round
trip, ordered by key. The picker and /aliases both use it, so neither
grows a round trip per saved alias. Bound the inline handler at 3s: an
answer later than that is rejected anyway, and giving up frees the worker
for the fresher query behind it. When Telegram does reject an answer, the
error now carries how long it took, which separates a slow handler from a
query that was already stale on arrival.

The 50-result cap now counts results the picker can show, so a video-note
alias — which has no cached inline type — no longer consumes a slot.
This commit is contained in:
tiennm99 committed 2026-09-08 16:02:39 +07:00
1 parent cb86ec4cf3
commit 8260d2860b
9 files changed
+278 -63

No files matched your search

+31 -8
View File
@@ -66,7 +66,30 @@ is the payoff for storing a `file_id` rather than bytes.
**Video-note aliases do not appear inline.** Telegram defines no
`InlineQueryResultCachedVideoNote`, and substituting a plain video would change
what was saved. They stay reachable through `/insert` and `/<name>`.
what was saved. They stay reachable through `/insert` and `/<name>`. The 50-cap
counts results the picker can show, so a skipped kind does not eat a slot.
**Answering is a race, and losing it is silent.** Telegram expires an inline
query and then rejects the answer with
> Bad Request: query is too old and response timeout expired or query ID is
> invalid
Every keystroke opens a *new* query, and the bot dispatches updates one at a
time, so one slow answer also delays the queries queued behind it — each ageing
while it waits. One slow read can therefore expire a whole burst of typing.
Two rules keep that from happening:
- The handler reads the store **once** per query (`Scan`), never a name listing
followed by a read per name. An N+1 read costs a round trip per saved alias,
on every keystroke.
- It runs under a 3-second deadline, not the 10 seconds the commands get. An
answer that late is rejected anyway; abandoning it frees the worker for the
fresher query behind it.
When the rejection does appear, the log line carries how long the answer took —
`answer 4 results after 12.4s: ...` — which separates a slow handler from a
query that was already stale on arrival.
**Two things gate inline mode, and both fail silently.**
@@ -146,11 +169,10 @@ One line per alias, showing what the name holds, with the invocation in a
Names list in their folded (lowercase) form, which is exactly what `/insert`
takes.
This costs one store read per *listed* alias — `DocStore` has no bulk get and
the kind lives in the document. The reads stop as soon as the message is full,
so the cost is bounded by what fits in one reply rather than by how many
aliases exist. The list is trimmed to Telegram's 4096-character limit and ends
with `…and N more.`; the count at the top is always the true total.
This is one store read total — `DocStore.Scan` returns the names with their
documents, and the kind that labels each line lives in the document. The list is
trimmed to Telegram's 4096-character limit and ends with `…and N more.`; the
count at the top is always the true total.
## Behaviour worth knowing
@@ -197,5 +219,6 @@ It reports **shape, never content**: field names, lengths and counts, but no
message text. The line lands in stdout and whatever ships it, so aliased
messages must not travel with it; a test asserts nothing leaks.
Both handlers run under a 10-second deadline. The bot processes updates one at a
time, so that bound is what keeps a slow store from stalling other users.
Both command handlers run under a 10-second deadline; the inline handler under
3. The bot processes updates one at a time, so those bounds are what keep a slow
store from stalling other users.
+41 -28
View File
@@ -2,8 +2,9 @@ package alias
import (
"context"
"sort"
"fmt"
"strings"
"time"
"github.com/go-telegram/bot"
"github.com/go-telegram/bot/models"
@@ -19,6 +20,17 @@ const (
// query. Kept short because the namespace is shared and writable: a name
// saved now should show up in the picker within seconds, not minutes.
inlineCacheSeconds = 5
// inlineTimeout bounds one answer, far tighter than the 10s handlerTimeout
// the command handlers use.
//
// An inline query has a shelf life: Telegram invalidates the id and rejects
// the answer as "query is too old". Every keystroke opens a new query, and
// updates are dispatched one at a time, so a single slow answer also holds
// up the queries queued behind it — which are themselves ageing while they
// wait. Giving up early loses one result set; running long loses the whole
// burst.
inlineTimeout = 3 * time.Second
)
// handleInline answers "@botname <prefix>" with the matching aliases.
@@ -32,51 +44,52 @@ func (s *state) handleInline(ctx context.Context, b *bot.Bot, update *models.Upd
return nil
}
names, err := s.store.List(ctx, "")
ctx, cancel := context.WithTimeout(ctx, inlineTimeout)
defer cancel()
started := time.Now()
// One round trip for names *and* values. Reading the names and then each
// alias would cost a round trip per saved name on every keystroke, which is
// exactly the latency an expiring query cannot absorb.
docs, err := s.store.Scan(ctx, "")
if err != nil {
log.Error("alias_inline_list", "err", err)
log.Error("alias_inline_scan", "err", err)
// Answer with nothing rather than leaving the client spinning. An empty
// answer is also what Telegram expects when a query has no matches.
return s.answer(ctx, b, query.ID, nil)
return s.answer(ctx, b, query.ID, nil, started)
}
// Scan returns key order, so the picker is sorted with no sort here.
prefix := strings.ToLower(strings.TrimSpace(query.Query))
matches := make([]string, 0, len(names))
for _, name := range names {
if prefix == "" || strings.HasPrefix(name, prefix) {
matches = append(matches, name)
}
}
sort.Strings(matches)
if len(matches) > maxInlineResults {
matches = matches[:maxInlineResults]
}
results := make([]models.InlineQueryResult, 0, len(matches))
for _, name := range matches {
entry, found, err := s.get(ctx, name)
if err != nil {
// One unreadable record must not blank the whole picker.
log.Error("alias_inline_get", "name", name, "err", err)
results := make([]models.InlineQueryResult, 0, min(len(docs), maxInlineResults))
for _, doc := range docs {
if prefix != "" && !strings.HasPrefix(doc.ID, prefix) {
continue
}
if !found {
continue // deleted between the List and this read
}
if r := inlineResult(name, entry); r != nil {
if r := inlineResult(doc.ID, doc.Val); r != nil {
results = append(results, r)
}
if len(results) == maxInlineResults {
break // Telegram's cap; the rest stay reachable by name
}
}
return s.answer(ctx, b, query.ID, results)
return s.answer(ctx, b, query.ID, results, started)
}
func (s *state) answer(ctx context.Context, b *bot.Bot, queryID string, results []models.InlineQueryResult) error {
func (s *state) answer(ctx context.Context, b *bot.Bot, queryID string, results []models.InlineQueryResult, started time.Time) error {
_, err := b.AnswerInlineQuery(ctx, &bot.AnswerInlineQueryParams{
InlineQueryID: queryID,
Results: results,
CacheTime: inlineCacheSeconds,
})
return err
if err != nil {
// The elapsed time is the whole diagnosis when Telegram rejects the
// answer as stale: it separates "this handler was slow" from "the query
// was already old when it reached us".
return fmt.Errorf("answer %d results after %s: %w",
len(results), time.Since(started).Round(time.Millisecond), err)
}
return nil
}
// inlineResult maps one alias to the inline result type that carries it, or nil
@@ -181,3 +181,25 @@ func TestInline_CapsAtFiftyResults(t *testing.T) {
func uniqueName(i int) string {
return "n" + strings.Repeat("x", i/10) + string(rune('a'+i%10))
}
// The cap counts results the picker can actually show. A kind with no cached
// inline type must not consume a slot, or a handful of video notes would
// shrink an otherwise full answer.
func TestInline_VideoNotesDoNotConsumeCapSlots(t *testing.T) {
rb := installAlias(t)
for i := 0; i < 10; i++ {
rb.Bot.ProcessUpdate(context.Background(),
aliasCmd(uniqueName(i), &models.Message{VideoNote: &models.VideoNote{FileID: "note-id"}}))
}
for i := 10; i < 70; i++ {
rb.Bot.ProcessUpdate(context.Background(),
aliasCmd(uniqueName(i), &models.Message{Text: "x"}))
}
rb.Reset()
rb.Bot.ProcessUpdate(context.Background(), inlineQuery(7, ""))
if got := len(inlineResults(t, rb)); got != 50 {
t.Errorf("returned %d results, want the full 50 despite the skipped video notes", got)
}
}
+13 -27
View File
@@ -6,7 +6,6 @@ import (
"fmt"
"html"
"regexp"
"sort"
"strings"
"time"
@@ -263,52 +262,39 @@ func (s *state) handleAliases(ctx context.Context, b *bot.Bot, update *models.Up
return nil
}
names, err := s.store.List(ctx, "")
// Scan, not List plus a read per name: the kind lives in the document, so a
// name-only listing would cost a round trip per alias to label each line.
// Scan also returns key order, which is what makes the same command twice
// produce the same list.
docs, err := s.store.Scan(ctx, "")
if err != nil {
log.Error("alias_list", "err", err)
return chathelper.Reply(ctx, b, msg, genericFailure)
}
if len(names) == 0 {
if len(docs) == 0 {
return chathelper.Reply(ctx, b, msg,
"No aliases saved yet. Reply to a message with /alias <name> to save one.")
}
// List gives no ordering guarantee, and an unstable list is unreadable when
// it is the same command run twice.
sort.Strings(names)
return chathelper.ReplyHTML(ctx, b, msg, s.renderNames(ctx, names))
return chathelper.ReplyHTML(ctx, b, msg, renderNames(docs))
}
// renderNames formats the list as Telegram HTML, trimmed to fit one message.
//
// One line per alias, each showing what the name holds, with the invocation
// wrapped in <code> so tapping it copies a command ready to send.
//
// This costs one store read per *listed* alias: DocStore has no bulk get and
// the kind lives in the document. The reads stop as soon as the message is
// full, so the cost is bounded by what fits in one reply rather than by how
// many aliases exist.
func (s *state) renderNames(ctx context.Context, names []string) string {
func renderNames(docs []storage.Doc[Alias]) string {
var sb strings.Builder
fmt.Fprintf(&sb, "%d aliases:", len(names))
for i, name := range names {
entry, found, err := s.get(ctx, name)
if err != nil {
// One unreadable record must not blank the whole list.
log.Error("alias_list_get", "name", name, "err", err)
} else if !found {
continue // deleted between the List above and this read
}
fmt.Fprintf(&sb, "%d aliases:", len(docs))
for i, doc := range docs {
// Escaped despite parseName restricting names to [a-zA-Z0-9_]: the
// validation and the rendering are far apart, and a later relaxation of
// the name rules must not silently become an HTML injection.
line := "\n<code>/" + html.EscapeString(name) + "</code> — " + kindLabel(entry.Kind)
line := "\n<code>/" + html.EscapeString(doc.ID) + "</code> — " + kindLabel(doc.Val.Kind)
// Reserve room for the "…and N more" tail before committing to a line,
// so the trim can never be what pushes the message over the limit.
if sb.Len()+len(line) > maxListBytes {
fmt.Fprintf(&sb, "\n…and %d more.", len(names)-i)
fmt.Fprintf(&sb, "\n…and %d more.", len(docs)-i)
return sb.String()
}
sb.WriteString(line)
@@ -320,7 +306,7 @@ func (s *state) renderNames(ctx context.Context, names []string) string {
//
// Differs from describe in exactly one case: text reads as "message" in a
// sentence ("send that message") but as "text" in a column of kinds, next to
// sticker and video. An empty kind means the read above failed.
// sticker and video. An empty kind means a stored document without one.
func kindLabel(kind string) string {
switch kind {
case kindText:
+23
View File
@@ -35,12 +35,32 @@ var ErrInvalidModuleName = errors.New("storage: invalid module name")
// - PutVersioned writes only if the stored version equals expectedVersion,
// then bumps it. expectedVersion == 0 means "create (or adopt a not-yet-
// versioned key)". A mismatch returns ErrConflict.
// - List returns the keys under a prefix; Scan returns those keys with their
// values, ordered by key ascending.
type DocStore[T any] interface {
Get(ctx context.Context, id string) (val T, version int64, err error)
Put(ctx context.Context, id string, val T) error
PutVersioned(ctx context.Context, id string, expectedVersion int64, val T) error
Delete(ctx context.Context, id string) error
List(ctx context.Context, prefix string) ([]string, error)
// Scan reads every document under prefix in one round trip, ordered by key
// ascending. An empty prefix reads the whole collection.
//
// It exists so a caller that needs the values — not just the names — never
// has to follow List with a Get per key. That N+1 shape costs one network
// round trip per document, and handlers run inline on the bot's single
// update worker: a read that scales with the document count turns ordinary
// store latency into requests that expire before they are answered.
Scan(ctx context.Context, prefix string) ([]Doc[T], error)
}
// Doc is one key paired with its value, as returned by Scan. The version is
// deliberately absent: a caller that intends to write back should re-read the
// key with Get so the version it locks on is the one it just observed.
type Doc[T any] struct {
ID string
Val T
}
// Provider yields a per-module Collection handle. Implementations decide how
@@ -154,3 +174,6 @@ func (s invalidDocStore[T]) Delete(context.Context, string) error { return s.err
func (s invalidDocStore[T]) List(context.Context, string) ([]string, error) {
return nil, s.err("List")
}
func (s invalidDocStore[T]) Scan(context.Context, string) ([]Doc[T], error) {
return nil, s.err("Scan")
}
+56
View File
@@ -129,3 +129,59 @@ func TestCheckReservedFields(t *testing.T) {
t.Fatalf("string payload flagged: %v", err)
}
}
// Scan is the batched read: keys and values together, in key order, so a
// caller never follows List with a Get per key.
func TestMemoryDocStore_ScanReturnsValuesInKeyOrder(t *testing.T) {
ctx := context.Background()
s := memStore("alias")
_ = s.Put(ctx, "cheese", testPayload{Name: "second", Count: 2})
_ = s.Put(ctx, "boo", testPayload{Name: "first", Count: 1})
_ = s.Put(ctx, "other", testPayload{Name: "third", Count: 3})
docs, err := s.Scan(ctx, "")
if err != nil {
t.Fatalf("Scan: %v", err)
}
want := []Doc[testPayload]{
{ID: "boo", Val: testPayload{Name: "first", Count: 1}},
{ID: "cheese", Val: testPayload{Name: "second", Count: 2}},
{ID: "other", Val: testPayload{Name: "third", Count: 3}},
}
if len(docs) != len(want) {
t.Fatalf("Scan returned %d docs, want %d: %+v", len(docs), len(want), docs)
}
for i := range want {
if docs[i] != want[i] {
t.Errorf("docs[%d] = %+v, want %+v", i, docs[i], want[i])
}
}
}
func TestMemoryDocStore_ScanPrefix(t *testing.T) {
ctx := context.Background()
s := memStore("coin")
_ = s.Put(ctx, "game:1", testPayload{Name: "a"})
_ = s.Put(ctx, "game:2", testPayload{Name: "b"})
_ = s.Put(ctx, "stats:1", testPayload{Name: "c"})
docs, err := s.Scan(ctx, "game:")
if err != nil {
t.Fatalf("Scan: %v", err)
}
if len(docs) != 2 || docs[0].ID != "game:1" || docs[1].ID != "game:2" {
t.Fatalf("Scan game: = %+v", docs)
}
if docs[0].Val.Name != "a" || docs[1].Val.Name != "b" {
t.Errorf("Scan lost payloads: %+v", docs)
}
}
// An empty collection yields no docs and no error — the picker path answers
// with nothing rather than treating it as a failure.
func TestMemoryDocStore_ScanEmpty(t *testing.T) {
docs, err := memStore("alias").Scan(context.Background(), "")
if err != nil || len(docs) != 0 {
t.Fatalf("Scan on empty store = %+v, err %v", docs, err)
}
}
+22
View File
@@ -140,3 +140,25 @@ func (s *memoryDocStore[T]) List(_ context.Context, prefix string) ([]string, er
sort.Strings(keys)
return keys, nil
}
func (s *memoryDocStore[T]) Scan(_ context.Context, prefix string) ([]Doc[T], error) {
if err := validatePrefix(prefix); err != nil {
return nil, err
}
s.c.mu.RLock()
defer s.c.mu.RUnlock()
docs := make([]Doc[T], 0, len(s.c.rows))
for k, row := range s.c.rows {
if !strings.HasPrefix(k, prefix) {
continue
}
var val T
if err := json.Unmarshal(row.data, &val); err != nil {
return nil, err
}
docs = append(docs, Doc[T]{ID: k, Val: val})
}
// Sorted for parity with the Mongo store, which sorts on _id in the query.
sort.Slice(docs, func(i, j int) bool { return docs[i].ID < docs[j].ID })
return docs, nil
}
+33
View File
@@ -183,3 +183,36 @@ func (s *mongoDocStore[T]) List(ctx context.Context, prefix string) ([]string, e
}
return keys, nil
}
// Scan reads every document under prefix in one Find, sorted by _id so the
// caller inherits key order from the index instead of sorting afterwards.
//
// Unlike List it takes no projection: the payload is the point. The prefix uses
// the same half-open range as List, so the read stays an index scan.
func (s *mongoDocStore[T]) Scan(ctx context.Context, prefix string) ([]Doc[T], error) {
if err := validatePrefix(prefix); err != nil {
return nil, err
}
filter := bson.M{}
if prefix != "" {
filter[mongoIDField] = bson.M{"$gte": prefix, "$lt": prefixSuccessor(prefix)}
}
cur, err := s.coll.Find(ctx, filter, options.Find().SetSort(bson.D{{Key: mongoIDField, Value: 1}}))
if err != nil {
return nil, fmt.Errorf("mongo scan %s prefix=%q: %w", s.module, prefix, err)
}
defer func() { _ = cur.Close(ctx) }()
docs := make([]Doc[T], 0)
for cur.Next(ctx) {
var out storedDoc[T]
if err := cur.Decode(&out); err != nil {
return nil, fmt.Errorf("mongo scan %s prefix=%q: decode: %w", s.module, prefix, err)
}
docs = append(docs, Doc[T]{ID: out.ID, Val: out.Payload})
}
if err := cur.Err(); err != nil {
return nil, fmt.Errorf("mongo scan %s prefix=%q: %w", s.module, prefix, err)
}
return docs, nil
}
+37
View File
@@ -198,3 +198,40 @@ func TestMongoDocStore_WrappedScalarAndArray(t *testing.T) {
t.Errorf("array root field subscribers = %T, want bson.A", doc["subscribers"])
}
}
// Scan reads a whole prefix in one query, sorted by _id, with the payload
// fields decoded from the document root.
func TestMongoDocStore_ScanReturnsPayloadsInKeyOrder(t *testing.T) {
store, _, cleanup := mongoStore[portfolioLike](t, "coin")
defer cleanup()
ctx := context.Background()
if err := store.Put(ctx, "user:2", portfolioLike{USD: 2}); err != nil {
t.Fatalf("Put user:2: %v", err)
}
if err := store.Put(ctx, "user:1", portfolioLike{USD: 1}); err != nil {
t.Fatalf("Put user:1: %v", err)
}
if err := store.Put(ctx, "other", portfolioLike{USD: 9}); err != nil {
t.Fatalf("Put other: %v", err)
}
docs, err := store.Scan(ctx, "user:")
if err != nil {
t.Fatalf("Scan: %v", err)
}
if len(docs) != 2 {
t.Fatalf("Scan user: returned %d docs, want 2: %+v", len(docs), docs)
}
if docs[0].ID != "user:1" || docs[1].ID != "user:2" {
t.Errorf("Scan out of key order: %s, %s", docs[0].ID, docs[1].ID)
}
if docs[0].Val.USD != 1 || docs[1].Val.USD != 2 {
t.Errorf("Scan lost hoisted payload fields: %+v", docs)
}
all, err := store.Scan(ctx, "")
if err != nil || len(all) != 3 {
t.Fatalf("Scan all = %d docs, err %v", len(all), err)
}
}