Раздел 32 · Системное программирование: Zig, ассемблер, Verilog
Потоки
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Потоки
В прошлом уроке эхо-сервер научился обслуживать многих клиентов двумя способами. Процесс на соединение пишется просто, но у каждого клиента своя память, и
forkстоит дорого. Событийный сервер делит всё, но каждый клиент в нём превращается в конечный автомат, и одна блокирующая запись встаёт поперёк всех. Потоки берут лучшее от обоих: код обслуживания пишется прямой линией, как у процесса, а память общая, как у событийного сервера. Сегодня разберём, что такое поток с точки зрения ядра, чемstd.Threadотличается отpthread_createи отclone, как не отдать потоку чужую переменную, соберём эхо-сервер и TINY на потоках, аzt profнаучим различать потоки и подписывать каждый стек именем того, кто жёг процессор.
Цели урока
- Назвать, что у потока своё (стек, регистры, номер, маска сигналов), а что общее со всеми потоками процесса (код, данные, куча, дескрипторы, обработчики сигналов), и чем это отличается от процесса после
fork. - Запускать потоки через
std.Thread.spawn, дожидаться их черезjoinили отпускать черезdetach, и знать, чтоreturnизmainубивает всех. - Увидеть, во что превращается
spawn:pthread_create,mmapпод стек иcloneс набором флагов. - Найти и исправить гонку на аргументе потока: адрес переменной цикла против её копии.
- Построить эхо-сервер и TINY с потоком на соединение и объяснить, почему медленный клиент больше не держит остальных.
- Научить
zt profраскладывать сэмплы по потокам:gettidв обработчикеSIGPROF, таблица границ стеков, подменаpthread_create, имя потока в корне свёрнутого стека.
Идея: второй поток управления в том же адресном пространстве
Процесс в уроке про fork был парой: поток управления и адресное пространство. fork копирует обе половины сразу, и ребёнок живёт в своей копии памяти. Поток разрезает эту пару. Новый поток получает свой поток управления, то есть свои регистры, свой счётчик команд и свой стек, а адресное пространство остаётся одно на всех.
Из этого следует почти всё остальное в уроке. Если один поток записал число в глобальную переменную, другой прочитает его без всяких пайпов и сокетов: это одна и та же ячейка памяти. Если один поток открыл файл, дескриптор виден всем: таблица дескрипторов тоже одна. Если один поток упал с SIGSEGV, умирает процесс целиком со всеми потоками. И если один поток прочитал переменную, пока другой её пишет, результат зависит от того, кто успел первым.
Вот что у потока своё и что общее, по разделу 12.4 книги и man 7 pthreads:
| Своё у каждого потока | Общее у всех потоков процесса |
|---|---|
номер потока (tid в ядре, pthread_t в libc) | номер процесса pid и родитель |
| регистры, счётчик команд, флаги | код, глобальные и статические данные |
| стек (а значит, локальные переменные) | куча и всё, что выделил malloc |
маска сигналов, errno | обработчики сигналов (диспозиции) |
| указатель на локальную память потока | таблица дескрипторов, текущий каталог, umask |
Стек своё у каждого только по договорённости. Ядро его никак не защищает: стек потока это кусок общего адресного пространства, и поток, у которого есть указатель на чужую локальную переменную, прочитает и перепишет её. На этом построена главная ловушка урока, до неё дойдём через раздел.
Среди потоков нет родителя и детей. Поток, с которого начался процесс, называют главным, и отличается он только тем, что был первым. Любой поток может дождаться любого другого, и любой может завершить процесс целиком. Планирует потоки ядро, каждый по отдельности: на Linux поток это такая же задача планировщика, как процесс, только с общей памятью. Многозадачность вытесняющая: таймер может отнять процессор у потока между любыми двумя инструкциями, это мы видели в уроке про исключения. Переключение между потоками одного процесса дешевле, чем между процессами: таблица страниц та же, и кэш трансляций из урока про TLB сбрасывать не нужно.
Если тема потоков для тебя новая, общую картину с точки зрения планировщика дают уроки про устройство машины, а многопоточность в Node (worker_threads, SharedArrayBuffer) разобрана в разделе про асинхронность. Здесь мы спускаемся на уровень ниже: какими системными вызовами это сделано и что из этого следует для кода на Zig.
Первая программа: чей стек, чей pid
Проверим таблицу руками. Программа заводит два потока, и каждый поток, включая главный, печатает номер процесса, свой номер потока, адрес своей локальной переменной, адрес глобальной переменной и дескриптор, который главный поток открыл до их рождения:
//! Что у потоков своё, а что общее: номер процесса, номер потока в ядре,
//! адрес локальной переменной (стек) и глобальной (данные), дескриптор.
const std = @import("std");
const c = std.c;
var global: u32 = 0;
fn report(who: []const u8, fd: c.fd_t) void {
var local: u32 = 0;
std.debug.print("{s:<7} pid {d} tid {d} &local 0x{x} &global 0x{x} fd {d}\n", .{
who,
c.getpid(),
std.Thread.getCurrentId(),
@intFromPtr(&local),
@intFromPtr(&global),
fd,
});
}
pub fn main() !void {
const fd = c.dup(1);
report("главный", fd);
const t1 = try std.Thread.spawn(.{}, report, .{ "T1", fd });
const t2 = try std.Thread.spawn(.{}, report, .{ "T2", fd });
t1.join();
t2.join();
}
std.Thread.spawn(.{}, report, .{ "T1", fd }) заводит поток, который выполнит report("T1", fd). Первый аргумент это настройки (размер стека), второй функция, третий кортеж её аргументов. join ждёт, пока функция потока вернётся, и освобождает всё, что поток занимал. std.Thread.getCurrentId() это номер потока, который знает ядро: на Linux gettid, на macOS pthread_threadid_np.
$ zig run ids.zig
главный pid 46462 tid 82874450 &local 0x16b5cda9c &global 0x1049dfb44 fd 3
T1 pid 46462 tid 82874460 &local 0x16c5d6dcc &global 0x1049dfb44 fd 3
T2 pid 46462 tid 82874461 &local 0x16d5e2dcc &global 0x1049dfb44 fd 3
А это Linux в контейнере (Debian 12, arm64), с -lc:
$ zig build-exe ids.zig -lc && ./ids
главный pid 143 tid 143 &local 0xfffff165b354 &global 0x11c32d0 fd 3
T1 pid 143 tid 144 &local 0xffff8d30c724 &global 0x11c32d0 fd 3
T2 pid 143 tid 145 &local 0xffff8c2fc724 &global 0x11c32d0 fd 3
Читаем по столбцам. pid у всех один: это один процесс. Номера потоков разные, а на Linux номер главного потока совпадает с pid: процесс в ядре Linux это группа задач, и номер группы это номер её первой задачи. Глобальная переменная лежит по одному адресу у всех, это одна ячейка. Дескриптор 3, который открыл главный, виден потокам без всякой передачи. А локальные переменные разбросаны: у каждого потока свой стек, и стеки лежат далеко друг от друга. Разница адресов T1 и T2 на Linux это 0x1010000, 16 МиБ стека плюс 64 КиБ, на macOS 0x100c000, 16 МиБ плюс 48 КиБ. Откуда 16 МиБ, увидим в следующем разделе.
Что делает spawn
std.Thread это тонкая обёртка, прочитай её исходник целиком: около тысячи семисот строк на все системы. Какую реализацию выбрать, решается при компиляции:
pub const use_pthreads = native_os != .windows and native_os != .wasi and builtin.link_libc;
const Impl = if (native_os == .windows)
WindowsThreadImpl
else if (use_pthreads)
PosixThreadImpl
else if (native_os == .linux)
LinuxThreadImpl
// ...
Программа, которая линкуется с libc (все наши эталоны, и любая программа на macOS), получает PosixThreadImpl, то есть обычный pthread_create. Программа без libc на Linux получает LinuxThreadImpl: Zig сам выделяет стек через mmap и сам зовёт clone. Посмотрим сначала на первую ветку, в ней самое интересное для урока:
fn spawn(config: SpawnConfig, comptime f: anytype, args: anytype) !Impl {
const Args = @TypeOf(args);
const allocator = std.heap.c_allocator;
const Instance = struct {
fn entryFn(raw_arg: ?*anyopaque) callconv(.c) ?*anyopaque {
const args_ptr: *Args = @ptrCast(@alignCast(raw_arg));
defer allocator.destroy(args_ptr);
return callFn(f, args_ptr.*);
}
};
const args_ptr = try allocator.create(Args);
args_ptr.* = args;
// ... pthread_attr: размер стека и защитная страница ...
switch (c.pthread_create(&handle, &attr, Instance.entryFn, @ptrCast(args_ptr))) {
pthread_create умеет передать потоку ровно один аргумент, нетипизированный указатель. Книга поэтому складывает аргументы в структуру и передаёт указатель на неё. Zig делает то же самое за тебя: кортеж аргументов копируется в кучу (allocator.create(Args), args_ptr.* = args), новый поток начинает со сгенерированной entryFn, распаковывает кортеж, освобождает память и зовёт твою функцию. Запомни эту строку args_ptr.* = args: всё, что ты передал в кортеже, поток получает копией, и копия живёт в куче, а не на стеке того, кто создавал. Это прямо относится к гонке, которую мы разберём ниже.
Стек задаёт SpawnConfig.stack_size, по умолчанию 16 * 1024 * 1024. Это больше, чем стек главного потока (ulimit -s на Linux обычно 8 МиБ), и больше, чем стек по умолчанию у самого pthread_create. Бояться этих мегабайт не надо: память под стек только зарезервирована, настоящие страницы появятся при первом касании, как в уроке про отображение памяти. Под стеком стоит защитная страница: поток, который переполнил стек, получит SIGSEGV, а не испортит соседа.
Функция потока может вернуть void, u8, noreturn или !void. Ошибку из !void поток никому не вернёт: callFn напечатает её имя с трассой и молча закончит поток. Результат работы поток отдаёт только через память: указатель на ячейку в аргументах, как out: *usize в тестах эталона.
spawn глазами ядра
Посмотрим на spawn снаружи, своим zt strace из урока про системные вызовы. В контейнере с Linux запускаем ids из прошлого раздела. Строки, которые печатает сама программа, смешались с трассой (у них общий stderr), мы их убрали:
$ zt strace ./ids
...
dup(1) = 3
getpid() = 360
gettid() = 360
...
mmap(NULL, 16842752, 0x0, 0x20022, -1, 0x0) = 0xffff9b58e000
mprotect(0xffff9b59e000, 16777216, 0x3) = 0
rt_sigprocmask(0, 0xffff9c73cc70, 0xffffdb9cfd48, 8) = 0
clone(0x3d0f00, 0xffff9c55ca40, 0xffff9c55d250, 0xffff9c55d8c0, 0xffff9c55d250) = 361
rt_sigprocmask(2, 0xffffdb9cfd48, NULL, 8) = 0
mmap(NULL, 16842752, 0x0, 0x20022, -1, 0x0) = 0xffff9a57e000
mprotect(0xffff9a58e000, 16777216, 0x3) = 0
rt_sigprocmask(0, 0xffff9c73cc70, 0xffffdb9cfd48, 8) = 0
clone(0x3d0f00, 0xffff9b54ca40, 0xffff9b54d250, 0xffff9b54d8c0, 0xffff9b54d250) = 362
rt_sigprocmask(2, 0xffffdb9cfd48, NULL, 8) = 0
futex(0xffff9b54d250, 265, 362, NULL, NULL, 4294967295) = 0
exit_group(0) = ?
+++ exited with 0 +++
Один spawn это три системных вызова. mmap на 16842752 байт без прав (0x0), то есть 16 МиБ и ещё 64 КиБ. mprotect открывает на чтение и запись (0x3) ровно 16 МиБ, начиная с отступа 0x10000, а первые 64 КиБ отображения остаются без прав. Стек растёт вниз, так что эта полоса лежит прямо под ним: это и есть защитная полоса, 64 КиБ без прав. И clone, который возвращает номер нового потока, 361. Вокруг clone glibc блокирует все сигналы (rt_sigprocmask), чтобы новый поток не получил сигнал раньше, чем у него настроена локальная память.
Всё, что делает поток потоком, лежит в первом аргументе clone, 0x3d0f00. Разложим его по битам:
| Флаг | Бит | Что значит |
|---|---|---|
CLONE_VM | 0x100 | общее адресное пространство, а не копия, как у fork |
CLONE_FS | 0x200 | общий текущий каталог, корень и umask |
CLONE_FILES | 0x400 | общая таблица дескрипторов, а не копия |
CLONE_SIGHAND | 0x800 | общие обработчики сигналов |
CLONE_THREAD | 0x10000 | та же группа задач: тот же pid, getpid вернёт номер главного |
CLONE_SYSVSEM | 0x40000 | общие отмены семафоров System V |
CLONE_SETTLS | 0x80000 | свой указатель на локальную память потока (tpidr_el0 на aarch64, %fs на x86-64) |
CLONE_PARENT_SETTID | 0x100000 | записать номер потока в память создателя |
CLONE_CHILD_CLEARTID | 0x200000 | обнулить эту ячейку и разбудить ждущих, когда поток умрёт |
Сумма этих битов и есть 0x3d0f00. Сравни с прошлым уроком про песочницу: там clone получал CLONE_NEWPID, CLONE_NEWNS и другие флаги пространств имён, чтобы ребёнок видел меньше общего. Здесь флаги говорят обратное: делить с создателем всё. В ядре Linux нет отдельного понятия «поток»: процесс и поток это одна и та же задача, и разница только в том, какие части состояния она делит с создателем. fork это clone без этих флагов.
Последний флаг объясняет join. Ядро по смерти потока обнуляет ячейку child_tid и будит всех, кто ждёт на ней через futex. pthread_join и есть такое ожидание: futex(0xffff9b54d250, ..., 362, ...) в трассе ждёт, пока в ячейке по этому адресу лежит 362. Первого join в трассе не видно: судя по всему, T1 к тому моменту уже закончил, и ждать было нечего. Ветка LinuxThreadImpl в std.Thread делает то же самое руками: тот же набор флагов, тот же CLONE_CHILD_CLEARTID, и join крутит futex на ячейке сам.
Жизненный цикл: join, detach и конец процесса
Поток рождается в spawn и умирает, когда его функция вернулась. Но ресурсы, стек и запись в ядре, живут дольше, и о них кто-то должен позаботиться. Способов два. Первый это join: кто-то ждёт поток и освобождает его ресурсы. Второй это detach: поток объявляется отсоединённым, его никто не ждёт, и ресурсы он вернёт сам, когда закончит. Поток, который никто не ждёт и который не отсоединён, после конца висит, как зомби-процесс до waitpid: работы нет, а стек на 16 МиБ и запись в ядре заняты.
Переключай сценарии и смотри, кто сколько живёт:
- главныйspawn ×2своя работаjoin(T1) ждётитогjoin(T2)exit
- T1work(1)присоединён
- T2work(2)ждёт joinприсоединён
const t1 = try std.Thread.spawn(.{}, work, .{1});
const t2 = try std.Thread.spawn(.{}, work, .{2});
mainWork();
t1.join(); // блокирует, пока T1 не вернётся из work
t2.join(); // T2 давно закончил: возврат сразуГлавный блокируется в join, пока T1 не вернётся из своей функции. T2 закончил раньше, но его стек и запись в ядре живут до join, как зомби-процесс до waitpid. Второй join возвращается сразу и освобождает их.
- главныйspawn, detachцикл accept, никого не ждёт
- T1обслуживает клиентаосвободил себя
- T2обслуживает клиентаосвободил себя
const t = try std.Thread.spawn(.{}, serve, .{connfd});
t.detach(); // join больше нельзя, ресурсы освободит сам потокОтсоединённый поток никто не ждёт: закончив, он сам возвращает стек и запись в ядре. Так устроен сервер с потоком на соединение: главный цикл принимает следующего клиента и не копит завершённые потоки. Цена: узнать, когда поток закончил, или забрать его результат уже нечем.
- главныйspawn ×2работаreturn из main
- T1work(1)не доработалубит
- T2work(2)не доработалубит
_ = try std.Thread.spawn(.{}, work, .{1});
_ = try std.Thread.spawn(.{}, work, .{2});
return; // из main: это exit, процесс умирает вместе с T1 и T2Возврат из main означает exit, а exit завершает процесс целиком со всеми потоками посреди их работы. detach от этого не спасает: отсоединённый поток живёт, пока жив процесс. Буферы, которые потоки не успели сбросить, пропадают.
- работает
- заблокирован в join
- закончил, ресурсы ждут join
- оборван выходом процесса
Три правила, которые показывает виджет.
joinждёт конкретный поток. В отличие отwaitpid(-1, ...), дождаться любого из потоков нельзя. Если T2 закончил первым, а главный ждёт T1, T2 так и висит до своегоjoin. Стивенс считал это ошибкой стандарта, но так устроены pthreads.detachдля потоков, результат которых не нужен. Сервер с потоком на соединение отсоединяет каждый поток: главный цикл не должен копить закончившиеся потоки. Уstd.Threadиjoin, иdetachпоглощают значение: после любого из них трогатьThreadнельзя, это неопределённое поведение.- Конец
mainэто конец всех.returnизmainзовётexit, аexitзавершает процесс со всеми потоками посреди их работы.detachот этого не спасает: отсоединённый поток живёт, пока жив процесс. Тот же эффект уstd.process.exitиз любого потока и у необработанногоSIGSEGVв любом потоке.
Из последнего правила следует упражнение книги, где поток засыпает на секунду и ничего не печатает. Мы вернёмся к нему в домашнем задании.
Гонка на аргументе
Теперь ловушка. Программа заводит четыре потока в цикле и хочет, чтобы каждый напечатал свой номер. Вариант ptr отдаёт потоку адрес переменной цикла, вариант val её значение:
//! Упражнение про гонку на аргументе: каждый поток должен напечатать свой
//! номер. Вариант `ptr` отдаёт потоку адрес переменной цикла, вариант `val`
//! её копию.
const std = @import("std");
const n = 4;
fn byPointer(id: *const usize) void {
std.debug.print("{d} ", .{id.*});
}
fn byValue(id: usize) void {
std.debug.print("{d} ", .{id});
}
pub fn main(init: std.process.Init) !void {
const args = try init.minimal.args.toSlice(init.arena.allocator());
const by_pointer = args.len > 1 and std.mem.eql(u8, args[1], "ptr");
var threads: [n]std.Thread = undefined;
var i: usize = 0;
while (i < n) : (i += 1) {
threads[i] = if (by_pointer)
try std.Thread.spawn(.{}, byPointer, .{&i})
else
try std.Thread.spawn(.{}, byValue, .{i});
}
for (threads) |t| t.join();
std.debug.print("\n", .{});
}
Собираем с -O ReleaseSafe и запускаем по несколько раз. macOS 26, Apple M4 Max:
$ zig build-exe race.zig -O ReleaseSafe
$ for k in 1 2 3 4 5 6; do ./race ptr; done
2 4 2 4
3 4 4 4
4 4 4 4
2 4 4 4
3 4 4 4
2 4 4 4
$ for k in 1 2 3; do ./race val; done
1 0 3 2
0 3 2 1
0 1 3 2
В контейнере с Linux:
$ for k in 1 2 3 4 5 6; do ./race ptr; done
1 2 3 4
1 2 3 4
1 2 3 4
1 2 3 4
1 2 3 4
1 2 3 4
$ for k in 1 2 3; do ./race val; done
0 1 2 3
0 1 2 3
0 1 2 3
Вариант ptr сломан на обеих системах, только по-разному. На macOS потоки чаще всего видят 4, то есть значение, которое i получила после конца цикла. Номера 4 среди потоков вообще нет. На Linux картина стабильная и тоже неверная: каждый поток видит номер на единицу больше своего, и никто не видит 0.
Что произошло. spawn честно скопировал кортеж .{&i}, но в кортеже лежит указатель, и копия указателя смотрит на ту же ячейку i на стеке главного потока. Поток читает id.* не в момент spawn, а когда ядро даст ему процессор. К этому времени главный поток уже сделал i += 1, а может быть, прошёл весь цикл. Результат зависит от того, кто успел первым: это и есть гонка. Похоже, на Linux в контейнере новый поток просыпается тогда, когда главный успел сделать ровно один i += 1, отсюда стабильный сдвиг на единицу. Стабильность обманчива: на другой машине или под нагрузкой ответ будет другим.
Вариант val верен, хотя порядок строк в нём случаен. Порядок и не обещан: потоки печатают, когда их запустит планировщик. Зато каждый напечатал свой номер, потому что i скопирована в кортеж в момент spawn, а кортеж скопирован в кучу нового потока.
У версии ptr есть и вторая проблема, которую этот запуск не показал. Указатель смотрит на стек главного потока. Здесь это безопасно, потому что main ждёт все потоки через join, прежде чем вернуться. Функция, которая заводит поток с указателем на свою локальную переменную и возвращается, оставляет поток с висячим указателем на чужой стек, который уже переписан следующими вызовами.
Отсюда правило, которым мы будем пользоваться весь блок: в поток передавай значения, а не адреса своих переменных. Если нужно передать адрес, объект по этому адресу должен жить дольше потока (в куче или на стеке того, кто гарантированно дождётся join), и никто не должен его менять, пока поток читает. Книга решает ту же задачу в эхо-сервере, выделяя под каждый дескриптор соединения отдельный int через malloc. В Zig это не нужно: spawn уже выделил память под копию кортежа.
Как эту гонку проверить в тесте? Запуски выше показывают, что поймать её случайно можно, но не гарантированно: на каждой машине и под каждой нагрузкой своя картина. Эталон делает гонку детерминированной: потоки стоят на воротах Io.Event и читают значение, только когда главный поток прошёл весь цикл и открыл ворота. Тогда вариант с адресом всегда видит конечное значение, а вариант с копией всегда свой номер. Это raceDemo в шаге проекта ниже.
Потоки и процессы: что общего и что нет
Соберём сравнение в одну таблицу. Слева процесс, созданный fork, справа поток, созданный spawn:
Процесс после fork | Поток после spawn | |
|---|---|---|
| память | копия при записи, изменения не видны друг другу | одна на всех, изменения видны сразу |
| дескрипторы | копия таблицы, описания файлов общие | одна таблица на всех |
| обработчики сигналов | копия | общие |
| номер | свой pid | общий pid, свой tid |
| кто кого ждёт | родитель ребёнка, waitpid(-1) ждёт любого | любой любого, join ждёт конкретного |
падение с SIGSEGV | умирает один процесс | умирает весь процесс |
| обмен данными | пайпы, сокеты, общая память через mmap | любая переменная, до которой дотянулся указатель |
| цена ошибки | другой процесс не испортить | любой поток портит память всех |
И цена рождения. Программа ниже две тысячи раз создаёт поток с пустой функцией и ждёт его, потом две тысячи раз делает fork с немедленным _exit в ребёнке и ждёт через waitpid. Аргумент задаёт, сколько мегабайт кучи процесс тронул перед замером:
//! Цена рождения: создать и дождаться поток против создать и дождаться
//! процесс. Работы у обоих нет, меряется только механизм.
const std = @import("std");
const c = std.c;
const rounds = 2000;
fn nothing() void {}
fn now() u64 {
var ts: c.timespec = undefined;
_ = c.clock_gettime(.MONOTONIC, &ts);
return @as(u64, @intCast(ts.sec)) * std.time.ns_per_s + @as(u64, @intCast(ts.nsec));
}
pub fn main(init: std.process.Init) !void {
const args = try init.minimal.args.toSlice(init.arena.allocator());
// Сколько мегабайт тронутой кучи у процесса: fork копирует таблицы
// страниц, и цена растёт вместе с памятью.
const mb = if (args.len > 1) try std.fmt.parseInt(usize, args[1], 10) else 0;
const heap = try std.heap.page_allocator.alloc(u8, @max(mb, 1) << 20);
defer std.heap.page_allocator.free(heap);
@memset(heap, 1);
const t0 = now();
for (0..rounds) |_| {
const t = try std.Thread.spawn(.{}, nothing, .{});
t.join();
}
const t1 = now();
for (0..rounds) |_| {
const pid = c.fork();
if (pid == 0) c._exit(0);
var status: c_int = 0;
_ = c.waitpid(pid, &status, 0);
}
const t2 = now();
std.debug.print("куча {d:>4} МБ · поток {d:>5.1} мкс · процесс {d:>6.1} мкс\n", .{
mb,
@as(f64, @floatFromInt(t1 - t0)) / rounds / 1000,
@as(f64, @floatFromInt(t2 - t1)) / rounds / 1000,
});
}
$ zig build-exe cost.zig -O ReleaseFast
$ for m in 0 64 512; do ./cost $m; done
| Куча, МБ | macOS: поток | macOS: процесс | Linux: поток | Linux: процесс |
|---|---|---|---|---|
| 0 | 32 до 41 мкс | 704 до 721 мкс | 58 до 69 мкс | 164 до 171 мкс |
| 64 | 29 до 33 мкс | 628 до 655 мкс | 60 до 61 мкс | 943 до 1036 мкс |
| 512 | 29 мкс | 655 до 695 мкс | 36 до 48 мкс | 7048 до 7610 мкс |
Сняты два прогона на каждую строку: Apple M4 Max (12 производительных и 4 энергоэффективных ядра), macOS 26.6.2, и контейнер runner-zig:dev на той же машине (OrbStack, ядро Linux 7.0.14, arm64, 16 vCPU). Машину в это время грузили соседние задачи, load average около 18 на 16 ядрах, так что смотри на порядок, а не на третий знак.
Поток стоит десятки микросекунд, и эта цена не зависит от размера процесса: новому потоку копировать нечего. fork на Linux растёт вместе с памятью: копия при записи не копирует страницы, но таблицы страниц копирует, и на 512 МБ тронутой кучи это уже семь миллисекунд на каждое соединение. На пустом процессе Linux делает fork всего в два или три раза дольше потока, на процессе с 512 МБ уже в сто с лишним раз. На macOS fork дорог с самого начала (около 0,7 мс, в двадцать раз дороже потока) и в этом замере почти не растёт с кучей: ядро XNU управляет памятью по-своему, и выводы Linux на него не переносятся. Для сервера вывод один: чем больше процесс, тем дороже ему каждое соединение через fork, а цена потока от размера не зависит.
Дешевизна потока не бесплатна. Всё, что процесс получал от ядра даром (изоляцию памяти, отдельные дескрипторы), поток должен обеспечить сам, аккуратностью программиста. Об этом следующие три урока блока.
Шаг проекта: сервер на потоках
В прошлом уроке эхо-сервер получил процессы и select, poll, epoll. Шаг 69 эталона our-tiny добавляет третью модель: echoservert, поток на соединение. Заодно поток на соединение получает и TINY из урока 66, а для самих потоков появится первая вычислительная задача: умножение матриц полосами строк. Новые файлы src/conc/echo_threads.zig и src/conc/matmul.zig, в src/http/tiny.zig одна новая функция, остальное это подключение и тесты.
echo_threads.zig
//! `echoservert` из главы 12: поток на соединение. Главный поток только
//! принимает, каждое соединение обслуживает свой поток, отсоединённый
//! через `detach`: его никто не ждёт, ресурсы вернёт система, когда он
//! кончится. Здесь же TINY на потоках и демонстрация гонки на аргументе.
const std = @import("std");
const c = std.c;
const posix = std.posix;
const Io = std.Io;
const socket = @import("../net/socket.zig");
const fdio = @import("../net/fdio.zig");
const echo = @import("../net/echo.zig");
const tiny = @import("../http/tiny.zig");
/// `accept` с повтором на `EINTR`, без имени клиента.
pub fn acceptRetry(listenfd: c.fd_t) !c.fd_t {
while (true) {
const fd = c.accept(listenfd, null, null);
if (fd >= 0) return fd;
if (posix.errno(fd) != .INTR) return error.AcceptFailed;
}
}
/// Эхо-сервер: поток на соединение. `connfd` уходит в поток по значению,
/// копией внутри `std.Thread.spawn`. Книга передаёт в `pthread_create`
/// указатель на `malloc`-нутую копию ровно ради этого. Вот так нельзя:
///
/// var connfd = try acceptRetry(listenfd);
/// _ = try std.Thread.spawn(.{}, sessionPtr, .{&connfd});
///
/// Поток читает `connfd.*` когда-нибудь потом, а главный поток к тому
/// времени мог принять следующего клиента и переписать ту же переменную:
/// два потока обслуживают один дескриптор, а первый клиент брошен.
/// `raceDemo` ниже показывает ту же ошибку без сети.
pub fn serve(listenfd: c.fd_t, max_clients: usize, quiet: bool) !void {
var served: usize = 0;
while (max_clients == 0 or served < max_clients) : (served += 1) {
const connfd = try acceptRetry(listenfd);
const thread = std.Thread.spawn(.{}, session, .{ connfd, quiet }) catch {
socket.close(connfd);
continue;
};
thread.detach();
}
}
fn session(connfd: c.fd_t, quiet: bool) void {
defer socket.close(connfd);
// Свой писатель лога у каждого потока: общий буфер испортили бы все
// сразу. `echo` сбрасывает его после каждой строки, одна строка это
// один `write`, и строки разных потоков в stderr не перемешиваются.
var log_buf: [128]u8 = undefined;
var stderr: fdio.Writer = .init(2, &log_buf);
var discard: Io.Writer.Discarding = .init(&.{});
echo.echo(connfd, if (quiet) &discard.writer else &stderr.interface) catch {};
}
/// TINY на потоках: тот же `Server`, тот же `handle`, только каждое
/// соединение в своём потоке. Имя клиента копируется в поток по значению
/// по той же причине, что и `connfd`: буфер `accept` переиспользуется.
pub fn serveTiny(server: *tiny.Server, max_clients: usize) !void {
var served: usize = 0;
while (max_clients == 0 or served < max_clients) : (served += 1) {
var name: Peer = .{};
const conn = try socket.accept(server.listenfd, &name.buf);
name.len = conn.name.len;
const thread = std.Thread.spawn(.{}, tinySession, .{ server, conn.fd, name }) catch {
socket.close(conn.fd);
continue;
};
thread.detach();
}
}
pub const Peer = struct {
buf: [socket.name_max_len]u8 = undefined,
len: usize = 0,
pub fn slice(p: *const Peer) []const u8 {
return p.buf[0..p.len];
}
};
fn tinySession(server: *tiny.Server, connfd: c.fd_t, peer: Peer) void {
defer socket.close(connfd);
server.handle(connfd, peer.slice());
}
/// Упражнение 12.5 без сети: `n` потоков получают номер итерации.
/// `by_pointer = true` передаёт адрес переменной цикла, как в `race.c`
/// книги. Потоки стоят на воротах, пока цикл не кончится, так что гонка
/// проигрывается всегда: каждый прочтёт уже конечное значение `n`.
/// `by_pointer = false` передаёт копию, и каждый видит свой номер.
pub fn raceDemo(io: Io, results: []usize, by_pointer: bool) !void {
var gate: Io.Event = .unset;
var threads: [64]std.Thread = undefined;
const n = results.len;
std.debug.assert(n <= threads.len);
var i: usize = 0;
while (i < n) : (i += 1) {
threads[i] = if (by_pointer)
try std.Thread.spawn(.{}, readShared, .{ io, &gate, &i, &results[i] })
else
try std.Thread.spawn(.{}, readOwn, .{ io, &gate, i, &results[i] });
}
gate.set(io);
for (threads[0..n]) |t| t.join();
}
fn readShared(io: Io, gate: *Io.Event, index: *const usize, out: *usize) void {
gate.waitUncancelable(io);
out.* = index.*;
}
fn readOwn(io: Io, gate: *Io.Event, index: usize, out: *usize) void {
gate.waitUncancelable(io);
out.* = index;
}
Разберём по функциям.
acceptRetry. accept может вернуть EINTR, если сигнал прервал ожидание. В сервере на процессах так приходил SIGCHLD, у потоков своих сигналов нет, но повтор ничего не стоит и спасает от чужих.
serve. Цикл книжного echoservert, только без malloc. connfd уходит в кортеж по значению, и после spawn главный поток волен переписать свою переменную следующим accept: у потока уже своя копия. Если spawn не смог завести поток (кончились потоки у пользователя, ThreadQuotaExceeded, или память), соединение закрывается, и сервер идёт дальше: один отказ не должен ронять всех. Сразу после spawn поток отсоединяется. В книге поток отсоединяет себя сам первой строкой, pthread_detach(pthread_self()). Разница только в том, кто это делает, результат тот же.
session. Поток закрывает соединение сам, в defer. Главный поток дескриптор не закрывает, и это не случайность: вспомни, что сервер на процессах закрывал connfd в обоих процессах. Почему здесь иначе, спрашивает первое упражнение урока.
Теперь лог. Эхо из урока про сокеты пишет строку server received N bytes в писатель, который ему дали. В итеративном сервере это был один общий писатель на stderr. Отдай один буферизованный писатель нескольким потокам, и они будут одновременно дописывать в один буфер и двигать одно поле длины: строки перемешаются, а то и потеряются. Поэтому у каждого потока свой писатель на 128 байт поверх дескриптора 2, а echo сбрасывает его после каждой строки. Одна строка лога это один write, а запись в дескриптор ядро выполняет целиком, так что строки разных потоков не рвут друг друга.
serveTiny и Peer. TINY хочет имя клиента host:port для лога. socket.accept пишет его в буфер, который ему дали, и отдаёт срез в этот буфер. Срез в поток отдать нельзя: это тот же адрес переменной, только в профиль. Следующий accept перепишет буфер, и поток напишет в лог имя чужого клиента. Peer это буфер вместе с длиной, и он уходит в поток целиком, по значению, как и connfd. Кортеж .{ server, conn.fd, name } при этом копирует около восьмидесяти байт Peer, и это правильная цена.
А вот server уходит указателем, и здесь это верно: Server живёт в main всё время работы сервера, и все потоки только читают его поля. Почему это безопасно, разберём в следующем разделе.
raceDemo. Упражнение про гонку на аргументе, только без сети и без случайности. gate это Io.Event, флаг, на котором можно ждать: waitUncancelable блокирует поток, пока кто-то не вызовет set. Потоки стартуют и сразу встают на ворота, главный доводит цикл до конца и открывает их. После этого каждый поток с адресом читает i, в которой уже лежит n, а каждый поток с копией читает свой номер. Примитивы из std.Io подробно разберём в следующем уроке, здесь нам хватит того, что ворота делают порядок событий известным заранее.
tiny.zig: транзакция на принятом соединении
Итеративный TINY делал всё в serveOnce: accept, транзакция, close. Потоку нужна средняя часть отдельно, от уже принятого соединения. Выносим её в handle:
defer socket.close(conn.fd);
+ s.handle(conn.fd, conn.name);
+ }
+
+ /// Одна транзакция на уже принятом соединении. Дескриптор закрывает
+ /// вызывающий. Конкурентные версии из урока 69 и 71 зовут это из своих
+ /// потоков: у `Server` нет изменяемого состояния, лог пишется одним
+ /// `write` на строку, поэтому один сервер обслуживает много потоков.
+ pub fn handle(s: *Server, connfd: c.fd_t, peer: []const u8) void {
var arena_state: std.heap.ArenaAllocator = .init(s.gpa);
defer arena_state.deinit();
- s.doit(arena_state.allocator(), conn.fd, conn.name) catch |err| s.log("{s} transaction failed: {t}", .{ conn.name, err });
+ s.doit(arena_state.allocator(), connfd, peer) catch |err| s.log("{s} transaction failed: {t}", .{ peer, err });
}
Может ли один Server обслуживать много потоков сразу? Пройдём по тому, что handle трогает.
- Поля
Server(gpa,options,listenfd) послеinitтолько читаются. Чтение одного и того же из разных потоков безопасно. - Арена своя у каждой транзакции: она создаётся на стеке потока в
handleи там же умирает. gpaобщий. В отладочной сборке этоDebugAllocator, у которого в многопоточной программе есть свой мьютекс, в релизной с libc этоmalloc, потокобезопасный по стандарту. Оба можно звать из многих потоков сразу.- Буферы чтения и записи в
doitлокальные, на стеке потока. logформатирует строку в локальный буфер и отдаёт её однимwriten, как лог эхо-сервера.
Изменяемого общего состояния нет, значит, и гонок нет. Так проще всего писать потоковый код: не защищать общее, а не заводить его.
Остаётся CGI, и здесь потоки встречаются с fork. serveDynamic делает fork из потока, который обслуживает запрос. В ребёнке после fork живёт только один поток, копия того, кто позвал fork. Все остальные потоки родителя в ребёнке исчезают, и если кто-то из них в этот момент держал блокировку внутри malloc, в ребёнке она останется закрытой навсегда. Поэтому ребёнок многопоточного процесса до execve может звать только async-signal-safe функции, и поэтому serveDynamic ещё в уроке 66 собирает argv и окружение до fork. Тогда это была предосторожность, теперь это условие работы.
И ещё одна деталь, уже без аккуратного решения. Ребёнок наследует все открытые дескрипторы процесса, а у многопоточного сервера это соединения всех остальных потоков. Сокеты открываются без флага CLOEXEC, поэтому CGI, запущенный из одного потока, получает в подарок чужие соединения и держит их открытыми, пока не выйдет. adder выходит за миллисекунды, и на ответах это не видно. CGI, который работает десять секунд, задержит на десять секунд конец чужих ответов: клиент HTTP/1.0 ждёт закрытия соединения, а у соединения теперь два владельца. Как это увидеть и как лечить, в упражнениях. В уроке 71 та же утечка найдётся в zbox.
matmul.zig: потоки для вычислений
Сервер это конкурентность ради ожидания: потоки большую часть времени спят в read. Вторая причина заводить потоки это параллельность ради скорости: разложить вычисление на ядра. Самый чистый пример из домашних заданий книги это умножение матриц. Строка i произведения зависит только от строки i левой матрицы и всей правой, значит, строки можно раздать потокам:
//! Параллельное умножение матриц (домашние 12.16 и 12.17 в духе): матрица
//! режется на полосы строк, каждая полоса своему потоку. Потоки пишут в
//! разные строки результата и только читают `a` и `b`, поэтому им не
//! нужны ни блокировки, ни атомики: общие данные есть, общих записей нет.
//! Матрицы квадратные `n x n`, по строкам в одном срезе.
const std = @import("std");
/// Строки `lo..hi` результата. Порядок `i-k-j`: внутренний цикл идёт по
/// строке `b` и строке `out` подряд, как в уроке про кэши. Порядок сложений
/// у каждого элемента один и тот же, поэтому параллельная версия совпадает
/// с последовательной бит в бит, а не «примерно».
fn rows(a: []const f64, b: []const f64, out: []f64, n: usize, lo: usize, hi: usize) void {
for (lo..hi) |i| {
const row = out[i * n ..][0..n];
@memset(row, 0);
for (0..n) |k| {
const aik = a[i * n + k];
for (row, b[k * n ..][0..n]) |*o, bkj| o.* += aik * bkj;
}
}
}
pub fn multiply(a: []const f64, b: []const f64, out: []f64, n: usize) void {
rows(a, b, out, n, 0, n);
}
/// `threads` потоков, у каждого полоса примерно в `n / threads` строк.
/// Остаток от деления достаётся первым полосам по строке.
pub fn multiplyParallel(a: []const f64, b: []const f64, out: []f64, n: usize, threads: usize) !void {
std.debug.assert(a.len == n * n and b.len == n * n and out.len == n * n);
var handles: [64]std.Thread = undefined;
const t_count = @min(threads, handles.len, @max(n, 1));
var lo: usize = 0;
for (0..t_count) |t| {
const hi = lo + n / t_count + @intFromBool(t < n % t_count);
handles[t] = try std.Thread.spawn(.{}, rows, .{ a, b, out, n, lo, hi });
lo = hi;
}
for (handles[0..t_count]) |h| h.join();
}
В этом коде три важные вещи.
Никакой синхронизации, кроме join. Все потоки читают a и b, но ни один их не пишет. Каждый пишет в out, но только в свои строки, полосы не пересекаются. Общие данные есть, общих записей нет, поэтому гонки нет. join в конце нужен не для корректности вычисления, а чтобы multiplyParallel не вернулась раньше, чем результат готов.
Остаток раздаётся по строке. Если n не делится на число потоков, первые n % t_count полос получают на строку больше. Так полосы идут подряд без дыр, покрывают все строки, а их высоты отличаются не больше чем на единицу. t_count при этом не больше n: на матрице в три строки восемь потоков не нужны.
Результат бит в бит. Порядок сложений у каждого элемента одинаков в последовательной и параллельной версиях: rows одна и та же функция, поток отличается только границами. Сложение чисел с плавающей точкой не ассоциативно (помнишь урок про float?), и если бы потоки делили работу по k, а не по строкам, ответы расходились бы в последних битах. Порядок i-k-j из урока про код, дружественный кэшу: внутренний цикл идёт по строке b и строке out подряд.
Задача урока в разделе «Практика» это та же функция в первом, наивном варианте, где спрятаны две ошибки. Одна из них это гонка из прошлого раздела.
main.zig, root.zig, build.zig
В src/main.zig у подкоманды tiny появляется флаг --threads, добавляется подкоманда echoservert, а conc-bench из прошлого урока получает модель threads в список по умолчанию:
@@
\\Использование:
- \\ tiny tiny <port> [root] веб-сервер, корень статики и cgi-bin (по умолчанию .)
+ \\ tiny tiny [--threads] <port> [root] веб-сервер: итеративный или поток на соединение
\\ tiny hostinfo [--std] <name> все адреса имени через getaddrinfo или std.Io.net
@@
\\ tiny echoservers --select|--poll|--epoll <port> эхо: события (--epoll это kqueue на macOS)
- \\ tiny conc-bench [--clients 100,1000] [--active A] [--lines K] [--models procs,select,poll,epoll] [--json]
+ \\ tiny echoservert <port> эхо: поток на соединение
+ \\ tiny conc-bench [--clients 100,1000] [--active A] [--lines K] [--models procs,threads,select,poll,epoll] [--json]
\\
@@
.{ "echoservers", cmdEchoservers },
+ .{ "echoservert", cmdEchoservert },
.{ "conc-bench", cmdConcBench },
@@
fn cmdTiny(ctx: Ctx, rest: []const [:0]const u8) !void {
- if (rest.len < 1) return fail(ctx.out, usage);
- const port = try parsePort(ctx, rest[0]);
- const root = if (rest.len > 1) rest[1] else ".";
+ var args = rest;
+ var mode: enum { iterative, threads } = .iterative;
+ if (args.len > 0 and std.mem.eql(u8, args[0], "--threads")) {
+ mode = .threads;
+ args = args[1..];
+ }
+ if (args.len < 1) return fail(ctx.out, usage);
+ const port = try parsePort(ctx, args[0]);
+ const root = if (args.len > 1) args[1] else ".";
var server: tiny.Server = try .init(ctx.gpa, ctx.io, .{ .port = port, .root = root });
defer server.deinit();
- try ctx.log.print("tiny: listening on port {d}, root {s}\n", .{ server.port, root });
+ try ctx.log.print("tiny: listening on port {d}, root {s}, {t}\n", .{ server.port, root, mode });
try ctx.log.flush();
- try server.serveForever();
+ switch (mode) {
+ .iterative => try server.serveForever(),
+ .threads => try conc.echo_threads.serveTiny(&server, 0),
+ }
}
@@
+fn cmdEchoservert(ctx: Ctx, rest: []const [:0]const u8) !void {
+ if (rest.len != 1) return fail(ctx.out, usage);
+ const listenfd = try listenOn(ctx, "echoservert", try parsePort(ctx, rest[0]));
+ defer tiny.socket.close(listenfd);
+ try conc.echo_threads.serve(listenfd, 0, false);
+}
+
fn cmdConcBench(ctx: Ctx, rest: []const [:0]const u8) !void {
var clients: []const usize = &.{ 100, 1000 };
- var models: []const conc.bench.Model = &.{ .procs, .select, .poll, .epoll };
+ var models: []const conc.bench.Model = &.{ .procs, .threads, .select, .poll, .epoll };
var cfg: conc.bench.Config = .{ .clients = 0 };
В src/root.zig два новых модуля в пространстве conc, в build.zig номер шага:
@@
pub const bench = @import("conc/bench.zig");
+ pub const echo_threads = @import("conc/echo_threads.zig");
+ pub const matmul = @import("conc/matmul.zig");
};
@@
/// соответствует файл `tests/step_NN.zig`, и все они зелёные на финале.
-const project_steps = [_]u8{ 63, 64, 65, 66, 68 };
+const project_steps = [_]u8{ 63, 64, 65, 66, 68, 69 };
Замер conc-bench учится запускать echoservert как ещё одну модель. Сервер работает в ребёнке после fork, поэтому ветке threads хватает одной строки: тот же serve, что у подкоманды, без лога. Тест шага 68 гоняет новую модель вместе с остальными:
@@
const echo_epoll = @import("echo_epoll.zig");
+const echo_threads = @import("echo_threads.zig");
@@
procs,
+ threads,
select,
@@
.procs => echo_procs.serve(listenfd, 0, &discard.writer),
+ .threads => echo_threads.serve(listenfd, 0, true),
.select => if (echo_select.serve(gpa, listenfd, opts)) |_| {} else |err| err,
@@
_ = conc.bench.raiseFdLimit(1024);
- inline for (.{ .procs, .select, .poll, .epoll }) |model| {
+ inline for (.{ .procs, .threads, .select, .poll, .epoll }) |model| {
const r = try conc.bench.run(testing.allocator, testing.io, model, .{ .clients = 8, .lines = 5 });
Строку threads в таблицах замера из прошлого урока теперь можно снять самому: tiny conc-bench --models procs,threads,epoll.
Прогон
Эхо-сервер и два клиента. Первый говорит, молчит две секунды и говорит снова, второй приходит, пока первый молчит:
$ tiny echoservert 15269
echoservert: listening on port 15269
server received 13 bytes
server received 13 bytes
server received 20 bytes
$ (printf 'первый\n'; sleep 2; printf 'первый ещё\n') | tiny echoclient 127.0.0.1 15269 &
$ printf 'второй\n' | tiny echoclient 127.0.0.1 15269
второй
первый
первый ещё
Второй клиент получил своё эхо сразу, пока первый ещё держал соединение: итеративный эхо-сервер из урока 64 заставил бы его ждать две секунды. Байты в логе считаются в UTF-8: шесть кириллических букв по два байта и перевод строки дают 13.
Пока клиенты подключены, у процесса три потока. На macOS их покажет ps -M:
$ ps -M -p 3844
USER PID TT %CPU STAT PRI STIME UTIME COMMAND
bondiano 3844 ?? 0.0 S 31T 0:00.01 0:00.01 zig-out/bin/tiny echoservert 15270
3844 0.0 S 31T 0:00.00 0:00.00
3844 0.0 S 31T 0:00.00 0:00.00
На Linux каждый поток это каталог в /proc/PID/task, и у каждого там свой файл comm с именем:
$ ls /proc/215/task
215
224
225
$ cat /proc/215/task/*/comm
tiny
tiny
tiny
Имя у всех одно: новый поток наследует имя создателя, а std.Thread сам его не меняет. Запомни эту строку, она понадобится в шаге zt.
Теперь TINY. Клиент открывает соединение через nc и три секунды молчит, а curl в это время просит страницу с лимитом в две секунды. Сначала итеративный сервер, потом тот же с --threads:
$ tiny tiny 8069 zig-out
tiny: listening on port 8069, root zig-out, iterative
$ sleep 3 | nc 127.0.0.1 8069 &
$ curl -s -o /dev/null --max-time 2 -w 'home %{http_code} %{time_total}\n' http://127.0.0.1:8069/home.html
home 000 2.011364
$ tiny tiny --threads 8069 zig-out
tiny: listening on port 8069, root zig-out, threads
$ sleep 3 | nc 127.0.0.1 8069 &
$ curl -s -o /dev/null --max-time 2 -w 'home %{http_code} %{time_total}\n' http://127.0.0.1:8069/home.html
home 200 0.000539
$ curl -s --max-time 2 'http://127.0.0.1:8069/cgi-bin/adder?15&213' | grep answer
<p>The answer is: 15 + 213 = 228
Итеративный сервер сидит в read молчащего клиента, и curl уходит ни с чем через две секунды. Сервер на потоках отвечает за полмиллисекунды: молчащий клиент держит только свой поток. CGI работает из потока так же, как из главного.
Тесты шага
//! Шаг 69: поток на соединение (эхо и TINY), гонка на адресе переменной
//! цикла, параллельное умножение матриц полосами строк.
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 c = std.c;
fn connect(port: u16) !c.fd_t {
const fd = try socket.openClientfd("127.0.0.1", port);
const tv: c.timeval = .{ .sec = 5, .usec = 0 };
_ = c.setsockopt(fd, c.SOL.SOCKET, c.SO.RCVTIMEO, &tv, @sizeOf(c.timeval));
return fd;
}
fn threadsServe(listenfd: c.fd_t, n: usize) void {
conc.echo_threads.serve(listenfd, n, true) catch |err| std.debug.print("echo_threads: {t}\n", .{err});
}
test "echoservert: три клиента одновременно, ответы в обратном порядке" {
const listenfd = try socket.openListenfd(0);
defer socket.close(listenfd);
const port = socket.localPort(listenfd).?;
const thread = try std.Thread.spawn(.{}, threadsServe, .{ listenfd, 3 });
var fds: [3]c.fd_t = undefined;
for (&fds) |*fd| fd.* = try connect(port);
thread.join();
for (0..3) |k| {
const i = 2 - k;
var buf: [32]u8 = undefined;
const line = try std.fmt.bufPrint(&buf, "поток {d}\n", .{i});
try fdio.writen(fds[i], line);
var reply: [32]u8 = undefined;
const n = try fdio.readn(fds[i], reply[0..line.len]);
try testing.expectEqualStrings(line, reply[0..n]);
}
for (fds) |fd| socket.close(fd);
}
test "упражнение 12.5: по адресу все потоки видят конечное значение, по значению каждый свой" {
var by_pointer: [8]usize = undefined;
try conc.echo_threads.raceDemo(testing.io, &by_pointer, true);
for (by_pointer) |v| try testing.expectEqual(@as(usize, 8), v);
var by_value: [8]usize = undefined;
try conc.echo_threads.raceDemo(testing.io, &by_value, false);
for (by_value, 0..) |v, i| try testing.expectEqual(i, v);
}
fn tinyThreads(server: *tiny.Server, n: usize) void {
conc.echo_threads.serveTiny(server, n) catch |err| std.debug.print("serveTiny: {t}\n", .{err});
}
test "tiny --threads: медленный клиент не держит остальных" {
var arena_state: std.heap.ArenaAllocator = .init(testing.allocator);
defer arena_state.deinit();
const arena = arena_state.allocator();
var root: support.Root = try .create(arena);
defer root.cleanup();
var server: tiny.Server = try .init(testing.allocator, testing.io, .{ .port = 0, .root = root.path, .quiet = true });
defer server.deinit();
const thread = try std.Thread.spawn(.{}, tinyThreads, .{ &server, 3 });
// Первый клиент прислал полстроки запроса и молчит. Итеративный TINY
// ждал бы его вечно; на потоках его ждёт только свой поток.
const slow = try connect(server.port);
defer socket.close(slow);
try fdio.writen(slow, "GET /hello.txt HT");
const home = try support.get(arena, server.port, "/home.html");
try testing.expectEqual(200, home.status);
const cgi = try support.get(arena, server.port, "/cgi-bin/adder?1&2");
try testing.expect(std.mem.indexOf(u8, cgi.body, "1 + 2 = 3") != null);
try fdio.writen(slow, "TP/1.0\r\n\r\n");
const reply = try support.parseReply(try support.readAll(arena, slow));
try testing.expectEqual(200, reply.status);
thread.join();
}
test "matmul: полосы строк на 1, 3 и 8 потоках совпадают с последовательным бит в бит" {
const n = 67;
const gpa = testing.allocator;
const a = try gpa.alloc(f64, n * n);
defer gpa.free(a);
const b = try gpa.alloc(f64, n * n);
defer gpa.free(b);
const seq = try gpa.alloc(f64, n * n);
defer gpa.free(seq);
const par = try gpa.alloc(f64, n * n);
defer gpa.free(par);
var prng: std.Random.DefaultPrng = .init(69);
for (a) |*x| x.* = prng.random().float(f64) - 0.5;
for (b) |*x| x.* = prng.random().float(f64) - 0.5;
conc.matmul.multiply(a, b, seq, n);
for ([_]usize{ 1, 3, 8, 100 }) |t| {
@memset(par, std.math.nan(f64));
try conc.matmul.multiplyParallel(a, b, par, n, t);
try testing.expectEqualSlices(f64, seq, par);
}
// И это действительно произведение: единичная матрица ничего не меняет.
@memset(b, 0);
for (0..n) |i| b[i * n + i] = 1;
conc.matmul.multiply(a, b, seq, n);
try testing.expectEqualSlices(f64, a, seq);
}
- «три клиента одновременно» открывает три соединения, дожидается, пока сервер примет все три (
thread.joinна циклеserveсmax_clients = 3), и говорит с ними в обратном порядке. Итеративный сервер на этом тесте завис бы: он сел бы читать первого клиента, а тест ещё ничего не пишет, он ждёт, пока сервер примет всех троих. У каждого клиента таймаут чтения в пять секунд (SO_RCVTIMEO), так что сломанный сервер даст падение теста, а не вечное ожидание. - «упражнение 12.5» это
raceDemoв обоих вариантах. Сравни, как он устроен, сrace.zigвыше: одна и та же ошибка, но проверяется не случайно, а всегда. - «медленный клиент не держит остальных» присылает половину строки запроса и молчит, пока двое других получают статику и CGI. Потом дописывает запрос и тоже получает свои 200.
- «matmul» сверяет параллельную версию с последовательной бит в бит на 1, 3, 8 и 100 потоках, при
n = 67, которое не делится ни на 3, ни на 8. Перед каждой сверкой результат забит NaN: строку, которую никто не посчитал, видно сразу. В конце проверка, что это вообще произведение: единичная матрица ничего не меняет.
$ zig build test -Dstep=69 --summary all
Build Summary: 6/6 steps succeeded; 4/4 tests passed
Зелёные и на macOS, и в контейнере с Linux, вместе с тестами всех прошлых шагов.
Шаг zt: профилировщик различает потоки
В уроке про сигналы мы написали zt prof: таймер ITIMER_PROF, обработчик SIGPROF, который идёт по цепочке указателей кадров и складывает адреса в кольцо, и свёрнутые стеки для flame graph из урока про профилирование. Там же стояла заметка: стеки собираются только у потока, который завёл таймер. Пора её закрыть. Возьмём программу с тремя потоками, heavy, medium и light, с нагрузкой 3 к 2 к 1 (её исходник ниже, в фикстурах шага), и отдадим профилировщику в том виде, в каком он вышел из урока 49. Контейнер с Linux, arm64:
$ zt prof ./threads
heavy 211 ms
medium 137 ms
light 71 ms
[libc.so.6];[libc.so.6];__clock_gettime 1
[libc.so.6];[libc.so.6];run;spin 339
zt prof: сэмплов 340, частота 997 Гц, затёрто 0
Сэмплы есть, и стеки даже целые, но все три потока слились в одну строку. Кто из них сколько жёг, из профиля не узнать. А целые стеки здесь везение. Обходу кадров нужна верхняя граница стека прерванного потока, а у сэмплера из урока 49 граница одна, главного потока. Стеки потоков на этой машине лежат ниже стека главного, и граница главного случайно их накрывает. Будь раскладка другой, обход обрывался бы на первом кадре или шёл бы по мусору до первой проверки.
Шаг 69 делает три вещи: обработчик узнаёт, какой поток он прервал, у каждого потока своя граница стека, и каждый сэмпл помечен номером потока, а отчёт по флагу --threads ставит имя потока первым кадром.
Кому достаётся SIGPROF
Таймер ITIMER_PROF один на процесс и считает процессорное время всех потоков вместе. Когда он истекает, ядро шлёт SIGPROF процессу, а доставляет сигнал одному потоку. Какому? Таймер проверяется на тиках планировщика, и истекает он на тике того потока, который в этот момент на процессоре. Ядро отдаёт сигнал ему же, если он не заблокировал SIGPROF. Значит, сэмпл достаётся тому, кто жёг процессор, а нам это и нужно.
Обработчик сигнала выполняется в прерванном потоке и на его стеке. Отсюда план: первым делом обработчик спрашивает у ядра номер своего потока через gettid. Это системный вызов без побочных эффектов, он безопасен в обработчике. По номеру он находит поток в таблице и берёт оттуда верхнюю границу стека, нижняя граница, как и раньше, это кадр самого обработчика.
Как граница попадает в таблицу? В самом обработчике её узнать нечем: pthread_getattr_np зовёт malloc, а malloc в обработчике нарушает правила async-signal-safe из урока 49. Значит, поток должен записать себя в таблицу сам, когда стартует, до того как начнёт работу. А чтобы программу не пришлось переписывать под профилировщик, libzt-prof.so подменит pthread_create тем же приёмом, каким zt ltrace подменял malloc в уроке про подмену символов.
threads.zig: таблица потоков
//! Таблица потоков профилируемой программы: у каждого потока свой стек и
//! своё имя.
//!
//! Обработчику `SIGPROF` нужна верхняя граница стека того потока, который
//! он прервал: без неё обход кадров не знает, где остановиться. Узнать её
//! в самом обработчике нечем, `pthread_getattr_np` зовёт `malloc`. Поэтому
//! поток записывается в таблицу сам, когда стартует, а обработчик только
//! ищет в ней свой `tid`.
//!
//! Таблица статическая, запись в неё это атомарное присваивание: пишут
//! сразу несколько потоков, а читает обработчик сигнала.
const std = @import("std");
/// Столько потоков помещается в таблицу. Поток сверх этого числа всё равно
/// попадёт в профиль, но его сэмплы будут из одного счётчика команд.
// ponytail: места не переиспользуются. Сервер, который заводит поток на
// каждое соединение, упрётся в предел; тогда освобождай место на выходе
// потока и ищи свободное.
pub const capacity = 256;
/// Имя потока в Linux не длиннее 15 байт плюс нулевой байт, как у `comm`.
pub const name_size = 16;
pub const Thread = struct {
/// Ноль значит, что место ещё не занято.
tid: std.atomic.Value(u32) = .init(0),
/// Первый адрес над стеком потока.
stack_high: usize = 0,
name_buffer: [name_size]u8 = @splat(0),
name_length: u8 = 0,
/// Свой таймер потока, если сэмплер работает с таймером на поток.
timer: ?i32 = null,
pub fn name(thread: *const Thread) []const u8 {
return thread.name_buffer[0..thread.name_length];
}
pub fn setName(thread: *Thread, text: []const u8) void {
const length = @min(text.len, name_size);
@memcpy(thread.name_buffer[0..length], text[0..length]);
thread.name_length = @intCast(length);
}
};
pub const Table = struct {
threads: [capacity]Thread = @splat(.{}),
/// Сколько мест выдано, может перевалить за `capacity`.
used: std.atomic.Value(usize) = .init(0),
/// Записывает поток. Сначала граница стека, потом `tid` с release:
/// обработчик, который увидел `tid`, видит и границу.
pub fn register(table: *Table, tid: u32, stack_high: usize) ?*Thread {
const index = table.used.fetchAdd(1, .monotonic);
if (index >= capacity) return null;
const thread = &table.threads[index];
thread.stack_high = stack_high;
thread.tid.store(tid, .release);
return thread;
}
/// Поток по `tid`. Ищем с конца: ядро раздаёт номера вышедших потоков
/// заново, и у самой свежей записи верная граница стека.
pub fn find(table: *Table, tid: u32) ?*Thread {
var index = table.count();
while (index > 0) {
index -= 1;
if (table.threads[index].tid.load(.acquire) == tid) return &table.threads[index];
}
return null;
}
/// Граница стека потока или ноль, если поток не записан. Годится для
/// обработчика сигнала: только чтение памяти.
pub fn stackHigh(table: *Table, tid: u32) usize {
const thread = table.find(tid) orelse return 0;
return thread.stack_high;
}
pub fn count(table: *const Table) usize {
return @min(table.used.load(.acquire), capacity);
}
pub fn all(table: *Table) []Thread {
return table.threads[0..table.count()];
}
pub fn clear(table: *Table) void {
for (table.all()) |*thread| thread.* = .{};
table.used.store(0, .release);
}
};
test "поток находится по tid" {
var table: Table = .{};
_ = table.register(100, 0x7000);
const worker = table.register(101, 0x9000).?;
worker.setName("worker-1");
try std.testing.expectEqual(0x9000, table.stackHigh(101));
try std.testing.expectEqual(0x7000, table.stackHigh(100));
try std.testing.expectEqual(0, table.stackHigh(102));
try std.testing.expectEqualStrings("worker-1", table.find(101).?.name());
}
test "номер вышедшего потока достался новому: берётся свежая запись" {
var table: Table = .{};
_ = table.register(200, 0x1000);
_ = table.register(200, 0x2000);
try std.testing.expectEqual(0x2000, table.stackHigh(200));
}
test "длинное имя обрезается, лишние потоки не пишутся" {
var table: Table = .{};
const first = table.register(1, 0x10).?;
first.setName("a-very-long-thread-name");
try std.testing.expectEqual(name_size, first.name().len);
for (1..capacity) |i| _ = table.register(@intCast(i + 1), 0x10);
try std.testing.expectEqual(null, table.register(9999, 0x10));
try std.testing.expectEqual(capacity, table.count());
try std.testing.expectEqual(0, table.stackHigh(9999));
}
Таблица статическая, как кольцо сэмплов: обработчик сигнала не может ничего выделять. Записывают в неё потоки при старте, и записывают одновременно, а читает обработчик сигнала на любом из потоков. Поэтому все поля, которые делятся между потоками, атомарные:
used.fetchAdd(1, .monotonic)раздаёт места. Два потока, которые стартуют одновременно, получат разные номера мест, даже если их инструкции выполняются в одну и ту же наносекунду на разных ядрах. Обычноеused += 1это чтение, сложение и запись, и два потока могли бы прочитать одно и то же значение.registerсначала пишетstack_high, потомtidс порядком.release, аfindчитаетtidс.acquire. Эта пара гарантирует: если обработчик увидел в записи нужныйtid, он видит и границу стека, записанную до него. Без неё процессор или компилятор вправе переставить две записи, и обработчик получил бы поток с нулевой границей.
Что такое порядки памяти и почему без них перестановка законна, подробно разберём в уроке про разделяемые переменные, а с точки зрения Rust это урок про атомики. Сейчас хватит правила: данные пишутся до флага с .release, читаются после флага с .acquire.
find ищет с конца. Ядро раздаёт номера потоков заново: поток 361 вышел, и следующий новый поток вполне может получить тот же 361. У самой свежей записи верная граница стека.
ring.zig: два обработчика одновременно
В однопоточной программе SIGPROF не мог прийти, пока предыдущий обработчик работает: сигнал заблокирован на время своего обработчика. Теперь потоков несколько, и обработчики на разных ядрах работают одновременно. Кольцо из урока 49 двигало счётчик written += 1 последней строкой, и два обработчика могли записать сэмпл в одно и то же место, а счётчик сдвинулся бы на один. Это первая настоящая гонка в нашем собственном коде, и лечится она так же, как в таблице: место берётся атомарным сложением.
//! по `dropped`.
+//!
+//! Писателей бывает несколько: в многопоточной программе обработчики
+//! `SIGPROF` на разных ядрах работают одновременно. Место под сэмпл каждый
+//! берёт атомарным `fetchAdd`, так что два обработчика не пишут в одно место.
@@
pub const Sample = struct {
+ /// Номер потока, который прервал таймер. Ноль, если он неизвестен.
+ tid: u32,
depth: u32,
@@
- /// Только присваивания и копирование: годится для обработчика сигнала.
- pub fn push(ring: *Self, addresses: []const u64) void {
+ /// Только атомарное сложение и копирование: годится для обработчика сигнала.
+ // ponytail: счётчик растёт раньше, чем сэмпл записан. Читатель,
+ // который читает кольцо на ходу, может застать сэмпл недописанным;
+ // сэмплер читает его после остановки таймера.
+ pub fn push(ring: *Self, tid: u32, addresses: []const u64) void {
const depth = @min(addresses.len, max_depth);
- const slot = &ring.samples[@intCast(ring.written % capacity)];
+ const number = @atomicRmw(u64, &ring.written, .Add, 1, .monotonic);
+ const slot = &ring.samples[@intCast(number % capacity)];
@memcpy(slot.addresses[0..depth], addresses[0..depth]);
+ slot.tid = tid;
slot.depth = @intCast(depth);
- // Счётчик двигаем последним: сэмпл виден читателю уже целым.
- ring.written += 1;
}
@@
var ring: Ring(4) = .{};
- ring.push(&.{ 0x10, 0x20 });
- ring.push(&.{0x30});
+ ring.push(7, &.{ 0x10, 0x20 });
+ ring.push(8, &.{0x30});
try std.testing.expectEqual(2, ring.count());
@@
try std.testing.expectEqualSlices(u64, &.{0x30}, ring.get(1).stack());
+ try std.testing.expectEqual(8, ring.get(1).tid);
}
@@
var ring: Ring(4) = .{};
- for (1..7) |i| ring.push(&.{i});
+ for (1..7) |i| ring.push(0, &.{i});
try std.testing.expectEqual(4, ring.count());
@@
for (&deep, 0..) |*address, i| address.* = i;
- ring.push(&deep);
+ ring.push(0, &deep);
try std.testing.expectEqual(max_depth, ring.get(0).depth);
@atomicRmw(u64, &ring.written, .Add, 1, .monotonic) возвращает старое значение и увеличивает счётчик одной неделимой инструкцией. У каждого обработчика свой номер и своё место. Цена этого в заметке: счётчик растёт раньше, чем сэмпл записан. Читатель, который читал бы кольцо на ходу, мог бы застать место занятым, но пустым. Наш читатель, .fini_array, читает кольцо после остановки таймера, так что это ограничение нас не касается, но записано честно.
sampler.zig
//! ни форматирования, ни вывода.
+//!
+//! `ITIMER_PROF` один на процесс, а сигнал от него ядро отдаёт тому потоку,
+//! на чьём такте таймер истёк, то есть тому, кто жёг процессор. Поэтому
+//! обработчик первым делом спрашивает `gettid` и ищет свой поток в таблице
+//! (`threads.zig`): там граница его стека и потом его имя.
+//!
+//! Таймер процесса проверяется на тиках планировщика и за проверку отдаёт
+//! не больше одного сигнала. Когда потоков на ядрах несколько, сэмплы
+//! делятся между ними пропорционально, но неточно (замеры в README). Для
+//! точного деления есть второй режим, `Clock.thread`: у каждого потока свой
+//! таймер `timer_create(CLOCK_THREAD_CPUTIME_ID)`, и ядро адресует сигнал
+//! именно ему через `SIGEV_THREAD_ID`. Только Linux.
@@
const ring_mod = @import("ring.zig");
+const threads = @import("threads.zig");
@@
var ring: Samples = .{};
+var table: threads.Table = .{};
var current_hz: u32 = 0;
-/// Верхняя граница стека потока, который вызвал `start`.
-var stack_high: usize = 0;
+var current_clock: Clock = .process;
+
+/// Чьё процессорное время отмеряет таймер.
+pub const Clock = enum {
+ /// Один `setitimer(ITIMER_PROF)` на процесс.
+ process,
+ /// Свой таймер у каждого записанного потока.
+ thread,
+};
@@
-/// Включает сэмплирование. Стеки собираются у потока, который вызвал
-/// `start`: у сэмпла с другого потока останется один счётчик команд.
-// ponytail: один поток. Границы стека на каждый поток и его имя в префиксе
-// стека это `zt prof --threads` из урока про потоки.
-pub fn start(hz: u32) Error!void {
+/// Включает сэмплирование и записывает в таблицу поток, который вызвал
+/// `start`. Остальные потоки записываются сами через `registerThread`:
+/// у сэмпла с незаписанного потока останется один счётчик команд.
+pub fn start(hz: u32, clock: Clock) Error!void {
if (hz == 0 or hz > 10_000) return error.BadFrequency;
- stack_high = stackHigh(@frameAddress()) orelse return error.NoStackBounds;
+ if (clock == .thread and builtin.os.tag != .linux) return error.TimerFailed;
+ const high = stackHigh(@frameAddress()) orelse return error.NoStackBounds;
ring.clear();
+ table.clear();
current_hz = hz;
+ current_clock = clock;
@@
posix.sigaction(.PROF, &action, null);
- try setTimer(1_000_000 / hz);
+ const main = table.register(currentTid(), high);
+ switch (clock) {
+ .process => try setTimer(1_000_000 / hz),
+ .thread => if (main) |entry| {
+ entry.timer = try threadTimer(hz);
+ },
+ }
}
-/// Выключает таймер. Обработчик остаётся на месте: сигнал, который ядро
+/// Выключает таймеры. Обработчик остаётся на месте: сигнал, который ядро
/// уже отправило, иначе убил бы программу.
pub fn stop() void {
- setTimer(0) catch {};
+ switch (current_clock) {
+ .process => setTimer(0) catch {},
+ .thread => for (table.all()) |*thread| deleteTimer(thread),
+ }
}
@@
+/// Записывает текущий поток и в режиме `Clock.thread` заводит ему таймер.
+/// Зовётся в самом потоке, до его работы, и не из обработчика сигнала.
+pub fn registerThread(stack_high: usize) ?*threads.Thread {
+ const thread = table.register(currentTid(), stack_high) orelse return null;
+ if (current_clock == .thread and current_hz != 0) thread.timer = threadTimer(current_hz) catch null;
+ return thread;
+}
+
+/// Поток закончил работу: запоминаем имя и снимаем его таймер.
+pub fn threadExited(thread: *threads.Thread, name: []const u8) void {
+ thread.setName(name);
+ deleteTimer(thread);
+}
+
+pub fn threadTable() *threads.Table {
+ return &table;
+}
+
pub fn frequency() u32 {
@@
var addresses: [ring_mod.max_depth]u64 = undefined;
- // Обработчик работает на том же стеке, ниже прерванного кода. Значит,
- // его собственный кадр это нижняя граница для кадров программы.
- const depth = frames.walk(context.getPc(), context.getFp(), @frameAddress(), stack_high, &addresses);
- ring.push(addresses[0..depth]);
+ const tid = currentTid();
+ // Обработчик работает на стеке прерванного потока, ниже его кода. Значит,
+ // собственный кадр обработчика это нижняя граница для кадров программы.
+ // У незаписанного потока верхняя граница ноль, и обход не начнётся.
+ const depth = frames.walk(context.getPc(), context.getFp(), @frameAddress(), table.stackHigh(tid), &addresses);
+ ring.push(tid, addresses[0..depth]);
+}
+
+/// Номер потока в ядре. `gettid` это системный вызов без побочных
+/// эффектов, его можно звать из обработчика.
+pub fn currentTid() u32 {
+ return switch (builtin.os.tag) {
+ .linux => @bitCast(linux.gettid()),
+ .macos => macos: {
+ var id: u64 = 0;
+ _ = std.c.pthread_threadid_np(null, &id);
+ break :macos @truncate(id);
+ },
+ else => 0,
+ };
}
@@
+/// `struct sigevent` из ядра: как сообщить об истечении таймера. Нам нужен
+/// вариант `SIGEV_THREAD_ID`: сигнал `signo` получает поток `tid`.
+const Sigevent = extern struct {
+ value: usize = 0,
+ signo: i32,
+ notify: i32,
+ tid: i32,
+ padding: [44]u8 = @splat(0),
+};
+
+const sigev_thread_id = 4;
+const clock_thread_cputime_id = 3;
+
+/// Таймер процессорного времени текущего потока, сигнал приходит ему же.
+fn threadTimer(hz: u32) Error!i32 {
+ var event: Sigevent = .{ .signo = @intFromEnum(posix.SIG.PROF), .notify = sigev_thread_id, .tid = linux.gettid() };
+ var timer: i32 = 0;
+ if (linux.errno(linux.syscall3(.timer_create, clock_thread_cputime_id, @intFromPtr(&event), @intFromPtr(&timer))) != .SUCCESS) return error.TimerFailed;
+ const period: linux.timespec = .{ .sec = 0, .nsec = @intCast(1_000_000_000 / hz) };
+ const value: linux.itimerspec = .{ .it_interval = period, .it_value = period };
+ if (linux.errno(linux.syscall4(.timer_settime, @intCast(timer), 0, @intFromPtr(&value), 0)) != .SUCCESS) {
+ _ = linux.syscall1(.timer_delete, @intCast(timer));
+ return error.TimerFailed;
+ }
+ return timer;
+}
+
+fn deleteTimer(thread: *threads.Thread) void {
+ const timer = thread.timer orelse return;
+ thread.timer = null;
+ if (builtin.os.tag == .linux) _ = linux.syscall1(.timer_delete, @intCast(timer));
+}
+
extern "c" fn pthread_self() ?*anyopaque;
Главное в обработчике. currentTid() на Linux это gettid, на macOS pthread_threadid_np (сэмплер на macOS не собирается, но таблица и её тесты должны компилироваться везде). table.stackHigh(tid) отдаёт ноль для потока, которого нет в таблице, и тогда frames.walk сразу остановится: в сэмпле останется только счётчик команд. Лучше короткий честный стек, чем обход по чужой границе.
start записывает в таблицу главный поток, остальные записываются сами через registerThread. threadExited запоминает имя потока и снимает его таймер.
Второй режим, Clock.thread, нужен для точности. Общий таймер процесса проверяется на тиках и за одну проверку отдаёт не больше одного сигнала. Когда на ядрах работают три потока сразу, их процессорное время на тике складывается, а сигнал всё равно один: сэмплов меньше, чем должно быть, и недосчитывается сильнее тот, кто работает только вместе с остальными. timer_create(CLOCK_THREAD_CPUTIME_ID) заводит таймер процессорного времени одного потока, а SIGEV_THREAD_ID велит ядру слать сигнал этому потоку и никому другому. У каждого свой таймер и свои сэмплы. std.os.linux в 0.16 обёрток для timer_create не даёт, поэтому threadTimer зовёт системные вызовы через syscall3 и syscall4 и сам описывает struct sigevent из ядра. Это только Linux.
preload.zig: подмена pthread_create
//!
-//! Ни одной функции программы она не подменяет. Ей нужны два момента:
-//! до `main`, чтобы завести таймер, и после него, чтобы сбросить сэмплы в
-//! файл. Оба даёт сам формат ELF: загрузчик вызывает функции из секции
-//! `.init_array` каждой библиотеки до передачи управления программе, а
-//! функции из `.fini_array` на выходе через `exit` или возврат из `main`.
+//! Ей нужны два момента: до `main`, чтобы завести таймер, и после него,
+//! чтобы сбросить сэмплы в файл. Оба даёт сам формат ELF: загрузчик вызывает
+//! функции из секции `.init_array` каждой библиотеки до передачи управления
+//! программе, а функции из `.fini_array` на выходе через `exit` или возврат
+//! из `main`.
+//!
+//! Подменяет она одну функцию, `pthread_create`, как `zt ltrace` подменяет
+//! `malloc`. Новый поток начинает не с функции программы, а с нашей
+//! `zt_thread_start`: она записывает границу стека потока в таблицу
+//! сэмплера, зовёт настоящую функцию, а на выходе запоминает имя потока.
+//! Имя к этому моменту программа уже поставила через `pthread_setname_np`,
+//! а после `pthread_join` каталога `/proc/self/task/<tid>` уже нет.
//!
@@
extern "c" fn getenv(name: [*:0]const u8) ?[*:0]const u8;
+extern "c" fn dlsym(handle: ?*anyopaque, name: [*:0]const u8) ?*anyopaque;
+extern "c" fn pthread_getattr_np(thread: std.c.pthread_t, attr: *std.c.pthread_attr_t) c_int;
+extern "c" fn pthread_attr_getstack(attr: *const std.c.pthread_attr_t, low: *usize, size: *usize) c_int;
+extern "c" fn pthread_attr_destroy(attr: *std.c.pthread_attr_t) c_int;
+
+/// Особое значение вместо описателя библиотеки: искать в следующих по порядку.
+const rtld_next: ?*anyopaque = @ptrFromInt(@as(usize, @bitCast(@as(isize, -1))));
+
+const Start = *const fn (?*anyopaque) callconv(.c) ?*anyopaque;
+const Create = fn (*std.c.pthread_t, ?*const std.c.pthread_attr_t, Start, ?*anyopaque) callconv(.c) c_int;
+
+/// Что передать новому потоку: настоящую функцию и её аргумент.
+const Launch = struct {
+ start: Start,
+ argument: ?*anyopaque,
+};
+
+var real_create: ?*const Create = null;
+
+export fn pthread_create(
+ thread: *std.c.pthread_t,
+ attr: ?*const std.c.pthread_attr_t,
+ start: Start,
+ argument: ?*anyopaque,
+) c_int {
+ const create = real_create orelse create: {
+ const found = dlsym(rtld_next, "pthread_create") orelse return @intFromEnum(linux.E.AGAIN);
+ real_create = @ptrCast(@alignCast(found));
+ break :create real_create.?;
+ };
+ // Пара живёт в куче: стек того, кто создаёт поток, может уйти раньше,
+ // чем новый поток её прочтёт. Освобождает её сам новый поток.
+ const launch: *Launch = @ptrCast(@alignCast(std.c.malloc(@sizeOf(Launch)) orelse return @intFromEnum(linux.E.AGAIN)));
+ launch.* = .{ .start = start, .argument = argument };
+ const result = create(thread, attr, &zt_thread_start, launch);
+ if (result != 0) std.c.free(launch);
+ return result;
+}
+
+/// Первая функция каждого нового потока. Экспортируется, чтобы в профиле
+/// у неё было понятное имя, а не внутреннее имя Zig.
+export fn zt_thread_start(raw: ?*anyopaque) ?*anyopaque {
+ const launch: *Launch = @ptrCast(@alignCast(raw.?));
+ const start = launch.start;
+ const argument = launch.argument;
+ std.c.free(launch);
+
+ const thread = prof.sampler.registerThread(stackHigh());
+ const result = start(argument);
+ // ponytail: поток, который вышел через pthread_exit, сюда не вернётся,
+ // и его имя не запомнится: в профиле он будет [tid N]. Нужны имена и
+ // таких, перехватывай ещё pthread_setname_np.
+ if (thread) |entry| {
+ var name: [prof.threads.name_size]u8 = @splat(0);
+ _ = linux.prctl(@intFromEnum(linux.PR.GET_NAME), @intFromPtr(&name), 0, 0, 0);
+ prof.sampler.threadExited(entry, std.mem.sliceTo(&name, 0));
+ }
+ return result;
+}
+
+/// Первый адрес над стеком текущего потока.
+fn stackHigh() usize {
+ var attr: std.c.pthread_attr_t = undefined;
+ if (pthread_getattr_np(std.c.pthread_self(), &attr) != 0) return 0;
+ defer _ = pthread_attr_destroy(&attr);
+ var low: usize = 0;
+ var size: usize = 0;
+ if (pthread_attr_getstack(&attr, &low, &size) != 0) return 0;
+ return low + size;
+}
@@
if (getenv(prof.frequency_variable)) |text| hz = std.fmt.parseInt(u32, std.mem.span(text), 10) catch hz;
- prof.sampler.start(hz) catch {};
+ const clock: prof.sampler.Clock = if (getenv(prof.clock_variable)) |text|
+ std.meta.stringToEnum(prof.sampler.Clock, std.mem.span(text)) orelse .process
+ else
+ .process;
+ prof.sampler.start(hz, clock) catch {};
}
@@
writeAll(fd, prof.sampler.readSelfMaps());
+ writeAll(fd, "threads\n");
+ for (prof.sampler.threadTable().all()) |*thread| {
+ const tid = thread.tid.load(.acquire);
+ // Живой поток мог переименоваться, у вышедшего имя уже записано.
+ var comm: [prof.threads.name_size + 1]u8 = undefined;
+ const name = liveName(tid, &comm) orelse thread.name();
+ writeAll(fd, std.fmt.bufPrint(&buffer, "{d} {s}\n", .{ tid, name }) catch continue);
+ }
writeAll(fd, "samples\n");
for (0..samples.count()) |index| {
- writeAll(fd, prof.dump.formatSample(&buffer, samples.get(index).stack()));
+ const sample = samples.get(index);
+ writeAll(fd, prof.dump.formatSample(&buffer, sample.tid, sample.stack()));
}
@@
+/// Имя живого потока из `/proc/self/task/<tid>/comm`, без перевода строки.
+fn liveName(tid: u32, buffer: []u8) ?[]const u8 {
+ var path_buffer: [64]u8 = undefined;
+ const path = std.fmt.bufPrintZ(&path_buffer, "/proc/self/task/{d}/comm", .{tid}) catch return null;
+ const opened = linux.open(path, .{ .ACCMODE = .RDONLY }, 0);
+ if (linux.errno(opened) != .SUCCESS) return null;
+ const fd: i32 = @intCast(opened);
+ defer _ = linux.close(fd);
+ const got = linux.read(fd, buffer.ptr, buffer.len);
+ if (linux.errno(got) != .SUCCESS or got == 0) return null;
+ return std.mem.trimEnd(u8, buffer[0..got], "\n");
+}
+
fn writeAll(fd: i32, bytes: []const u8) void {
Путь нового потока теперь такой.
- Программа зовёт
pthread_create(&t, attr, run, arg). Динамический загрузчик находит этот символ первым вlibzt-prof.so, потому что она подгружена черезLD_PRELOAD. - Наша
pthread_createодин раз находит настоящую черезdlsym(RTLD_NEXT, "pthread_create"), кладёт пару «функция, аргумент» в кучу и зовёт настоящую, но стартовой функцией нового потока ставитzt_thread_start. - Новый поток начинает с
zt_thread_start. Она забирает пару, освобождает память, узнаёт границу своего стека черезpthread_getattr_npиpthread_attr_getstackи записывается в таблицу. Здесьmallocможно: мы не в обработчике сигнала, а в обычном коде потока. - Потом зовёт настоящую функцию программы.
- Когда та вернулась, берёт имя потока через
prctl(PR_GET_NAME)и отдаёт его вthreadExited.
Почему пара лежит в куче, а не на стеке pthread_create? Ровно по правилу из раздела про гонку: наша pthread_create вернётся раньше, чем новый поток прочтёт аргумент, и её стек к тому времени будет занят другими вызовами. Это тот же приём, что у std.Thread.spawn внутри, только руками.
Почему имя берётся на выходе, а не при запуске потока? Программа ставит имя изнутри потока, pthread_setname_np в начале функции потока, то есть уже после zt_thread_start. А после того как поток вышел и его дождались через pthread_join, каталога /proc/self/task/<tid> уже нет, и имя взять неоткуда. Поэтому его снимает сам поток на выходе. Потокам, живым в момент дампа, имя берётся из /proc/self/task/<tid>/comm: fini работает в обычном коде, не в обработчике, так что открыть файл можно.
В дамп добавляется часть threads: номер и имя каждого потока из таблицы, а у каждого сэмпла впереди номер потока с двоеточием.
dump.zig, folded.zig, symbolize.zig
Формат дампа растёт, но старый дамп тоже должен читаться: у его сэмплов номер потока ноль.
//!
-//! Это текст из двух частей. Карта памяти, копия `/proc/self/maps`, и сэмплы,
-//! по строке на сэмпл: адреса в hex через пробел, от прерванной инструкции
-//! к `main`. Имён функций тут нет: искать их внутри программы, да ещё на
-//! выходе из неё, незачем. Адреса превратит в имена `zt prof`, а карта
-//! нужна ему затем, что при каждом запуске библиотеки лежат по новым адресам.
+//! Это текст из трёх частей. Карта памяти, копия `/proc/self/maps`. Потоки,
+//! по строке на поток: номер и имя. И сэмплы, по строке на сэмпл: номер
+//! потока с двоеточием, потом адреса в hex через пробел, от прерванной
+//! инструкции к корню. Имён функций тут нет: искать их внутри программы, да
+//! ещё на выходе из неё, незачем. Адреса превратит в имена `zt prof`, а
+//! карта нужна ему затем, что при каждом запуске библиотеки лежат по новым
+//! адресам.
//!
@@
//! 00400000-00401000 r--p 00000000 08:01 42 /tmp/spin
+//! threads
+//! 4242 spin
//! samples
-//! 401136 40117b 4011a2
+//! 4242: 401136 40117b 4011a2
+//!
+//! Части `threads` и номера потоков появились в уроке про потоки. Дамп без
+//! них тоже читается: у его сэмплов номер потока ноль.
@@
/// Строка дампа для одного сэмпла. Пишет в готовый буфер, без аллокаций.
-pub fn formatSample(buffer: []u8, stack: []const u64) []const u8 {
+pub fn formatSample(buffer: []u8, tid: u32, stack: []const u64) []const u8 {
var out: std.Io.Writer = .fixed(buffer);
- for (stack, 0..) |address, position| {
- if (position > 0) out.writeByte(' ') catch break;
- out.print("{x}", .{address}) catch break;
- }
+ out.print("{d}:", .{tid}) catch {};
+ for (stack) |address| out.print(" {x}", .{address}) catch break;
out.writeByte('\n') catch {};
@@
-/// Места хватит на самый глубокий стек: 16 цифр и разделитель на адрес.
-pub const sample_line_size = ring_mod.max_depth * 17 + 1;
+/// Места ровно на самый длинный сэмпл: номер потока с двоеточием, по
+/// 16 цифр с разделителем на адрес и перевод строки.
+pub const sample_line_size = 11 + ring_mod.max_depth * 17 + 1;
+
+pub const Thread = struct {
+ tid: u32,
+ /// Срез смотрит в исходный текст.
+ name: []const u8,
+};
@@
regions: []pmap.Region,
+ threads: []Thread,
/// Стеки лежат подряд в `addresses`, у каждого свой срез.
stacks: [][]const u64,
+ /// Номер потока для каждого стека, по тем же индексам.
+ tids: []u32,
addresses: []u64,
@@
gpa.free(dump.regions);
+ gpa.free(dump.threads);
gpa.free(dump.stacks);
+ gpa.free(dump.tids);
gpa.free(dump.addresses);
}
+
+ /// Имя потока по номеру, если поток есть в дампе.
+ pub fn threadName(dump: Dump, tid: u32) ?[]const u8 {
+ for (dump.threads) |thread| {
+ if (thread.tid == tid) return thread.name;
+ }
+ return null;
+ }
};
@@
errdefer regions.deinit(gpa);
+ var threads: std.ArrayList(Thread) = .empty;
+ errdefer threads.deinit(gpa);
+ var tids: std.ArrayList(u32) = .empty;
+ errdefer tids.deinit(gpa);
var addresses: std.ArrayList(u64) = .empty;
@@
- const Part = enum { header, maps, samples };
+ const Part = enum { header, maps, threads, samples };
var part: Part = .header;
@@
part = .maps;
+ } else if (std.mem.eql(u8, line, "threads")) {
+ part = .threads;
} else if (std.mem.eql(u8, line, "samples")) {
@@
.maps => try regions.append(gpa, try pmap.parseLine(line)),
+ // Имя это остаток строки: в имени потока пробел законен.
+ .threads => {
+ const space = std.mem.indexOfScalar(u8, line, ' ') orelse line.len;
+ const tid = std.fmt.parseInt(u32, line[0..space], 10) catch return error.BadSample;
+ try threads.append(gpa, .{ .tid = tid, .name = if (space < line.len) line[space + 1 ..] else "" });
+ },
.samples => {
@@
var words = std.mem.tokenizeScalar(u8, line, ' ');
+ var tid: u32 = 0;
+ const first = words.peek() orelse "";
+ if (cutSuffix(first, ":")) |number| {
+ tid = std.fmt.parseInt(u32, number, 10) catch return error.BadSample;
+ _ = words.next();
+ }
+ try tids.append(gpa, tid);
while (words.next()) |word| {
@@
const stacks = try gpa.alloc([]const u64, lengths.items.len);
+ errdefer gpa.free(stacks);
var start: usize = 0;
@@
}
+ const owned_regions = try regions.toOwnedSlice(gpa);
+ errdefer gpa.free(owned_regions);
+ const owned_threads = try threads.toOwnedSlice(gpa);
+ errdefer gpa.free(owned_threads);
return .{
@@
.dropped = dropped,
- .regions = try regions.toOwnedSlice(gpa),
+ .regions = owned_regions,
+ .threads = owned_threads,
.stacks = stacks,
+ .tids = try tids.toOwnedSlice(gpa),
.addresses = all,
@@
+fn cutSuffix(word: []const u8, suffix: []const u8) ?[]const u8 {
+ return if (std.mem.endsWith(u8, word, suffix)) word[0 .. word.len - suffix.len] else null;
+}
+
test "строка сэмпла" {
var buffer: [sample_line_size]u8 = undefined;
- try std.testing.expectEqualStrings("401136 40117b 4011a2\n", formatSample(&buffer, &.{ 0x401136, 0x40117b, 0x4011a2 }));
+ try std.testing.expectEqualStrings("4242: 401136 40117b 4011a2\n", formatSample(&buffer, 4242, &.{ 0x401136, 0x40117b, 0x4011a2 }));
const deepest = [_]u64{std.math.maxInt(u64)} ** ring_mod.max_depth;
- try std.testing.expectEqual(sample_line_size - 1, formatSample(&buffer, &deepest).len);
+ try std.testing.expectEqual(sample_line_size, formatSample(&buffer, std.math.maxInt(u32), &deepest).len);
}
@@
try std.testing.expectEqualSlices(u64, &.{0x401150}, dump.stacks[1]);
+ // Дамп до урока про потоки: номеров нет, у всех сэмплов ноль.
+ try std.testing.expectEqualSlices(u32, &.{ 0, 0 }, dump.tids);
+}
+
+test "потоки в дампе" {
+ const gpa = std.testing.allocator;
+ const text = magic ++
+ \\threads
+ \\4242 pool
+ \\4243 worker 1
+ \\samples
+ \\4243: 401136 40117b
+ \\4242: 401150
+ \\
+ ;
+ const dump = try parse(gpa, text);
+ defer dump.deinit(gpa);
+
+ try std.testing.expectEqualSlices(u32, &.{ 4243, 4242 }, dump.tids);
+ try std.testing.expectEqualSlices(u64, &.{ 0x401136, 0x40117b }, dump.stacks[0]);
+ try std.testing.expectEqualStrings("worker 1", dump.threadName(4243).?);
+ try std.testing.expectEqualStrings("pool", dump.threadName(4242).?);
+ try std.testing.expectEqual(null, dump.threadName(1));
}
sample_line_size посчитан заново: у самой длинной строки теперь впереди до десяти цифр номера и двоеточие. Буфер для строки дампа выделен заранее, и если бы размер не сошёлся, самый глубокий стек молча обрезался бы. Тест «строка сэмпла» проверяет ровно это: самый длинный номер и самый глубокий стек занимают буфер целиком.
Свёрнутый стек получает необязательный корень:
//! main;work;hot_b 104
+//!
+//! С `zt prof --threads` у стека появляется корень над всеми: имя потока.
+//! Тогда у flame graph столько подножий, сколько потоков, и видно, кто
+//! сколько жёг.
+//!
+//! heavy;start_thread;run;spin 290
+//! light;start_thread;run;spin 98
@@
- /// `stack` идёт от листа к корню, как его записал сэмплер.
- pub fn add(folded: *Folded, gpa: std.mem.Allocator, symbolizer: *const symbolize.Symbolizer, stack: []const u64) !void {
+ /// `stack` идёт от листа к корню, как его записал сэмплер. `root`, если
+ /// он есть, встаёт первым кадром, например имя потока.
+ pub fn add(
+ folded: *Folded,
+ gpa: std.mem.Allocator,
+ symbolizer: *const symbolize.Symbolizer,
+ root: ?[]const u8,
+ stack: []const u64,
+ ) !void {
var line: std.ArrayList(u8) = .empty;
defer line.deinit(gpa);
+ if (root) |name| try line.appendSlice(gpa, name);
И одна правка символизации, которую подсказал первый же профиль с потоками. В .dynsym glibc у многих функций два имени на одном адресе, например __pthread_mutex_lock и pthread_mutex_lock. Какое из них попадёт в профиль, решал порядок сортировки. Теперь среди псевдонимов с одним адресом последним встаёт имя с наименьшим числом подчёркиваний в начале, а поиск берёт последнее:
+ /// По адресу, а у псевдонимов с одним адресом (`__pthread_mutex_lock` и
+ /// `pthread_mutex_lock`) последним встаёт имя без подчёркиваний: поиск
+ /// берёт последнюю функцию с подходящим адресом.
fn before(_: void, a: Function, b: Function) bool {
- return a.address < b.address;
+ if (a.address != b.address) return a.address < b.address;
+ return underscores(a.name) > underscores(b.name);
+ }
+
+ fn underscores(name: []const u8) usize {
+ return std.mem.indexOfNone(u8, name, "_") orelse name.len;
}
prof.zig и main.zig
report получает настройки, а корень стека берётся из дампа по номеру потока. Поток без имени становится [tid N], сэмпл без номера (старый дамп) остаётся без корня:
pub const ring = @import("prof/ring.zig");
+pub const threads = @import("prof/threads.zig");
pub const frames = @import("prof/frames.zig");
@@
pub const frequency_variable = "ZT_PROF_HZ";
+/// `process` или `thread`, см. `sampler.Clock`.
+pub const clock_variable = "ZT_PROF_CLOCK";
@@
hz: u32,
+ clock: sampler.Clock,
) !void {
@@
try environment.put(frequency_variable, text);
+ try environment.put(clock_variable, @tagName(clock));
}
@@
+pub const Options = struct {
+ /// Первым кадром каждого стека ставить имя потока.
+ threads: bool = false,
+};
+
/// Сырой дамп в свёрнутые стеки.
-pub fn report(gpa: std.mem.Allocator, out: *std.Io.Writer, text: []const u8, loader: Loader) !Stats {
+pub fn report(gpa: std.mem.Allocator, out: *std.Io.Writer, text: []const u8, loader: Loader, options: Options) !Stats {
// Таблицы символов смотрят в байты файлов, поэтому файлы живут в арене
@@
defer result.deinit(gpa);
- for (parsed.stacks) |stack| try result.add(gpa, &symbolizer, stack);
+ for (parsed.stacks, parsed.tids) |stack, tid| {
+ var buffer: [32]u8 = undefined;
+ const root = if (options.threads) threadRoot(parsed, &buffer, tid) else null;
+ try result.add(gpa, &symbolizer, root, stack);
+ }
try result.print(out);
@@
+/// Имя потока для корня стека. Поток без имени в дампе остаётся номером,
+/// дамп без номеров остаётся без корня.
+fn threadRoot(parsed: dump.Dump, buffer: []u8, tid: u32) ?[]const u8 {
+ if (tid == 0) return null;
+ if (parsed.threadName(tid)) |name| {
+ if (name.len > 0) return name;
+ }
+ return std.fmt.bufPrint(buffer, "[tid {d}]", .{tid}) catch null;
+}
+
test "библиотека ищется рядом с программой" {
@@
_ = ring;
+ _ = threads;
_ = frames;
В main.zig два флага: --threads включает корень с именем потока, --thread-timers переключает сэмплер на таймер у каждого потока. Режим таймера уходит программе через переменную окружения ZT_PROF_CLOCK, как частота через ZT_PROF_HZ:
\\ zt pmap <pid>
- \\ zt prof [--hz=N] [--raw=файл] [-o файл] <программа> [аргументы...]
+ \\ zt prof [--threads] [--thread-timers] [--hz=N] [--raw=файл] [-o файл] <программа> [аргументы...]
\\
@@
\\Флаги prof (только Linux):
+ \\ --threads первым кадром каждого стека имя потока
+ \\ --thread-timers свой таймер у каждого потока вместо общего на процесс
\\ --hz=N сэмплов в секунду процессорного времени, по умолчанию 997
@@
var hz: u32 = zt.prof.sampler.default_hz;
+ var options: zt.prof.Options = .{};
+ var clock: zt.prof.sampler.Clock = .process;
var raw_path: ?[]const u8 = null;
@@
if (hz == 0 or hz > 10_000) return fail(out, "zt prof: --hz ждёт число от 1 до 10000\n");
+ } else if (std.mem.eql(u8, arg, "--threads")) {
+ options.threads = true;
+ } else if (std.mem.eql(u8, arg, "--thread-timers")) {
+ clock = .thread;
} else if (std.mem.startsWith(u8, arg, "--raw=")) {
@@
var environment = try init.environ_map.clone(gpa);
- try zt.prof.prepareEnvironment(gpa, &environment, library, dump_path, hz);
+ try zt.prof.prepareEnvironment(gpa, &environment, library, dump_path, hz, clock);
@@
var writer = file.writerStreaming(init.io, &file_buf);
- const stats = try zt.prof.report(gpa, &writer.interface, text, loader);
+ const stats = try zt.prof.report(gpa, &writer.interface, text, loader, options);
try writer.interface.flush();
break :stats stats;
- } else try zt.prof.report(gpa, out, text, loader);
+ } else try zt.prof.report(gpa, out, text, loader, options);
try out.flush();
Фикстура и build.zig
Подопытная программа на C, чтобы стеки не зависели от того, как Zig называет свои функции. Три потока с именами и нагрузкой 3 к 2 к 1. Каждый поток на выходе меряет своё процессорное время через CLOCK_THREAD_CPUTIME_ID: с этими числами тест сверит, сколько сэмплов досталось каждому.
/* Цель для zt prof --threads: три потока с именами и нагрузкой 3 к 2 к 1.
* Главный поток только ждёт их. Каждый поток на выходе печатает своё
* процессорное время по CLOCK_THREAD_CPUTIME_ID: с ним сверяется, сколько
* сэмплов профилировщик выдал каждому потоку.
* Собирается как spin.c: -O0 и кадр у каждой функции. */
#define _GNU_SOURCE
#include <pthread.h>
#include <stdio.h>
#include <stdlib.h>
#include <time.h>
struct job {
const char *name;
unsigned long rounds;
double cpu_ms;
};
__attribute__((noinline)) void spin(unsigned long rounds) {
volatile unsigned long sink = 0;
for (unsigned long i = 0; i < rounds; i++) sink += i * 3;
}
__attribute__((noinline)) void *run(void *argument) {
struct job *job = argument;
pthread_setname_np(pthread_self(), job->name);
spin(job->rounds);
struct timespec used;
clock_gettime(CLOCK_THREAD_CPUTIME_ID, &used);
job->cpu_ms = used.tv_sec * 1e3 + used.tv_nsec / 1e6;
return NULL;
}
int main(int argc, char **argv) {
/* Миллионов оборотов на единицу нагрузки, по умолчанию 100. */
unsigned long unit = (argc > 1 ? strtoul(argv[1], NULL, 10) : 100) * 1000000UL;
struct job jobs[] = {
{"heavy", 3 * unit, 0},
{"medium", 2 * unit, 0},
{"light", 1 * unit, 0},
};
pthread_t threads[3];
for (int i = 0; i < 3; i++) pthread_create(&threads[i], NULL, run, &jobs[i]);
for (int i = 0; i < 3; i++) pthread_join(threads[i], NULL);
for (int i = 0; i < 3; i++) printf("%s %.0f ms\n", jobs[i].name, jobs[i].cpu_ms);
return 0;
}
noinline и сборка с -O0 и указателями кадров держат в стеке кадры run и spin. Обрати внимание: job передаётся в поток адресом элемента массива jobs на стеке main. Это законно по правилу из раздела про гонку: main не вернётся, пока не дождётся всех трёх потоков, и никто не меняет jobs[i], пока поток его читает. Поле cpu_ms пишет только свой поток, а главный читает его после join.
В build.zig сборка spin из урока 49 становится циклом по фикстурам, а шаг 69 получает свои тесты:
/// соответствует файл `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 };
+const project_steps = [_]u8{ 6, 7, 9, 10, 11, 12, 13, 14, 17, 18, 42, 43, 44, 45, 46, 47, 49, 54, 69 };
@@
- // Цель для сквозного теста `zt prof`. Всегда Debug: это -O0, без
- // встраивания и с кадром у каждой функции.
- const spin = b.addExecutable(.{
- .name = "spin",
- .root_module = b.createModule(.{
- .target = target,
- .optimize = .Debug,
- .link_libc = true,
- .sanitize_c = .off,
- }),
- });
- spin.root_module.addCSourceFile(.{
- .file = b.path("fixtures/prof/spin.c"),
- // Второй флаг нужен листовым функциям: без него на aarch64 лист
- // не заводит кадра, и тот, кто его вызвал, выпадает из стека.
- .flags = &.{ "-fno-omit-frame-pointer", "-mno-omit-leaf-frame-pointer" },
- });
- const install_spin = b.addInstallArtifact(spin, .{ .dest_dir = .{ .override = .{ .custom = "fixtures" } } });
- b.getInstallStep().dependOn(&install_spin.step);
+ // Цели для сквозных тестов `zt prof`: spin для шага 49, threads для
+ // шага 69. Всегда Debug: это -O0, без встраивания и с кадром у каждой
+ // функции.
+ inline for (.{ "spin", "threads" }) |name| {
+ const fixture = b.addExecutable(.{
+ .name = name,
+ .root_module = b.createModule(.{
+ .target = target,
+ .optimize = .Debug,
+ .link_libc = true,
+ .sanitize_c = .off,
+ }),
+ });
+ fixture.root_module.addCSourceFile(.{
+ .file = b.path("fixtures/prof/" ++ name ++ ".c"),
+ // Второй флаг нужен листовым функциям: без него на aarch64 лист
+ // не заводит кадра, и тот, кто его вызвал, выпадает из стека.
+ .flags = &.{ "-fno-omit-frame-pointer", "-mno-omit-leaf-frame-pointer" },
+ });
+ const install = b.addInstallArtifact(fixture, .{ .dest_dir = .{ .override = .{ .custom = "fixtures" } } });
+ b.getInstallStep().dependOn(&install.step);
+ }
}
@@
paths.addOption([]const u8, "hello_dyn", b.pathFromRoot("fixtures/dyn/hello_dyn"));
- paths.addOption([]const u8, "spin", b.getInstallPath(.{ .custom = "fixtures" }, "spin"));
+ inline for (.{ "spin", "threads" }) |name| {
+ paths.addOption([]const u8, name, b.getInstallPath(.{ .custom = "fixtures" }, name));
+ }
@@
// Шаги со сквозными тестами: они запускают собранный `zt`.
- if (number == 46 or number == 47 or number == 49) {
+ if (number == 46 or number == 47 or number == 49 or number == 69) {
step_tests.root_module.addOptions("paths", paths);
Прогон
Та же программа, что в начале шага, теперь с --threads. Контейнер runner-zig:dev (Linux 7.0.14, arm64, 16 vCPU) на Apple M4 Max:
$ zig build
$ zig-out/bin/zt prof --threads zig-out/fixtures/threads
heavy 405 ms
medium 278 ms
light 144 ms
heavy;[libc.so.6];[libc.so.6];zt_thread_start;run;spin 372
light;[libc.so.6];[libc.so.6];zt_thread_start;run;spin 112
medium;[libc.so.6];[libc.so.6];zt_thread_start;run;spin 231
zt prof: сэмплов 715, частота 997 Гц, затёрто 0
У каждого потока своё подножие, и числа идут в пропорции нагрузки. Два безымянных [libc.so.6] под zt_thread_start это start_thread и clone внутри glibc: их нет в .dynsym, и имён взять неоткуда. Сам zt_thread_start в профиле виден под своим именем, потому что экспортирован. Без --threads всё сливается, как раньше:
$ zig-out/bin/zt prof zig-out/fixtures/threads
heavy 392 ms
medium 263 ms
light 133 ms
[libc.so.6];[libc.so.6];zt_thread_start;run;spin 689
memset 1
zt prof: сэмплов 690, частота 997 Гц, затёрто 0
Насколько точно общий таймер делит сэмплы? При 997 Гц на миллисекунду процессорного времени приходится почти один сэмпл. Замеры из README эталона: пять прогонов threads 150, в ячейке отношение числа сэмплов к числу миллисекунд процессора потока.
| Таймер | heavy | medium | light |
|---|---|---|---|
общий ITIMER_PROF | 0,89 до 0,93 | 0,79 до 0,88 | 0,74 до 0,85 |
--thread-timers | 0,85 до 0,91 | 0,79 до 0,89 | 0,82 до 0,98 |
Сняты на той же машине в контейнере arm64, load average от 25 до 57. С общим таймером лёгкий поток недосчитывается сильнее всех. Он работает только тогда, когда на ядрах все трое, и сигналы на тиках теряются чаще всего в это время, а тяжёлый успевает поработать и один. Свой таймер на поток этот перекос убирает. Общие потери в 10 до 15 процентов остаются и там: таймеры процессорного времени проверяются на тех же тиках.
Под x86-64 (контейнер runner-zig:dev-amd64 под Rosetta) стеки на кадр короче, потому что libc собрана без указателей кадров, и один безымянный кадр под zt_thread_start выпадает. Пропорции те же:
| Стек | Сэмплов | Процессорное время потока |
|---|---|---|
heavy;[libc.so.6];zt_thread_start;run;spin | 598 | 649 мс |
medium;[libc.so.6];zt_thread_start;run;spin | 447 | 469 мс |
light;[libc.so.6];zt_thread_start;run;spin | 224 | 230 мс |
Свёрнутые стеки в том же формате, что в уроке 36, так что flamegraph.pl рисует по ним SVG, где у каждого потока своё подножие. Настоящий многопоточный сервер под нагрузкой мы профилируем в уроке про пул потоков: туда же приедет дамп по Ctrl-C, без которого сервер, который сам не выходит, профиля не оставит.
Тесты шага
//! Шаг 69: `zt prof --threads`, стеки по потокам.
//!
//! Таблица потоков и префикс с именем потока это чистый код, он проверяется
//! везде. Сквозные тесты запускают `fixtures/prof/threads.c` (три потока с
//! нагрузкой 3 к 2 к 1) через `LD_PRELOAD`, поэтому только на Linux.
const std = @import("std");
const builtin = @import("builtin");
const paths = @import("paths");
const zt = @import("zt");
const prof = zt.prof;
const spin_elf = @embedFile("prof/spin");
test "таблица потоков: граница стека и имя по tid" {
var table: prof.threads.Table = .{};
_ = table.register(4242, 0x7fff0000);
const worker = table.register(4243, 0x7ffe0000).?;
worker.setName("worker-1");
try std.testing.expectEqual(0x7ffe0000, table.stackHigh(4243));
try std.testing.expectEqualStrings("worker-1", table.find(4243).?.name());
// Незаписанный поток: границы нет, обход кадров не начнётся.
try std.testing.expectEqual(0, table.stackHigh(1));
}
test "кольцо помнит, с какого потока сэмпл" {
var ring: prof.ring.Ring(4) = .{};
ring.push(4242, &.{0x1000});
ring.push(4243, &.{ 0x2000, 0x3000 });
try std.testing.expectEqual(4243, ring.get(1).tid);
var buffer: [prof.dump.sample_line_size]u8 = undefined;
try std.testing.expectEqualStrings("4243: 2000 3000\n", prof.dump.formatSample(&buffer, 4243, ring.get(1).stack()));
}
fn loadEmbedded(_: *anyopaque, _: std.mem.Allocator, path: []const u8) ?[]const u8 {
return if (std.mem.eql(u8, path, "/tmp/spin")) spin_elf else null;
}
/// Три сэмпла в hot_a из work (адреса как в тесте шага 49): с потока с
/// именем, с потока без имени и из дампа без номера потока.
const threaded_dump = prof.dump.magic ++
\\maps
\\01001000-01002000 r-xp 00000000 00:102 1 /tmp/spin
\\threads
\\4242 heavy
\\samples
\\4242: 1001540 10015d5
\\4243: 1001540 10015d5
\\1001540 10015d5
\\
;
fn render(text: []const u8, options: prof.Options) ![]const u8 {
const S = struct {
var buffer: [512]u8 = undefined;
};
var out: std.Io.Writer = .fixed(&S.buffer);
var unused: u8 = 0;
_ = try prof.report(std.testing.allocator, &out, text, .{ .context = &unused, .load = loadEmbedded }, options);
return out.buffered();
}
test "с --threads первый кадр это имя потока" {
try std.testing.expectEqualStrings(
\\[tid 4243];work;hot_a 1
\\heavy;work;hot_a 1
\\work;hot_a 1
\\
, try render(threaded_dump, .{ .threads = true }));
}
test "без --threads потоки сливаются в один стек" {
try std.testing.expectEqualStrings("work;hot_a 3\n", try render(threaded_dump, .{}));
}
const Thread = struct {
name: []const u8,
/// Процессорное время, которое поток напечатал сам.
cpu_ms: u64,
samples: u64 = 0,
/// Из них со всем стеком до листа: `zt_thread_start;run;spin`.
in_spin: u64 = 0,
};
/// Прогоняет `zt prof --threads [флаги] threads 50` и раскладывает сэмплы
/// по потокам.
fn profileThreads(extra: []const []const u8) ![3]Thread {
const gpa = std.testing.allocator;
var argv: std.ArrayList([]const u8) = .empty;
defer argv.deinit(gpa);
try argv.appendSlice(gpa, &.{ paths.zt, "prof", "--threads" });
try argv.appendSlice(gpa, extra);
try argv.appendSlice(gpa, &.{ paths.threads, "50" });
const result = try std.process.run(gpa, std.testing.io, .{ .argv = argv.items });
defer gpa.free(result.stdout);
defer gpa.free(result.stderr);
try std.testing.expectEqual(std.process.Child.Term{ .exited = 0 }, result.term);
var threads: [3]Thread = .{
.{ .name = "heavy", .cpu_ms = 0 },
.{ .name = "medium", .cpu_ms = 0 },
.{ .name = "light", .cpu_ms = 0 },
};
var lines = std.mem.tokenizeScalar(u8, result.stdout, '\n');
while (lines.next()) |line| {
for (&threads) |*thread| {
// Строка программы: `heavy 218 ms`.
if (std.mem.startsWith(u8, line, thread.name) and std.mem.endsWith(u8, line, " ms")) {
const number = line[thread.name.len + 1 .. line.len - " ms".len];
thread.cpu_ms = try std.fmt.parseInt(u64, number, 10);
}
// Свёрнутый стек потока: `heavy;...;zt_thread_start;run;spin 187`.
if (std.mem.startsWith(u8, line, thread.name) and line.len > thread.name.len and line[thread.name.len] == ';') {
const space = std.mem.lastIndexOfScalar(u8, line, ' ').?;
const count = try std.fmt.parseInt(u64, line[space + 1 ..], 10);
thread.samples += count;
if (std.mem.indexOf(u8, line, ";zt_thread_start;run;spin ") != null) thread.in_spin += count;
}
}
}
// Почти все сэмплы потока со стеком до листа: обход прошёл через границу
// стека этого потока. Остальное это pthread_setname_np и clock_gettime.
for (threads) |thread| {
if (thread.samples == 0 or thread.in_spin * 10 < thread.samples * 8) {
std.debug.print("профиль потоков:\n{s}", .{result.stdout});
return error.ThreadStackLost;
}
}
return threads;
}
test "zt prof --threads: у каждого потока своё подножие" {
if (builtin.os.tag != .linux) return error.SkipZigTest;
const threads = try profileThreads(&.{});
for (threads) |thread| try std.testing.expect(thread.samples > 0);
// Нагрузка 3 к 1. Общий таймер процесса делит сэмплы неточно, но
// порядок сохраняет с большим запасом.
try std.testing.expect(threads[0].samples > threads[2].samples);
}
test "zt prof --thread-timers: сэмплов не больше, чем процессорного времени" {
if (builtin.os.tag != .linux) return error.SkipZigTest;
const threads = try profileThreads(&.{"--thread-timers"});
for (threads) |thread| {
// 997 Гц это почти сэмпл на миллисекунду. На свободной машине замеры
// дают от 0,8 до 1,0, на перегруженной бывает и меньше половины:
// таймер проверяется на тиках, а поток на ядре бывает урывками.
// Поэтому снизу проверяем только порядок, а сверху предел.
const expected = thread.cpu_ms * prof.sampler.default_hz / 1000;
try std.testing.expect(thread.samples < expected * 5 / 4 + 10);
}
try std.testing.expect(threads[0].samples > threads[2].samples);
}
Первые четыре теста чистые и идут везде: таблица, номер потока в кольце и в строке дампа, отчёт с --threads и без него на синтетическом дампе с адресами из теста шага 49. Два сквозных запускают zt prof на фикстуре и только на Linux.
Тесты шага 49 зовут push, report, prepareEnvironment и start по-старому, и без правки шаг 69 не соберётся. Правка механическая: номер потока 0, пустые настройки отчёта и общий таймер процесса.
var ring: prof.ring.Ring(3) = .{};
- for (0..5) |i| ring.push(&.{ 0x1000 + i, 0x2000 });
+ for (0..5) |i| ring.push(0, &.{ 0x1000 + i, 0x2000 });
try std.testing.expectEqual(3, ring.count());
@@
- const stats = try prof.report(gpa, &out, spin_raw, .{ .context = &unused, .load = loadEmbedded });
+ const stats = try prof.report(gpa, &out, spin_raw, .{ .context = &unused, .load = loadEmbedded }, .{});
try std.testing.expectEqual(198, stats.samples);
@@
var unused: u8 = 0;
- _ = try prof.report(gpa, &out, text, .{ .context = &unused, .load = loadEmbedded });
+ _ = try prof.report(gpa, &out, text, .{ .context = &unused, .load = loadEmbedded }, .{});
try std.testing.expectEqualStrings("work;hot_a 1\n", out.buffered());
@@
defer environment.deinit();
- try prof.prepareEnvironment(gpa, &environment, "/opt/zt/lib/libzt-prof.so", "/tmp/out.raw", 250);
+ try prof.prepareEnvironment(gpa, &environment, "/opt/zt/lib/libzt-prof.so", "/tmp/out.raw", 250, .process);
try std.testing.expectEqualStrings("/opt/zt/lib/libzt-prof.so", environment.get("LD_PRELOAD").?);
@@
const wanted = 60;
- try prof.sampler.start(prof.sampler.default_hz);
+ try prof.sampler.start(prof.sampler.default_hz, .process);
outer(wanted);
- С общим таймером у каждого потока должны быть сэмплы, у
heavyбольше, чем уlight, и не меньше 80 процентов сэмплов каждого потока со стеком до самого листа,zt_thread_start;run;spin. Последнее и проверяет границы стеков: обход прошёл по стеку этого потока, а не чужого. - С
--thread-timersсэмплов не больше, чем разрешает процессорное время потока, с запасом в четверть. Нижней границы нет сознательно: на перегруженной машине (load average около 50 на 16 ядрах) этот режим выдавал меньше половины ожидаемого. Поток попадает на ядро урывками короче тика, и проверка таймера реже застаёт его на процессоре. Тест, который падает от соседской нагрузки, хуже, чем тест, который проверяет меньше.
$ zig build test -Dstep=69 --summary all # macOS: сквозные тесты пропущены
Build Summary: 7/7 steps succeeded; 4/6 tests passed (2 skipped)
$ zig build test -Dstep=69 --summary all # контейнер с Linux
Build Summary: 15/15 steps succeeded; 6/6 tests passed
Весь набор тестов проекта в контейнере с Linux на состоянии шага 69 зелёный: 249 из 252, три пропущены (это тесты, которым нужна не та система).
На macOS
Всё, что в уроке про std.Thread, работает на macOS напрямую: ids.zig, race.zig, cost.zig, эхо-сервер и TINY на потоках, тесты шага our-tiny. Выводы выше сняты на macOS 26.6.2 (Apple M4 Max) и повторены в контейнере ghcr.io/bondiano/runner-zig:dev с Debian 12 (linux/arm64). На macOS std.Thread всегда идёт через pthreads, потому что libc там линкуется всегда.
Разница в деталях. gettid и /proc/PID/task есть только в Linux. На macOS номер потока в ядре даёт pthread_threadid_np, он 64-битный и с pid не совпадает, а потоки процесса показывает ps -M. pthread_setname_np на macOS принимает одно имя и называет только вызвавший поток, а на Linux принимает поток и имя. Поэтому фикстура threads.c написана под Linux. std.Thread.setName в Zig принимает поток и имя на обеих системах, но на macOS умеет назвать только текущий поток, для чужого вернёт error.Unsupported. zt strace работает только в Linux: на macOS системные вызовы смотрит dtruss, которому нужны права root и ослабленная защита системы.
Шаг zt целиком Linux: LD_PRELOAD, /proc/self/maps, доставка ITIMER_PROF потоку и timer_create с SIGEV_THREAD_ID. На macOS zt prof честно отвечает «работает только на Linux», а чистые тесты шага (таблица, кольцо, дамп, отчёт) проходят. Сквозные тесты и прогоны гоняй в контейнере так же, как в уроке 49:
docker run --rm --platform linux/arm64 -v "$PWD":/w:ro ghcr.io/bondiano/runner-zig:dev sh -c '
mkdir -p /tmp/p && cd /w && cp -r build.zig build.zig.zon src fixtures tests /tmp/p/ && cd /tmp/p &&
zig build test -Dstep=69 --summary all --cache-dir /tmp/p/.zig-cache --global-cache-dir /tmp/g'
Замер cost.zig на macOS и на Linux отличается не только числами, но и формой: fork в XNU не растёт с размером процесса так, как в Linux. Не переноси выводы одной системы на другую, меряй на той, где будет жить сервер.
Практика
Параллельное умножение матриц, но в том виде, в каком его пишут в первый раз. Потоки получают полосы строк, и в стартере две ошибки. Первая арифметическая: высота полосы n / threads, и если n не делится на число потоков, последние строки никому не достаются. Вторая это гонка из урока: все потоки получают указатель на одну и ту же переменную с описанием работы, а цикл переписывает её, пока потоки ещё не проснулись. Тесты сверяют результат с последовательной версией на размерах, которые делятся и не делятся на число потоков, при потоках больше, чем строк, и двадцать раз подряд на восьми потоках. Перед сверкой результат забит NaN, так что пропущенная строка видна сразу.
Упражнения
Итоги
- Поток это свой поток управления (регистры, стек, номер, маска сигналов) в общем адресном пространстве процесса. Код, данные, куча, дескрипторы и обработчики сигналов общие.
- Стек потока свой только по договорённости: ядро его не охраняет, и указатель на чужую локальную переменную читает и пишет её.
std.Thread.spawnс libc этоpthread_create, а тот на Linux этоmmapпод стек,mprotectс защитной страницей иcloneс флагамиCLONE_VM,CLONE_FILES,CLONE_SIGHAND,CLONE_THREAD,CLONE_SETTLSи другими. Без libc Zig делает то же сам. Поток и процесс в Linux различаются только флагамиclone.spawnкопирует кортеж аргументов в кучу нового потока. Значение из кортежа у потока своё, указатель из кортежа смотрит туда же, куда и у создателя.- Каждый поток либо ждут через
join, либо отпускают черезdetach.joinждёт конкретный поток.returnизmain,exitиз любого потока и смертельный сигнал в любом потоке завершают все потоки сразу. - Гонка на аргументе: поток читает переменную цикла, когда ему дадут процессор, а не в момент
spawn. Лекарство: передавать значения, а адреса только объектов, которые живут дольше потока и не меняются, пока он читает. - Поток рождается за десятки микросекунд независимо от размера процесса,
forkна Linux дорожает вместе с таблицами страниц. Цена дешевизны: изоляции нет, и порядок доступа к общему обеспечивает программист. - Сервер с потоком на соединение пишется так же прямо, как сервер на процессах: медленный клиент держит только свой поток.
Serverбез изменяемого состояния можно отдать всем потокам без блокировок. forkиз многопоточного процесса оставляет в ребёнке один поток, и доexecveребёнку можно только async-signal-safe функции. Дескрипторы безCLOEXECутекают в детей всех потоков.ITIMER_PROFотдаётSIGPROFпотоку, который жёг процессор. Обработчик узнаёт поток черезgettid, границу его стека берёт из таблицы, которую потоки заполнили при старте через подменённыйpthread_create. Место в общем кольце берётся атомарным сложением: обработчики на разных ядрах работают одновременно.
Дальше
Сегодня потоки обслуживали клиентов и считали матрицы, и нигде не встретились на одной записываемой ячейке, кроме двух мест в нашем профилировщике. Там мы обошлись атомарным сложением и парой .release и .acquire, почти не объясняя, почему без них нельзя. В следующем уроке потоки начнут писать в общее всерьёз. Разберём, какие переменные на самом деле общие, почему cnt++ в двух потоках теряет прибавления, как граф выполнения показывает опасные траектории, и соберём из семафоров буфер производителя и потребителя, читателей и писателей. А куча zl станет общей для нескольких потоков, и сборщику мусора придётся останавливать их всех.
домашка