Files
tiennm99bot/internal/storage/memory_provider.go
T
tiennm99 8260d2860b 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.
2026-09-08 16:02:39 +07:00

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
}