mirror of
https://github.com/tiennm99/tiennm99bot.git
synced 2026-10-11 03:13:46 +00:00
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.
165 lines
4.2 KiB
Go
165 lines
4.2 KiB
Go
package storage
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"sort"
|
|
"strings"
|
|
"sync"
|
|
)
|
|
|
|
// MemoryProvider is a Provider backed by per-module in-process maps. Intended
|
|
// for tests and local no-database runs (MODULES= / no MONGO_URL); state is lost
|
|
// when the process exits. Production uses MongoProvider.
|
|
type MemoryProvider struct {
|
|
mu sync.Mutex
|
|
cols map[string]*memoryCollection
|
|
}
|
|
|
|
// NewMemoryProvider returns a fresh in-process provider.
|
|
func NewMemoryProvider() *MemoryProvider {
|
|
return &MemoryProvider{cols: make(map[string]*memoryCollection)}
|
|
}
|
|
|
|
// Collection returns the per-module in-memory collection, creating it on first
|
|
// use. An invalid module name yields an invalidCollection so Typed produces a
|
|
// store whose every op errors — mirroring MongoProvider's defense in depth.
|
|
func (p *MemoryProvider) Collection(module string) Collection {
|
|
if !collectionNameRe.MatchString(module) {
|
|
return invalidCollection{name: module}
|
|
}
|
|
p.mu.Lock()
|
|
defer p.mu.Unlock()
|
|
c, ok := p.cols[module]
|
|
if !ok {
|
|
c = &memoryCollection{rows: make(map[string]memoryRow)}
|
|
p.cols[module] = c
|
|
}
|
|
return c
|
|
}
|
|
|
|
// memoryCollection holds one module's rows. Multiple typed views (different
|
|
// payload types under disjoint key prefixes) share the same map.
|
|
type memoryCollection struct {
|
|
mu sync.RWMutex
|
|
rows map[string]memoryRow
|
|
}
|
|
|
|
func (*memoryCollection) isCollection() {}
|
|
|
|
// memoryRow stores the JSON-encoded value plus its version. Encoding to bytes
|
|
// isolates stored state from caller mutation, matching production where Mongo
|
|
// serializes the value on write.
|
|
type memoryRow struct {
|
|
version int64
|
|
data []byte
|
|
}
|
|
|
|
// memoryDocStore is a typed DocStore over one memoryCollection.
|
|
type memoryDocStore[T any] struct {
|
|
c *memoryCollection
|
|
}
|
|
|
|
func (s *memoryDocStore[T]) Get(_ context.Context, id string) (T, int64, error) {
|
|
var out T
|
|
if err := validateKey(id); err != nil {
|
|
return out, 0, err
|
|
}
|
|
s.c.mu.RLock()
|
|
defer s.c.mu.RUnlock()
|
|
row, ok := s.c.rows[id]
|
|
if !ok {
|
|
return out, 0, ErrNotFound
|
|
}
|
|
if err := json.Unmarshal(row.data, &out); err != nil {
|
|
return out, 0, err
|
|
}
|
|
return out, row.version, nil
|
|
}
|
|
|
|
func (s *memoryDocStore[T]) Put(_ context.Context, id string, val T) error {
|
|
if err := validateKey(id); err != nil {
|
|
return err
|
|
}
|
|
data, err := json.Marshal(val)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
s.c.mu.Lock()
|
|
defer s.c.mu.Unlock()
|
|
// A plain Put bumps the version too, so a concurrent versioned writer that
|
|
// read the old version correctly sees a conflict.
|
|
s.c.rows[id] = memoryRow{version: s.c.rows[id].version + 1, data: data}
|
|
return nil
|
|
}
|
|
|
|
func (s *memoryDocStore[T]) PutVersioned(_ context.Context, id string, expectedVersion int64, val T) error {
|
|
if err := validateKey(id); err != nil {
|
|
return err
|
|
}
|
|
data, err := json.Marshal(val)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
s.c.mu.Lock()
|
|
defer s.c.mu.Unlock()
|
|
row, ok := s.c.rows[id]
|
|
if expectedVersion == 0 {
|
|
if ok {
|
|
return ErrConflict
|
|
}
|
|
} else if !ok || row.version != expectedVersion {
|
|
return ErrConflict
|
|
}
|
|
s.c.rows[id] = memoryRow{version: row.version + 1, data: data}
|
|
return nil
|
|
}
|
|
|
|
func (s *memoryDocStore[T]) Delete(_ context.Context, id string) error {
|
|
if err := validateKey(id); err != nil {
|
|
return err
|
|
}
|
|
s.c.mu.Lock()
|
|
defer s.c.mu.Unlock()
|
|
delete(s.c.rows, id)
|
|
return nil
|
|
}
|
|
|
|
func (s *memoryDocStore[T]) List(_ context.Context, prefix string) ([]string, error) {
|
|
if err := validatePrefix(prefix); err != nil {
|
|
return nil, err
|
|
}
|
|
s.c.mu.RLock()
|
|
defer s.c.mu.RUnlock()
|
|
keys := make([]string, 0)
|
|
for k := range s.c.rows {
|
|
if strings.HasPrefix(k, prefix) {
|
|
keys = append(keys, k)
|
|
}
|
|
}
|
|
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
|
|
}
|