Этап 2: чат 1:1 — устройства, очередь, SSE, шифрование сообщений

Сервер: регистрация устройств и X-Device, hub с одним потоком на устройство,
очередь per-device с фан-аутом без эха отправителю, POST /api/messages
с проверками в порядке protocol.md, ACK, SSE с воспроизведением очереди,
ready и пингом раз в 20 секунд, контакты в обе стороны при первом сообщении,
лимит 30 сообщений в минуту.

Клиент: ULID, ключ 1:1 из ECDH через HKDF, шифрование конверта с AAD,
sync.js как единственный писатель в IndexedDB, ACK строго после записи,
список чатов, экран чата по эталону, разделители дат и «новые»,
pending и failed с повтором, полоса «нет соединения».

ADR-033: у неотправленного есть текст отказа — clock_skew стало видно.
ADR-034: входящее с известным id не перезаписывает запись. Собеседник знает
открытый id конверта и подменял им чужое сообщение в чужой истории — вплоть
до стирания своего присланного, чего «удалить у всех не существует» не допускает.
ADR-035: один поток событий на браузерный профиль (locks + BroadcastChannel):
две вкладки отбирали поток друг у друга и оставались без живой доставки.
ADR-036: повтор отправки сохраняет ULID, пока он в пределах окна часов, —
иначе потерянный ответ давал у собеседника два сообщения вместо одного.

