From c9d69594ccdebed132552ac856a82f67cd13c1b6 Mon Sep 17 00:00:00 2001 From: Hermes Date: Fri, 14 Aug 2026 20:33:04 +0500 Subject: [PATCH] =?UTF-8?q?internal/chat:=20=D0=BC=D1=83=D0=BB=D1=8C=D1=82?= =?UTF-8?q?=D0=B8=D0=BA=D0=B0=D0=BD=D0=B0=D0=BB=D1=8C=D0=BD=D1=8B=D0=B9=20?= =?UTF-8?q?Router=20(UserID-=D1=81=D0=B5=D1=81=D1=81=D0=B8=D0=B8,=20routin?= =?UTF-8?q?g,=20pending=20Ask/=D0=BE=D1=82=D0=B2=D0=B5=D1=82)=20+=20Channe?= =?UTF-8?q?l=20=D0=B8=D0=BD=D1=82=D0=B5=D1=80=D1=84=D0=B5=D0=B9=D1=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/chat/errors.go | 20 ++++ internal/chat/fake_channel_test.go | 65 ++++++++++++ internal/chat/questionid.go | 18 ++++ internal/chat/router.go | 123 +++++++++++++++++++++++ internal/chat/router_test.go | 156 +++++++++++++++++++++++++++++ internal/chat/types.go | 70 +++++++++++++ 6 files changed, 452 insertions(+) create mode 100644 internal/chat/errors.go create mode 100644 internal/chat/fake_channel_test.go create mode 100644 internal/chat/questionid.go create mode 100644 internal/chat/router.go create mode 100644 internal/chat/router_test.go create mode 100644 internal/chat/types.go diff --git a/internal/chat/errors.go b/internal/chat/errors.go new file mode 100644 index 0000000..8b9012d --- /dev/null +++ b/internal/chat/errors.go @@ -0,0 +1,20 @@ +package chat + +import "errors" + +// Классы ошибок (M1..M5) — по конвенции проекта. Определены как отдельные +// sentinel-ошибки, чтобы вызывающий мог отличить их errors.Is(). + +var ( + // M1 — нет маршрута для UserID. Send → no-op; Ask → ошибка. + ErrRouteNotFound = errors.New("chat: route not found for user") + + // M2 — канал не смог доставить исходящее. + ErrSendFailed = errors.New("chat: send failed") + + // M3 — Ask при уже открытом pending-вопросе (не затираем). + ErrWaitingAnswer = errors.New("chat: pending question already open") + + // M5 — Run-цикл канала упал фатально (Router перезапускает канал). + ErrChannelFatal = errors.New("chat: channel fatal") +) diff --git a/internal/chat/fake_channel_test.go b/internal/chat/fake_channel_test.go new file mode 100644 index 0000000..0289817 --- /dev/null +++ b/internal/chat/fake_channel_test.go @@ -0,0 +1,65 @@ +package chat + +import ( + "context" + "sync" +) + +// fakeChannel — тестовый канал: захватывает исходящие и умеет эмулировать +// входящие события. Не запускает реальный цикл. +type fakeChannel struct { + mu sync.Mutex + handler Handler + sent []SendCall + addr Address + runErr error +} + +type SendCall struct { + To Address + Msg Message +} + +func newFakeChannel(addr Address) *fakeChannel { + return &fakeChannel{addr: addr} +} + +func (f *fakeChannel) OnMessage(h Handler) { + f.mu.Lock() + defer f.mu.Unlock() + f.handler = h +} + +func (f *fakeChannel) Run(ctx context.Context) error { return f.runErr } + +func (f *fakeChannel) Send(_ context.Context, to Address, m Message) error { + f.mu.Lock() + defer f.mu.Unlock() + f.sent = append(f.sent, SendCall{To: to, Msg: m}) + return nil +} + +func (f *fakeChannel) Ask(_ context.Context, to Address, m Message) error { + f.mu.Lock() + defer f.mu.Unlock() + f.sent = append(f.sent, SendCall{To: to, Msg: m}) + return nil +} + +func (f *fakeChannel) Close() error { return nil } + +// emit доставляет входящее через зарегистрированный handler. +func (f *fakeChannel) emit(uid UserID, addr Address, text string) { + f.mu.Lock() + h := f.handler + f.mu.Unlock() + if h != nil { + h(Incoming{UserID: uid, Address: addr, Msg: Message{Text: text}, Channel: f}) + } +} + +func (f *fakeChannel) sentCount() int { + f.mu.Lock() + defer f.mu.Unlock() + return len(f.sent) +} diff --git a/internal/chat/questionid.go b/internal/chat/questionid.go new file mode 100644 index 0000000..efa3085 --- /dev/null +++ b/internal/chat/questionid.go @@ -0,0 +1,18 @@ +package chat + +import ( + "crypto/rand" + "encoding/hex" + "time" +) + +// newQuestionID генерирует уникальный идентификатор вопроса (блок 4: +// идентификатор в каждом Message, чтобы связать ответ с вопросом). +func newQuestionID() string { + b := make([]byte, 8) + if _, err := rand.Read(b); err != nil { + // крипто-rand недоступен — fallback по времени (уникален на практике) + return "q_" + hex.EncodeToString([]byte(time.Now().Format("150405.000000000"))) + } + return "q_" + hex.EncodeToString(b) +} diff --git a/internal/chat/router.go b/internal/chat/router.go new file mode 100644 index 0000000..94fa097 --- /dev/null +++ b/internal/chat/router.go @@ -0,0 +1,123 @@ +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) +} + +// NewRouter создаёт роутер. onUserMsg — колбэк обработки входящего. +func NewRouter(onUserMsg func(Incoming)) *Router { + if onUserMsg == nil { + onUserMsg = func(Incoming) {} + } + return &Router{ + sessions: map[UserID]any{}, + routes: map[UserID]Route{}, + pending: map[UserID]PendingQ{}, + onUserMsg: onUserMsg, + } +} + +// 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() + r.onUserMsg(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 +} diff --git a/internal/chat/router_test.go b/internal/chat/router_test.go new file mode 100644 index 0000000..6f04b7e --- /dev/null +++ b/internal/chat/router_test.go @@ -0,0 +1,156 @@ +package chat + +import ( + "context" + "errors" + "testing" +) + +const ( + uidA UserID = "u-a" + uidB UserID = "u-b" + tg Address = "tg://123" + tui Address = "tui://local" +) + +// fakeOnMsg — тест-колбэк, копящий входящие. +type fakeOnMsg struct{ got []Incoming } + +func (f *fakeOnMsg) h(inc Incoming) { f.got = append(f.got, inc) } + +func TestRouter_AttachAndIncoming(t *testing.T) { + cb := &fakeOnMsg{} + r := NewRouter(cb.h) + + tgCh := newFakeChannel(tg) + if err := r.Attach(tgCh); err != nil { + t.Fatalf("Attach: %v", err) + } + + tgCh.emit(uidA, tg, "привет") + if len(cb.got) != 1 { + t.Fatalf("handler got %d, want 1", len(cb.got)) + } + got := cb.got[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 := &fakeOnMsg{} + r := NewRouter(cb.h) + tgCh := newFakeChannel(tg) + + _ = r.Attach(tgCh) + tgCh.emit(uidA, tg, "hi") // устанавливает маршрут + + 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 := &fakeOnMsg{} + 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 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 := &fakeOnMsg{} + r := NewRouter(cb.h) + tgCh := newFakeChannel(tg) + _ = r.Attach(tgCh) + tgCh.emit(uidA, tg, "hi") + + 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 _, 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.got[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 := &fakeOnMsg{} + 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 должен остаться (ответ из другого адреса)") + } +} diff --git a/internal/chat/types.go b/internal/chat/types.go new file mode 100644 index 0000000..ceb5391 --- /dev/null +++ b/internal/chat/types.go @@ -0,0 +1,70 @@ +// Package chat — мультиканальное общение с пользователем. +// +// Единый Router поверх любых Channel (Telegram / TUI / GUI...). Сессии +// ключуются по UserID («кто»), доставка — по Address («где»), а pull/push +// разница (polling/scan/подписка) скрыта внутри каждого Channel. +// +// Примитивы обмена: Send (уведомить), Ask (спросить, асинхронно), +// + входящие события через OnMessage. Ask привязан к адресу вопроса. +package chat + +import "context" + +// Address — «где» пользователь: конкретный канал+адрес (tg://123, tui://local). +type Address string + +// UserID — «кто» владелец сессии (один на человека, канал неважен). +type UserID string + +// Option — выбор (кнопка в TG, пронумерованный пункт в TUI). +type Option struct { + ID string + Label string +} + +// Message — один исходящий/входящий контент. +type Message struct { + Text string + Options []Option // Confirm = Ask с да/нет; выбор = список + QuestionID string // блок 4: идентификатор вопроса в каждом Message +} + +// Incoming — входящее событие, канал сам заполняет UserID/Address. +type Incoming struct { + UserID UserID + Address Address + Msg Message + Channel Channel +} + +// Handler — единый вход всех каналов в Router. +type Handler func(Incoming) + +// Channel — любой канал связи. pull/push разница скрыта внутри Run. +type Channel interface { + // Run запускает цикл канала (polling/scan/подписка). Может вернуть + // M5 ChannelFatal — тогда Router перезапустит канал. + Run(ctx context.Context) error + // OnMessage регистрирует единый обработчик входящих (Router). + OnMessage(h Handler) + // Send — fire-and-forget уведомление. M1 если нет адреса, M2 если не доставлено. + Send(ctx context.Context, to Address, m Message) error + // Ask — неблокирующий вопрос: регистрирует pending и возвращается сразу. + // M3 если pending уже открыт. Ответ придёт через OnMessage с тем же QuestionID. + Ask(ctx context.Context, to Address, prompt Message) error + // Close закрывает канал. + Close() error +} + +// PendingQ — незакрытый вопрос (блок 4: один на сессию). +type PendingQ struct { + UserID UserID + AnswerAddr Address + Msg Message +} + +// Route — куда доставить следующий ответ пользователю. +type Route struct { + Address Address + Channel Channel +}