From a64e3d6cc378eaed1fd1a8f2a336b84a0d81c193 Mon Sep 17 00:00:00 2001 From: Hermes Date: Tue, 18 Aug 2026 19:55:46 +0500 Subject: [PATCH] =?UTF-8?q?feat(chat):=20=D0=B0=D1=81=D0=B8=D0=BD=D1=85?= =?UTF-8?q?=D1=80=D0=BE=D0=BD=D0=BD=D0=B0=D1=8F=20=D0=BE=D0=B1=D1=80=D0=B0?= =?UTF-8?q?=D0=B1=D0=BE=D1=82=D0=BA=D0=B0=20=D0=B2=D1=85=D0=BE=D0=B4=D1=8F?= =?UTF-8?q?=D1=89=D0=B8=D1=85=20=E2=80=94=20=D0=BE=D1=82=D0=B7=D1=8B=D0=B2?= =?UTF-8?q?=D1=87=D0=B8=D0=B2=D1=8B=D0=B9=20=D0=B8=D0=BD=D1=82=D0=B5=D1=80?= =?UTF-8?q?=D1=84=D0=B5=D0=B9=D1=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Router теперь обрабатывает входящие в воркер-горутине (FIFO-очередь с буфером 256) вместо синхронного вызова onUserMsg из long-poll цикла канала. Долгий вызов аналитика (Decide) больше не блокирует приём новых сообщений от Telegram: цикл getUpdates продолжает работать. - router.go: NewRouter запускает processLoop; handleIncoming кладёт событие в канал и возвращается; маршрутизация + pending по-прежнему обновляются синхронно под мьютексом. - router_test.go: fakeOnMsg стал потокобезопасным с ожиданием числа входящих (wait), т.к. обработка теперь асинхронная. Преимущества: интерфейс не замирает на время анализа; порядок входящих сохраняется (FIFO). Ограничение: воркер один — при очень долгом аналитике следующие сообщения ждут в очереди, но канал их продолжает принимать. --- internal/chat/router.go | 26 ++++++++++- internal/chat/router_test.go | 87 ++++++++++++++++++++++++++++++------ 2 files changed, 97 insertions(+), 16 deletions(-) diff --git a/internal/chat/router.go b/internal/chat/router.go index 94fa097..b62a9ce 100644 --- a/internal/chat/router.go +++ b/internal/chat/router.go @@ -18,6 +18,13 @@ type Router struct { // Hook, вызываемый на каждое входящее событие (обычно → process_turn). onUserMsg func(Incoming) + + // Асинхронная обработка входящих: handleIncoming кладёт событие в канал, + // воркер-горутина последовательно вызывает onUserMsg. Благодаря этому + // long-poll цикл канала (Telegram) не блокируется на время долгого + // вызова аналитика и продолжает принимать новые сообщения. + incoming chan Incoming + wg sync.WaitGroup } // NewRouter создаёт роутер. onUserMsg — колбэк обработки входящего. @@ -25,11 +32,22 @@ func NewRouter(onUserMsg func(Incoming)) *Router { if onUserMsg == nil { onUserMsg = func(Incoming) {} } - return &Router{ + 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.wg.Done() } } @@ -59,7 +77,11 @@ func (r *Router) handleIncoming(inc Incoming) { delete(r.pending, inc.UserID) } r.mu.Unlock() - r.onUserMsg(inc) + + // Асинхронная обработка: кладём событие в очередь воркера и сразу + // возвращаемся, не блокируя вызывающий long-poll цикл канала. + r.wg.Add(1) + r.incoming <- inc } // Send уведомляет пользователя через текущий маршрут. M1 (нет маршрута) — no-op, diff --git a/internal/chat/router_test.go b/internal/chat/router_test.go index 6f04b7e..855212f 100644 --- a/internal/chat/router_test.go +++ b/internal/chat/router_test.go @@ -3,7 +3,9 @@ package chat import ( "context" "errors" + "sync" "testing" + "time" ) const ( @@ -13,13 +15,58 @@ const ( tui Address = "tui://local" ) -// fakeOnMsg — тест-колбэк, копящий входящие. -type fakeOnMsg struct{ got []Incoming } +// fakeOnMsg — тест-колбэк, копящий входящие. Т.к. Router теперь обрабатывает +// входящие асинхронно (воркер-горутина), доступ потокобезопасный, а ожидание +// нужного числа сообщений — через wait. +type fakeOnMsg struct { + mu sync.Mutex + ch chan struct{} // сигнал о появлении каждого нового входящего + got []Incoming +} -func (f *fakeOnMsg) h(inc Incoming) { f.got = append(f.got, inc) } +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 := &fakeOnMsg{} + cb := newFakeOnMsg() r := NewRouter(cb.h) tgCh := newFakeChannel(tg) @@ -28,10 +75,10 @@ func TestRouter_AttachAndIncoming(t *testing.T) { } tgCh.emit(uidA, tg, "привет") - if len(cb.got) != 1 { - t.Fatalf("handler got %d, want 1", len(cb.got)) + if !cb.wait(1) { + t.Fatal("handler не получил входящее за таймаут") } - got := cb.got[0] + got := cb.get(0) if got.UserID != uidA || got.Address != tg || got.Msg.Text != "привет" { t.Errorf("incoming = %+v", got) } @@ -53,12 +100,15 @@ func TestRouter_Send_NoRoute(t *testing.T) { } func TestRouter_Send_UsesCurrentRoute(t *testing.T) { - cb := &fakeOnMsg{} + 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) @@ -72,7 +122,7 @@ func TestRouter_Send_UsesCurrentRoute(t *testing.T) { } func TestRouter_SwitchChannel_Continues(t *testing.T) { - cb := &fakeOnMsg{} + cb := newFakeOnMsg() r := NewRouter(cb.h) tgCh := newFakeChannel(tg) tuiCh := newFakeChannel(tui) @@ -83,6 +133,9 @@ func TestRouter_SwitchChannel_Continues(t *testing.T) { 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 ничего не шлём") } @@ -98,11 +151,14 @@ func TestRouter_SwitchChannel_Continues(t *testing.T) { } func TestRouter_Ask_PendingThenAnswer(t *testing.T) { - cb := &fakeOnMsg{} + 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 { @@ -118,13 +174,16 @@ func TestRouter_Ask_PendingThenAnswer(t *testing.T) { // ответ с того же адреса потребляет pending tgCh.emit(uidA, tg, "Анна") + if !cb.wait(2) { + t.Fatal("ответ не обработан") + } if _, ok := r.Pending(uidA); ok { t.Fatal("pending должен быть закрыт после ответа") } - if len(cb.got) != 2 { - t.Fatalf("handler got %d, want 2 (hi + ответ)", len(cb.got)) + if cb.count() != 2 { + t.Fatalf("handler got %d, want 2 (hi + ответ)", cb.count()) } - if cb.got[1].Msg.QuestionID == "" { + if cb.get(1).Msg.QuestionID == "" { t.Error("ответ должен нести QuestionID вопроса") } } @@ -137,7 +196,7 @@ func TestRouter_Ask_NoRoute(t *testing.T) { } func TestRouter_Ask_PendingNotConsumedFromOtherAddr(t *testing.T) { - cb := &fakeOnMsg{} + cb := newFakeOnMsg() r := NewRouter(cb.h) tgCh := newFakeChannel(tg) tuiCh := newFakeChannel(tui)