Files
ratatoskr-go/internal/chat/router.go
ki.sagidullin cd0619926e
Some checks failed
CI / test (push) Failing after 1m15s
CI / build-and-package (amd64, linux) (push) Failing after 58s
CI / build-and-package (amd64, windows) (push) Successful in 30s
perf(chat,update): пул воркеров per-user вместо сериальной очереди + HEAD-проба обновлений
- chat.Router: ограниченный пул chatWorkers=4 воркеров + FIFO-очереди
  per-user (userState/workerLoop/runUser). Порядок сообщений одного UserID
  сохраняется; разные пользователи обрабатываются параллельно (до 4
  одновременных LLM-вызовов), long-poll Telegram не блокируется чужим
  аналитиком. Backpressure по jobs — только на перегруженного пользователя.
- app.FreeChat: sessions под sync.Mutex (защита от data race при параллельных
  воркерах роутера).
- update: ResolveLatest проверяет наличие бинаря HEAD-пробой без скачивания
  тела (fallback GET Range 0-0 при 405/501), сортировка версий по id убыв.;
  один общий http.Client (keep-alive) вместо нового на каждый запрос.
- тесты: порядок/параллелизм per-user в router, HEAD-без-тела и фоллбэк на
  версию без бинаря в update.
- память Serena: инварианты Router/update, примечания по форматированию на Windows.
2026-08-22 11:44:34 +05:00