Приёмка на боевом сервере: два аккаунта, пять устройств, живая доставка,
копия на второе устройство, очередь офлайн-устройству, ACK, подмена from
игнорируется, чужой deviceId и запрос без Origin отбиваются, плейнтекста
в базе и WAL ноль вхождений.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_015DbCjVfTFq4ZFG8juD45YJ
This commit is contained in:
2026-08-22 18:18:04 +03:00
co-authored by Claude Opus 5
parent 597c55301c
commit 63a7a1ef52
39 changed files with 5126 additions and 46 deletions
+123
View File
@@ -0,0 +1,123 @@
// Package hub держит открытые SSE-потоки устройств (ADR-004, ADR-017).
//
// На устройство приходится один поток: новое соединение закрывает
// предыдущее. Потерянное живое событие не теряет сообщения — оно лежит
// в очереди до ACK и выдаётся заново при следующем подключении
// (docs/protocol.md, «События»).
package hub
import "sync"
// buffer — сколько событий поток держит, пока обработчик их не разобрал.
const buffer = 32
// Event — одно событие SSE: имя и готовый JSON. Данные — строка: её
// нельзя изменить после того, как она ушла в несколько потоков сразу.
type Event struct {
Name string
Data string
}
// Hub — карта «устройство → открытый поток». Пуст, пока никто не подключён.
type Hub struct {
mu sync.Mutex
streams map[string]*Stream
}
// New заводит пустой hub.
func New() *Hub { return &Hub{streams: make(map[string]*Stream)} }
// Stream — поток одного устройства. Обработчик читает Events до тех пор,
// пока не закроется Done или не уйдёт клиент.
type Stream struct {
hub *Hub
device string
events chan Event
done chan struct{}
once sync.Once
}
// Open открывает поток устройства и закрывает предыдущий, если он был:
// одно соединение на устройство (docs/protocol.md, «События»).
func (h *Hub) Open(device string) *Stream {
s := &Stream{
hub: h,
device: device,
events: make(chan Event, buffer),
done: make(chan struct{}),
}
h.mu.Lock()
prev := h.streams[device]
h.streams[device] = s
h.mu.Unlock()
if prev != nil {
prev.stop()
}
return s
}
// Send отдаёт событие подключённому устройству. Устройство не подключено —
// молча ничего: конверт уже лежит в его очереди.
func (h *Hub) Send(device string, ev Event) {
h.mu.Lock()
s := h.streams[device]
h.mu.Unlock()
if s == nil {
return
}
select {
case s.events <- ev:
default:
// Клиент не успевает читать. Закрываем поток: переподключение
// выдаст очередь целиком, а копить события в памяти сервера —
// не его дело (ADR-008).
h.drop(s)
}
}
// Close закрывает поток устройства: устройство удалили (docs/protocol.md,
// «Устройства»).
func (h *Hub) Close(device string) {
h.mu.Lock()
s := h.streams[device]
delete(h.streams, device)
h.mu.Unlock()
if s != nil {
s.stop()
}
}
// CloseAll закрывает все потоки: сервер останавливается. Без этого
// остановка ждала бы, пока клиенты уйдут сами.
func (h *Hub) CloseAll() {
h.mu.Lock()
streams := h.streams
h.streams = make(map[string]*Stream)
h.mu.Unlock()
for _, s := range streams {
s.stop()
}
}
// drop снимает регистрацию именно этого потока и закрывает его. Если
// устройство успело подключиться заново, новый поток остаётся на месте.
func (h *Hub) drop(s *Stream) {
h.mu.Lock()
if h.streams[s.device] == s {
delete(h.streams, s.device)
}
h.mu.Unlock()
s.stop()
}
// Events — события, пришедшие потоку.
func (s *Stream) Events() <-chan Event { return s.events }
// Done закрывается, когда поток закрыт: новым соединением того же
// устройства, удалением устройства или остановкой сервера.
func (s *Stream) Done() <-chan struct{} { return s.done }
// Close закрывает поток — его зовёт обработчик, когда клиент ушёл.
func (s *Stream) Close() { s.hub.drop(s) }
func (s *Stream) stop() { s.once.Do(func() { close(s.done) }) }
+170
View File
@@ -0,0 +1,170 @@
package hub
import (
"strconv"
"sync"
"testing"
"time"
)
const wait = 2 * time.Second
// next ждёт событие потока.
func next(t *testing.T, s *Stream) Event {
t.Helper()
select {
case ev := <-s.Events():
return ev
case <-time.After(wait):
t.Fatal("событие не пришло")
}
return Event{}
}
// closed ждёт закрытия потока.
func closed(t *testing.T, s *Stream) {
t.Helper()
select {
case <-s.Done():
case <-time.After(wait):
t.Fatal("поток не закрылся")
}
}
func open(t *testing.T, s *Stream) {
t.Helper()
select {
case <-s.Done():
t.Fatal("поток закрыт")
default:
}
}
func TestSend(t *testing.T) {
h := New()
s := h.Open("device")
h.Send("device", Event{Name: "msg", Data: `{"id":"1"}`})
if ev := next(t, s); ev.Name != "msg" || ev.Data != `{"id":"1"}` {
t.Errorf("событие: %+v", ev)
}
// Неподключённое устройство — молча ничего: конверт лежит в очереди.
h.Send("другое", Event{Name: "msg"})
open(t, s)
}
// Одно соединение на устройство: новое закрывает предыдущее.
func TestOpenClosesPrevious(t *testing.T) {
h := New()
first := h.Open("device")
second := h.Open("device")
closed(t, first)
open(t, second)
h.Send("device", Event{Name: "msg"})
if ev := next(t, second); ev.Name != "msg" {
t.Errorf("событие ушло не в тот поток: %+v", ev)
}
}
// Close закрывает поток устройства: устройство удалили.
func TestCloseDevice(t *testing.T) {
h := New()
s := h.Open("device")
h.Close("device")
closed(t, s)
// Второе закрытие и закрытие неизвестного устройства — не беда.
h.Close("device")
h.Close("другое")
}
// CloseAll — остановка сервера.
func TestCloseAll(t *testing.T) {
h := New()
first := h.Open("first")
second := h.Open("second")
h.CloseAll()
closed(t, first)
closed(t, second)
}
// Клиент, который не читает, теряет поток, а не память сервера:
// переподключение выдаст очередь целиком.
func TestOverflowDropsStream(t *testing.T) {
h := New()
s := h.Open("device")
for i := 0; i < buffer+1; i++ {
h.Send("device", Event{Name: "msg", Data: strconv.Itoa(i)})
}
closed(t, s)
// Место в карте освободилось: следующее подключение начинает с нуля.
fresh := h.Open("device")
h.Send("device", Event{Name: "msg", Data: "снова"})
if ev := next(t, fresh); ev.Data != "снова" {
t.Errorf("событие: %+v", ev)
}
}
// Закрытие потока обработчиком не трогает уже открытый новый.
func TestStreamCloseKeepsNewer(t *testing.T) {
h := New()
first := h.Open("device")
second := h.Open("device")
first.Close()
h.Send("device", Event{Name: "msg"})
if ev := next(t, second); ev.Name != "msg" {
t.Errorf("событие: %+v", ev)
}
}
// Доставки идут из разных горутин: hub обязан это переживать.
func TestConcurrent(t *testing.T) {
h := New()
done := make(chan struct{})
var wg sync.WaitGroup
for i := 0; i < 4; i++ {
wg.Add(1)
go func(n int) {
defer wg.Done()
device := "device" + strconv.Itoa(n%2)
for {
select {
case <-done:
return
default:
}
h.Send(device, Event{Name: "msg"})
h.Open(device)
}
}(i)
}
// Читатель, чтобы буфер не переполнялся мгновенно.
wg.Add(1)
go func() {
defer wg.Done()
s := h.Open("device0")
for {
select {
case <-done:
return
case <-s.Events():
case <-s.Done():
s = h.Open("device0")
}
}
}()
time.Sleep(50 * time.Millisecond)
close(done)
wg.Wait()
h.CloseAll()
}