booblik-net
1. Зона ответственности
Провод. Модуль владеет тем, как запрос превращается в байты и обратно, и тем, как эти байты доезжают до сокета: кадрирование, разбор PRODUCE и FETCH, приём соединений, готовность сокета, отдача тела ответа мимо кучи.
Чем не занимается:
- не знает, как лог устроен на диске — только просит
PartitionLogиPartitionWriter; - не создаёт топики и партиции. Их набор задаётся при старте
Broker.open, всё остальное — ошибкаUNKNOWN_TOPIC_OR_PARTITION. Слоя метаданных нет: договариваться о состоянии кластера — ровно та часть Kafka, без которой проект и затевался; - не хранит позиции потребителей: оффсет помнит клиент;
- не шифрует. TLS и zero-copy физически несовместимы (ресёрч §1.2), и это не «пока не сделали».
Главный инвариант: ответы на одном соединении выходят строго в порядке запросов. Клиенту
разрешено держать несколько запросов в полёте и сопоставлять ответы по correlationId — если
сервер начнёт отвечать вразнобой, сломается ровно тот сценарий, ради которого конвейеризация
существует.
2. Контракт
protocol-wire — формат кадра, оба запроса, коды ошибок.
2а. Ключевые файлы (якоря кода)
| Файл | Что там |
|---|---|
src/main/kotlin/.../net/wire/Protocol.kt | константы формата, ApiKey, ErrorCode |
.../net/wire/Requests.kt | разбор кадра; здесь же все проверки длин, пришедших с провода |
.../net/wire/ResponseEncoder.kt | сборка ответов; заголовок FETCH отдельно от тела |
.../net/nio/SelectorLoop.kt | движок готовности: интерес → ожидание → возобновление |
.../net/nio/Connection.kt | два транспорта за одним интерфейсом |
.../net/Session.kt | протокольный цикл, один на оба транспорта |
.../net/BooblikServer.kt | acceptor, конфигурация, реестр партиций |
.../net/Broker.kt | топики и партиции: каталоги, логи, писатели, retention |
.../net/wire/RequestEncoder.kt | сборка запросов; общий для обоих клиентов |
.../net/client/BooblikConnection.kt | конвейерное соединение: сопоставление ответов |
.../net/client/Producer.kt | накопление записей в батчи |
.../net/client/Consumer.kt | чтение вперёд с собственной позицией |
.../net/client/BooblikClient.kt | блокирующий клиент для тестов и стенда |
3. Как устроено
Свой селектор, а не Ktor. Решение Р3. Публичный API Ktor отдаёт ByteChannel, а
FileChannel.transferTo нужен настоящий SocketChannel (ресёрч §1.3). SelectorLoop —
минимальный движок, который это даёт: зарегистрировать интерес, приостановить корутину,
возобновить по готовности.
Три вещи внутри держат всю конструкцию, и каждая ломается тихо:
- изменения интереса ставятся в очередь и применяются потоком селектора. Вызов
interestOpsиз другого потока во времяselect()на разных платформах либо блокируется, либо молча не срабатывает; выглядит это как зависший клиент под нагрузкой; - интерес снимается до возобновления ожидающего. Оставленный, он на level-triggered селекторе крутит цикл на полной скорости по каналу, которого никто не ждёт;
- регистрация тоже идёт через очередь.
registerконкурирует сselect()и на части реализаций JDK это дедлок, а не задержка.
Один цикл сессии на два транспорта. Connection абстрагирует чтение, запись и transferTo.
Это не красота ради красоты: M-36 сравнивает селектор с виртуальными потоками, и сравнение
имело бы смысл, только если по обе стороны исполняется один и тот же протокольный код.
Заголовок FETCH пишется отдельно от тела. Тело не существует ни в одном буфере — оно едет из page cache в сокет. Поэтому кодируется только заголовок, а длина в нём обещает байты, которые сервер отправит следом.
Срез удерживается на время всего ответа. PartitionLog.openFetch держит сегмент от момента,
когда посчитали длину, до момента, когда отправили последний байт. Иначе retention может увести
сегмент между обещанием и отправкой, и ответ окажется короче собственного префикса длины.
3а. Клиент
Соединение конвейерное, а сопоставление — очередь, не словарь. Брокер отвечает строго в
порядке запросов, поэтому ожидающие лежат в FIFO, а correlationId сверяется с головой этой
очереди. Словарь молча пережил бы нарушение порядка; очередь падает громко. Порядок — обещание
протокола, и клиент, неспособный заметить его нарушение, однажды отдаст не ту запись не тому.
Пишет в сокет одна корутина. Несколько вызывающих отправляют одновременно, а сокет не терпит двух переплетённых записей. Вместо лока — канал и один потребитель: заодно это бесплатно даёт тот самый порядок, на который опирается очередь.
Producer копит записи, и это не оптимизация, а суть. M-14 намерил, что единица записи стоит
54× — больше любого другого решения в проекте. Клиент, отправляющий записи по одной, выбрасывает
бо́льшую часть брокера. Окно linger отсчитывается от первой записи батча: от последней —
и ровный ручеёк запросов откладывал бы отправку бесконечно, превращая ограничение задержки
в надежду на неё.
Consumer помнит позицию сам. Это половина причины, по которой здесь нет ни групп, ни
координатора, ни хранения оффсетов: оффсет — число, которое читатель и так знает.
4. Зависимости
| Тип | Имя | Для чего |
|---|---|---|
| Module | :booblik-client | клиент и общий кодек; сервер зависит от клиента, а не наоборот (M-80) |
| Module | :booblik-core | лог, партиция, актор записи |
| Library | kotlinx-coroutines-core | сессии как корутины |
5. Конфигурация
ServerConfig:
| Параметр | Дефолт | Смысл |
|---|---|---|
port | 0 | 0 = любой свободный; фактический адрес возвращает start() |
transport | SELECTOR | VIRTUAL_THREADS — опорная точка для замеров, не рабочий режим |
fetchMode | ZERO_COPY | HEAP — контроль в эксперименте M-35 |
tcpNoDelay | true | брокер шлёт целые ответы; накопление добавляет только задержку |
backlog | 1024 |
6. Инфраструктура и деплой
Пока нечего: сервер поднимается из кода, исполняемого артефакта нет. Появится вместе с конфигурацией (M-50).
7. Локальный запуск
./gradlew :booblik-net:testНагрузка по сети — benchmarking, проба probeLoad.
8. Сознательные ограничения / грабли
- Привязка к wildcard может тихо уступить порт чужому процессу.
SO_REUSEADDRвключён по умолчанию, и на BSD-системах привязка к[::]:Pпроходит, даже если кто-то уже слушает[::1]:P; соединения при этом достаются более специфичному слушателю. Брокер стартует, печатает порт, принимает ноль соединений и не сообщает ни о чём — исключений нет, счётчики по нулям, поток селектора жив. Разбор — ресёрч §1.16, лечение —ServerConfig.bindAddress(booblik.bind.address). Тесты привязываются к127.0.0.1явно. - Цикл приёма соединений переживает ошибку, а не умирает от неё. Он один на весь брокер и
живёт корутиной под
SupervisorJob: раньше единственное исключение прекращало приём навсегда и молча, оставляя процесс живым и порт занятым. Отказ одного соединения восстановим, отказ всех будущих — нет. - Сессия, умершая на исключении, попадает в метрики (
sessionFailures,lastSessionFailure). Клиент в этот момент видит обрыв вместо ответа, и раньше причина выбрасывалась в пустоту — из-за чего M-64 диагностировалась по косвенным признакам. connectionsAcceptedиconnectionsOpened— разные числа. Между «сокет снят с очереди» и «сессия началась» лежит настройка сокета и регистрация в селекторе; сокет может умереть там, и без двух счётчиков этот участок невидим.correlationIdв конвейерном клиенте выдаётся атомарно. Обычныйvar i++отдавал двум одновременным вызовам один и тот же идентификатор, после чего ответ доставался не тому, кто его ждал, — ровно та поломка, ради предотвращения которой correlationId и существует. ПойманоClientTest.ackPolicy = NONEне получает ответа вообще. Не пустой ответ, а тишина: оффсета не существует, пока актор не дошёл до батча. Клиент в этом режиме не узнаёт и о перегрузке.- Ответ на ошибку несёт
correlationId = 0только тогда, когда его действительно негде взять — то есть если кадр короче заголовка. Если заголовок разобрался, идентификатор возвращается, даже когда всё остальное сломано: клиент сопоставляет ответы по нему, и ответ с чужим идентификатором не теряет информацию, а разрешает чужой запрос. - Неизвестный
apiKeyилиapiVersion— этоUNSUPPORTED_VERSION, а неCORRUPT_REQUEST. Кадр безупречен, брокер просто не умеет того, что просят. - Абсурдная длина кадра стоит соединения. Длина приходит с провода; аллоцировать столько,
сколько попросили, — это способ убить брокер одним пакетом, тем более при куче в 64 МиБ.
Потолок —
Protocol.MAX_FRAME_BYTES. - FETCH никогда не пересекает границу сегмента. Один вызов — один файл. Потребитель просит ещё раз с тем оффсетом, до которого дочитал.
- Ответ обрывается по границе байтов, а не записи. Хвост незавершённой записи отбрасывает клиент. Брокер батч не разбирает — именно этого разбора избегает zero-copy.
VIRTUAL_THREADS— линейка, а не альтернатива.transferToблокирует на дисковом вводе-выводе, а файловый ввод-вывод виртуальные потоки не виртуализуют: страничный промах внутриsendfileпиннит несущий поток (ресёрч §1.4). Замер M-36 это подтвердил числом.- Одна сессия обслуживается последовательно. Конвейеризация даёт несколько запросов в полёте, но сервер отвечает по одному — иначе порядок ответов поедет.