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 должен остаться (ответ из другого адреса)") } } // TestRouter_PerUserOrdering проверяет, что сообщения одного пользователя // обрабатываются строго в порядке поступления (пул воркеров не перемешивает). func TestRouter_PerUserOrdering(t *testing.T) { cb := newFakeOnMsg() r := NewRouter(cb.h) tgCh := newFakeChannel(tg) _ = r.Attach(tgCh) for _, txt := range []string{"1", "2", "3"} { tgCh.emit(uidA, tg, txt) } if !cb.wait(3) { t.Fatal("сообщения не обработаны за таймаут") } for i, want := range []string{"1", "2", "3"} { if got := cb.get(i).Msg.Text; got != want { t.Errorf("порядок обработки нарушен: idx %d = %q, want %q", i, got, want) } } } // TestRouter_ParallelismAcrossUsers проверяет, что пока обработчик одного // пользователя заблокирован (долгий LLM-вызов), сообщение другого пользователя // обрабатывается в другом воркере, а второе сообщение того же пользователя — // ждёт своей очереди (per-user порядок). func TestRouter_ParallelismAcrossUsers(t *testing.T) { r := NewRouter(nil) tgCh := newFakeChannel(tg) tuiCh := newFakeChannel(tui) _ = r.Attach(tgCh) _ = r.Attach(tuiCh) blocked := make(chan struct{}) release := make(chan struct{}) var muLocal sync.Mutex seen := make([]string, 0, 3) signal := make(chan struct{}, 8) h := func(inc Incoming) { if inc.Msg.Text == "block" { close(blocked) <-release // держим воркера, пока не отпустим } muLocal.Lock() seen = append(seen, inc.Msg.Text) muLocal.Unlock() signal <- struct{}{} } r.onUserMsg = h snapshot := func() []string { muLocal.Lock() defer muLocal.Unlock() return append([]string(nil), seen...) } waitFor := func(n int) bool { deadline := time.After(2 * time.Second) for len(snapshot()) < n { select { case <-signal: case <-deadline: return false } } return true } // первое сообщение A блокирует своего воркера tgCh.emit(uidA, tg, "block") <-blocked // B обрабатывается параллельно, пока A висит tuiCh.emit(uidB, tui, "B1") select { case <-signal: case <-time.After(100 * time.Millisecond): t.Fatal("B не обработан, пока блокирован A — чат-путь снова сериализован") } if got := snapshot(); len(got) != 1 || got[0] != "B1" { t.Fatalf("ожидали обработку B1, got %v", got) } // второе сообщение A НЕ обрабатывается, пока занят воркер A (порядок per-user) tgCh.emit(uidA, tg, "a2") select { case <-signal: t.Fatal("сообщение A обработано ДО освобождения A — нарушен per-user порядок") case <-time.After(80 * time.Millisecond): } // отпускаем A → дообрабатывается a2 close(release) if !waitFor(3) { t.Fatal("итоговые сообщения не обработаны") } got := snapshot() if got[2] != "a2" { t.Errorf("порядок персональной очереди нарушен: pos2 = %q, want a2", got[2]) } }