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.
This commit is contained in:
@@ -21,36 +21,90 @@ type Router struct {
|
||||
// Hook, вызываемый на каждое входящее событие (обычно → process_turn).
|
||||
onUserMsg func(Incoming)
|
||||
|
||||
// Асинхронная обработка входящих: handleIncoming кладёт событие в канал,
|
||||
// воркер-горутина последовательно вызывает onUserMsg. Благодаря этому
|
||||
// long-poll цикл канала (Telegram) не блокируется на время долгого
|
||||
// вызова аналитика и продолжает принимать новые сообщения.
|
||||
incoming chan 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{},
|
||||
sessions: map[UserID]any{},
|
||||
routes: map[UserID]Route{},
|
||||
pending: map[UserID]PendingQ{},
|
||||
onUserMsg: onUserMsg,
|
||||
incoming: make(chan Incoming, 256),
|
||||
jobs: make(chan *userState, chatWorkers),
|
||||
users: make(map[UserID]*userState),
|
||||
}
|
||||
for i := 0; i < chatWorkers; i++ {
|
||||
go r.workerLoop()
|
||||
}
|
||||
go r.processLoop()
|
||||
return r
|
||||
}
|
||||
|
||||
// processLoop — воркер асинхронной обработки входящих (FIFO).
|
||||
func (r *Router) processLoop() {
|
||||
for inc := range r.incoming {
|
||||
// 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)
|
||||
}
|
||||
@@ -99,9 +153,19 @@ func (r *Router) handleIncoming(inc Incoming) {
|
||||
}
|
||||
r.mu.Unlock()
|
||||
|
||||
// Асинхронная обработка: кладём событие в очередь воркера и сразу
|
||||
// возвращаемся, не блокируя вызывающий long-poll цикл канала.
|
||||
r.incoming <- inc
|
||||
// Асинхронная обработка: кладём событие в 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,
|
||||
|
||||
Reference in New Issue
Block a user