230 lines
8.9 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package chat
import (
"context"
"fmt"
"sync"
"sync/atomic"
"time"
)
// Router — единый диспетчер входящих из всех каналов и маршрутизатор исходящих.
// Владеет сессиями (по UserID), routing-таблицей (UserID→Address+Channel)
// и pending-вопросами (один незакрытый на сессию).
type Router struct {
mu sync.Mutex
channels []Channel
sessions map[UserID]any // заглушка: реальные сессии подключаются через Hook
routes map[UserID]Route
pending map[UserID]PendingQ
// Hook, вызываемый на каждое входящее событие (обычно → process_turn).
onUserMsg func(Incoming)
// Асинхронная обработка входящих ограниченным пулом воркеров с
// упорядоченными очередями per-user (см. chatWorkers, userState).
// Благодаря этому long-poll цикл канала (Telegram) не блокируется на время
// долгого вызова аналитика, а сообщения разных пользователей не сериализуются
// друг за другом: каждый активный пользователь занимает своего воркера.
jobs chan *userState
users map[UserID]*userState
userMu sync.Mutex
// processed — число обработанных воркером событий (для синхронизации
// тестов с асинхронной очередью: WaitProcessed ждёт обработку события).
processed atomic.Int64
}
// userState — FIFO-очередь входящих одного пользователя. В каждый момент
// для пользователя активен ровно один воркер (scheduled), поэтому порядок
// обработки его сообщений сохраняется, а параллелизм достигается между
// разными пользователями.
type userState struct {
mu sync.Mutex
pending []Incoming
scheduled bool
}
// chatWorkers — число воркеров обработки входящих. Ограничивает количество
// одновременных тяжёлых LLM-вызовов (аналитик/свободный чат), чтобы поток
// каналов не упирался в один долгий вызов.
const chatWorkers = 4
// NewRouter создаёт роутер. onUserMsg — колбэк обработки входящего.
func NewRouter(onUserMsg func(Incoming)) *Router {
if onUserMsg == nil {
onUserMsg = func(Incoming) {}
}
r := &Router{
sessions: map[UserID]any{},
routes: map[UserID]Route{},
pending: map[UserID]PendingQ{},
onUserMsg: onUserMsg,
jobs: make(chan *userState, chatWorkers),
users: make(map[UserID]*userState),
}
for i := 0; i < chatWorkers; i++ {
go r.workerLoop()
}
return r
}
// userStateOf возвращает очередь пользователя (создаёт при первом сообщении).
// Очереди живут вечно — по одной маленькой структуре на пользователя/вкладку.
func (r *Router) userStateOf(uid UserID) *userState {
r.userMu.Lock()
defer r.userMu.Unlock()
st, ok := r.users[uid]
if !ok {
st = &userState{}
r.users[uid] = st
}
return st
}
// workerLoop — воркер пула: берёт пользователя из общей очереди и дренит его.
func (r *Router) workerLoop() {
for st := range r.jobs {
r.runUser(st)
}
}
// runUser обрабатывает все накопленные сообщения пользователя по порядку.
// По исчерпании очереди снимает scheduled — следующий handleIncoming вновь
// поставит пользователя в jobs. Возвращается в workerLoop, чтобы тот взял
// следующего пользователя из общей очереди.
func (r *Router) runUser(st *userState) {
for {
st.mu.Lock()
if len(st.pending) == 0 {
st.scheduled = false
st.mu.Unlock()
return
}
inc := st.pending[0]
st.pending = st.pending[1:]
st.mu.Unlock()
r.onUserMsg(inc)
r.processed.Add(1)
}
}
// Processed возвращает число обработанных воркером входящих событий.
func (r *Router) Processed() int64 { return r.processed.Load() }
// WaitProcessed ждёт, пока воркер обработает не меньше target событий
// (для синхронизации с асинхронной очередью в тестах).
func (r *Router) WaitProcessed(target int64) bool {
deadline := time.Now().Add(30 * time.Second)
for r.processed.Load() < target {
if time.Now().After(deadline) {
return false
}
time.Sleep(2 * time.Millisecond)
}
return true
}
// Attach регистрирует канал и подключает его к обработчику входящих.
// Возвращает ошибку только при пустом канале (nil).
func (r *Router) Attach(ch Channel) error {
if ch == nil {
return fmt.Errorf("chat: attach nil channel")
}
r.mu.Lock()
defer r.mu.Unlock()
ch.OnMessage(r.handleIncoming)
r.channels = append(r.channels, ch)
return nil
}
// handleIncoming — единый вход входящих из любого канала: обновляет маршрут
// пользователя. Если у пользователя открыт pending-вопрос и входящее пришло
// с адреса вопроса — это ответ, pending закрывается (блок 4.1: «с того же
// Address = ответ»). Затем событие уходит в обработчик.
func (r *Router) handleIncoming(inc Incoming) {
r.mu.Lock()
r.routes[inc.UserID] = Route{Address: inc.Address, Channel: inc.Channel}
if p, ok := r.pending[inc.UserID]; ok && p.AnswerAddr == inc.Address {
// ответ на вопрос: консумируем pending
inc.Msg.QuestionID = p.Msg.QuestionID
delete(r.pending, inc.UserID)
}
r.mu.Unlock()
// Асинхронная обработка: кладём событие в FIFO-очередь пользователя и
// сразу возвращаемся, не блокируя вызывающий long-poll цикл канала.
// Если пользователь ещё не обрабатывается — ставим его в общую очередь
// пула воркеров. Backpressure по jobs блокирует только перегруженного
// пользователя (его собственную горутину канала), не весь роутер.
st := r.userStateOf(inc.UserID)
st.mu.Lock()
st.pending = append(st.pending, inc)
if !st.scheduled {
st.scheduled = true
r.jobs <- st
}
st.mu.Unlock()
}
// Send уведомляет пользователя через текущий маршрут. M1 (нет маршрута) — no-op,
// возвращаем nil (M1-подобный деградированный путь). M2 при сбое канала.
func (r *Router) Send(ctx context.Context, uid UserID, m Message) error {
route, ok := r.routeOf(uid)
if !ok {
// M1: маршрута нет — тихо пропускаем.
return nil
}
if err := route.Channel.Send(ctx, route.Address, m); err != nil {
return fmt.Errorf("%w: %v", ErrSendFailed, err)
}
return nil
}
// Ask открывает pending-вопрос (не блокирует). M3 если уже есть открытый.
func (r *Router) Ask(ctx context.Context, uid UserID, prompt Message) error {
route, ok := r.routeOf(uid)
if !ok {
return ErrRouteNotFound
}
r.mu.Lock()
if _, exists := r.pending[uid]; exists {
r.mu.Unlock()
return ErrWaitingAnswer
}
prompt.QuestionID = newQuestionID()
r.pending[uid] = PendingQ{UserID: uid, AnswerAddr: route.Address, Msg: prompt}
r.mu.Unlock()
if err := route.Channel.Send(ctx, route.Address, prompt); err != nil {
// откат pending, если вопрос не ушёл
r.mu.Lock()
delete(r.pending, uid)
r.mu.Unlock()
return fmt.Errorf("%w: %v", ErrSendFailed, err)
}
return nil
}
// ClosePending закрывает открытый вопрос (например, диалог сам отменил его).
func (r *Router) ClosePending(uid UserID) {
r.mu.Lock()
delete(r.pending, uid)
r.mu.Unlock()
}
// Pending возвращает незакрытый вопрос пользователя (блок 4).
func (r *Router) Pending(uid UserID) (PendingQ, bool) {
r.mu.Lock()
defer r.mu.Unlock()
p, ok := r.pending[uid]
return p, ok
}
func (r *Router) routeOf(uid UserID) (Route, bool) {
r.mu.Lock()
defer r.mu.Unlock()
rt, ok := r.routes[uid]
return rt, ok
}