feat(chat): асинхронная обработка входящих — отзывчивый интерфейс
Router теперь обрабатывает входящие в воркер-горутине (FIFO-очередь с буфером 256) вместо синхронного вызова onUserMsg из long-poll цикла канала. Долгий вызов аналитика (Decide) больше не блокирует приём новых сообщений от Telegram: цикл getUpdates продолжает работать. - router.go: NewRouter запускает processLoop; handleIncoming кладёт событие в канал и возвращается; маршрутизация + pending по-прежнему обновляются синхронно под мьютексом. - router_test.go: fakeOnMsg стал потокобезопасным с ожиданием числа входящих (wait), т.к. обработка теперь асинхронная. Преимущества: интерфейс не замирает на время анализа; порядок входящих сохраняется (FIFO). Ограничение: воркер один — при очень долгом аналитике следующие сообщения ждут в очереди, но канал их продолжает принимать.
This commit is contained in:
@@ -18,6 +18,13 @@ type Router struct {
|
|||||||
|
|
||||||
// Hook, вызываемый на каждое входящее событие (обычно → process_turn).
|
// Hook, вызываемый на каждое входящее событие (обычно → process_turn).
|
||||||
onUserMsg func(Incoming)
|
onUserMsg func(Incoming)
|
||||||
|
|
||||||
|
// Асинхронная обработка входящих: handleIncoming кладёт событие в канал,
|
||||||
|
// воркер-горутина последовательно вызывает onUserMsg. Благодаря этому
|
||||||
|
// long-poll цикл канала (Telegram) не блокируется на время долгого
|
||||||
|
// вызова аналитика и продолжает принимать новые сообщения.
|
||||||
|
incoming chan Incoming
|
||||||
|
wg sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewRouter создаёт роутер. onUserMsg — колбэк обработки входящего.
|
// NewRouter создаёт роутер. onUserMsg — колбэк обработки входящего.
|
||||||
@@ -25,11 +32,22 @@ func NewRouter(onUserMsg func(Incoming)) *Router {
|
|||||||
if onUserMsg == nil {
|
if onUserMsg == nil {
|
||||||
onUserMsg = func(Incoming) {}
|
onUserMsg = func(Incoming) {}
|
||||||
}
|
}
|
||||||
return &Router{
|
r := &Router{
|
||||||
sessions: map[UserID]any{},
|
sessions: map[UserID]any{},
|
||||||
routes: map[UserID]Route{},
|
routes: map[UserID]Route{},
|
||||||
pending: map[UserID]PendingQ{},
|
pending: map[UserID]PendingQ{},
|
||||||
onUserMsg: onUserMsg,
|
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)
|
delete(r.pending, inc.UserID)
|
||||||
}
|
}
|
||||||
r.mu.Unlock()
|
r.mu.Unlock()
|
||||||
r.onUserMsg(inc)
|
|
||||||
|
// Асинхронная обработка: кладём событие в очередь воркера и сразу
|
||||||
|
// возвращаемся, не блокируя вызывающий long-poll цикл канала.
|
||||||
|
r.wg.Add(1)
|
||||||
|
r.incoming <- inc
|
||||||
}
|
}
|
||||||
|
|
||||||
// Send уведомляет пользователя через текущий маршрут. M1 (нет маршрута) — no-op,
|
// Send уведомляет пользователя через текущий маршрут. M1 (нет маршрута) — no-op,
|
||||||
|
|||||||
@@ -3,7 +3,9 @@ package chat
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -13,13 +15,58 @@ const (
|
|||||||
tui Address = "tui://local"
|
tui Address = "tui://local"
|
||||||
)
|
)
|
||||||
|
|
||||||
// fakeOnMsg — тест-колбэк, копящий входящие.
|
// fakeOnMsg — тест-колбэк, копящий входящие. Т.к. Router теперь обрабатывает
|
||||||
type fakeOnMsg struct{ got []Incoming }
|
// входящие асинхронно (воркер-горутина), доступ потокобезопасный, а ожидание
|
||||||
|
// нужного числа сообщений — через 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) {
|
func TestRouter_AttachAndIncoming(t *testing.T) {
|
||||||
cb := &fakeOnMsg{}
|
cb := newFakeOnMsg()
|
||||||
r := NewRouter(cb.h)
|
r := NewRouter(cb.h)
|
||||||
|
|
||||||
tgCh := newFakeChannel(tg)
|
tgCh := newFakeChannel(tg)
|
||||||
@@ -28,10 +75,10 @@ func TestRouter_AttachAndIncoming(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
tgCh.emit(uidA, tg, "привет")
|
tgCh.emit(uidA, tg, "привет")
|
||||||
if len(cb.got) != 1 {
|
if !cb.wait(1) {
|
||||||
t.Fatalf("handler got %d, want 1", len(cb.got))
|
t.Fatal("handler не получил входящее за таймаут")
|
||||||
}
|
}
|
||||||
got := cb.got[0]
|
got := cb.get(0)
|
||||||
if got.UserID != uidA || got.Address != tg || got.Msg.Text != "привет" {
|
if got.UserID != uidA || got.Address != tg || got.Msg.Text != "привет" {
|
||||||
t.Errorf("incoming = %+v", got)
|
t.Errorf("incoming = %+v", got)
|
||||||
}
|
}
|
||||||
@@ -53,12 +100,15 @@ func TestRouter_Send_NoRoute(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func TestRouter_Send_UsesCurrentRoute(t *testing.T) {
|
func TestRouter_Send_UsesCurrentRoute(t *testing.T) {
|
||||||
cb := &fakeOnMsg{}
|
cb := newFakeOnMsg()
|
||||||
r := NewRouter(cb.h)
|
r := NewRouter(cb.h)
|
||||||
tgCh := newFakeChannel(tg)
|
tgCh := newFakeChannel(tg)
|
||||||
|
|
||||||
_ = r.Attach(tgCh)
|
_ = r.Attach(tgCh)
|
||||||
tgCh.emit(uidA, tg, "hi") // устанавливает маршрут
|
tgCh.emit(uidA, tg, "hi") // устанавливает маршрут
|
||||||
|
if !cb.wait(1) {
|
||||||
|
t.Fatal("маршрут не установился за таймаут")
|
||||||
|
}
|
||||||
|
|
||||||
if err := r.Send(context.Background(), uidA, Message{Text: "отв"}); err != nil {
|
if err := r.Send(context.Background(), uidA, Message{Text: "отв"}); err != nil {
|
||||||
t.Fatalf("Send: %v", err)
|
t.Fatalf("Send: %v", err)
|
||||||
@@ -72,7 +122,7 @@ func TestRouter_Send_UsesCurrentRoute(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func TestRouter_SwitchChannel_Continues(t *testing.T) {
|
func TestRouter_SwitchChannel_Continues(t *testing.T) {
|
||||||
cb := &fakeOnMsg{}
|
cb := newFakeOnMsg()
|
||||||
r := NewRouter(cb.h)
|
r := NewRouter(cb.h)
|
||||||
tgCh := newFakeChannel(tg)
|
tgCh := newFakeChannel(tg)
|
||||||
tuiCh := newFakeChannel(tui)
|
tuiCh := newFakeChannel(tui)
|
||||||
@@ -83,6 +133,9 @@ func TestRouter_SwitchChannel_Continues(t *testing.T) {
|
|||||||
tgCh.emit(uidA, tg, "hi")
|
tgCh.emit(uidA, tg, "hi")
|
||||||
// продолжил в GUI
|
// продолжил в GUI
|
||||||
tuiCh.emit(uidA, tui, "продолжаю тут")
|
tuiCh.emit(uidA, tui, "продолжаю тут")
|
||||||
|
if !cb.wait(2) {
|
||||||
|
t.Fatal("входящие не обработаны за таймаут")
|
||||||
|
}
|
||||||
if tgCh.sentCount() != 0 || tuiCh.sentCount() != 0 {
|
if tgCh.sentCount() != 0 || tuiCh.sentCount() != 0 {
|
||||||
t.Fatal("до Send ничего не шлём")
|
t.Fatal("до Send ничего не шлём")
|
||||||
}
|
}
|
||||||
@@ -98,11 +151,14 @@ func TestRouter_SwitchChannel_Continues(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func TestRouter_Ask_PendingThenAnswer(t *testing.T) {
|
func TestRouter_Ask_PendingThenAnswer(t *testing.T) {
|
||||||
cb := &fakeOnMsg{}
|
cb := newFakeOnMsg()
|
||||||
r := NewRouter(cb.h)
|
r := NewRouter(cb.h)
|
||||||
tgCh := newFakeChannel(tg)
|
tgCh := newFakeChannel(tg)
|
||||||
_ = r.Attach(tgCh)
|
_ = r.Attach(tgCh)
|
||||||
tgCh.emit(uidA, tg, "hi")
|
tgCh.emit(uidA, tg, "hi")
|
||||||
|
if !cb.wait(1) {
|
||||||
|
t.Fatal("первое входящее не обработано")
|
||||||
|
}
|
||||||
|
|
||||||
prompt := Message{Text: "Как зовут?", Options: []Option{{ID: "a", Label: "Анна"}}}
|
prompt := Message{Text: "Как зовут?", Options: []Option{{ID: "a", Label: "Анна"}}}
|
||||||
if err := r.Ask(context.Background(), uidA, prompt); err != nil {
|
if err := r.Ask(context.Background(), uidA, prompt); err != nil {
|
||||||
@@ -118,13 +174,16 @@ func TestRouter_Ask_PendingThenAnswer(t *testing.T) {
|
|||||||
|
|
||||||
// ответ с того же адреса потребляет pending
|
// ответ с того же адреса потребляет pending
|
||||||
tgCh.emit(uidA, tg, "Анна")
|
tgCh.emit(uidA, tg, "Анна")
|
||||||
|
if !cb.wait(2) {
|
||||||
|
t.Fatal("ответ не обработан")
|
||||||
|
}
|
||||||
if _, ok := r.Pending(uidA); ok {
|
if _, ok := r.Pending(uidA); ok {
|
||||||
t.Fatal("pending должен быть закрыт после ответа")
|
t.Fatal("pending должен быть закрыт после ответа")
|
||||||
}
|
}
|
||||||
if len(cb.got) != 2 {
|
if cb.count() != 2 {
|
||||||
t.Fatalf("handler got %d, want 2 (hi + ответ)", len(cb.got))
|
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 вопроса")
|
t.Error("ответ должен нести QuestionID вопроса")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -137,7 +196,7 @@ func TestRouter_Ask_NoRoute(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func TestRouter_Ask_PendingNotConsumedFromOtherAddr(t *testing.T) {
|
func TestRouter_Ask_PendingNotConsumedFromOtherAddr(t *testing.T) {
|
||||||
cb := &fakeOnMsg{}
|
cb := newFakeOnMsg()
|
||||||
r := NewRouter(cb.h)
|
r := NewRouter(cb.h)
|
||||||
tgCh := newFakeChannel(tg)
|
tgCh := newFakeChannel(tg)
|
||||||
tuiCh := newFakeChannel(tui)
|
tuiCh := newFakeChannel(tui)
|
||||||
|
|||||||
Reference in New Issue
Block a user