Раздел 32 · Системное программирование: Zig, ассемблер, Verilog
Проект: async в Zig и кэширующий прокси
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Проект: async в Zig и кэширующий прокси
В прошлом уроке мы разобрали, как потоки ломают друг другу данные: четыре класса небезопасных функций, блокировка с копированием, взаимоблокировка и правило порядка захвата. До этого было пять уроков про то, как вообще обслужить много клиентов сразу: процессы,
selectиpoll, поток на соединение, пул. Каждая модель заставляла переписывать сервер по-своему: у событийного сервера из урока 68 код сессии разрезан на куски вокруг циклаpoll, у пула из урока 71 есть очередь и воркеры. Сегодня последний урок раздела, и в нём Zig 0.16 показывает свой ответ: код пишется один раз, как обычный блокирующий, а модель выбирается при запуске, как аллокатор. Мы разберёмstd.Ioцеликом (io.async,Group,Queue,Select, отмену), честно скажем, что из этого в 0.16.0 работает, переведём на него эхо-сервер и раннерzboxс отменой по крайнему сроку, а в финале соберём proxylab: кэширующий прокси с пулом потоков и LRU-кэшем под блокировкой читателей-писателей.
Цели урока
- Объяснить, что такое «цвет функции», почему async/await в JavaScript красит вызывающих и как Zig 0.16 обходится без этого:
std.Ioпередаётся параметром, какAllocator. - Различать
io.asyncиio.concurrent: первое разрешает одновременность, второе её требует. Предсказать, что сделает каждое подIo.Threadedи под однопоточной реализацией. - Пользоваться
Future,Io.Group,Io.QueueиIo.Select, знать, чем группа отличается от брошенного потока. - Понимать отмену в
std.Io: точки отмены,error.Canceled, какIo.Threadedбудит поток, застрявший в системном вызове, и почему долгий цикл без ввода-вывода отменить нельзя. - Сопоставить реализации
std.Ioс моделями книги: процессы, мультиплексирование, потоки, пул. Знать, чтоIo.Eventedв 0.16.0 умеет и чего пока не умеет. - Написать кэширующий прокси по мотивам proxylab: разбор абсолютного URI, переписанные заголовки, пересылка ответа кусками, кэш только для удачных ответов.
- Спроектировать LRU-кэш под
std.Io.RwLock, в котором попадания не мешают друг другу, хотя чтение тоже меняет порядок LRU. - Довести
zboxдо модели наstd.Io: задача гонится с таймером, проигравший отменяется, отмена доходит до чужой программы в её группе процессов.
Идея: четыре модели и одна функция
Вспомни, как выглядел эхо-сервер в каждой модели главы 12.
| Модель | Урок | Кто ждёт ввода-вывода | Как выглядит код сессии |
|---|---|---|---|
| процесс на соединение | 68 | ядро усыпляет процесс в read | обычный цикл read, write |
| мультиплексирование | 68 | один цикл на poll, epoll или kqueue | конечный автомат: «что пришло на этот дескриптор, в каком состоянии клиент» |
| поток на соединение | 69 | ядро усыпляет поток в read | обычный цикл, но память общая |
| пул потоков | 71 | то же, потоков фиксированное число | обычный цикл плюс очередь и воркеры вокруг |
Три из четырёх моделей оставляют код сессии простым: читай строку, пиши строку. Ломается он только в событийной модели, и именно она самая экономная: один поток, никаких переключений контекста ядром, тысяча соединений стоит тысячу маленьких структур, а не тысячу стеков. Прикладные языки хотели и экономию событийного цикла, и простой код, и пришли к async/await. Ты видел это в уроке про Promise и async/await: функция с await внутри сама становится async, возвращает обещание, и её вызывающий тоже обязан либо await, либо работать с обещанием.
Это свойство называют цветом функции. Красная функция заражает всех, кто её вызывает. Библиотека, которая хочет работать и там и там, пишется дважды: fs.readFileSync и fs.readFile, requests и aiohttp, std::fs и tokio::fs в Rust (как async fn превращается в конечный автомат, разобрано в уроке про рассахаривание async). Хуже того, цвет выбирается автором библиотеки, а не автором программы.
Zig 0.16 делает шаг, который уже знаком тебе по аллокаторам. В уроке 05 функция, которой нужна память, получала std.mem.Allocator параметром и не знала, арена это, DebugAllocator или page_allocator. С вводом-выводом теперь так же: функция, которой надо ждать сеть, диск или таймер, получает параметр io: std.Io. Это интерфейс из указателя на состояние и таблицы функций (vtable), такой же, как у Allocator и у Io.Reader из урока 61. Под ним может быть пул потоков ОС (Io.Threaded) или цикл событий с файберами (Io.Evented). Функция одна и та же, никакого async в сигнатуре. Асинхронность стала свойством того, как функцию вызвали, а не того, как её написали.
Бесцветная функция в деле
Вот функция, которая «скачивает» страницу: ждёт заданное время и возвращает размер. Обычная функция с обычным возвратом. И функция, которая скачивает две страницы так, чтобы загрузки могли идти одновременно.
const std = @import("std");
const Io = std.Io;
/// «Скачать» страницу: полсекунды ожидания и длина ответа. Обычная функция
/// с обычным возвратом, никакого `async` в сигнатуре.
fn fetch(io: Io, ms: i64, size: usize) Io.Cancelable!usize {
try io.sleep(.fromMilliseconds(ms), .awake);
return size;
}
/// Две загрузки, которые могут идти одновременно. Могут, а не обязаны:
/// решает реализация `io`.
fn fetchBoth(io: Io) !usize {
var a = io.async(fetch, .{ io, 200, 100 });
defer _ = a.cancel(io) catch 0;
var b = io.async(fetch, .{ io, 200, 20 });
defer _ = b.cancel(io) catch 0;
return try a.await(io) + try b.await(io);
}
fn elapsedMs(io: Io, comptime f: anytype) !struct { usize, i64 } {
const start = Io.Clock.awake.now(io);
const result = try f(io);
return .{ result, start.durationTo(Io.Clock.awake.now(io)).toMilliseconds() };
}
test "один код, две реализации Io" {
// Пул потоков: загрузки идут одновременно, около 200 мс.
var threaded: Io.Threaded = .init(std.testing.allocator, .{});
defer threaded.deinit();
const par = try elapsedMs(threaded.io(), fetchBoth);
// Один поток: `async` выполняет функцию на месте, около 400 мс.
var single: Io.Threaded = .init_single_threaded;
const seq = try elapsedMs(single.io(), fetchBoth);
try std.testing.expectEqual(120, par[0]);
try std.testing.expectEqual(120, seq[0]);
std.debug.print("\nThreaded {d} мс, single-threaded {d} мс\n", .{ par[1], seq[1] });
try std.testing.expect(par[1] < 350);
try std.testing.expect(seq[1] >= 400);
}
test "concurrent требует настоящей одновременности" {
var single: Io.Threaded = .init_single_threaded;
const io = single.io();
try std.testing.expectError(error.ConcurrencyUnavailable, io.concurrent(fetch, .{ io, 10, 1 }));
}
Тест прогоняет один и тот же fetchBoth под двумя реализациями Io. Io.Threaded.init это пул потоков, init_single_threaded это та же реализация, которой запрещено заводить потоки. Прогон на Apple M4 Max (16 ядер, macOS 26.6.2, Zig 0.16.0):
$ zig test colorless.zig
1/2 colorless.test.один код, две реализации Io...
Threaded 201 мс, single-threaded 408 мс
OK
2/2 colorless.test.concurrent требует настоящей одновременности...OK
All 2 tests passed.
Две загрузки по 200 мс в пуле заняли 201 мс, в одном потоке 408. Результат одинаковый, код одинаковый, и fetch понятия не имеет, какой из двух случаев сейчас идёт. В контейнере с Linux (runner-zig:dev, arm64) тесты те же и тоже зелёные.
async разрешает, concurrent требует
Во втором тесте спрятано главное различие std.Io. Функций запуска две.
io.async(f, args)говорит: «fможно выполнять одновременно со мной, но не обязательно». Реализация вправе вызватьfпрямо здесь и вернуть готовыйFuture. Однопоточная реализация так и делает, поэтомуfetchBothв ней шёл 400 мс: первая загрузка выполнилась целиком внутри первогоasync, вторая внутри второго.io.concurrent(f, args)говорит: «fобязана идти одновременно со мной». Если реализация не может выделить ей свою единицу исполнения (поток или файбер), вызов возвращаетerror.ConcurrencyUnavailable. Однопоточная реализация не может никогда.
Зачем два вида. Представь производителя и потребителя, которые общаются через очередь на один элемент. Если потребителя запустить через async, а реализация выполнит его на месте, он будет вечно ждать элемента, которого производитель ещё не положил, потому что производитель ждёт возврата из async. Это взаимоблокировка из прошлого урока, только без мьютексов. Где корректность программы зависит от одновременности, нужен concurrent. Где одновременность только ускоряет (две независимые загрузки), хватает async, и такой код работает под любой реализацией, даже под однопоточной во встраиваемой системе.
Io.Threaded тоже может выполнить async на месте. Число потоков под async у него ограничено полем async_limit, по умолчанию «ядер минус один» (на 16-ядерной машине 15). Когда все заняты, очередной async выполняется в вызывающем потоке. Это видно прямо в Threaded.zig:
if (busy_count >= @intFromEnum(t.async_limit)) {
mutexUnlock(&t.mutex);
future.destroy(gpa);
start(context.ptr, result.ptr);
return null;
}
У concurrent свой предел concurrent_limit, по умолчанию без ограничения: пул растёт, пока ОС даёт потоки. Мы ещё встретим это различие в zbox, где таймер крайнего срока запущенный через async мог бы стартовать только после задачи, которую он должен прервать.
Future: результат, которого ещё нет
io.async и io.concurrent возвращают Future(T), где T это тип возврата функции. У него два метода.
future.await(io)ждёт завершения и отдаёт результат.future.cancel(io)просит задачу остановиться, ждёт её завершения и тоже отдаёт результат. Если задача успела закончиться сама, это её честный результат; если нет, скорее всегоerror.Canceled.
Оба идемпотентны: второй вызов отдаёт тот же результат без ожидания. Поэтому идиома такая же, как с памятью: сразу после запуска defer future.cancel(io). В fetchBoth это defer _ = a.cancel(io) catch 0;. Если a.await вернул ошибку, мы выходим, и defer отменяет вторую загрузку, а не бросает её висеть. Если всё прошло хорошо, cancel после await ничего не делает. Забыть await или cancel у Future это утечка того же рода, что забыть free: ресурсы задачи (поток пула, память под аргументы и результат) не вернутся.
Под капотом: потоки и файберы
Io.Threaded: пул потоков ОС
Io.Threaded это пул из урока 71, только написанный за тебя. async кладёт задачу в очередь выполнения и, если свободных потоков нет, а предел не достигнут, заводит новый поток, который остаётся в пуле навсегда. Операции ввода-вывода (read, accept, nanosleep, futex) это обычные блокирующие системные вызовы: поток засыпает в ядре, как в модели «поток на соединение». Примитивы синхронизации из урока 70 (Io.Mutex, Io.Semaphore, Io.Condition, Io.RwLock) построены над futex и потому тоже идут через io.
Отсюда практичное следствие, которым пользовались уроки 70 до 72: под Io.Threaded примитивы std.Io можно звать из потоков, созданных std.Thread.spawn. Эталоны так и делают, тестовый std.testing.io это Io.Threaded.
Io.Evented: файберы над циклом событий
Io.Evented устроен как событийный сервер из урока 68, но цикл событий спрятан внутрь. Каждая задача это файбер со своим стеком. Когда файбер доходит до ввода-вывода, реализация ставит операцию в очередь ядра (io_uring в Linux, libdispatch поверх kqueue в macOS), сохраняет регистры файбера и прыгает в другой файбер, готовый продолжить. Когда ядро сообщает о завершении, файбер снова попадает в очередь готовых. Код сессии при этом остаётся блокирующим с виду: reader.takeDelimiterInclusive('\n') «просто ждёт».
Переключение файбера в Zig 0.16 это десяток инструкций. Вот сердце std/Io/fiber.zig для aarch64:
ldp x0, x2, [x1] // x0 = куда сохранить себя, x2 = откуда взять другого
ldr x3, [x2, #16] // x3 = pc другого файбера
mov x4, sp
stp x4, fp, [x0] // сохранить свои sp и fp
adr x5, 0f
ldp x4, fp, [x2] // загрузить sp и fp другого
str x5, [x0, #16] // свой pc: метка 0 ниже
mov sp, x4
br x3 // прыжок в другой файбер
0:
Сохраняются только три значения: sp, fp, pc. Остальные регистры объявлены в asm volatile как испорченные (длинный список .x19 = true, ...), и компилятор сам кладёт живые значения на стек файбера перед переключением. Это то соглашение о вызовах из урока 14 и список испорченных регистров из урока 18, а сама идея та же, что у переключения процессов в ядре Y86 из урока 50. Разница в том, что здесь нет ни прерывания, ни смены привилегий, ни таблицы страниц: файберы живут в одном адресном пространстве и уступают друг другу сами. Это кооперативная многозадачность, в отличие от вытесняющей у потоков ОС.
Что из этого работает в Zig 0.16.0
Интерфейс std.Io в 0.16.0 стабилен, реализации нет. Эталоны курса проверили каждую руками, на macOS 26 и в контейнере с Linux 7.0, и вот честная сводка.
| Реализация | Где | Что работает | Чего нет |
|---|---|---|---|
Io.Threaded | везде | всё: файлы, процессы, сеть std.Io.net, Group, Select, отмена | ничего из нужного нам |
Io.Evented = Io.Dispatch | macOS | файберы, Group, sleep, файлы, запуск процессов, чтение пайпов | сеть: listen, accept, read и write у std.Io.net отвечают error.NetworkDown; deinit не компилируется; отмена задачи, ждущей в futex, роняет процесс |
Io.Evented = Io.Uring | Linux | ничего | не компилируется: dirOpenDir возвращает ошибку, которой нет в объявленном множестве; сеть там такая же заглушка |
Io.Kqueue | BSD | ничего | не компилируется: таблица функций отстала от Io.VTable |
Файберы доступны только на aarch64, riscv64 и x86_64 (std.Io.fiber.supported). Поэтому всё, что в этом уроке работает с сетью, работает на Io.Threaded. Там, где можно обойти сеть std.Io, мы покажем, как тот же код идёт под Io.Evented на macOS. И это главный аргумент в пользу бесцветного кода: когда Evented получит сеть, наши серверы заработают на нём без единой правки, потому что они написаны против интерфейса, а не против реализации.
Group: задачи, которые не теряются
Один Future на задачу удобен, когда задач две. Сервер же заводит задачу на каждое соединение, и хранить тысячу Future незачем. Для этого есть Io.Group: неупорядоченное множество задач, которое ждут или отменяют только целиком.
group.async(io, f, args)иgroup.concurrent(io, f, args)запускают задачу в группе. Функция обязана возвращать то, что приводится кCancelable!void: результат группа не хранит, ошибкуCanceledпроглатывает.group.await(io)ждёт, пока закончатся все.group.cancel(io)отменяет всех и ждёт, пока все действительно закончатся.
Ресурсы задачи возвращаются, как только она закончилась, а не когда закончилась вся группа. Поэтому группа на весь срок жизни сервера, в которую без конца добавляются сессии, не течёт.
Главное свойство группы видно по тому, чего нельзя. В уроке 69 эхо-сервер делал detach потоку сессии, и после этого о потоке никто не знал: сервер не мог ни дождаться его, ни остановить. Группа так не умеет. Задача в группе живёт не дольше блока, в котором объявлена группа, если на выходе стоит defer group.cancel(io). Это называют структурной конкурентностью, и ты её уже видел у дочерних файберов в Effect: у конкурентности появляется форма, совпадающая с формой кода.
Queue: sbuf, который уже в std
В уроке 70 мы написали sbuf: кольцевой буфер производителя-потребителя на трёх семафорах. В std.Io он есть готовый: Io.Queue(T). Буфер даёшь ты, как и у sbuf, а сверху две вещи, которых у нас не было.
- Закрытие.
q.close(io)запрещает новыеput, аgetотдаёт оставшиеся элементы и только потомerror.Closed. В пулеecho_preиз урока 71 воркеры останавливались «ядовитой пилюлей», дескриптором-1, по одному на поток. С закрытием пилюли не нужны. - Пачки и отмена.
put(io, elements, min)иget(io, buffer, min)работают сразу с несколькими элементами, а ожидание в них это точка отмены. Для одного элемента естьputOneиgetOne, для кода, который отмену не принимает, варианты...Uncancelable.
const std = @import("std");
const Io = std.Io;
/// Производитель: числа от 1 до n в очередь, потом закрыть её.
fn produce(io: Io, q: *Io.Queue(u32), n: u32) Io.Cancelable!void {
defer q.close(io);
var i: u32 = 1;
while (i <= n) : (i += 1) q.putOne(io, i) catch |err| switch (err) {
error.Closed => return,
error.Canceled => |e| return e,
};
}
/// Потребитель: забирает, пока очередь не закрыта и не пуста.
fn consume(io: Io, q: *Io.Queue(u32), sum: *std.atomic.Value(u64)) Io.Cancelable!void {
while (true) {
const item = q.getOne(io) catch |err| switch (err) {
error.Closed => return,
error.Canceled => |e| return e,
};
_ = sum.fetchAdd(item, .monotonic);
}
}
test "Io.Queue и Io.Group: sbuf из урока 70 уже в std" {
const io = std.testing.io;
var buffer: [4]u32 = undefined;
var q: Io.Queue(u32) = .init(&buffer);
var sum: std.atomic.Value(u64) = .init(0);
var group: Io.Group = .init;
defer group.cancel(io);
try group.concurrent(io, produce, .{ io, &q, 1000 });
for (0..3) |_| try group.concurrent(io, consume, .{ io, &q, &sum });
try group.await(io);
try std.testing.expectEqual(1000 * 1001 / 2, sum.load(.monotonic));
}
Один производитель, три потребителя, очередь на четыре места, тысяча чисел. Производитель закрывает очередь в defer, потребители выходят по error.Closed, группа ждёт всех. Обрати внимание, что все четыре задачи запущены через concurrent: производитель и потребители должны идти одновременно, иначе полная очередь остановит производителя навсегда.
Отмена
В книге отмены нет. Если потоку сервера надо прекратить работу, книга предлагает pthread_cancel одной строкой и дальше не идёт, и правильно: отмена потока в POSIX это минное поле. В std.Io отмена это часть интерфейса, и устроена она аккуратно.
Точки отмены
future.cancel(io) и group.cancel(io) не убивают задачу. Они ставят флаг, и задача получает error.Canceled из ближайшей точки отмены. Дальше она ведёт себя как при любой ошибке: errdefer и defer убирают ресурсы, ошибка поднимается вверх, и задача заканчивается сама. Никакого состояния, брошенного посередине, никакого мьютекса, оставшегося запертым.
const std = @import("std");
const Io = std.Io;
/// Долгая задача: десять раз по секунде. Каждый `sleep` это точка отмены.
fn slow(io: Io, steps: *u32) Io.Cancelable!u32 {
for (0..10) |_| {
try io.sleep(.fromSeconds(1), .awake);
steps.* += 1;
}
return 42;
}
test "cancel: задача выходит с Canceled на ближайшей точке отмены" {
const io = std.testing.io;
var steps: u32 = 0;
var future = try io.concurrent(slow, .{ io, &steps });
try io.sleep(.fromMilliseconds(50), .awake);
const start = Io.Clock.awake.now(io);
try std.testing.expectError(error.Canceled, future.cancel(io));
const waited = start.durationTo(Io.Clock.awake.now(io)).toMilliseconds();
try std.testing.expectEqual(0, steps);
try std.testing.expect(waited < 100);
}
const Race = union(enum) {
work: Io.Cancelable!u32,
timeout: Io.Cancelable!void,
};
/// Задача против таймера: кто первый, тот и ответ, второго отменяем.
fn withTimeout(io: Io, steps: *u32, ms: i64) !u32 {
var buffer: [2]Race = undefined;
var select = Io.Select(Race).init(io, &buffer);
defer select.cancelDiscard();
try select.concurrent(.work, slow, .{ io, steps });
try select.concurrent(.timeout, Io.sleep, .{ io, .fromMilliseconds(ms), .awake });
return switch (try select.await()) {
.work => |result| result,
.timeout => error.Timeout,
};
}
test "Select: таймер обогнал задачу, задачу отменили" {
var steps: u32 = 0;
try std.testing.expectError(error.Timeout, withTimeout(std.testing.io, &steps, 1500));
try std.testing.expectEqual(1, steps);
}
Первый тест: задача из десяти секундных sleep, отмена через 50 мс. Задача не успела сделать ни шага, cancel вернул error.Canceled меньше чем за 100 мс, хотя sleep просил секунду.
Три правила, которые из этого следуют.
- Не глотай
error.Canceled. Отмена сигналится один раз: следующая точка отмены уже не вернёт ошибку. Если поймать её и продолжить, задача доработает до конца, а тот, кто отменял, будет ждать. Если отмену надо отложить (закончить запись в файл), естьio.recancel(), который взводит её снова, иio.swapCancelProtection(.blocked), который временно отключает точки отмены. - Цикл без ввода-вывода не отменяется. Если задача считает
psumиз урока 71 миллиард итераций, точек отмены в ней нет, иcancelбудет ждать конца счёта. Для таких задач естьio.checkCancel(): ничего не делает, кроме проверки флага. Звать её раз в какое-то число итераций. - Uncancelable там, где отмена невозможна по смыслу.
lockUncancelableу мьютекса,getOneUncancelableу очереди: для кода, который вызывают из обычных потоковstd.Thread(у них нет задачи, которую можно отменить) или который не имеет права остановиться на полпути.
Как Io.Threaded будит спящий поток
Под Io.Evented отмена проста: файбер спит в очереди реализации, его можно разбудить с ошибкой. Под Io.Threaded поток задачи спит в ядре, в nanosleep или read, и флаг в памяти он не увидит. Тут пригодится урок про сигналы. При старте Io.Threaded.init ставит на SIGIO пустой обработчик без SA_RESTART. Отмена посылает потоку задачи SIGIO через pthread_kill (в Linux без pthreads через tgkill). Сигнал прерывает системный вызов, тот возвращает EINTR, обёртка std.Io видит EINTR, проверяет флаг отмены и возвращает error.Canceled.
И тут же гонка из урока 49: сигнал может прийти за мгновение до того, как поток вошёл в системный вызов. Тогда обработчик отработает впустую, а поток уснёт и не проснётся. Урок 49 лечил такую гонку маской и sigsuspend; Io.Threaded лечит её повтором. Пока поток отмечен как «заблокирован в вызове», отменяющий шлёт сигнал снова, с экспоненциальной паузой от 1 мкс. Комментарий в исходнике честно говорит, что на практике второй сигнал нужен редко.
Отсюда видно и ограничение. std.Io умеет прервать только то, что сам вызвал. Если задача запустила дочерний процесс, отмена задачи не отменяет процесс: его надо убить самому, в errdefer. Это ровно то, что понадобится zbox.
Select: кто первый
Io.Select(U) запускает несколько задач с результатами разных типов и отдаёт первый готовый. U это union(enum), где у каждого поля тип результата своей задачи. select.concurrent(.work, f, args) запускает f и кладёт её результат в поле work. select.await() ждёт первый результат, select.cancelDiscard() отменяет всех остальных и ждёт их конца. Внутри это Group плюс Queue(U) на буфере, который даёшь ты.
Второй тест в cancel.zig это тайм-аут в самом общем виде: задача на десять секунд против таймера на 1,5 с. Таймер пришёл первым, withTimeout вернул error.Timeout, а defer select.cancelDiscard() отменил задачу, которая успела сделать ровно один шаг. Так в zbox будет устроен крайний срок задачи.
Книга против std.Io
Теперь можно положить std.Io рядом с моделями главы 12.
| Модель книги | Что в std.Io | Цена |
|---|---|---|
| процесс на соединение | прямого аналога нет; процесс это граница изоляции, а не конкурентности (так в zbox) | fork, своя память, общение через пайпы |
событийный сервер на poll | Io.Evented: цикл событий внутри, файбер на соединение, код сессии блокирующий с виду | один поток на ядро, стек файбера на соединение, кооперативность: файбер, который считает, держит свой поток |
| поток на соединение | Io.Threaded и group.concurrent на соединение | поток ОС на соединение, пока не упрёмся в concurrent_limit |
| пул потоков | Io.Threaded с async и пределом async_limit; Io.Queue вместо sbuf | потоков не больше предела, лишние задачи ждут или выполняются на месте |
Код сессии во всех строках справа один и тот же. Это и есть обещание урока, и сейчас мы его проверим на эхо-сервере.
Шаг tinylab: эхо на std.Io
Проект tiny из уроков 63 до 72 живёт в одном каталоге: модуль tiny (src/root.zig) и программа с подкомандами (src/main.zig). Сегодня в нём три новых файла: эхо-сервер на std.Io и две части прокси. Начнём с эха: это самый короткий способ увидеть Group и отмену на живой сети.
echo_io.zig
Сессия пишется один раз, echoLines над Io.Reader и Io.Writer. Ей всё равно, откуда они: из std.Io.net.Stream, из Io.File поверх сокета или из буфера в тесте. Серверов два, потому что сеть std.Io в 0.16.0 есть только у Io.Threaded.
//! Эхо-сервер на `std.Io` из Zig 0.16: код не знает, потоки под ним или
//! файберы. Соединение уходит в `Io.Group` через `concurrent`, а кто его
//! обслужит, решает реализация `Io`, которую выбрали в `main`:
//!
//! * `Io.Threaded`: пул потоков ОС, блокирующие системные вызовы. Это модель
//! «поток на соединение» из урока 69, только потоки берёт `Io`.
//! * `Io.Evented`: файберы поверх цикла событий (`Dispatch` на macOS, `Uring`
//! на Linux). Это событийный сервер из урока 68, только цикл событий
//! спрятан в `Io`, а код сессии пишется как обычный блокирующий.
//!
//! Что в 0.16.0 работает на самом деле (сверено с std и запуском):
//!
//! * `Io.Threaded` умеет всё, в том числе `std.Io.net`. Это `serve`.
//! * `Io.Evented` на macOS (`Dispatch`) поднимается, файберы, `Group`, `sleep`
//! и чтение и запись файлов работают, а вся сеть (`netListenIp`,
//! `netAccept`, `netRead`, `netWrite`) заглушена и отвечает `NetworkDown`.
//! Поэтому есть `serveFd`: слушающий сокет и `accept` из libc, а сессия
//! читает и пишет сокет как `Io.File`. Её код один и тот же под обеими
//! реализациями. Ещё две ловушки 0.16.0: `Evented.deinit` не компилируется
//! (освобождает стек цикла как массив, а не срез), а с аллокатором под
//! мьютексом (`backing_allocator_needs_mutex = true`, по умолчанию) файбер
//! падает при освобождении. Отсюда `runEvented` ниже.
//! * `Io.Evented` на Linux (`Uring`) в 0.16.0 не компилируется вовсе
//! (несовпадение множеств ошибок в `Uring.zig`), а сеть там тоже заглушена.
//!
//! Отмена: `Group.cancel` (и `defer` с ним в `serve`) будит каждую сессию на
//! ближайшей точке отмены, то есть на чтении из сокета, и та выходит с
//! `error.Canceled`. `Future.cancel` на задаче с `serve` делает то же для
//! всего сервера: `accept` прерывается, сессии отменяются, сокеты закрываются.
const std = @import("std");
const builtin = @import("builtin");
const c = std.c;
const posix = std.posix;
const Io = std.Io;
const net = Io.net;
/// Эхо строк, пока клиент не закроет соединение. Одна функция на все
/// варианты: ей всё равно, откуда `Reader` и `Writer`.
pub fn echoLines(r: *Io.Reader, w: *Io.Writer) error{ ReadFailed, WriteFailed }!usize {
var total: usize = 0;
while (r.takeDelimiterInclusive('\n')) |line| {
try w.writeAll(line);
try w.flush();
total += line.len;
} else |err| switch (err) {
error.EndOfStream => {},
error.StreamTooLong => {
// Строка длиннее буфера: отдаём, что есть, и дальше по кругу.
try w.writeAll(r.buffered());
try w.flush();
r.tossBuffered();
return total + try echoLines(r, w);
},
error.ReadFailed => return error.ReadFailed,
}
return total;
}
/// Сервер целиком на `std.Io.net`: `accept`, сессия в группе, повтор.
/// `max_clients` 0 значит вечно; иначе после стольких `accept` ждём группу.
pub fn serve(io: Io, server: *net.Server, max_clients: usize) !void {
var group: Io.Group = .init;
// Выход по ошибке или по отмене: не бросаем сессии, а отменяем их.
defer group.cancel(io);
var served: usize = 0;
while (max_clients == 0 or served < max_clients) : (served += 1) {
const stream = try server.accept(io);
group.concurrent(io, netSession, .{ io, stream }) catch |err| {
stream.close(io);
return err;
};
}
try group.await(io);
}
fn netSession(io: Io, stream: net.Stream) Io.Cancelable!void {
defer stream.close(io);
var in_buf: [4096]u8 = undefined;
var out_buf: [4096]u8 = undefined;
var reader = stream.reader(io, &in_buf);
var writer = stream.writer(io, &out_buf);
_ = echoLines(&reader.interface, &writer.interface) catch {
// `Io.Reader` умеет сказать только `ReadFailed`; настоящая причина
// лежит в `reader.err`. Отмену надо вернуть наверх, остальное нет.
if (reader.err) |err| if (err == error.Canceled) return error.Canceled;
};
}
/// Тот же сервер, но сеть мимо `std.Io.net`: слушающий сокет из libc
/// (например, `socket.openListenfd`), сессия пишет и читает `Io.File`.
/// Работает и под `Io.Evented` на macOS, где `std.Io.net` заглушена.
pub fn serveFd(io: Io, listenfd: c.fd_t, max_clients: usize) !void {
setNonblocking(listenfd, true);
defer setNonblocking(listenfd, false);
var group: Io.Group = .init;
defer group.cancel(io);
var served: usize = 0;
while (max_clients == 0 or served < max_clients) : (served += 1) {
const connfd = try acceptPolling(io, listenfd);
const file: Io.File = .{ .handle = connfd, .flags = .{ .nonblocking = false } };
group.concurrent(io, fileSession, .{ io, file }) catch |err| {
file.close(io);
return err;
};
}
try group.await(io);
}
// ponytail: `accept` опрашивается раз в миллисекунду. У `Io` в 0.16 нет
// операции «ждать готовности дескриптора», а блокирующий `accept` в файбере
// занял бы поток цикла событий. Пока `Io.Evented` не умеет сеть, это
// самый короткий честный путь; с сетью в `Evented` берите `serve`.
fn acceptPolling(io: Io, listenfd: c.fd_t) !c.fd_t {
while (true) {
const fd = c.accept(listenfd, null, null);
if (fd >= 0) {
// На macOS принятый сокет наследует O_NONBLOCK слушающего, на
// Linux нет. Сессии нужен блокирующий: `Io` сам решает, ждать
// ли ему в потоке или в цикле событий.
setNonblocking(fd, false);
return fd;
}
switch (posix.errno(fd)) {
.AGAIN, .INTR => try io.sleep(.fromMilliseconds(1), .awake),
else => return error.AcceptFailed,
}
}
}
fn fileSession(io: Io, file: Io.File) Io.Cancelable!void {
defer file.close(io);
var in_buf: [4096]u8 = undefined;
var out_buf: [4096]u8 = undefined;
var reader = file.readerStreaming(io, &in_buf);
var writer = file.writerStreaming(io, &out_buf);
_ = echoLines(&reader.interface, &writer.interface) catch {
if (reader.err) |err| if (err == error.Canceled) return error.Canceled;
};
}
fn setNonblocking(fd: c.fd_t, on: bool) void {
const flags = c.fcntl(fd, c.F.GETFL);
if (flags < 0) return;
const bit: c_int = @bitCast(c.O{ .NONBLOCK = true });
_ = c.fcntl(fd, c.F.SETFL, if (on) flags | bit else flags & ~bit);
}
/// Есть ли рабочий `Io.Evented` в этой сборке: только macOS, см. шапку.
pub const evented_supported = builtin.os.tag == .macos and Io.Evented != void;
/// Запускает `serveFd` под `Io.Evented`. Текущий поток становится главным
/// файбером цикла событий. Звать только из `main`: после первой же уступки
/// главный файбер продолжает на любом потоке пула `Dispatch`, и вернуть его
/// на родной поток умеет только `deinit`, а он в 0.16.0 не компилируется.
/// Из `std.Thread` возврат получается, а `join` этого потока виснет навсегда.
/// Аллокатор файберов `smp_allocator` без внешнего мьютекса: с мьютексом
/// 0.16.0 падает. Ресурсы цикла вернёт ОС при выходе процесса.
pub fn runEvented(listenfd: c.fd_t, max_clients: usize) !void {
if (comptime !evented_supported) {
return error.EventedUnavailable;
} else {
var ev: Io.Evented = undefined;
try ev.init(std.heap.smp_allocator, .{ .backing_allocator_needs_mutex = false });
try serveFd(ev.io(), listenfd, max_clients);
}
}
Разберём то, чего не было в эхо-серверах уроков 68 до 71.
serveэто весь сервер.server.accept(io)вместоacceptиз libc,group.concurrentвместоstd.Thread.spawnиdetach. Поток на соединение заводитIo.Threaded, и он же вернёт его в пул, когда сессия закончится. Сокетыstd.Io.netоткрываются с close-on-exec, об этом ниже вzbox.- Почему
concurrent, а неasync. Сессия это ожидание клиента, который может молчать сколько угодно. Если быasyncвыполнил её на месте (пул занят), циклacceptвстал бы до конца этой сессии, и сервер превратился бы в итеративный из урока 64. Корректность зависит от одновременности, значитconcurrent. Если ОС не даст потока,concurrentвернёт ошибку, мы закроем соединение и выйдем, аdefer group.cancel(io)отменит остальных. defer group.cancel(io). Выход изserveпо любой причине, по ошибкеacceptили по отмене всего сервера, отменяет все сессии и ждёт их. Ни одного брошенного потока.- Где настоящая ошибка.
Io.Readerсообщает толькоerror.ReadFailed, как в уроке 61. Настоящую причину реализация кладёт в полеerrсвоего читателя. Отмену надо вернуть наверх, иначе группа будет ждать сессию, которая решила, что всё в порядке. Остальные ошибки чтения (клиент сбросил соединение) для эха не ошибка. serveFdдляIo.Evented. Слушающий сокет иacceptиз libc, а принятый сокет заворачивается вIo.File: чтение и запись файлов уIo.Dispatchработают.acceptнеблокирующий и опрашивается раз в миллисекунду черезio.sleep: блокирующийacceptзанял бы поток цикла событий целиком, а операции «жди готовности дескриптора» уstd.Io0.16 нет. Это отмечено какponytail:с условием, когда убрать.runEventedтолько изmain. Поток, позвавшийEvented.init, становится главным файбером. После первой уступки этот файбер может продолжить на любом потоке пула libdispatch, и вернуть его на родной поток умеет толькоdeinit, который в 0.16.0 не компилируется. Изstd.Threadвозврат получается, аjoinэтого потока потом виснет. Поэтому--eventedэто отдельный процесс, и тест запускает его как дочерний.
Подкоманда echoserver-io
@@
\\ tiny echoservert-pre [--workers N] <port> эхо: пул потоков и sbuf
+ \\ tiny echoserver-io --threaded|--evented <port> эхо на std.Io с Io.Group
\\ tiny conc-bench [--clients 100,1000] [--active A] [--lines K] [--models procs,threads,pool,select,poll,epoll] [--workers N] [--json]
@@
\\ tiny loadgen <host> <port> <conns> <reqs> <path> нагрузка на HTTP-сервер
+ \\ tiny proxy [--workers N] <port> кэширующий прокси
\\
@@
.{ "echoservert-pre", cmdEchoservertPre },
+ .{ "echoserver-io", cmdEchoserverIo },
.{ "conc-bench", cmdConcBench },
@@
.{ "loadgen", cmdLoadgen },
+ .{ "proxy", cmdProxy },
};
@@
try conc.echo_pre.serve(ctx.gpa, ctx.io, listenfd, .{ .workers = workers }, conc.echo_pre.EchoCtx{ .io = ctx.io }, conc.echo_pre.echoCnt);
}
+fn cmdEchoserverIo(ctx: Ctx, rest: []const [:0]const u8) !void {
+ if (rest.len != 2) return fail(ctx.out, usage);
+ const port = try parsePort(ctx, rest[1]);
+ if (std.mem.eql(u8, rest[0], "--threaded")) {
+ const address: std.Io.net.IpAddress = .{ .ip4 = .loopback(port) };
+ var server = try address.listen(ctx.io, .{ .reuse_address = true });
+ defer server.deinit(ctx.io);
+ try ctx.log.print("echoserver-io (Io.Threaded, std.Io.net): listening on {f}\n", .{server.socket.address});
+ try ctx.log.flush();
+ try conc.echo_io.serve(ctx.io, &server, 0);
+ } else if (std.mem.eql(u8, rest[0], "--evented")) {
+ if (!conc.echo_io.evented_supported) {
+ try ctx.log.writeAll("echoserver-io: Io.Evented в Zig 0.16.0 работает только на macOS (Dispatch); на Linux Uring не компилируется\n");
+ try ctx.log.flush();
+ std.process.exit(1);
+ }
+ const listenfd = try listenOn(ctx, "echoserver-io (Io.Evented, libc accept, Io.File)", port);
+ defer tiny.socket.close(listenfd);
+ try conc.echo_io.runEvented(listenfd, 0);
+ } else return fail(ctx.out, usage);
+}
+
fn cmdConcBench(ctx: Ctx, rest: []const [:0]const u8) !void {
@@
try ctx.out.flush();
}
+fn cmdProxy(ctx: Ctx, rest: []const [:0]const u8) !void {
+ var args = rest;
+ var workers: usize = 8;
+ if (args.len > 1 and std.mem.eql(u8, args[0], "--workers")) {
+ workers = try parseCount(ctx, args[1]);
+ args = args[2..];
+ }
+ if (args.len != 1) return fail(ctx.out, usage);
+ tiny.server.ignoreSigpipe();
+ const listenfd = try listenOn(ctx, "proxy", try parsePort(ctx, args[0]));
+ defer tiny.socket.close(listenfd);
+ var cache: tiny.proxy.cache.Cache = .init(ctx.gpa, ctx.io, .{});
+ defer cache.deinit();
+ var proxy: tiny.proxy.server.Proxy = .{ .gpa = ctx.gpa, .cache = &cache };
+ try conc.echo_pre.serve(ctx.gpa, ctx.io, listenfd, .{ .workers = workers }, &proxy, tiny.proxy.server.handle);
+}
+
fn listenOn(ctx: Ctx, name: []const u8, port: u16) !std.c.fd_t {
Здесь сразу и подкоманда прокси, до которого мы дойдём через раздел. ctx.io приходит из main: std.process.Init в Zig 0.16 даёт программе готовый Io.Threaded, так же как готовый gpa.
Прогон
$ tiny echoserver-io --threaded 15218
echoserver-io (Io.Threaded, std.Io.net): listening on 127.0.0.1:15218
$ tiny echoserver-io --evented 15219
echoserver-io (Io.Evented, libc accept, Io.File): listening on port 15219
$ printf 'fiber\n' | tiny echoclient 127.0.0.1 15219
fiber
Второй сервер обслуживает соединение файбером под libdispatch. На Linux (runner-zig:dev) --evented честно отказывается с кодом 1, а --threaded работает так же.
Финал: кэширующий прокси
Proxylab, последняя лабораторная книги, это программа-посредник. Браузер, которому сказали ходить через прокси, шлёт ему не GET /home.html, а запрос с полным адресом. Вот что прислал curl -x слушающему nc -l (macOS 26.6.2, curl 8.7.1):
GET http://localhost:8000/home.html HTTP/1.1
Host: localhost:8000
User-Agent: curl/8.7.1
Accept: */*
Proxy-Connection: Keep-Alive
Прокси разбирает URI, сам подключается к localhost:8000, отправляет серверу обычный запрос с путём и пересылает ответ обратно клиенту. Задание лабораторной идёт в три этапа, и мы пройдём их все.
- Последовательный прокси. Разбор абсолютного URI, свой запрос к серверу на HTTP/1.0 с заданными заголовками, пересылка ответа.
- Конкурентный. Много клиентов сразу. Книга предлагает поток на запрос; у нас есть пул из урока 71, и прокси встаёт на него одной строкой.
- Кэширующий. Недавние ответы лежат в памяти, повторный запрос сервер не видит. Лимиты лабораторной: весь кэш не больше 1 049 000 байт, один объект не больше 102 400 байт, вытесняется давно не читанный объект. Кэш общий на все потоки, читают его много чаще, чем пишут: блокировка читателей-писателей из урока 70.
С точки зрения раздела про сети прокси это посредник из RFC 9110, «прямой» (forward) прокси: его выбирает клиент. Обратный прокси вроде nginx перед приложением выбирает сервер. Устройство одно и то же: сервер для клиента и клиент для сервера в одном процессе.
proxy.zig
//! Кэширующий HTTP/1.0-прокси из proxylab. Браузер шлёт прокси абсолютный
//! URI (`GET http://host:port/path HTTP/1.1`), прокси разбирает его, сам
//! подключается к `host:port`, пересылает запрос с путём и переписанными
//! заголовками, а ответ по кусочкам отдаёт клиенту. Ответ не больше
//! `MAX_OBJECT_SIZE` заодно копится и после конца ложится в кэш; следующий
//! запрос того же URI сервер уже не увидит.
//!
//! Конкурентность берётся из урока 71: `handle` вешается на пул потоков
//! `echo_pre.serve`, кэш общий на все потоки (см. `cache.zig`).
const std = @import("std");
const c = std.c;
const Io = std.Io;
const socket = @import("../net/socket.zig");
const fdio = @import("../net/fdio.zig");
const request = @import("../http/request.zig");
const tiny = @import("../http/tiny.zig");
const Cache = @import("cache.zig").Cache;
/// Заголовок `User-Agent` из заготовки proxylab.
pub const user_agent = "Mozilla/5.0 (X11; Linux x86_64; rv:10.0.3) Gecko/20120305 Firefox/10.0.3";
pub const Target = struct {
host: []const u8,
port: u16,
/// Путь с `?args`, начинается с `/`.
path: []const u8,
};
/// `http://host[:port][/path]`. Без порта 80, без пути `/`.
pub fn parseUri(uri: []const u8) error{BadUri}!Target {
const scheme = "http://";
if (uri.len < scheme.len or !std.ascii.eqlIgnoreCase(uri[0..scheme.len], scheme)) return error.BadUri;
const rest = uri[scheme.len..];
const slash = std.mem.indexOfScalar(u8, rest, '/') orelse rest.len;
const authority = rest[0..slash];
const path = if (slash < rest.len) rest[slash..] else "/";
const colon = std.mem.indexOfScalar(u8, authority, ':');
const host = if (colon) |i| authority[0..i] else authority;
if (host.len == 0) return error.BadUri;
const port = if (colon) |i| std.fmt.parseInt(u16, authority[i + 1 ..], 10) catch return error.BadUri else 80;
return .{ .host = host, .port = port, .path = path };
}
/// Запрос к серверу. `Host` клиента, если он его прислал, иначе собираем
/// из URI. `User-Agent`, `Connection` и `Proxy-Connection` всегда наши:
/// прокси говорит на HTTP/1.0 и держать соединение не умеет. Остальные
/// заголовки клиента уходят как есть.
pub fn writeRequest(w: *Io.Writer, target: Target, headers: *const request.Headers) Io.Writer.Error!void {
try w.print("GET {s} HTTP/1.0\r\n", .{target.path});
if (headers.get("Host")) |host| {
try w.print("Host: {s}\r\n", .{host});
} else if (target.port == 80) {
try w.print("Host: {s}\r\n", .{target.host});
} else {
try w.print("Host: {s}:{d}\r\n", .{ target.host, target.port });
}
try w.print("User-Agent: {s}\r\n", .{user_agent});
try w.writeAll("Connection: close\r\nProxy-Connection: close\r\n");
var it = headers.iterator();
while (it.next()) |h| {
const own = [_][]const u8{ "Host", "User-Agent", "Connection", "Proxy-Connection" };
for (own) |name| {
if (std.ascii.eqlIgnoreCase(h.name, name)) break;
} else try w.print("{s}: {s}\r\n", .{ h.name, h.value });
}
try w.writeAll("\r\n");
}
pub const Proxy = struct {
gpa: std.mem.Allocator,
cache: *Cache,
quiet: bool = false,
/// Одна транзакция. Зовётся из воркера пула, дескриптор закрывает пул.
pub fn handle(p: *Proxy, connfd: c.fd_t, peer: []const u8) void {
var arena_state: std.heap.ArenaAllocator = .init(p.gpa);
defer arena_state.deinit();
p.doit(arena_state.allocator(), connfd, peer) catch |err| p.log("{s} proxy failed: {t}", .{ peer, err });
}
fn doit(p: *Proxy, arena: std.mem.Allocator, connfd: c.fd_t, peer: []const u8) !void {
var in_buf: [8192]u8 = undefined;
var out_buf: [1024]u8 = undefined;
var reader: fdio.Reader = .init(connfd, &in_buf);
var client: fdio.Writer = .init(connfd, &out_buf);
const out = &client.interface;
defer out.flush() catch {};
const raw_line = reader.interface.takeDelimiterInclusive('\n') catch |err| switch (err) {
error.EndOfStream => return,
else => return tiny.clienterror(out, 400, "request line too long"),
};
var words = std.mem.tokenizeScalar(u8, std.mem.trimEnd(u8, raw_line, "\r\n"), ' ');
const method = words.next() orelse return tiny.clienterror(out, 400, "malformed request line");
const uri = words.next() orelse return tiny.clienterror(out, 400, "malformed request line");
var headers: request.Headers = .{};
request.readRequestHeaders(&reader.interface, &headers) catch return tiny.clienterror(out, 400, "malformed headers");
if (!std.mem.eql(u8, method, "GET")) {
p.log("{s} \"{s} {s}\" 501", .{ peer, method, uri });
return tiny.clienterror(out, 501, "Proxy does not implement this method");
}
const target = parseUri(uri) catch {
p.log("{s} \"{s} {s}\" 400", .{ peer, method, uri });
return tiny.clienterror(out, 400, "proxy wants an absolute http:// URI");
};
if (try p.cache.get(uri, arena)) |object| {
try fdio.writen(connfd, object);
p.log("{s} \"GET {s}\" HIT {d}", .{ peer, uri, object.len });
return;
}
const serverfd = socket.openClientfd(target.host, target.port) catch {
p.log("{s} \"GET {s}\" 502", .{ peer, uri });
return tiny.clienterror(out, 502, "proxy could not reach the server");
};
defer socket.close(serverfd);
var req_buf: [8192]u8 = undefined;
var server: fdio.Writer = .init(serverfd, &req_buf);
try writeRequest(&server.interface, target, &headers);
try server.interface.flush();
// Ответ кусками к клиенту; копия для кэша, пока влезает в лимит.
var object: std.ArrayList(u8) = .empty;
var cacheable = true;
var total: usize = 0;
var chunk: [8192]u8 = undefined;
while (true) {
const n = try fdio.readSome(serverfd, &chunk);
if (n == 0) break;
try fdio.writen(connfd, chunk[0..n]);
total += n;
if (cacheable and total <= p.cache.limits.object) {
try object.appendSlice(arena, chunk[0..n]);
} else cacheable = false;
}
// Кэшируем только удачные ответы: 404 или 502 завтра может стать 200.
const stored = cacheable and isOk(object.items) and try p.cache.put(uri, object.items);
p.log("{s} \"GET {s}\" MISS {d}{s}", .{ peer, uri, total, if (stored) " cached" else "" });
}
fn log(p: *Proxy, comptime fmt: []const u8, args: anytype) void {
if (p.quiet) return;
var buf: [1024]u8 = undefined;
const text = std.fmt.bufPrint(&buf, fmt ++ "\n", args) catch blk: {
@memcpy(buf[buf.len - 4 ..], "...\n");
break :blk buf[0..];
};
fdio.writen(2, text) catch {};
}
};
/// Строка статуса `HTTP/1.x 200 ...`.
pub fn isOk(response: []const u8) bool {
return response.len >= 13 and std.mem.startsWith(u8, response, "HTTP/1.") and std.mem.eql(u8, response[8..13], " 200 ");
}
/// Для пула: `echo_pre.serve(..., proxy, proxy_mod.handle)`.
pub fn handle(p: *Proxy, connfd: c.fd_t, peer: []const u8) void {
p.handle(connfd, peer);
}
Что здесь решено и почему.
parseUriстрогий. Толькоhttp://, хост не пустой, порт помещается вu16.https://прокси безCONNECTне обслужит, поэтому это 400, а не попытка. Путь с?argsуходит серверу как есть: для CGI из урока 66 аргументы и есть запрос.- HTTP/1.0 и
Connection: closeвсегда. Прокси не держит соединения открытыми: ответ сервера читается до конца потока, и конец потока означает конец ответа. С keep-alive пришлось бы разбиратьContent-Lengthиchunked, чтобы понять, где кончается ответ. ЗаголовкиConnectionиProxy-Connectionклиента поэтому заменяются нашими,User-Agentберётся из заготовки лабораторной, аHostклиента сохраняется: по нему сервер с несколькими сайтами понимает, какой из них нужен. - Остальные заголовки уходят как есть. Цикл с
for ... elseиз Zig:elseуforвыполняется, если цикл не прервалсяbreak, то есть имя не нашлось среди наших четырёх. - Ответ кусками.
readSomeчитает, сколько пришло (до 8 КиБ), и тут же отдаёт клиенту черезwritenиз урока 64. Картинка на мегабайт не копится в памяти целиком, а клиент начинает получать байты раньше, чем сервер закончил. Параллельно те же куски копятся вobject, пока не превышен лимит объекта; превысили, копить перестаём, но пересылать продолжаем. - Кэшируется только
200. Ответ 404 завтра может стать 200, когда файл появится, а 502 значит, что сервер лежал. Кэшировать их значит запомнить ошибку. CGI с ответом 200 кэшируется, как требует лабораторная, хотя настоящий прокси смотрел бы наCache-Control(это упражнение). - Арена на соединение.
handleзаводит арену, всё, что нужно транзакции (заголовки, копия объекта), берётся из неё и уходит однимdeinit. Копия объекта из кэша тоже живёт в арене:cache.getполучает её аллокатор. - Лог одной записью.
logформатирует строку в буфер на стеке и отдаёт её однимwriteв stderr. Восемь воркеров пишут одновременно, иwriteна одну строку не даёт строкам перемешаться, как было с логом TINY в уроке 69.
Прокси пользуется двумя мелкими правками старых файлов. readSome в fdio.zig: один read с повтором на EINTR, в отличие от readn, который дочитывает до конца буфера.
@@
return got;
}
+/// Один `read` с повтором на `EINTR`: сколько пришло, 0 на конце потока.
+/// Прокси (урок 73) пересылает ответ сервера такими кусками.
+pub fn readSome(fd: c.fd_t, buf: []u8) ReadError!usize {
+ while (true) {
+ const n = c.read(fd, buf.ptr, buf.len);
+ if (n >= 0) return @intCast(n);
+ if (posix.errno(n) != .INTR) return error.ReadFailed;
+ }
+}
+
/// Читатель поверх дескриптора. Всё, что выше `vtable` (строки, `peek`,
И текст статуса 502 в request.zig, чтобы clienterror из TINY умел отвечать «Bad Gateway», когда сервер недоступен:
@@
500 => "Internal Server Error",
501 => "Not Implemented",
+ 502 => "Bad Gateway",
else => "Unknown",
Кэш: где LRU встречает читателей-писателей
Прежде чем читать код кэша, сформулируем задачу честно. Попадание в кэш должно:
- найти объект по URI;
- отдать его байты;
- отметить, что объект только что читали, чтобы LRU не вытеснил его первым.
Первые два шага только читают, и их хочется делать под общей блокировкой, чтобы восемь воркеров отдавали горячую страницу одновременно. Третий шаг пишет. Классический LRU это двусвязный список, где прочитанный элемент переезжает в голову. Переставить элемент в списке под общей блокировкой нельзя: два читателя переставят его одновременно и порвут список. Значит, либо каждое попадание берёт исключительную блокировку (и RwLock ничем не лучше мьютекса), либо список стерегут отдельным мьютексом (и каждое попадание снова проходит через один общий замок).
Выход в том, чтобы запись, которую делает читатель, была одной атомарной операцией, не мешающей другим читателям. Вместо списка у каждого объекта метка «когда читали» (last_used), а время это атомарный счётчик clock. Попадание берёт clock.fetchAdd(1) и кладёт результат в метку объекта через store. Два читателя одного объекта одновременно запишут две метки, и победит любая: обе «сейчас», порядок LRU от этого не пострадает. Порядок живёт в метках, а не в структуре, и структура под общей блокировкой не меняется.
Цена переезжает к писателю: чтобы найти жертву, он проходит по всем объектам и ищет наименьшую метку, O(n). Это честный размен. Проход идёт только на промахе, который вытесняет, а промах и так ходит в сеть, это миллисекунды против микросекунд прохода по сотне объектов. При лимитах лабораторной (мегабайт на объекты от сотен байт) объектов сотни или тысячи.
Попробуй это на модели. Три клиента делят кэш на 32 КиБ; каждый клиент сперва входит читателем, при промахе выходит из чтения, качает объект с сервера без блокировки и встаёт в очередь на запись. Сценарий «горячая страница» показывает трёх читателей внутри одновременно, «вытеснение» показывает, как /photo.jpg выталкивает объект с хвоста, «большой объект» пролетает мимо кэша, а «писатель голодает» переключает политику блокировки: при приоритете читателей писатель ждёт, пока поток читателей не иссякнет.
Сценарий «холодный кэш» стоит пройти по шагам: два клиента промахиваются по одному адресу, оба качают его с сервера, и второй писатель находит объект уже в кэше. Запомни эту картину, в прогоне прокси под нагрузкой мы увидим её в восьмикратном размере.
cache.zig
//! Кэш веб-объектов прокси из proxylab: ключ это URI запроса, значение
//! весь ответ сервера как есть (строка статуса, заголовки, тело). Лимиты
//! книги: весь кэш не больше `MAX_CACHE_SIZE` (1 049 000 байт), один объект
//! не больше `MAX_OBJECT_SIZE` (102 400 байт). Вытесняется давно не читанный
//! объект (LRU).
//!
//! Блокировка читателей-писателей: `get` берёт замок на чтение, и попадания
//! в кэш из разных потоков идут параллельно; `put` берёт замок на запись.
//! Но LRU просит, чтобы и чтение что-то писало: время последнего доступа.
//! Здесь это атомарная метка `last_used` у объекта, а часы это атомарный
//! счётчик `clock`. Читатели пишут её `store`-ом под замком на чтение и
//! друг другу не мешают; писатель читает метки под замком на запись, когда
//! выбирает жертву. Отдельный мьютекс на список LRU тоже работал бы, но
//! превратил бы каждое попадание в захват одного общего мьютекса, то есть
//! свёл бы на нет замок на чтение.
const std = @import("std");
const Io = std.Io;
pub const max_cache_size = 1_049_000;
pub const max_object_size = 102_400;
pub const Limits = struct {
cache: usize = max_cache_size,
object: usize = max_object_size,
};
const Entry = struct {
data: []u8,
last_used: std.atomic.Value(u64),
};
pub const Cache = struct {
gpa: std.mem.Allocator,
io: Io,
limits: Limits,
lock: Io.RwLock = .init,
/// Ключи и данные принадлежат кэшу.
map: std.StringHashMapUnmanaged(*Entry) = .empty,
/// Сумма длин всех объектов.
size: usize = 0,
clock: std.atomic.Value(u64) = .init(0),
/// Счётчики для тестов и лога: попадания, промахи, вытеснения.
hits: std.atomic.Value(usize) = .init(0),
misses: std.atomic.Value(usize) = .init(0),
evictions: usize = 0,
pub fn init(gpa: std.mem.Allocator, io: Io, limits: Limits) Cache {
return .{ .gpa = gpa, .io = io, .limits = limits };
}
pub fn deinit(cache: *Cache) void {
var it = cache.map.iterator();
while (it.next()) |kv| {
cache.gpa.free(kv.key_ptr.*);
cache.gpa.free(kv.value_ptr.*.data);
cache.gpa.destroy(kv.value_ptr.*);
}
cache.map.deinit(cache.gpa);
}
/// Копия объекта в `out` (память вызывающего) или `null`. Копируем под
/// замком: как только он отпущен, писатель вправе вытеснить объект и
/// освободить его байты.
pub fn get(cache: *Cache, key: []const u8, out: std.mem.Allocator) !?[]u8 {
cache.lock.lockSharedUncancelable(cache.io);
defer cache.lock.unlockShared(cache.io);
const entry = cache.map.get(key) orelse {
_ = cache.misses.fetchAdd(1, .monotonic);
return null;
};
entry.last_used.store(cache.clock.fetchAdd(1, .monotonic), .monotonic);
_ = cache.hits.fetchAdd(1, .monotonic);
return try out.dupe(u8, entry.data);
}
/// Кладёт копию `data`. Объект больше лимита не кэшируется (`false`).
/// Пока не хватает места, вытесняется объект с самой старой меткой.
pub fn put(cache: *Cache, key: []const u8, data: []const u8) !bool {
if (data.len > cache.limits.object or data.len > cache.limits.cache) return false;
cache.lock.lockUncancelable(cache.io);
defer cache.lock.unlock(cache.io);
// Два потока могли промахнуться по одному URI одновременно: второй
// просто освежает метку.
if (cache.map.get(key)) |entry| {
entry.last_used.store(cache.clock.fetchAdd(1, .monotonic), .monotonic);
return true;
}
while (cache.size + data.len > cache.limits.cache) cache.evictOne();
const owned_key = try cache.gpa.dupe(u8, key);
errdefer cache.gpa.free(owned_key);
const entry = try cache.gpa.create(Entry);
errdefer cache.gpa.destroy(entry);
entry.* = .{ .data = try cache.gpa.dupe(u8, data), .last_used = .init(cache.clock.fetchAdd(1, .monotonic)) };
errdefer cache.gpa.free(entry.data);
try cache.map.put(cache.gpa, owned_key, entry);
cache.size += data.len;
return true;
}
// ponytail: жертва ищется проходом по всем объектам, O(n) на вытеснение.
// При лимитах книги объектов сотни, а проход идёт только на промахе с
// вытеснением; двусвязный список дал бы O(1), но его пришлось бы
// переставлять и на попадании, то есть под замком на запись.
fn evictOne(cache: *Cache) void {
var victim: ?std.StringHashMapUnmanaged(*Entry).Entry = null;
var oldest: u64 = std.math.maxInt(u64);
var it = cache.map.iterator();
while (it.next()) |kv| {
const stamp = kv.value_ptr.*.last_used.load(.monotonic);
if (stamp < oldest) {
oldest = stamp;
victim = kv;
}
}
const kv = victim orelse unreachable;
const key = kv.key_ptr.*;
const entry = kv.value_ptr.*;
cache.size -= entry.data.len;
cache.evictions += 1;
_ = cache.map.remove(key);
cache.gpa.free(key);
cache.gpa.free(entry.data);
cache.gpa.destroy(entry);
}
/// Есть ли ключ, без побочных эффектов на LRU. Для тестов.
pub fn contains(cache: *Cache, key: []const u8) bool {
cache.lock.lockSharedUncancelable(cache.io);
defer cache.lock.unlockShared(cache.io);
return cache.map.contains(key);
}
};
Три места, которые легко написать неправильно.
- Копия под замком.
getвозвращает не срез на байты кэша, а копию в памяти вызывающего, и делает её доunlockShared. Стоит отпустить замок, писатель вправе вытеснить объект и освободить его байты, и срез станет висячим. Это блокировка с копированием из прошлого урока, та же, что уgethostbyname_ts. - Повторная вставка. Два воркера промахнулись по одному URI, оба скачали объект, оба зовут
put. Второй найдёт ключ и только освежит метку, иначе в карте было бы два ключа с разной памятью, аsizeсчитал бы объект дважды. - Проверка лимита до замка. Слишком большой объект отбрасывается, не трогая замок: незачем останавливать всех читателей ради отказа. По той же причине в идеале стоит и копировать байты до захвата исключительной блокировки, пока держишь её, ждут все. Эталон копирует под замком ради простоты, а решение код-задачи урока делает правильно.
Счётчики hits и misses атомарные, их меняют читатели. evictions и size меняются только под исключительной блокировкой, им атомарность не нужна.
std.Io.RwLock в 0.16 отдаёт приоритет писателю. Внутри у него слово состояния с числом читателей и числом ждущих писателей, мьютекс и семафор. Читатель входит без очереди одной cmpxchg, только пока нет ни пишущего, ни ждущего писателя. Как только писатель встал в очередь, новые читатели идут через мьютекс, который держит писатель, и ждут его. Это второй вариант читателей-писателей из урока 70, и голодания писателя из сценария виджета у нашего прокси нет.
get и put зовут ...Uncancelable варианты: их вызывают воркеры пула echo_pre, это потоки std.Thread, и отменять там некого. В код-задаче урока кэш вызывается из задач std.Io, поэтому там обычные lockShared(io) и lock(io) с try.
Модуль и сборка
@@
pub const deadlock = @import("conc/deadlock.zig");
+ pub const echo_io = @import("conc/echo_io.zig");
};
+/// Кэширующий прокси, урок 73.
+pub const proxy = struct {
+ pub const cache = @import("proxy/cache.zig");
+ pub const server = @import("proxy/proxy.zig");
+};
+
// Публичное API для других пакетов: our-runner добавляет маршрут POST /run.
@@
-const project_steps = [_]u8{ 63, 64, 65, 66, 68, 69, 70, 71, 72 };
+const project_steps = [_]u8{ 63, 64, 65, 66, 68, 69, 70, 71, 72, 73 };
Тесты шага
Тесты идут тремя слоями: эхо на std.Io с отменой, чистые части прокси (разбор URI, запрос к серверу, кэш, в том числе восемь потоков против одного кэша) и прокси перед настоящим TINY. Сервер и прокси крутятся в потоках того же процесса, как в tests/support.zig из урока 66.
//! Шаг 73: эхо на `std.Io` с `Io.Group` и отменой; кэширующий прокси:
//! разбор URI, переписанные заголовки, LRU-кэш под читателями-писателями,
//! прокси перед TINY с конкурентными клиентами и ответами из кэша после
//! остановки TINY.
const std = @import("std");
const tiny = @import("tiny");
const support = @import("support.zig");
const conc = tiny.conc;
const proxy = tiny.proxy.server;
const Cache = tiny.proxy.cache.Cache;
const socket = tiny.socket;
const fdio = tiny.fdio;
const testing = std.testing;
const io = testing.io;
const c = std.c;
const net = std.Io.net;
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;
}
fn roundtrip(fd: c.fd_t, text: []const u8) !void {
try fdio.writen(fd, text);
var buf: [256]u8 = undefined;
const n = try fdio.readn(fd, buf[0..text.len]);
try testing.expectEqualStrings(text, buf[0..n]);
}
// ---- echo на std.Io ----
fn netServe(server: *net.Server, n: usize) void {
conc.echo_io.serve(io, server, n) catch |err| std.debug.print("echo_io.serve: {t}\n", .{err});
}
test "echoserver-io на Io.Threaded: сессии в Io.Group идут одновременно" {
const address: net.IpAddress = .{ .ip4 = .loopback(0) };
var server = try address.listen(io, .{ .reuse_address = true });
defer server.deinit(io);
const port = server.socket.address.getPort();
const thread = try std.Thread.spawn(.{}, netServe, .{ &server, 2 });
// Первый клиент подключился и молчит; второй всё равно получает эхо.
const idle = try connect(port);
const active = try connect(port);
try roundtrip(active, "группа\n");
try roundtrip(idle, "теперь я\n");
socket.close(idle);
socket.close(active);
thread.join();
}
test "отмена: Future.cancel прерывает accept и сессии, клиент видит конец потока" {
const address: net.IpAddress = .{ .ip4 = .loopback(0) };
var server = try address.listen(io, .{ .reuse_address = true });
defer server.deinit(io);
const port = server.socket.address.getPort();
var future = try io.concurrent(conc.echo_io.serve, .{ io, &server, 0 });
const fd = try connect(port);
defer socket.close(fd);
try roundtrip(fd, "до отмены\n");
try testing.expectError(error.Canceled, future.cancel(io));
// Сессия отменена и закрыла сокет: read отдаёт 0, а не висит.
var buf: [8]u8 = undefined;
try testing.expectEqual(@as(usize, 0), try fdio.readSome(fd, &buf));
}
fn fdServe(listenfd: c.fd_t, n: usize) void {
conc.echo_io.serveFd(io, listenfd, n) catch |err| std.debug.print("echo_io.serveFd: {t}\n", .{err});
}
test "serveFd на Io.Threaded: libc accept, сессия пишет и читает Io.File" {
const listenfd = try socket.openListenfd(0);
defer socket.close(listenfd);
const port = socket.localPort(listenfd).?;
const thread = try std.Thread.spawn(.{}, fdServe, .{ listenfd, 2 });
const idle = try connect(port);
const active = try connect(port);
try roundtrip(active, "файл\n");
try roundtrip(idle, "сокет\n");
socket.close(idle);
socket.close(active);
thread.join();
}
test "serveFd на Io.Evented (только macOS): тот же код на файберах" {
if (!conc.echo_io.evented_supported) return error.SkipZigTest;
// Отдельный процесс: поток, позвавший `Evented.init`, становится главным
// файбером и без `deinit` (в 0.16.0 он не компилируется) обратно не
// возвращается. Порт сервер называет первой строкой stderr.
var child = try std.process.spawn(io, .{
.argv = &.{ support.tiny_exe, "echoserver-io", "--evented", "0" },
.stdin = .ignore,
.stdout = .ignore,
.stderr = .pipe,
});
defer child.kill(io);
var log_buf: [512]u8 = undefined;
var log = child.stderr.?.readerStreaming(io, &log_buf);
const first = try log.interface.takeDelimiterInclusive('\n');
const marker = "listening on port ";
const at = std.mem.indexOf(u8, first, marker).? + marker.len;
const port = try std.fmt.parseInt(u16, std.mem.trimEnd(u8, first[at..], "\n"), 10);
const idle = try connect(port);
defer socket.close(idle);
const active = try connect(port);
defer socket.close(active);
try roundtrip(active, "файбер\n");
try roundtrip(idle, "второй файбер\n");
}
// ---- прокси: чистые части ----
test "parseUri: хост, порт, путь; без порта 80, без пути /" {
const t = try proxy.parseUri("http://www.cmu.edu:8080/hub/index.html?x=1");
try testing.expectEqualStrings("www.cmu.edu", t.host);
try testing.expectEqual(@as(u16, 8080), t.port);
try testing.expectEqualStrings("/hub/index.html?x=1", t.path);
const d = try proxy.parseUri("HTTP://example.com");
try testing.expectEqualStrings("example.com", d.host);
try testing.expectEqual(@as(u16, 80), d.port);
try testing.expectEqualStrings("/", d.path);
try testing.expectError(error.BadUri, proxy.parseUri("/home.html"));
try testing.expectError(error.BadUri, proxy.parseUri("https://example.com/"));
try testing.expectError(error.BadUri, proxy.parseUri("http://:80/"));
try testing.expectError(error.BadUri, proxy.parseUri("http://h:99999/"));
}
test "writeRequest: HTTP/1.0, свои User-Agent и Connection, чужие заголовки как есть" {
var headers: tiny.request.Headers = .{};
var reader: std.Io.Reader = .fixed("Host: example.com:8080\r\nUser-Agent: curl/8\r\nConnection: keep-alive\r\nProxy-Connection: keep-alive\r\nAccept: */*\r\n\r\n");
try tiny.request.readRequestHeaders(&reader, &headers);
var buf: [1024]u8 = undefined;
var w: std.Io.Writer = .fixed(&buf);
try proxy.writeRequest(&w, try proxy.parseUri("http://example.com:8080/a"), &headers);
try testing.expectEqualStrings("GET /a HTTP/1.0\r\n" ++
"Host: example.com:8080\r\n" ++
"User-Agent: " ++ proxy.user_agent ++ "\r\n" ++
"Connection: close\r\n" ++
"Proxy-Connection: close\r\n" ++
"Accept: */*\r\n\r\n", w.buffered());
// Без Host у клиента: из URI, порт только если не 80.
var empty: tiny.request.Headers = .{};
w = .fixed(&buf);
try proxy.writeRequest(&w, try proxy.parseUri("http://example.com/"), &empty);
try testing.expect(std.mem.indexOf(u8, w.buffered(), "Host: example.com\r\n") != null);
}
test "isOk: в кэш идут только ответы 200" {
try testing.expect(proxy.isOk("HTTP/1.0 200 OK\r\n\r\n"));
try testing.expect(proxy.isOk("HTTP/1.1 200 \r\n"));
try testing.expect(!proxy.isOk("HTTP/1.0 404 Not Found\r\n"));
try testing.expect(!proxy.isOk("HTTP/1.0 2000 X"));
try testing.expect(!proxy.isOk("HTTP/1.0 20"));
}
// ---- кэш ----
test "кэш: попадание отдаёт копию, объект больше лимита не кладётся" {
var cache: Cache = .init(testing.allocator, io, .{ .cache = 1000, .object = 100 });
defer cache.deinit();
try testing.expect(try cache.put("http://h/a", "alpha"));
const got = (try cache.get("http://h/a", testing.allocator)).?;
defer testing.allocator.free(got);
try testing.expectEqualStrings("alpha", got);
try testing.expectEqual(null, try cache.get("http://h/none", testing.allocator));
const big = [_]u8{'x'} ** 101;
try testing.expect(!try cache.put("http://h/big", &big));
try testing.expect(!cache.contains("http://h/big"));
try testing.expectEqual(@as(usize, 5), cache.size);
try testing.expectEqual(@as(usize, 1), cache.hits.load(.monotonic));
try testing.expectEqual(@as(usize, 1), cache.misses.load(.monotonic));
}
test "кэш: вытесняется давно не читанный, чтение освежает метку" {
var cache: Cache = .init(testing.allocator, io, .{ .cache = 300, .object = 100 });
defer cache.deinit();
const obj = [_]u8{'o'} ** 100;
_ = try cache.put("a", &obj);
_ = try cache.put("b", &obj);
_ = try cache.put("c", &obj);
// `a` самый старый, но его прочли: теперь самый старый `b`.
testing.allocator.free((try cache.get("a", testing.allocator)).?);
_ = try cache.put("d", &obj);
try testing.expect(cache.contains("a"));
try testing.expect(!cache.contains("b"));
try testing.expect(cache.contains("c") and cache.contains("d"));
try testing.expectEqual(@as(usize, 300), cache.size);
try testing.expectEqual(@as(usize, 1), cache.evictions);
}
fn hammer(cache: *Cache, id: usize, bad: *std.atomic.Value(usize)) void {
var key_buf: [16]u8 = undefined;
for (0..2000) |i| {
const k = (i * 7 + id) % 40;
const key = std.fmt.bufPrint(&key_buf, "k{d}", .{k}) catch unreachable;
if (cache.get(key, testing.allocator) catch null) |data| {
defer testing.allocator.free(data);
// Объект ключа kN это N+1 байт со значением N: чужие байты
// или обрезанный объект значат, что замок не держит.
if (data.len != k + 1 or data[0] != k or data[data.len - 1] != k) _ = bad.fetchAdd(1, .monotonic);
} else {
var obj: [64]u8 = undefined;
@memset(obj[0 .. k + 1], @intCast(k));
_ = cache.put(key, obj[0 .. k + 1]) catch {};
}
}
}
test "кэш: восемь потоков читают и пишут, лимит соблюдён, байты не перепутаны" {
var cache: Cache = .init(testing.allocator, io, .{ .cache = 400, .object = 64 });
defer cache.deinit();
var bad: std.atomic.Value(usize) = .init(0);
var threads: [8]std.Thread = undefined;
for (&threads, 0..) |*t, id| t.* = try std.Thread.spawn(.{}, hammer, .{ &cache, id, &bad });
for (threads) |t| t.join();
try testing.expectEqual(@as(usize, 0), bad.load(.monotonic));
try testing.expect(cache.size <= 400);
try testing.expect(cache.evictions > 0);
}
// ---- прокси перед TINY ----
/// Маршрут TINY, который возвращает строку запроса и заголовки, какими
/// их прислал прокси.
fn echoHeaders(arena: std.mem.Allocator, req: *const tiny.Request) anyerror!tiny.Response {
var out: std.ArrayList(u8) = .empty;
try out.print(arena, "{s} {s} {s}\n", .{ req.method_text, req.uri, req.version });
var it = req.headers.iterator();
while (it.next()) |h| try out.print(arena, "{s}: {s}\n", .{ h.name, h.value });
return .{ .body = out.items };
}
const Fixture = struct {
arena_state: std.heap.ArenaAllocator,
root: support.Root,
web: support.Background,
cache: Cache,
proxy: proxy.Proxy,
listenfd: c.fd_t,
port: u16,
/// Порт TINY: после `web.join` сервера уже нет, а номер нужен.
web_port: u16,
thread: std.Thread,
fn arena(f: *Fixture) std.mem.Allocator {
return f.arena_state.allocator();
}
};
fn proxyServe(f: *Fixture, n: usize) void {
conc.echo_pre.serve(testing.allocator, io, f.listenfd, .{ .workers = 4, .max_clients = n }, &f.proxy, proxy.handle) catch |err| std.debug.print("proxy: {t}\n", .{err});
}
/// GET через прокси: абсолютный URI и свой заголовок, ответ целиком.
fn viaProxy(f: *Fixture, path: []const u8) !support.Reply {
const raw = try std.fmt.allocPrint(f.arena(), "GET http://127.0.0.1:{d}{s} HTTP/1.1\r\nHost: 127.0.0.1:{d}\r\nX-Test: 73\r\nConnection: keep-alive\r\n\r\n", .{ f.web_port, path, f.web_port });
return support.exchange(f.arena(), f.port, raw);
}
const Fetch = struct {
f: *Fixture,
path: []const u8,
reply: support.Reply = undefined,
ok: bool = false,
fn run(job: *Fetch) void {
job.reply = viaProxy(job.f, job.path) catch return;
job.ok = true;
}
};
test "прокси перед TINY: заголовки переписаны, конкурентные клиенты, кэш после остановки TINY" {
const f = try testing.allocator.create(Fixture);
defer testing.allocator.destroy(f);
f.arena_state = .init(testing.allocator);
defer f.arena_state.deinit();
f.root = try .create(f.arena());
defer f.root.cleanup();
// Файл больше лимита объекта: его прокси отдаст, но не закэширует.
const big = try f.arena().alloc(u8, 3000);
@memset(big, 'B');
try f.root.tmp.dir.writeFile(io, .{ .sub_path = "big.txt", .data = big });
// TINY обслужит ровно пять соединений и остановится.
const routes = [_]tiny.Route{.{ .method = .GET, .path = "/headers", .handle = echoHeaders }};
try f.web.start(testing.allocator, .{ .port = 0, .root = f.root.path, .routes = &routes, .quiet = true }, 5);
f.web_port = f.web.server.port;
f.cache = .init(testing.allocator, io, .{ .cache = 8192, .object = 1024 });
defer f.cache.deinit();
f.proxy = .{ .gpa = testing.allocator, .cache = &f.cache, .quiet = true };
f.listenfd = try socket.openListenfd(0);
defer socket.close(f.listenfd);
f.port = socket.localPort(f.listenfd).?;
const proxy_conns = 1 + 4 + 5;
f.thread = try std.Thread.spawn(.{}, proxyServe, .{ f, proxy_conns });
// 1. Что видит сервер за прокси.
const seen = try viaProxy(f, "/headers");
try testing.expectEqual(200, seen.status);
try testing.expect(std.mem.startsWith(u8, seen.body, "GET /headers HTTP/1.0\n"));
try testing.expect(std.mem.indexOf(u8, seen.body, "User-Agent: " ++ proxy.user_agent ++ "\n") != null);
try testing.expect(std.mem.indexOf(u8, seen.body, "Connection: close\n") != null);
try testing.expect(std.mem.indexOf(u8, seen.body, "Proxy-Connection: close\n") != null);
try testing.expect(std.mem.indexOf(u8, seen.body, "X-Test: 73\n") != null);
try testing.expect(std.mem.indexOf(u8, seen.body, "keep-alive") == null);
// 2. Четыре клиента одновременно, каждый за своим файлом.
var jobs = [_]Fetch{
.{ .f = f, .path = "/home.html" },
.{ .f = f, .path = "/dot.png" },
.{ .f = f, .path = "/hello.txt" },
.{ .f = f, .path = "/big.txt" },
};
var threads: [jobs.len]std.Thread = undefined;
for (&threads, &jobs) |*t, *job| t.* = try std.Thread.spawn(.{}, Fetch.run, .{job});
for (threads) |t| t.join();
for (jobs) |job| {
try testing.expect(job.ok);
try testing.expectEqual(200, job.reply.status);
}
const png = try std.Io.Dir.cwd().readFileAlloc(io, support.www_dir ++ "/dot.png", f.arena(), .limited(1 << 16));
try testing.expectEqualSlices(u8, png, jobs[1].reply.body);
try testing.expectEqual(@as(usize, 3000), jobs[3].reply.body.len);
// 3. TINY отработал свои пять соединений и закрыл сокет.
f.web.join();
for ([_][]const u8{ "/home.html", "/dot.png", "/hello.txt", "/headers" }, 0..) |path, i| {
const again = try viaProxy(f, path);
try testing.expectEqual(200, again.status);
if (i < 3) try testing.expectEqualSlices(u8, jobs[i].reply.body, again.body);
}
// Большой файл не кэшировался, а сервера больше нет: 502.
const gone = try viaProxy(f, "/big.txt");
try testing.expectEqual(502, gone.status);
f.thread.join();
try testing.expectEqual(@as(usize, 4), f.cache.hits.load(.monotonic));
}
Последний тест стоит прочитать внимательно, он и есть приёмка лабораторной. TINY обслуживает ровно пять соединений и останавливается. Первое соединение показывает, что сервер за прокси видит: HTTP/1.0, наш User-Agent, Connection: close, свой заголовок клиента X-Test и ни следа keep-alive. Четыре клиента одновременно забирают три маленьких файла и один на 3000 байт, больше лимита объекта в этом тесте. Потом TINY уже нет, а три маленьких файла и /headers прокси отдаёт из кэша байт в байт, большой файл получает 502. Четыре попадания, ни одного лишнего.
Тест с восемью потоками против кэша на 400 байт проверяет то, что нельзя увидеть глазами: объект ключа kN это N + 1 байт со значением N, и чужие или обрезанные байты значили бы, что копия не защищена замком.
$ zig build test -Dstep=73 --summary all
Build Summary: 6/6 steps succeeded; 11/11 tests passed
test success
+- run test 11 pass (11 total) 7s MaxRSS:129M
На macOS все 11 зелёные. В контейнере с Linux (runner-zig:dev, arm64) 10 зелёных и один пропущен: Io.Evented там не собирается.
Прогон
Прокси перед TINY на потоках, в двух терминалах плюс третий для curl. curl -x говорит «ходи через этот прокси»:
$ tiny tiny --threads 8000 zig-out
$ tiny proxy 15300
$ curl -si -x http://localhost:15300 http://localhost:8000/home.html | head -5
HTTP/1.0 200 OK
Server: Tiny Web Server
Connection: close
Content-length: 114
Content-type: text/html
$ curl -s -x http://localhost:15300 'http://localhost:8000/cgi-bin/adder?1&2' | sed -n 2p
<p>The answer is: 1 + 2 = 3
$ curl -si -x http://localhost:15300 http://localhost:8000/nope.html | head -1
HTTP/1.0 404 Not Found
# TINY остановлен
$ curl -si -x http://localhost:15300 http://localhost:8000/home.html | head -1
HTTP/1.0 200 OK
$ curl -si -x http://localhost:15300 http://localhost:8000/nope.html | head -1
HTTP/1.0 502 Bad Gateway
Лог прокси в это время:
proxy: listening on port 15300
127.0.0.1:56664 "GET http://localhost:8000/home.html" MISS 223 cached
127.0.0.1:56668 "GET http://localhost:8000/home.html" HIT 223
127.0.0.1:56670 "GET http://localhost:8000/cgi-bin/adder?1&2" MISS 211 cached
127.0.0.1:56674 "GET http://localhost:8000/nope.html" MISS 266
127.0.0.1:56678 "GET http://localhost:8000/home.html" HIT 223
127.0.0.1:56680 "GET http://localhost:8000/nope.html" 502
127.0.0.1:56684 "GET http://localhost:8000/dot.png" 502
home.html пережил остановку TINY, nope.html нет: ответ 404 не кэшировался, и повтор пошёл к серверу, которого уже нет.
Под нагрузкой
Теперь генератор нагрузки из урока 71. loadgen шлёт GET <path>, и если в path положить абсолютный URI, он становится клиентом прокси. 16 соединений, по 500 запросов на соединение (на прямой CGI по 20: он медленный). TINY на пуле из восьми потоков, прокси с восемью воркерами, сборка ReleaseFast. Apple M4 Max, 16 ядер, macOS 26.6.2, Zig 0.16.0, три прогона подряд; машину в это время грузили посторонние процессы, load average около 18, поэтому смотри на порядок, а не на третью цифру.
| Что | Прогон 1 | Прогон 2 | Прогон 3 |
|---|---|---|---|
| статика напрямую, запросов в секунду | 27288 | 19736 | 22219 |
| статика через прокси, из кэша | 24300 | 22258 | 19778 |
CGI adder напрямую | 1542 | 1618 | 1663 |
CGI adder через прокси, из кэша | 15518 | 19687 | 18015 |
Статика из кэша прокси идёт с той же скоростью, что прямо из TINY: TINY отдаёт маленький файл через mmap за те же микросекунды, за которые прокси копирует объект из кэша. А ответ CGI из кэша на порядок быстрее, чем у самого сервера: fork, execve и waitpid на каждый запрос прокси больше не делает. Медиана задержки CGI упала с 9,5 мс до 0,8 мс. Кэш окупается там, где ответ дорого делать, а не там, где его дорого передавать.
И обещанная картина из виджета. В каждом прогоне лог прокси показал по восемь промахов на каждый из двух адресов:
127.0.0.1:49730 "GET http://127.0.0.1:18000/home.html" MISS 223 cached
127.0.0.1:49734 "GET http://127.0.0.1:18000/home.html" MISS 223 cached
127.0.0.1:49732 "GET http://127.0.0.1:18000/home.html" MISS 223 cached
127.0.0.1:49733 "GET http://127.0.0.1:18000/home.html" MISS 223 cached
127.0.0.1:49736 "GET http://127.0.0.1:18000/home.html" MISS 223 cached
127.0.0.1:49731 "GET http://127.0.0.1:18000/home.html" MISS 223 cached
127.0.0.1:49735 "GET http://127.0.0.1:18000/home.html" MISS 223 cached
127.0.0.1:49737 "GET http://127.0.0.1:18000/home.html" MISS 223 cached
Восемь воркеров одновременно взяли по первому запросу, все восемь заглянули в пустой кэш, все восемь пошли на сервер за одним и тем же файлом. Это называют давкой в кэше. Корректности это не вредит: put второго и дальнейших только освежает метку. Но серверу пришлось восемь раз сделать работу, которую прокси обещал делать один раз, и для дорогого CGI это заметно. Лечение, одиночный полёт на Io.Condition, в домашнем задании.
Шаг zbox: раннер на std.Io с крайним сроком
zbox в уроке 71 научился держать нагрузку: главный поток делает accept с close-on-exec, кладёт дескриптор в очередь на Io.Mutex и двух Io.Condition, восемь потоков пула разбирают запросы, семафор Slots пускает к песочнице не больше --jobs задач, а сама задача это дочерний процесс zbox job, потому что run() держит состояние в глобальных переменных и сигналах процесса. У этой модели нет одного: задачу нельзя остановить. Если чужая программа зависла в сборке или ждёт ввода, у неё есть только лимит времени самой программы, а у сервера нет способа сказать «эта задача мне больше не нужна».
Сегодня тот же сервис переезжает на std.Io, и отмена появляется бесплатно, вместе с крайним сроком на всю задачу.
zbox serve PORT --io threaded [--jobs N] [--queue N] [--deadline-ms N]
zbox serve PORT --io evented # код 3 и объяснение
Что меняется
- Своих потоков и очереди нет.
acceptизstd.Io.net, каждое соединение это член одногоIo.Groupчерезgroup.concurrent. Очередь из шага 71 превратилась в счётчикinflight: соединений, которые ждут места или выполняются, не большеjobs + queue, остальным тот же 503 сRetry-After. Slotsтот же. Семафор наIo.Semaphoreиз шага 71 с самого начала принималioи возвращалCancelable!voidизacquire: ожидание места это точка отмены.- Задача гонится с таймером.
Io.Selectс двумя полями:doneэто результат задачи,deadlineэтоIo.sleepна крайний срок. Первый результат и есть ответ, второго отменяетcancelDiscard. - Отмена доходит до программы. Через четыре звена, о них ниже.
server.zig
Файл целиком: модель шага 71 (Service, runSlot, пул, acceptCloexec) ты уже видел, новое в нём это поля timeouts и inflight, развилка по модели в route, Race и runDeadline, счётчик timeouts в ответе /stats, serveIo с ioConnection и eventedVerdict в конце.
//! `zbox serve` под нагрузкой: две модели на одном маршруте `POST /run`.
//!
//! Шаг 71, `servePool`: главный поток делает `accept` и кладёт дескриптор
//! в очередь, потоки пула разбирают запрос и запускают задачу. Задач
//! одновременно не больше `jobs` (семафор `Slots`), соединений в очереди
//! не больше `queue`; сверх этого главный поток сам отвечает 503.
//!
//! Шаг 73, `serveIo`: то же на `std.Io`. Соединение это член `Io.Group`,
//! задача гонится с таймером через `Io.Select`, и проигравшего отменяют.
//! Никаких своих потоков и очереди: их роль играет реализация `Io`.
const std = @import("std");
const builtin = @import("builtin");
const c = std.c;
const Io = std.Io;
const net = Io.net;
const tiny = @import("tiny");
const serve = @import("../box/serve.zig");
const http = @import("http.zig");
const job = @import("job.zig");
const pool_mod = @import("pool.zig");
/// Клиент, который молчит дольше, теряет соединение.
const socket_timeout_ms = 30_000;
/// Общее для обеих моделей: как выполнить `POST /run` и что отдать
/// на `GET /stats`.
pub const Service = struct {
io: Io,
/// Путь к самому zbox: задачи идут дочерними `zbox job`.
exe: []const u8,
args: serve.ServeArgs,
slots: pool_mod.Slots,
/// Отказы 503: очередь полна или некому принять соединение.
rejected: std.atomic.Value(u64) = .init(0),
/// Задачи, отменённые по крайнему сроку (только `--io`).
timeouts: std.atomic.Value(u64) = .init(0),
/// Соединения, которые приняты и ещё не закрыты (только `--io`).
inflight: std.atomic.Value(u32) = .init(0),
pub fn init(io: Io, exe: []const u8, args: serve.ServeArgs) Service {
return .{ .io = io, .exe = exe, .args = args, .slots = .init(args.jobs) };
}
fn route(svc: *Service, arena: std.mem.Allocator, req: http.Request) anyerror!http.Response {
if (req.method == .GET and std.mem.eql(u8, req.path, "/stats")) return svc.stats(arena);
if (req.method != .POST or !std.mem.eql(u8, req.path, "/run")) {
return .{ .status = 404, .body = "{\"error\":\"NoSuchRoute\"}\n" };
}
// Кривой запрос отбиваем до очереди к песочнице: 400 не ждёт места.
const request = serve.parseRequest(arena, req.body) catch |err| switch (err) {
error.OutOfMemory => return err,
else => return .{ .status = 400, .body = try std.fmt.allocPrint(arena, "{{\"error\":\"{t}\"}}\n", .{err}) },
};
return switch (svc.args.model) {
.pool => svc.runSlot(arena, req.body),
else => svc.runDeadline(arena, req.body, request.deadline_ms orelse svc.args.deadline_ms),
};
}
/// Шаг 71: занять место и выполнить задачу до конца.
fn runSlot(svc: *Service, arena: std.mem.Allocator, body: []const u8) !http.Response {
try svc.slots.acquire(svc.io);
defer svc.slots.release(svc.io);
return svc.runJob(arena, body);
}
const Race = union(enum) {
done: anyerror!http.Response,
deadline: Io.Cancelable!void,
};
/// Шаг 73: задача и таймер стартуют вместе, кто первый, тот и ответ.
/// Второго отменяем. Отменённая задача получает `error.Canceled` в
/// ближайшей точке отмены (чтение ответа или `wait`), и `job.run`
/// убивает свой `zbox job`. Крайний срок считаем с момента, когда
/// задача получила место: ожидание в очереди в него не входит.
fn runDeadline(svc: *Service, arena: std.mem.Allocator, body: []const u8, deadline_ms: u32) !http.Response {
try svc.slots.acquire(svc.io);
defer svc.slots.release(svc.io);
var buffer: [2]Race = undefined;
var select = Io.Select(Race).init(svc.io, &buffer);
// `concurrent`, а не `async`: `async` в `Io.Threaded` при занятых
// потоках выполнит задачу прямо здесь, и таймер стартует после неё.
try select.concurrent(.done, runJob, .{ svc, arena, body });
select.concurrent(.deadline, Io.sleep, .{ svc.io, .fromMilliseconds(deadline_ms), .awake }) catch |err| {
select.cancelDiscard();
return err;
};
const first = try select.await();
// Проигравший ещё идёт: отменяем и ждём, пока он действительно
// закончится. Результат задачи после отмены не нужен.
select.cancelDiscard();
return switch (first) {
.done => |result| result,
.deadline => blk: {
_ = svc.timeouts.fetchAdd(1, .monotonic);
break :blk .{ .status = 504, .body = try std.fmt.allocPrint(arena, "{{\"error\":\"timeout\",\"deadline_ms\":{d}}}\n", .{deadline_ms}) };
},
};
}
fn runJob(svc: *Service, arena: std.mem.Allocator, body: []const u8) anyerror!http.Response {
return .{ .body = try job.run(svc.io, arena, svc.exe, svc.args.isolate, body) };
}
fn stats(svc: *Service, arena: std.mem.Allocator) !http.Response {
return .{ .body = try std.fmt.allocPrint(arena, "{{\"running\":{d},\"peak\":{d},\"done\":{d},\"rejected\":{d},\"timeouts\":{d}}}\n", .{
svc.slots.running.load(.acquire),
svc.slots.peak.load(.acquire),
svc.slots.done.load(.acquire),
svc.rejected.load(.acquire),
svc.timeouts.load(.acquire),
}) };
}
};
// Шаг 71: пул потоков.
const Pool = pool_mod.Pool(c.fd_t, *Service, connection);
/// Обслужить одно соединение в потоке пула. Своя арена на соединение:
/// общего аллокатора между потоками здесь нет, кроме `gpa` под ареной.
fn connection(svc: *Service, fd: c.fd_t) void {
defer tiny.socket.close(fd);
http.setTimeouts(fd, socket_timeout_ms);
var arena_state: std.heap.ArenaAllocator = .init(std.heap.page_allocator);
defer arena_state.deinit();
var in_buf: [8192]u8 = undefined;
var out_buf: [8192]u8 = undefined;
var reader: tiny.fdio.Reader = .init(fd, &in_buf);
var writer: tiny.fdio.Writer = .init(fd, &out_buf);
http.transaction(arena_state.allocator(), &reader.interface, &writer.interface, serve.max_body, svc, Service.route) catch {};
}
pub fn servePool(svc: *Service, gpa: std.mem.Allocator, log: *Io.Writer) !void {
tiny.server.ignoreSigpipe();
const listenfd = try tiny.socket.openListenfd(svc.args.port);
defer tiny.socket.close(listenfd);
setCloexec(listenfd);
const queue_buffer = try gpa.alloc(c.fd_t, svc.args.queue);
defer gpa.free(queue_buffer);
const threads = try gpa.alloc(std.Thread, svc.args.workers);
defer gpa.free(threads);
var pool: Pool = undefined;
try pool.start(svc.io, svc, queue_buffer, threads);
defer pool.stop();
try announce(log, tiny.socket.localPort(listenfd) orelse svc.args.port, svc.args);
while (true) {
const fd = acceptCloexec(listenfd) orelse continue;
if (pool.submit(fd)) continue;
_ = svc.rejected.fetchAdd(1, .monotonic);
http.reject(fd, http.busy);
tiny.socket.close(fd);
}
}
/// Принятый сокет обязан закрываться при `exec`. Иначе каждый `zbox job`,
/// запущенный соседним потоком, унесёт с собой копию чужого соединения,
/// и `close` в сервере не закроет его: клиент дождётся конца ответа только
/// вместе с последней задачей пачки. `std.Io.net` ставит флаг сам, а
/// `accept` из libc нет.
///
/// В Linux флаг ставит сам `accept4(SOCK_CLOEXEC)`, атомарно. В macOS
/// `accept4` нет, поэтому `accept` и следом `fcntl`.
/// ponytail: в macOS между ними окно: `fork` соседа в этот миг унесёт копию.
fn acceptCloexec(listenfd: c.fd_t) ?c.fd_t {
if (comptime builtin.os.tag == .linux) {
const fd = c.accept4(listenfd, null, null, c.SOCK.CLOEXEC);
return if (fd < 0) null else fd;
}
const fd = c.accept(listenfd, null, null);
if (fd < 0) return null;
setCloexec(fd);
return fd;
}
fn setCloexec(fd: c.fd_t) void {
_ = c.fcntl(fd, c.F.SETFD, @as(c_int, c.FD_CLOEXEC));
}
// Шаг 73: `std.Io`.
pub fn serveIo(svc: *Service, log: *Io.Writer) !void {
const io = svc.io;
tiny.server.ignoreSigpipe();
const address: net.IpAddress = .{ .ip4 = .unspecified(svc.args.port) };
var server = try address.listen(io, .{ .reuse_address = true });
defer server.deinit(io);
try announce(log, server.socket.address.getPort(), svc.args);
// Все соединения в одной группе. Выход из цикла отменяет всех разом.
var group: Io.Group = .init;
defer group.cancel(io);
// Очередь из 71 превратилась в счётчик: соединений, которые ждут места
// или выполняются, не больше `jobs + queue`, остальным 503.
const limit: u32 = @as(u32, svc.args.jobs) + svc.args.queue;
while (true) {
const stream = server.accept(io) catch |err| switch (err) {
error.Canceled => return err,
else => continue,
};
if (svc.inflight.fetchAdd(1, .acq_rel) < limit) {
if (group.concurrent(io, ioConnection, .{ svc, stream })) continue else |_| {}
}
_ = svc.inflight.fetchSub(1, .acq_rel);
_ = svc.rejected.fetchAdd(1, .monotonic);
http.reject(stream.socket.handle, http.busy);
stream.close(io);
}
}
fn ioConnection(svc: *Service, stream: net.Stream) void {
const io = svc.io;
defer _ = svc.inflight.fetchSub(1, .acq_rel);
defer stream.close(io);
http.setTimeouts(stream.socket.handle, socket_timeout_ms);
var arena_state: std.heap.ArenaAllocator = .init(std.heap.page_allocator);
defer arena_state.deinit();
var in_buf: [8192]u8 = undefined;
var out_buf: [8192]u8 = undefined;
var reader = stream.reader(io, &in_buf);
var writer = stream.writer(io, &out_buf);
http.transaction(arena_state.allocator(), &reader.interface, &writer.interface, serve.max_body, svc, Service.route) catch {};
}
/// Порт в stderr первой строкой, как у шага 66: тесты его оттуда читают.
fn announce(log: *Io.Writer, port: u16, args: serve.ServeArgs) !void {
try log.print("zbox: listening on port {d} ({t}, jobs {d}, queue {d})\n", .{ port, args.model, args.jobs, args.queue });
try log.flush();
}
/// `serve --io evented`: `Io.Evented` в Zig 0.16.0 есть, но сервер на нём
/// не поднять. Проверяем честно: пробуем открыть сокет и печатаем ответ.
pub fn eventedVerdict(gpa: std.mem.Allocator, port: u16, log: *Io.Writer) !void {
if (comptime builtin.os.tag == .macos) {
// Dispatch поверх GCD. `deinit` в 0.16.0 не компилируется
// (`free` многоэлементного указателя), поэтому не зовём его:
// процесс сразу после проверки выходит.
var evented: Io.Evented = undefined;
try evented.init(gpa, .{});
const address: net.IpAddress = .{ .ip4 = .unspecified(port) };
if (address.listen(evented.io(), .{ .reuse_address = true })) |_| {
try log.writeAll("zbox: Io.Evented открыл сокет, эта версия std уже умеет сеть\n");
} else |err| {
try log.print("zbox: Io.Evented (Dispatch) в Zig 0.16.0 без сети: listen ответил {t}\n", .{err});
}
} else {
// Uring в 0.16.0 не компилируется: `dirOpenDir` возвращает ошибку
// `ReadOnlyFileSystem`, которой нет в объявленном множестве. А сеть
// у него та же заглушка `NetworkDown`, что и у Dispatch.
try log.writeAll("zbox: Io.Evented (Uring) в Zig 0.16.0 не собирается, а сети в нём нет\n");
}
try log.flush();
}
Самое интересное здесь runDeadline, в нём три решения.
concurrent, а неasync, у обеих гонщиц.asyncвIo.Threadedпри занятых потоках выполнитrunJobпрямо здесь, иselect.async(.deadline, ...)начнётся только после того, как задача закончится сама. Таймер, который стартует после задачи, её не прервёт. Это тот самый случай «корректность зависит от одновременности». Если втораяconcurrentне удалась, первую надо отменить до выхода, отсюдаcancelDiscardв ветке ошибки.- Крайний срок считается с получения места.
slots.acquireстоит доSelect. Задача, которая ждала в очереди, не должна проиграть таймеру, не начавшись. Это выбор, а не закон: крайний срок «с прихода запроса» тоже осмыслен, тогда ожидание места стоит внутри гонки. cancelDiscardсразу после первого результата. Проигравший ещё идёт. Если это таймер, отмена разбудит егоsleep. Если это задача, отмена пройдёт по цепочке до чужой программы.cancelDiscardждёт, пока проигравший действительно закончится, так что к ответу 504 процесс задачи уже мёртв и временный каталог убран.
И serveIo: на отказ 503 здесь нет своего потока, соединение сверх лимита отбивает сам цикл accept. fetchAdd возвращает старое значение, поэтому проверка «меньше лимита» и увеличение счётчика это одна атомарная операция: два соединения одновременно не пройдут на последнее место. Если group.concurrent не дал потока, соединение тоже получает 503, а счётчик откатывается.
Цепочка отмены
std.Io ничего не знает о процессах, которые задача запустила. Отмена для него это error.Canceled из ближайшей точки отмены в самой задаче. Дальше за дело берутся наши errdefer и сигналы из урока 49.
select.cancelDiscard()просит отмену у задачиrunJob.- Поток задачи спит в ядре: в
readиз пайпа (ждёт ответzbox job) или вwait4.Io.Threadedшлёт емуSIGIO, вызов возвращаетEINTR,std.Ioпревращает его вerror.Canceled. И чтение, иwaitидут черезio, потому чтоjob.runиз шага 71 запускает процесс черезstd.process.spawn(io, ...), а не черезforkиз libc. Это решение шага 71 окупилось сегодня. errdefer child.kill(io)вjob.runшлётzbox jobSIGTERM и ждёт его.- Чужая программа живёт в своей группе процессов (так её отделил лимит времени в уроке 49), SIGTERM её не касается. Поэтому
zbox jobставит обработчик SIGTERM, который убивает группу программы.run()видит обычное завершение ребёнка по сигналу, убирает временный каталог,zbox jobвыходит.
Код job.run менять не пришлось, только его комментарии: теперь они объясняют, почему процесс запущен через Io и кто убьёт программу.
@@
-//! Задача сервера это отдельный процесс `zbox job` (шаг 71).
+//! Задача сервера это отдельный процесс `zbox job` (шаги 71 и 73).
//!
@@
//! глобальные переменные, свои сигналы и свой будильник.
+//!
+//! Процесс запускаем через `std.process.spawn(io, ...)`: и чтение ответа,
+//! и `wait` проходят через `Io`, поэтому они точки отмены (шаг 73).
@@
/// Запускает `exe job [--isolate ROOTFS]`, отдаёт ему тело запроса в stdin
-/// и возвращает его stdout: готовую строку JSON с полем `stage`. Ошибка
-/// по дороге убивает ребёнка: `kill` шлёт ему SIGTERM и ждёт.
+/// и возвращает его stdout: готовую строку JSON с полем `stage`.
+///
+/// Отмена (`error.Canceled` из записи, чтения или `wait`) убивает ребёнка:
+/// `kill` шлёт ему SIGTERM и ждёт. Программу в её собственной группе
+/// процессов сигнал не заденет, поэтому `zbox job` ловит SIGTERM и сам
+/// убивает группу программы (`watchdog.catchTerm`).
pub fn run(io: Io, arena: std.mem.Allocator, exe: []const u8, rootfs: ?[]const u8, body: []const u8) ![]u8 {
Звено 4 это правка в watchdog.zig:
@@
var deadline_ms: std.atomic.Value(u64) = .init(0);
+/// Пришёл SIGTERM: задачу отменили снаружи (шаг 73). Флаг не сбрасывается.
+var terminated: std.atomic.Value(bool) = .init(false);
@@
killGroup();
+}
+
+fn onTerm(_: posix.SIG) callconv(.c) void {
+ terminated.store(true, .seq_cst);
+ killGroup();
+}
+
+/// Шаг 73: `zbox job` ловит SIGTERM. Сервер отменяет задачу и шлёт сигнал
+/// только своему ребёнку, а программа сидит в своей группе процессов и
+/// сама бы не узнала. Обработчик убивает группу, `run` возвращается как
+/// обычно, временный каталог убирается. Отмена до fork не теряется:
+/// `startTimer` проверяет флаг, как только группа известна.
+pub fn catchTerm() void {
+ posix.sigaction(.TERM, &.{
+ .handler = .{ .handler = onTerm },
+ .mask = posix.sigemptyset(),
+ .flags = 0,
+ }, null);
}
@@
group.store(child, .seq_cst);
+ if (terminated.load(.seq_cst)) killGroup();
Обработчик делает только разрешённое в обработчике: атомарную запись флага и kill, это правила async-signal-safe из урока 49. Последняя строка закрывает ещё одну гонку того же урока: SIGTERM может прийти, пока run() только собирается делать fork. Тогда группы ещё нет и killGroup ничего не убьёт. Флаг запоминает отмену, и как только группа известна, startTimer добивает её сам.
Запрос и флаги
Поле deadline_ms в запросе переопределяет --deadline-ms сервера. Ноль запрещён, как и нулевые лимиты: «крайний срок ноль» скорее опечатка, чем желание.
@@
argv: []const []const u8 = &.{},
+ /// Крайний срок всей задачи на сервере `--io` (шаг 73), миллисекунды.
+ /// Сборка плюс запуск; по истечении задачу отменяют, ответ 504.
+ deadline_ms: ?u32 = null,
};
@@
if (request.source.len == 0 and request.argv.len == 0) return error.EmptySource;
- if (request.time_ms == 0 or request.mem_mb == 0) return error.BadLimit;
+ if (request.time_ms == 0 or request.mem_mb == 0 or request.deadline_ms == 0) return error.BadLimit;
return request;
@@
/// Как `zbox serve` держит нагрузку. Шаг 66: итеративно, по запросу за раз.
-/// Шаг 71: главный поток принимает, пул потоков обслуживает.
-pub const Model = enum { iterative, pool };
+/// Шаг 71: главный поток принимает, пул потоков обслуживает. Шаг 73: то же
+/// на `std.Io`, задача с крайним сроком и отменой.
+pub const Model = enum { iterative, pool, io_threaded, io_evented };
@@
queue: u16 = 16,
+ /// Крайний срок задачи по умолчанию для `--io`, миллисекунды.
+ deadline_ms: u32 = 90_000,
};
-/// `serve <port> [--isolate ROOTFS] [--jobs N] [--workers N] [--queue N]`.
-/// Любой из флагов пула включает пул.
+/// `serve <port> [--isolate ROOTFS] [--jobs N] [--workers N] [--queue N]
+/// [--io threaded|evented] [--deadline-ms N]`. Любой из флагов пула
+/// включает пул, `--io` включает модель шага 73.
pub fn parseServe(argv: []const []const u8) ArgsError!ServeArgs {
@@
result.isolate = value;
+ } else if (std.mem.eql(u8, flag, "--io")) {
+ result.model = if (std.mem.eql(u8, value, "threaded")) .io_threaded else if (std.mem.eql(u8, value, "evented")) .io_evented else return error.BadUsage;
+ } else if (std.mem.eql(u8, flag, "--deadline-ms")) {
+ result.deadline_ms = try positive(u32, value);
} else {
По умолчанию крайний срок 90 секунд: больше, чем лимит компилятора (60 с), иначе медленная сборка на нагруженной машине получала бы 504 раньше, чем свой честный ответ.
@@
\\ пул потоков: N задач одновременно, очередь соединений, сверх неё 503
+ \\ zbox serve <port> --io threaded|evented [--jobs N] [--queue N] [--deadline-ms N]
+ \\ то же на std.Io: задача с крайним сроком, по нему отмена и 504
\\ zbox job [--isolate ROOTFS]
@@
.iterative => serve(init, serve_args.port),
- .pool => servePool(init, serve_args),
+ .io_evented => {
+ try zbox.server.eventedVerdict(init.gpa, serve_args.port, out);
+ std.process.exit(3);
+ },
+ .pool, .io_threaded => serveConcurrent(init, serve_args),
};
@@
-/// Шаг 71: пул потоков. Задачи идут дочерними `zbox job`, поэтому
-/// серверу нужен путь к самому себе.
-fn servePool(init: std.process.Init, args: zbox.serve.ServeArgs) !void {
+/// Шаги 71 и 73: пул потоков или `std.Io`. Задачи идут дочерними
+/// `zbox job`, поэтому серверу нужен путь к самому себе.
+fn serveConcurrent(init: std.process.Init, args: zbox.serve.ServeArgs) !void {
const exe = try std.process.executablePathAlloc(init.io, init.arena.allocator());
@@
var stderr = std.Io.File.stderr().writerStreaming(init.io, &err_buf);
- return zbox.server.servePool(&service, init.gpa, &stderr.interface);
+ return switch (args.model) {
+ .pool => zbox.server.servePool(&service, init.gpa, &stderr.interface),
+ else => zbox.server.serveIo(&service, &stderr.interface),
+ };
}
@@
/// `stage` в stdout. Сервер уже проверил запрос, но процесс мог позвать и
-/// человек, поэтому проверяем ещё раз.
+/// человек, поэтому проверяем ещё раз. SIGTERM от сервера (отмена по
+/// крайнему сроку) убивает программу вместе с её группой процессов.
fn job(init: std.process.Init, args: []const []const u8, out: *std.Io.Writer) !void {
@@
};
+ zbox.watchdog.catchTerm();
var in_buf: [4096]u8 = undefined;
Мелочи вокруг: номер шага в сборке, комментарии и поле timeouts в разборе /stats для тестов.
@@
/// соответствует файл `tests/step_NN.zig`, и все они зелёные на финале.
-const project_steps = [_]u8{ 48, 49, 54, 61, 64, 66, 67, 71 };
+const project_steps = [_]u8{ 48, 49, 54, 61, 64, 66, 67, 71, 73 };
@@
const tiny = b.dependency("tiny", .{ .target = target, .optimize = optimize }).module("tiny");
- // Пул потоков (шаг 71) разбирает HTTP частями TINY.
+ // Пул и сервер на std.Io (шаги 71 и 73) разбирают HTTP частями TINY.
zbox.addImport("tiny", tiny);
@@
- // Нагрузка на `zbox serve`: N клиентов разом (шаг 71). В тесты не входит.
+ // Нагрузка для таблицы трёх моделей (README, шаг 73). В тесты не входит.
const load = b.addExecutable(.{
@@
-// Шаг 71: сервер под нагрузкой.
+// Шаги 71 и 73: сервер под нагрузкой.
pub const queue = @import("pool/queue.zig");
@@
-pub const Stats = struct { peak: u32, done: u64, rejected: u64 };
-
-/// `GET /stats` сервера шага 71.
+pub const Stats = struct { peak: u32, done: u64, rejected: u64, timeouts: u64 };
+
+/// `GET /stats` серверов шагов 71 и 73.
pub fn stats(arena: std.mem.Allocator, port: u16) !Stats {
Io.Evented: честный отказ
--io evented не притворяется. eventedVerdict на macOS поднимает Io.Evented (это Io.Dispatch) и пробует открыть слушающий сокет. В 0.16.0 listen отвечает error.NetworkDown, и zbox печатает это и выходит с кодом 3. На Linux Io.Uring не компилируется, поэтому ветка Linux выбрана на этапе компиляции и отвечает без него. Когда Evented получит сеть, eventedVerdict сам напишет «эта версия std уже умеет сеть», а serveIo заработает на нём без правок: Service, job.run и Slots принимают любой std.Io.
Проверяли и обходной путь, как serveFd у эха: сокеты из libc, а задачи на файберах. Он упирается в третью ловушку из таблицы выше: отмена задачи, которая ждёт в futex (а Select с таймером ждёт именно там), роняет процесс на Io.Dispatch с Segmentation fault at address 0xaaaaaaaaaaaaaa08 в DoublyLinkedList.remove. Отладочная сборка Zig заполняет undefined байтами 0xaa, так что адрес 0xaaaa...08 значит: реализация прочла поле по указателю, который так и остался undefined, в списке ожидающих futex. Раннер без отмены нам не нужен, так что честный отказ лучше полурабочего сервера.
Тесты шага
//! Шаг 73: тот же сервис на `std.Io`. Задача и таймер гонятся через
//! `Io.Select`, проигравшего отменяют; отмена доходит до процесса
//! задачи через SIGTERM, а до программы через обработчик в `zbox job`.
const std = @import("std");
const builtin = @import("builtin");
const build_options = @import("build_options");
const zbox = @import("zbox");
const support = @import("support.zig");
const testing = std.testing;
const io = testing.io;
test "parseServe: --io и крайний срок" {
const threaded = try zbox.serve.parseServe(&.{ "serve", "0", "--io", "threaded", "--jobs", "3", "--deadline-ms", "500" });
try testing.expectEqual(.io_threaded, threaded.model);
try testing.expectEqual(3, threaded.jobs);
try testing.expectEqual(500, threaded.deadline_ms);
// Порядок флагов не важен: --jobs после --io не возвращает пул.
const evented = try zbox.serve.parseServe(&.{ "serve", "0", "--jobs", "2", "--io", "evented" });
try testing.expectEqual(.io_evented, evented.model);
try testing.expectError(error.BadUsage, zbox.serve.parseServe(&.{ "serve", "0", "--io", "fibers" }));
try testing.expectError(error.BadUsage, zbox.serve.parseServe(&.{ "serve", "0", "--deadline-ms", "0" }));
}
/// Жив ли процесс. `kill` с нулевым сигналом ничего не шлёт, только
/// проверяет, что такой pid есть.
fn alive(pid: std.c.pid_t) bool {
return std.c.kill(pid, @enumFromInt(0)) == 0;
}
/// Pid программы, которая записала его в файл первой командой.
fn readPid(path: []const u8) !std.c.pid_t {
var buf: [32]u8 = undefined;
const text = try std.Io.Dir.cwd().readFile(io, path, &buf);
return std.fmt.parseInt(std.c.pid_t, std.mem.trimEnd(u8, text, "\n"), 10);
}
fn victimJson(arena: std.mem.Allocator, pid_path: []const u8, deadline_ms: ?u32) ![]const u8 {
const script = try std.fmt.allocPrint(arena, "echo $$ > {s}; exec sleep 30", .{pid_path});
return std.json.Stringify.valueAlloc(arena, .{
.argv = &[_][]const u8{ "sh", "-c", script },
.time_ms = 60_000,
.deadline_ms = deadline_ms,
}, .{ .emit_null_optional_fields = false });
}
test "zbox job: тело запроса в stdin, ответ со stage в stdout" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const arena = arena_state.allocator();
const reply = try zbox.job.run(io, arena, build_options.zbox_exe, null, "{\"argv\":[\"echo\",\"из job\"]}");
const parsed = try support.parseReply(arena, reply);
try testing.expectEqualStrings("из job\n", parsed.reply.stdout);
try testing.expect(std.mem.endsWith(u8, reply, "\"stage\":\"run\"}\n"));
}
test "отмена future: job.run получает Canceled, программа убита вместе с группой" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const arena = arena_state.allocator();
const pid_path = "/tmp/zbox-test-73-future.pid";
std.Io.Dir.cwd().deleteFile(io, pid_path) catch {};
defer std.Io.Dir.cwd().deleteFile(io, pid_path) catch {};
var future = try io.concurrent(zbox.job.run, .{ io, arena, build_options.zbox_exe, null, try victimJson(arena, pid_path, null) });
// Ждём, пока программа запишет свой pid: значит, она запущена.
const pid = while (true) {
try io.sleep(.fromMilliseconds(20), .awake);
break readPid(pid_path) catch continue;
};
try testing.expect(alive(pid));
const started = std.Io.Clock.awake.now(io);
try testing.expectError(error.Canceled, future.cancel(io));
try testing.expect(started.durationTo(std.Io.Clock.awake.now(io)).toMilliseconds() < 2000);
try testing.expect(!alive(pid));
}
test "zbox serve --io threaded: десять параллельных запросов, не больше трёх сразу" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const arena = arena_state.allocator();
var server = try support.Serve.start(&.{ "--io", "threaded", "--jobs", "3" });
defer server.stop();
try support.warmup(server.port);
var shots: [10]support.Shot = undefined;
const elapsed = try support.volley(server.port, &shots);
defer support.freeVolley(&shots);
for (&shots, 0..) |*shot, i| {
if (shot.err) |err| return err;
try testing.expectEqual(200, shot.status);
try testing.expectEqualStrings(try std.fmt.allocPrint(arena, "задача {d}\n", .{i}), shot.stdout);
}
const s = try support.stats(arena, server.port);
try testing.expectEqual(3, s.peak);
try testing.expectEqual(11, s.done);
// Десять задач по 200 мс на трёх местах: четыре волны, не меньше 800 мс.
// Одна за другой они шли бы сумму своих времён, не меньше 2000 мс.
const serial = support.serialMs(&shots);
if (elapsed < 750 or elapsed >= serial) {
std.debug.print("залп {d} мс, сумма задач {d} мс\n", .{ elapsed, serial });
return error.NotParallel;
}
}
const Victim = struct {
port: u16,
arena: std.mem.Allocator,
json: []const u8,
reply: ?support.HttpReply = null,
err: ?anyerror = null,
elapsed_ms: i64 = 0,
fn fire(v: *Victim) void {
const started = std.Io.Clock.awake.now(io);
v.reply = support.postRun(v.arena, v.port, v.json) catch |err| {
v.err = err;
return;
};
v.elapsed_ms = started.durationTo(std.Io.Clock.awake.now(io)).toMilliseconds();
}
};
test "крайний срок: задача отменена и убита, ответ 504, соседи по группе целы" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const arena = arena_state.allocator();
const pid_path = "/tmp/zbox-test-73-deadline.pid";
std.Io.Dir.cwd().deleteFile(io, pid_path) catch {};
defer std.Io.Dir.cwd().deleteFile(io, pid_path) catch {};
var server = try support.Serve.start(&.{ "--io", "threaded", "--jobs", "4" });
defer server.stop();
try support.warmup(server.port);
// Жертва (sleep 30 при сроке 400 мс) стартует вместе с пачкой обычных задач.
var victim: Victim = .{ .port = server.port, .arena = arena, .json = try victimJson(arena, pid_path, 400) };
const victim_thread = try std.Thread.spawn(.{}, Victim.fire, .{&victim});
var shots: [6]support.Shot = undefined;
_ = try support.volley(server.port, &shots);
defer support.freeVolley(&shots);
victim_thread.join();
if (victim.err) |err| return err;
const reply = victim.reply.?;
try testing.expectEqual(504, reply.status);
try testing.expectEqualStrings("{\"error\":\"timeout\",\"deadline_ms\":400}\n", reply.body);
try testing.expect(victim.elapsed_ms < 3000);
try testing.expect(!alive(try readPid(pid_path)));
for (&shots, 0..) |*shot, i| {
if (shot.err) |err| return err;
try testing.expectEqual(200, shot.status);
try testing.expectEqualStrings(try std.fmt.allocPrint(arena, "задача {d}\n", .{i}), shot.stdout);
}
const s = try support.stats(arena, server.port);
try testing.expectEqual(1, s.timeouts);
try testing.expectEqual(8, s.done);
}
test "zbox serve --io threaded: 503 сверх jobs + queue" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const arena = arena_state.allocator();
var server = try support.Serve.start(&.{ "--io", "threaded", "--jobs", "1", "--queue", "1" });
defer server.stop();
try support.warmup(server.port);
var shots: [6]support.Shot = undefined;
_ = try support.volley(server.port, &shots);
defer support.freeVolley(&shots);
var busy: usize = 0;
for (&shots) |*shot| {
if (shot.err) |err| return err;
if (shot.status == 503) busy += 1 else try testing.expectEqual(200, shot.status);
}
try testing.expect(busy >= 1 and busy <= 4);
try testing.expectEqual(busy, (try support.stats(arena, server.port)).rejected);
}
test "zbox serve --io evented: в Zig 0.16.0 сервера на Io.Evented не поднять" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const result = try support.zboxRaw(arena_state.allocator(), &.{ "serve", "0", "--io", "evented" });
try testing.expectEqual(std.process.Child.Term{ .exited = 3 }, result.term);
try testing.expect(std.mem.indexOf(u8, result.stdout, "Io.Evented") != null);
if (builtin.os.tag == .macos) try testing.expect(std.mem.indexOf(u8, result.stdout, "NetworkDown") != null);
}
test "сокеты соединений закрываются при exec и не утекают в задачи" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
var server = try support.Serve.start(&.{ "--io", "threaded", "--jobs", "2" });
defer server.stop();
try testing.expectEqual(0, try support.leakedSockets(arena_state.allocator(), server.port));
}
Два теста стоит разобрать.
- Отмена future без HTTP.
job.runзапускается черезio.concurrent, программа задачи (sh -c 'echo $$ > файл; exec sleep 30') пишет свой pid. Тест ждёт этот файл, проверяет, что процесс жив, отменяет future и проверяет три вещи:cancelвернулerror.Canceled, на всё ушло меньше двух секунд (на macOS около 90 мс), программы больше нет.kill(pid, 0)ничего не шлёт, только проверяет, что такой процесс есть. - Крайний срок посреди нагрузки. Жертва с
sleep 30и сроком 400 мс стартует вместе с шестью обычными задачами на четырёх местах. Жертва получает 504 меньше чем за 3 с, её программа мертва, а шесть соседей по группе отвечают 200 каждый своим выводом. Отмена одной задачи не задела ни одну другую.timeoutsв/statsравен 1,doneравен 8 (прогрев, жертва и шесть задач).
Тест на утечку сокетов такой же, как в шаге 71: программа задачи считает сокеты у себя до и после того, как рядом открыто молчащее соединение. У шага 73 этой ошибки не было с самого начала, std.Io.net ставит close-on-exec сам. Тест сторожит, чтобы так и осталось.
$ zig build test -Dstep=73 --summary all
Build Summary: 7/7 steps succeeded; 8/8 tests passed
test success
+- run test 8 pass (8 total) 3s MaxRSS:5M
Это macOS. Эталон проверен и в контейнере с Linux (runner-zig:dev, arm64, --privileged): там весь набор шагов 48 до 73 даёт 98 зелёных из 101 при трёх пропущенных, и под эмуляцией x86-64 шаги 66, 71 и 73 тоже зелёные.
Прогон
$ zbox serve 8081 --io threaded --jobs 2
zbox: listening on port 8081 (io_threaded, jobs 2, queue 16)
$ curl -s -d '{"argv":["sh","-c","sleep 0.2; echo готово"]}' localhost:8081/run
{"exit_code":0,...,"wall_ms":212,"reason":"exited",...,"stdout":"готово\n","stderr":"","net":"host","stage":"run"}
$ curl -s -i -d '{"argv":["sleep","30"],"deadline_ms":500}' localhost:8081/run
HTTP/1.0 504 Gateway Timeout
Connection: close
Content-Type: application/json
Content-Length: 38
{"error":"timeout","deadline_ms":500}
$ curl -s localhost:8081/stats
{"running":0,"peak":1,"done":2,"rejected":0,"timeouts":1}
$ zbox serve 0 --io evented
zbox: Io.Evented (Dispatch) в Zig 0.16.0 без сети: listen ответил NetworkDown
sleep 30 получил 504 через полсекунды, а /stats показывает, что задача учтена и как выполненная (done), и как отменённая (timeouts).
Три модели одного сервиса
Итог zbox: одна и та же нагрузка против трёх моделей. Десять клиентов одновременно, у каждого {"argv":["sleep","0.2"]}, три прогона после прогрева. Пул и std.Io с --jobs 4, у пула 8 потоков, сборка ReleaseSafe. «Строк» это код без пустых строк и комментариев, который модель добавляет сама поверх общего. Память это RSS самого сервера без дочерних процессов, опрос раз в 20 мс.
Load average в момент замера не записан, а машину в тот день грузили соседние процессы (в других замерах того же дня от 13 до 57). Время здесь почти целиком сон задач, так что смотри на волны и порядок моделей, а не на отдельные миллисекунды.
Apple M4 Max, 16 ядер, macOS 26.6.2:
| Модель | Строк | RSS покоя, КБ | RSS пик, КБ | Пачка, мс | Медиана, мс | Худшая, мс |
|---|---|---|---|---|---|---|
| 66, итеративный | 32 | 2512 | 2992 | 2161 до 2229 | 1307 до 1338 | 2161 до 2228 |
| 71, пул потоков | 296 | 2464 | 2928 | 658 до 703 | 441 до 472 | 657 до 703 |
73, std.Io (Io.Threaded) | 223 | 2000 | 3344 | 675 до 697 | 455 до 460 | 675 до 696 |
Linux 7.0.14 (OrbStack), aarch64, 16 ядер, тот же M4 Max, контейнер runner-zig:dev:
| Модель | RSS покоя, КБ | VmHWM, КБ | Пачка, мс | Медиана, мс | Худшая, мс |
|---|---|---|---|---|---|
| 66, итеративный | 2324 | 2404 | 2295 до 2448 | 1478 до 1510 | 2291 до 2441 |
| 71, пул потоков | 4304 | 4612 | 837 до 1140 | 583 до 614 | 832 до 1055 |
73, std.Io (Io.Threaded) | 2772 | 7852 | 712 до 770 | 461 до 536 | 672 до 767 |
Как это читать.
- Итеративный сервер это десять задач подряд: пачка чуть больше 10 × 200 мс, медиана на середине очереди.
- Пул и
std.Ioупираются в--jobs 4: три волны по 4, 4 и 2 задачи, пачка около 600 мс плюс запуск процессов. Медиана около 450 мс, половина клиентов попала во вторую волну. Разница между 71 и 73 в пределах разброса: время съедают задачи, а не способ их раздавать. - Пул держит восемь потоков всегда (в Linux это видно по RSS покоя),
Io.Threadedзаводит потоки по требованию: в покое меньше, в пике больше, потому что у каждого соединения под нагрузкой ещё и поток таймера. - Код на
std.Ioкороче пула на 73 строки и при этом умеет то, чего пул не умеет: отменять задачу. Очередь, воркеры и ядовитые пилюли стали реализациейIo, а не нашим кодом.
Последний вывод и есть ответ на вопрос урока. Книга показывает три модели как три разных программы. std.Io делает их тремя реализациями одного интерфейса, и выбор между ними переезжает из кода в main. Сегодня выбор ограничен тем, что Io.Evented не умеет сеть; сервис к этому выбору уже готов.
На macOS
Этот урок на Mac идёт почти целиком, и в одном месте Mac даже впереди Linux.
Что работает напрямую:
- все демонстрации
std.Io(colorless.zig,queue.zig,cancel.zig):zig testна любой из двух систем; - шаг
tinyцеликом,zig build test -Dstep=73даёт 11 зелёных из 11; прокси,curl -x,loadgenчерез прокси; tiny echoserver-io --evented:Io.Eventedв 0.16.0 собирается только на macOS, этоIo.Dispatchповерх libdispatch. На Linux его реализацияIo.Uringне компилируется, и в контейнере этот тест пропускается;zbox serve --io threadedи все тесты шага 73zbox, кроме изоляции:--isolateостаётся только для Linux, как в уроке 67;zbox serve --io eventedна Mac доходит до настоящей проверки (listenотвечаетNetworkDown), на Linux отвечает без неё.
Что через контейнер ghcr.io/bondiano/runner-zig:dev (рецепт в уроках 66 и 67):
zbox serve --io threaded --isolate ROOTFS: изоляция, cgroups и пространства имён,--privileged;- ThreadSanitizer. В Zig 0.16.0 он работает только на Linux (на macOS 26 программа с
-fsanitize-threadпадает при запуске, это разобрано в уроке 72). Тесты шагаtinyпод ним:zig build test -Dstep=73 -Dtsanв контейнере arm64 даёт 10 зелёных и один пропуск, без отчётов о гонках, включая тест с восемью потоками против кэша.
Про epoll и kqueue. В уроке 68 ты писал обёртку Poller над epoll и kqueue руками. std.Io прячет эту разницу глубже: Io.Dispatch на macOS живёт на libdispatch, которая сама опирается на kqueue, а Io.Uring на Linux вообще не про готовность дескрипторов, а про очередь завершённых операций io_uring. Это другая модель, чем epoll: ядро не говорит «можно читать», а сразу отдаёт прочитанное. Для кода на std.Io разницы нет никакой, в этом и смысл.
gettid и /proc/<pid>/task из уроков 69 и 71 здесь не нужны: std.Io своих потоков не показывает, а пул echo_pre под прокси называет воркеров через setName, который в macOS срабатывает только для текущего потока.
Практика
Напиши кэш прокси сам, в том виде, в каком его просит лабораторная. Отличий от эталона три: лимитов три (max_bytes, max_objects, max_object_size), записи лежат в std.ArrayList и ищутся линейно, а блокировка берётся отменяемыми lockShared(io) и lock(io), потому что кэш зовут из задач std.Io.
Готово: Cache с полями entries, bytes, clock, lock, методы init, deinit, tick() (следующая отметка по атомарным часам) и find(key). Написать надо два метода.
get(io, key, out) !?usize: под общей блокировкой найти запись, скопировать значение вout(не влезает,error.NoSpace), освежить её атомарную отметкуstampи вернуть длину; промах этоnull.put(io, key, value) !void: слишком большой объект не кэшировать; под исключительной блокировкой заменить значение или добавить запись, затем выбрасывать запись с самой старой отметкой, пока нарушен хоть один лимит. Копии ключа и значения кэш держит свои, и копировать их лучше до захвата блокировки.
Тесты проверяют попадания, замену и независимость кэша от буферов вызывающего, порядок вытеснения по числу объектов и по объёму, пропуск больших объектов. Отдельные тесты проверяют сам вид блокировки: get обязан пройти, пока другой читатель держит общую блокировку, а put обязан ждать. В конце два писателя и четыре читателя по 3000 операций над двенадцатью ключами: ни одного испорченного значения, лимиты соблюдены, утечек нет.
Упражнения
Итоги
- Цвет функции это свойство async/await: асинхронная функция заражает всех вызывающих. Zig 0.16 передаёт
std.Ioпараметром, какAllocator, и одна функция работает и на потоках, и на цикле событий. io.asyncразрешает одновременность и может выполнить задачу на месте;io.concurrentтребует её и возвращаетerror.ConcurrencyUnavailable, если не может. Где корректность зависит от одновременности (сессия сервера, таймер против задачи, производитель против потребителя), нуженconcurrent.Futureждут черезawaitили отменяют черезcancel, оба идемпотентны;defer future.cancel(io)сразу после запуска.Io.Groupдержит задачи без результата и не даёт им пережить свой блок:defer group.cancel(io).Io.Queueэтоsbufс закрытием и отменой,Io.Selectэто гонка задач с первым результатом иcancelDiscardдля проигравших.- Отмена приходит
error.Canceledиз ближайшей точки отмены, то есть вызоваstd.IoсCanceledв множестве ошибок. Её не глотают; долгий цикл без ввода-вывода зовётio.checkCancel().Io.Threadedбудит поток в системном вызове сигналомSIGIOс повтором на случай гонки. Процессы, запущенные задачей, отмена не трогает: их убивает нашerrdefer. Io.Threadedэто пул из урока 71,Io.Eventedэто событийный сервер из урока 68 с файберами, переключение файбера сохраняетsp,fpиpc. В 0.16.0 уIo.Eventedнет сети,Io.Uringне собирается, поэтому сетевой код работает наIo.Threadedи готов кEventedбез правок.- Прокси разбирает абсолютный URI, говорит с сервером на HTTP/1.0 с
Connection: close, пересылает ответ кусками и кэширует только200не больше лимита объекта. Кэш окупается там, где ответ дорого делать: CGI через кэш на порядок быстрее, чем напрямую. - LRU под
RwLock: порядок живёт в атомарных метках, а не в списке, поэтому попадание пишет одну метку под общей блокировкой; вытеснение ищет наименьшую метку заO(n)под исключительной. Копия объекта делается доunlockShared. - Одновременные промахи по одному ключу все идут к серверу (давка в кэше); лечится одиночным полётом.
zbox --io threadedгонит задачу с таймером черезIo.Select, отмена доходит до программы черезerrdefer child.killи обработчик SIGTERM, который убивает группу. Код короче пула и умеет отменять задачу.
Дальше
Это последний урок раздела. В первом уроке ты написал hello на Zig и смотрел, как текст превращается в процесс. Сегодня твой прокси на пуле потоков отдаёт ответы CGI из кэша на порядок быстрее, чем их делает сервер, а твой раннер собирает чужой исходник в песочнице и снимает его по крайнему сроку. Между этими точками книга Брайанта и О’Халларона пройдена целиком, глава за главой, а в каждой главе мы шли на шаг глубже, чем она.
Четыре сквозных проекта
| Проект | Где начался | Чем закончился |
|---|---|---|
zt, свои binutils | hexdump в уроке 06 | дизассемблер x86-64, выросший по классу инструкций за урок в блоке про машинный уровень; readelf, nm и компоновщик, который в уроке 44 собрал запускаемый бинарник из двух .o; ldd, ltrace, strace на ptrace, pmap; сэмплирующий профилировщик, который в уроке 71 рисует flame graph сервера с пулом потоков по потокам |
zl, Lisp от интерпретатора до JIT | eval Маккарти в уроке 05 | три исполнителя одной программы (дерево, байткод на двух диспетчерах, JIT в машинный код из урока 19), тегированные указатели, своя куча ячеек на аллокаторе из malloclab и Mark&Sweep, а в уроке 70 общая куча на несколько потоков с остановкой мира |
| Y86-64 растёт в компьютер | ISA и кодировка в уроке 20 | процессор на Verilog от вентилей до конвейера, затем в симуляторе кэш, бит привилегий, таблица исключений, таймер, ядро на ассемблере Y86 с переключением процессов из урока 50, страничная память с TLB и устройства с отображением в память из урока 60 |
zbox, песочница и раннер | zbox run на fork и execve в уроке 48 | лимиты времени и памяти, захват вывода, выключенная сеть, HTTP-фасад на TINY, шесть пространств имён, pivot_root и seccomp из урока 67, пул воркеров и сегодняшняя модель на std.Io с отменой |
Каждый из четырёх отвечает на один из тезисов, с которых начинался раздел. zt про то, что программа это байты по протоколу: ты читаешь и пишешь эти байты своими инструментами. zl про цену абстракции: одна программа, три исполнителя, и цена каждого измерена в тактах. Y86 про то, что процессор это провода, а ОС это исключения плюс виртуальная память, и это проверено на машине, собранной своими руками. zbox про ОС как набор механизмов изоляции, каждый из которых это один или два системных вызова.
Шесть лабораторных
| Лабораторная | Урок | Что осталось у тебя |
|---|---|---|
| y86lab | 20 до 30 | ассемблер и симулятор Y86-64 на Zig, SEQ и PIPE на Verilog с тестбенчами |
| perflab | 31 до 36 | harness для CPE, combine1 до combine7, развёртка и несколько аккумуляторов, SIMD-вариант |
| cachelab | 39 и 40 | симулятор кэша по трассе, транспонирование блоками |
| shlab | 52 | оболочка с job control, пайпами и переадресацией |
| malloclab | 58 | три std.mem.Allocator с замером utilization и throughput |
| tinylab и proxylab | 66 и этот урок | веб-сервер TINY (итеративный, на потоках, на пуле), эхо-серверы во всех моделях главы 12 и кэширующий прокси |
Куда идти
- Конкурентность с другой стороны. В разделе про асинхронность те же очереди, семафоры и гонки с точки зрения прикладного кода, а урок про конкурентность и параллелизм снова приводит к закону Амдала из урока 71. Язык, где гонки ловит компилятор: в Rust потоки,
SendиSync,MutexиArc, атомики и порядки памяти (здесь мы всюду писали.monotonicи.acq_rel, там объяснено, что они значат для процессора), а устройство tokio и асинхронный сервер покажут цикл событий, который в Zig 0.16 пока спрятан вIo.Evented. Модель горутин из Go: конкурентность вглубь, планировщик M:N поверх тех же потоков ОС. - Файберы и структурная отмена. Файберы Effect это та же идея, что
Io.GroupиFuture.cancel, доведённая до библиотеки для прикладного кода. - Дальше книги.
Io.Eventedбудет получать сеть, аIo.Uringкомпилироваться: следи заlib/std/Io/в репозитории Zig на Codeberg и попробуй перевестиecho_io.serveиzbox serve --io eventedнаEvented, когда это случится. Ради этого код и писался против интерфейса. Про io_uring стоит прочитать работу Йенса Аксбо «Efficient IO with io_uring», про кэши HTTP сам RFC 9111. - Своё. Четыре репозитория лежат у тебя в руках. Раннер курса, в котором ты весь раздел гонял задачи, теперь не чёрный ящик:
runner/bs-runnerчитается строчка за строчкой. Лучший следующий шаг это взять один из четырёх проектов и довести его до того, чем ты пользуешься каждый день.
домашка