Files
bare/internal/api/events.go
T
mayatnikovandClaude Opus 5 63a7a1ef52 Этап 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
2026-08-22 18:18:04 +03:00

119 lines
4.1 KiB
Go

package api
import (
"fmt"
"io"
"net/http"
"time"
"github.com/xmatic-squad/bare/internal/auth"
)
// pingEvery — период комментария-пинга: он держит соединение живым
// через прокси и показывает клиенту, что поток цел (docs/protocol.md).
const pingEvery = 20 * time.Second
// GET /api/events?device= — поток событий устройства (ADR-004).
// Устройство передаётся в query: EventSource не умеет заголовки.
//
// Last-Event-ID игнорируется: механизм восстановления — не докрутка
// по идентификатору, а повторная выдача очереди при каждом подключении.
func (s *server) events(w http.ResponseWriter, r *http.Request) {
sess, _ := auth.From(r)
device := r.URL.Query().Get("device")
if !validID(device) {
unknownDevice(w)
return
}
owned, err := s.st.DeviceOwned(r.Context(), device, sess.Nick)
if err != nil {
s.internal(w, r, err)
return
}
if !owned {
unknownDevice(w)
return
}
// Порядок — docs/protocol.md, «События»: сначала push_pending и
// last_seen, потом поток, потом очередь. Сорвавшаяся запись last_seen
// не должна рвать исправный поток устройства, а он был бы уже закрыт
// открытием нового.
now := time.Now().UnixMilli()
if err := s.st.TouchDevice(r.Context(), device, now); err != nil {
s.internal(w, r, err)
return
}
// Поток открывается до чтения очереди: конверт, попавший в очередь
// между выборкой и подпиской, иначе пролежал бы там до следующего
// подключения. Обратная крайность — дубль, а его клиент сливает по id
// (ADR-017). Открытие закрывает прежний поток этого устройства.
stream := s.hub.Open(device)
defer stream.Close()
queued, err := s.st.Queue(r.Context(), device)
if err != nil {
s.internal(w, r, err)
return
}
head := w.Header()
head.Set("Content-Type", "text/event-stream")
head.Set("Cache-Control", "no-cache")
// nginx буферизует ответы проксируемых приложений; для потока это
// означало бы, что события копятся и не уходят (docs/deploy.md).
head.Set("X-Accel-Buffering", "no")
w.WriteHeader(http.StatusOK)
send := sender(w)
for _, envelope := range queued {
if !send("msg", envelope) {
return
}
}
if !send("ready", "{}") {
return
}
ping := time.NewTicker(pingEvery)
defer ping.Stop()
for {
select {
case <-r.Context().Done():
// Клиент ушёл.
return
case <-stream.Done():
// Поток закрыли: новое соединение того же устройства,
// удаление устройства или остановка сервера.
return
case ev := <-stream.Events():
if !send(ev.Name, ev.Data) {
return
}
case <-ping.C:
if !write(w, ": ping\n\n") {
return
}
}
}
}
// sender собирает функцию записи события. Данные — компактный JSON
// без переводов строки, поэтому кадр SSE собирается одной строкой data.
// Ответ false означает, что писать больше некуда: соединение оборвалось.
func sender(w http.ResponseWriter) func(name, data string) bool {
return func(name, data string) bool {
return write(w, fmt.Sprintf("event: %s\ndata: %s\n\n", name, data))
}
}
func write(w http.ResponseWriter, frame string) bool {
if _, err := io.WriteString(w, frame); err != nil {
return false
}
// Без Flush кадр остался бы в буфере net/http до конца ответа,
// а конца у потока нет.
return http.NewResponseController(w).Flush() == nil
}