Раздел 32 · Системное программирование: Zig, ассемблер, Verilog
Пул потоков и параллельные вычисления
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Пул потоков и параллельные вычисления
В прошлом уроке ты собрал из семафоров ограниченный буфер
sbufи увидел, как один общий счётчик под двумя потоками теряет половину прибавлений. Сегодняsbufполучает настоящую работу: он станет очередью соединений между главным потоком и пулом заранее созданных воркеров. На этом пуле встанут три сервера сразу: эхо из книги, TINY иzbox serve, который наконец перестанет гонять задачи студентов по одной. Вторая половина урока про параллельные вычисления: сумма ряда в пяти версиях, от мьютекса на каждом сложении до локального аккумулятора, ускорение и эффективность на шестнадцати ядрах, закон Амдала числом и ложное разделение, которое прячется в соседних ячейках массива. А в концеzt profпокажет, во что блокировка обходится в настоящем сервере.
Цели урока
- Объяснить, чем пул заранее созданных потоков лучше потока на соединение, и где у пула предел.
- Собрать
echoservert_preкниги: главный поток принимает соединения и кладёт их вsbuf, воркеры забирают по одному; остановить воркеры отравленной пилюлей. - Сделать однократную инициализацию без
pthread_once: быстрый путь на атомике, медленный под мьютексом. - Посадить TINY на пул и нагрузить его генератором
loadgen. - Написать
psumв пяти версиях и объяснить по замерам, почему мьютекс и атомик на каждой итерации медленнее одного потока. - Посчитать ускорение и эффективность, различать сильное и слабое масштабирование, оценить последовательную долю по закону Амдала и по Карпу и Флэтт.
- Найти ложное разделение в замере и вылечить его выравниванием по линии кэша.
- Дать
zbox serveпул воркеров, очередь с отказом 503 и семафор на число одновременных задач; найти и закрыть утечку сокетов черезexec. - Снять профиль пула и сервера под нагрузкой через
zt prof --threadsи прочитать, сколько стоит блокировка.
Идея: потоки заранее
В уроке про потоки сервер заводил поток на каждое соединение. Это прямолинейно и честно: медленный клиент держит только свой поток, остальные его не ждут. Но у модели две цены, и обе видны на замере.
Первая цена: поток создаётся на каждый запрос. pthread_create это системный вызов clone, выделение стека через mmap и запись в таблицы ядра, а на выходе всё то же в обратную сторону. Когда запрос короче этой бухгалтерии, сервер тратит больше на рождение и смерть потоков, чем на работу. Это показал замер TINY в прошлых уроках: на статике tiny --threads в контейнере Linux отдал 7343 запроса в секунду против 15716 у итеративного сервера.
Вторая цена: число потоков ничем не ограничено. Тысяча клиентов это тысяча потоков, у каждого стек, каждый борется за ядра, и планировщик тратит время на переключения. Сервер не выбирает, сколько работы брать, её выбирают клиенты.
Обе цены снимает пул потоков. Главный поток только принимает соединения и кладёт их в ограниченную очередь. Заранее созданные воркеры, их фиксированное число, забирают соединения из очереди по одному, обслуживают и возвращаются за следующим. Очередь у нас уже есть: это sbuf из прошлого урока, производитель у него главный поток, потребители воркеры.
sbuf на 16 слотов
accept ──► [ fd | fd | fd | | | ... ] ──► воркер 1: handle(fd), close(fd)
(главный) insert ждёт свободный слот ──► воркер 2
──► воркер 3
──► воркер 4
У ограниченной очереди есть свойство, которое легко не заметить: обратное давление. Когда все воркеры заняты и очередь полна, insert в главном потоке засыпает на семафоре slots, и сервер перестаёт звать accept. Новые клиенты копятся в очереди ядра listen (её длину задал backlog в уроке про сокеты), а когда переполнится и она, ядро перестанет отвечать на их SYN. Лишняя работа не копится в памяти процесса. Цена у этого тоже есть: клиент, которому не хватило воркера, ждёт. В zbox ниже мы выберем другое поведение: отказать сразу, кодом 503.
С точки зрения прикладного кода тот же приём ты видел в разделе про асинхронность, в разговоре об очередях и пулах: там пул ограничивал число одновременных промисов. Здесь те же идеи, только исполнители это настоящие потоки ядра, а очередь стоит на семафорах.
Шаг проекта: TINY на пуле
Эталон our-tiny получает в этом шаге три файла: пул src/conc/echo_pre.zig (эхо-сервер книги плюс общий каркас, на котором сидят TINY и прокси из урока 73), генератор нагрузки src/conc/loadgen.zig и psum в src/conc/psum.zig, о котором вторая половина урока. sbuf берём из прошлого урока как есть, Server.handle из урока про потоки тоже.
echo_pre.zig: пул, отравленная пилюля и Once
//! `echoservert_pre` из главы 12: предварительно созданный пул потоков.
//! Главный поток только принимает соединения и кладёт их в `sbuf`,
//! `workers` потоков забирают их оттуда по одному. Потоки не создаются на
//! каждое соединение, а их число ограничено: тысяча клиентов не превращается
//! в тысячу потоков. Цена: клиент, которому не хватило воркера, ждёт в очереди.
//!
//! `serve` не знает, что делать с соединением: это решает `handle`. Эхо,
//! TINY (`tiny --pool N`) и прокси из урока 73 сидят на одном пуле.
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 sbuf_mod = @import("sbuf.zig");
const Peer = @import("echo_threads.zig").Peer;
const tiny = @import("../http/tiny.zig");
pub const Options = struct {
workers: usize = 4,
/// Слотов в `sbuf`. Когда очередь полна, `accept` ждёт свободного слота.
queue: usize = 16,
/// Сколько соединений принять и выйти, дождавшись воркеров. 0 значит вечно.
max_clients: usize = 0,
};
/// Соединение в очереди: дескриптор и имя клиента по значению.
/// Дескриптор -1 это сигнал воркеру выйти.
pub const Job = struct {
fd: c.fd_t,
peer: Peer = .{},
};
/// `handle(ctx, connfd, peer)` обслуживает соединение; закрывает его пул.
pub fn serve(gpa: std.mem.Allocator, io: Io, listenfd: c.fd_t, options: Options, ctx: anytype, comptime handle: fn (@TypeOf(ctx), c.fd_t, []const u8) void) !void {
var sbuf: sbuf_mod.Sbuf(Job) = try .init(gpa, io, options.queue);
defer sbuf.deinit(gpa);
const threads = try gpa.alloc(std.Thread, options.workers);
defer gpa.free(threads);
const Worker = struct {
fn run(sp: *sbuf_mod.Sbuf(Job), cx: @TypeOf(ctx)) void {
while (true) {
const job = sp.remove();
if (job.fd < 0) return;
handle(cx, job.fd, job.peer.slice());
socket.close(job.fd);
}
}
};
for (threads, 0..) |*t, i| {
t.* = try std.Thread.spawn(.{}, Worker.run, .{ &sbuf, ctx });
// Имя потока видно в `top -H`, `/proc/<pid>/task/*/comm` и у `zt prof
// --threads` основанием своего столбца. Linux даёт назвать чужой поток,
// macOS только текущий: там вызов молча не срабатывает.
var name_buf: [std.Thread.max_name_len]u8 = undefined;
const name = std.fmt.bufPrint(&name_buf, "worker-{d}", .{i}) catch "worker";
t.setName(io, name) catch {};
}
var served: usize = 0;
while (options.max_clients == 0 or served < options.max_clients) : (served += 1) {
var job: Job = .{ .fd = -1 };
const conn = socket.accept(listenfd, &job.peer.buf) catch continue;
job.fd = conn.fd;
job.peer.len = conn.name.len;
sbuf.insert(job);
}
for (threads) |_| sbuf.insert(.{ .fd = -1 });
for (threads) |t| t.join();
}
/// TINY на пуле: `tiny --pool N`.
pub fn tinyHandle(server: *tiny.Server, connfd: c.fd_t, peer: []const u8) void {
server.handle(connfd, peer);
}
/// `pthread_once` своими руками: функция зовётся ровно один раз, сколько
/// бы потоков ни пришло одновременно. Быстрый путь без блокировки, когда
/// всё уже сделано; медленный под мьютексом с повторной проверкой.
pub const Once = struct {
done: std.atomic.Value(bool) = .init(false),
mutex: Io.Mutex = .init,
pub fn call(o: *Once, io: Io, comptime f: fn () void) void {
if (o.done.load(.acquire)) return;
o.mutex.lockUncancelable(io);
defer o.mutex.unlock(io);
if (o.done.raw) return;
f();
o.done.store(true, .release);
}
};
// Счётчик байтов всех соединений, `echo_cnt` книги. Общий для всех воркеров,
// поэтому под мьютексом. Книга инициализирует мьютекс через `pthread_once`,
// потому что `sem_init` нельзя позвать статически. В Zig мьютекс и так
// инициализируется статически (`.init`), и `Once` здесь стережёт только
// обнуление счётчика: так видно, как устроен `pthread_once`.
var byte_cnt: usize = 0;
var byte_mutex: Io.Mutex = .init;
var byte_once: Once = .{};
fn initEchoCnt() void {
byte_cnt = 0;
}
pub fn byteCount(io: Io) usize {
byte_mutex.lockUncancelable(io);
defer byte_mutex.unlock(io);
return byte_cnt;
}
pub const EchoCtx = struct {
io: Io,
quiet: bool = false,
};
/// `echo_cnt`: эхо с общим счётчиком байтов.
pub fn echoCnt(ctx: EchoCtx, connfd: c.fd_t, _: []const u8) void {
byte_once.call(ctx.io, initEchoCnt);
var buf: [8192]u8 = undefined;
var reader: fdio.Reader = .init(connfd, &buf);
while (true) {
const line = reader.interface.takeDelimiterInclusive('\n') catch |err| switch (err) {
error.StreamTooLong => reader.interface.buffered(),
else => return,
};
byte_mutex.lockUncancelable(ctx.io);
byte_cnt += line.len;
const total = byte_cnt;
byte_mutex.unlock(ctx.io);
if (!ctx.quiet) {
var log_buf: [128]u8 = undefined;
var stderr: fdio.Writer = .init(2, &log_buf);
stderr.interface.print("server received {d} ({d} total) bytes on fd {d}\n", .{ line.len, total, connfd }) catch {};
stderr.interface.flush() catch {};
}
fdio.writen(connfd, line) catch return;
if (line.len == buf.len) reader.interface.tossBuffered();
}
}
Что здесь стоит разобрать.
Соединение передаётся по значению. В sbuf лежит не указатель, а Job целиком: дескриптор и имя клиента в своём буфере. insert копирует его в слот, remove копирует обратно в стек воркера. Главный поток после insert сразу переиспользует свою переменную job для следующего accept, и воркер этого не видит. Это та же ошибка, что и в упражнении про гонку на аргументе потока: указатель на локальную переменную главного потока означал бы, что два воркера могут прочитать один и тот же дескриптор.
Воркеры останавливаются отравленной пилюлей. Когда serve принял max_clients соединений (это нужно тестам; настоящий сервер крутится вечно), он кладёт в очередь по одному Job с fd = -1 на каждого воркера. Воркер, достав такой Job, выходит из цикла. Пилюли идут после всех настоящих соединений, sbuf раздаёт по порядку, поэтому ни одно принятое соединение не бросается. Остановить чужой поток снаружи в POSIX можно через pthread_cancel, и тогда он умирает в точке отмены, которую выбрал не он, возможно, с захваченным мьютексом или недописанной структурой. Пилюля даёт потоку выйти в той точке, которую выбрал он сам.
handle приходит параметром времени компиляции. Эхо, TINY и прокси отличаются только тем, что делать с принятым соединением. Контекст ctx передаётся в каждый воркер по значению (anytype), поэтому для TINY это указатель *tiny.Server, а для эха маленькая структура с io. У Server после урока про потоки нет изменяемого состояния, и один сервер обслуживают все воркеры сразу.
Воркеры названы. Главный поток сразу после spawn даёт каждому имя worker-N через setName. Имя видно в top -H, в /proc/<pid>/task/*/comm и, главное для нас, у zt prof --threads в конце урока: там у каждого воркера будет своё подножие на flame graph. Linux разрешает назвать чужой поток, macOS только текущий, поэтому на Mac вызов молча не срабатывает, и ошибка намеренно проглатывается.
Once. Книга в echo_cnt инициализирует счётчик и мьютекс через pthread_once: первый пришедший воркер делает работу, остальные ждут и больше её не повторяют. В std 0.16 ни std.once, ни обёртки над pthread_once нет, поэтому Once написан руками, и это классическая схема с двойной проверкой. Быстрый путь это одна атомарная загрузка done с порядком .acquire: если там true, всё, что сделала f, уже видно этому потоку. Медленный путь берёт мьютекс и проверяет done ещё раз: пока поток ждал мьютекса, первый мог закончить. Запись done идёт с .release и только после f. Если переставить её перед вызовом f, второй поток по быстрому пути увидит true и побежит пользоваться счётчиком, который ещё не обнулён. Порядки памяти подробно разобраны в уроке про атомики в Rust; здесь хватит правила «запись после работы с release, чтение перед использованием с acquire».
Сам мьютекс в Zig 0.16 живёт в std.Io.Mutex и инициализируется статически (.init), так что Once в эталоне стережёт только обнуление счётчика. Это сделано нарочно, чтобы показать устройство, а не потому, что иначе нельзя.
Подкоманды
\\
\\Использование:
- \\ tiny tiny [--threads] <port> [root] веб-сервер: итеративный или поток на соединение
+ \\ tiny tiny [--threads | --pool N] <port> [root] веб-сервер: итеративный, поток на соединение или пул
\\ tiny hostinfo [--std] <name> все адреса имени через getaddrinfo или std.Io.net
\\ tiny echoserver [--std] <port> эхо-сервер, одно соединение за раз
@@
\\ tiny echoservers --select|--poll|--epoll <port> эхо: события (--epoll это kqueue на macOS)
\\ tiny echoservert <port> эхо: поток на соединение
- \\ tiny conc-bench [--clients 100,1000] [--active A] [--lines K] [--models procs,threads,select,poll,epoll] [--json]
+ \\ tiny echoservert-pre [--workers N] <port> эхо: пул потоков и sbuf
+ \\ tiny conc-bench [--clients 100,1000] [--active A] [--lines K] [--models procs,threads,pool,select,poll,epoll] [--workers N] [--json]
\\ tiny badcnt [niters] гонка на счётчике и три исправления
+ \\ tiny psum-bench [--n N] [--n-slow N] [--threads 1,2,4,8] [--repeat R] [--json]
+ \\ tiny loadgen <host> <port> <conns> <reqs> <path> нагрузка на HTTP-сервер
\\
;
@@
.{ "echoservers", cmdEchoservers },
.{ "echoservert", cmdEchoservert },
+ .{ "echoservert-pre", cmdEchoservertPre },
.{ "conc-bench", cmdConcBench },
.{ "badcnt", cmdBadcnt },
+ .{ "psum-bench", cmdPsumBench },
+ .{ "loadgen", cmdLoadgen },
};
inline for (commands) |entry| {
@@
fn cmdTiny(ctx: Ctx, rest: []const [:0]const u8) !void {
var args = rest;
- var mode: enum { iterative, threads } = .iterative;
+ var mode: enum { iterative, threads, pool } = .iterative;
+ var workers: usize = 0;
if (args.len > 0 and std.mem.eql(u8, args[0], "--threads")) {
mode = .threads;
args = args[1..];
+ } else if (args.len > 1 and std.mem.eql(u8, args[0], "--pool")) {
+ mode = .pool;
+ workers = try parseCount(ctx, args[1]);
+ args = args[2..];
}
if (args.len < 1) return fail(ctx.out, usage);
@@
.iterative => try server.serveForever(),
.threads => try conc.echo_threads.serveTiny(&server, 0),
+ .pool => try conc.echo_pre.serve(ctx.gpa, ctx.io, server.listenfd, .{ .workers = workers }, &server, conc.echo_pre.tinyHandle),
}
}
@@
}
+fn cmdEchoservertPre(ctx: Ctx, rest: []const [:0]const u8) !void {
+ var args = rest;
+ var workers: usize = 4;
+ 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);
+ const listenfd = try listenOn(ctx, "echoservert-pre", try parsePort(ctx, args[0]));
+ defer tiny.socket.close(listenfd);
+ try conc.echo_pre.serve(ctx.gpa, ctx.io, listenfd, .{ .workers = workers }, conc.echo_pre.EchoCtx{ .io = ctx.io }, conc.echo_pre.echoCnt);
+}
+
fn cmdConcBench(ctx: Ctx, rest: []const [:0]const u8) !void {
var clients: []const usize = &.{ 100, 1000 };
- var models: []const conc.bench.Model = &.{ .procs, .threads, .select, .poll, .epoll };
+ var models: []const conc.bench.Model = &.{ .procs, .threads, .pool, .select, .poll, .epoll };
var cfg: conc.bench.Config = .{ .clients = 0 };
var json = false;
@@
} else if (std.mem.eql(u8, flag, "--active")) {
cfg.active = try parseCount(ctx, value);
+ } else if (std.mem.eql(u8, flag, "--workers")) {
+ cfg.workers = try parseCount(ctx, value);
} else if (std.mem.eql(u8, flag, "--models")) {
var list: std.ArrayList(conc.bench.Model) = .empty;
@@
}
+fn cmdPsumBench(ctx: Ctx, rest: []const [:0]const u8) !void {
+ const cpus = std.Thread.getCpuCount() catch 1;
+ var threads: std.ArrayList(usize) = .empty;
+ for ([_]usize{ 1, 2, 4, 8 }) |t| if (t <= cpus) try threads.append(ctx.arena, t);
+ if (cpus > 8 and cpus <= conc.psum.max_threads) try threads.append(ctx.arena, cpus);
+ var cfg: conc.psum.BenchConfig = .{ .threads = threads.items };
+ var json = false;
+ var i: usize = 0;
+ while (i < rest.len) : (i += 1) {
+ const flag = rest[i];
+ if (std.mem.eql(u8, flag, "--json")) {
+ json = true;
+ continue;
+ }
+ if (i + 1 >= rest.len) return fail(ctx.out, usage);
+ i += 1;
+ if (std.mem.eql(u8, flag, "--n")) {
+ cfg.n = try parseCount(ctx, rest[i]);
+ } else if (std.mem.eql(u8, flag, "--repeat")) {
+ cfg.repeat = try parseCount(ctx, rest[i]);
+ } else if (std.mem.eql(u8, flag, "--n-slow")) {
+ cfg.n_slow = try parseCount(ctx, rest[i]);
+ } else if (std.mem.eql(u8, flag, "--threads")) {
+ cfg.threads = try parseList(ctx, rest[i]);
+ } else return fail(ctx.out, usage);
+ }
+ const samples = try ctx.arena.alloc(conc.psum.Sample, cfg.versions.len * cfg.threads.len);
+ const got = try conc.psum.bench(ctx.io, cfg, samples);
+ if (json) try conc.psum.writeJson(ctx.out, got) else try conc.psum.writeTable(ctx.out, got);
+ try ctx.out.flush();
+}
+
+fn cmdLoadgen(ctx: Ctx, rest: []const [:0]const u8) !void {
+ if (rest.len != 5) return fail(ctx.out, usage);
+ const r = try conc.loadgen.run(ctx.gpa, ctx.io, .{
+ .host = rest[0],
+ .port = try parsePort(ctx, rest[1]),
+ .conns = try parseCount(ctx, rest[2]),
+ .reqs = try parseCount(ctx, rest[3]),
+ .path = rest[4],
+ });
+ try conc.loadgen.write(ctx.out, r);
+ try ctx.out.flush();
+}
+
fn listenOn(ctx: Ctx, name: []const u8, port: u16) !std.c.fd_t {
const listenfd = try tiny.socket.openListenfd(port);
tiny tiny --pool N это тот же TINY из урока про веб-сервер, только соединения обслуживает пул из N воркеров, а echoservert-pre это эхо книги с общим счётчиком байтов. psum-bench и loadgen понадобятся ниже. Подключение модулей и номер шага:
pub const rw = @import("conc/rw.zig");
pub const badcnt = @import("conc/badcnt.zig");
+ pub const echo_pre = @import("conc/echo_pre.zig");
+ pub const loadgen = @import("conc/loadgen.zig");
+ pub const psum = @import("conc/psum.zig");
};
/// Номера уроков курса, на которых проект вырос. Каждому шагу
/// соответствует файл `tests/step_NN.zig`, и все они зелёные на финале.
-const project_steps = [_]u8{ 63, 64, 65, 66, 68, 69, 70 };
+const project_steps = [_]u8{ 63, 64, 65, 66, 68, 69, 70, 71 };
pub fn build(b: *std.Build) void {
conc-bench: пул рядом с остальными моделями
В уроке про мультиплексирование conc-bench сравнивал процессы и события, урок 69 добавил поток на соединение. Пул встаёт туда же моделью pool, со своим числом воркеров. Сервер бенча живёт в дочернем процессе после fork, а Io родителя через fork не переносится, поэтому ребёнок заводит свой Io.Threaded для семафоров sbuf.
-//! `conc-bench`: одна нагрузка на все модели эхо-сервера урока 68. Сервер
-//! работает в отдельном процессе (`fork` после `openListenfd`), клиенты в
-//! этом: у сервера свои дескрипторы, и `select` на тысяче клиентов
-//! укладывается в `FD_SETSIZE`.
+//! `conc-bench`: одна нагрузка на все модели эхо-сервера урока 68 (и 69,
+//! 71 для сравнения). Сервер работает в отдельном процессе (`fork` после
+//! `openListenfd`), клиенты в этом: у сервера свои дескрипторы, и `select`
+//! на тысяче клиентов укладывается в `FD_SETSIZE`.
//!
//! Нагрузка: `clients` соединений открываются заранее, потом каждое шлёт
@@
//! немногие. Тут и видна разница между `select`/`poll`, которые на каждом
//! вызове проходят все дескрипторы, и `epoll`/`kqueue`, которые отдают
-//! только готовые.
+//! только готовые. Говорят первые `active` подключившихся: пул принимает
+//! их раньше молчащих, иначе молчуны заняли бы все воркеры навсегда.
const std = @import("std");
@@
const echo_epoll = @import("echo_epoll.zig");
const echo_threads = @import("echo_threads.zig");
+const echo_pre = @import("echo_pre.zig");
pub const Model = enum {
procs,
threads,
+ pool,
select,
poll,
@@
lines: usize = 100,
line_len: usize = 64,
+ /// Воркеров у модели `pool`.
+ workers: usize = 8,
};
@@
if (pid < 0) return error.ForkFailed;
if (pid == 0) {
- runServer(model, listenfd);
+ runServer(model, listenfd, cfg);
c._exit(0);
}
@@
}
-fn runServer(model: Model, listenfd: c.fd_t) void {
- // Свежий процесс: свой аллокатор libc, родительский через `fork` не живёт.
+fn runServer(model: Model, listenfd: c.fd_t, cfg: Config) void {
+ // Свежий процесс: свой аллокатор libc и своя реализация `Io` для
+ // семафоров пула. Состояние родительского `Io` через `fork` не живёт.
const gpa = std.heap.c_allocator;
+ var threaded: Io.Threaded = .init(gpa, .{});
+ const io = threaded.io();
var discard: Io.Writer.Discarding = .init(&.{});
const opts: echo_select.Options = .{};
@@
.procs => echo_procs.serve(listenfd, 0, &discard.writer),
.threads => echo_threads.serve(listenfd, 0, true),
+ .pool => echo_pre.serve(gpa, io, listenfd, .{ .workers = cfg.workers, .queue = cfg.clients }, echo_pre.EchoCtx{ .io = io, .quiet = true }, echo_pre.echoCnt),
.select => if (echo_select.serve(gpa, listenfd, opts)) |_| {} else |err| err,
.poll => if (echo_poll.serve(gpa, listenfd, opts)) |_| {} else |err| err,
test "conc-bench: каждая модель отвечает на всю нагрузку" {
_ = conc.bench.raiseFdLimit(1024);
- inline for (.{ .procs, .threads, .select, .poll, .epoll }) |model| {
- const r = try conc.bench.run(testing.allocator, testing.io, model, .{ .clients = 8, .lines = 5 });
+ inline for (.{ .procs, .threads, .pool, .select, .poll, .epoll }) |model| {
+ const r = try conc.bench.run(testing.allocator, testing.io, model, .{ .clients = 8, .lines = 5, .workers = 2 });
try testing.expectEqual(@as(usize, 8), r.clients);
try testing.expect(r.throughput > 0);
Тысяча клиентов, все говорят, по 300 строк; контейнер Linux на той же машине (из README эталона):
model clients active lines seconds lines/s p50 us p99 us max us
procs 1000 1000 300 4.310 69601 13772.3 29659.3 91487
threads 1000 1000 300 3.780 79373 11356.4 30797.3 102971
pool 1000 1000 300 2.128 140948 47.6 299.6 2114356
epoll 1000 1000 300 2.037 147293 5597.1 21041.9 75457
У пула лучшая медиана и худший максимум, и это одно и то же свойство. Восемь клиентов обслуживаются сразу и без очереди, поэтому строка идёт туда и обратно за десятки микросекунд. Остальные 992 клиента ждут в sbuf, пока воркер не закончит с кем-то целую сессию из трёхсот строк, и это ожидание достаётся первой строке каждого клиента: две секунды. Таких строк одна на клиента, в p99 они не попадают, их видно только в max. Пул справедлив к тем, кого уже взял, и несправедлив к тем, кто ждёт. Для коротких запросов вроде HTTP это хорошая сделка, для долгих соединений вроде чата плохая: там нужны события из урока 68.
Прогон эха
Сервер с двумя воркерами, в соседнем терминале два tiny echoclient по очереди:
$ tiny echoservert-pre --workers 2 15217
echoservert-pre: listening on port 15217
server received 4 (4 total) bytes on fd 5
server received 5 (9 total) bytes on fd 5
Счётчик в скобках общий для всех воркеров. Третий клиент, открытый при двух занятых воркерах, получит соединение (его примет главный поток), но ответ только тогда, когда один из первых двух уйдёт: его Job лежит в sbuf. Этот сценарий проверяет первый тест шага.
loadgen.zig: нагрузка на HTTP
Чтобы сравнить три версии TINY и чтобы было что профилировать в конце урока, нужен генератор нагрузки. loadgen запускает conns потоков, и каждый делает reqs запросов подряд, по соединению на запрос (TINY закрывает соединение после ответа). Время каждого запроса идёт в общий массив, у каждого потока своя полоса, так что писать в него можно без блокировки. Счётчики ошибок и байтов атомарные.
//! `loadgen <host> <port> <conns> <reqs> <path>`: генератор нагрузки для
//! TINY. `conns` потоков, каждый делает `reqs` запросов `GET path HTTP/1.0`
//! подряд, по соединению на запрос (TINY закрывает соединение после ответа).
//! Урок 71 гоняет им TINY с пулом под `zt prof`, чтобы на flame graph было
//! что смотреть.
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 bench = @import("bench.zig");
pub const Config = struct {
host: []const u8,
port: u16,
conns: usize,
reqs: usize,
path: []const u8,
};
pub const Result = struct {
requests: usize,
/// Не подключились, оборвались или ответ не `200`.
errors: usize,
bytes: usize,
seconds: f64,
rps: f64,
p50_us: f64,
p99_us: f64,
};
const Shared = struct {
io: Io,
cfg: Config,
request: []const u8,
latencies: []u64,
errors: std.atomic.Value(usize) = .init(0),
bytes: std.atomic.Value(usize) = .init(0),
};
pub fn run(gpa: std.mem.Allocator, io: Io, cfg: Config) !Result {
const request = try std.fmt.allocPrint(gpa, "GET {s} HTTP/1.0\r\nHost: {s}\r\nUser-Agent: loadgen\r\n\r\n", .{ cfg.path, cfg.host });
defer gpa.free(request);
const latencies = try gpa.alloc(u64, cfg.conns * cfg.reqs);
defer gpa.free(latencies);
var shared: Shared = .{ .io = io, .cfg = cfg, .request = request, .latencies = latencies };
const threads = try gpa.alloc(std.Thread, cfg.conns);
defer gpa.free(threads);
const t0 = Io.Timestamp.now(io, .awake);
for (threads, 0..) |*t, i| t.* = try std.Thread.spawn(.{}, worker, .{ &shared, i });
for (threads) |t| t.join();
const seconds = @as(f64, @floatFromInt(t0.untilNow(io, .awake).nanoseconds)) / 1e9;
std.mem.sort(u64, latencies, {}, std.sort.asc(u64));
return .{
.requests = latencies.len,
.errors = shared.errors.load(.monotonic),
.bytes = shared.bytes.load(.monotonic),
.seconds = seconds,
.rps = @as(f64, @floatFromInt(latencies.len)) / seconds,
.p50_us = bench.percentile(latencies, 50),
.p99_us = bench.percentile(latencies, 99),
};
}
fn worker(s: *Shared, id: usize) void {
for (s.latencies[id * s.cfg.reqs ..][0..s.cfg.reqs]) |*lat| {
const t0 = Io.Timestamp.now(s.io, .awake);
const got = once(s) catch 0;
lat.* = @intCast(t0.untilNow(s.io, .awake).nanoseconds);
if (got == 0) _ = s.errors.fetchAdd(1, .monotonic) else _ = s.bytes.fetchAdd(got, .monotonic);
}
}
/// Один запрос: ответ читается до закрытия соединения. 0 значит ошибка.
fn once(s: *Shared) !usize {
const fd = try socket.openClientfd(s.cfg.host, s.cfg.port);
defer socket.close(fd);
try fdio.writen(fd, s.request);
var buf: [16384]u8 = undefined;
var total: usize = 0;
var head: [12]u8 = undefined;
while (true) {
const n = c.read(fd, &buf, buf.len);
if (n < 0) return error.ReadFailed;
if (n == 0) break;
const got: usize = @intCast(n);
if (total < head.len) {
const k = @min(head.len - total, got);
@memcpy(head[total..][0..k], buf[0..k]);
}
total += got;
}
if (total < head.len or !std.mem.eql(u8, head[9..12], "200")) return 0;
return total;
}
pub fn write(out: *Io.Writer, r: Result) !void {
try out.print("requests {d}, errors {d}, bytes {d}\n", .{ r.requests, r.errors, r.bytes });
try out.print("time {d:.3} s, {d:.0} req/s, p50 {d:.1} us, p99 {d:.1} us\n", .{ r.seconds, r.rps, r.p50_us, r.p99_us });
}
bench.percentile пришёл из урока про мультиплексирование, где тот же расчёт делал conc-bench. Ответ считается успешным, только если в строке статуса стоит 200: TINY на ошибку тоже отвечает быстро, и без этой проверки сломанный сервер выглядел бы рекордсменом.
Три версии TINY под нагрузкой
Сборка ReleaseFast, 16 соединений, статика /home.html по 500 запросов на соединение, CGI adder по 50. Apple M4 Max (12 производительных и 4 энергоэффективных ядра), macOS 26.6.2, Zig 0.16.0; столбец «Linux» снят на той же машине в контейнере runner-zig:dev (linux/arm64, OrbStack, ядро 7.0.14, 16 vCPU). Машину в это время грузили посторонние процессы (load average от 13 до 30), так что абсолютные числа гуляют между прогонами в полтора раза, а порядок держится.
$ tiny loadgen 127.0.0.1 8210 16 500 /home.html
requests 8000, errors 0, bytes 1784000
time 0.263 s, 30468 req/s, p50 484.1 us, p99 1381.6 us
macOS Linux (контейнер)
статика, tiny 28467 req/s, p99 0.8 мс 15716 req/s, p99 6.7 мс
статика, tiny --threads 23920 req/s, p99 1.4 мс 7343 req/s, p99 10.0 мс
статика, tiny --pool 8 30468 req/s, p99 1.4 мс 28255 req/s, p99 3.5 мс
CGI, tiny 252 req/s, p99 83 мс 705 req/s, p99 53 мс
CGI, tiny --threads 1261 req/s, p99 24 мс 1802 req/s, p99 26 мс
CGI, tiny --pool 8 2093 req/s, p99 20 мс 3005 req/s, p99 11 мс
На CGI, где каждый запрос это fork, execve и waitpid, картина из книги на обеих системах. Итеративный сервер ждёт каждый CGI целиком, пул обслуживает восемь сразу и даёт больше в 8,3 раза на macOS и в 4,3 раза в Linux, поток на соединение между ними: создать поток на каждый запрос тоже стоит. На статике запрос короче рождения потока, и поток на соединение медленнее итеративного сервера, в Linux вдвое. Пул на статике выигрывает меньше: в Linux в 1,8 раза, на macOS почти ничего. Во что он упирается, догадками решать не будем: в конце урока его покажет профиль.
Параллельные вычисления: psum
До сих пор потоки нужны были для конкурентности: сервер обслуживает многих клиентов, и большую часть времени каждый поток ждёт сеть. Вторая причина завести потоки это параллелизм: одна большая вычислительная задача, у машины шестнадцать ядер, хочется в шестнадцать раз быстрее. Книга разбирает это на самой простой задаче из возможных, сумме 0 + 1 + ... + (n - 1). Ответ известен заранее, n(n - 1)/2, поэтому любую версию легко проверить, а вся разница между версиями сводится к одному вопросу: куда поток складывает свою часть суммы.
Раскладка у всех версий одна. t потоков, поток id складывает полосу [id * n / t, (id + 1) * n / t). Полосы идут подряд и вместе покрывают 0..n без дыр, даже если n не делится на t: округление вниз в обеих границах сдвигает остаток по полосам. Дальше пять вариантов:
psum-mutex: одна общая сумма, мьютекс вокруг каждого сложения. Так пишут, когда впервые узнали про гонки.psum-atomic: одна общая сумма, атомарноеfetchAddна каждой итерации. Мьютекса нет, гонки нет.psum-array: у каждого потока своя ячейкаpsum[id]общего массива, прибавление прямо в память на каждой итерации. Общих данных нет, синхронизация не нужна.psum-array-padded: то же, но каждая ячейка выровнена на свою линию кэша.psum-local: сумма в локальной переменной, то есть в регистре, и одна запись вpsum[id]в самом конце.
Книга показывает первую, третью и пятую. Атомик и выравнивание добавлены потому, что на них видны два эффекта, ради которых всё затевается.
//! `psum` из главы 12: сумма `0 + 1 + ... + (n - 1)` на `t` потоках, каждый
//! складывает свою полосу `[id * n / t, (id + 1) * n / t)`. Версии отличаются
//! только тем, куда поток копит сумму, и в этом весь урок:
//!
//! * `mutex`, `psum-mutex`: общая сумма под мьютексом на каждой итерации.
//! Потоков больше, а работает по-прежнему один, плюс плата за замок.
//! * `atomic`: общая сумма через `fetchAdd` на каждой итерации. Без замка,
//! но одна линия кэша прыгает между ядрами на каждом сложении.
//! * `array`, `psum-array`: у каждого потока свой элемент `psum[id]`, запись
//! в память на каждой итерации. Элементы соседние, лежат в одной линии
//! кэша, и ядра всё равно отнимают её друг у друга: ложное разделение.
//! * `padded`: то же, но каждый элемент выровнен на свою линию кэша.
//! * `local`, `psum-local`: сумма в регистре, в память один раз в конце.
//!
//! Запись «в память на каждой итерации» делается через `volatile`, как
//! в книге: иначе компилятор держал бы `psum[id]` в регистре, и `array` не
//! отличалась бы от `local`. А `local` защищена от обратного: без
//! `doNotOptimizeAway` LLVM заменил бы цикл формулой `n * (n - 1) / 2`.
const std = @import("std");
const Io = std.Io;
pub const Version = enum {
mutex,
atomic,
array,
padded,
local,
pub fn name(v: Version) []const u8 {
return switch (v) {
.mutex => "psum-mutex",
.atomic => "psum-atomic",
.array => "psum-array",
.padded => "psum-array-padded",
.local => "psum-local",
};
}
};
pub const max_threads = 64;
/// Элемент `psum` на отдельной линии кэша (128 байт на Apple M, 64 на x86-64).
const Padded = struct {
value: u64 align(std.atomic.cache_line) = 0,
};
const Shared = struct {
io: Io,
n: u64,
nthreads: u64,
mutex: Io.Mutex = .init,
gsum: u64 = 0,
asum: std.atomic.Value(u64) = .init(0),
array: [max_threads]u64 = @splat(0),
padded: [max_threads]Padded = @splat(.{}),
};
fn worker(comptime version: Version, s: *Shared, id: u64) void {
const lo = id * s.n / s.nthreads;
const hi = (id + 1) * s.n / s.nthreads;
switch (version) {
.mutex => for (lo..hi) |i| {
s.mutex.lockUncancelable(s.io);
s.gsum += i;
s.mutex.unlock(s.io);
},
.atomic => for (lo..hi) |i| {
_ = s.asum.fetchAdd(i, .monotonic);
},
.array => {
const slot: *volatile u64 = &s.array[id];
for (lo..hi) |i| slot.* += i;
},
.padded => {
const slot: *volatile u64 = &s.padded[id].value;
for (lo..hi) |i| slot.* += i;
},
.local => {
var acc: u64 = 0;
for (lo..hi) |i| {
acc += i;
std.mem.doNotOptimizeAway(acc);
}
s.array[id] = acc;
},
}
}
/// Сумма `0..n-1` на `nthreads` потоках версией `version`.
pub fn sum(io: Io, version: Version, n: u64, nthreads: usize) !u64 {
std.debug.assert(nthreads >= 1 and nthreads <= max_threads);
const s = try std.heap.page_allocator.create(Shared);
defer std.heap.page_allocator.destroy(s);
s.* = .{ .io = io, .n = n, .nthreads = nthreads };
var threads: [max_threads]std.Thread = undefined;
for (threads[0..nthreads], 0..) |*t, id| t.* = try switch (version) {
inline else => |v| std.Thread.spawn(.{}, worker, .{ v, s, @as(u64, id) }),
};
for (threads[0..nthreads]) |t| t.join();
return switch (version) {
.mutex => s.gsum,
.atomic => s.asum.load(.monotonic),
.array, .local => total(s.array[0..nthreads]),
.padded => blk: {
var acc: u64 = 0;
for (s.padded[0..nthreads]) |p| acc += p.value;
break :blk acc;
},
};
}
fn total(values: []const u64) u64 {
var acc: u64 = 0;
for (values) |v| acc += v;
return acc;
}
pub fn expected(n: u64) u64 {
return if (n == 0) 0 else n * (n - 1) / 2;
}
pub const BenchConfig = struct {
/// n для `array`, `padded`, `local`: 2^31, как в книге.
n: u64 = 1 << 31,
/// n для `mutex` и `atomic`: у них каждая итерация идёт через общую
/// линию кэша, и на 2^31 замер длился бы минуты.
n_slow: u64 = 1 << 24,
threads: []const usize,
/// Сколько раз мерить каждую точку; в таблицу идёт лучшее время.
/// Минимум честнее среднего: соседние процессы могут только замедлить.
repeat: usize = 3,
versions: []const Version = &.{ .mutex, .atomic, .array, .padded, .local },
};
pub const Sample = struct {
version: Version,
n: u64,
threads: usize,
seconds: f64,
/// T1 / Tt: во сколько раз быстрее одного потока той же версии.
speedup: f64,
/// speedup / t: какая доля каждого потока пошла в дело.
efficiency: f64,
};
/// Замеры по всем версиям и числам потоков. Каждый результат сверяется с
/// формулой: неверная сумма это ошибка, а не строчка в таблице.
pub fn bench(io: Io, cfg: BenchConfig, out: []Sample) ![]Sample {
var count: usize = 0;
for (cfg.versions) |version| {
const n = if (version == .mutex or version == .atomic) cfg.n_slow else cfg.n;
var t1: f64 = 0;
for (cfg.threads) |t| {
var seconds: f64 = std.math.inf(f64);
for (0..@max(cfg.repeat, 1)) |_| {
const start = Io.Timestamp.now(io, .awake);
const got = try sum(io, version, n, t);
seconds = @min(seconds, @as(f64, @floatFromInt(start.untilNow(io, .awake).nanoseconds)) / 1e9);
if (got != expected(n)) return error.WrongSum;
}
if (t1 == 0) t1 = seconds * @as(f64, @floatFromInt(t));
const speedup = t1 / seconds;
out[count] = .{ .version = version, .n = n, .threads = t, .seconds = seconds, .speedup = speedup, .efficiency = speedup / @as(f64, @floatFromInt(t)) };
count += 1;
}
}
return out[0..count];
}
pub fn writeTable(w: *Io.Writer, samples: []const Sample) !void {
try w.print("{s:<18} {s:>11} {s:>7} {s:>9} {s:>8} {s:>10}\n", .{ "version", "n", "threads", "seconds", "speedup", "efficiency" });
for (samples) |s| {
try w.print("{s:<18} {d:>11} {d:>7} {d:>9.3} {d:>8.2} {d:>10.2}\n", .{ s.version.name(), s.n, s.threads, s.seconds, s.speedup, s.efficiency });
}
}
pub fn writeJson(w: *Io.Writer, samples: []const Sample) !void {
const builtin = @import("builtin");
const cpus = std.Thread.getCpuCount() catch 0;
try w.print("{{\"os\":\"{t}\",\"arch\":\"{t}\",\"cpus\":{d},\"samples\":[", .{ builtin.os.tag, builtin.cpu.arch, cpus });
for (samples, 0..) |s, i| {
if (i > 0) try w.writeAll(",\n");
try w.print("{{\"version\":\"{s}\",\"n\":{d},\"threads\":{d},\"seconds\":{d:.4},\"speedup\":{d:.3},\"efficiency\":{d:.3}}}", .{ s.version.name(), s.n, s.threads, s.seconds, s.speedup, s.efficiency });
}
try w.writeAll("]}\n");
}
Две ловушки компилятора, обе знакомые по уроку про пределы компилятора. В ReleaseFast LLVM видит, что psum[id] пишет только этот поток, держит ячейку в регистре и пишет в память один раз, после цикла; psum-array превратилась бы в psum-local. Книга с gcc -O2 этого не встречает, потому что пишет через глобальный массив, а наш *volatile возвращает запись на каждую итерацию. Обратная ловушка у psum-local: цикл acc += i LLVM заменит формулой n(n - 1)/2, и замер покажет ноль. doNotOptimizeAway(acc) на каждой итерации заставляет компилятор считать, что значение кому-то нужно, и цикл остаётся циклом.
В main.zig шаг добавил подкоманду psum-bench (она в диффе выше): потоки 1, 2, 4, 8 и число ядер, лучшее из --repeat прогонов на точку, таблица или JSON.
Ускорение и эффективность
Пусть T₁ это время на одном ядре, а Tₚ на p ядрах. Тогда ускорение это Sₚ = T₁ / Tₚ, а эффективность Eₚ = Sₚ / p. Идеал это Sₚ = p и Eₚ = 100 %: вдвое больше ядер, вдвое меньше время.
Тонкость в том, что брать за T₁. Если это та же параллельная программа, запущенная на одном ядре, ускорение называют относительным. Если лучшая последовательная программа без всякой синхронизации, абсолютным. Относительное ускорение легко приукрасить: версия с мьютексом на одном ядре тоже платит за захват на каждой итерации, её T₁ раздуто, и деление на раздутое число даёт красивую цифру. Честнее абсолютное, но для него нужна вторая, последовательная программа, а её не всегда есть смысл писать.
Второе различие про то, что держим постоянным. Сильное масштабирование: задача та же, n = 2³¹ при любом числе потоков, вопрос «насколько быстрее». Слабое масштабирование: работа на поток постоянна, при удвоении потоков удваивается и n, вопрос «какую задачу можно решить за то же время». Слабое ближе к тому, зачем покупают большие машины (больше данных, мельче сетка в физической модели), сильное к задачам, где объём задан извне, как обработка сигнала с одного датчика в реальном времени. psum-bench меряет сильное. А замер общей кучи zl в прошлом уроке был слабым: каждый поток считал свой (fib 23), работа росла вместе с потоками, и «ускорение» там означало работу в единицу времени. Слабое масштабирование psum остаётся тебе в упражнении.
Замер
Все пять версий на одной машине, n = 2³¹ для array, padded и local и n = 2²⁴ для mutex и atomic (на 2³¹ они шли бы минутами). Лучшее из пяти прогонов на точку, ReleaseFast. Apple M4 Max, 12 производительных и 4 энергоэффективных ядра, 64 ГБ, macOS 26.6.2. Машину в это время грузили соседние процессы: load average 6,9 на старте прогона и 23,8 в конце, так что на многих потоках время может быть завышено; повторный прогон дал кривые той же формы:
$ tiny psum-bench --threads 1,2,4,8,12,16 --repeat 5
version n threads seconds speedup efficiency
psum-mutex 16777216 1 0.035 1.00 1.00
psum-mutex 16777216 2 0.231 0.15 0.07
psum-mutex 16777216 4 0.160 0.22 0.06
psum-mutex 16777216 8 0.139 0.25 0.03
psum-mutex 16777216 12 0.121 0.29 0.02
psum-mutex 16777216 16 0.123 0.28 0.02
psum-atomic 16777216 1 0.034 1.00 1.00
psum-atomic 16777216 2 0.071 0.47 0.24
psum-atomic 16777216 4 0.168 0.20 0.05
psum-atomic 16777216 8 0.424 0.08 0.01
psum-atomic 16777216 12 0.602 0.06 0.01
psum-atomic 16777216 16 0.529 0.06 0.00
psum-array 2147483648 1 1.848 1.00 1.00
psum-array 2147483648 2 1.687 1.10 0.55
psum-array 2147483648 4 0.898 2.06 0.52
psum-array 2147483648 8 0.552 3.35 0.42
psum-array 2147483648 12 0.425 4.35 0.36
psum-array 2147483648 16 0.393 4.70 0.29
psum-array-padded 2147483648 1 2.341 1.00 1.00
psum-array-padded 2147483648 2 1.223 1.91 0.96
psum-array-padded 2147483648 4 0.539 4.34 1.09
psum-array-padded 2147483648 8 0.231 10.15 1.27
psum-array-padded 2147483648 12 0.137 17.15 1.43
psum-array-padded 2147483648 16 0.143 16.34 1.02
psum-local 2147483648 1 0.711 1.00 1.00
psum-local 2147483648 2 0.353 2.02 1.01
psum-local 2147483648 4 0.174 4.09 1.02
psum-local 2147483648 8 0.096 7.39 0.92
psum-local 2147483648 12 0.072 9.81 0.82
psum-local 2147483648 16 0.063 11.33 0.71
Ускорение в таблице относительное: каждая версия к самой себе на одном потоке. Виджет ниже строит этот же прогон, взятый в JSON (tiny psum-bench --json); у каждой версии своё n, поэтому сравнивается время на элемент. Переключи базу на «лучшая последовательная», и ускорение станет абсолютным: все версии делятся на один поток psum-local. Кривая Амдала с ползунком понадобится через раздел.
Одна строка таблицы выглядит как чудо: у psum-array-padded ускорение больше числа потоков, 10,15 на восьми и 17,15 на двенадцати. Это та самая ловушка относительного ускорения из раздела выше. Один поток padded медленнее, чем у array (2,34 с против 1,85), база раздута, и деление на неё даёт красивую цифру. Переключи виджет на лучшую последовательную базу: там padded на двенадцати потоках даёт 5,2, а psum-local 9,8. Ложное разделение к тому же зависит от того, как планировщик разложил потоки по ядрам, а этого мы не контролируем, так что читать надо форму кривых, а не третий знак.
Мьютекс и атомик медленнее одного потока. На двух потоках psum-mutex почти в семь раз медленнее, чем на одном, psum-atomic вдвое, и быстрее одного потока они не становятся ни на каком числе потоков. Сумма одна на всех, и работы, которую можно делать параллельно, нет вовсе: каждое сложение ждёт предыдущего, чьим бы оно ни было. Зато появилась цена. Линия кэша с общей суммой (а у мьютекса ещё и линия с его словом) переезжает между ядрами на каждой итерации: ядро, которое хочет писать, обязано получить линию в исключительное владение, и копия у соседа становится недействительной. На одном потоке линия всё время живёт в L1 своего ядра, и lock с unlock стоят несколько наносекунд. Атомик быстрее мьютекса только потому, что одна инструкция вместо двух операций и возможного сна в futex; линия у него скачет точно так же.
Ложное разделение. У psum-array общих данных нет, гонки нет, и всё-таки она масштабируется хуже padded, а на двух потоках почти не быстрее одного (1,10). Восемь ячеек u64 лежат в одной линии (на Apple M линия 128 байт, std.atomic.cache_line, помнишь по уроку про устройство кэша), и каждая запись одного потока отнимает линию у соседа, который пишет в свою ячейку той же линии. Для протокола когерентности это одна переменная, хотя в коде их восемь. Мы уже ловили это двумя счётчиками в уроке про код, дружественный кэшу; Padded с align(std.atomic.cache_line) лечит так же: у каждой ячейки своя линия, и на восьми потоках padded считает за 0,23 с, а array за 0,55. Обрати внимание, что на одном потоке padded, наоборот, медленнее array (2,34 с против 1,85): соседние ячейки на одном потоке линию не делят ни с кем, так что эта разница не про когерентность. Её причину одним замером не установить, и мы не будем делать вид, что можем.
Локальный аккумулятор выигрывает у всех. psum-local в два с половиной раза быстрее array уже на одном потоке, потому что сложение в регистре, а не через память, и быстрее всех на любом числе потоков: 11,33 на шестнадцати. Урок книги в одной строке: минимальное, на первый взгляд, изменение, где копить сумму, меняет время на порядки. Синхронизация на каждой итерации это не мелкая накладная, это отказ от параллелизма.
Почему не шестнадцать. Эффективность psum-local держится от 92 до 102 процентов до восьми потоков (чуть больше ста это шум замера) и падает до 82 на двенадцати и 71 на шестнадцати. Четыре ядра из шестнадцати энергоэффективные, и поток, попавший на такое ядро, считает свою полосу медленнее; полосы у нас равные, так что все ждут самого медленного. Книга на четырёхъядерной машине видела другое: после четырёх потоков время перестаёт падать, потому что лишние потоки делят те же ядра и платят за переключения. Отсюда и совет писать параллельные программы на один поток на ядро.
Закон Амдала числом
В уроке про профилирование закон Амдала отвечал на вопрос, какую функцию ускорять. Для параллельной программы он звучит так же: если доля s времени последовательная, а остальное делится на p ядер без потерь, то
S(p) = 1 / (s + (1 - s) / p) предел при p → ∞: 1 / s
При s = 5 % на шестнадцати ядрах S = 1 / (0,05 + 0,95 / 16) ≈ 9,1, а бесконечность ядер даст не больше 20. Пять процентов последовательного кода, и треть машины уже простаивает. Ты уже видел это вживую: на маленькой куче zl из прошлого урока восемь потоков стояли в остановке мира 26,3 мс из 37,5, сборка шла последовательно, и ускорение застряло около 2,4. Поставь в виджете ползунок последовательной доли на 5 процентов и сравни пунктир Амдала с кривыми: psum-local идёт выше, padded с раздутой базой ещё выше, остальные глубоко внизу.
Закон можно повернуть обратно: по замеру найти, какая последовательная доля объяснила бы полученное ускорение. Это метрика Карпа и Флэтт:
e = (1/S - 1/p) / (1 - 1/p)
Для psum-local на 16 потоках e = (1/11,33 - 1/16) / (1 - 1/16) = (0,0883 - 0,0625) / 0,9375 ≈ 0,028. Как будто 2,8 процента программы последовательные. Последняя колонка таблицы в виджете считает то же для каждой версии. Но у psum последовательной части почти нет: запуск потоков, join и сложение шестнадцати чисел. Эти почти три процента это энергоэффективные ядра и шум нагруженной машины, записанные в виде последовательной доли. Метрика полезна другим: если e растёт вместе с p, программа теряет не на последовательном коде (он от p не зависит), а на накладных расходах, которые растут с числом потоков, то есть на синхронизации и обмене данными. У psum-mutex и psum-atomic она уходит за единицу, и это честный диагноз: ускорения нет, есть только расходы.
Тесты шага TINY
//! Шаг 71: предварительно созданный пул потоков (эхо и TINY), `Once`,
//! генератор нагрузки и `psum` в пяти версиях.
const std = @import("std");
const tiny = @import("tiny");
const support = @import("support.zig");
const conc = tiny.conc;
const socket = tiny.socket;
const fdio = tiny.fdio;
const testing = std.testing;
const io = testing.io;
const c = std.c;
fn connect(port: u16, timeout_ms: u31) !c.fd_t {
const fd = try socket.openClientfd("127.0.0.1", port);
const tv: c.timeval = .{ .sec = timeout_ms / 1000, .usec = @as(i32, timeout_ms % 1000) * 1000 };
_ = c.setsockopt(fd, c.SOL.SOCKET, c.SO.RCVTIMEO, &tv, @sizeOf(c.timeval));
return fd;
}
fn poolServe(listenfd: c.fd_t, n: usize) void {
conc.echo_pre.serve(testing.allocator, io, listenfd, .{ .workers = 2, .queue = 4, .max_clients = n }, conc.echo_pre.EchoCtx{ .io = io, .quiet = true }, conc.echo_pre.echoCnt) catch |err| std.debug.print("echo_pre: {t}\n", .{err});
}
test "echoservert_pre: два воркера, третий клиент ждёт в sbuf, пока один не освободится" {
const listenfd = try socket.openListenfd(0);
defer socket.close(listenfd);
const port = socket.localPort(listenfd).?;
const before = conc.echo_pre.byteCount(io);
const thread = try std.Thread.spawn(.{}, poolServe, .{ listenfd, 3 });
var buf: [16]u8 = undefined;
const a = try connect(port, 5000);
const b = try connect(port, 5000);
try fdio.writen(a, "a\n");
try testing.expectEqual(@as(usize, 2), try fdio.readn(a, buf[0..2]));
try fdio.writen(b, "b\n");
try testing.expectEqual(@as(usize, 2), try fdio.readn(b, buf[0..2]));
// Оба воркера заняты: третьему никто не ответит, пока `a` не уйдёт.
const third = try connect(port, 200);
try fdio.writen(third, "c\n");
try testing.expectError(error.ReadFailed, fdio.readn(third, buf[0..2]));
socket.close(a);
const tv: c.timeval = .{ .sec = 5, .usec = 0 };
_ = c.setsockopt(third, c.SOL.SOCKET, c.SO.RCVTIMEO, &tv, @sizeOf(c.timeval));
try testing.expectEqual(@as(usize, 2), try fdio.readn(third, buf[0..2]));
try testing.expectEqualStrings("c\n", buf[0..2]);
socket.close(b);
socket.close(third);
thread.join();
// Счётчик общий на все воркеры: шесть байтов от трёх клиентов.
try testing.expectEqual(before + 6, conc.echo_pre.byteCount(io));
}
var once_calls: std.atomic.Value(usize) = .init(0);
fn onceTarget() void {
_ = once_calls.fetchAdd(1, .monotonic);
}
fn callOnce(once: *conc.echo_pre.Once) void {
once.call(io, onceTarget);
}
test "Once: восемь потоков, одна инициализация" {
var once: conc.echo_pre.Once = .{};
var threads: [8]std.Thread = undefined;
for (&threads) |*t| t.* = try std.Thread.spawn(.{}, callOnce, .{&once});
for (threads) |t| t.join();
once.call(io, onceTarget);
try testing.expectEqual(@as(usize, 1), once_calls.load(.monotonic));
}
fn tinyPool(server: *tiny.Server, n: usize) void {
conc.echo_pre.serve(testing.allocator, io, server.listenfd, .{ .workers = 3, .max_clients = n }, server, conc.echo_pre.tinyHandle) catch |err| std.debug.print("tiny pool: {t}\n", .{err});
}
test "tiny --pool 3 под loadgen: 4 соединения по 5 запросов, все 200" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
var root: support.Root = try .create(arena_state.allocator());
defer root.cleanup();
var server: tiny.Server = try .init(testing.allocator, io, .{ .port = 0, .root = root.path, .quiet = true });
defer server.deinit();
const thread = try std.Thread.spawn(.{}, tinyPool, .{ &server, 20 });
const r = try conc.loadgen.run(testing.allocator, io, .{ .host = "127.0.0.1", .port = server.port, .conns = 4, .reqs = 5, .path = "/home.html" });
thread.join();
try testing.expectEqual(@as(usize, 20), r.requests);
try testing.expectEqual(@as(usize, 0), r.errors);
// Каждый ответ это заголовок плюс 114 байт home.html.
try testing.expect(r.bytes > 20 * 114);
try testing.expect(r.p50_us <= r.p99_us);
}
test "psum: пять версий на 1, 3 и 4 потоках дают n(n-1)/2" {
const n = 100_003;
for ([_]conc.psum.Version{ .mutex, .atomic, .array, .padded, .local }) |v| {
for ([_]usize{ 1, 3, 4 }) |t| {
try testing.expectEqual(conc.psum.expected(n), try conc.psum.sum(io, v, n, t));
}
}
}
test "psum bench: ускорение одного потока 1, эффективность это ускорение на поток" {
var samples: [4]conc.psum.Sample = undefined;
const got = try conc.psum.bench(io, .{ .n = 1 << 20, .n_slow = 1 << 12, .threads = &.{ 1, 2 }, .versions = &.{ .local, .atomic } }, &samples);
try testing.expectEqual(@as(usize, 4), got.len);
try testing.expectApproxEqAbs(@as(f64, 1), got[0].speedup, 1e-9);
try testing.expectApproxEqAbs(got[1].speedup / 2, got[1].efficiency, 1e-9);
}
Тест пула детерминирован без единого sleep: два воркера заняты двумя клиентами, и третий клиент с таймаутом чтения 200 мс обязан получить ReadFailed. Как только a закрыт, воркер освобождается, Job третьего выходит из sbuf, и ответ приходит. Тест psum гоняет все версии на 1, 3 и 4 потоках при нечётном n, так что остаток в последней полосе тоже проверен. Бенч в тестах идёт на маленьком n и проверяет только арифметику: у одного потока ускорение ровно 1, эффективность это ускорение на поток. Скорость в тестах не проверяется: на нагруженной машине любой порог ложно падал бы.
$ zig build test -Dstep=71 --summary all
test success
+- run test 5 pass (5 total) 1s MaxRSS:5M
Шаг проекта: zbox serve под нагрузкой
zbox serve из урока про песочницу итеративный: пока одна задача компилируется десять секунд, остальные студенты ждут в очереди listen. Раннеру курса нужно другое: несколько задач сразу, но не больше, чем выдержит машина, а тем, кому места не хватило, честный отказ, а не зависание. План шага:
- Главный поток делает
acceptи кладёт дескриптор в очередь. Очередь полна: главный поток сам отвечает503 Service Unavailableс заголовкомRetry-After: 1и закрывает соединение. - Потоки пула (
--workers, по умолчанию 8) забирают соединения, читают HTTP-запрос, отвечают наGET /stats,404и400сразу. POST /runсначала занимает место у семафора (--jobs, по умолчанию 4), и только потом запускает задачу. Потоков больше, чем мест, чтобы разбор запроса и отказ на кривой JSON не стояли в очереди к песочнице.- Задача это дочерний процесс
zbox job, которому тело запроса уходит в stdin, а ответ JSON приходит из stdout.
Почему процесс, а не run() прямо в потоке пула, стоит сказать до кода. run из урока про сигналы держит состояние в глобальных переменных сторожа: группу процессов программы, крайний срок, флаги «ребёнок вышел» и «будильник сработал». Будильник setitimer один на процесс, маска и обработчики сигналов тоже. Два run в двух потоках одновременно перепутали бы чужие SIGCHLD и SIGALRM: один поток убил бы программу соседа по его будильнику. Это функция, которая небезопасна для потоков потому, что держит общее состояние; разбор всех классов таких функций ждёт в следующем уроке. Лечится она здесь самым грубым и самым надёжным способом: у каждой задачи свой процесс, а значит свои глобальные переменные, свои сигналы и свой будильник.
queue.zig: очередь с отказом
Это sbuf из прошлого урока, только на мьютексе и двух условных переменных. Семафоры не умеют того, что нужно пулу сверх sbuf: положить без ожидания и сказать «места нет» (tryPut для ответа 503) и закрыть очередь, разбудив всех, кто ждёт.
//! Ограниченная очередь задач между главным потоком и пулом (шаг 71).
//! Это `sbuf` из урока 70, только на `Mutex` и двух `Condition`, а не на
//! трёх семафорах: так проще закрыть очередь и разбудить всех разом.
//!
//! В Zig 0.16 `std.Thread.Mutex` и `std.Thread.Condition` больше нет:
//! примитивы живут в `std.Io` и принимают `io`. С `Io.Threaded` это те же
//! futex, и звать их можно из потоков `std.Thread.spawn`.
const std = @import("std");
const Io = std.Io;
pub fn Queue(comptime T: type) type {
return struct {
const Self = @This();
io: Io,
buffer: []T,
/// Индекс первого элемента и число элементов в кольце.
head: usize = 0,
len: usize = 0,
closed: bool = false,
mutex: Io.Mutex = .init,
/// Сигнал «появился элемент» для потребителей.
not_empty: Io.Condition = .init,
/// Сигнал «освободилось место» для блокирующего `put`.
not_full: Io.Condition = .init,
pub fn init(io: Io, buffer: []T) Self {
return .{ .io = io, .buffer = buffer };
}
/// Положить без ожидания. `false`, если места нет или очередь
/// закрыта: сервер ответит клиенту 503, а не заставит его ждать.
pub fn tryPut(q: *Self, item: T) bool {
q.mutex.lockUncancelable(q.io);
defer q.mutex.unlock(q.io);
if (q.closed or q.len == q.buffer.len) return false;
q.push(item);
return true;
}
/// Положить, дождавшись места. `false`, если очередь закрыли.
pub fn put(q: *Self, item: T) bool {
q.mutex.lockUncancelable(q.io);
defer q.mutex.unlock(q.io);
while (!q.closed and q.len == q.buffer.len) q.not_full.waitUncancelable(q.io, &q.mutex);
if (q.closed) return false;
q.push(item);
return true;
}
/// Взять, дождавшись элемента. `null` значит, что очередь закрыта
/// и пуста: потоку пула пора выходить. Оставшееся после `close`
/// ещё раздаётся, принятые соединения не бросаем.
pub fn get(q: *Self) ?T {
q.mutex.lockUncancelable(q.io);
defer q.mutex.unlock(q.io);
while (!q.closed and q.len == 0) q.not_empty.waitUncancelable(q.io, &q.mutex);
if (q.len == 0) return null;
const item = q.buffer[q.head];
q.head = (q.head + 1) % q.buffer.len;
q.len -= 1;
q.not_full.signal(q.io);
return item;
}
/// После `close` новые элементы не принимаются, а все ждущие
/// просыпаются: `broadcast`, а не `signal`, иначе проснётся один.
pub fn close(q: *Self) void {
q.mutex.lockUncancelable(q.io);
defer q.mutex.unlock(q.io);
q.closed = true;
q.not_empty.broadcast(q.io);
q.not_full.broadcast(q.io);
}
/// Вызывается под мьютексом.
fn push(q: *Self, item: T) void {
q.buffer[(q.head + q.len) % q.buffer.len] = item;
q.len += 1;
q.not_empty.signal(q.io);
}
};
}
Условие ждётся в цикле while, а не в if, по той же причине, что в SbufCond прошлого урока: проснуться можно без повода, и слот может успеть забрать другой поток между signal и моментом, когда проснувшийся снова держит мьютекс. get после close продолжает раздавать то, что уже лежит в кольце, и возвращает null только на пустой закрытой очереди: принятые соединения не бросаются. Везде lockUncancelable и waitUncancelable: в Zig 0.16 ожидание у примитивов std.Io это точка отмены, а потоки std.Thread отменять некому.
pool.zig: потоки и семафор на запуски
//! Пул потоков (шаг 71): `echoservert_pre` из главы 12. Потоки создаются
//! один раз при старте, главный поток кладёт работу в очередь, потоки
//! пула забирают её по одной. Отдельно `Slots`: семафор на число
//! одновременных запусков, общий для пула и для модели на `std.Io`.
const std = @import("std");
const Io = std.Io;
const Queue = @import("queue.zig").Queue;
/// Счётчик свободных мест для запуска задачи плюс статистика: сколько
/// выполняется сейчас, сколько выполнялось одновременно в пике, сколько
/// всего выполнено. Пик и есть проверка того, что семафор работает.
pub const Slots = struct {
semaphore: Io.Semaphore,
running: std.atomic.Value(u32) = .init(0),
peak: std.atomic.Value(u32) = .init(0),
done: std.atomic.Value(u64) = .init(0),
pub fn init(jobs: u16) Slots {
return .{ .semaphore = .{ .permits = jobs } };
}
/// Занять место. Ждёт, пока освободится, отмена прерывает ожидание.
pub fn acquire(s: *Slots, io: Io) Io.Cancelable!void {
try s.semaphore.wait(io);
const now = s.running.fetchAdd(1, .acq_rel) + 1;
_ = s.peak.fetchMax(now, .acq_rel);
}
pub fn release(s: *Slots, io: Io) void {
_ = s.running.fetchSub(1, .acq_rel);
_ = s.done.fetchAdd(1, .acq_rel);
s.semaphore.post(io);
}
};
/// `workers` потоков, каждый в цикле берёт задачу из очереди и зовёт
/// `handle`. Задача передаётся по значению: у каждого потока своя копия,
/// общей переменной, которую перезапишет следующий `accept`, нет.
pub fn Pool(comptime Task: type, comptime Context: type, comptime handle: fn (Context, Task) void) type {
return struct {
const Self = @This();
queue: Queue(Task),
context: Context,
threads: []std.Thread,
/// Сколько задач не влезло в очередь и получило отказ.
rejected: std.atomic.Value(u64) = .init(0),
/// Память под очередь и потоки даёт вызывающий. Потоки стартуют
/// сразу; `self` не должен переезжать, пока пул жив.
pub fn start(self: *Self, io: Io, context: Context, queue_buffer: []Task, threads: []std.Thread) std.Thread.SpawnError!void {
self.* = .{ .queue = .init(io, queue_buffer), .context = context, .threads = threads };
for (threads, 0..) |*thread, started| {
thread.* = std.Thread.spawn(.{}, worker, .{self}) catch |err| {
self.queue.close();
for (threads[0..started]) |t| t.join();
return err;
};
}
}
/// Отдать задачу без ожидания. `false`, если очередь полна.
pub fn submit(self: *Self, task: Task) bool {
if (self.queue.tryPut(task)) return true;
_ = self.rejected.fetchAdd(1, .monotonic);
return false;
}
/// Закрыть очередь, дать потокам доделать принятое и дождаться их.
pub fn stop(self: *Self) void {
self.queue.close();
for (self.threads) |t| t.join();
}
fn worker(self: *Self) void {
while (self.queue.get()) |task| handle(self.context, task);
}
};
}
Pool это echoservert_pre в общем виде: тип задачи, контекст и обработчик приходят параметрами времени компиляции, и тесты подставляют в него поддельную задачу вместо соединения. Память под очередь и потоки даёт вызывающий, а сам Pool лежит у него на стеке, поэтому start принимает self указателем и потоки получают этот указатель. Если spawn посередине вернёт ошибку, уже запущенные потоки дождутся закрытой очереди и выйдут, и только после их join ошибка уйдёт наверх: иначе они продолжали бы читать очередь из кадра, которого уже нет.
Slots это семафор с разрешениями по числу мест плюс три атомарных счётчика для /stats. peak обновляется через fetchMax: атомарное «записать, если больше». Два потока, которые одновременно увидели running 3 и 4, оба зовут fetchMax, и в peak останется 4 при любом порядке. Обычное «прочитал, сравнил, записал» здесь было бы гонкой: второй поток мог бы записать своё 3 поверх чужого 4. Пик и есть проверка того, что семафор работает: при --jobs 3 он обязан быть ровно 3.
job.zig: задача в своём процессе
//! Задача сервера это отдельный процесс `zbox job` (шаг 71).
//!
//! Почему не `run()` прямо в потоке пула: `run` держит состояние в
//! глобальных переменных сторожа (группа, крайний срок, флаги сигналов),
//! а будильник `setitimer` и маска сигналов общие на весь процесс. Два
//! потока с `run` одновременно перепутали бы чужие SIGCHLD и SIGALRM. Это
//! первый класс небезопасных для потоков функций из урока 72. Лечится он
//! здесь отдельным процессом на задачу: у каждого `zbox job` свои
//! глобальные переменные, свои сигналы и свой будильник.
const std = @import("std");
const Io = std.Io;
/// Ответ `zbox job` больше этого не бывает: два потока по 64 КБ, каждый
/// байт в худшем случае экранируется в шесть, плюс поля.
const max_reply = 1 << 20;
/// Запускает `exe job [--isolate ROOTFS]`, отдаёт ему тело запроса в stdin
/// и возвращает его stdout: готовую строку JSON с полем `stage`. Ошибка
/// по дороге убивает ребёнка: `kill` шлёт ему SIGTERM и ждёт.
pub fn run(io: Io, arena: std.mem.Allocator, exe: []const u8, rootfs: ?[]const u8, body: []const u8) ![]u8 {
const argv: []const []const u8 = if (rootfs) |dir| &.{ exe, "job", "--isolate", dir } else &.{ exe, "job" };
var child = try std.process.spawn(io, .{ .argv = argv, .stdin = .pipe, .stdout = .pipe, .stderr = .inherit });
errdefer child.kill(io);
var in_buf: [4096]u8 = undefined;
var stdin = child.stdin.?.writerStreaming(io, &in_buf);
stdin.interface.writeAll(body) catch return stdin.err.?;
stdin.interface.flush() catch return stdin.err.?;
child.stdin.?.close(io);
child.stdin = null;
var out_buf: [4096]u8 = undefined;
var stdout = child.stdout.?.readerStreaming(io, &out_buf);
const reply = stdout.interface.allocRemaining(arena, .limited(max_reply)) catch |err| switch (err) {
error.ReadFailed => return stdout.err.?,
else => |e| return e,
};
const term = try child.wait(io);
if (term != .exited or term.exited != 0 or reply.len == 0) return error.JobFailed;
return reply;
}
std.process.spawn в Zig 0.16 принимает io и пайпы для stdin и stdout. Порядок важен: сначала всё тело в stdin и close, иначе zbox job будет ждать конца ввода, пока мы ждём его вывода, и оба зависнут. errdefer child.kill(io) на любой ошибке по дороге шлёт ребёнку SIGTERM и дожидается его, чтобы не оставить зомби.
http.zig: транзакция и отказ
TINY умеет разбирать запрос, но его doit приватный, а ответ он пишет сам и заголовка Retry-After не знает. Поэтому транзакцию пишем свою, а разбор строки запроса и заголовков берём у TINY. Функция работает над любыми *Io.Reader и *Io.Writer: над дескриптором через tiny.fdio сейчас и над потоком std.Io.net в уроке 73, поэтому код ответа 504 тоже уже здесь.
//! Одна HTTP-транзакция над любыми `*Io.Reader` и `*Io.Writer`: над
//! дескриптором (`tiny.fdio`, пул шага 71) и над `std.Io.net.Stream`
//! (шаг 73). Разбор строки запроса и заголовков берём у TINY, ответ
//! пишем сами: TINY не умеет заголовок `Retry-After`.
const std = @import("std");
const c = std.c;
const Io = std.Io;
const tiny = @import("tiny");
pub const Request = struct {
method: tiny.Method,
path: []const u8,
body: []const u8,
};
pub const Response = struct {
status: u16 = 200,
body: []const u8 = "",
/// Через сколько секунд клиенту стоит повторить: 503 при перегрузке.
retry_after: ?u32 = null,
};
/// Строка запроса, заголовки, тело по `Content-Length`, ответ `route`.
/// Ответ уходит из буфера на выходе любым путём, в том числе ранний 400.
pub fn transaction(
arena: std.mem.Allocator,
in: *Io.Reader,
out: *Io.Writer,
max_body: usize,
context: anytype,
comptime route: fn (@TypeOf(context), std.mem.Allocator, Request) anyerror!Response,
) !void {
defer out.flush() catch {};
const raw_line = in.takeDelimiterInclusive('\n') catch |err| switch (err) {
error.EndOfStream => return,
else => return reply(out, .{ .status = 400, .body = "{\"error\":\"BadRequestLine\"}\n" }),
};
const line = tiny.request.parseRequestLine(raw_line) catch
return reply(out, .{ .status = 400, .body = "{\"error\":\"BadRequestLine\"}\n" });
var headers: tiny.request.Headers = .{};
tiny.request.readRequestHeaders(in, &headers) catch
return reply(out, .{ .status = 400, .body = "{\"error\":\"BadHeaders\"}\n" });
const length = headers.content_length orelse 0;
if (length > max_body) return reply(out, .{ .status = 413, .body = "{\"error\":\"TooLarge\"}\n" });
const body = try arena.alloc(u8, length);
in.readSliceAll(body) catch return reply(out, .{ .status = 400, .body = "{\"error\":\"ShortBody\"}\n" });
const path = line.uri[0 .. std.mem.indexOfScalar(u8, line.uri, '?') orelse line.uri.len];
const response = route(context, arena, .{ .method = line.method, .path = path, .body = body }) catch |err|
Response{ .status = 500, .body = try std.fmt.allocPrint(arena, "{{\"error\":\"{t}\"}}\n", .{err}) };
try reply(out, response);
}
pub fn reply(out: *Io.Writer, response: Response) Io.Writer.Error!void {
try out.print("HTTP/1.0 {d} {s}\r\n", .{ response.status, statusText(response.status) });
try out.writeAll("Connection: close\r\nContent-Type: application/json\r\n");
if (response.retry_after) |seconds| try out.print("Retry-After: {d}\r\n", .{seconds});
try out.print("Content-Length: {d}\r\n\r\n", .{response.body.len});
try out.writeAll(response.body);
}
fn statusText(status: u16) []const u8 {
return switch (status) {
503 => "Service Unavailable",
504 => "Gateway Timeout",
else => tiny.request.statusText(status),
};
}
/// Ответ 503 прямо из главного потока, без очереди и без разбора запроса.
pub const busy: Response = .{ .status = 503, .body = "{\"error\":\"busy\"}\n", .retry_after = 1 };
/// Отказ соединению, которому не нашлось места. Сначала вычитываем то,
/// что клиент уже прислал: `close` с непрочитанными байтами в сокете шлёт
/// RST вместо FIN, и клиент может потерять наш ответ.
/// ponytail: ждём запрос не дольше 50 мс; кто прислал позже, может получить RST.
pub fn reject(fd: c.fd_t, response: Response) void {
var poll_fd = [_]c.pollfd{.{ .fd = fd, .events = c.POLL.IN, .revents = 0 }};
_ = c.poll(&poll_fd, 1, 50);
const flags = c.fcntl(fd, c.F.GETFL);
_ = c.fcntl(fd, c.F.SETFL, flags | @as(c_int, @bitCast(c.O{ .NONBLOCK = true })));
var sink: [4096]u8 = undefined;
while (c.read(fd, &sink, sink.len) > 0) {}
_ = c.fcntl(fd, c.F.SETFL, flags);
var buf: [256]u8 = undefined;
var out: Io.Writer = .fixed(&buf);
reply(&out, response) catch return;
tiny.fdio.writen(fd, out.buffered()) catch {};
}
/// Таймауты на чтение и запись сокета: клиент, который открыл соединение
/// и молчит, не должен навсегда занять поток.
pub fn setTimeouts(fd: c.fd_t, ms: u32) void {
const tv: c.timeval = .{ .sec = @intCast(ms / 1000), .usec = @intCast((ms % 1000) * 1000) };
_ = c.setsockopt(fd, c.SOL.SOCKET, c.SO.RCVTIMEO, std.mem.asBytes(&tv), @sizeOf(c.timeval));
_ = c.setsockopt(fd, c.SOL.SOCKET, c.SO.SNDTIMEO, std.mem.asBytes(&tv), @sizeOf(c.timeval));
}
reject решает мелкую, но настоящую проблему. Главный поток отвечает 503, не прочитав запроса. Если закрыть сокет, в приёмном буфере которого лежат непрочитанные байты, ядро пошлёт клиенту не FIN, а RST, и клиент может потерять наш ответ, даже если он уже ушёл. Поэтому reject ждёт запрос не дольше 50 мс, вычитывает всё без блокировки и только потом отвечает. setTimeouts ставит SO_RCVTIMEO и SO_SNDTIMEO на соединение: клиент, который открыл сокет и молчит, иначе навсегда занял бы поток пула. Это та же атака медленного клиента, что и в уроке про мультиплексирование, только теперь она съедает не весь сервер, а одного воркера на тридцать секунд.
server.zig: главный цикл
//! `zbox serve` под нагрузкой (шаг 71): главный поток делает `accept` и
//! кладёт дескриптор в очередь, потоки пула разбирают запрос и запускают
//! задачу. Задач одновременно не больше `jobs` (семафор `Slots`),
//! соединений в очереди не больше `queue`; сверх этого главный поток сам
//! отвечает 503.
const std = @import("std");
const builtin = @import("builtin");
const c = std.c;
const Io = std.Io;
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),
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 не ждёт места.
_ = 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 svc.runSlot(arena, req.body);
}
/// Занять место и выполнить задачу до конца.
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 .{ .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}}}\n", .{
svc.slots.running.load(.acquire),
svc.slots.peak.load(.acquire),
svc.slots.done.load(.acquire),
svc.rejected.load(.acquire),
}) };
}
};
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));
}
/// Порт в 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();
}
Разбор POST /run идёт до acquire: кривой JSON получает 400 сразу, даже если все четыре места заняты. У каждого соединения своя арена поверх page_allocator, общего аллокатора с блокировкой между потоками здесь нет. Функцию acceptCloexec разберём отдельно ниже: вокруг неё главная история шага.
Мелкие правки
Запрос научился нести готовую команду вместо исходника. Нагрузочному тесту нужна задача ровно на 200 мс, а не компилятор на несколько секунд и сотни мегабайт; стадия у такой задачи одна, run, с теми же лимитами и в том же временном каталоге. Заодно разбор флагов serve переехал в parseServe: любой флаг пула включает модель pool.
/// Что программа прочитает из stdin.
stdin: []const u8 = "",
+ /// Готовая команда вместо исходника (шаг 71): запускается сразу, без
+ /// стадии сборки. Нагрузочному тесту нужна задача на 200 мс, а не
+ /// компилятор на секунды и сотни мегабайт.
+ argv: []const []const u8 = &.{},
};
@@
else => return error.BadJson,
};
- if (request.source.len == 0) return error.EmptySource;
+ 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;
return request;
@@
pub const Error = error{ TempDirFailed, WriteFailed } || run_mod.Error;
-/// Кладёт исходник во временный каталог, собирает и запускает. Каталог
+/// Кладёт исходник во временный каталог, собирает и запускает (или
+/// запускает готовую команду из `argv`). Каталог
/// удаляется на выходе. `zig` берётся из `ZBOX_ZIG` или из `PATH`.
pub fn runSource(arena: std.mem.Allocator, io: std.Io, request: Request) Error!Reply {
@@
dir.writeFile(io, .{ .sub_path = "main.zig", .data = request.source }) catch return error.WriteFailed;
dir.writeFile(io, .{ .sub_path = "stdin", .data = request.stdin }) catch return error.WriteFailed;
+ const stdin_path = try std.fmt.allocPrint(arena, "{s}/stdin", .{dir_path});
+
+ // Готовая команда: одна стадия, те же лимиты и тот же каталог.
+ if (request.argv.len > 0) return .{ .stage = .run, .outcome = try run_mod.run(arena, .{
+ .time_ms = request.time_ms,
+ .mem_mb = request.mem_mb,
+ .stdin_path = stdin_path,
+ .cwd = dir_path,
+ .isolate = rootfs,
+ .work = dir_path,
+ .argv = request.argv,
+ }) };
// Под изоляцией программа видит каталог по другому пути.
@@
.time_ms = request.time_ms,
.mem_mb = request.mem_mb,
- .stdin_path = try std.fmt.allocPrint(arena, "{s}/stdin", .{dir_path}),
+ .stdin_path = stdin_path,
.cwd = dir_path,
.isolate = rootfs,
@@
return std.fmt.parseInt(u16, argv[1], 10) catch error.BadUsage;
}
+
+/// Как `zbox serve` держит нагрузку. Шаг 66: итеративно, по запросу за раз.
+/// Шаг 71: главный поток принимает, пул потоков обслуживает.
+pub const Model = enum { iterative, pool };
+
+pub const ServeArgs = struct {
+ port: u16,
+ model: Model = .iterative,
+ /// Сборка и запуск в песочнице из этого каталога корневой ФС.
+ isolate: ?[]const u8 = null,
+ /// Сколько задач выполняется одновременно: семафор на запуск.
+ jobs: u16 = 4,
+ /// Потоков пула. Их больше, чем `jobs`: разбор запроса, ответ 400 и
+ /// `GET /stats` не должны стоять в очереди к песочнице.
+ workers: u16 = 8,
+ /// Сколько принятых соединений ждёт свободного потока. Сверх этого 503.
+ queue: u16 = 16,
+};
+
+/// `serve <port> [--isolate ROOTFS] [--jobs N] [--workers N] [--queue N]`.
+/// Любой из флагов пула включает пул.
+pub fn parseServe(argv: []const []const u8) ArgsError!ServeArgs {
+ if (argv.len < 2) return error.BadUsage;
+ var result: ServeArgs = .{ .port = try parsePort(argv[0..2]) };
+ var rest = argv[2..];
+ while (rest.len > 0) : (rest = rest[2..]) {
+ if (rest.len < 2) return error.BadUsage;
+ const flag = rest[0];
+ const value = rest[1];
+ if (std.mem.eql(u8, flag, "--isolate")) {
+ if (value.len == 0) return error.BadUsage;
+ result.isolate = value;
+ } else {
+ const number = try positive(u16, value);
+ if (std.mem.eql(u8, flag, "--jobs")) {
+ result.jobs = number;
+ } else if (std.mem.eql(u8, flag, "--workers")) {
+ result.workers = number;
+ } else if (std.mem.eql(u8, flag, "--queue")) {
+ result.queue = number;
+ } else return error.BadUsage;
+ if (result.model == .iterative) result.model = .pool;
+ }
+ }
+ return result;
+}
+
+fn positive(comptime T: type, text: []const u8) ArgsError!T {
+ const number = std.fmt.parseInt(T, text, 10) catch return error.BadUsage;
+ return if (number == 0) error.BadUsage else number;
+}
main.zig получает две ветки: serve с флагами пула и подкоманду zbox job. Серверу нужен путь к самому себе, чтобы запускать zbox job, его даёт executablePathAlloc.
\\ zbox serve <port> --isolate ROOTFS
\\ то же, но сборка и запуск в песочнице из каталога ROOTFS (только Linux)
+ \\ zbox serve <port> --jobs N [--workers N] [--queue N]
+ \\ пул потоков: N задач одновременно, очередь соединений, сверх неё 503
+ \\ zbox job [--isolate ROOTFS]
+ \\ одна задача сервера: тело POST /run в stdin, ответ в stdout
\\
\\Флаги:
@@
if (args.len > 1 and std.mem.eql(u8, args[1], "serve")) {
- // `serve <port> --isolate ROOTFS`: обе стадии в песочнице.
- if (args.len == 5 and std.mem.eql(u8, args[3], "--isolate")) server_rootfs = args[4];
- const port = zbox.serve.parsePort(if (server_rootfs != null) args[1..3] else args[1..]) catch {
+ const serve_args = zbox.serve.parseServe(args[1..]) catch {
try out.writeAll(usage);
try out.flush();
std.process.exit(2);
};
+ server_rootfs = serve_args.isolate;
if (server_rootfs) |rootfs| {
// Корневую ФС проверяем при старте, а не на первом запросе:
@@
};
}
- return serve(init, port);
+ return switch (serve_args.model) {
+ .iterative => serve(init, serve_args.port),
+ .pool => servePool(init, serve_args),
+ };
}
+ if (args.len > 1 and std.mem.eql(u8, args[1], "job")) return job(init, args[2..], out);
+
const command = zbox.args.parse(args[1..]) catch {
try out.writeAll(usage);
@@
return .{ .content_type = "application/json", .body = body.written() };
}
+
+/// Шаг 71: пул потоков. Задачи идут дочерними `zbox job`, поэтому
+/// серверу нужен путь к самому себе.
+fn servePool(init: std.process.Init, args: zbox.serve.ServeArgs) !void {
+ const exe = try std.process.executablePathAlloc(init.io, init.arena.allocator());
+ var service: zbox.server.Service = .init(init.io, exe, args);
+ var err_buf: [256]u8 = undefined;
+ var stderr = std.Io.File.stderr().writerStreaming(init.io, &err_buf);
+ return zbox.server.servePool(&service, init.gpa, &stderr.interface);
+}
+
+/// `zbox job [--isolate ROOTFS]`: тело `POST /run` из stdin, ответ с полем
+/// `stage` в stdout. Сервер уже проверил запрос, но процесс мог позвать и
+/// человек, поэтому проверяем ещё раз.
+fn job(init: std.process.Init, args: []const []const u8, out: *std.Io.Writer) !void {
+ const arena = init.arena.allocator();
+ const rootfs: ?[]const u8 = if (args.len == 2 and std.mem.eql(u8, args[0], "--isolate")) args[1] else if (args.len == 0) null else {
+ try out.writeAll(usage);
+ try out.flush();
+ std.process.exit(2);
+ };
+ var in_buf: [4096]u8 = undefined;
+ var stdin = std.Io.File.stdin().readerStreaming(init.io, &in_buf);
+ const body = try stdin.interface.allocRemaining(arena, .limited(zbox.serve.max_body + 1));
+ const request = zbox.serve.parseRequest(arena, body) catch |err| {
+ try out.print("{{\"error\":\"{t}\"}}\n", .{err});
+ try out.flush();
+ std.process.exit(2);
+ };
+ // До fork в буфер stdout ничего не пишем, как и в `zbox run`.
+ const reply = try zbox.serve.runSourceIn(arena, init.io, request, rootfs);
+ try zbox.serve.writeJson(out, reply);
+ try out.flush();
+}
pub const serve = @import("box/serve.zig");
pub const watchdog = @import("box/watchdog.zig");
+
+// Шаг 71: сервер под нагрузкой.
+pub const queue = @import("pool/queue.zig");
+pub const pool = @import("pool/pool.zig");
+pub const http = @import("pool/http.zig");
+pub const job = @import("pool/job.zig");
+pub const server = @import("pool/server.zig");
Модулю zbox теперь нужен tiny (разбор HTTP и fdio), а рядом с тестами встаёт программа нагрузки.
/// Номера уроков курса, на которых проект вырос. Каждому шагу
/// соответствует файл `tests/step_NN.zig`, и все они зелёные на финале.
-const project_steps = [_]u8{ 48, 49, 54, 61, 64, 66, 67 };
+const project_steps = [_]u8{ 48, 49, 54, 61, 64, 66, 67, 71 };
pub fn build(b: *std.Build) void {
@@
// эталона our-tiny, подключён по пути в build.zig.zon.
const tiny = b.dependency("tiny", .{ .target = target, .optimize = optimize }).module("tiny");
+ // Пул потоков (шаг 71) разбирает HTTP частями TINY.
+ zbox.addImport("tiny", tiny);
const exe = b.addExecutable(.{
@@
options.addOptionPath("netcheck_exe", netcheck.getEmittedBin());
+ // Нагрузка на `zbox serve`: N клиентов разом (шаг 71). В тесты не входит.
+ const load = b.addExecutable(.{
+ .name = "load",
+ .root_module = b.createModule(.{
+ .root_source_file = b.path("bench/load.zig"),
+ .target = target,
+ .optimize = optimize,
+ .link_libc = true,
+ .imports = &.{.{ .name = "tiny", .module = tiny }},
+ }),
+ });
+ b.installArtifact(load);
+ const load_cmd = b.addRunArtifact(load);
+ if (b.args) |args| load_cmd.addArgs(args);
+ b.step("load", "Десять параллельных POST /run: zig build load -- PORT [N] [SECONDS]").dependOn(&load_cmd.step);
+
// zig build test прогоняет все шаги, zig build test -Dstep=48 только один.
const only_step = b.option(u8, "step", "Прогнать тесты одного шага, например -Dstep=48");
//! Нагрузка на `zbox serve`: N клиентов одновременно
//! шлют `POST /run` с программой на 200 мс, каждый в своём потоке.
//! Печатает задержку каждого запроса (медиана и худшая) и время пачки.
//!
//! zig build load -- PORT [N] [SECONDS]
const std = @import("std");
const tiny = @import("tiny");
const Client = struct {
port: u16,
body: []const u8,
latency_ms: i64 = 0,
status: u16 = 0,
fn fire(client: *Client, io: std.Io) void {
const started = std.Io.Clock.awake.now(io);
client.status = post(client.port, client.body) catch 0;
client.latency_ms = started.durationTo(std.Io.Clock.awake.now(io)).toMilliseconds();
}
};
/// Запрос целиком и ответ до закрытия; из ответа нужен только код.
fn post(port: u16, body: []const u8) !u16 {
const fd = try tiny.socket.openClientfd("127.0.0.1", port);
defer tiny.socket.close(fd);
var head: [128]u8 = undefined;
try tiny.fdio.writen(fd, try std.fmt.bufPrint(&head, "POST /run HTTP/1.0\r\nContent-Length: {d}\r\n\r\n", .{body.len}));
try tiny.fdio.writen(fd, body);
var reply: [64 * 1024]u8 = undefined;
const n = try tiny.fdio.readn(fd, &reply);
if (n < 12) return error.ShortReply;
return std.fmt.parseInt(u16, reply[9..12], 10);
}
pub fn main(init: std.process.Init) !void {
const arena = init.arena.allocator();
const args = try init.minimal.args.toSlice(arena);
if (args.len < 2) return error.Usage;
const port = try std.fmt.parseInt(u16, args[1], 10);
const n = if (args.len > 2) try std.fmt.parseInt(usize, args[2], 10) else 10;
const seconds = if (args.len > 3) args[3] else "0.2";
const body = try std.fmt.allocPrint(arena, "{{\"argv\":[\"sleep\",\"{s}\"]}}", .{seconds});
const clients = try arena.alloc(Client, n);
const threads = try arena.alloc(std.Thread, n);
const started = std.Io.Clock.awake.now(init.io);
for (clients, threads) |*client, *thread| {
client.* = .{ .port = port, .body = body };
thread.* = try std.Thread.spawn(.{}, Client.fire, .{ client, init.io });
}
for (threads) |t| t.join();
const total = started.durationTo(std.Io.Clock.awake.now(init.io)).toMilliseconds();
const latencies = try arena.alloc(i64, n);
var ok: usize = 0;
for (clients, latencies) |client, *latency| {
latency.* = client.latency_ms;
if (client.status == 200) ok += 1;
}
std.mem.sort(i64, latencies, {}, std.sort.asc(i64));
std.debug.print("{d} запросов, {d} ответов 200, пачка {d} мс, задержка медиана {d} мс, худшая {d} мс\n", .{ n, ok, total, latencies[n / 2], latencies[n - 1] });
}
Находка: сокеты утекают в задачи
Первая версия пула принимала соединения обычным accept из libc, и нагрузочный тест проходил. Но цифры были странные. Десять клиентов, каждый шлёт задачу sleep 0.2, мест четыре: задачи идут тремя волнами (4, 4, 2), и первые четыре клиента должны получить ответ примерно через 200 мс, а последние два через 600. Медиана задержки должна быть заметно меньше худшей. Вот что показал load на macOS (Apple M4 Max, 16 ядер, macOS 26.6.2, сборка ReleaseSafe, load average около 20):
$ zig build load -- 8172 10 0.2 # accept из libc
10 запросов, 10 ответов 200, пачка 655 мс, задержка медиана 647 мс, худшая 655 мс
10 запросов, 10 ответов 200, пачка 641 мс, задержка медиана 632 мс, худшая 641 мс
Медиана равна худшей: все десять клиентов получили ответ одновременно, вместе с последней задачей. Сервер при этом ответил каждому вовремя: transaction записала ответ и close закрыл дескриптор. Но клиент видит конец ответа только по FIN, а FIN ядро шлёт, когда закрыта последняя копия сокета во всей системе. Копии были. Дескриптор, который вернул accept, наследуется при fork и переживает execve (это правило из урока про разделение дескрипторов). Пока поток пула обслуживает клиента 1, соседний поток запускает zbox job для клиента 2, и этот zbox job, а за ним и программа задачи после своего execvp, уносят с собой копии всех открытых в эту секунду соединений. Соединение клиента 1 закроется по-настоящему только тогда, когда умрёт последняя задача, которая успела его унаследовать. На пачке из десяти это последняя задача пачки.
Программа задачи видит чужое соединение и сама. Попросим её посчитать свои сокеты:
$ curl -s -d '{"argv":["sh","-c","n=0; for f in /dev/fd/*; do [ -S \"$f\" ] && n=$((n+1)); done; echo $n"]}' localhost:8172/run
{..."stdout":"2\n",...} # accept из libc
{..."stdout":"1\n",...} # acceptCloexec
Один сокет остаётся и после исправления: в macOS getaddrinfo держит свой служебный сокет, и программа его видит при любом сервере. Поэтому тест шага сравнивает разность: сколько сокетов видит задача до и после того, как рядом открыто молчащее соединение.
Лечит это флаг close-on-exec. В Linux его ставит сам accept4(SOCK_CLOEXEC) атомарно, в момент создания дескриптора. В macOS accept4 нет, и приходится делать accept, а следом fcntl(F_SETFD, FD_CLOEXEC). Между двумя вызовами остаётся окно: если соседний поток сделает fork ровно в этот миг, копия всё равно утечёт. Окно крошечное, и в коде оно помечено; настоящее лечение без окна только в Linux. std.Io.net из урока 73 ставит флаг сам, поэтому там этой ошибки не будет. Слушающий сокет тоже получает флаг: ему в задаче делать нечего.
После исправления:
$ zig build load -- 8171 10 0.2
10 запросов, 10 ответов 200, пачка 652 мс, задержка медиана 433 мс, худшая 652 мс
10 запросов, 10 ответов 200, пачка 653 мс, задержка медиана 436 мс, худшая 653 мс
Пачка не изменилась, её держат три волны по 200 мс плюс запуск процессов. Изменилась медиана: теперь клиенты получают ответ, когда закончилась их задача, а не последняя. Мораль шире zbox: любой многопоточный сервер, который запускает дочерние процессы, обязан открывать всё с close-on-exec, и ошибка видна не в корректности, а только в задержках.
Прогон
$ zbox serve 8171 --jobs 4
zbox: listening on port 8171 (pool, jobs 4, queue 16)
$ zig build load -- 8171 10 0.2
10 запросов, 10 ответов 200, пачка 652 мс, задержка медиана 433 мс, худшая 652 мс
$ zig build load -- 8171 10 0.2
10 запросов, 10 ответов 200, пачка 653 мс, задержка медиана 436 мс, худшая 653 мс
$ curl -s localhost:8171/stats
{"running":0,"peak":4,"done":20,"rejected":0}
$ curl -s -d '{"argv":["sh","-c","sleep 0.2; echo готово"]}' localhost:8171/run
{"exit_code":0,"signal":null,"timed_out":false,"cpu_user_ms":2,"cpu_sys_ms":4,"max_rss_kb":2032,"wall_ms":222,"reason":"exited","vm_peak_kb":null,"vm_hwm_kb":null,"cgroup_peak_kb":null,"oom_kills":null,"stdout_truncated":false,"stderr_truncated":false,"stdout":"готово\n","stderr":"","net":"host","stage":"run"}
$ curl -s -d '{}' localhost:8171/run
{"error":"EmptySource"}
Двадцать задач, в пике ровно четыре одновременно, отказов нет: очередь на 16 соединений десять клиентов не переполнили. Итеративный zbox serve из урока 66 на той же пачке шёл бы десять задач подряд, больше двух секунд.
Тесты шага
Тестам нужен настоящий zbox serve дочерним процессом и клиенты в потоках, поэтому общие помощники из шага 66 переехали в tests/support.zig (в tests/step_66.zig их копии удалены, вызовы идут через support.) и подросли: Serve.start принимает флаги, volley стреляет пачкой клиентов, каждый в своём потоке и со своей ареной (арена не потокобезопасна), stats читает /stats, leakedSockets считает утёкшие сокеты разностью, как в прогоне выше.
const std = @import("std");
const build_options = @import("build_options");
+const tiny = @import("tiny");
+const zbox_mod = @import("zbox");
/// Поля ответа, которые проверяют тесты. Лишние поля JSON пропускаем:
@@
return .{ .reply = reply, .program_stdout = reply.stdout };
}
+
+/// Настоящий `zbox serve 0 [флаги]` дочерним процессом. Порт он выбирает
+/// сам и печатает первой строкой в stderr.
+pub const Serve = struct {
+ child: std.process.Child,
+ port: u16,
+
+ pub fn start(flags: []const []const u8) !Serve {
+ var argv_buf: [16][]const u8 = undefined;
+ const head = [_][]const u8{ build_options.zbox_exe, "serve", "0" };
+ @memcpy(argv_buf[0..head.len], &head);
+ @memcpy(argv_buf[head.len..][0..flags.len], flags);
+ var child = try std.process.spawn(std.testing.io, .{
+ .argv = argv_buf[0 .. head.len + flags.len],
+ .stdin = .ignore,
+ .stdout = .ignore,
+ .stderr = .pipe,
+ });
+ errdefer child.kill(std.testing.io);
+ var buf: [256]u8 = undefined;
+ var reader = child.stderr.?.readerStreaming(std.testing.io, &buf);
+ const line = try reader.interface.takeDelimiterExclusive('\n');
+ const prefix = "zbox: listening on port ";
+ if (!std.mem.startsWith(u8, line, prefix)) return error.NoPort;
+ const digits = std.mem.sliceTo(line[prefix.len..], ' ');
+ return .{ .child = child, .port = try std.fmt.parseInt(u16, digits, 10) };
+ }
+
+ pub fn stop(s: *Serve) void {
+ s.child.kill(std.testing.io);
+ }
+};
+
+pub const HttpReply = struct { status: u16, head: []const u8, body: []const u8 };
+
+/// Сырой запрос по сокету из модуля `tiny`, ответ читается до закрытия.
+/// Таймаут на сокете: сломанный сервер валит тест, а не вешает его.
+pub fn exchange(arena: std.mem.Allocator, port: u16, raw: []const u8) !HttpReply {
+ const fd = try tiny.socket.openClientfd("127.0.0.1", port);
+ defer tiny.socket.close(fd);
+ zbox_mod.http.setTimeouts(fd, 20_000);
+ try tiny.fdio.writen(fd, raw);
+ var response: std.ArrayList(u8) = .empty;
+ var chunk: [4096]u8 = undefined;
+ while (true) {
+ const n = try tiny.fdio.readn(fd, &chunk);
+ try response.appendSlice(arena, chunk[0..n]);
+ if (n < chunk.len) break;
+ }
+ const split = std.mem.indexOf(u8, response.items, "\r\n\r\n") orelse return error.NoHeaderEnd;
+ // "HTTP/1.0 200 OK": код стоит с 9-го по 12-й байт.
+ if (split < 12) return error.BadStatusLine;
+ return .{
+ .status = try std.fmt.parseInt(u16, response.items[9..12], 10),
+ .head = response.items[0..split],
+ .body = response.items[split + 4 ..],
+ };
+}
+
+pub fn postRun(arena: std.mem.Allocator, port: u16, json: []const u8) !HttpReply {
+ const raw = try std.fmt.allocPrint(arena, "POST /run HTTP/1.0\r\nContent-Type: application/json\r\nContent-Length: {d}\r\n\r\n{s}", .{ json.len, json });
+ return exchange(arena, port, raw);
+}
+
+/// Один клиент нагрузочного теста шагов 71 и 73.
+pub const Shot = struct {
+ port: u16,
+ index: usize,
+ /// У каждого клиента своя арена: арена не потокобезопасна.
+ arena: std.heap.ArenaAllocator,
+ status: u16 = 0,
+ stdout: []const u8 = "",
+ /// Сколько задача шла сама, по часам zbox.
+ wall_ms: u64 = 0,
+ head: []const u8 = "",
+ err: ?anyerror = null,
+
+ fn fire(shot: *Shot) void {
+ shot.post() catch |err| {
+ shot.err = err;
+ };
+ }
+
+ /// `sleep` на 200 мс и номер запроса в stdout: так видно, что каждое
+ /// соединение получило свой ответ, а не соседский.
+ fn post(shot: *Shot) !void {
+ const arena = shot.arena.allocator();
+ const script = try std.fmt.allocPrint(arena, "sleep 0.2; echo задача {d}", .{shot.index});
+ const json = try std.json.Stringify.valueAlloc(arena, .{ .argv = &[_][]const u8{ "sh", "-c", script } }, .{});
+ const reply = try postRun(arena, shot.port, json);
+ shot.status = reply.status;
+ shot.head = reply.head;
+ if (reply.status != 200) return;
+ const run = try parseReply(arena, reply.body);
+ shot.stdout = run.reply.stdout;
+ shot.wall_ms = run.reply.wall_ms;
+ }
+};
+
+/// Один запрос той же задачи до залпа: прогревает кэши (под Rosetta
+/// первый запуск бинарника это трансляция), чтобы залп мерил задачи.
+pub fn warmup(port: u16) !void {
+ var shot: Shot = .{ .port = port, .index = 0, .arena = .init(std.heap.page_allocator) };
+ defer shot.arena.deinit();
+ shot.fire();
+ if (shot.err) |err| return err;
+ if (shot.status != 200) return error.UnexpectedStatus;
+}
+
+/// Сколько шли бы задачи залпа одна за другой: сумма их собственных
+/// времён. Сравнивать с ней честнее, чем с константой: на нагруженной
+/// машине (соседние тесты собирают Zig) растут оба числа.
+pub fn serialMs(shots: []const Shot) u64 {
+ var total: u64 = 0;
+ for (shots) |shot| total += shot.wall_ms;
+ return total;
+}
+
+/// Клиенты одновременно, каждый в своём потоке. Возвращает время всей
+/// пачки в миллисекундах. Ответы живут до `freeVolley`.
+pub fn volley(port: u16, shots: []Shot) !i64 {
+ var threads: [16]std.Thread = undefined;
+ const started = std.Io.Clock.awake.now(std.testing.io);
+ for (shots, 0..) |*shot, i| {
+ shot.* = .{ .port = port, .index = i, .arena = .init(std.heap.page_allocator) };
+ threads[i] = try std.Thread.spawn(.{}, Shot.fire, .{shot});
+ }
+ for (threads[0..shots.len]) |t| t.join();
+ return started.durationTo(std.Io.Clock.awake.now(std.testing.io)).toMilliseconds();
+}
+
+pub fn freeVolley(shots: []Shot) void {
+ for (shots) |*shot| shot.arena.deinit();
+}
+
+pub const Stats = struct { peak: u32, done: u64, rejected: u64 };
+
+/// `GET /stats` сервера шага 71.
+pub fn stats(arena: std.mem.Allocator, port: u16) !Stats {
+ const reply = try exchange(arena, port, "GET /stats HTTP/1.0\r\n\r\n");
+ return std.json.parseFromSliceLeaky(Stats, arena, reply.body, .{ .ignore_unknown_fields = true });
+}
+
+/// Сколько сокетов видит программа задачи. Считает сама программа:
+/// оболочка обходит свои дескрипторы и проверяет тип (`-S`).
+fn socketsSeen(arena: std.mem.Allocator, port: u16) !usize {
+ const script = "n=0; for f in /dev/fd/*; do [ -S \"$f\" ] && n=$((n+1)); done; echo $n";
+ const json = try std.json.Stringify.valueAlloc(arena, .{ .argv = &[_][]const u8{ "sh", "-c", script } }, .{});
+ const reply = try postRun(arena, port, json);
+ if (reply.status != 200) return error.UnexpectedStatus;
+ const out = (try parseReply(arena, reply.body)).reply.stdout;
+ return std.fmt.parseInt(usize, std.mem.trimEnd(u8, out, "\n"), 10);
+}
+
+/// Утекают ли в задачу чужие соединения. Замер до и после того, как
+/// рядом открыто молчащее соединение: сервер его принял, и если сокет
+/// без close-on-exec, программа увидит на один сокет больше. Разность,
+/// а не абсолютное число: в macOS `getaddrinfo` оставляет свой
+/// служебный сокет, и он виден программе при любом сервере.
+pub fn leakedSockets(arena: std.mem.Allocator, port: u16) !usize {
+ const before = try socketsSeen(arena, port);
+ const idle = try tiny.socket.openClientfd("127.0.0.1", port);
+ defer tiny.socket.close(idle);
+ // Главный поток должен успеть принять молчащее соединение.
+ try std.testing.io.sleep(.fromMilliseconds(100), .awake);
+ const after = try socketsSeen(arena, port);
+ return after -| before;
+}
//! Шаг 71: очередь, пул потоков, семафор на запуски и `zbox serve --jobs`.
//! Нагрузочный тест: десять параллельных `POST /run` с программой на 200 мс.
const std = @import("std");
const builtin = @import("builtin");
const zbox = @import("zbox");
const support = @import("support.zig");
const testing = std.testing;
const io = testing.io;
const Queue = zbox.queue.Queue;
test "очередь: FIFO, tryPut на полной очереди, close отдаёт остаток" {
var buffer: [3]u32 = undefined;
var q: Queue(u32) = .init(io, &buffer);
try testing.expect(q.tryPut(1));
try testing.expect(q.put(2));
try testing.expect(q.tryPut(3));
try testing.expect(!q.tryPut(4));
try testing.expectEqual(1, q.get());
try testing.expect(q.tryPut(4));
q.close();
try testing.expect(!q.tryPut(5));
try testing.expect(!q.put(5));
try testing.expectEqual(2, q.get());
try testing.expectEqual(3, q.get());
try testing.expectEqual(4, q.get());
try testing.expectEqual(null, q.get());
}
const Tally = struct {
seen: [4000]std.atomic.Value(u8) = @splat(.init(0)),
};
fn produce(q: *Queue(u32), base: u32) void {
for (0..1000) |i| _ = q.put(base + @as(u32, @intCast(i)));
}
fn consume(q: *Queue(u32), tally: *Tally) void {
while (q.get()) |item| _ = tally.seen[item].fetchAdd(1, .monotonic);
}
test "очередь: четыре производителя, три потребителя, каждый элемент ровно один раз" {
var buffer: [8]u32 = undefined;
var q: Queue(u32) = .init(io, &buffer);
var tally: Tally = .{};
var consumers: [3]std.Thread = undefined;
for (&consumers) |*t| t.* = try std.Thread.spawn(.{}, consume, .{ &q, &tally });
var producers: [4]std.Thread = undefined;
for (&producers, 0..) |*t, i| t.* = try std.Thread.spawn(.{}, produce, .{ &q, @as(u32, @intCast(i)) * 1000 });
for (producers) |t| t.join();
q.close();
for (consumers) |t| t.join();
for (&tally.seen) |*count| try testing.expectEqual(1, count.load(.monotonic));
}
/// Поддельная задача: занять место, поспать 30 мс, освободить.
const Fake = struct {
slots: zbox.pool.Slots,
finished: std.atomic.Value(u32) = .init(0),
fn handle(fake: *Fake, _: u32) void {
fake.slots.acquire(io) catch return;
defer fake.slots.release(io);
io.sleep(.fromMilliseconds(30), .awake) catch {};
_ = fake.finished.fetchAdd(1, .monotonic);
}
};
test "пул: восемь потоков, но задач одновременно не больше трёх" {
const FakePool = zbox.pool.Pool(u32, *Fake, Fake.handle);
var fake: Fake = .{ .slots = .init(3) };
var queue_buffer: [16]u32 = undefined;
var threads: [8]std.Thread = undefined;
var pool: FakePool = undefined;
try pool.start(io, &fake, &queue_buffer, &threads);
for (0..12) |i| try testing.expect(pool.submit(@intCast(i)));
pool.stop();
try testing.expectEqual(12, fake.finished.load(.monotonic));
try testing.expectEqual(3, fake.slots.peak.load(.monotonic));
try testing.expectEqual(0, fake.slots.running.load(.monotonic));
try testing.expectEqual(12, fake.slots.done.load(.monotonic));
}
test "parseServe: модель по флагам" {
const serve = zbox.serve;
const plain = try serve.parseServe(&.{ "serve", "8081" });
try testing.expectEqual(.iterative, plain.model);
const pooled = try serve.parseServe(&.{ "serve", "0", "--jobs", "3", "--workers", "10", "--queue", "2" });
try testing.expectEqual(.pool, pooled.model);
try testing.expectEqual(3, pooled.jobs);
try testing.expectEqual(10, pooled.workers);
try testing.expectEqual(2, pooled.queue);
const isolated = try serve.parseServe(&.{ "serve", "0", "--isolate", "rootfs/root", "--jobs", "2" });
try testing.expectEqualStrings("rootfs/root", isolated.isolate.?);
try testing.expectEqual(.pool, isolated.model);
try testing.expectError(error.BadUsage, serve.parseServe(&.{ "serve", "0", "--jobs", "0" }));
try testing.expectError(error.BadUsage, serve.parseServe(&.{ "serve", "0", "--jobs" }));
try testing.expectError(error.BadUsage, serve.parseServe(&.{ "serve", "0", "--threads", "4" }));
}
test "parseRequest: argv вместо исходника" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const request = try zbox.serve.parseRequest(arena_state.allocator(), "{\"argv\":[\"sleep\",\"0.2\"],\"time_ms\":1000}");
try testing.expectEqual(2, request.argv.len);
try testing.expectEqualStrings("sleep", request.argv[0]);
try testing.expectEqualStrings("", request.source);
}
test "runSource: argv идёт одной стадией, stdin из запроса" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const reply = try zbox.serve.runSource(arena_state.allocator(), io, .{ .argv = &.{ "tr", "a-z", "A-Z" }, .stdin = "abc\n" });
try testing.expectEqual(.run, reply.stage);
try testing.expectEqual(0, reply.outcome.exit_code);
try testing.expectEqualStrings("ABC\n", reply.outcome.stdout);
}
test "zbox serve --jobs 3: десять параллельных запросов, свой ответ каждому, не больше трёх сразу" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const arena = arena_state.allocator();
var server = try support.Serve.start(&.{ "--jobs", "3", "--workers", "10" });
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;
}
}
test "zbox serve: переполненная очередь отвечает 503 с Retry-After" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const arena = arena_state.allocator();
var server = try support.Serve.start(&.{ "--jobs", "1", "--workers", "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 ok: usize = 0;
var busy: usize = 0;
for (&shots) |*shot| {
if (shot.err) |err| return err;
switch (shot.status) {
200 => ok += 1,
503 => {
busy += 1;
try testing.expect(std.mem.indexOf(u8, shot.head, "Retry-After: 1") != null);
},
else => return error.UnexpectedStatus,
}
}
// Одно место в пуле и одно в очереди на шесть клиентов: кому-то точно
// не хватило. Сколько именно дошло, зависит от того, успел ли поток
// пула забрать первое соединение из очереди до прихода второго.
try testing.expect(ok >= 1);
try testing.expect(busy >= 1);
try testing.expectEqual(6, ok + busy);
try testing.expectEqual(busy, (try support.stats(arena, server.port)).rejected);
}
test "zbox serve --jobs: 400 и 404 не ждут песочницу" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const arena = arena_state.allocator();
var server = try support.Serve.start(&.{ "--jobs", "1" });
defer server.stop();
const bad = try support.postRun(arena, server.port, "{}");
try testing.expectEqual(400, bad.status);
try testing.expectEqualStrings("{\"error\":\"EmptySource\"}\n", bad.body);
try testing.expectEqual(404, (try support.exchange(arena, server.port, "GET /run HTTP/1.0\r\n\r\n")).status);
const big = try std.fmt.allocPrint(arena, "POST /run HTTP/1.0\r\nContent-Length: {d}\r\n\r\n", .{zbox.serve.max_body + 1});
try testing.expectEqual(413, (try support.exchange(arena, server.port, big)).status);
}
test "linux: zbox serve --jobs --isolate, задача пула в песочнице" {
if (builtin.os.tag != .linux) return error.SkipZigTest;
const rootfs = std.mem.span(std.c.getenv("ZBOX_ROOTFS") orelse {
std.debug.print("пропуск: нет ZBOX_ROOTFS, собери корневую ФС через rootfs/build.sh\n", .{});
return error.SkipZigTest;
});
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const arena = arena_state.allocator();
var server = try support.Serve.start(&.{ "--jobs", "2", "--isolate", rootfs });
defer server.stop();
const reply = try support.postRun(arena, server.port, "{\"argv\":[\"sh\",\"-c\",\"echo $$\"]}");
try testing.expectEqual(200, reply.status);
const run = try support.parseReply(arena, reply.body);
if (run.reply.exit_code == zbox.run.limit_failed_code) {
std.debug.print("пропуск: {s}", .{run.reply.stderr});
return error.SkipZigTest;
}
// Оболочка это PID 1 своего пространства, сети нет.
try testing.expectEqualStrings("1\n", run.reply.stdout);
try testing.expectEqualStrings("none", run.reply.net);
}
test "сокеты соединений закрываются при exec и не утекают в задачи" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
var server = try support.Serve.start(&.{ "--jobs", "2" });
defer server.stop();
try testing.expectEqual(0, try support.leakedSockets(arena_state.allocator(), server.port));
}
Нагрузочный тест не сравнивает время с константой. На машине, где соседние тесты собирают Zig, растут и время пачки, и время каждой задачи, поэтому порог это сумма собственных wall_ms задач: параллельный сервер обязан уложиться меньше, чем они шли бы подряд. Нижняя граница в 750 мс проверяет обратное: при трёх местах десять задач по 200 мс не могут пройти быстрее четырёх волн, иначе семафор не работает. В тесте 503 одно место, один поток и одна ячейка очереди на шесть клиентов: сколько именно дойдёт, зависит от того, успел ли поток пула забрать первое соединение до прихода второго, поэтому проверяется только «хотя бы один прошёл, хотя бы одному отказали, всего шесть».
$ zig build test -Dstep=71 --summary all # macOS
+- run test 10 pass, 1 skip (11 total) 3s MaxRSS:4M
Пропущен тест с --isolate: пространства имён есть только в Linux. В контейнере runner-zig:dev с --privileged и корневой ФС из ZBOX_ROOTFS проходят все одиннадцать.
Шаг проекта: zt prof смотрит на блокировку
В уроке про потоки zt prof --threads научился раскладывать сэмплы по потокам и подписывать стек именем потока. Сегодня последний шаг профилировщика в разделе: снять профиль пула под нагрузкой и увидеть на flame graph, сколько стоит блокировка. Для этого нужны две вещи. Подопытный, у которого цена блокировки заведомо есть и которую можно убрать одним параметром. И умение профилировать сервер: он сам не выходит, его останавливают Ctrl-C, а дамп сэмплера до сих пор писался только из .fini_array, то есть на обычном exit.
pool.c: очередь под одним мьютексом
Подопытный написан на C с pthread, как и остальные фикстуры zt prof: так кадры функций предсказуемы, -O0 и указатель кадра у каждой, и имена pthread_mutex_lock приходят из настоящей libc.
/* Цель для zt prof --threads: пул воркеров над общей очередью задач.
*
* pool shared [воркеров] [тысяч задач] каждую задачу берут под мьютексом
* pool batched [воркеров] [тысяч задач] берут пачкой по 256 задач
*
* Очередь это счётчик следующей задачи под одним мьютексом, задача это
* короткий цикл. В режиме shared мьютекс захватывается на каждую задачу, и
* воркеры толкаются за него: в профиле видны pthread_mutex_lock и ожидание
* в ядре (futex). В режиме batched захватов в 256 раз меньше.
* Собирается как spin.c: -O0 и кадр у каждой функции. */
#define _GNU_SOURCE
#include <pthread.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
static pthread_mutex_t lock = PTHREAD_MUTEX_INITIALIZER;
static unsigned long next_task;
static unsigned long total_tasks;
static unsigned long batch = 1;
struct worker {
int number;
unsigned long done;
};
/* Берёт из очереди до batch задач подряд, возвращает первую, число в *count. */
__attribute__((noinline)) unsigned long take(unsigned long *count) {
pthread_mutex_lock(&lock);
unsigned long first = next_task;
unsigned long left = total_tasks - next_task;
*count = left < batch ? left : batch;
next_task += *count;
pthread_mutex_unlock(&lock);
return first;
}
/* Сама задача: пара сотен наносекунд счёта. */
__attribute__((noinline)) unsigned long work(unsigned long task) {
volatile unsigned long sink = task;
for (int i = 0; i < 100; i++) sink = sink * 31 + i;
return sink;
}
__attribute__((noinline)) void *worker(void *argument) {
struct worker *self = argument;
char name[16];
snprintf(name, sizeof name, "worker-%d", self->number);
pthread_setname_np(pthread_self(), name);
for (;;) {
unsigned long count;
unsigned long first = take(&count);
if (count == 0) break;
for (unsigned long task = first; task < first + count; task++) work(task);
self->done += count;
}
return NULL;
}
int main(int argc, char **argv) {
if (argc < 2 || (strcmp(argv[1], "shared") != 0 && strcmp(argv[1], "batched") != 0)) {
fprintf(stderr, "usage: pool shared|batched [workers] [thousands of tasks]\n");
return 2;
}
if (strcmp(argv[1], "batched") == 0) batch = 256;
int count = argc > 2 ? atoi(argv[2]) : 4;
total_tasks = (argc > 3 ? strtoul(argv[3], NULL, 10) : 2000) * 1000UL;
if (count < 1 || count > 64) count = 4;
pthread_t threads[64];
struct worker workers[64] = {0};
for (int i = 0; i < count; i++) {
workers[i].number = i + 1;
pthread_create(&threads[i], NULL, worker, &workers[i]);
}
unsigned long done = 0;
for (int i = 0; i < count; i++) {
pthread_join(threads[i], NULL);
done += workers[i].done;
}
printf("%s: %d workers, %lu tasks\n", argv[1], count, done);
return 0;
}
Очередь здесь это счётчик следующей задачи под мьютексом, задача это сотня умножений. В режиме shared воркер берёт по одной задаче, то есть захватывает мьютекс на каждые несколько сотен наносекунд работы. В режиме batched берёт по 256 штук, и захватов в 256 раз меньше. Работа одна и та же, отличается только то, как часто воркеры толкаются за мьютекс. Это psum-mutex против psum-local в миниатюре, только с настоящей очередью задач, как в пуле.
Дамп по Ctrl-C
const linux = std.os.linux;
+const posix = std.posix;
const prof = zt.prof;
@@
.process;
prof.sampler.start(hz, clock) catch {};
+
+ // Сервер сам не выходит, его останавливают через Ctrl-C или kill. До
+ // .fini_array такая смерть не доходит, поэтому дамп пишем и из
+ // обработчика. Сигнал, который программа или родитель уже настроили
+ // (например, игнор у фоновой задачи), не трогаем.
+ for ([_]posix.SIG{ .INT, .TERM }) |signal| {
+ var old: posix.Sigaction = undefined;
+ posix.sigaction(signal, null, &old);
+ if (old.handler.handler != posix.SIG.DFL) continue;
+ const action: posix.Sigaction = .{ .handler = .{ .handler = onStop }, .mask = posix.sigemptyset(), .flags = 0 };
+ posix.sigaction(signal, &action, null);
+ }
+}
+
+/// Дамп, потом тот же сигнал с действием по умолчанию: программа умирает
+/// так же, как умерла бы без профилировщика, и родитель видит ту же причину.
+fn onStop(signal: posix.SIG) callconv(.c) void {
+ writeDump();
+ const default: posix.Sigaction = .{ .handler = .{ .handler = posix.SIG.DFL }, .mask = posix.sigemptyset(), .flags = 0 };
+ posix.sigaction(signal, &default, null);
+ // Пока обработчик работает, этот сигнал заблокирован. Он придёт сразу
+ // после возврата и убьёт процесс.
+ _ = linux.kill(linux.getpid(), signal);
}
-// ponytail: программа, которая вышла через _exit или убита сигналом, дампа
-// не оставит: .fini_array до неё не дойдёт. Нужен профиль и таких программ,
-// пиши дамп из обработчика SIGTERM или отображай кольцо в файл через mmap.
fn fini() callconv(.c) void {
+ writeDump();
+}
+
+var dumped: std.atomic.Value(bool) = .init(false);
+
+// ponytail: программа, которая вышла через _exit или убита SIGKILL, дампа не
+// оставит. Нужен профиль и таких, отображай кольцо в файл через mmap.
+/// Пишет дамп один раз, из `.fini_array` или из обработчика сигнала. Всё
+/// здесь годится для обработчика: системные вызовы и готовые буферы.
+fn writeDump() void {
+ if (dumped.swap(true, .acq_rel)) return;
prof.sampler.stop();
const path = getenv(prof.output_variable) orelse return;
Обработчик пишет дамп и тут же возвращает сигналу действие по умолчанию, а потом шлёт его себе ещё раз. Пока обработчик работает, этот сигнал заблокирован, поэтому повторный придёт сразу после возврата и убьёт процесс ровно так, как убил бы без профилировщика: родитель увидит смерть от SIGINT, а не нормальный выход. Сигналы, которые кто-то уже настроил (фоновая задача оболочки получает SIGINT в игноре), не трогаем, а программа, поставившая свой обработчик после старта, перетрёт наш.
Из обработчика можно звать только то, что разрешают шесть правил урока про сигналы, и writeDump писался под них с самого начала: open, write и close напрямую через std.os.linux, формат в буфер на стеке, никакого аллокатора. Новое только одно: дамп теперь могут попросить дважды, из обработчика и из .fini_array (или из двух обработчиков в двух потоках). Флаг dumped взводится атомарным swap: кто первым увидел false, тот и пишет, остальные выходят.
Второй конец той же истории в самом zt. Ctrl-C в терминале получает вся группа процессов переднего плана: и программа, и zt, который её ждёт. Программа по сигналу сбросит дамп и умрёт, а zt должен дожить и прочесть дамп.
try out.flush();
+ // Ctrl-C в терминале получают и программа, и сам zt. Программа по нему
+ // сбросит дамп и умрёт, а zt должен дожить и прочесть дамп. Пустой
+ // обработчик, а не игнор: игнор унаследовала бы программа.
+ const keep_alive: std.posix.Sigaction = .{ .handler = .{ .handler = ignoreSignal }, .mask = std.posix.sigemptyset(), .flags = std.posix.SA.RESTART };
+ std.posix.sigaction(.INT, &keep_alive, null);
var child = std.process.spawn(init.io, .{ .argv = program, .environ_map = &environment }) catch |err| {
try out.print("zt prof: не удалось запустить {s} ({t})\n", .{ program[0], err });
@@
}
+fn ignoreSignal(_: std.posix.SIG) callconv(.c) void {}
+
/// Читает файлы из карты памяти с диска. Библиотека, которой уже нет на
/// месте, просто останется без имён функций.
Пустой обработчик, а не SIG_IGN, и это важно. Игнорирование сигнала наследуется через execve, и программа, запущенная из zt, тоже перестала бы реагировать на Ctrl-C. Обработчик же при execve сбрасывается в действие по умолчанию: адреса функции zt в новой программе нет. Это тоже правило из урока про сигналы, и здесь оно работает на нас. SA_RESTART просит ядро перезапускать системные вызовы zt, прерванные этим сигналом, а не возвращать из них EINTR.
/// Номера уроков курса, на которых проект вырос. Каждому шагу
/// соответствует файл `tests/step_NN.zig`, и все они зелёные на финале.
-const project_steps = [_]u8{ 6, 7, 9, 10, 11, 12, 13, 14, 17, 18, 42, 43, 44, 45, 46, 47, 49, 54, 69 };
+const project_steps = [_]u8{ 6, 7, 9, 10, 11, 12, 13, 14, 17, 18, 42, 43, 44, 45, 46, 47, 49, 54, 69, 71 };
pub fn build(b: *std.Build) void {
@@
b.installArtifact(sampler);
- // Цели для сквозных тестов `zt prof`: spin для шага 49, threads для
- // шага 69. Всегда Debug: это -O0, без встраивания и с кадром у каждой
- // функции.
- inline for (.{ "spin", "threads" }) |name| {
+ // Цели для сквозных тестов `zt prof`: spin для шага 49, threads и
+ // pool для шагов 69 и 71. Всегда Debug: это -O0, без встраивания и
+ // с кадром у каждой функции.
+ inline for (.{ "spin", "threads", "pool" }) |name| {
const fixture = b.addExecutable(.{
.name = name,
@@
paths.addOption([]const u8, "zt", b.getInstallPath(.bin, "zt"));
paths.addOption([]const u8, "hello_dyn", b.pathFromRoot("fixtures/dyn/hello_dyn"));
- inline for (.{ "spin", "threads" }) |name| {
+ inline for (.{ "spin", "threads", "pool" }) |name| {
paths.addOption([]const u8, name, b.getInstallPath(.{ .custom = "fixtures" }, name));
}
@@
const run_tests = b.addRunArtifact(step_tests);
// Шаги со сквозными тестами: они запускают собранный `zt`.
- if (number == 46 or number == 47 or number == 49 or number == 69) {
+ if (number == 46 or number == 47 or number == 49 or number == 69 or number == 71) {
step_tests.root_module.addOptions("paths", paths);
// Программа и библиотека-перехватчик должны быть собраны до теста.
Прогон: shared против batched
Контейнер runner-zig:dev (linux/arm64, OrbStack, ядро 7.0.14, 16 vCPU на Apple M4 Max, glibc 2.36), четыре воркера, миллион задач, 997 Гц:
zig build run -- prof --threads -o shared.folded zig-out/fixtures/pool shared 4 1000
zig build run -- prof --threads -o batched.folded zig-out/fixtures/pool batched 4 1000
Стеки одного воркера из shared:
worker-1;[libc.so.6];[libc.so.6];zt_thread_start;worker;take;[libc.so.6] 29
worker-1;[libc.so.6];[libc.so.6];zt_thread_start;worker;take;pthread_mutex_lock 4
worker-1;[libc.so.6];[libc.so.6];zt_thread_start;worker;take;pthread_mutex_lock;[libc.so.6] 20
worker-1;[libc.so.6];[libc.so.6];zt_thread_start;worker;work 56
И тот же воркер в batched:
worker-1;[libc.so.6];[libc.so.6];zt_thread_start;worker;work 16
По пяти прогонам (x86-64 снят в контейнере runner-zig:dev-amd64 под Rosetta; машину в это время грузили соседние процессы, load average от 25 до 57):
| сэмплов воркеров | доля в блокировке | доля вне work | |
|---|---|---|---|
aarch64, shared | 591 до 765 | 0,51 до 0,54 | 0,53 до 0,57 |
aarch64, batched | 122 до 188 | 0 до 0,01 | 0,01 до 0,03 |
x86-64, shared | 593 до 739 | 0,08 до 0,16 | 0,33 до 0,62 |
x86-64, batched | 246 до 292 | 0 | 0,03 до 0,09 |
Одна и та же работа в shared стоит в 2,5 (x86-64) до 4 (aarch64) раз больше процессора, и на aarch64 больше половины этого процессора уходит на мьютекс. Flame graph рисуется тем же flamegraph.pl, что в уроке про профилирование: формат свёрнутых стеков у zt prof такой же, как у stackcollapse-perf.pl. На картинке shared на aarch64 у каждого воркера два плато почти одной ширины, take и work; у batched остаётся одно.
Как читать эти стеки, чтобы не обмануться:
pthread_mutex_lockназван по имени. Адрес внутри libc переведён в пару «файл, смещение» картой памяти (zt pmap, шаг 54) и найден в.dynsymlibc.- Безымянный
[libc.so.6]над ним почти наверняка__lll_lock_wait, где поток уходит в системный вызовfutex. Это внутренняя функция glibc, в.dynsymеё нет, и имени взять неоткуда. - Какие кадры выпадают. Обход идёт по указателю кадра, а функция без своего кадра теряет того, кто её позвал. На x86-64 libc собрана без указателя кадра, поэтому выпадает
take, иpthread_mutex_lockвисит прямо наworker. На aarch64 чаще выпадает самpthread_mutex_lock(строкаtake;[libc.so.6]). Поэтому тест шага считает блокировкой любой из трёх кадров:take,pthread_mutex_lockилиpthread_mutex_unlock. - Профиль видит процессор, а не ожидание. Поток, который спит в
futex, процессор не жжёт и сэмплов не получает. В профиле видна цена блокировки: захват, освобождение, системные вызовы, прыжки линии с мьютексом между ядрами. Сколько воркер простоял в очереди за мьютексом, показывает только профиль ожидания, и это другой инструмент. - На нагруженной машине толкотни меньше. Воркеры реже работают одновременно и реже сталкиваются на мьютексе. Замер блокировки на занятой машине занижает её цену.
Прогон: TINY с пулом под нагрузкой
Теперь то, ради чего всё затевалось: настоящий сервер. TINY с пулом из четырёх воркеров, сборка ReleaseSafe (указатели кадров на месте, проверки тоже), под zt prof, в соседнем терминале loadgen на 16 соединений по 3000 запросов, потом Ctrl-C:
(cd ../our-tiny && zig build -Doptimize=ReleaseSafe)
zig build
zig-out/bin/zt prof --threads -o tiny.folded \
../our-tiny/zig-out/bin/tiny tiny --pool 4 8080 ../our-tiny/zig-out
# в соседнем терминале
../our-tiny/zig-out/bin/tiny loadgen 127.0.0.1 8080 16 3000 /home.html
# Ctrl-C в первом
FlameGraph/flamegraph.pl tiny.folded > tiny.svg
Контейнер runner-zig:dev (linux/arm64, 16 vCPU, ядро 7.0.14 на Apple M4 Max, load average около 5), два прогона:
requests 48000, errors 0, bytes 10704000
time 0.851 s, 56416 req/s, p50 248.1 us, p99 1093.1 us
zt prof: сэмплов 1918, частота 997 Гц, затёрто 0
Корни свёрнутых стеков это имена потоков, и нагрузка между воркерами ровная:
tiny 113 главный поток: accept и insert в sbuf
worker-0 455
worker-1 467
worker-2 446
worker-3 437
Сводка по кадрам: доля сэмплов, в стеке которых есть кадр (первый прогон, во втором 1767 сэмплов, доли те же с точностью до нескольких процентов):
http.tiny.Server.handle 89 % вся транзакция в воркерах
http.tiny.Server.serveStatic 50 %
memset 36 % из них 25 % прямо в Server.handle
write 36 % ответ клиенту и журнал
http.tiny.Server.access 24 % строка журнала доступа в stderr
close 7 %
accept 5 % главный поток
sbuf, Semaphore, futex 0,4 % 8 сэмплов из 1918
Блокировки в профиле TINY почти нет: восемь сэмплов на весь сервер, четыре в главном потоке (insert и post) и четыре в воркерах (ожидание в remove и пробуждение соседа). Транзакция воркера это сотни микросекунд системных вызовов, а за sbuf он берётся один раз на соединение. Это batched, а не shared. Главный поток тоже не узкое место: он занят меньше чем на шесть процентов сэмплов, пока четыре воркера заняты почти целиком. Деньги уходят в другое:
write. Ответ клиенту и строка журнала доступа, которую TINY пишет в stderr на каждый запрос: четверть всего процессора сервера уходит на то, чтобы сообщить о запросе, а не на сам запрос. Отключи журнал, и сервер станет заметно быстрее.memsetна ровном месте. ВServer.handleи глубже лежат буферыvar in_buf: [8192]u8 = undefinedи соседи. В сборках с проверками,DebugиReleaseSafe, Zig заполняетundefinedбайтами0xaa, чтобы чтение неинициализированной памяти было видно сразу. Проверь сам: функция с таким буфером в-O ReleaseSafeсодержит вызовmemset, в-O ReleaseFastнет. Больше двадцати килобайт на запрос, десятки тысяч запросов в секунду, и около трети процессора сервера занято заполнением буферов, в которые тут же пишут настоящие данные. ВReleaseFastэтого нет, но и проверок нет; выбор режима сборки для сервера это решение, а не мелочь.
Имена воркеров появились потому, что пул зовёт setName сразу после spawn. Без этого все стеки начинались бы с tiny, и flame graph слил бы воркеров в одну гору.
Мораль шага: пул в TINY стоит почти ничего, и оптимизировать синхронизацию здесь бессмысленно, профиль показывает на журнал и режим сборки. Цену мьютекса видно только там, где работа на один захват сравнима с самим захватом, как в pool shared и в psum-mutex. Сначала профиль, потом оптимизация, как в уроке 36.
Тесты шага
//! Шаг 71: профиль пула потоков, где видна цена блокировки.
//!
//! Фикстура `fixtures/prof/pool.c` раздаёт задачи четырём воркерам из общей
//! очереди под одним мьютексом. В режиме `shared` задача берётся по одной,
//! в режиме `batched` пачкой по 256. Работа одна и та же, разница только в
//! том, как часто воркеры толкаются за мьютекс. Тест снимает оба профиля
//! через `zt prof --threads` и сравнивает две доли сэмплов воркеров: в
//! блокировке (кадр `take`, `pthread_mutex_lock` или `pthread_mutex_unlock`)
//! и вообще вне работы (нет кадра `work`).
//!
//! Кадров смотрим несколько: на x86-64 libc собрана без указателя кадра, и
//! `take`, которая позвала мьютекс, из стека выпадает, а на aarch64 чаще
//! выпадает сам `pthread_mutex_lock`. Только Linux: профилировщик входит
//! через `LD_PRELOAD`.
const std = @import("std");
const builtin = @import("builtin");
const paths = @import("paths");
const Profile = struct {
/// Сэмплы потоков `worker-N`.
workers: u64 = 0,
/// Из них в функции `work` или глубже.
working: u64 = 0,
/// Из них в блокировке.
locking: u64 = 0,
/// Из них с `pthread_mutex_lock` в стеке: имя из `.dynsym` libc.
in_mutex_lock: u64 = 0,
fn share(profile: Profile, part: u64) f64 {
return @as(f64, @floatFromInt(part)) / @as(f64, @floatFromInt(profile.workers));
}
fn overhead(profile: Profile) f64 {
return profile.share(profile.workers - profile.working);
}
fn lockShare(profile: Profile) f64 {
return profile.share(profile.locking);
}
};
fn profilePool(mode: []const u8, dump: *std.ArrayList(u8)) !Profile {
const gpa = std.testing.allocator;
const result = try std.process.run(gpa, std.testing.io, .{
.argv = &.{ paths.zt, "prof", "--threads", paths.pool, mode, "4", "1000" },
});
defer gpa.free(result.stdout);
defer gpa.free(result.stderr);
try std.testing.expectEqual(std.process.Child.Term{ .exited = 0 }, result.term);
try dump.appendSlice(gpa, result.stdout);
var profile: Profile = .{};
var lines = std.mem.tokenizeScalar(u8, result.stdout, '\n');
while (lines.next()) |line| {
if (!std.mem.startsWith(u8, line, "worker-")) continue;
const space = std.mem.lastIndexOfScalar(u8, line, ' ') orelse continue;
const count = try std.fmt.parseInt(u64, line[space + 1 ..], 10);
const stack = line[0..space];
profile.workers += count;
if (hasFrame(stack, "work")) profile.working += count;
if (hasFrame(stack, "pthread_mutex_lock")) profile.in_mutex_lock += count;
if (hasFrame(stack, "take") or hasFrame(stack, "pthread_mutex_lock") or hasFrame(stack, "pthread_mutex_unlock")) {
profile.locking += count;
}
}
return profile;
}
/// Есть ли в свёрнутом стеке кадр ровно с таким именем.
fn hasFrame(stack: []const u8, name: []const u8) bool {
var frames = std.mem.splitScalar(u8, stack, ';');
while (frames.next()) |frame| {
if (std.mem.eql(u8, frame, name)) return true;
}
return false;
}
test "общая очередь под мьютексом: заметная доля сэмплов уходит на блокировку" {
if (builtin.os.tag != .linux) return error.SkipZigTest;
const gpa = std.testing.allocator;
var output: std.ArrayList(u8) = .empty;
defer output.deinit(gpa);
const shared = try profilePool("shared", &output);
const batched = try profilePool("batched", &output);
// Замеры в README: доля в блокировке у shared от 0,08 (x86-64 под
// Rosetta) до 0,54 (aarch64), у batched ноль или один сэмпл. На
// загруженной машине воркеры реже работают одновременно, и толкотни
// меньше. Пороги с запасом на это.
const ok = shared.workers >= 50 and batched.workers >= 50 and
shared.locking >= 5 and shared.lockShare() > 4 * batched.lockShare() and
shared.overhead() > batched.overhead() and
// Символизация дошла до libc: функция блокировки названа по имени.
shared.in_mutex_lock > 0;
if (!ok) {
std.debug.print("shared {any}\nbatched {any}\n{s}", .{ shared, batched, output.items });
return error.LockCostNotVisible;
}
}
test "программа, остановленная Ctrl-C, оставляет профиль" {
if (builtin.os.tag != .linux) return error.SkipZigTest;
const gpa = std.testing.allocator;
const io = std.testing.io;
var path_buffer: [64]u8 = undefined;
const path = try std.fmt.bufPrint(&path_buffer, "/tmp/zt-step71-{d}.folded", .{std.os.linux.getpid()});
defer std.Io.Dir.cwd().deleteFile(io, path) catch {};
// Сервер сам не выходит, его останавливают. Своя группа процессов для zt
// и программы: сигнал всей группе это то, что терминал делает по Ctrl-C.
// Без сигнала потоки крутились бы секунд двадцать, тест упал бы, но не завис.
var child = try std.process.spawn(io, .{
.argv = &.{ paths.zt, "prof", "--threads", "-o", path, paths.threads, "5000" },
.pgid = 0,
.stderr = .ignore,
});
try io.sleep(.fromMilliseconds(500), .awake);
_ = std.os.linux.kill(-child.id.?, .INT);
const term = try child.wait(io);
// Программа умерла от SIGINT, zt дожил, прочёл дамп и вернул 1.
try std.testing.expectEqual(std.process.Child.Term{ .exited = 1 }, term);
const folded = try std.Io.Dir.cwd().readFileAlloc(io, path, gpa, .limited(1 << 20));
defer gpa.free(folded);
// Потоки ещё жили, их имена взяты из /proc/self/task/<tid>/comm.
for ([_][]const u8{ "heavy;", "medium;", "light;" }) |root| {
try std.testing.expect(std.mem.indexOf(u8, folded, root) != null);
}
}
Первый тест сравнивает не имя функции, а две доли, и с большим запасом: на x86-64 и на aarch64 из стека выпадают разные кадры, а на загруженной машине толкотни меньше. Второй проверяет всю цепочку Ctrl-C: своя группа процессов (pgid = 0) для zt и программы, сигнал всей группе, как делает терминал, zt выходит с кодом 1 (программа умерла от сигнала), а в файле свёрнутых стеков есть корни всех трёх живых потоков. Их имена взяты из /proc/self/task/<tid>/comm в момент дампа, потому что до join дело не дошло.
Оба теста только для Linux: профилировщик входит через LD_PRELOAD. На macOS шаг даёт два пропуска, в контейнере runner-zig:dev оба зелёные, и все тесты прошлых шагов тоже (251 из 254, три пропуска).
На macOS
Потоки, мьютексы и семафоры std.Io, атомики и сокеты работают на Mac напрямую, поэтому почти весь урок проходится без контейнера:
our-tiny:zig build test -Dstep=71даёт пять зелёных,tiny --pool,echoservert-pre,loadgenиpsum-benchзапускаются как есть. Замерыpsumв уроке сняты на macOS.zbox serve --jobs: десять зелёных и один пропуск (тест с--isolate, пространства имён есть только в Linux). Утечку сокетов черезexecты увидишь и на Mac, а исправление там с окном:accept4в macOS нет, флаг ставится вторым вызовомfcntl.- Код-задача урока чистая и идёт в песочнице курса как обычно.
Что только в Linux, то есть в контейнере ghcr.io/bondiano/runner-zig:dev (linux/arm64, на Apple Silicon работает нативно):
zt profцеликом: он входит в программу черезLD_PRELOADи.init_arrayELF, а в macOS другой формат и другой загрузчик.gettid,/proc/self/task/<tid>/commи--thread-timersтоже только Linux.accept4(SOCK_CLOEXEC)без окна и тестzbox serve --jobs --isolate(нужен--privilegedи корневая ФС, рецепт в прошлом уроке).- Под эмуляцией x86-64 (
runner-zig:dev-amd64, Rosetta)zt profработает, и таблица выше снята и там, но запуск процесса в ней стоит около 110 мс вместо единиц миллисекунд, так что нагрузочные замерыzboxпод Rosetta про эмулятор, а не про пул.
Одно отличие касается чисел. Потоки на Mac нельзя привязать к ядрам, и планировщик сам решает, какой поток поедет на энергоэффективное ядро; отсюда падение эффективности psum-local на шестнадцати потоках. В Linux для чистого замера есть taskset и pthread_setaffinity_np, в контейнере OrbStack ядра виртуальные, и привязка там тоже не гарантирует, какое физическое ядро достанется.
Практика
Перепиши psum-mutex в psum-local. В стартере функция psumLocal(io, n, partials) уже даёт верный ответ, но каждый поток складывает в общую сумму под мьютексом на каждой итерации, и четыре потока выходят медленнее одного. Нужно, чтобы поток t складывал свою полосу band(n, partials.len, t) в переменную на своём стеке и один раз, в конце, записал итог в partials[t]. Главный поток дожидается всех через join и складывает слоты. Функция band готова; это та же нарезка полос, что в задаче про умножение матриц из урока 69.
Две детали, о которых урок говорил выше. Аргументы spawn копируются в кортеж, поэтому полосу и указатель на свой слот можно передать прямо туда, без общей структуры. И если очередной spawn вернул ошибку, уже запущенные потоки надо дождаться до возврата, иначе они запишут в partials из кадра, которого больше нет: errdefer с join по уже запущенным.
Тесты сверяют сумму с формулой на числе потоков от 1 до 8, проверяют, что в каждом слоте лежит ровно сумма его полосы (слоты заранее забиты maxInt), краевые случаи (n = 0, потоков больше, чем чисел) и n = 2²⁴ на восьми потоках. Последний тест замеряет: лучший из трёх прогонов твоей версии на 2¹⁸ чисел и четырёх потоках должен быть хотя бы втрое быстрее psum-mutex. На M4 Max в отладочной сборке разница выходит около тридцати раз (6100 мкс против 192), в песочнице курса с половиной ядра от семи раз.
Упражнения
Итоги
- Поток на соединение платит за рождение и смерть потока на каждом запросе и не ограничивает их число. Пул создаёт потоки один раз, главный поток кладёт соединения в ограниченный
sbuf, воркеры забирают. Полная очередь останавливаетaccept(или, как вzbox, отвечает 503): это обратное давление. - Соединение передаётся воркеру по значению, воркеры останавливаются отравленной пилюлей, по одной на воркер, после всей настоящей работы.
pthread_onceв Zig пишется двойной проверкой: быстрый путь на атомике сacquire, медленный под мьютексом, флаг сreleaseтолько после работы.- Ускорение
S = T₁ / Tₚ, эффективностьE = S / p. Абсолютное ускорение считают от лучшей последовательной программы, относительное от той же параллельной на одном ядре, и относительное легко приукрасить. Сильное масштабирование держит постоянной задачу, слабое работу на ядро. - Синхронизация на каждой итерации это отказ от параллелизма:
psum-mutexиpsum-atomicна любом числе потоков медленнее одного. Общая линия кэша переезжает между ядрами на каждой записи. - Ложное разделение делает независимые ячейки массива одной переменной для протокола когерентности; лечится выравниванием на
std.atomic.cache_line(128 байт на Apple M, 64 на x86-64). Лучше всего локальный аккумулятор и одна запись в конце. - Закон Амдала ограничивает ускорение единицей, делённой на последовательную долю; метрика Карпа и Флэтт по замеру показывает, растут ли потери вместе с числом потоков.
zbox serve --jobs N: очередь сtryPutи 503, потоков пула больше, чем мест, семафорSlotsдержит не большеNзадач, задача в своём процессеzbox job, потому что сторож на сигналах и будильнике один на процесс.- Многопоточный сервер, который запускает процессы, обязан открывать сокеты с close-on-exec: иначе соседний
forkуносит копию, и клиент ждёт FIN до конца чужой задачи. Ошибка видна только в задержках. zt profвидит процессорное время: цену блокировки, но не ожидание. В пуле с мелкими задачами блокировка съедает больше половины процессора воркеров, в TINY меньше процента, а деньги уходят в журнал и вmemsetбуферовundefined, которые вReleaseSafeзаполняются0xaa.
Дальше
У нас три сервера на одном пуле и честные числа о том, где параллелизм работает, а где синхронизация его съедает. Но все наши потоки до сих пор трогали либо свои данные, либо данные под мьютексом, и одна история уже намекнула, что так бывает не всегда: run из zbox нельзя звать из двух потоков, потому что он держит состояние в глобальных переменных. Следующий урок про то, как такие функции распознать и починить: четыре класса потоконебезопасных функций, gethostbyname и блокировка с копированием, реентерабельность, гонки, которые ловит ThreadSanitizer, и взаимоблокировка двух мьютексов, захваченных в разном порядке.
домашка