Назад к блогу
Architecture12 min read

Архитектура peer- и channel-чата

В этой статье end-to-end объясняется, как работает offline-first чат PlainApp: как сообщение проходит от нажатия в UI до другого устройства через peer-транспорт, как групповые каналы распределяют сообщения множеству участников и как система сохраняет устойчивость при исчезновении сетей. Сопряжение (доверие и обмен ключами, бутстрапящий два устройства) рассмотрено в отдельной статье «Pairing Flow».

Содержание

Высокоуровневая архитектура {#high-level-architecture}

Чат PlainApp — serverless. На каждом устройстве работает встроенный Ktor HTTP-сервер, и устройства общаются напрямую через локальную сеть, Wi-Fi Aware (NAN) или Bluetooth Low Energy. Ретрансляционного сервера нет, облачного inbox нет, идентификации по номеру телефона нет. Устройства идентифицируются по самостоятельно сгенерированному clientId и аутентифицируются через рукопожатие Ed25519 + ECDH, выполняемое во время сопряжения.

Существуют два вида разговоров:

ТипКонстантаОписание
PEERChatTargetType.PEERПрямой чат 1-к-1 между двумя сопряжёнными устройствами.
CHANNELChatTargetType.CHANNELМногопартийный групповой чат, владельцем которого является одно устройство; участники распределяют сообщения друг другу.

Специальный "local" target — это блокнот самого устройства (заметки самому себе) — отправка в него является no-op на проводе.

Карта компонентов

Diagram 1
1

Архитектура намеренно многоуровневая:

  1. Точки входа UI / GraphQL никогда не касаются транспортов или БД напрямую.
  2. ChatManager — фасад; каждый вызывающий (UI, GraphQL-резолвер, peer-приёмник) проходит через него.
  3. ChatSender — диспетчер, ветвящийся по ChatTargetType и делегирующий отправителю пира или канала.
  4. Транспортный слой — это подключаемая цепочка стратегий с circuit breaking, поэтому нестабильный канал Wi-Fi Aware никогда не блокирует сообщение, которое могло бы пойти через BLE.

Модель данных {#data-model}

ChatTarget

Наименьшая единица маршрутизации — ChatTarget, пара (toId, type), где typePEER или CHANNEL. Он предоставляет encodedToId (peer:<id> или channel:<id>), который UI использует как стабильный ключ маршрутизации (например, TempData.activeToId, чтобы приёмник знал, нужно ли эмитить уведомление), проверку isLocal() (toId == "local") и компаньон parseId, восстанавливающий target из сохранённой строки.

Таблицы базы данных

Вся персистентность использует Room. Для чата важны три таблицы:

ТаблицаСущностьНазначение
chatsDChatОдна строка на сообщение (текст / изображение / файл).
chat_channelsDChatChannelОдна строка на групповой канал.
peersDPeerОдна строка на известное устройство (сопряжённое или только channel).

Diagram 2
2

Стоит отметить несколько моментов:

  • Идентификация — clientId, никогда MAC. Android рандомизирует BLE MAC при каждом соединении, поэтому база использует стабильный 13-символьный самосгенерированный id. Только 8-байтовый префикс SHA-256 (shortId) транслируется через BLE для обеспечения обнаружения.
  • Пиры со status="channel" — это участники канала, с которыми это устройство никогда напрямую не сопрягалось. Их поле key пусто — они аутентифицируются ключом канала вместо парного общего ключа.
  • owner="me" — сигнальное значение, позволяющее свежеустановленному устройству выступать владельцем до того, как его clientId станет стабильным; isOwnedByMe() принимает и "me", и TempData.clientId.

Поверхность GraphQL API {#graphql-api-surface}

PlainApp предоставляет две GraphQL-схемы:

  1. Web GraphQL (addChatChannelSchema + addChatMessageSchema) — обслуживается локальным Ktor-сервером для браузерного UI и для харнесса apitest/. Аутентифицируется ChaCha20-зашифрованным токеном.
  2. Peer GraphQL (PeerGraphQLService.applyPeerSchema) — предоставляется на /peer_graphql для других устройств через зашифрованный peer-транспорт. Аутентифицируется подписью Ed25519 + ChaCha20-шифрованием тела.

Обе схемы используют одни и те же синглтоны бизнес-логики (ChannelManager, ChatMessageReceiver, …), но предоставляют разные поверхности, поскольку модель доверия различается: web GraphQL доверяет локальному UI, тогда как peer GraphQL доверяет только криптографически аутентифицированным пирам.

Поверхность Web GraphQL (чат)

Запросы: chatChannels (список всех каналов), chatItems(id) (сообщения для target — id это "local", peer:<id> или channel:<id>), и latestChatItems (предпросмотр по всем чатам).

Мутации чата: sendChatItem(toId, content), deleteChatItem(id), deleteChatItems(query) и retryChatItem(id).

Мутации канала: createChatChannel(name), updateChatChannel(id, name), deleteChatChannel(id), leaveChatChannel(id), addChatChannelMember(id, peerId), removeChatChannelMember(id, peerId), acceptChatChannelInvite(id) и declineChatChannelInvite(id).

Поверхность Peer GraphQL (транспорт)

Предоставляется на /peer_graphql и аутентифицируется подписью Ed25519 + шифрованием тела ChaCha20. Границу транспорта пересекают только три мутации: createChatItem(content) (входящее peer-сообщение), channelSystemMessage(type, payload) (события жизненного цикла канала вроде invite/leave) и startAware (nudge, просящий пир запустить свой Wi-Fi Aware-сервис, чтобы более быстрый транспорт мог вступить).

HTTP-заголовок c-id несёт clientId отправителя; заголовок c-cid несёт id канала, когда запрос ограничен каналом (чтобы приёмник взял ключ канала, а не парный peer-ключ для расшифровки).

Peer-чат: отправка сообщения {#peer-chat-sending-a-message}

Когда пользователь нажимает Send в peer-разговоре, цепочка вызовов:

Diagram 3
3

Ключевые инварианты, обеспечиваемые на каждом шаге:

  1. ChatManager.createChatItem всегда сначала вставляет строку, затем отправляет. Это значит, что UI сразу видит пузырь «pending», и сообщение переживёт краш приложения, даже если доставка ещё не произошла.
  2. PeerGraphQLClient.buildSignedRequest строит конверт вида signature|timestamp|requestJson. Подпись Ed25519 поверх "$timestamp$requestJson", привязывая метку времени к телу, поэтому её нельзя переиграть с новой меткой.
  3. PeerTransportRouter.send перебирает транспорты в порядке Lan → WifiAware → Ble. Каждый транспорт может выбросить TransportUnavailable, чтобы маршрутизатор попробовал следующий.
  4. На принимающей стороне PeerChatParser.decrypt проверяет, что метка времени в ±5 min, и проверяет подпись Ed25519 до того, как мутация GraphQL вообще выполнится.
  5. ChatMessageReceiver.receive держит множество seenSignatures с ключом "$fromPeerId|$signature|$timestamp" и выбрасывает ReplayedMessageException при дубликатах — это существенно, поскольку транспорт может доставить одну и ту же полезную нагрузку дважды (LAN + BLE).

Если PeerChatSender.send возвращает непустую строку ошибки, ChatSender вызывает triggerPeerRediscovery(peerId), который запускает направленный, зашифрованный широковещательный DISCOVER, чтобы пир мог повторно анонсировать свой текущий IP/порт.

Peer-чат: получение сообщения {#peer-chat-receiving-a-message}

Входящие запросы попадают на маршрут /peer_graphql локального Ktor-сервера, обрабатываемый PeerGraphQLService:

Diagram 4
4

Уведомления

emitNotificationIfNeeded — финальный шаг. Он подавляет уведомление, когда TempData.activeToId == targetId (т. е. пользователь сейчас смотрит этот разговор), или когда canShowNotifications() равно false. Уведомления канала префиксируются именем отправителя.

Channel-чат: выбор лидера и fan-out {#channel-chat-leader-election--fan-out}

Каналы многопартийны, но serverless. Чтобы каждый участник не рассылал одно и то же сообщение N раз, отправляющая сторона выбирает одного лидера, чья задача — широковещательно разослать всем вступившим участникам.

Алгоритм выбора лидера (DChatChannel.electLeader)

  1. Фильтр до вступивших участников, которые сейчас онлайн (локальное устройство всегда считается онлайн).
  2. Если владелец среди онлайн-вступивших участников → владелец — лидер.
  3. Иначе лидер — онлайн-вступивший участник с наименьшим clientId (детерминированный тай-брейк, без координации).
  4. Возвращает null, если онлайн-вступивших участников нет.

Поток отправки

Diagram 5
5

Зачем вообще лидер?

Представьте 5-участниковый канал, где каждый вещает каждому: одно сообщение породило бы 20 сетевых round-trip и 4 дублирующих копии, приходящих каждому участнику. Выбрав одного лидера, fan-out делает только это устройство — отправитель либо сам выполняет fan-out (если он лидер), либо передаёт одну копию лидеру, который затем fan-out'ит.

Если лидер оффлайн, отправчик откатывается к Result.NoLeader, запускает rediscovery пира (чтобы IP лидера был найден) и очищает статус, чтобы пользователь мог повторить.

Маршрутизация по ключу канала

Сообщения канала шифруются ключом ChaCha20 канала, а не парным peer-ключом. Это и позволяет участнику, который познакомился с остальными только через канал (никогда не сопрягаясь 1-к-1), получать сообщения — его строка peers имеет status="channel" и key="". Отправитель ставит HTTP-заголовок c-cid в id канала; приёмник ищет ChannelCacher.getKeyBytes(channelId) вместо парного ключа.

Повторная попытка по получателю

Каждый sendToMember возвращает DMessageDeliveryResult. Агрегированный DMessageStatusData сохраняется как JSON status_data chat item. UI показывает «Доставлено Alice, Bob; Сбой для Carol» и позволяет пользователю нажать Retry конкретно для Carol — ChatManager.sendToChannelMembers перезапускает sendToRecipients для подмножества повторных попыток и слияет новые результаты с существующими, заменяя только повторяемых пиров.

Системные сообщения канала {#channel-system-messages}

Управляющие сообщения канала (invite, accept, decline, update, kick, leave) обмениваются через peer-GraphQL мутацию channelSystemMessage. Это JSON-полезные нагрузки, типизированные строкой type:

ТипНаправлениеПодписано?Назначение
channel_inviteOwner → inviteeДаПриглашение пира; несёт ключ канала + участников.
channel_invite_acceptInvitee → ownerНетПринятие; несёт публичный ключ принимающего.
channel_invite_declineInvitee → ownerНетОтклонение; владелец удаляет участника.
channel_updateOwner → all membersДаШироковещательное оповещение об изменении состава/имени.
channel_kickOwner → kicked peerДаЦелевой кик; также broadcast при удалении канала.
channel_leaveMember → ownerНетИнициированное участником уведомление о выходе.

Формат подписанной полезной нагрузки

Три подписанных типа (invite, update, kick) используют каноническую pipe-разделённую строку: "$channelId|$version|$action|$target", где action — один из invite, update, kick, а target — id приглашаемого/кикаемого пира (пусто для широковещательного kick).

Владелец подписывает эту строку своим ключом Ed25519. Приёмники отбрасывают любое сообщение, где channel.owner != fromId, до проверки подписи, и отбрасывают ChannelUpdate-полезные нагрузки, чей version локальной версии (защита от устаревшей версии против out-of-order доставки).

Diagram 6
6

Ленивая гидратация пиров

ChannelInvite и ChannelUpdate несут список memberPeers: List<MemberPeerInfo> — лёгкая информация о пире (id, name, publicKey, deviceType, ip, port) для каждого участника. ensureChannelPeer приёмника создаёт строку DPeer со status="channel" для любого участника, которого он раньше не видел. Это критично, поскольку маршрутизация fan-out требует peer-записи каждого участника для отправки сообщений.

Жизненный цикл канала {#channel-lifecycle}

Diagram 7
7

Слой peer-транспорта (LAN → Wi-Fi Aware → BLE) {#peer-transport-layer-lan--wi-fi-aware--ble}

PeerTransportRouterцепочка стратегий с circuit breaking. Упорядоченный список транспортов:

  1. LanTransport — первый выбор. Использует OkHttp с ChaCha20 crypto-интерсептором поверх HTTPS. Полностью пропускается, когда peer.ip пуст (peer в другой подсети, ещё не обнаруженный).
  2. WifiAwareTransport (только Android 13+) — использует Wi-Fi Aware (NAN) data paths. Быстрый skip, когда флаг awareRunning пира равен false (обновляется BLE-prewarmer сканированием). IPv6 пира разрешается через кастомный DNS, отображающий hostname plain-aware-peer на link-local-адрес.
  3. BleTransport — гарантированный fallback для любого сопряжённого пира. Стримит чанковый RPC через GATT. Медленнее, но работает без какой-либо IP-связности.

Diagram 8
8

Почему такой порядок?

  • LAN самый быстрый (один HTTPS round trip, ~10 мс тайм-аут).
  • Wi-Fi Aware средний (setup data-path ~5 с, затем ~10 мс round trip) и работает кросс-подсеть (например, одно устройство на guest Wi-Fi, другое на IoT Wi-Fi). Настроен на быстрый skip, когда Aware-сервис пира не запущен, избегая 10-секундного тайм-аута buildLink.
  • BLE самый медленный, но работает без какой-либо IP-связности — даже без Wi-Fi сообщение всё равно проходит. Используется как гарантированный fallback для сопряжённых пиров.

Circuit breaker обеспечивает, что нестабильный транспорт (особенно Wi-Fi Aware при сетевых колебаниях) пропускается на 30 с после 2 сбоев, поэтому fallback происходит быстро, не дожидаясь повторных 10-секундных тайм-аутов.

Рукопожатие Wi-Fi Aware

AwareSession делает двухсообщенийное рукопожатие перед открытием data path:

  • MSG_HELLO (subscriber → publisher): «Я вижу тебя, вот мой peer handle.»
  • MSG_READY (publisher → subscriber): «Я зарегистрировал свой network specifier, ты можешь делать requestNetwork сейчас.»

Это синхронизирует вызовы connectivityManager.requestNetwork(...) обеих сторон в пределах ~500 мс окна Android-фреймворка. Subscriber — сторона с меньшим clientId (детерминированное разделение ролей — обе стороны согласны без координации), и именно он владеет циклом повторных попыток.

Статус и presence пира {#peer-status--presence}

Presence отслеживается через длинноживущие WebSocket-соединения. Только одна сторона каждой пары открывает сокет — решается детерминированным правилом TempData.clientId < peer.id. Другая сторона принимает входящее соединение на /peer_status.

Diagram 9
9

PeerCacher.onlineMap — источник истины для presence. Он предоставляется как onlinePeerIds: StateFlow<Set<String>>, который потребляется алгоритмом выбора лидера канала (electLeader(onlinePeerIds, myId)).

Слой кэширования {#caching-layer}

Два кэша зеркалируют таблицы БД в памяти и предоставляют StateFlow, которые Compose собирает напрямую:

Diagram 10
10

Почему copy-on-write?

MutableStateFlow.distinctUntilChanged в Kotlin использует структурное равенство. Если бы мы мутировали DPeer на месте, производный список pairedPeers содержал бы ту же ссылку DPeer до и после, и distinctUntilChanged не увидел бы разницы и подавил эмиссию. Копируя сущность сначала, мутируя копию и заменяя запись в карте на новый PeerRuntime/ChannelRuntime, производный список получает новый список из новых ссылок, и flow срабатывает.

Загрузка файлов {#file-downloads}

Входящие файлы/изображения загружаются автоматически ограниченным пулом воркеров. Каждая загрузка стримится через любой доступный транспорт (PeerTransportRouter.downloadFile) и пишется во временный файл, затем импортируется в медиа-хранилище приложения и патчит поле uri chat item.

Diagram 11
11

Транспортно-агностичный стриминг

Абстракция DownloadedResponse(status, ByteReadChannel, onClose): AutoCloseable позволяет LAN и Wi-Fi Aware стримить живое HTTP-тело, тогда как BLE стримит чанковый RPC (чанки 16 KiB через GET /fs?id=…&offset=…&length=…) через тот же ByteReadChannel. Колбэк onClose позволяет BLE отменить фоновую корутину загрузки, когда потребитель закрывает ответ раньше (например, при паузе).

Сводка паттернов проектирования {#design-patterns-recap}

ПаттернГдеПочему
ФасадChatManagerЕдиная точка входа; вызывающие никогда не касаются БД/транспорта напрямую.
Strategy + Chain of Resp.PeerTransportRouter + LanTransport/WifiAwareTransport/BleTransportПодключаемые транспорты с TransportUnavailable как сигналом проваливания.
Circuit BreakerPeerCircuitBreaker2 сбоя / 30 с открывают leg (peer, transport), чтобы Wi-Fi Aware не блокировал fallback.
Конечный автоматPeerStatusManager.PeerState, AwarePeerLink.LinkStateЯвные переходы для жизненного цикла сокета и жизненного цикла NDP-линка.
Producer/Consumer + PoolDownloadQueue (3 воркера, Channel.BUFFERED)Ограниченная параллельность для загрузок файлов.
Observer / ReactiveStateFlow вездеCompose собирает напрямую; без ручного обновления.
Защита от повтораChatMessageReceiver.seenSignatures, PeerChatParser.MAX_TIMESTAMP_DIFF_MSОтбрасывать дубликаты от двойной доставки LAN+BLE; отвергать out-of-window метки времени.
Экспоненциальный backoffPeerStatusManager.scheduleReconnectmin(60 s, 1 s × 2^min(n-1, 6)) — кэп на 64 с.
Copy-on-WritePeerCacher.mutatePeer, ChannelCacher.mutateChannelЗаставляет StateFlow.distinctUntilChanged срабатывать на каждой мутации.
Подписанный конвертPeerGraphQLClient.buildSignedRequestsignature|timestamp|body — привязывает метку времени к телу для защиты от повторов.
Детерминированное разделение ролейTempData.clientId < peer.idРешает WebSocket-клиент против сервера и Wi-Fi Aware subscriber против publisher.
Ленивая гидратацияensureChannelPeer на invite/updateСоздаёт строки peers для невиданных участников канала, чтобы работала fan-out маршрутизация.
Шифрованная идентификацияLANDiscoverManager.discoverSpecificDeviceНаправленный DISCOVER шифрует целевой id ключом пира — только target его распознаёт.

Дополнительная литература

  • Pairing Flow — как два устройства устанавливают доверие и обмениваются общим ключом ChaCha20, используемым каждым транспортом в этой статье.
  • apitest/groups/chat-messages.sh и apitest/groups/chat-channels.sh — исполняемый тест-план, проверяющий каждую GraphQL-мутацию end-to-end.