Раздел 32 · Системное программирование: Zig, ассемблер, Verilog
Сервер на процессах и мультиплексирование ввода-вывода
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Сервер на процессах и мультиплексирование ввода-вывода
В прошлом уроке
zboxстал песочницей, но обслуживает по одному: пока один студент собирает свою задачу, остальные стоят в очередиlisten. Так же устроены эхо-сервер из урока про сокеты и TINY. Сегодня научим сервер не ждать одного клиента, пока другие готовы. Сначала способом, который у нас уже в руках:forkна каждое соединение. Потом совсем другим: один процесс, один поток, и ядро само говорит, какой из тысячи сокетов готов. Пройдёмselect,poll,epollиkqueue, и в конце замерим их на тысяче соединений, чтобы увидеть, гдеselectне хуже остальных, а где проигрывает в шесть раз.
Цели урока
- Объяснить, почему итеративный сервер держит всех за одним клиентом, и назвать три модели конкурентности главы 12 с ценой каждой.
- Написать эхо-сервер на процессах: кто закрывает какой дескриптор,
SIGCHLDи циклwaitpid,EINTRуaccept. - Понимать
select: битовая карта,FD_SETSIZE, первый аргументmaxfd + 1, набор, который надо собирать перед каждым вызовом. - Построить событийный сервер как пул конечных автоматов и не дать неполной строке остановить цикл.
- Перейти с
selectнаpollи сказать, что это даёт, а чего не даёт. - Объяснить, почему
epollиkqueueдешевле на тысяче молчащих соединений, и спрятать их за одной обёрткой. - Снять замер и прочитать таблицу: когда
selectдержится наравне сepoll, а когда проигрывает в разы.
Идея: пока один клиент думает
Вспомни эхо-сервер из урока про сокеты. Подключись к нему через nc и ничего не пиши. Второй клиент подключится (рукопожатие делает ядро), отправит строку и будет ждать ответа, пока первый не закроет соединение. Сервер не занят работой, он спит в read на первом сокете, а второй сокет давно готов. Беда не в скорости, а в том, что сервер ждёт не того.
Глава 12 книги называет программу конкурентной, если её логические потоки управления перекрываются во времени. Для сервера логический поток это сессия одного клиента. Способов перекрыть сессии три, и все три строятся на том, что мы уже прошли.
| Модель | Кто переключает сессии | Что у сессий общее | Главная цена |
|---|---|---|---|
| процессы | ядро | только файлы, открытые до fork | fork на соединение, отдельная память, обмен данными через ядро |
| мультиплексирование ввода-вывода | сама программа, в своём цикле | всё: один процесс и один поток | сложный код, одно ядро процессора, долгий клиент держит всех |
| потоки | ядро | всё адресное пространство | гонки на общих данных |
Сегодня первые две модели. Потоки придут в следующем уроке, и гонкам на них уйдут ещё три.
Вторую модель ты уже видел, только сверху. Цикл событий Node из урока про колбеки и цикл событий это она и есть: один поток, очередь готовых событий, обработчики, которым нельзя надолго занимать цикл. Здесь мы спустимся на уровень системного вызова, на котором этот цикл стоит.
Виджет: шесть моделей под одной нагрузкой
Прежде чем писать код, посмотри на все модели сразу. Виджет это симуляция, а не замер: цены fork, создания потока и вызова select в нём ориентировочные, и он сам их подписывает внизу. Клиенты приходят по одному, каждый запрос стоит немного процессора, а медленный клиент (плохая сеть) забирает ответ долго, и всё это время кто-то ждёт его на записи. Главная метрика это p99 быстрых клиентов: ждут ли они за медленными.
Что попробовать:
- Оставь сто клиентов и десять процентов медленных. У итеративного сервера p99 быстрых около 450 мс: каждый ждёт всех медленных перед собой. У процесса на соединение около 5 мс: родитель делает
forkсам, по очереди, и очередь заforkвидна в задержке. Уselectиepollоколо 50 мкс. - Подвинь молчащие соединения keep-alive на десять. Итеративный сервер встаёт насовсем: он висит в
readна первом молчуне. Пул из шестнадцати потоков ещё отвечает, но десять воркеров из шестнадцати заняты молчунами, и p99 быстрых вырастает до десятков миллисекунд. - Подвинь молчащих на тысячу. Процессов становится больше тысячи, и память сервера в модели переходит за гигабайт.
selectтратит внутри себя миллисекунды,epoll_waitвсё те же сто микросекунд. - Сдвинь ещё на одно деление, до 1010.
selectсдаётся: номер дескриптора не помещается вfd_set. Почему именно там, разберём ниже.
Поток на соединение и пул в виджете для сравнения, их мы напишем в уроках 69 и 71.
Сервер на процессах
Самый короткий путь к конкурентности: пусть каждого клиента обслуживает отдельный процесс. Родитель крутит accept и на каждое новое соединение делает fork из урока про процессы. Ребёнок получает копию всех дескрипторов, обслуживает клиента и выходит, а родитель сразу возвращается к accept.
Весь фокус в том, кто что закрывает. После fork на один и тот же сокет смотрят два дескриптора, в родителе и в ребёнке, и оба ведут в одну запись таблицы открытых файлов из урока про разделяемые файлы. Ядро закрывает соединение, только когда счётчик ссылок на эту запись падает до нуля.
- Родитель закрывает
connfdсразу послеfork. Забудет, и соединение не закроется никогда: ребёнок давно вышел, а ссылку держит родитель. Клиент будет ждать конца потока вечно, а родитель через несколько сотен соединений получитEMFILEнаaccept. - Ребёнок закрывает
listenfd. Ему слушающий сокет не нужен, а пока он открыт хоть у одного ребёнка, порт занят: перезапущенный сервер получитEADDRINUSEнаbind.
Второй вопрос, кто собирает детей. Вышедший ребёнок становится зомби, пока родитель не заберёт его статус через waitpid. Родитель занят accept и ждать детей сам не может, поэтому их собирает обработчик SIGCHLD. Из урока про сигналы помним две вещи. Сигналы не копятся: три ребёнка, вышедшие одновременно, могут дать один SIGCHLD, поэтому обработчик зовёт waitpid(-1, WNOHANG) в цикле, пока тот возвращает кого-то. И обработчик не портит errno тому, кого прервал: сохраняет его на входе и возвращает на выходе.
Третий вопрос, что станет с accept, пока приходит сигнал. Родитель большую часть времени спит в accept. Сигнал будит его, отрабатывает обработчик, и дальше выбор за флагом SA_RESTART: с ним ядро перезапустит accept само, без него accept вернёт -1 с EINTR. Мы ставим обработчик без SA_RESTART намеренно, чтобы увидеть EINTR своими глазами: он считается в interrupted и ведёт на следующий круг цикла. Флаг SA.NOCLDSTOP просит не присылать SIGCHLD, когда ребёнка всего лишь остановили.
//! `echoserverp` из главы 12: процесс на соединение. Родитель принимает
//! соединение и сразу закрывает свою копию `connfd`, ребёнок закрывает
//! `listenfd`, обслуживает клиента и выходит. Детей собирает обработчик
//! `SIGCHLD`: `waitpid(-1, WNOHANG)` в цикле, потому что сигналы не
//! копятся, и три умерших ребёнка могут прийти одним `SIGCHLD`.
const std = @import("std");
const c = std.c;
const posix = std.posix;
const Io = std.Io;
const socket = @import("../net/socket.zig");
const echo = @import("../net/echo.zig");
/// Сколько детей собрал обработчик. Нужен тестам: сервер вернулся, а
/// дети ещё дорабатывают своих клиентов.
pub var reaped: std.atomic.Value(usize) = .init(0);
/// Сколько раз `accept` прервался сигналом. Обёртка `Signal` из csapp.c
/// ставит `SA_RESTART`, а мы намеренно нет: так `accept` возвращает `EINTR`,
/// и видно, сколько раз `SIGCHLD` прервал ожидание соединения.
pub var interrupted: std.atomic.Value(usize) = .init(0);
fn onChild(_: posix.SIG) callconv(.c) void {
// errno общий для всего потока: обработчик не должен его портить
// тому, кого прервал.
const saved = c._errno().*;
defer c._errno().* = saved;
while (c.waitpid(-1, null, c.W.NOHANG) > 0) _ = reaped.fetchAdd(1, .monotonic);
}
/// Принимает `max_clients` соединений (0 значит вечно), каждое в своём
/// процессе. Возвращается сразу после последнего `accept`: дети живут
/// дальше и уходят сами, обработчик `SIGCHLD` остаётся стоять.
pub fn serve(listenfd: c.fd_t, max_clients: usize, log: *Io.Writer) !void {
posix.sigaction(.CHLD, &.{
.handler = .{ .handler = onChild },
.mask = posix.sigemptyset(),
.flags = posix.SA.NOCLDSTOP,
}, null);
var served: usize = 0;
while (max_clients == 0 or served < max_clients) : (served += 1) {
const connfd = while (true) {
const fd = c.accept(listenfd, null, null);
if (fd >= 0) break fd;
if (posix.errno(fd) != .INTR) return error.AcceptFailed;
_ = interrupted.fetchAdd(1, .monotonic);
};
const pid = c.fork();
if (pid < 0) {
socket.close(connfd);
return error.ForkFailed;
}
if (pid == 0) {
// Ребёнок: слушающий сокет ему не нужен. Не закроет, и порт
// останется занятым, пока жив хоть один ребёнок.
socket.close(listenfd);
echo.echo(connfd, log) catch {};
socket.close(connfd);
// `_exit`, а не возврат: в тестах ребёнок иначе побежал бы
// дальше по коду тестового раннера.
c._exit(0);
}
// Родитель: без этого `close` соединение не закроется никогда,
// счётчик ссылок на сокет держит и наша копия.
socket.close(connfd);
}
}
Сессию ребёнка делает echo.echo из урока 64 без единой правки: строки через буферизованный читатель до конца потока. Ребёнок выходит через _exit, а не возвратом из serve: иначе в тестах он побежал бы дальше по коду тестового раннера, и тестов стало бы два комплекта.
Сколько это стоит
Модель на процессах подкупает тем, что сессии не мешают друг другу. Упавший ребёнок не роняет сервер и соседей, медленный клиент держит только свой процесс. Именно поэтому CGI в TINY и программы в zbox запускаются отдельными процессами: там нужна изоляция, а не скорость.
Цена у изоляции такая:
forkне бесплатен даже с копированием при записи из урока про отображение памяти: ядро копирует таблицы страниц, дескрипторы, структуры процесса, и делает это родитель, по одному соединению за раз.- У каждого процесса свой набор таблиц страниц и свой стек. Тысяча клиентов это тысяча процессов в
ps. - Общего состояния нет. Посчитать, сколько байтов сервер вернул всем клиентам вместе, нельзя без канала через ядро: у каждого ребёнка своя копия счётчика. Даже строки
server receivedв логе печатают дети, каждый из своей копии буфера писателя.
Живой прогон на macOS:
$ tiny echoserverp 15241
echoserverp: listening on port 15241
server received 6 bytes
server received 27 bytes
$ printf 'hello\n' | tiny echoclient localhost 15241
hello
$ printf 'привет, сервер\n' | tiny echoclient 127.0.0.1 15241
привет, сервер
select: спросить ядро, кто готов
Итеративный сервер ошибается в одном: он спрашивает ядро про один дескриптор («прочитай из этого сокета») и спит, пока этот сокет не ответит. Можно спросить иначе: «вот набор дескрипторов, разбуди меня, когда хоть один из них станет готов, и скажи, какой». Такой вопрос называется мультиплексированием ввода-вывода, и первый системный вызов для него это select:
int select(int nfds, fd_set *readfds, fd_set *writefds, fd_set *exceptfds, struct timeval *timeout);
fd_set это битовая карта: бит номер fd поднят, значит дескриптор нас интересует. Карта фиксированного размера, 1024 бита, это константа FD_SETSIZE. nfds это не число дескрипторов, а граница: ядро смотрит биты от 0 до nfds - 1, поэтому передают наибольший дескриптор плюс один. Три набора это готовность к чтению, к записи и исключительные условия, нам хватит первого. timeout равный null значит ждать сколько угодно.
Главная особенность select в том, что наборы это одновременно вопрос и ответ. На входе ты поднимаешь биты тех, за кем следить, а на выходе ядро оставляет поднятыми только биты готовых, остальные гасит. Возвращает вызов число готовых дескрипторов.
Готовность к чтению для обычного сокета значит, что read не заснёт: пришли данные, либо клиент закрыл соединение (тогда read вернёт 0), либо на сокете ошибка. Для слушающего сокета готовность значит, что в очереди есть соединение и accept не заснёт.
В std Zig 0.16 нет ни select, ни fd_set, поэтому объявим их сами. Битовую карту удобно держать словами по 64 бита: на little-endian машине бит fd % 8 в байте fd / 8 у нас лежит там же, где его ищет libc. В glibc набор это шестнадцать long, в macOS тридцать два int32_t, и раскладка у всех трёх одинаковая. Функцию объявляем через extern "c": линковщик найдёт её в libc.
Вот первая программа из раздела 12.2 книги своими словами. Она слушает два источника сразу: клавиатуру (дескриптор 0) и слушающий сокет. Строка с клавиатуры печатается, новое соединение обслуживается эхом.
//! Два источника событий в одном цикле, как `select.c` из книги:
//! клавиатура (дескриптор 0) и слушающий сокет. `select` ждёт, пока
//! хотя бы один из них не станет готов, и говорит, какой именно.
const std = @import("std");
const c = std.c;
/// `fd_set` как битовая карта на 1024 дескриптора (`FD_SETSIZE`).
const FdSet = extern struct {
words: [1024 / 64]u64 = @splat(0),
fn set(s: *FdSet, fd: c.fd_t) void {
const i: usize = @intCast(fd);
s.words[i / 64] |= @as(u64, 1) << @intCast(i % 64);
}
fn clear(s: *FdSet, fd: c.fd_t) void {
const i: usize = @intCast(fd);
s.words[i / 64] &= ~(@as(u64, 1) << @intCast(i % 64));
}
fn isSet(s: *const FdSet, fd: c.fd_t) bool {
const i: usize = @intCast(fd);
return s.words[i / 64] & (@as(u64, 1) << @intCast(i % 64)) != 0;
}
};
// В std 0.16 нет ни `select`, ни `fd_set`: объявляем функцию libc сами.
extern "c" fn select(nfds: c_int, readfds: ?*FdSet, writefds: ?*FdSet, exceptfds: ?*FdSet, timeout: ?*c.timeval) c_int;
fn openListenfd(port: u16) !c.fd_t {
const fd = c.socket(c.AF.INET, c.SOCK.STREAM, 0);
if (fd < 0) return error.SocketFailed;
const one: c_int = 1;
_ = c.setsockopt(fd, c.SOL.SOCKET, c.SO.REUSEADDR, &one, @sizeOf(c_int));
const addr: c.sockaddr.in = .{ .port = std.mem.nativeToBig(u16, port), .addr = 0 };
if (c.bind(fd, @ptrCast(&addr), @sizeOf(c.sockaddr.in)) != 0) return error.BindFailed;
if (c.listen(fd, 1024) != 0) return error.ListenFailed;
return fd;
}
/// Строка с клавиатуры. `false`: конец ввода (Ctrl+D).
fn command() bool {
var buf: [256]u8 = undefined;
const n = c.read(0, &buf, buf.len);
if (n <= 0) return false;
std.debug.print("command: {s}", .{buf[0..@intCast(n)]});
return true;
}
/// Эхо до конца соединения, как `echo` из урока про сокеты. Пока клиент
/// не уйдёт, цикл `select` стоит, и клавиатура ждёт.
fn echo(connfd: c.fd_t) void {
var buf: [256]u8 = undefined;
while (true) {
const n = c.read(connfd, &buf, buf.len);
if (n <= 0) return;
std.debug.print("server received {d} bytes\n", .{n});
_ = c.write(connfd, &buf, @intCast(n));
}
}
pub fn main() !void {
const listenfd = try openListenfd(15213);
var read_set: FdSet = .{};
read_set.set(0);
read_set.set(listenfd);
while (true) {
// `select` пишет ответ поверх запроса, поэтому ждём на копии.
var ready_set = read_set;
if (select(listenfd + 1, &ready_set, null, null, null) < 0) return error.SelectFailed;
if (ready_set.isSet(0) and !command()) {
std.debug.print("stdin closed\n", .{});
read_set.clear(0);
}
if (ready_set.isSet(listenfd)) {
const connfd = c.accept(listenfd, null, null);
if (connfd < 0) continue;
echo(connfd);
_ = c.close(connfd);
}
}
}
Собирается отдельно: zig build-exe select_stdin.zig на macOS, на Linux добавь -lc. Прогон ниже сделал скрипт: он пишет программе в stdin и подключается к ней по часам, слева секунды от старта. ls набран на 0.3 с, клиент подключился и прислал hello на 0.6 с, status набран на 0.9 с, клиент ушёл на 2.9 с, stdin закрыт на 3.2 с, второй клиент прислал bye на 3.5 с:
0.38 с command: ls
0.62 с server received 6 bytes
2.93 с command: status
3.24 с stdin closed
3.54 с server received 4 bytes
Строка status ждала две секунды. select честно ответил, что слушающий сокет готов, но дальше программа ушла в echo и сидела там, пока клиент не закрыл соединение. Мультиплексирование не спасает, если после ответа ядра программа сама засыпает на одном дескрипторе. Правило событийного сервера отсюда: на одно событие делать ровно столько работы, сколько можно сделать без ожидания, и возвращаться в цикл.
Событийный сервер: клиент как конечный автомат
Чтобы вернуться в цикл посреди сессии, сессию надо уметь остановить и продолжить. Всё, что сервер помнит о клиенте между событиями, это его состояние, а событие «сокет готов к чтению» переводит клиента из одного состояния в другое. Так каждый клиент становится конечным автоматом, а сервер крутит сразу много автоматов в одном цикле.
У эхо-сервера состояние клиента простое: байты, которые пришли, но ещё не сложились в строку. Переход тоже: один read в хвост буфера, все целые строки отправить обратно, остаток сохранить.
Почему ровно один read? Готовность обещает, что один read не заснёт. Второй уже может: данные кончились, а клиент ещё не закрыл соединение. Книжный check_clients читает клиента через rio_readlineb, который читает, пока не встретит \n. Клиент, приславший половину строки, усыпит такой сервер вместе со всеми остальными клиентами: это одно из домашних упражнений главы. Наш Conn этой дыры не имеет, и тест шага проверяет это прямо.
//! Общее для событийных эхо-серверов урока 68 (`select`, `poll`, `epoll`
//! и `kqueue`): буфер строк одного клиента, параметры и счётчики сервера.
//! Событийный сервер не имеет права ждать: дескриптор готов к чтению,
//! значит один `read` не заблокирует, а второй уже может. Поэтому клиент
//! читается ровно одним `read` за событие, а неполная строка ждёт в буфере
//! следующего. Книга читает строку через `rio_readlineb` и на неполной
//! строке зависает вместе со всеми клиентами; здесь этой дыры нет.
const std = @import("std");
const c = std.c;
const posix = std.posix;
const Io = std.Io;
const fdio = @import("../net/fdio.zig");
pub const Options = struct {
/// Сколько клиентов принять, прежде чем закрыть приём и дождаться,
/// пока уйдут принятые. 0 значит работать вечно.
max_clients: usize = 0,
/// Строка на каждую принятую строку, как `echoservers` книги.
log: ?*Io.Writer = null,
};
pub const Stats = struct {
/// Сколько соединений принято.
clients: usize = 0,
/// Сколько байтов вернулось эхом всем клиентам (`byte_cnt` книги).
bytes: usize = 0,
};
pub const Conn = struct {
fd: c.fd_t,
buf: [4096]u8 = undefined,
len: usize = 0,
/// Один `read` в хвост буфера, все целые строки уходят обратно.
/// `false`: клиент закрыл соединение или оно сломалось, пора закрывать.
/// Строка длиннее буфера уходит кусками, как в `echo.zig`.
pub fn onReadable(conn: *Conn, stats: *Stats, log: ?*Io.Writer) bool {
const n = c.read(conn.fd, conn.buf[conn.len..].ptr, conn.buf.len - conn.len);
if (n < 0) return posix.errno(n) == .INTR;
if (n == 0) return false;
conn.len += @intCast(n);
var start: usize = 0;
while (std.mem.indexOfScalarPos(u8, conn.buf[0..conn.len], start, '\n')) |nl| {
if (!conn.reply(conn.buf[start .. nl + 1], stats, log)) return false;
start = nl + 1;
}
if (start == 0 and conn.len == conn.buf.len) {
if (!conn.reply(conn.buf[0..conn.len], stats, log)) return false;
start = conn.len;
}
std.mem.copyForwards(u8, conn.buf[0 .. conn.len - start], conn.buf[start..conn.len]);
conn.len -= start;
return true;
}
// ponytail: ответ пишется блокирующим `writen`. Клиент, который шлёт и
// не читает, забьёт свой сокетный буфер и остановит весь цикл событий;
// лечится очередью на запись и ожиданием POLLOUT, эхо на строках её не просит.
fn reply(conn: *Conn, line: []const u8, stats: *Stats, log: ?*Io.Writer) bool {
fdio.writen(conn.fd, line) catch return false;
stats.bytes += line.len;
if (log) |w| {
w.print("Server received {d} ({d} total) bytes on fd {d}\n", .{ line.len, stats.bytes, conn.fd }) catch {};
w.flush() catch {};
}
return true;
}
};
Две детали. Строка длиннее буфера (4 КБ) уходит кусками: иначе клиент с бесконечной строкой заполнил бы буфер и застрял навсегда. И остаток после последнего \n сдвигается в начало через copyForwards: области источника и приёмника могут перекрываться, а copyForwards на это рассчитан.
Комментарий с меткой ponytail: честно говорит, где модель пока врёт: ответ уходит блокирующим writen. Клиент, который шлёт строки и не читает ответы, заполнит свой буфер сокета в ядре, write заснёт, и цикл встанет. Лечится это очередью на запись и ожиданием готовности к записи, это первое домашнее задание урока.
Пул книги на select
Книга хранит все автоматы в структуре pool, и наш Pool повторяет её поле в поле:
| Книга | У нас | Что это |
|---|---|---|
maxfd | maxfd | наибольший дескриптор в read_set, для первого аргумента select |
read_set | read_set | за кем следим: слушающий сокет и все клиенты |
ready_set | ready_set | копия read_set, которую select урезает до готовых |
nready | nready | сколько готово после последнего select |
maxi | maxi | последний занятый слот в clientfd |
clientfd[] | clientfd | дескрипторы клиентов по слотам, -1 это свободный слот |
clientrio[] | conns | состояние клиента: у книги буфер RIO, у нас Conn |
//! `echoservers` из главы 12: событийный эхо-сервер на `select`. Состояние
//! сервера живёт в `Pool`, как в книге: `read_set` это все дескрипторы,
//! которые мы слушаем, `ready_set` его копия, которую ядро урезает до
//! готовых, `clientfd[]` клиенты по слотам, `maxi` последний занятый слот.
//! Один поток, ни одного `fork`: каждый шаг цикла обслуживает тех, кто
//! готов, и возвращается к `select`.
const std = @import("std");
const c = std.c;
const posix = std.posix;
const socket = @import("../net/socket.zig");
const conn_mod = @import("conn.zig");
pub const Options = conn_mod.Options;
pub const Stats = conn_mod.Stats;
const Conn = conn_mod.Conn;
/// `FD_SETSIZE`: `select` знает дескрипторы только меньше этого числа.
/// Дескриптор 1024 в набор уже не положить, и это главная беда `select`.
pub const setsize = 1024;
/// `fd_set` как битовая карта. В glibc это `long[16]`, в macOS `int32_t[32]`:
/// на little-endian обе раскладки дают один и тот же бит `fd % 8` в байте
/// `fd / 8`, поэтому одно определение годится для обеих систем.
pub const FdSet = extern struct {
words: [setsize / 64]u64 = @splat(0),
pub fn set(s: *FdSet, fd: c.fd_t) void {
const i: usize = @intCast(fd);
s.words[i / 64] |= @as(u64, 1) << @intCast(i % 64);
}
pub fn clear(s: *FdSet, fd: c.fd_t) void {
const i: usize = @intCast(fd);
s.words[i / 64] &= ~(@as(u64, 1) << @intCast(i % 64));
}
pub fn isSet(s: *const FdSet, fd: c.fd_t) bool {
const i: usize = @intCast(fd);
return s.words[i / 64] & (@as(u64, 1) << @intCast(i % 64)) != 0;
}
};
// В std 0.16 нет ни `select`, ни `fd_set`: объявляем сами. На macOS x86-64
// символ назывался бы `select$1050`, на arm64 и в glibc это просто `select`.
extern "c" fn select(nfds: c_int, readfds: ?*FdSet, writefds: ?*FdSet, exceptfds: ?*FdSet, timeout: ?*c.timeval) c_int;
pub const Pool = struct {
/// Наибольший дескриптор в `read_set`: `select` просит `maxfd + 1`.
maxfd: c.fd_t,
read_set: FdSet = .{},
ready_set: FdSet = .{},
/// Сколько дескрипторов готово после последнего `select`.
nready: c_int = 0,
/// Последний занятый слот в `clientfd`, -1 значит клиентов нет.
maxi: isize = -1,
/// Слот -1 свободен. `conns[i]` это буфер строк клиента `clientfd[i]`.
clientfd: [setsize]c.fd_t = @splat(-1),
conns: [setsize]?*Conn = @splat(null),
/// Сколько клиентов сейчас подключено.
active: usize = 0,
/// `init_pool`: слушаем только `listenfd`.
pub fn init(listenfd: c.fd_t) Pool {
var p: Pool = .{ .maxfd = listenfd };
p.read_set.set(listenfd);
return p;
}
/// `add_client`: первый свободный слот, дескриптор в `read_set`.
pub fn addClient(p: *Pool, gpa: std.mem.Allocator, connfd: c.fd_t) !void {
p.nready -= 1;
if (connfd >= setsize) {
socket.close(connfd);
return error.TooManyClients;
}
for (&p.clientfd, &p.conns, 0..) |*slot, *conn, i| {
if (slot.* >= 0) continue;
conn.* = try gpa.create(Conn);
conn.*.?.* = .{ .fd = connfd };
slot.* = connfd;
p.read_set.set(connfd);
p.maxfd = @max(p.maxfd, connfd);
p.maxi = @max(p.maxi, @as(isize, @intCast(i)));
p.active += 1;
return;
}
socket.close(connfd);
return error.TooManyClients;
}
/// `check_clients`: обходим слоты до `maxi`, пока не кончатся готовые.
pub fn checkClients(p: *Pool, gpa: std.mem.Allocator, stats: *Stats, log: ?*std.Io.Writer) void {
var i: usize = 0;
while (@as(isize, @intCast(i)) <= p.maxi and p.nready > 0) : (i += 1) {
const fd = p.clientfd[i];
if (fd < 0 or !p.ready_set.isSet(fd)) continue;
p.nready -= 1;
if (p.conns[i].?.onReadable(stats, log)) continue;
socket.close(fd);
p.read_set.clear(fd);
gpa.destroy(p.conns[i].?);
p.conns[i] = null;
p.clientfd[i] = -1;
p.active -= 1;
}
}
pub fn deinit(p: *Pool, gpa: std.mem.Allocator) void {
for (p.clientfd[0..@intCast(p.maxi + 1)], p.conns[0..@intCast(p.maxi + 1)]) |fd, conn| {
if (fd < 0) continue;
socket.close(fd);
gpa.destroy(conn.?);
}
}
};
/// Цикл книги: копия набора, `select`, новый клиент, готовые клиенты.
/// `Pool` весит около 12 КБ, поэтому живёт в куче, а не на стеке потока.
pub fn serve(gpa: std.mem.Allocator, listenfd: c.fd_t, options: Options) !Stats {
const pool = try gpa.create(Pool);
defer gpa.destroy(pool);
pool.* = .init(listenfd);
defer pool.deinit(gpa);
var stats: Stats = .{};
var listening = true;
while (listening or pool.active > 0) {
pool.ready_set = pool.read_set;
pool.nready = select(pool.maxfd + 1, &pool.ready_set, null, null, null);
if (pool.nready < 0) {
if (posix.errno(pool.nready) == .INTR) continue;
return error.SelectFailed;
}
if (listening and pool.ready_set.isSet(listenfd)) {
const connfd = c.accept(listenfd, null, null);
if (connfd >= 0) {
stats.clients += 1;
pool.addClient(gpa, connfd) catch {};
// Лимит тестов: больше не принимаем, дорабатываем принятых.
if (options.max_clients != 0 and stats.clients >= options.max_clients) {
listening = false;
pool.read_set.clear(listenfd);
}
}
}
pool.checkClients(gpa, &stats, options.log);
}
return stats;
}
Пройдём один круг serve. read_set копируется в ready_set, копия уходит в select. Если среди готовых слушающий сокет, accept не заснёт, и новый клиент идёт в addClient: первый свободный слот, бит в read_set, maxfd и maxi при необходимости растут. Потом checkClients обходит слоты до maxi, у готовых зовёт onReadable, закрывшихся выбрасывает: close, бит из read_set долой, слот освобождается. Обход прекращается, как только готовые кончились (nready дошёл до нуля), до конца массива идти незачем.
Дескриптор 1024 и дальше в набор положить нельзя: бит за пределами карты это запись в чужую память. В C макрос FD_SET из glibc со включённым _FORTIFY_SOURCE в таком случае роняет программу, наш FdSet в отладочной сборке поймает выход за массив проверкой границ. Поэтому addClient такого клиента сразу закрывает. У maxfd, как и в книге, нет обратного хода: клиент с большим номером ушёл, а граница осталась. Это стоит лишних битов в каждом вызове, но не ошибка.
Pool весит около 12 КБ (два набора по 128 байт, массивы на 1024 слота) и создаётся в куче, а не на стеке функции.
Что даёт событийная модель и что отнимает
Книга подводит баланс, и с ним стоит согласиться:
- Процесс один, поток один. Отладчик видит всю программу, синхронизации нет, общее состояние общее по-настоящему:
stats.bytesсчитает байты всех клиентов одним сложением, без каналов через ядро. - Программа сама решает, кого обслужить первым. Можно, например, отвечать платным клиентам раньше бесплатных, в модели с процессами порядок решает планировщик ядра.
- Нет цены на соединение, кроме буфера и слота: ни
fork, ни переключения контекста между клиентами. - Код сложнее. Эхо на процессах это двадцать строк сессии, эхо на событиях это автомат, буферы и пул. У настоящего протокола состояний больше, и автомат растёт вместе с ним.
- Зернистость. Всё, что обработчик делает долго (тяжёлое вычисление, блокирующая запись, чтение файла с диска), останавливает всех. Это та же беда, что у долгого синхронного кода в Node.
- Ядро процессора одно. Процессы и потоки разойдутся по всем ядрам сами, событийному серверу для этого нужно несколько циклов.
poll: массив вместо битовой карты
poll решает две бытовые беды select. Вместо битовой карты у него массив структур:
struct pollfd { int fd; short events; short revents; };
int poll(struct pollfd *fds, nfds_t nfds, int timeout);
В events ты пишешь, чего ждёшь (POLLIN, POLLOUT), а ядро пишет ответ в отдельное поле revents. Запрос не портится, поэтому копия перед каждым вызовом не нужна. Номер дескриптора лежит в поле, а не в номере бита, поэтому предела FD_SETSIZE нет: сколько откроет процесс, столько и можно ждать. Отрицательный fd ядро пропускает и ставит ему revents = 0, это удобный способ временно выключить элемент, не сдвигая массив.
Кроме POLLIN ядро ставит в revents два флага без спроса: POLLHUP (соединение закрыто) и POLLERR (ошибка на сокете). И то и другое повод прочитать: read вернёт 0 или ошибку, и клиент уйдёт по обычной дороге.
Главного poll не меняет: каждый вызов отдаёт ядру весь массив, ядро проходит его целиком, программа после вызова проходит его ещё раз. Цена вызова растёт с числом подключённых, а не готовых.
//! Тот же событийный эхо-сервер на `poll`. Вместо битовой карты массив
//! `pollfd`: ни предела `FD_SETSIZE`, ни копии набора перед каждым
//! вызовом (ядро пишет ответ в `revents`, запрос в `events` остаётся).
//! Цена та же, что у `select`: каждый вызов отдаёт ядру весь массив, и
//! ядро проходит его целиком, даже если готов один клиент из тысячи.
const std = @import("std");
const c = std.c;
const posix = std.posix;
const socket = @import("../net/socket.zig");
const conn_mod = @import("conn.zig");
pub const Options = conn_mod.Options;
pub const Stats = conn_mod.Stats;
const Conn = conn_mod.Conn;
/// `fds[0]` это слушающий сокет, `fds[i]` для `i > 0` клиент `conns[i - 1]`.
pub fn serve(gpa: std.mem.Allocator, listenfd: c.fd_t, options: Options) !Stats {
var fds: std.ArrayList(c.pollfd) = .empty;
defer fds.deinit(gpa);
var conns: std.ArrayList(*Conn) = .empty;
defer {
for (conns.items) |conn| {
socket.close(conn.fd);
gpa.destroy(conn);
}
conns.deinit(gpa);
}
try fds.append(gpa, .{ .fd = listenfd, .events = c.POLL.IN, .revents = 0 });
var stats: Stats = .{};
while (fds.items[0].fd >= 0 or conns.items.len > 0) {
if (c.poll(fds.items.ptr, @intCast(fds.items.len), -1) < 0) {
if (posix.errno(-1) == .INTR) continue;
return error.PollFailed;
}
// Клиентов обходим с конца: `swapRemove` ставит на место ушедшего
// последнего, а его мы уже посмотрели.
var i = fds.items.len;
while (i > 1) {
i -= 1;
if (fds.items[i].revents & (c.POLL.IN | c.POLL.HUP | c.POLL.ERR) == 0) continue;
const conn = conns.items[i - 1];
if (conn.onReadable(&stats, options.log)) continue;
socket.close(conn.fd);
gpa.destroy(conn);
_ = fds.swapRemove(i);
_ = conns.swapRemove(i - 1);
}
// Новый клиент после обхода: его `revents` ещё не заполнены.
if (fds.items[0].revents & c.POLL.IN != 0) {
const connfd = c.accept(listenfd, null, null);
if (connfd >= 0) {
stats.clients += 1;
const conn = try gpa.create(Conn);
conn.* = .{ .fd = connfd };
try conns.append(gpa, conn);
try fds.append(gpa, .{ .fd = connfd, .events = c.POLL.IN, .revents = 0 });
// Отрицательный fd `poll` пропускает: так выключаем приём.
if (options.max_clients != 0 and stats.clients >= options.max_clients) fds.items[0].fd = -1;
}
}
}
return stats;
}
Здесь пул устроен не так, как у книги: без дыр. Ушедший клиент вынимается через swapRemove, на его место встаёт последний элемент массива. Поэтому обход идёт с конца: элемент, переставленный на место ушедшего, мы уже посмотрели. Новый клиент добавляется после обхода, у него ещё нет заполненного revents. Приём выключается отрицательным fd в элементе слушающего сокета: тот же приём, которым poll пропускает пустые слоты.
В практике урока ты соберёшь другой вариант, ближе к книге: фиксированные слоты, дыры с fd = -1, maxi, который съезжает влево. Полезно видеть оба.
epoll и kqueue: ядро помнит, за кем следить
Почему select и poll платят за каждого подключённого? Посмотри, что делает ядро на каждом вызове. Копирует набор из памяти программы. Для каждого дескриптора спрашивает сокет, готов ли он, и вешает процесс в очередь ожидания этого сокета. Засыпает. Проснувшись, снова проходит все дескрипторы, снимает процесс со всех очередей, копирует ответ обратно. Программа проходит набор ещё раз, чтобы найти готовых. При тысяче соединений, из которых говорят десять, девяносто девять процентов этой работы впустую, и повторяется она на каждом круге цикла.
Выход в том, чтобы разделить два действия, которые select делает вместе: «вот за кем следить» и «кто готов». Первое меняется редко (клиент пришёл, клиент ушёл), второе нужно на каждом круге. Linux для этого даёт epoll, BSD и macOS дают kqueue.
epoll это три вызова:
epoll_create1(flags)создаёт объект epoll и возвращает его дескриптор.epoll_ctl(epfd, EPOLL_CTL_ADD, fd, &event)один раз добавляетfdв набор интереса.eventговорит, чего ждать (EPOLLIN), и несёт полеdata, которое ядро вернёт вместе с событием: туда кладут номер дескриптора или указатель на состояние клиента.epoll_wait(epfd, events, maxevents, timeout)ждёт и возвращает только готовые события.
Набор интереса ядро хранит у себя, а сокет, в котором появились данные, сам добавляет себя в список готовых. epoll_wait просто отдаёт этот список. Цена вызова зависит от числа готовых, а не подключённых.
kqueue устроен похоже, но экономнее на вызовах. kqueue() создаёт очередь, а единственный вызов kevent(kq, changes, nchanges, events, nevents, timeout) и регистрирует изменения, и ждёт событий, можно за один раз. Событие описывается парой из дескриптора и фильтра: EVFILT_READ, EVFILT_WRITE, и дальше то, чего у epoll нет вовсе, EVFILT_PROC (процесс вышел), EVFILT_SIGNAL, EVFILT_TIMER. В поле data для сокета ядро кладёт, сколько байтов можно прочитать.
У обоих есть два режима уведомления. По умолчанию по уровню: пока в сокете есть непрочитанные данные, каждое ожидание будет возвращать его снова. Так же ведут себя select и poll, и на этом держится наш Conn с одним read на событие. Флаг EPOLLET у epoll и EV_CLEAR у kqueue включают срабатывание по фронту: событие приходит один раз, когда пришли новые данные. Тогда читать надо до EAGAIN на неблокирующем сокете, иначе остаток пролежит до следующего пакета. Это на одно ожидание экономнее и на порядок проще ошибиться; у нас режим по уровню.
Снимать дескриптор из набора перед close не нужно. В kqueue события дескриптора исчезают при его закрытии. В epoll с оговоркой: ядро помнит не номер дескриптора, а запись в таблице открытых файлов, и убирает её из набора, когда закрыты все дескрипторы, которые на неё смотрят. Если сокет успел попасть в ребёнка через fork или получить копию через dup, close одной копии не снимет его с учёта. Наш событийный сервер не делает ни того ни другого.
Другое дело дескриптор, который остаётся открытым, а читать его мы перестали. Так бывает со слушающим сокетом, когда сервер принял max_clients соединений и закрыл приём. Пока в его очереди ждёт хоть одно непринятое соединение, по уровню он готов, и wait возвращал бы его на каждом круге: занятый цикл, пока не уйдут принятые клиенты. Поэтому у Poller есть remove (EPOLL_CTL_DEL или EV_DELETE), и сервер зовёт его в тот момент, когда закрывает приём. У select то же самое делает read_set.clear, у poll запись fd = -1.
Две системы отличаются типами и константами, а сервер над ними одинаковый. Поэтому разницу прячем в Poller, а какой механизм взять, решает компилятор по целевой ОС:
//! Событийный эхо-сервер дальше книги: `epoll` на Linux, `kqueue` на macOS.
//! Интерес регистрируется в ядре один раз (`epoll_ctl`, `kevent` со списком
//! изменений), а ожидание возвращает только готовые дескрипторы. Цена
//! вызова зависит от числа готовых, а не от числа подключённых: на тысяче
//! соединений, из которых говорят несколько, это и есть выигрыш.
//! `Poller` прячет разницу двух систем; сервер над ним один.
const std = @import("std");
const builtin = @import("builtin");
const c = std.c;
const posix = std.posix;
const linux = std.os.linux;
const socket = @import("../net/socket.zig");
const conn_mod = @import("conn.zig");
pub const Options = conn_mod.Options;
pub const Stats = conn_mod.Stats;
const Conn = conn_mod.Conn;
pub const Backend = enum { epoll, kqueue };
pub const backend: Backend = switch (builtin.os.tag) {
.linux => .epoll,
.macos, .freebsd, .netbsd, .openbsd, .dragonfly => .kqueue,
else => @compileError("нужен epoll или kqueue"),
};
pub const Poller = struct {
fd: c.fd_t,
pub const Event = switch (backend) {
.epoll => linux.epoll_event,
.kqueue => c.Kevent,
};
pub fn init() !Poller {
const fd = switch (backend) {
.epoll => c.epoll_create1(linux.EPOLL.CLOEXEC),
.kqueue => c.kqueue(),
};
if (fd < 0) return error.PollerFailed;
return .{ .fd = fd };
}
pub fn deinit(p: *Poller) void {
socket.close(p.fd);
}
/// Следить за чтением `fd`. Снимать закрытый не нужно: `close` убирает
/// дескриптор и из `epoll`, и из `kqueue` сам.
pub fn add(p: *Poller, fd: c.fd_t) !void {
const rc = switch (backend) {
.epoll => blk: {
var ev: linux.epoll_event = .{ .events = linux.EPOLL.IN, .data = .{ .fd = fd } };
break :blk c.epoll_ctl(p.fd, linux.EPOLL.CTL_ADD, fd, &ev);
},
.kqueue => blk: {
const change: c.Kevent = .{
.ident = @intCast(fd),
.filter = c.EVFILT.READ,
.flags = c.EV.ADD,
.fflags = 0,
.data = 0,
.udata = 0,
};
break :blk c.kevent(p.fd, (&change)[0..1], 1, &[0]c.Kevent{}, 0, null);
},
};
if (rc < 0) return error.PollerFailed;
}
/// Перестать следить за `fd`, который остаётся открытым. Иначе готовый
/// дескриптор, который мы больше не читаем, будил бы `wait` на каждом
/// круге: оба механизма здесь работают по уровню.
pub fn remove(p: *Poller, fd: c.fd_t) !void {
const rc = switch (backend) {
.epoll => c.epoll_ctl(p.fd, linux.EPOLL.CTL_DEL, fd, null),
.kqueue => blk: {
const change: c.Kevent = .{
.ident = @intCast(fd),
.filter = c.EVFILT.READ,
.flags = c.EV.DELETE,
.fflags = 0,
.data = 0,
.udata = 0,
};
break :blk c.kevent(p.fd, (&change)[0..1], 1, &[0]c.Kevent{}, 0, null);
},
};
if (rc < 0) return error.PollerFailed;
}
/// Ждать, пока что-нибудь не станет готово; вернуть готовые события.
pub fn wait(p: *Poller, events: []Event) ![]Event {
while (true) {
const n = switch (backend) {
.epoll => c.epoll_wait(p.fd, events.ptr, @intCast(events.len), -1),
.kqueue => c.kevent(p.fd, &[0]c.Kevent{}, 0, events.ptr, @intCast(events.len), null),
};
if (n >= 0) return events[0..@intCast(n)];
if (posix.errno(n) != .INTR) return error.PollerFailed;
}
}
pub fn eventFd(ev: Event) c.fd_t {
return switch (backend) {
.epoll => ev.data.fd,
.kqueue => @intCast(ev.ident),
};
}
};
pub fn serve(gpa: std.mem.Allocator, listenfd: c.fd_t, options: Options) !Stats {
var poller: Poller = try .init();
defer poller.deinit();
try poller.add(listenfd);
var conns: std.AutoHashMapUnmanaged(c.fd_t, *Conn) = .empty;
defer {
var it = conns.valueIterator();
while (it.next()) |conn| {
socket.close(conn.*.fd);
gpa.destroy(conn.*);
}
conns.deinit(gpa);
}
var stats: Stats = .{};
var listening = true;
var events: [256]Poller.Event = undefined;
while (listening or conns.count() > 0) {
for (try poller.wait(&events)) |ev| {
const fd = Poller.eventFd(ev);
if (fd == listenfd) {
if (!listening) continue;
const connfd = c.accept(listenfd, null, null);
if (connfd < 0) continue;
stats.clients += 1;
const conn = try gpa.create(Conn);
conn.* = .{ .fd = connfd };
try conns.put(gpa, connfd, conn);
try poller.add(connfd);
if (options.max_clients != 0 and stats.clients >= options.max_clients) {
// Как `read_set.clear` у select и `fd = -1` у poll: очередь
// на приём больше не наша забота.
try poller.remove(listenfd);
listening = false;
}
continue;
}
const conn = conns.get(fd) orelse continue;
if (conn.onReadable(&stats, options.log)) continue;
_ = conns.remove(fd);
socket.close(fd);
gpa.destroy(conn);
}
}
return stats;
}
backend это константа времени компиляции, и switch (backend) оставляет в бинарнике только одну ветку: на Linux код про kqueue даже не анализируется, и наоборот. Тип Event тоже выбирается при компиляции: linux.epoll_event или c.Kevent. Слотов и maxi здесь нет: wait отдаёт сами готовые события, и клиента по номеру дескриптора находит хеш-таблица.
Этот приём не наш. libuv под Node, сетевой планировщик Go, mio под Tokio устроены так же: одна обёртка, под ней epoll, kqueue или механизм Windows. В курсе Rust ты собирал такой реактор руками.
Прогон на macOS, --epoll там называется kqueue:
$ tiny echoservers --epoll 15242
echoservers (kqueue): listening on port 15242
Server received 6 (6 total) bytes on fd 6
Server received 6 (12 total) bytes on fd 6
Server received 13 (25 total) bytes on fd 6
$ printf 'hello\nworld\n' | tiny echoclient 127.0.0.1 15242
hello
world
$ printf 'второй\n' | tiny echoclient 127.0.0.1 15242
второй
Первые две строки прислал один клиент, третью второй. У обоих fd 6: первый успел закрыться, и ядро выдало тот же номер. У того же сервера на select клиент получает fd 5 (прогон ниже, в шаге проекта): здесь на один дескриптор больше, его занимает сама очередь kqueue.
Замер: тысяча соединений
Теорию про «O от числа подключённых» проверим числами. conc-bench запускает сервер нужной модели в ребёнке (fork после openListenfd), а в родителе открывает N клиентских соединений. Каждый клиент шлёт строку в 64 байта, ждёт эхо целиком и шлёт следующую. Задержка это время одной такой пары, пропускная способность это все строки, делённые на всё время. Клиентов ведёт один поток на том же Poller, так что клиентская сторона одинаковая для всех моделей.
Флаг --active делает картину, которую называют C10K: соединений тысяча, а говорят из них немногие, как у чата или сервера уведомлений, где большинство клиентов подключены и молчат.
//! `conc-bench`: одна нагрузка на все модели эхо-сервера урока 68. Сервер
//! работает в отдельном процессе (`fork` после `openListenfd`), клиенты в
//! этом: у сервера свои дескрипторы, и `select` на тысяче клиентов
//! укладывается в `FD_SETSIZE`.
//!
//! Нагрузка: `clients` соединений открываются заранее, потом каждое шлёт
//! `lines` строк по `line_len` байт, следующую только после эха предыдущей.
//! Задержка это время от отправки строки до полного эха, пропускная
//! способность это все строки, делённые на время от первой отправки до
//! последнего эха. Клиенты ведёт один поток на `epoll` или `kqueue`, так что
//! клиентская сторона одинаковая для всех моделей.
//!
//! `active` меньше `clients` даёт картину C10K: соединений много, говорят
//! немногие. Тут и видна разница между `select`/`poll`, которые на каждом
//! вызове проходят все дескрипторы, и `epoll`/`kqueue`, которые отдают
//! только готовые.
const std = @import("std");
const builtin = @import("builtin");
const c = std.c;
const posix = std.posix;
const Io = std.Io;
const socket = @import("../net/socket.zig");
const fdio = @import("../net/fdio.zig");
const echo_procs = @import("echo_procs.zig");
const echo_select = @import("echo_select.zig");
const echo_poll = @import("echo_poll.zig");
const echo_epoll = @import("echo_epoll.zig");
pub const Model = enum {
procs,
select,
poll,
/// `epoll` на Linux, `kqueue` на macOS.
epoll,
pub fn label(m: Model) []const u8 {
return switch (m) {
.epoll => @tagName(echo_epoll.backend),
else => @tagName(m),
};
}
};
pub const Config = struct {
clients: usize,
/// Сколько из них шлют строки; 0 значит все.
active: usize = 0,
lines: usize = 100,
line_len: usize = 64,
};
pub const Result = struct {
model: Model,
clients: usize,
active: usize,
lines: usize,
seconds: f64,
/// Строк в секунду по всем клиентам.
throughput: f64,
p50_us: f64,
p99_us: f64,
/// Худшая задержка: у пула это первая строка клиента, который ждал
/// свободного воркера. В p99 она не попадает, таких строк мало.
max_us: f64,
};
/// Поднимает мягкий предел дескрипторов до жёсткого (не больше `want`):
/// на macOS по умолчанию 256, а тысяче клиентов нужна тысяча сокетов.
pub fn raiseFdLimit(want: u64) u64 {
var lim: c.rlimit = undefined;
if (c.getrlimit(.NOFILE, &lim) != 0) return 0;
const target = @min(want, lim.max);
if (lim.cur < target) {
lim.cur = target;
_ = c.setrlimit(.NOFILE, &lim);
_ = c.getrlimit(.NOFILE, &lim);
}
return lim.cur;
}
pub fn run(gpa: std.mem.Allocator, io: Io, model: Model, cfg: Config) !Result {
const listenfd = try socket.openListenfd(0);
const port = socket.localPort(listenfd).?;
const pid = c.fork();
if (pid < 0) return error.ForkFailed;
if (pid == 0) {
runServer(model, listenfd);
c._exit(0);
}
socket.close(listenfd);
defer {
_ = c.kill(pid, .KILL);
_ = c.waitpid(pid, null, 0);
}
return drive(gpa, io, model, port, cfg);
}
fn runServer(model: Model, listenfd: c.fd_t) void {
// Свежий процесс: свой аллокатор libc, родительский через `fork` не живёт.
const gpa = std.heap.c_allocator;
var discard: Io.Writer.Discarding = .init(&.{});
const opts: echo_select.Options = .{};
_ = switch (model) {
.procs => echo_procs.serve(listenfd, 0, &discard.writer),
.select => if (echo_select.serve(gpa, listenfd, opts)) |_| {} else |err| err,
.poll => if (echo_poll.serve(gpa, listenfd, opts)) |_| {} else |err| err,
.epoll => if (echo_epoll.serve(gpa, listenfd, opts)) |_| {} else |err| err,
} catch |err| std.debug.print("bench server {t}: {t}\n", .{ model, err });
}
const Client = struct {
fd: c.fd_t,
sent: usize = 0,
got: usize = 0,
started_ns: i96 = 0,
};
fn now(io: Io) i96 {
return Io.Timestamp.now(io, .awake).nanoseconds;
}
fn drive(gpa: std.mem.Allocator, io: Io, model: Model, port: u16, cfg: Config) !Result {
const clients = try gpa.alloc(Client, cfg.clients);
defer gpa.free(clients);
var opened: usize = 0;
defer for (clients[0..opened]) |cl| if (cl.fd >= 0) socket.close(cl.fd);
var max_fd: usize = 0;
for (clients) |*cl| {
cl.* = .{ .fd = try socket.openClientfd("127.0.0.1", port) };
opened += 1;
max_fd = @max(max_fd, @as(usize, @intCast(cl.fd)));
}
const by_fd = try gpa.alloc(u32, max_fd + 1);
defer gpa.free(by_fd);
var poller: echo_epoll.Poller = try .init();
defer poller.deinit();
for (clients, 0..) |cl, i| {
by_fd[@intCast(cl.fd)] = @intCast(i);
try poller.add(cl.fd);
}
const active = if (cfg.active == 0) cfg.clients else @min(cfg.active, cfg.clients);
const latencies = try gpa.alloc(u64, active * cfg.lines);
defer gpa.free(latencies);
var recorded: usize = 0;
const line = try gpa.alloc(u8, cfg.line_len);
defer gpa.free(line);
@memset(line, 'x');
line[line.len - 1] = '\n';
const t0 = now(io);
for (clients[0..active]) |*cl| try send(io, cl, line);
var done: usize = 0;
var events: [256]echo_epoll.Poller.Event = undefined;
var sink: [4096]u8 = undefined;
while (done < active) {
for (try poller.wait(&events)) |ev| {
const cl = &clients[by_fd[@intCast(echo_epoll.Poller.eventFd(ev))]];
const n = c.read(cl.fd, &sink, sink.len);
if (n <= 0) return error.ServerClosed;
cl.got += @intCast(n);
if (cl.got < cfg.line_len) continue;
// Эхо строки пришло целиком: записываем задержку, шлём следующую.
cl.got -= cfg.line_len;
latencies[recorded] = @intCast(now(io) - cl.started_ns);
recorded += 1;
if (cl.sent < cfg.lines) {
try send(io, cl, line);
} else {
// Отработавший клиент уходит сразу: воркер пула свободен
// только после конца сессии, и очередь ждёт именно этого.
socket.close(cl.fd);
cl.fd = -1;
done += 1;
}
}
}
const seconds = @as(f64, @floatFromInt(now(io) - t0)) / 1e9;
std.mem.sort(u64, latencies[0..recorded], {}, std.sort.asc(u64));
return .{
.model = model,
.clients = cfg.clients,
.active = active,
.lines = cfg.lines,
.seconds = seconds,
.throughput = @as(f64, @floatFromInt(recorded)) / seconds,
.p50_us = percentile(latencies[0..recorded], 50),
.p99_us = percentile(latencies[0..recorded], 99),
.max_us = if (recorded == 0) 0 else @as(f64, @floatFromInt(latencies[recorded - 1])) / 1000.0,
};
}
fn send(io: Io, cl: *Client, line: []const u8) !void {
cl.started_ns = now(io);
cl.sent += 1;
try fdio.writen(cl.fd, line);
}
pub fn percentile(sorted: []const u64, p: usize) f64 {
if (sorted.len == 0) return 0;
const idx = @min(sorted.len - 1, sorted.len * p / 100);
return @as(f64, @floatFromInt(sorted[idx])) / 1000.0;
}
pub fn writeTable(out: *Io.Writer, results: []const Result) !void {
try out.print("{s:<8} {s:>7} {s:>6} {s:>6} {s:>9} {s:>12} {s:>9} {s:>9} {s:>9}\n", .{ "model", "clients", "active", "lines", "seconds", "lines/s", "p50 us", "p99 us", "max us" });
for (results) |r| {
try out.print("{s:<8} {d:>7} {d:>6} {d:>6} {d:>9.3} {d:>12.0} {d:>9.1} {d:>9.1} {d:>9.0}\n", .{ r.model.label(), r.clients, r.active, r.lines, r.seconds, r.throughput, r.p50_us, r.p99_us, r.max_us });
}
}
pub fn writeJson(out: *Io.Writer, results: []const Result) !void {
const cpus = std.Thread.getCpuCount() catch 0;
try out.print("{{\"os\":\"{t}\",\"arch\":\"{t}\",\"cpus\":{d},\"results\":[", .{ builtin.os.tag, builtin.cpu.arch, cpus });
for (results, 0..) |r, i| {
if (i > 0) try out.writeAll(",");
try out.print("{{\"model\":\"{s}\",\"clients\":{d},\"active\":{d},\"lines\":{d},\"seconds\":{d:.4},\"throughput\":{d:.0},\"p50_us\":{d:.1},\"p99_us\":{d:.1},\"max_us\":{d:.0}}}", .{ r.model.label(), r.clients, r.active, r.lines, r.seconds, r.throughput, r.p50_us, r.p99_us, r.max_us });
}
try out.writeAll("]}\n");
}
raiseFdLimit поднимает мягкий предел дескрипторов до жёсткого. На macOS мягкий предел по умолчанию 256, а тысяче клиентов нужна тысяча сокетов в родителе и столько же в сервере. Ребёнок после fork берёт аллокатор libc: состояние аллокатора родителя через fork переносится копией, и у каждого процесса оно своё.
Замеры сняты на одной машине: Apple M4 Max, 16 ядер (12 производительных и 4 энергоэффективных), 64 ГБ, macOS 26.6.2, Zig 0.16.0, сборка -Doptimize=ReleaseFast. Строки Linux сняты там же, в контейнере runner-zig (linux/arm64, нативно, OrbStack, ядро 7.0.14, 16 виртуальных ядер). Машину в это время грузили посторонние процессы (средняя загрузка, load average, от 13 до 30 на 16 ядрах), поэтому абсолютные числа между прогонами гуляют в полтора раза. Порядок моделей при повторах сохранялся, на него и смотрим. Строки threads и pool из той же таблицы появятся в уроках 69 и 71.
Все клиенты говорят, по 300 строк (--clients 100,1000 --lines 300), macOS:
model clients active lines seconds lines/s p50 us p99 us max us
procs 100 100 300 0.383 78251 993.6 3467.0 70759
select 100 100 300 0.336 89217 965.7 2305.1 3491
poll 100 100 300 0.447 67128 1119.3 4288.0 36267
kqueue 100 100 300 0.448 66896 1125.4 7152.1 25547
procs 1000 1000 300 3.896 76994 11628.5 32610.3 119199
select 1000 1000 300 3.181 94310 10669.1 12995.9 22443
poll 1000 1000 300 3.826 78402 11349.9 36275.9 105963
kqueue 1000 1000 300 3.287 91274 10789.7 18653.7 27138
Linux:
model clients active lines seconds lines/s p50 us p99 us max us
procs 100 100 300 0.242 124066 530.6 3506.0 6860
select 100 100 300 0.182 164492 488.3 1690.8 2268
poll 100 100 300 0.191 157218 513.9 1866.9 2576
epoll 100 100 300 0.161 186115 453.2 1486.9 2553
procs 1000 1000 300 4.310 69601 13772.3 29659.3 91487
select 1000 1000 300 1.634 183595 4033.7 11582.8 745479
poll 1000 1000 300 1.663 180385 4030.1 12559.6 740850
epoll 1000 1000 300 2.037 147293 5597.1 21041.9 75457
Когда говорят все, select и poll не проигрывают epoll. На каждом вызове готовых много, и проход по всему набору окупается: работы ровно столько, сколько готовых. На Linux при тысяче клиентов событийные модели обгоняют процессы больше чем вдвое: на каждый круг строк процессы платят тысячей переключений контекста. На macOS все модели в пределах четверти друг от друга. Задержка у всех растёт с числом клиентов: это очередь из тысячи строк к одному серверу.
Теперь C10K: тысяча соединений, говорят десять, каждый по 10 000 строк (--clients 1000 --active 10 --lines 10000), macOS:
model clients active lines seconds lines/s p50 us p99 us max us
procs 1000 10 10000 1.218 82081 106.0 379.8 2995
select 1000 10 10000 4.050 24694 262.6 1717.8 32002
poll 1000 10 10000 7.293 13711 504.6 3345.3 22385
kqueue 1000 10 10000 1.315 76018 108.5 496.3 3484
Linux:
model clients active lines seconds lines/s p50 us p99 us max us
procs 1000 10 10000 1.865 53611 48.3 1806.3 34143
select 1000 10 10000 3.131 31936 202.3 1880.9 24530
poll 1000 10 10000 2.874 34798 203.0 1606.2 26922
epoll 1000 10 10000 0.554 180648 46.5 204.5 4253
Вот где epoll и kqueue платят за себя. select и poll на каждом вызове отдают ядру тысячу дескрипторов ради десяти готовых и проигрывают: на Linux в пять с лишним раз, на macOS в три (select) и в пять с половиной (poll, который на macOS ещё и почти вдвое медленнее select). p99 у epoll в девять раз ниже, чем у select. А процессы на молчащих соединениях почти ничего не стоят: девятьсот девяносто детей спят в read, и ядро их не будит. На macOS процессы в этом сценарии даже чуть впереди kqueue, на Linux втрое позади epoll. Их цена здесь не время, а тысяча процессов в памяти.
Прогони conc-bench у себя. Числа будут другими, а порядок должен совпасть. Если не совпал, это интереснее любой таблицы из урока.
Шаг проекта: our-tiny учится конкурентности
Эталон our-tiny растёт с урока 63, и сегодня у него появляется подпакет src/conc/. Файлы шага:
| Файл | Что в нём |
|---|---|
src/conc/echo_procs.zig | echoserverp: процесс на соединение, SIGCHLD (листинг выше) |
src/conc/conn.zig | состояние клиента событийного сервера, Options, Stats (листинг выше) |
src/conc/echo_select.zig | пул книги на select, свой FdSet (листинг выше) |
src/conc/echo_poll.zig | тот же сервер на poll (листинг выше) |
src/conc/echo_epoll.zig | Poller над epoll и kqueue, сервер над ним (листинг выше) |
src/conc/bench.zig | conc-bench (листинг выше) |
src/main.zig | подкоманды echoserverp, echoservers, conc-bench |
src/root.zig | пространство имён conc |
build.zig | номер шага 68 |
tests/step_68.zig | тесты шага |
main.zig: таблица подкоманд
До сих пор main был цепочкой if/else if по имени подкоманды. К концу блока подкоманд станет пятнадцать, и цепочка превратилась бы в страницу. Поэтому сегодня main становится таблицей: кортеж пар «имя, функция», по которому inline for проходит при компиляции. Каждая подкоманда получила свою функцию cmdXxx с одинаковой сигнатурой, а всё, что им нужно (аллокаторы, Io, писатели stdout и stderr), лежит в одной структуре Ctx. Разбор числа и открытие слушающего сокета с строкой в логе переехали в помощники parsePort, parseCount, parseList и listenOn. Следующие уроки будут добавлять по строке в таблицу и по функции.
//! Точка входа: подкоманды глав 11 и 12. Вся работа живёт в модуле `tiny`,
//! здесь только разбор командной строки.
const std = @import("std");
const tiny = @import("tiny");
const conc = tiny.conc;
const usage =
\\tiny, веб-сервер из главы 11 и его конкурентные версии из главы 12
\\
\\Использование:
\\ tiny tiny <port> [root] веб-сервер, корень статики и cgi-bin (по умолчанию .)
\\ tiny hostinfo [--std] <name> все адреса имени через getaddrinfo или std.Io.net
\\ tiny echoserver [--std] <port> эхо-сервер, одно соединение за раз
\\ tiny echoclient <host> <port> эхо-клиент: строки из stdin туда и обратно
\\ tiny echoserverp <port> эхо: процесс на соединение
\\ tiny echoservers --select|--poll|--epoll <port> эхо: события (--epoll это kqueue на macOS)
\\ tiny conc-bench [--clients 100,1000] [--active A] [--lines K] [--models procs,select,poll,epoll] [--json]
\\
;
const Ctx = struct {
gpa: std.mem.Allocator,
arena: std.mem.Allocator,
io: std.Io,
out: *std.Io.Writer,
log: *std.Io.Writer,
};
pub fn main(init: std.process.Init) !void {
const arena = init.arena.allocator();
const io = init.io;
const args = try init.minimal.args.toSlice(arena);
var out_buf: [4096]u8 = undefined;
var stdout = std.Io.File.stdout().writerStreaming(io, &out_buf);
const out = &stdout.interface;
var err_buf: [4096]u8 = undefined;
var stderr = std.Io.File.stderr().writerStreaming(io, &err_buf);
const log = &stderr.interface;
const ctx: Ctx = .{ .gpa = init.gpa, .arena = arena, .io = io, .out = out, .log = log };
if (args.len < 2) return fail(out, usage);
const command = args[1];
const rest = args[2..];
const commands = .{
.{ "tiny", cmdTiny },
.{ "hostinfo", cmdHostinfo },
.{ "echoserver", cmdEchoserver },
.{ "echoclient", cmdEchoclient },
.{ "echoserverp", cmdEchoserverp },
.{ "echoservers", cmdEchoservers },
.{ "conc-bench", cmdConcBench },
};
inline for (commands) |entry| {
if (std.mem.eql(u8, command, entry[0])) return entry[1](ctx, rest);
}
return fail(out, usage);
}
fn cmdTiny(ctx: Ctx, rest: []const [:0]const u8) !void {
if (rest.len < 1) return fail(ctx.out, usage);
const port = try parsePort(ctx, rest[0]);
const root = if (rest.len > 1) rest[1] else ".";
var server: tiny.Server = try .init(ctx.gpa, ctx.io, .{ .port = port, .root = root });
defer server.deinit();
try ctx.log.print("tiny: listening on port {d}, root {s}\n", .{ server.port, root });
try ctx.log.flush();
try server.serveForever();
}
fn cmdHostinfo(ctx: Ctx, rest: []const [:0]const u8) !void {
const use_std = rest.len > 0 and std.mem.eql(u8, rest[0], "--std");
const names = if (use_std) rest[1..] else rest;
if (names.len != 1) return fail(ctx.out, usage);
if (use_std) {
try tiny.hostinfo.lookupStd(ctx.io, names[0], ctx.out);
} else {
tiny.hostinfo.hostinfo(names[0], ctx.out) catch |err| switch (err) {
error.LookupFailed => {
try ctx.out.flush();
std.process.exit(1);
},
else => return err,
};
}
try ctx.out.flush();
}
fn cmdEchoserver(ctx: Ctx, rest: []const [:0]const u8) !void {
const use_std = rest.len > 0 and std.mem.eql(u8, rest[0], "--std");
const ports = if (use_std) rest[1..] else rest;
if (ports.len != 1) return fail(ctx.out, usage);
const port = try parsePort(ctx, ports[0]);
if (use_std) {
var server = try tiny.echo_std.listen(ctx.io, port);
defer server.deinit(ctx.io);
try ctx.log.print("echoserver (std): listening on {f}\n", .{server.socket.address});
try ctx.log.flush();
try tiny.echo_std.serve(ctx.io, &server, 0, ctx.log);
} else {
const listenfd = try listenOn(ctx, "echoserver", port);
defer tiny.socket.close(listenfd);
try tiny.echo.serve(listenfd, 0, ctx.log);
}
}
fn cmdEchoclient(ctx: Ctx, rest: []const [:0]const u8) !void {
if (rest.len != 2) return fail(ctx.out, usage);
const port = try parsePort(ctx, rest[1]);
var in_buf: [8192]u8 = undefined;
var stdin = std.Io.File.stdin().readerStreaming(ctx.io, &in_buf);
try tiny.echo.client(rest[0], port, &stdin.interface, ctx.out);
}
fn cmdEchoserverp(ctx: Ctx, rest: []const [:0]const u8) !void {
if (rest.len != 1) return fail(ctx.out, usage);
const listenfd = try listenOn(ctx, "echoserverp", try parsePort(ctx, rest[0]));
defer tiny.socket.close(listenfd);
try conc.echo_procs.serve(listenfd, 0, ctx.log);
}
fn cmdEchoservers(ctx: Ctx, rest: []const [:0]const u8) !void {
if (rest.len != 2) return fail(ctx.out, usage);
const port = try parsePort(ctx, rest[1]);
const flag = rest[0];
if (!std.mem.startsWith(u8, flag, "--") or flag.len < 3) return fail(ctx.out, usage);
const kind = if (std.mem.eql(u8, flag, "--epoll")) conc.bench.Model.epoll.label() else flag[2..];
const listenfd = try listenOn(ctx, try std.fmt.allocPrint(ctx.arena, "echoservers ({s})", .{kind}), port);
defer tiny.socket.close(listenfd);
const opts: conc.conn.Options = .{ .log = ctx.log };
_ = if (std.mem.eql(u8, flag, "--select"))
try conc.echo_select.serve(ctx.gpa, listenfd, opts)
else if (std.mem.eql(u8, flag, "--poll"))
try conc.echo_poll.serve(ctx.gpa, listenfd, opts)
else if (std.mem.eql(u8, flag, "--epoll"))
try conc.echo_epoll.serve(ctx.gpa, listenfd, opts)
else
return fail(ctx.out, usage);
}
fn cmdConcBench(ctx: Ctx, rest: []const [:0]const u8) !void {
var clients: []const usize = &.{ 100, 1000 };
var models: []const conc.bench.Model = &.{ .procs, .select, .poll, .epoll };
var cfg: conc.bench.Config = .{ .clients = 0 };
var json = false;
var i: usize = 0;
while (i < rest.len) : (i += 1) {
const flag = rest[i];
if (std.mem.eql(u8, flag, "--json")) {
json = true;
continue;
}
if (i + 1 >= rest.len) return fail(ctx.out, usage);
i += 1;
const value = rest[i];
if (std.mem.eql(u8, flag, "--clients")) {
clients = try parseList(ctx, value);
} else if (std.mem.eql(u8, flag, "--lines")) {
cfg.lines = try parseCount(ctx, value);
} else if (std.mem.eql(u8, flag, "--active")) {
cfg.active = try parseCount(ctx, value);
} else if (std.mem.eql(u8, flag, "--models")) {
var list: std.ArrayList(conc.bench.Model) = .empty;
var it = std.mem.tokenizeScalar(u8, value, ',');
while (it.next()) |name| {
const m = std.meta.stringToEnum(conc.bench.Model, name) orelse if (std.mem.eql(u8, name, "kqueue")) conc.bench.Model.epoll else return fail(ctx.out, usage);
try list.append(ctx.arena, m);
}
models = list.items;
} else return fail(ctx.out, usage);
}
const max_clients = std.mem.max(usize, clients);
const limit = conc.bench.raiseFdLimit(max_clients * 2 + 64);
if (limit < max_clients * 2 + 32) {
try ctx.log.print("conc-bench: предел дескрипторов {d}, на {d} клиентов нужно больше (ulimit -n)\n", .{ limit, max_clients });
try ctx.log.flush();
std.process.exit(1);
}
var results: std.ArrayList(conc.bench.Result) = .empty;
for (clients) |n| for (models) |m| {
cfg.clients = n;
const r = try conc.bench.run(ctx.gpa, ctx.io, m, cfg);
try results.append(ctx.arena, r);
if (!json) {
try conc.bench.writeTable(ctx.log, results.items[results.items.len - 1 ..]);
try ctx.log.flush();
}
};
if (json) try conc.bench.writeJson(ctx.out, results.items) else try conc.bench.writeTable(ctx.out, results.items);
try ctx.out.flush();
}
fn listenOn(ctx: Ctx, name: []const u8, port: u16) !std.c.fd_t {
const listenfd = try tiny.socket.openListenfd(port);
try ctx.log.print("{s}: listening on port {d}\n", .{ name, tiny.socket.localPort(listenfd) orelse port });
try ctx.log.flush();
return listenfd;
}
fn parsePort(ctx: Ctx, text: []const u8) !u16 {
return std.fmt.parseInt(u16, text, 10) catch {
try fail(ctx.out, "tiny: порт это число от 0 до 65535\n");
unreachable;
};
}
fn parseCount(ctx: Ctx, text: []const u8) !usize {
return std.fmt.parseInt(usize, text, 10) catch {
try fail(ctx.out, "tiny: ожидалось целое число\n");
unreachable;
};
}
fn parseList(ctx: Ctx, text: []const u8) ![]const usize {
var list: std.ArrayList(usize) = .empty;
var it = std.mem.tokenizeScalar(u8, text, ',');
while (it.next()) |item| try list.append(ctx.arena, try parseCount(ctx, item));
if (list.items.len == 0) {
try fail(ctx.out, usage);
unreachable;
}
return list.items;
}
fn fail(out: *std.Io.Writer, text: []const u8) !void {
try out.writeAll(text);
try out.flush();
std.process.exit(2);
}
cmdEchoservers берёт вариант из флага: --select, --poll или --epoll, и для последнего пишет в лог настоящее имя механизма через Model.label. cmdConcBench понимает список числа клиентов через запятую, число говорящих, число строк и список моделей, а kqueue принимает как синоним epoll.
root.zig и build.zig
//! Корень модуля `tiny`. Программа, тесты шагов и our-runner берут части
//! сервера отсюда. Шаги: 63 `addr` и `hostinfo`, 64 `socket`, `echo`,
//! `echo_std` и `fdio`, 65 `request`, 66 `server`. Глава 12 живёт в `conc`.
pub const addr = @import("net/addr.zig");
pub const hostinfo = @import("net/hostinfo.zig");
pub const fdio = @import("net/fdio.zig");
pub const socket = @import("net/socket.zig");
pub const echo = @import("net/echo.zig");
pub const echo_std = @import("net/echo_std.zig");
pub const request = @import("http/request.zig");
pub const server = @import("http/tiny.zig");
/// Конкурентность, глава 12.
pub const conc = struct {
pub const conn = @import("conc/conn.zig");
pub const echo_procs = @import("conc/echo_procs.zig");
pub const echo_select = @import("conc/echo_select.zig");
pub const echo_poll = @import("conc/echo_poll.zig");
pub const echo_epoll = @import("conc/echo_epoll.zig");
pub const bench = @import("conc/bench.zig");
};
// Публичное API для других пакетов: our-runner добавляет маршрут POST /run.
pub const Server = server.Server;
pub const Request = server.Request;
pub const Response = server.Response;
pub const Route = server.Route;
pub const Options = server.Options;
pub const Method = request.Method;
/// Номера уроков курса, на которых проект вырос. Каждому шагу
/// соответствует файл `tests/step_NN.zig`, и все они зелёные на финале.
-const project_steps = [_]u8{ 63, 64, 65, 66 };
+const project_steps = [_]u8{ 63, 64, 65, 66, 68 };
pub fn build(b: *std.Build) void {
Прогон
Сервер на select, два клиента по очереди:
$ tiny echoservers --select 15224
echoservers (select): listening on port 15224
Server received 6 (6 total) bytes on fd 5
Server received 6 (12 total) bytes on fd 5
Server received 13 (25 total) bytes on fd 5
Первые две строки прислал один клиент (hello, world), третью второй (второй\n, 13 байт: кириллица в UTF-8 по два байта). Строку лога книги Server received N (M total) bytes печатает Conn.reply, и total в ней общий для всех клиентов, чего сервер на процессах сделать не мог. --poll печатает то же самое.
conc-bench на своей машине:
$ zig build -Doptimize=ReleaseFast
$ zig-out/bin/tiny conc-bench --clients 1000 --active 10 --lines 10000
Без ReleaseFast ты замеришь отладочную сборку со всеми проверками, и сравнение моделей утонет в них.
Тесты шага
//! Шаг 68: эхо-сервер на процессах, событийный эхо-сервер на `select`,
//! `poll` и `epoll`/`kqueue`, бенч моделей. Сервер живёт в потоке этого
//! процесса (или в ребёнке для процессов и бенча), клиенты на libc-сокетах.
const std = @import("std");
const tiny = @import("tiny");
const conc = tiny.conc;
const socket = tiny.socket;
const fdio = tiny.fdio;
const testing = std.testing;
const c = std.c;
/// Клиент с таймаутом на чтение: сломанный сервер валит тест, а не вешает.
pub fn connect(port: u16) !c.fd_t {
const fd = try socket.openClientfd("127.0.0.1", port);
const tv: c.timeval = .{ .sec = 5, .usec = 0 };
_ = c.setsockopt(fd, c.SOL.SOCKET, c.SO.RCVTIMEO, &tv, @sizeOf(c.timeval));
return fd;
}
/// Пишет `text` и дочитывает ровно столько же байтов ответа.
pub fn roundtrip(fd: c.fd_t, text: []const u8, buf: []u8) ![]const u8 {
try fdio.writen(fd, text);
const n = try fdio.readn(fd, buf[0..text.len]);
return buf[0..n];
}
test "FdSet: бит fd % 64 в слове fd / 64, как FD_SET в glibc и macOS" {
var set: conc.echo_select.FdSet = .{};
set.set(3);
set.set(64);
set.set(1023);
try testing.expect(set.isSet(3) and set.isSet(64) and set.isSet(1023));
try testing.expect(!set.isSet(4));
try testing.expectEqual(@as(u64, 1 << 3), set.words[0]);
try testing.expectEqual(@as(u64, 1), set.words[1]);
set.clear(64);
try testing.expect(!set.isSet(64));
// Байтовая раскладка совпадает с битовой картой `fd_set`: fd 1023 это
// старший бит последнего байта.
const bytes: *const [128]u8 = @ptrCast(&set.words);
try testing.expectEqual(@as(u8, 0x80), bytes[127]);
}
/// Сервер на процессах живёт в своём процессе, созданном до клиентов.
/// Иначе дети `echoserverp` унаследовали бы клиентские сокеты теста, и
/// закрытие сокета в тесте не давало бы ребёнку конца потока. Процесс
/// выходит с 0, если все дети собраны и зомби не осталось.
fn procsServer(listenfd: c.fd_t, n: usize) noreturn {
var discard: std.Io.Writer.Discarding = .init(&.{});
conc.echo_procs.serve(listenfd, n, &discard.writer) catch c._exit(10);
socket.close(listenfd);
var waited: usize = 0;
while (conc.echo_procs.reaped.load(.monotonic) < n) : (waited += 1) {
if (waited == 500) c._exit(11);
const ts: c.timespec = .{ .sec = 0, .nsec = 10 * std.time.ns_per_ms };
_ = c.nanosleep(&ts, null);
}
// Собирать больше некого: waitpid отвечает ECHILD.
if (c.waitpid(-1, null, c.W.NOHANG) != -1 or std.posix.errno(@as(c_int, -1)) != .CHILD) c._exit(12);
c._exit(0);
}
test "echoserverp: три клиента сразу, каждого обслуживает свой процесс, зомби не остаётся" {
const listenfd = try socket.openListenfd(0);
const port = socket.localPort(listenfd).?;
const pid = c.fork();
if (pid == 0) procsServer(listenfd, 3);
socket.close(listenfd);
// Все трое подключены одновременно, отвечают в обратном порядке:
// итеративный сервер застрял бы на первом.
var fds: [3]c.fd_t = undefined;
for (&fds) |*fd| fd.* = try connect(port);
var buf: [64]u8 = undefined;
for (0..3) |k| {
const i = 2 - k;
const line = try std.fmt.bufPrint(&buf, "client {d}\n", .{i});
var reply: [64]u8 = undefined;
try testing.expectEqualStrings(line, try roundtrip(fds[i], line, &reply));
}
for (fds) |fd| socket.close(fd);
var status: c_int = 0;
try testing.expectEqual(pid, c.waitpid(pid, &status, 0));
try testing.expect(c.W.IFEXITED(@bitCast(status)));
try testing.expectEqual(@as(u8, 0), c.W.EXITSTATUS(@bitCast(status)));
}
fn eventServe(comptime serve: anytype, listenfd: c.fd_t, n: usize, stats: *conc.conn.Stats) void {
stats.* = serve(testing.allocator, listenfd, .{ .max_clients = n }) catch |err| {
std.debug.print("event serve: {t}\n", .{err});
return;
};
}
fn checkEventServer(comptime serve: anytype) !void {
const listenfd = try socket.openListenfd(0);
defer socket.close(listenfd);
const port = socket.localPort(listenfd).?;
var stats: conc.conn.Stats = .{};
const thread = try std.Thread.spawn(.{}, eventServe, .{ serve, listenfd, 3, &stats });
var fds: [3]c.fd_t = undefined;
for (&fds) |*fd| fd.* = try connect(port);
var reply: [8192]u8 = undefined;
// Неполная строка у первого клиента не должна задержать остальных.
try fdio.writen(fds[0], "полстро");
try testing.expectEqualStrings("второй\n", try roundtrip(fds[1], "второй\n", &reply));
try testing.expectEqualStrings("третий\n", try roundtrip(fds[2], "третий\n", &reply));
try fdio.writen(fds[0], "ки\n");
const n = try fdio.readn(fds[0], reply[0.."полстроки\n".len]);
try testing.expectEqualStrings("полстроки\n", reply[0..n]);
// Две строки одним write и строка длиннее буфера клиента.
try testing.expectEqualStrings("a\nb\n", try roundtrip(fds[1], "a\nb\n", &reply));
var long: [5000]u8 = undefined;
@memset(&long, 'z');
long[long.len - 1] = '\n';
try testing.expectEqualSlices(u8, &long, try roundtrip(fds[2], &long, &reply));
for (fds) |fd| socket.close(fd);
thread.join();
try testing.expectEqual(@as(usize, 3), stats.clients);
const total = "полстроки\n".len + "второй\n".len + "третий\n".len + 4 + long.len;
try testing.expectEqual(total, stats.bytes);
}
test "echoservers на select: пул книги, неполная строка никого не держит" {
try checkEventServer(conc.echo_select.serve);
}
test "echoservers на poll" {
try checkEventServer(conc.echo_poll.serve);
}
test "echoservers на epoll (Linux) или kqueue (macOS)" {
try checkEventServer(conc.echo_epoll.serve);
}
test "Poller.remove: снятый слушающий сокет больше не будит wait" {
const listenfd = try socket.openListenfd(0);
defer socket.close(listenfd);
var poller: conc.echo_epoll.Poller = try .init();
defer poller.deinit();
try poller.add(listenfd);
var events: [4]conc.echo_epoll.Poller.Event = undefined;
// Соединение ждёт в очереди на приём: слушающий сокет готов к чтению.
const client = try connect(socket.localPort(listenfd).?);
defer socket.close(client);
const first = try poller.wait(&events);
try testing.expectEqual(listenfd, conc.echo_epoll.Poller.eventFd(first[0]));
// После remove его не видно, хотя соединение так и не принято.
try poller.remove(listenfd);
var pipe: [2]c.fd_t = undefined;
try testing.expectEqual(@as(c_int, 0), c.pipe(&pipe));
defer for (pipe) |fd| socket.close(fd);
try fdio.writen(pipe[1], "x");
try poller.add(pipe[0]);
const ready = try poller.wait(&events);
try testing.expectEqual(@as(usize, 1), ready.len);
try testing.expectEqual(pipe[0], conc.echo_epoll.Poller.eventFd(ready[0]));
}
test "conc-bench: каждая модель отвечает на всю нагрузку" {
_ = conc.bench.raiseFdLimit(1024);
inline for (.{ .procs, .select, .poll, .epoll }) |model| {
const r = try conc.bench.run(testing.allocator, testing.io, model, .{ .clients = 8, .lines = 5 });
try testing.expectEqual(@as(usize, 8), r.clients);
try testing.expect(r.throughput > 0);
try testing.expect(r.p50_us <= r.p99_us);
}
}
Что проверяется:
FdSetкладёт бит туда же, кудаFD_SETв glibc и macOS: fd 1023 это старший бит последнего байта.- Сервер на процессах обслуживает трёх клиентов, подключённых одновременно, отвечая им в обратном порядке: итеративный сервер застрял бы на первом. Сервер живёт в своём процессе, созданном до клиентов, и выходит с кодом 0, только если все дети собраны и
waitpidотвечаетECHILD, то есть зомби нет. - Событийные серверы на
select,pollиepoll/kqueueпроходят одну и ту же функциюcheckEventServer: половина строки у первого клиента не задерживает второго и третьего, две строки однимwrite, строка длиннее буфера клиента, общий счётчик байтов. Poller.remove: снятый слушающий сокет больше не будитwait, хотя соединение в его очереди так и не принято.conc-benchна восьми клиентах: каждая модель отвечает на всю нагрузку.
У каждого клиента таймаут на чтение в пять секунд через SO_RCVTIMEO. Сломанный сервер валит тест, а не вешает его навсегда.
$ zig build test -Dstep=68 --summary all
Build Summary: 6/6 steps succeeded; 7/7 tests passed
Все шаги сразу (zig build test) дают на этом состоянии эталона 50 тестов, и все зелёные на macOS и в контейнере с Linux.
На macOS
Весь урок работает на macOS напрямую: fork, SIGCHLD, select, poll и kqueue там есть, выводы выше сняты на macOS 26 (Apple Silicon). Linux-вариант проверен в контейнере runner-zig (linux/arm64), рецепт из README эталона:
docker run --rm --platform linux/arm64 -v "$PWD":/w:ro ghcr.io/bondiano/runner-zig:dev sh -c '
mkdir -p /tmp/p && cd /w && cp -r build.zig build.zig.zon src cgi www tests /tmp/p/ && cd /tmp/p &&
zig build test --summary all --cache-dir /tmp/p/.zig-cache --global-cache-dir /tmp/g'
Разница такая:
epollесть только в Linux. На macOS его место занимаетkqueue, выбор делаетPollerпри компиляции, и флаг--epollв логе честно называетсяkqueue. Чтобы потрогать самepoll, собери и гоняй эталон в контейнере.- Мягкий предел дескрипторов на macOS 256.
conc-benchподнимает его до жёсткого сам, для своих программ поможетulimit -n 4096в том же терминале. - На Intel-маках заголовки libc переименовывают
selectв символselect$1050, на Apple Silicon и в glibc это простоselect. Наше объявление черезextern "c"рассчитано на эти два случая. - Чьи дескрипторы открыты у процесса, на Linux видно в
/proc/<pid>/fd, на macOS черезlsof -p <pid>. - Отдельный листинг из урока (
select_stdin.zig) на macOS собирается как есть, libc там линкуется всегда. На Linux добавь-lc, эталон делает это сам черезlink_libc = true. - Числа отличаются, и не только масштабом: на macOS
pollв сценарии C10K почти вдвое медленнееselect, на Linux они вровень.
Практика
В задаче ты соберёшь пул событийного сервера на poll так, как книга строит его на select: фиксированные слоты, дыры с fd = -1, maxi, nready. В std Zig 0.16 для этого всё есть: std.posix.pollfd с полями fd, events, revents, константы std.posix.POLL.IN, POLL.HUP, POLL.ERR и std.posix.poll(fds, timeout), который возвращает число готовых. Разница с echo_poll.zig из урока в том, что ушедший клиент оставляет дыру, а не сдвигает массив, и maxi должен съехать влево через все дыры сразу.
Чтение, запись и закрытие задача получает параметром io: anytype. Это любое значение с методами read, writeAll и close: настоящий сервер передал бы обёртку над системными вызовами, а тесты передают подделку и проверяют, кого и сколько раз ты прочитал. Есть ли у переданного значения эти методы, компилятор проверит при сборке, в месте вызова. Последний тест настоящий: пары UNIX-сокетов и std.posix.poll с таймаутом.
Упражнения
Итоги
- Итеративный сервер спит на одном дескрипторе, пока другие готовы. Три модели главы 12 перекрывают сессии по-разному: процессы и потоки отдают переключение ядру, мультиплексирование делает его само в одном цикле.
- Сервер на процессах:
forkна соединение, родитель закрываетconnfd, ребёнокlistenfd, иначе соединение не закроется или порт останется занят. Детей собирает обработчикSIGCHLDцикломwaitpid(-1, WNOHANG), потому что сигналы не копятся. БезSA_RESTARTacceptвозвращаетEINTR. - Процессы дают изоляцию и платят за неё
fork, памятью и отсутствием общего состояния. selectберёт битовые карты на 1024 дескриптора и переписывает их ответом, поэтому набор копируется перед каждым вызовом, а первый аргумент это наибольший дескриптор плюс один.- Событийный сервер держит каждого клиента как конечный автомат: состояние в буфере, событие это готовность, переход это ровно один
read. Неполная строка ждёт в буфере, а не вread. pollубирает пределFD_SETSIZEи копию набора, но не линейную цену: каждый вызов проходит всех подключённых.epollиkqueueзапоминают интерес в ядре один раз и отдают только готовые. На тысяче соединений, где говорят десять, это дало на Linux пятикратный выигрыш уselectиpoll. Когда говорят все, разницы почти нет.- По умолчанию оба уведомляют по уровню, как
select. Режим по фронту требует неблокирующих сокетов и чтения доEAGAIN. - Событийная модель не прощает долгих обработчиков: блокирующая запись или тяжёлое вычисление останавливают всех клиентов сразу.
Дальше
Сегодня сервер научился не ждать одного клиента двумя способами. Процессы дали изоляцию ценой fork и памяти, событийный цикл дал общий счётчик и копеечную цену соединения ценой автомата, который приходится держать в голове. А epoll и kqueue показали, что на тысяче молчащих клиентов цена ожидания решает больше, чем цена обработки.
Осталась третья модель. В следующем уроке сессии снова станут простым последовательным кодом, как у процессов, но в одном адресном пространстве, как у событийного сервера: это потоки. Мы разберём std.Thread поверх pthreads, классическую ошибку с адресом переменной цикла, эхо-сервер и TINY с потоком на соединение, добавим строку threads в conc-bench, а zt prof научится различать стеки разных потоков.
домашка