Router теперь обрабатывает входящие в воркер-горутине (FIFO-очередь с буфером 256) вместо синхронного вызова onUserMsg из long-poll цикла канала. Долгий вызов аналитика (Decide) больше не блокирует приём новых сообщений от Telegram: цикл getUpdates продолжает работать. - router.go: NewRouter запускает processLoop; handleIncoming кладёт событие в канал и возвращается; маршрутизация + pending по-прежнему обновляются синхронно под мьютексом. - router_test.go: fakeOnMsg стал потокобезопасным с ожиданием числа входящих (wait), т.к. обработка теперь асинхронная. Преимущества: интерфейс не замирает на время анализа; порядок входящих сохраняется (FIFO). Ограничение: воркер один — при очень долгом аналитике следующие сообщения ждут в очереди, но канал их продолжает принимать.
216 lines
5.8 KiB
Go
216 lines
5.8 KiB
Go
package chat
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"sync"
|
||
"testing"
|
||
"time"
|
||
)
|
||
|
||
const (
|
||
uidA UserID = "u-a"
|
||
uidB UserID = "u-b"
|
||
tg Address = "tg://123"
|
||
tui Address = "tui://local"
|
||
)
|
||
|
||
// fakeOnMsg — тест-колбэк, копящий входящие. Т.к. Router теперь обрабатывает
|
||
// входящие асинхронно (воркер-горутина), доступ потокобезопасный, а ожидание
|
||
// нужного числа сообщений — через wait.
|
||
type fakeOnMsg struct {
|
||
mu sync.Mutex
|
||
ch chan struct{} // сигнал о появлении каждого нового входящего
|
||
got []Incoming
|
||
}
|
||
|
||
func newFakeOnMsg() *fakeOnMsg {
|
||
return &fakeOnMsg{ch: make(chan struct{}, 64)}
|
||
}
|
||
|
||
func (f *fakeOnMsg) h(inc Incoming) {
|
||
f.mu.Lock()
|
||
f.got = append(f.got, inc)
|
||
f.mu.Unlock()
|
||
f.ch <- struct{}{}
|
||
}
|
||
|
||
// wait блокируется, пока не наберётся n входящих. Возвращает false по таймауту.
|
||
func (f *fakeOnMsg) wait(n int) bool {
|
||
deadline := time.After(2 * time.Second)
|
||
for {
|
||
f.mu.Lock()
|
||
got := len(f.got)
|
||
f.mu.Unlock()
|
||
if got >= n {
|
||
return true
|
||
}
|
||
select {
|
||
case <-f.ch:
|
||
case <-deadline:
|
||
return false
|
||
}
|
||
}
|
||
}
|
||
|
||
func (f *fakeOnMsg) get(i int) Incoming {
|
||
f.mu.Lock()
|
||
defer f.mu.Unlock()
|
||
return f.got[i]
|
||
}
|
||
|
||
func (f *fakeOnMsg) count() int {
|
||
f.mu.Lock()
|
||
defer f.mu.Unlock()
|
||
return len(f.got)
|
||
}
|
||
|
||
func TestRouter_AttachAndIncoming(t *testing.T) {
|
||
cb := newFakeOnMsg()
|
||
r := NewRouter(cb.h)
|
||
|
||
tgCh := newFakeChannel(tg)
|
||
if err := r.Attach(tgCh); err != nil {
|
||
t.Fatalf("Attach: %v", err)
|
||
}
|
||
|
||
tgCh.emit(uidA, tg, "привет")
|
||
if !cb.wait(1) {
|
||
t.Fatal("handler не получил входящее за таймаут")
|
||
}
|
||
got := cb.get(0)
|
||
if got.UserID != uidA || got.Address != tg || got.Msg.Text != "привет" {
|
||
t.Errorf("incoming = %+v", got)
|
||
}
|
||
}
|
||
|
||
func TestRouter_AttachNil(t *testing.T) {
|
||
r := NewRouter(nil)
|
||
if err := r.Attach(nil); err == nil {
|
||
t.Fatal("Attach(nil) должен вернуть ошибку")
|
||
}
|
||
}
|
||
|
||
func TestRouter_Send_NoRoute(t *testing.T) {
|
||
r := NewRouter(nil)
|
||
// M1: нет маршрута → Send no-op, nil
|
||
if err := r.Send(context.Background(), uidA, Message{Text: "x"}); err != nil {
|
||
t.Fatalf("Send без маршрута: %v", err)
|
||
}
|
||
}
|
||
|
||
func TestRouter_Send_UsesCurrentRoute(t *testing.T) {
|
||
cb := newFakeOnMsg()
|
||
r := NewRouter(cb.h)
|
||
tgCh := newFakeChannel(tg)
|
||
|
||
_ = r.Attach(tgCh)
|
||
tgCh.emit(uidA, tg, "hi") // устанавливает маршрут
|
||
if !cb.wait(1) {
|
||
t.Fatal("маршрут не установился за таймаут")
|
||
}
|
||
|
||
if err := r.Send(context.Background(), uidA, Message{Text: "отв"}); err != nil {
|
||
t.Fatalf("Send: %v", err)
|
||
}
|
||
if tgCh.sentCount() != 1 {
|
||
t.Fatalf("sent = %d, want 1", tgCh.sentCount())
|
||
}
|
||
if tgCh.sent[0].Msg.Text != "отв" {
|
||
t.Errorf("msg = %q", tgCh.sent[0].Msg.Text)
|
||
}
|
||
}
|
||
|
||
func TestRouter_SwitchChannel_Continues(t *testing.T) {
|
||
cb := newFakeOnMsg()
|
||
r := NewRouter(cb.h)
|
||
tgCh := newFakeChannel(tg)
|
||
tuiCh := newFakeChannel(tui)
|
||
_ = r.Attach(tgCh)
|
||
_ = r.Attach(tuiCh)
|
||
|
||
// начал в TG
|
||
tgCh.emit(uidA, tg, "hi")
|
||
// продолжил в GUI
|
||
tuiCh.emit(uidA, tui, "продолжаю тут")
|
||
if !cb.wait(2) {
|
||
t.Fatal("входящие не обработаны за таймаут")
|
||
}
|
||
if tgCh.sentCount() != 0 || tuiCh.sentCount() != 0 {
|
||
t.Fatal("до Send ничего не шлём")
|
||
}
|
||
|
||
// ответ должен уйти в последний канал (GUI)
|
||
_ = r.Send(context.Background(), uidA, Message{Text: "отв"})
|
||
if tuiCh.sentCount() != 1 {
|
||
t.Errorf("tui sent = %d, want 1 (последний маршрут)", tuiCh.sentCount())
|
||
}
|
||
if tgCh.sentCount() != 0 {
|
||
t.Errorf("tg sent = %d, want 0", tgCh.sentCount())
|
||
}
|
||
}
|
||
|
||
func TestRouter_Ask_PendingThenAnswer(t *testing.T) {
|
||
cb := newFakeOnMsg()
|
||
r := NewRouter(cb.h)
|
||
tgCh := newFakeChannel(tg)
|
||
_ = r.Attach(tgCh)
|
||
tgCh.emit(uidA, tg, "hi")
|
||
if !cb.wait(1) {
|
||
t.Fatal("первое входящее не обработано")
|
||
}
|
||
|
||
prompt := Message{Text: "Как зовут?", Options: []Option{{ID: "a", Label: "Анна"}}}
|
||
if err := r.Ask(context.Background(), uidA, prompt); err != nil {
|
||
t.Fatalf("Ask: %v", err)
|
||
}
|
||
if _, ok := r.Pending(uidA); !ok {
|
||
t.Fatal("pending не открыт")
|
||
}
|
||
// повторный Ask → M3
|
||
if err := r.Ask(context.Background(), uidA, prompt); !errors.Is(err, ErrWaitingAnswer) {
|
||
t.Fatalf("второй Ask err = %v, want ErrWaitingAnswer", err)
|
||
}
|
||
|
||
// ответ с того же адреса потребляет pending
|
||
tgCh.emit(uidA, tg, "Анна")
|
||
if !cb.wait(2) {
|
||
t.Fatal("ответ не обработан")
|
||
}
|
||
if _, ok := r.Pending(uidA); ok {
|
||
t.Fatal("pending должен быть закрыт после ответа")
|
||
}
|
||
if cb.count() != 2 {
|
||
t.Fatalf("handler got %d, want 2 (hi + ответ)", cb.count())
|
||
}
|
||
if cb.get(1).Msg.QuestionID == "" {
|
||
t.Error("ответ должен нести QuestionID вопроса")
|
||
}
|
||
}
|
||
|
||
func TestRouter_Ask_NoRoute(t *testing.T) {
|
||
r := NewRouter(nil)
|
||
if err := r.Ask(context.Background(), uidB, Message{Text: "q"}); !errors.Is(err, ErrRouteNotFound) {
|
||
t.Fatalf("Ask без маршрута err = %v, want ErrRouteNotFound", err)
|
||
}
|
||
}
|
||
|
||
func TestRouter_Ask_PendingNotConsumedFromOtherAddr(t *testing.T) {
|
||
cb := newFakeOnMsg()
|
||
r := NewRouter(cb.h)
|
||
tgCh := newFakeChannel(tg)
|
||
tuiCh := newFakeChannel(tui)
|
||
_ = r.Attach(tgCh)
|
||
_ = r.Attach(tuiCh)
|
||
tgCh.emit(uidA, tg, "hi")
|
||
|
||
if err := r.Ask(context.Background(), uidA, Message{Text: "q"}); err != nil {
|
||
t.Fatalf("Ask: %v", err)
|
||
}
|
||
// ответ из ДРУГОГО канала → это новый message, pending НЕ потребляется
|
||
tuiCh.emit(uidA, tui, "ответ из gui")
|
||
if _, ok := r.Pending(uidA); !ok {
|
||
t.Fatal("pending должен остаться (ответ из другого адреса)")
|
||
}
|
||
}
|