143 lines
5.2 KiB
Go
143 lines
5.2 KiB
Go
package chat
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
"sync"
|
||
)
|
||
|
||
// 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
|
||
}
|
||
|
||
// 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)
|
||
}
|
||
}
|
||
|
||
// 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
|
||
}
|