Files
ratatoskr-go/internal/chat/router.go
ki.sagidullin b5fb583c90
All checks were successful
CI / test (pull_request) Successful in 41s
CI / build-and-package (amd64, linux) (pull_request) Successful in 37s
CI / build-and-package (amd64, windows) (pull_request) Successful in 34s
test: починить тесты на Windows
- E2E (app): e2eFakeAPI переведён на v2 HTTP API opencode (/api/*) с
  определением агента по тексту промпта; Router получает processed-счётчик
  и WaitProcessed, e2eChannel.deliver ждёт асинхронную обработку — убирает
  гонку «запрос сразу после deliver» и коллатеральный 'database is closed'.
- app_test: одинарные YAML-кавычки для путей Windows (backslash-escape) +
  закрытие Store в TestNew/TestNew_RunCtxCancel/TestNew_UpdateWiring.
- config_test: абсолютный путь строится с корнем тома (C:\...) и одинарными
  кавычками YAML.
- opencode/server_test: fakeServeBin на Windows — .cmd с ping (#!/bin/sh
  не исполняется).
- worker_test: TestWorkerSemaphore поллит до целевого статуса вместо
  фиксированных sleep (git на Windows медленнее).
2026-08-19 10:47:55 +05:00

166 lines
6.1 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)
// Асинхронная обработка входящих: handleIncoming кладёт событие в канал,
// воркер-горутина последовательно вызывает onUserMsg. Благодаря этому
// long-poll цикл канала (Telegram) не блокируется на время долгого
// вызова аналитика и продолжает принимать новые сообщения.
incoming chan Incoming
// processed — число обработанных воркером событий (для синхронизации
// тестов с асинхронной очередью: WaitProcessed ждёт обработку события).
processed atomic.Int64
}
// 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,
incoming: make(chan Incoming, 256),
}
go r.processLoop()
return r
}
// processLoop — воркер асинхронной обработки входящих (FIFO).
func (r *Router) processLoop() {
for inc := range r.incoming {
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()
// Асинхронная обработка: кладём событие в очередь воркера и сразу
// возвращаемся, не блокируя вызывающий long-poll цикл канала.
r.incoming <- inc
}
// 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
}