- 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.
230 lines
8.9 KiB
Go
230 lines
8.9 KiB
Go
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
|
||
}
|