Раздел 32 · Системное программирование: Zig, ассемблер, Verilog
Разделяемые переменные и семафоры
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Разделяемые переменные и семафоры
В прошлом уроке потоки появились:
spawn,join,detach, сервер на потоках и гонка с индексом цикла, которую мы вылечили передачей по значению. Там гонку удалось обойти, потому что потокам не нужно было ничего общего. Сегодня общее появится нарочно. Разберём, какие переменные у потоков общие на самом деле, почемуcnt += 1в двух потоках теряет прибавления и как это видно на графе выполнения. Потом возьмём семафоры Дейкстры и соберём из них три классические конструкции: мьютекс, ограниченный буферsbufи замок читателей и писателей. В Zig 0.16 все эти примитивы переехали изstd.Threadвstd.Io, и мы посмотрим, из чего они сделаны. А в конце проектzlполучит общую кучу: несколько потоков вычисляют выражения на одних cons-ячейках, и сборщик мусора останавливает мир.
Цели урока
- Для любой переменной программы на Zig сказать, сколько у неё экземпляров и какие потоки до них дотягиваются: глобальная, переменная контейнера внутри функции, локальная,
threadlocal. - Объяснить, почему
cnt += 1не атомарно даже там, где это одна инструкция, и найти на графе выполнения небезопасную зонуbadcnt. - Определить P и V, инвариант семафора и запретную зону; закрыть критическую секцию двоичным семафором и объяснить, чем мьютекс отличается от семафора с одним разрешением.
- Пользоваться
std.Io.Mutex,Io.Semaphore,Io.ConditionиIo.RwLockв Zig 0.16 и понимать, зачем имioи чемwaitотличается отwaitUncancelable. - Написать
sbufна трёх семафорах и на мьютексе с двумя условными переменными и объяснить порядок захвата вinsert. - Различать две задачи о читателях и писателях: кто голодает в каждой и какую выбрала стандартная библиотека.
- Сделать общую кучу
zlдля нескольких потоков: ошибку примитива вthreadlocal, выдачу ячеек пачками под мьютексом и остановку мира для сборщика; прочитать в машинном коде, как пишется переменная потока.
Идея: общая память плюс порядок
Потоки одного процесса живут в одном адресном пространстве. Таблица страниц у них одна, та самая, которую мы разбирали в уроке про трансляцию адресов, поэтому адрес значит одно и то же в любом потоке. Отсюда и удобство, и беда. Удобство: чтобы передать потоку мегабайт данных, достаточно передать указатель. Беда: два потока, которые пишут по одному адресу без договорённости, получают результат, зависящий от того, как планировщик перемешал их инструкции.
Договорённостей бывает две. Первая про исключение: пока один поток внутри участка кода, второй туда не входит. Вторая про порядок: поток ждёт, пока другой не сделает что-то нужное, например не положит элемент в очередь. Семафор умеет и то и другое, поэтому глава 12 книги строит на нём всё. Мы пройдём тот же путь, только на Zig 0.16, где семафор, мьютекс и условная переменная лежат в std.Io.
С точки зрения прикладного кода эти же вещи уже встречались в уроке про инструменты конкурентности: там мьютекс и семафор на промисах защищали общий ресурс в одном потоке JavaScript. Здесь потоки настоящие, ядер несколько, и гонка идёт на уровне отдельных инструкций.
Что у потоков общее
Каждый поток получает свой контекст: идентификатор, стек, указатель стека, счётчик команд, регистры общего назначения и флаги. Всё остальное общее: код, данные, куча, разделяемые библиотеки, открытые файлы. Стек потока тоже лежит в общем адресном пространстве, и защиты у него нет. Поток, у которого есть указатель на переменную в чужом стеке, прочитает и запишет её так же, как свою.
Отсюда точное определение, которым пользуется книга. Переменная разделяемая, если хотя бы на один её экземпляр ссылается больше одного потока. Ключевое слово здесь «экземпляр». Сколько их у переменной, зависит от того, где она объявлена:
| Где объявлена | В C | В Zig | Экземпляров | Где лежит |
|---|---|---|---|---|
| на уровне файла | глобальная | var на уровне файла | один | .data или .bss |
| внутри функции, но живёт вечно | static внутри функции | var в struct внутри функции | один | .data или .bss |
| внутри функции | автоматическая | var или const в теле | по одному на каждый вызов в каждом потоке | стек своего потока |
| на уровне файла, с пометкой | _Thread_local | threadlocal var | по одному на поток | блок TLS потока |
Статических локальных переменных в Zig нет как отдельного понятия, но они не нужны: переменная, объявленная в безымянной структуре внутри функции, видна только из этой функции, а живёт в одном экземпляре на программу, как static в C. Последняя строка таблицы для нас новая. Переменная с пометкой threadlocal объявлена один раз, но у каждого потока своя копия, и никакой другой поток её не видит, если ему не дали указатель. Эта строка понадобится нам в шаге zl.
Проверим все четыре строки одной программой. Главный поток кладёт массив сообщений в свой стек и публикует указатель на него в глобальной переменной, а два потока по очереди читают сообщения через этот указатель:
//! Кто из потоков видит какую переменную. Главный поток кладёт массив
//! сообщений в свой стек и публикует указатель на него в глобальной
//! переменной; два потока читают сообщения через этот указатель, считают
//! вызовы в «статической» переменной и в переменной threadlocal.
const std = @import("std");
/// Глобальная: один экземпляр на программу, в секции данных.
var ptr: *const [2][]const u8 = undefined;
/// Своя у каждого потока: компилятор кладёт её в блок TLS потока.
threadlocal var calls: u32 = 0;
const Seen = struct { msg: []const u8, cnt: u32, calls: u32 };
fn peer(myid: usize, out: *Seen) void {
// Аналог static внутри функции C: переменная контейнера. Видна только
// из этой функции, но экземпляр у неё один на всю программу.
const S = struct {
var cnt: u32 = 0;
};
S.cnt += 1;
calls += 1;
out.* = .{ .msg = ptr[myid], .cnt = S.cnt, .calls = calls };
}
test "стек главного потока виден по указателю, static общая, threadlocal своя" {
// Локальная переменная главного потока, живёт в его стеке.
const msgs = [2][]const u8{ "привет от нулевого", "привет от первого" };
ptr = &msgs;
var seen: [2]Seen = undefined;
// Потоки идут строго по очереди: второй стартует после join первого,
// так что гонки за S.cnt здесь нет, и итог детерминирован.
for (0..2) |myid| {
const t = try std.Thread.spawn(.{}, peer, .{ myid, &seen[myid] });
t.join();
}
try std.testing.expectEqualStrings("привет от нулевого", seen[0].msg);
try std.testing.expectEqualStrings("привет от первого", seen[1].msg);
// Один экземпляр cnt на двоих: второй поток увидел единицу первого.
try std.testing.expectEqual(@as(u32, 1), seen[0].cnt);
try std.testing.expectEqual(@as(u32, 2), seen[1].cnt);
// У каждого потока свой calls, и у главного тоже свой, нетронутый.
try std.testing.expectEqual(@as(u32, 1), seen[0].calls);
try std.testing.expectEqual(@as(u32, 1), seen[1].calls);
try std.testing.expectEqual(@as(u32, 0), calls);
}
$ zig test sharing.zig
1/1 sharing.test.стек главного потока виден по указателю, static общая, threadlocal своя...OK
All 1 tests passed.
Разберём, кто что видит. ptr глобальная, её экземпляр один, и оба потока через него читают msgs, хотя msgs локальная переменная главного потока и лежит в его стеке. Значит, msgs.m (экземпляр в стеке главного потока) разделяемая, и это законно ровно до тех пор, пока главный поток не вышел из функции: после выхода кадр переиспользуется, и указатель будет смотреть на мусор. S.cnt одна на всех: второй поток увидел двойку. calls у каждого потока своя, а у главного осталась нулём, хотя в функции потока её увеличили дважды. myid у каждого потока свой экземпляр в своём стеке, и на него никто, кроме хозяина, не ссылается, так что он не разделяемый.
Потоки здесь запущены строго по очереди: второй стартует после join первого. Сделано это нарочно, чтобы тест был детерминирован. Запусти их одновременно, и S.cnt += 1 станет гонкой. Её мы сейчас и разберём.
badcnt: почему cnt += 1 теряет прибавления
Программа badcnt из книги запускает два потока, каждый niters раз прибавляет единицу к общему счётчику, и в конце счётчик должен стать 2 * niters. В модуле tiny рядом с ней лежат три исправления, к ним мы придём ниже. Весь файл:
//! `badcnt` из главы 12: два потока по `niters` раз прибавляют единицу к
//! общему счётчику. `cnt++` это три инструкции (загрузить, прибавить,
//! сохранить), и если потоки перемешают их на графе выполнения, одно из
//! прибавлений теряется. Рядом три исправления: семафор как мьютекс
//! (`goodcnt` книги), `std.Io.Mutex` и атомарный `fetchAdd`.
//!
//! Счётчик `badcnt` объявлен через `volatile`, как в книге: без него
//! компилятор в `ReleaseFast` сложит цикл в одно `cnt += niters`, и гонку
//! станет почти не видно. `volatile` не делает доступ атомарным, он только
//! запрещает выкидывать загрузки и сохранения.
const std = @import("std");
const Io = std.Io;
pub const Result = struct {
cnt: u64,
expected: u64,
pub fn ok(r: Result) bool {
return r.cnt == r.expected;
}
};
fn run(nthreads: usize, comptime worker: anytype, args: anytype) !void {
var threads: [16]std.Thread = undefined;
for (threads[0..nthreads]) |*t| t.* = try std.Thread.spawn(.{}, worker, args);
for (threads[0..nthreads]) |t| t.join();
}
fn countBad(cnt: *volatile u64, niters: u64) void {
for (0..niters) |_| cnt.* += 1;
}
/// Гонка: результат часто меньше `2 * niters`, но может и совпасть.
pub fn badcnt(niters: u64) !Result {
var cnt: u64 = 0;
try run(2, countBad, .{ @as(*volatile u64, &cnt), niters });
return .{ .cnt = cnt, .expected = 2 * niters };
}
fn countSem(io: Io, mutex: *Io.Semaphore, cnt: *u64, niters: u64) void {
for (0..niters) |_| {
mutex.waitUncancelable(io); // P(&mutex)
cnt.* += 1;
mutex.post(io); // V(&mutex)
}
}
/// `goodcnt` книги: двоичный семафор вокруг критической секции.
pub fn goodcnt(io: Io, niters: u64) !Result {
var cnt: u64 = 0;
var mutex: Io.Semaphore = .{ .permits = 1 };
try run(2, countSem, .{ io, &mutex, &cnt, niters });
return .{ .cnt = cnt, .expected = 2 * niters };
}
fn countMutex(io: Io, mutex: *Io.Mutex, cnt: *u64, niters: u64) void {
for (0..niters) |_| {
mutex.lockUncancelable(io);
cnt.* += 1;
mutex.unlock(io);
}
}
pub fn mutexcnt(io: Io, niters: u64) !Result {
var cnt: u64 = 0;
var mutex: Io.Mutex = .init;
try run(2, countMutex, .{ io, &mutex, &cnt, niters });
return .{ .cnt = cnt, .expected = 2 * niters };
}
fn countAtomic(cnt: *std.atomic.Value(u64), niters: u64) void {
// `monotonic` хватает: порядок относительно других переменных не нужен,
// нужна только неделимость самой операции.
for (0..niters) |_| _ = cnt.fetchAdd(1, .monotonic);
}
/// Атомарный счётчик. Результат `fetchAdd` отброшен, так что на x86-64 это
/// одна `lock addq`, на aarch64 с LSE (Apple M) одна `ldadd`, а без LSE
/// цикл `ldxr`/`stxr`, который повторяет попытку, пока запись не пройдёт.
pub fn atomiccnt(niters: u64) !Result {
var cnt: std.atomic.Value(u64) = .init(0);
try run(2, countAtomic, .{ &cnt, niters });
return .{ .cnt = cnt.load(.monotonic), .expected = 2 * niters };
}
Прогоним на macOS (Apple M4 Max, 16 ядер, macOS 26.6.2, отладочная сборка) три раза подряд:
$ tiny badcnt 1000000
badcnt cnt=1006147 expected=2000000 BOOM!
goodcnt (sem) cnt=2000000 expected=2000000 OK
mutexcnt cnt=2000000 expected=2000000 OK
atomiccnt cnt=2000000 expected=2000000 OK
$ tiny badcnt 1000000
badcnt cnt=1090022 expected=2000000 BOOM!
...
$ tiny badcnt 1000000
badcnt cnt=1000777 expected=2000000 BOOM!
...
Почти половина прибавлений пропала. В ReleaseFast картина та же: 1127359, 1107180, 1000542. А вот три прогона той же сборки ReleaseFast в контейнере с Linux (linux/arm64, 16 vCPU, та же машина):
badcnt cnt=1042489 expected=2000000 BOOM!
badcnt cnt=1037310 expected=2000000 BOOM!
badcnt cnt=2000000 expected=2000000 OK
Третий прогон дал верный ответ. Это самое неприятное свойство гонки: она не обязана проявиться. Программа, которая прошла тест, могла просто получить удачное расписание. Поэтому тест шага проверяет у badcnt только то, что счётчик не больше 2n, а у исправленных версий требует точного 2n каждый раз.
Одна строка, три действия
Почему теряются прибавления? Посмотрим, во что превращается цикл countBad. Для aarch64 (сборка ReleaseFast, zig build-obj -target aarch64-linux):
<countBad>:
cbz x1, <выход>
ldr x8, [x0] # L: загрузить cnt в регистр
subs x1, x1, #0x1 # счётчик цикла
add x8, x8, #0x1 # U: прибавить единицу в регистре
str x8, [x0] # S: сохранить обратно
b.ne <ldr>
Три инструкции на прибавление, как в книге: загрузить, изменить, сохранить. Книга обозначает их L, U и S, а то, что до и после (H, заголовок цикла, и T, хвост), в гонке не участвует. Если оба потока загрузили одно и то же значение, оба прибавили единицу к своей копии в регистре и оба сохранили, счётчик вырос на один, а не на два.
На x86-64 компилятор делает одну инструкцию:
<countBad>:
addq $0x1, (%rdi) # cnt += 1 прямо в памяти
addq $0x1, (%rdi) # цикл развёрнут по четыре
addq $0x1, (%rdi)
addq $0x1, (%rdi)
addq $-0x4, %rsi
jne <addq>
Кажется, что тут гонке негде взяться. Но addq с операндом в памяти внутри процессора это те же три шага: прочитать линию кэша, сложить, записать. Между чтением и записью другое ядро вправе прочитать ту же линию, и одно прибавление пропадёт так же. Неделимой инструкцию делает только префикс lock: он удерживает линию кэша за ядром до конца операции. Именно его компилятор ставит в атомарной версии:
<countAtomic>: # x86-64
lock addq $0x1, (%rdi)
<countAtomic>: # aarch64 с расширением LSE, как у Apple M
ldadd x8, x9, [x0]
Про volatile в объявлении счётчика. Без него ReleaseFast сложил бы весь цикл в одно cnt += niters, и гонку стало бы почти не видно. volatile запрещает компилятору выкидывать и склеивать обращения к памяти, но атомарности не даёт никакой. Это частая ошибка, пришедшая из C: volatile для памяти, которую меняет устройство, а для памяти, которую меняет другой поток, нужны атомики или блокировка.
Граф выполнения
Чтобы рассуждать о перемешивании инструкций строго, книга рисует граф выполнения. Для двух потоков это сетка: по горизонтали инструкции потока 1, по вертикали потока 2. Точка (i, j) означает, что первый поток выполнил i инструкций, а второй j. Выполнение программы это траектория из левого нижнего угла в правый верхний. Шаг вправо это инструкция первого потока, шаг вверх второго. По диагонали ходить нельзя: модель считает, что инструкции выполняются по одной.
Инструкции L, U, S одного потока образуют его критическую секцию. Пересечение двух критических секций на графе это прямоугольник, внутри которого оба потока уже загрузили cnt и ещё не сохранили. Это небезопасная зона. Траектория, которая обходит её, даёт верный ответ, траектория, которая через неё проходит, теряет прибавление.
Попробуй сам. Тяни траекторию мышью или шагай стрелками, строка под графом показывает регистры обоих потоков и cnt в каждой точке:
Виджет перебирает все траектории и считает, сколько из них безопасны. При одной итерации у каждого потока пять инструкций, и путей из угла в угол C(10, 5) = 252. Безопасных среди них всего 52: либо первый поток успевает сохранить cnt до того, как второй его загрузит, либо наоборот. Переключи на две итерации: путей 12 870, безопасных 230. Чем длиннее цикл, тем меньше шансов у случайного расписания обойти все небезопасные зоны, и при миллионе итераций гонка проявляется почти всегда.
Граф не описывает многоядерную машину буквально: на двух ядрах инструкции выполняются одновременно, а не по очереди. Но для такой программы результат одновременного выполнения совпадает с результатом какой-то траектории, поэтому граф честно отвечает на вопрос «может ли сломаться». Где эта модель перестаёт работать (процессор и компилятор переставляют независимые обращения к памяти), разбирает урок про атомики и порядки памяти в Rust.
Семафоры
Нужен способ запретить траекториям входить в небезопасную зону. Эдсгер Дейкстра предложил для этого семафор: неотрицательное целое s, с которым можно делать только две неделимые операции.
- P(s): если
sбольше нуля, уменьшить его на единицу и вернуться. Еслиsравен нулю, заснуть до тех пор, пока кто-нибудь не сделает V, а потом повторить попытку. - V(s): увеличить
sна единицу. Если кто-то спит в P, разбудить одного из них.
Названия голландские: P от proberen (проверить), V от verhogen (увеличить). Главное свойство называется инвариантом семафора: при правильной работе s никогда не становится отрицательным. Всё, что мы дальше построим, держится на этом инварианте.
Мьютекс как двоичный семафор
Заведём семафор s с начальным значением 1 и обернём критическую секцию: P(s) перед ней, V(s) после. Это goodcnt из листинга выше. Семафор, который может быть только нулём или единицей, называется двоичным, а используемый так, для взаимного исключения, называется мьютексом.
Что это меняет на графе? Посчитаем значение s в каждой точке: 1 минус число потоков, которые сделали P и ещё не сделали V. Точки, где оба потока между своими P и V, дали бы s = -1, а инвариант этого не допускает. Эти точки образуют запретную зону, и она целиком накрывает небезопасную. Ни одна траектория больше не может в неё войти: поток, который попытался бы, спит на P.
Тут виджет открыт сразу с семафором и двумя итерациями. Попробуй шагнуть в серую зону: поток откажется и уснёт на P. Все 486 траекторий, которые обходят запретную зону, дают cnt = 4, и тупиков среди них нет.
Мьютекс и двоичный семафор не одно и то же, хотя код выглядит одинаково. У мьютекса есть хозяин: отпустить его должен тот поток, который взял. У семафора хозяина нет, и V может сделать кто угодно. Это не недостаток, а другое назначение: мьютекс для исключения, семафор ещё и для сигнала «у меня готово». Мы увидим это в sbuf, где производитель делает V, а просыпается потребитель.
Примитивы Zig 0.16 живут в std.Io
Если ты писал на Zig 0.13 или 0.14, то помнишь std.Thread.Mutex, std.Thread.Semaphore, std.Thread.Condition и std.Thread.Pool. В 0.16 их там нет. Примитивы синхронизации переехали в std.Io, а в std.Thread остались сами потоки: spawn, join, detach, yield, getCpuCount. Вот что есть теперь:
| Примитив | Создать | Операции |
|---|---|---|
Io.Mutex | .init | lock(io), lockUncancelable(io), unlock(io), tryLock() |
Io.Semaphore | .{ .permits = n } | wait(io) это P, waitUncancelable(io), post(io) это V |
Io.Condition | .init | wait(io, &mutex), waitUncancelable(io, &mutex), signal(io), broadcast(io) |
Io.RwLock | .init | lock, unlock, lockShared, unlockShared и варианты Uncancelable |
std.atomic.Value(T) | .init(x) | load, store, fetchAdd, cmpxchgStrong с порядком памяти |
Зачем мьютексу io? Затем же, зачем функции чтения файла: ожидание это операция ввода-вывода в широком смысле, и выполняет её та реализация Io, которую выбрала программа. Io.Threaded усыпляет поток ядра. Io.Evented в той же ситуации усыпил бы лёгкий поток внутри цикла событий и отдал ядро другому. Код, который написан против Io, работает в обоих мирах, и об этом весь урок 73.
Отсюда же второе новшество: точка отмены. lock(io) и wait(io) возвращают Io.Cancelable!void: если задачу отменили, пока она ждала, вызов вернёт error.Canceled. Варианты lockUncancelable и waitUncancelable ждут до конца. Потоки из std.Thread.spawn никто не отменяет, поэтому эталон везде зовёт Uncancelable. В код-задаче урока наоборот: там сигнатуры честные, и ошибку отмены надо пробросить через try.
Из чего сделан мьютекс? Откроем lib/std/Io.zig. Состояний три: свободен, занят один раз, занят и кто-то ждёт:
pub fn lock(m: *Mutex, io: Io) Cancelable!void {
const initial_state = m.state.cmpxchgStrong(
.unlocked,
.locked_once,
.acquire,
.monotonic,
) orelse {
@branchHint(.likely);
return;
};
if (initial_state == .contended) {
try io.futexWait(State, &m.state.raw, .contended);
}
while (m.state.swap(.contended, .acquire) != .unlocked) {
try io.futexWait(State, &m.state.raw, .contended);
}
}
Быстрый путь это одна атомарная операция сравнения с обменом, без системного вызова: если мьютекс был свободен, он наш. Медленный путь помечает мьютекс как спорный и засыпает через io.futexWait. Io.Threaded превращает это в системный вызов futex на Linux и в __ulock_wait2 на macOS. unlock будит спящего, только если состояние было спорным. Пока потоки не дерутся, мьютекс стоит пару атомарных инструкций.
А семафор в lib/std/Io/Semaphore.zig оказывается совсем не объектом ядра, а счётчиком под мьютексом с условной переменной:
mutex: Io.Mutex = .init,
cond: Io.Condition = .init,
/// It is OK to initialize this field to any value.
permits: usize = 0,
pub fn waitUncancelable(s: *Semaphore, io: Io) void {
s.mutex.lockUncancelable(io);
defer s.mutex.unlock(io);
while (s.permits == 0) s.cond.waitUncancelable(io, &s.mutex);
s.permits -= 1;
if (s.permits > 0) s.cond.signal(io);
}
pub fn post(s: *Semaphore, io: Io) void {
s.mutex.lockUncancelable(io);
defer s.mutex.unlock(io);
s.permits += 1;
s.cond.signal(io);
}
Тут есть тонкость, которой нет в книге. По определению Дейкстры V будит одного спящего, и тот завершает свой P: разрешение переходит из рук в руки. post в Zig только увеличивает счётчик и подаёт сигнал. Разбуженный поток сначала должен снова взять мьютекс и проверить permits, и если за это время пришёл новый поток и успел сделать wait, разрешение достанется новичку, а разбуженный снова уснёт. Инвариант от этого не страдает, а вот справедливость страдает: семафор Zig не гарантирует, что ждущий дольше всех войдёт первым. Для goodcnt это неважно, для замка читателей и писателей ниже это важно.
Семафоры для порядка: производители и потребители
Мьютекс решает задачу исключения. Вторая задача, порядок, решается семафором, который считает ресурсы. Классический пример это ограниченный буфер: производители кладут элементы, потребители забирают, между ними кольцо на n слотов. Так устроены пайп в ядре, очередь сообщений, очередь соединений в сервере с пулом потоков из следующего урока.
Книга называет эту структуру sbuf и держит её на трёх семафорах:
mutexс начальным значением 1 защищает кольцо и индексы;slotsсчитает свободные слоты, в начале ихn;itemsсчитает занятые, в начале ноль.
Производитель делает P над slots (ждёт, пока есть куда положить), кладёт элемент под мьютексом и делает V над items (объявляет, что появилось что забрать). Потребитель зеркально: P над items, забирает под мьютексом, V над slots. Никто не крутится в цикле проверок: поток, которому нечего делать, спит в P.
Поиграй с буфером на три слота, двумя производителями и двумя потребителями. Шагай отдельными потоками, а потом включи «авто» и посмотри, как спящие потоки просыпаются от чужого V:
Виджет моделирует семафор по книге: V, у которого есть спящий, отдаёт разрешение прямо ему. Выше мы видели, что Io.Semaphore так не делает. На результат sbuf это не влияет: элемент всё равно достанется кому-то из потребителей, и ни один не потеряется.
Вот sbuf из эталона, в двух версиях с одним интерфейсом:
//! Ограниченный буфер `sbuf` из главы 12: кольцо на `n` слотов между
//! производителями и потребителями. Две версии с одним интерфейсом.
//!
//! `Sbuf` как в книге, на трёх семафорах: `mutex` (начальное значение 1)
//! даёт взаимное исключение, `slots` считает свободные слоты, `items`
//! занятые. В Zig 0.16 семафор живёт в `std.Io.Semaphore`, операции
//! P и V называются `wait` и `post`, и обеим нужен `Io`.
//!
//! `SbufCond` на мьютексе и двух условных переменных: так ограниченный
//! буфер пишут без семафоров. Условие перепроверяется в `while`: проснуться
//! можно и без повода, и после того, как слот успел забрать другой поток.
const std = @import("std");
const Io = std.Io;
pub fn Sbuf(comptime T: type) type {
return struct {
const Self = @This();
io: Io,
buf: []T,
/// `buf[(front + 1) % n]` первый элемент.
front: usize = 0,
/// `buf[rear % n]` последний элемент.
rear: usize = 0,
mutex: Io.Semaphore = .{ .permits = 1 },
slots: Io.Semaphore,
items: Io.Semaphore = .{ .permits = 0 },
pub fn init(gpa: std.mem.Allocator, io: Io, n: usize) !Self {
return .{ .io = io, .buf = try gpa.alloc(T, n), .slots = .{ .permits = n } };
}
pub fn deinit(sp: *Self, gpa: std.mem.Allocator) void {
gpa.free(sp.buf);
}
/// P(slots), P(mutex), положить, V(mutex), V(items).
pub fn insert(sp: *Self, item: T) void {
sp.slots.waitUncancelable(sp.io);
sp.mutex.waitUncancelable(sp.io);
sp.rear += 1;
sp.buf[sp.rear % sp.buf.len] = item;
sp.mutex.post(sp.io);
sp.items.post(sp.io);
}
/// P(items), P(mutex), забрать, V(mutex), V(slots).
pub fn remove(sp: *Self) T {
sp.items.waitUncancelable(sp.io);
sp.mutex.waitUncancelable(sp.io);
sp.front += 1;
const item = sp.buf[sp.front % sp.buf.len];
sp.mutex.post(sp.io);
sp.slots.post(sp.io);
return item;
}
};
}
pub fn SbufCond(comptime T: type) type {
return struct {
const Self = @This();
io: Io,
buf: []T,
front: usize = 0,
rear: usize = 0,
/// Сколько элементов в буфере сейчас.
count: usize = 0,
mutex: Io.Mutex = .init,
not_full: Io.Condition = .init,
not_empty: Io.Condition = .init,
pub fn init(gpa: std.mem.Allocator, io: Io, n: usize) !Self {
return .{ .io = io, .buf = try gpa.alloc(T, n) };
}
pub fn deinit(sp: *Self, gpa: std.mem.Allocator) void {
gpa.free(sp.buf);
}
pub fn insert(sp: *Self, item: T) void {
sp.mutex.lockUncancelable(sp.io);
defer sp.mutex.unlock(sp.io);
while (sp.count == sp.buf.len) sp.not_full.waitUncancelable(sp.io, &sp.mutex);
sp.rear += 1;
sp.buf[sp.rear % sp.buf.len] = item;
sp.count += 1;
sp.not_empty.signal(sp.io);
}
pub fn remove(sp: *Self) T {
sp.mutex.lockUncancelable(sp.io);
defer sp.mutex.unlock(sp.io);
while (sp.count == 0) sp.not_empty.waitUncancelable(sp.io, &sp.mutex);
sp.front += 1;
const item = sp.buf[sp.front % sp.buf.len];
sp.count -= 1;
sp.not_full.signal(sp.io);
return item;
}
};
}
Три вещи в первой версии заслуживают внимания.
Порядок захвата. В insert сначала P над slots, потом над mutex. Поменяй их местами, и получишь тупик на полном буфере: производитель возьмёт мьютекс и уснёт на slots, держа мьютекс, а потребитель, который мог бы освободить слот, уснёт на мьютексе. Каждый ждёт другого вечно. Это взаимоблокировка, и в уроке 72 мы нарисуем её на графе выполнения.
Индексы растут без конца. front и rear не заворачиваются, слот считается как rear % n, как в книге. Переполнения бояться не нужно: usize на 64 битах при миллиарде вставок в секунду переполнится через пятьсот восемьдесят лет. Код-задача урока заворачивает индексы сразу, это тоже правильно.
Мьютекс здесь семафор. Поле mutex в Sbuf имеет тип Io.Semaphore с одним разрешением, ровно как в книге. Настоящий Io.Mutex был бы дешевле: семафор Zig сам внутри берёт мьютекс, так что каждое P над нашим «мьютексом» стоит захвата и отпускания настоящего. Эталон оставил книжную форму, чтобы три семафора читались одинаково.
Условные переменные
Вторая версия, SbufCond, обходится без семафоров: один Io.Mutex и две условные переменные, not_full и not_empty. Семафор хранит число, условная переменная не хранит ничего: это место, где можно подождать, пока условие над общими данными не станет истинным. Условие проверяет сам поток, под мьютексом.
wait(io, &mutex) делает три вещи неделимо: отпускает мьютекс, засыпает, а проснувшись, берёт мьютекс обратно. Неделимость здесь главное. Если бы поток отпустил мьютекс и только потом лёг спать, сигнал, поданный между этими двумя шагами, ушёл бы в пустоту, и поток проспал бы его навсегда.
Почему проверка в while, а не в if? Причин две. Первая: между сигналом и тем моментом, когда разбуженный снова взял мьютекс, другой поток мог успеть забрать слот. Вторая: контракт условной переменной не обещает, что wait вернётся только по сигналу, адресованному именно тебе. POSIX прямо разрешает ложные пробуждения. Io.Condition в Zig прячет ложные пробуждения futex внутри своего цикла (в исходнике это место помечено комментарием про spurious wakeup), но сигнал, поданный для одного ждущего, может достаться другому. Поэтому правило без исключений: условие ожидания всегда проверяется в цикле.
signal будит одного ждущего, broadcast всех. В SbufCond хватает signal: одна вставка освобождает ровно одного потребителя. Если бы условие у разных ждущих было разным (например, одни ждут один элемент, другие десять), signal мог бы разбудить не того, и нужен был бы broadcast.
Что выбрать, семафор или условную переменную? По силе они равны: мы только что видели, что Io.Semaphore сам построен на мьютексе и условной переменной, а в обратную сторону условную переменную тоже можно собрать из семафоров. Семафор удобен, когда ждут ресурса, который можно посчитать: слоты, элементы, соединения. Условная переменная удобна, когда ждут произвольного условия над несколькими полями: «очередь не пуста и сервер не остановлен». В шаге zl ниже будет такой случай.
Читатели и писатели
Обобщим взаимное исключение. Общий объект одни потоки только читают, другие меняют. Писатель должен работать в одиночестве, а читатели могут читать одновременно сколько угодно: они друг другу не мешают. Так устроены кэш веб-прокси из урока 73, таблица маршрутов, конфигурация, которую перечитывают на лету.
У задачи два классических решения, и они различаются тем, кого пропускают вперёд.
Первая задача, приоритет читателей. Читатель не ждёт, если объект уже кто-то читает. Первый вошедший читатель берёт семафор w за всех, последний вышедший отпускает. Писатель просто берёт w. Пока читатели приходят непрерывно, w так и не освободится, и писатель будет ждать вечно. Это голодание.
Вторая задача, приоритет писателей. Её решение опубликовали Куртуа, Хейманс и Парнас в 1971 году. Пришедший писатель закрывает вход новым читателям семафором r, дожидается, пока выйдут те, кто уже внутри, и работает. Последний писатель в очереди открывает вход. Теперь голодать могут читатели, если писатели идут потоком.
Вот оба решения и третий замок, std.Io.RwLock из стандартной библиотеки, за одним интерфейсом:
//! Читатели и писатели из главы 12: общий объект читают многие сразу,
//! а пишет только один и в одиночестве. Три замка с одним интерфейсом
//! (`lockRead`, `unlockRead`, `lockWrite`, `unlockWrite`):
//!
//! * `ReadersFirst`, первая задача, листинг книги: читатель не ждёт, если
//! объект уже читают. Поток читателей может уморить писателя голодом.
//! * `WritersFirst`, вторая задача (Куртуа, Хейманс, Парнас, 1971): пришедший
//! писатель закрывает вход новым читателям. Голодать теперь могут они.
//! * `Std` над `std.Io.RwLock`: готовый замок из стандартной библиотеки.
//!
//! Все семафоры тут двоичные, начальное значение 1: `P` это `wait`, `V` это `post`.
const std = @import("std");
const Io = std.Io;
const Sem = Io.Semaphore;
fn P(io: Io, s: *Sem) void {
s.waitUncancelable(io);
}
fn V(io: Io, s: *Sem) void {
s.post(io);
}
pub const ReadersFirst = struct {
io: Io,
/// Сколько читателей сейчас внутри.
readcnt: usize = 0,
/// Защищает `readcnt`.
mutex: Sem = .{ .permits = 1 },
/// Держит первый вошедший читатель за всех или писатель.
w: Sem = .{ .permits = 1 },
pub fn init(io: Io) ReadersFirst {
return .{ .io = io };
}
pub fn lockRead(l: *ReadersFirst) void {
P(l.io, &l.mutex);
l.readcnt += 1;
if (l.readcnt == 1) P(l.io, &l.w);
V(l.io, &l.mutex);
}
pub fn unlockRead(l: *ReadersFirst) void {
P(l.io, &l.mutex);
l.readcnt -= 1;
if (l.readcnt == 0) V(l.io, &l.w);
V(l.io, &l.mutex);
}
pub fn lockWrite(l: *ReadersFirst) void {
P(l.io, &l.w);
}
pub fn unlockWrite(l: *ReadersFirst) void {
V(l.io, &l.w);
}
};
pub const WritersFirst = struct {
io: Io,
readcnt: usize = 0,
writecnt: usize = 0,
/// Защищает `readcnt`.
mutex_r: Sem = .{ .permits = 1 },
/// Защищает `writecnt`.
mutex_w: Sem = .{ .permits = 1 },
/// Турникет читателей: пропускает по одному, чтобы ждущий писатель
/// взял `r` раньше, чем очередь читателей.
turnstile: Sem = .{ .permits = 1 },
/// Вход для читателей. Первый писатель закрывает его, последний открывает.
r: Sem = .{ .permits = 1 },
/// Сам объект: либо писатель, либо первый читатель за всех.
w: Sem = .{ .permits = 1 },
pub fn init(io: Io) WritersFirst {
return .{ .io = io };
}
pub fn lockRead(l: *WritersFirst) void {
P(l.io, &l.turnstile);
P(l.io, &l.r);
P(l.io, &l.mutex_r);
l.readcnt += 1;
if (l.readcnt == 1) P(l.io, &l.w);
V(l.io, &l.mutex_r);
V(l.io, &l.r);
V(l.io, &l.turnstile);
}
pub fn unlockRead(l: *WritersFirst) void {
P(l.io, &l.mutex_r);
l.readcnt -= 1;
if (l.readcnt == 0) V(l.io, &l.w);
V(l.io, &l.mutex_r);
}
pub fn lockWrite(l: *WritersFirst) void {
P(l.io, &l.mutex_w);
l.writecnt += 1;
if (l.writecnt == 1) P(l.io, &l.r);
V(l.io, &l.mutex_w);
P(l.io, &l.w);
}
pub fn unlockWrite(l: *WritersFirst) void {
V(l.io, &l.w);
P(l.io, &l.mutex_w);
l.writecnt -= 1;
if (l.writecnt == 0) V(l.io, &l.r);
V(l.io, &l.mutex_w);
}
};
pub const Std = struct {
io: Io,
lock: Io.RwLock = .init,
pub fn init(io: Io) Std {
return .{ .io = io };
}
pub fn lockRead(l: *Std) void {
l.lock.lockSharedUncancelable(l.io);
}
pub fn unlockRead(l: *Std) void {
l.lock.unlockShared(l.io);
}
pub fn lockWrite(l: *Std) void {
l.lock.lockUncancelable(l.io);
}
pub fn unlockWrite(l: *Std) void {
l.lock.unlock(l.io);
}
};
Самое хитрое место во втором варианте это turnstile, турникет читателей. Без него все ждущие читатели стояли бы в очереди семафора r, и когда писатель, закончив, откроет r, семафор мог бы пропустить вперёд очередного читателя, а не следующего писателя. С турникетом у r ждёт не больше одного читателя, остальные стоят раньше, у турникета. Пришедший писатель встаёт в очередь r одновременно с этим единственным читателем и получает её раньше всех, кто стоит у турникета.
А какой вариант у стандартной библиотеки? Открой lib/std/Io/RwLock.zig. Читатель входит по быстрому пути, без мьютекса, только если в состоянии замка нет ни пишущего, ни ждущих писателей (is_writing | writer_mask). Писатель же сначала увеличивает счётчик ждущих писателей, и с этого момента новые читатели уходят на медленный путь через мьютекс, который писатель держит. Выходит второй вариант: писатели вперёд.
Проверить приоритет можно одним тестом, и он есть в тестах шага. Первый читатель берёт замок, приходит писатель, через 50 мс второй читатель, и первый отпускает. В ReadersFirst второй читатель входит раньше писателя: объект уже читают, значит можно. В WritersFirst писатель успел закрыть вход, и второй читатель ждёт.
Шаг проекта: семафоры в модуле tiny
Три файла шага мы уже прочитали целиком: badcnt.zig, sbuf.zig и rw.zig. Остались подключение и тесты.
pub const echo_threads = @import("conc/echo_threads.zig");
pub const matmul = @import("conc/matmul.zig");
+ pub const sbuf = @import("conc/sbuf.zig");
+ pub const rw = @import("conc/rw.zig");
+ pub const badcnt = @import("conc/badcnt.zig");
};
\\ tiny conc-bench [--clients 100,1000] [--active A] [--lines K] [--models procs,threads,select,poll,epoll] [--json]
+ \\ tiny badcnt [niters] гонка на счётчике и три исправления
\\
;
@@
.{ "conc-bench", cmdConcBench },
+ .{ "badcnt", cmdBadcnt },
};
@@
if (json) try conc.bench.writeJson(ctx.out, results.items) else try conc.bench.writeTable(ctx.out, results.items);
try ctx.out.flush();
}
+fn cmdBadcnt(ctx: Ctx, rest: []const [:0]const u8) !void {
+ const niters: u64 = if (rest.len > 0) try parseCount(ctx, rest[0]) else 1_000_000;
+ const b = conc.badcnt;
+ const rows = .{
+ .{ "badcnt", try b.badcnt(niters) },
+ .{ "goodcnt (sem)", try b.goodcnt(ctx.io, niters) },
+ .{ "mutexcnt", try b.mutexcnt(ctx.io, niters) },
+ .{ "atomiccnt", try b.atomiccnt(niters) },
+ };
+ inline for (rows) |row| {
+ const r = row[1];
+ try ctx.out.print("{s:<14} cnt={d} expected={d} {s}\n", .{ row[0], r.cnt, r.expected, if (r.ok()) "OK" else "BOOM!" });
+ }
+ try ctx.out.flush();
+}
+
fn listenOn(ctx: Ctx, name: []const u8, port: u16) !std.c.fd_t {
-const project_steps = [_]u8{ 63, 64, 65, 66, 68, 69 };
+const project_steps = [_]u8{ 63, 64, 65, 66, 68, 69, 70 };
Тесты шага
//! Шаг 70: `sbuf` на семафорах и на условных переменных, читатели и
//! писатели в двух вариантах и `std.Io.RwLock`, `badcnt` и исправления.
const std = @import("std");
const tiny = @import("tiny");
const conc = tiny.conc;
const testing = std.testing;
const io = testing.io;
const producers = 3;
const consumers = 3;
const per_producer = 2000;
const total = producers * per_producer;
fn produce(comptime Buf: type, sp: *Buf, id: u32) void {
for (0..per_producer) |k| sp.insert(id * per_producer + @as(u32, @intCast(k)));
}
fn consume(comptime Buf: type, sp: *Buf, seen: []std.atomic.Value(u8), count: usize) void {
for (0..count) |_| _ = seen[sp.remove()].fetchAdd(1, .monotonic);
}
fn checkBuffer(comptime Buf: type) !void {
var sp: Buf = try .init(testing.allocator, io, 4);
defer sp.deinit(testing.allocator);
var seen: [total]std.atomic.Value(u8) = @splat(.init(0));
var threads: [producers + consumers]std.Thread = undefined;
for (0..consumers) |i| threads[i] = try std.Thread.spawn(.{}, consume, .{ Buf, &sp, &seen, total / consumers });
for (0..producers) |i| threads[consumers + i] = try std.Thread.spawn(.{}, produce, .{ Buf, &sp, @as(u32, @intCast(i)) });
for (threads) |t| t.join();
// Каждый элемент дошёл ровно один раз: ни потерь, ни дублей.
for (seen) |s| try testing.expectEqual(@as(u8, 1), s.load(.monotonic));
// Кольцо пустое: вставили столько же, сколько забрали.
try testing.expectEqual(sp.front, sp.rear);
}
test "sbuf на трёх семафорах: 3 производителя, 3 потребителя, 4 слота" {
try checkBuffer(conc.sbuf.Sbuf(u32));
}
test "sbuf на мьютексе и двух условных переменных" {
try checkBuffer(conc.sbuf.SbufCond(u32));
}
fn Probe(comptime Lock: type) type {
return struct {
lock: *Lock,
readers: std.atomic.Value(u32) = .init(0),
writers: std.atomic.Value(u32) = .init(0),
violations: std.atomic.Value(u32) = .init(0),
max_readers: std.atomic.Value(u32) = .init(0),
value: u64 = 0,
fn reader(p: *@This()) void {
for (0..2000) |_| {
p.lock.lockRead();
const now = p.readers.fetchAdd(1, .acq_rel) + 1;
_ = p.max_readers.fetchMax(now, .monotonic);
if (p.writers.load(.acquire) != 0) _ = p.violations.fetchAdd(1, .monotonic);
std.mem.doNotOptimizeAway(p.value);
_ = p.readers.fetchSub(1, .acq_rel);
p.lock.unlockRead();
}
}
fn writer(p: *@This()) void {
for (0..500) |_| {
p.lock.lockWrite();
if (p.writers.fetchAdd(1, .acq_rel) != 0) _ = p.violations.fetchAdd(1, .monotonic);
if (p.readers.load(.acquire) != 0) _ = p.violations.fetchAdd(1, .monotonic);
p.value += 1;
_ = p.writers.fetchSub(1, .acq_rel);
p.lock.unlockWrite();
}
}
};
}
fn checkRw(comptime Lock: type) !void {
var lock: Lock = .init(io);
var p: Probe(Lock) = .{ .lock = &lock };
var threads: [6]std.Thread = undefined;
for (threads[0..4]) |*t| t.* = try std.Thread.spawn(.{}, Probe(Lock).reader, .{&p});
for (threads[4..6]) |*t| t.* = try std.Thread.spawn(.{}, Probe(Lock).writer, .{&p});
for (threads) |t| t.join();
// Писатель всегда один и всегда без читателей; ни одна запись не потеряна.
try testing.expectEqual(@as(u32, 0), p.violations.load(.monotonic));
try testing.expectEqual(@as(u64, 1000), p.value);
}
test "читатели-писатели: первый вариант, второй и std.Io.RwLock держат инвариант" {
try checkRw(conc.rw.ReadersFirst);
try checkRw(conc.rw.WritersFirst);
try checkRw(conc.rw.Std);
}
/// Кто когда вошёл: номер по порядку входа.
fn Order(comptime Lock: type) type {
return struct {
lock: *Lock,
seq: std.atomic.Value(u32) = .init(0),
writer_at: u32 = 0,
reader2_at: u32 = 0,
fn writer(o: *@This()) void {
o.lock.lockWrite();
o.writer_at = o.seq.fetchAdd(1, .acq_rel);
o.lock.unlockWrite();
}
fn reader2(o: *@This()) void {
o.lock.lockRead();
o.reader2_at = o.seq.fetchAdd(1, .acq_rel);
o.lock.unlockRead();
}
};
}
/// Первый читатель внутри, приходит писатель, за ним второй читатель.
/// Возвращает, кто вошёл раньше, когда первый читатель ушёл.
fn whoFirst(comptime Lock: type) !enum { writer, reader } {
var lock: Lock = .init(io);
var o: Order(Lock) = .{ .lock = &lock };
lock.lockRead();
const w = try std.Thread.spawn(.{}, Order(Lock).writer, .{&o});
try std.Io.sleep(io, .fromMilliseconds(50), .awake);
const r = try std.Thread.spawn(.{}, Order(Lock).reader2, .{&o});
try std.Io.sleep(io, .fromMilliseconds(50), .awake);
lock.unlockRead();
w.join();
r.join();
return if (o.writer_at < o.reader2_at) .writer else .reader;
}
test "приоритеты: в первом варианте второй читатель обгоняет писателя, во втором нет" {
try testing.expectEqual(.reader, try whoFirst(conc.rw.ReadersFirst));
try testing.expectEqual(.writer, try whoFirst(conc.rw.WritersFirst));
}
test "badcnt может потерять прибавления, goodcnt, мьютекс и атомик никогда" {
const niters = 1_000_000;
const bad = try conc.badcnt.badcnt(niters);
// Гонка не обязана случиться при каждом запуске, но больше 2n быть не может.
try testing.expect(bad.cnt <= bad.expected);
try testing.expect((try conc.badcnt.goodcnt(io, niters)).ok());
try testing.expect((try conc.badcnt.mutexcnt(io, niters)).ok());
try testing.expect((try conc.badcnt.atomiccnt(niters)).ok());
}
- Два теста
sbufгонят через буфер на четыре слота три производителя и три потребителя, 6000 элементов. Каждый элемент обязан прийти ровно один раз: массивseenиз атомарных счётчиков ловит и потерю (ноль), и дубль (двойка). Счётчики атомарные потому, что потребители пишут вseenодновременно, и обычный+= 1сам стал бы гонкой. - Тест на инвариант читателей и писателей даёт четырём читателям и двум писателям по нескольку тысяч заходов. Внутри замка поток сверяет атомарные счётчики: писатель обязан быть один и без читателей. Заодно
max_readersпоказывает, что читатели действительно входили одновременно. Все три замка проходят один и тот же тест. - Тест приоритетов разобран выше. Паузы по 50 мс нужны, чтобы писатель и второй читатель успели встать в очередь до того, как первый читатель уйдёт.
- Тест
badcntтребует точности только от исправленных версий.
$ zig build test -Dstep=70 --summary all
Build Summary: 6/6 steps succeeded; 5/5 tests passed
$ zig build test --summary all
Build Summary: 18/18 steps succeeded; 59/59 tests passed
Шаг 70 зелёный на macOS и в контейнере с Linux (linux/arm64), полный прогон проекта на macOS тоже зелёный.
Шаг zl: одна куча на несколько потоков
Теперь сквозной проект. До сих пор машина zl исполняла одну программу в одном потоке. Сделаем так, чтобы несколько потоков вычисляли каждый свою форму на одной машине: shared.evalParallel(vm, io, forms, results, options) запускает поток на форму и складывает значения в results. Общим у потоков становится всё: пул cons-ячеек, карта глобальных имён, таблица символов.
Прежде чем писать код, решим, что делим и с каким исполнителем.
- Глобальные имена и символы потоки только читают. Все
defineи всё чтение исходников делаются до запуска,evalParallelполучает готовые формы. Чтение без записи гонки не даёт, и блокировка не нужна. - Пул ячеек меняет каждое выделение. Это единственное место, где потокам нужна договорённость, и весь шаг про него.
- Исполнитель это обход дерева. Всё состояние вычисления у него лежит в кадрах Zig (
eval,apply,evalList), а стек и так читается консервативно, как мы сделали в уроке про сборку мусора. Потоку нечего регистрировать, кроме самого себя. Байткоду пришлось бы дать машину на поток со своими корнями, это домашнее задание. JIT не взят совсем: у машинного кода нет места, где поток согласился бы остановиться, и об этом ниже.
Гонок в шаге две. Первая сидела в машине с урока 14, и её легко не заметить.
Первая гонка: ошибка примитива в поле машины
В уроке 14 примитивы стали функциями с соглашением C, а такая функция не может вернуть ошибку Zig. Ошибку тогда положили в поле машины: примитив пишет vm.failure, вызывающий забирает её через takeFailure и гасит поле. В одном потоке это безупречно. В двух потоках на одной машине поле общее, и получается ровно badcnt: поток A пишет в поле NotAPair, поток B, у которого всё хорошо, вызывает takeFailure раньше A, забирает чужую ошибку и гасит поле. B падает на ровном месте, A молча продолжает с nil вместо ошибки.
Лечение: у каждого потока своя ошибка. Переменная с пометкой threadlocal даёт ровно это:
nil: SymbolId,
};
+/// Ошибка примитива, позванного по соглашению C. Вернуть ошибку Zig такой
+/// вызов не может, поэтому она ждёт здесь, пока вызывающий её не заберёт.
+/// До шага 70 это было поле машины; с шага 70 одну машину делят несколько
+/// потоков, и ошибка одного потока не должна достаться другому, поэтому у
+/// каждого потока своя.
+threadlocal var failure: ?errors.Error = null;
+
pub const Vm = struct {
gpa: std.mem.Allocator,
heap: Heap,
@@
/// Куда печатает примитив `print`. Пусто означает, что печать отключена:
/// так тесты гоняют программы, которым вывод не нужен.
out: ?*std.Io.Writer = null,
- /// Ошибка примитива, позванного по соглашению C. Вернуть ошибку Zig такой
- /// вызов не может, поэтому она ждёт здесь, пока вызывающий её не заберёт.
- failure: ?errors.Error = null,
pub fn init(gpa: std.mem.Allocator) std.mem.Allocator.Error!Vm {
return initWith(gpa, .{});
@@
/// Запомнить ошибку и вернуть nil: его отдаст наружу вход по соглашению C.
pub fn fail(self: *Vm, err: errors.Error) Value {
- self.failure = err;
+ _ = self;
+ failure = err;
return .nil;
}
- /// Забрать ошибку последнего вызова, если она была.
+ /// Забрать ошибку последнего вызова, если она была. Запись только когда
+ /// ошибка есть: примитив зовут на каждом шаге, и пустая запись в общую
+ /// переменную гоняла бы её линию кэша между ядрами.
pub fn takeFailure(self: *Vm) ?errors.Error {
- defer self.failure = null;
- return self.failure;
+ _ = self;
+ const err = failure orelse return null;
+ failure = null;
+ return err;
}
pub fn isTrue(v: Value) bool {
Сигнатуры fail и takeFailure не изменились, параметр self остался, чтобы не трогать ни одного вызывающего. Тест урока 14 читал поле напрямую, теперь читает через функцию:
// `primitives.call` забирает ошибку сам и возвращает её как ошибку Zig.
const idx = primitives.indexOf("+");
try testing.expectError(error.NotANumber, primitives.call(&vm, idx, try reader.readOne(&vm, "(a 1)")));
- try testing.expectEqual(@as(?zl.errors.Error, null), vm.failure);
+ try testing.expectEqual(@as(?zl.errors.Error, null), vm.takeFailure());
}
/// Функция с соглашением C, которая возвращает свой аргумент. Если значение
@@
try expectPrinted(&vm, try m.runSource("(+ 1 2)", .loop), "3");
try testing.expectError(error.NotAPair, m.runSource("(cdr 7)", .labeled));
try testing.expectError(error.NotAPair, zl.eval.evalSource(&vm, "(cdr 7)"));
- try testing.expectEqual(@as(?zl.errors.Error, null), vm.failure);
+ try testing.expectEqual(@as(?zl.errors.Error, null), vm.takeFailure());
}
test "кадр замыкания: функция под аргументами, результат на её месте" {
Посмотрим, что это значит для машинного кода. В уроке 14 мы разбирали вход примитива car и видели запись ошибки как movw %ax, 0xbc(%rdi): два байта по смещению от адреса машины. Соберём zl до шага 70 и после (zig build asm -Dtarget=x86_64-linux -Doptimize=ReleaseFast) и сравним тот же вход. До шага 70, путь ошибки:
<primitives.native__struct_32866.entry>: # car, до шага 70
...
movw %ax, 0x144(%rdi) # vm.failure = код ошибки, rdi это машина
xorl %eax, %eax # вернуть nil
popq %rbp
retq
Смещение выросло с 0xbc до 0x144: с урока 14 машина обросла полями. После шага 70:
<primitives.native__struct_32864.entry>: # car, после шага 70
pushq %rbp
movq %rsp, %rbp
testq %rsi, %rsi # args == nil?
sete %cl
testb $0xf, %sil # тег не пара?
setne %dl
movw $0x28, %ax # код ArityMismatch наготове
orb %cl, %dl
jne <ошибка>
cmpq $0x0, 0x8(%rsi) # cdr списка аргументов не nil: аргументов больше одного
jne <ошибка>
movq (%rsi), %rcx # сам аргумент
testq %rcx, %rcx
sete %dl
testb $0xf, %cl
setne %sil
movw $0x27, %ax # теперь наготове NotAPair
orb %dl, %sil
jne <ошибка>
movq (%rcx), %rax # его car
popq %rbp
retq
<ошибка>:
movw %ax, %fs:-0x4000a # failure потока = код ошибки
xorl %eax, %eax
popq %rbp
retq
%rdi больше не нужен вовсе: машина в пути ошибки не участвует. Запись идёт по адресу %fs:-0x4000a. На x86-64 в Linux сегментный регистр %fs у каждого потока указывает на его собственный блок управления, а переменные threadlocal исполняемого файла лежат прямо перед ним, по смещениям, которые компоновщик знает заранее. Это самая быстрая из четырёх моделей доступа к TLS, local exec: одна инструкция, без вызова. В разделяемой библиотеке смещение заранее неизвестно, и доступ идёт через __tls_get_addr, по той же схеме с GOT, что в уроке про динамическую компоновку.
Откуда такое большое смещение, 256 КБ с хвостиком? Спросим таблицу символов:
$ objdump -t zig-out/bin/zl-llvm | grep tbss | sort
0000000000000000 l .tbss 0000000000000008 Io.Threaded.Thread.current
0000000000000008 l .tbss 0000000000000008 debug.panic_stage
0000000000000010 l .tbss 0000000000000004 Thread.LinuxThreadImpl.tls_thread_id.0
0000000000000014 l .tbss 0000000000000001 Thread.LinuxThreadImpl.tls_thread_id.1
0000000000000016 l .tbss 0000000000000002 vm.failure
0000000000000018 l .tbss 0000000000040000 Thread.maybeAttachSignalStack.global.signal_stack
0000000000040018 l .tbss 0000000000000004 heap.SmpAllocator.thread_index
vm.failure лежит по смещению 0x16 в блоке, а целиком блок занимает 0x40020 байт, почти всё это signal_stack: стандартная библиотека даёт каждому потоку свой запасной стек для обработчиков сигналов на 256 КБ (тот самый sigaltstack из урока про сигналы). Указатель потока смотрит на конец блока, поэтому 0x16 - 0x40020 = -0x4000a. Всё сходится до байта.
На aarch64 (-Dtarget=aarch64-linux) указатель потока хранится в системном регистре TPIDR_EL0:
<ошибка>: # car, aarch64, после шага 70
mrs x9, TPIDR_EL0 # указатель потока
mov x0, xzr # вернуть nil
add x9, x9, #0x0, lsl #12 # смещение переменной, старшая часть
add x9, x9, #0x28 # младшая часть
strh w8, [x9] # failure потока = код ошибки
ldp x29, x30, [sp], #0x10
ret
Теперь сторона вызывающего, конец bytecode.Machine.call после вызова примитива. До шага 70 takeFailure читала поле и гасила его каждый раз:
callq *0x1017110(,%r14,8) # natives[index]
movq %rax, %rcx # результат примитива
movzwl 0x144(%r15), %eax # прочитать vm.failure
movw $0x0, 0x144(%r15) # и погасить, даже если там ноль
testw %ax, %ax
jne <ошибка>
После шага 70:
callq *0x10170e0(,%r14,8) # natives[index]
movq %rax, %rcx # результат примитива
movzwl %fs:-0x4000a, %eax # прочитать failure потока
testw %ax, %ax
je <дальше> # ноль: писать нечего
movw $0x0, %fs:-0x4000a # погасить только на пути ошибки
Это вторая половина правки: takeFailure больше не пишет ноль, если там и так ноль. Раньше каждый вызов примитива на любом потоке писал в поле машины, то есть в одну и ту же линию кэша. В общей машине это значило бы, что линия с полем failure на каждом car перебегает из кэша одного ядра в кэш другого: запись требует, чтобы ядро владело линией единолично, и все остальные копии становятся недействительными (как устроены линии, мы разбирали в уроке про организацию кэша). После переезда в threadlocal линия у каждого потока своя, и пустая запись уже никому не мешает, но и пользы от неё нет: запись по пути без ошибок была лишней с самого начала.
Урок 14 при этом остаётся верным: он описывает машину своего шага, где поток был один. Проект всегда хранит финал, и если собрать zl сейчас, ассемблер будет таким, как здесь.
Вторая гонка: свободный список
Свободный список ячеек мы построили в уроке про аллокатор: свободные ячейки связаны через cdr, голова лежит в куче. Выделение снимает голову. Это два действия: прочитать голову вместе со следующей и записать следующую на место головы. Посмотри, как это похоже на L и S из badcnt. Если два потока оба прочитали одну голову, оба запишут одну и ту же следующую и оба получат одну и ту же ячейку. Два cons в разных потоках положат в неё свои car и cdr, и один из них будет жить с чужими данными.
Чтобы показать гонку без надежды на планировщик, эталон раскладывает снятие головы на два шага:
pub const RacyList = struct {
head: std.atomic.Value(?*Cell),
pub const Snapshot = struct { cell: *Cell, next: ?*Cell };
/// Шаг первый: что сейчас в голове и что за ней.
pub fn load(self: *RacyList) ?Snapshot { ... }
/// Шаг второй: новая голова. Ячейку из снимка вызывающий считает своей.
pub fn commit(self: *RacyList, s: Snapshot) *Cell { ... }
};
Каждый шаг атомарен сам по себе, пара нет. Первый тест шага разыгрывает чередование руками: A читает, B читает, A пишет, B пишет. Оба получили cells[0], голова сдвинулась на одну ячейку, хотя сняли две. Это одна конкретная траектория через небезопасную зону, выбранная тестом, а не расписанием, поэтому тест детерминирован. На настоящих потоках гонку показывает бенч, ниже.
Заметь, что атомарность отдельных шагов ничего не спасает. Можно сделать неделимой пару целиком (сравнение с обменом: записать новую голову, только если голова всё ещё та, что я читал), но на этом пути ждёт знаменитая проблема ABA. Мы пойдём проще и быстрее.
Лечение: пачки под мьютексом
Самое прямое лечение это мьютекс вокруг каждого снятия головы. Оно верное и очень медленное: cons вызывается на каждом шаге интерпретатора, и все потоки будут стоять в очереди к одному мьютексу. Мы измерим это ниже.
Лечение получше: поток берёт у общего списка сразу пачку ячеек под мьютексом, а дальше выдаёт из своего локального пула, Mutator.local, вообще без блокировки и без атомарных операций. Локальный пул никто, кроме хозяина, не видит, значит, и гонки на нём нет. Размер пачки задаёт Options.batch: 256 по умолчанию, а пачка в одну ячейку это и есть мьютекс на каждое выделение, удобно для сравнения. Так устроены серьёзные многопоточные аллокаторы вроде tcmalloc: у каждого потока свой кэш свободных блоков, а общая структура трогается редко.
Первая версия эталона нарезала пачку прямо под блокировкой: шла по свободному списку на batch ячеек и ставила им биты занятости. Критическая секция длиной в пачку значит, что сумма всех критических секций пропорциональна числу выделенных ячеек, как бы ни была велика пачка, и блокировка сама становилась последовательной частью программы. Теперь пачки режутся заранее, при остановленном мире: голова пачки связана со следующей пачкой через car, ячейки внутри пачки через cdr, биты поставлены сразу. Снять пачку под мьютексом стоит O(1).
Остановка мира
Пул конечен, и когда пачки кончились, нужна сборка мусора. Сборщик из урока 59 помечает всё, до чего можно дотянуться от корней, включая слова стека, и освобождает остальное. С несколькими потоками к этому добавляется требование: пока идёт пометка, никто не должен менять кучу. Иначе поток может переложить единственную ссылку на ячейку из непомеченного места в уже просмотренное, и сборщик освободит живое. Поэтому сборка останавливает мир.
Останавливать поток посреди произвольной инструкции нельзя: в регистрах может лежать единственная ссылка на ячейку, которую сборщик не увидит. Поток останавливается сам, в безопасной точке, где всё, что он держит, лежит там, где сборщик это найдёт. У нас безопасная точка это вход в refill, медленный путь выделения. Механизм такой:
- Поток, которому не хватило пачки, поднимает флаг
stop, вычитает себя из счётчика работающихrunningи ждёт на условной переменной, покаrunningне станет нулём. - Остальные потоки проверяют флаг на каждом выделении, без блокировки. Заметив его, они входят в
refill. Функция помеченаnoinlineи первым делом сбрасывает сохраняемые регистры в массив в своём кадре, тем же приёмом, чтоspillRegistersв сборщике урока 59. Адрес массива становится вершиной стека потока для сборщика. Потом поток встаёт: уменьшаетrunning, будит ждущих и засыпает, пока флаг поднят. - Когда
runningдошёл до нуля, сборщик работает один. Он снимает с учёта невыданные ячейки (пачки и остатки локальных пулов, их цепочки черезcdrвыглядят как списки, и пометка сочла бы их живыми), метит от глобальных имён, явных корней и стеков всех подключённых потоков, каждый от сохранённой вершины до своего дна, делает sweep и заново режет свободное на пачки. - Флаг опускается,
broadcastбудит всех, потоки продолжают.
Тут видно, зачем нужна условная переменная, а не семафор. Ждать надо условия над счётчиком (running == 0) и условия над флагом (stop опущен), а будить всех сразу. Семафором это тоже можно собрать (сборщик делает P столько раз, сколько потоков, каждый вставший поток делает V), но потоки подключаются и отключаются посреди прогона, и число P пришлось бы угадывать заранее. С условной переменной и счётчиком под мьютексом такой задачи нет.
У безопасной точки в refill есть предел: поток, который долго не выделяет, мир не остановит никогда, и сборщик будет его ждать. Обходу дерева это не грозит, он выделяет на каждом вызове (окружение и список аргументов). А вот функция last в JIT идёт по списку через car и cdr в регистрах и вызывает сама себя, ни разу не выделяя. Такой поток остановки не дождётся. JVM и Go лечат это опросом флага в прологе каждой функции и на обратных переходах циклов, и ровно это пришлось бы добавить в наш JIT.
И последняя деталь: флаг stop лежит на своей линии кэша (Flag с align(std.atomic.cache_line)). Его читают все потоки на каждом выделении, а соседние поля Shared пишутся на каждой пачке. На общей линии каждая такая запись выбивала бы флаг из кэшей всех ядер. Это ложное разделение, и в следующем уроке мы измерим его цену.
shared.zig целиком
//! Общая куча для нескольких потоков: шаг 70.
//!
//! Одну машину `Vm` делят несколько потоков, каждый вычисляет свою форму
//! интерпретатором по дереву. Общее у них всё: пул ячеек, глобальные имена,
//! таблица символов. Глобальные имена и символы потоки только читают (все
//! `define` делаются до запуска), а вот пул меняет каждое выделение, и это
//! единственное место, где потокам нужна договорённость.
//!
//! Договорённость в три слоя:
//!
//! 1. Свободный список без блокировки теряет ячейки в гонке: два потока
//! читают одну голову, и одна ячейка уходит обоим (`RacyList`).
//! 2. Под мьютексом гонки нет, но мьютекс берётся на каждую ячейку. Лечение:
//! поток берёт у общего списка пачку (`Options.batch`) под блокировкой
//! и дальше выдаёт из своего локального пула без неё. Пачка в одну
//! ячейку это и есть «мьютекс на каждое выделение».
//! 3. Сборка останавливает мир. Поток, которому не хватило ячеек, поднимает
//! флаг `stop` и ждёт, пока остальные дойдут до безопасной точки. Точка
//! это вход в `refill`: поток сбрасывает там регистры в свой кадр, и
//! сборщик читает его стек от этого кадра до дна. Корни это глобальные
//! имена, явные списки и стеки всех подключённых потоков. После sweep
//! флаг опускается, и все продолжают.
//!
//! Флаг `stop` проверяется на каждом выделении без блокировки. Поток, который
//! долго не выделяет, мир не остановит никогда, но интерпретатор по дереву
//! выделяет на каждом вызове (окружение и список аргументов), так что
//! безопасная точка у него всегда близко.
const std = @import("std");
const errors = @import("errors.zig");
const eval = @import("eval.zig");
const gc = @import("gc.zig");
const heap_mod = @import("heap.zig");
const sys = @import("sys.zig");
const value = @import("value.zig");
const Vm = @import("vm.zig").Vm;
const Allocator = std.mem.Allocator;
const Cell = value.Cell;
const Heap = heap_mod.Heap;
const Io = std.Io;
const Value = value.Value;
/// Следующая ячейка свободной цепочки: она лежит в `cdr`.
fn next(cell: *const Cell) ?*Cell {
return if (cell.cdr.isCons()) cell.cdr.asCell() else null;
}
// --- 1. Гонка ---
/// Свободный список без блокировки. Снять голову это два шага: прочитать
/// голову вместе со следующей и записать следующую на место головы. Каждый
/// шаг атомарен сам по себе, пара нет: между ними другой поток успевает
/// прочитать ту же голову. Шаги открыты, чтобы тест разыграл чередование
/// руками, без надежды на планировщик.
pub const RacyList = struct {
head: std.atomic.Value(?*Cell),
pub const Snapshot = struct { cell: *Cell, next: ?*Cell };
pub fn init(first: ?*Cell) RacyList {
return .{ .head = .init(first) };
}
/// Шаг первый: что сейчас в голове и что за ней.
pub fn load(self: *RacyList) ?Snapshot {
const cell = self.head.load(.monotonic) orelse return null;
return .{ .cell = cell, .next = next(cell) };
}
/// Шаг второй: новая голова. Ячейку из снимка вызывающий считает своей.
pub fn commit(self: *RacyList, s: Snapshot) *Cell {
self.head.store(s.next, .monotonic);
return s.cell;
}
pub fn take(self: *RacyList) ?*Cell {
return self.commit(self.load() orelse return null);
}
};
// --- 2 и 3. Общий пул и остановка мира ---
pub const Options = struct {
/// Сколько ячеек поток берёт за одну блокировку. Единица значит мьютекс
/// на каждую ячейку.
batch: usize = 256,
/// Стек потока. Обход дерева рекурсивен на стеке Zig.
stack_size: usize = 64 << 20,
};
/// Поток-вычислитель, мутатор в терминах сборщиков: всё, что куча знает о
/// потоке. Живёт в кадре самого потока.
pub const Mutator = struct {
/// Локальный пул: пачка, снятая с общего списка. Цепочка через `cdr`.
local: ?*Cell = null,
/// Дно стека потока, старший адрес.
stack_bottom: usize = 0,
/// Вершина стека в безопасной точке: адрес слепка регистров в кадре
/// `refill`. Верна, пока поток стоит там.
stack_top: usize = 0,
/// Сколько ячеек поток выделил и сколько раз ходил за пачкой.
allocated: usize = 0,
refills: usize = 0,
};
/// Мутатор текущего потока. У потока, который не подключён, его нет, и
/// выделять из общей кучи он не может.
threadlocal var current: ?*Mutator = null;
/// Флаг на собственной линии кэша. Его читают все потоки на каждом
/// выделении, а поля рядом с ним пишутся на каждой пачке: на общей линии
/// каждая запись выбивала бы флаг из кэша всех ядер.
const Flag = struct {
raised: std.atomic.Value(bool) align(std.atomic.cache_line) = .init(false),
};
pub const Shared = struct {
heap: *Heap,
io: Io,
batch: usize,
/// Свободные ячейки, заранее нарезанные пачками. Пачка это цепочка через
/// `cdr`, а пачки связаны через `car` головы. Нарезка идёт при
/// остановленном миру, поэтому снять пачку под блокировкой стоит O(1).
chunks: ?*Cell = null,
mutex: Io.Mutex = .init,
/// Одна переменная условия на оба направления: сборщик ждёт, пока все
/// остановятся, остальные ждут конца сборки.
changed: Io.Condition = .init,
/// Просьба остановиться. Читается без блокировки на каждом выделении.
stop: Flag = .{},
/// Подключённые потоки. Меняется и читается под `mutex`.
mutators: std.ArrayList(*Mutator) = .empty,
/// Сколько подключённых потоков сейчас работает, а не стоит.
running: usize = 0,
/// Итоги ушедших потоков.
allocated: usize = 0,
refills: usize = 0,
/// Сколько всего стоял мир: от просьбы остановиться до отпускания.
stopped_ns: u64 = 0,
/// Забрать у кучи все свободные ячейки и нарезать пачками. Пока общий
/// режим включён, одинокий поток кучей пользоваться не может.
pub fn init(heap: *Heap, io: Io, batch: usize) Shared {
std.debug.assert(batch > 0);
var self: Shared = .{ .heap = heap, .io = io, .batch = batch };
self.carve();
return self;
}
/// Вернуть куче всё, что не выдано: пачки и пулы ушедших потоков.
pub fn deinit(self: *Shared) void {
std.debug.assert(self.mutators.items.len == 0);
self.unreserve();
self.heap.threadFreeList();
self.mutators.deinit(self.heap.backing);
}
/// Подключить текущий поток. Зовётся из самого потока: дно стека
/// узнаётся по адресу его локальной переменной.
pub fn attach(self: *Shared, m: *Mutator) Allocator.Error!void {
var probe: usize = 0;
m.* = .{ .stack_bottom = sys.stackBottom(@intFromPtr(&probe)) orelse return error.OutOfMemory };
self.mutex.lockUncancelable(self.io);
defer self.mutex.unlock(self.io);
try self.mutators.append(self.heap.backing, m);
self.running += 1;
current = m;
}
/// Отключить текущий поток. Остаток локального пула уходит обратно
/// пачкой. Всё, что поток вычислил и хочет отдать, к этому моменту
/// должно лежать в корне: его стек больше никто не читает.
pub fn detach(self: *Shared, m: *Mutator) void {
self.mutex.lockUncancelable(self.io);
defer self.mutex.unlock(self.io);
if (m.local) |head| self.pushChunk(head);
m.local = null;
for (self.mutators.items, 0..) |item, i| {
if (item == m) {
_ = self.mutators.swapRemove(i);
break;
}
}
self.allocated += m.allocated;
self.refills += m.refills;
self.running -= 1;
// Сборщик мог ждать именно этот поток.
self.changed.broadcast(self.io);
current = null;
}
/// Выделение в общем режиме. Быстрый путь без блокировки и без
/// атомарных операций: снять голову своего пула. В `refill` уходим,
/// когда пул пуст или просят остановиться.
pub fn take(self: *Shared) Allocator.Error!*Cell {
const m = current.?;
if (m.local == null or self.stop.raised.load(.monotonic)) try self.refill(m);
const cell = m.local.?;
m.local = next(cell);
m.allocated += 1;
return cell;
}
/// Медленный путь, он же безопасная точка. `noinline`: слепок регистров
/// обязан лежать в собственном кадре, который жив всё время ожидания.
noinline fn refill(self: *Shared, m: *Mutator) Allocator.Error!void {
var regs: [gc.saved_count]usize = undefined;
gc.spillRegisters(®s);
m.stack_top = @intFromPtr(®s);
self.mutex.lockUncancelable(self.io);
defer self.mutex.unlock(self.io);
while (true) {
if (self.stop.raised.load(.monotonic)) {
self.park();
continue;
}
if (m.local != null) return;
if (self.popChunk()) |head| {
m.local = head;
m.refills += 1;
return;
}
try self.stopTheWorld();
}
}
/// Встать в безопасной точке и дождаться конца сборки. Под `mutex`.
fn park(self: *Shared) void {
self.running -= 1;
self.changed.broadcast(self.io);
while (self.stop.raised.load(.monotonic)) self.changed.waitUncancelable(self.io, &self.mutex);
self.running += 1;
}
/// Остановить всех, собрать, отпустить. Под `mutex`; ожидание его
/// отпускает, и остальные потоки успевают войти в `refill` и встать.
fn stopTheWorld(self: *Shared) Allocator.Error!void {
const started = Io.Timestamp.now(self.io, .awake);
self.stop.raised.store(true, .monotonic);
self.running -= 1;
defer {
const stopped = started.durationTo(Io.Timestamp.now(self.io, .awake));
self.stopped_ns += @intCast(stopped.nanoseconds);
self.stop.raised.store(false, .monotonic);
self.running += 1;
self.changed.broadcast(self.io);
}
while (self.running > 0) self.changed.waitUncancelable(self.io, &self.mutex);
try self.collect();
if (self.chunks == null) return error.OutOfMemory;
}
/// Сборка при остановленном мире. Невыданные ячейки (пачки и остатки
/// пулов) сначала снимаются с учёта: они отмечены занятыми, и их
/// цепочки через `cdr` выглядят как списки, так что пометка сочла бы
/// их живыми. После sweep свободное снова режется на пачки.
fn collect(self: *Shared) Allocator.Error!void {
const h = self.heap;
self.unreserve();
try gc.markValues(h);
for (self.mutators.items) |m| try gc.markRange(h, m.stack_top, m.stack_bottom);
_ = gc.sweep(h);
self.carve();
}
/// Нарезать свободный список кучи на пачки по `batch` ячеек. Ячейки пачек
/// отмечаются занятыми сразу: тогда выдача из пула не трогает общую
/// битовую карту, и ей не нужна ни блокировка, ни атомарная запись.
/// `used` при этом не растёт, он считает только выданное.
fn carve(self: *Shared) void {
const h = self.heap;
var rest = h.free_list;
h.free_list = null;
while (rest) |head| {
var last = head;
h.live.set(h.indexOf(head));
var n: usize = 1;
while (n < self.batch) : (n += 1) {
last = next(last) orelse break;
h.live.set(h.indexOf(last));
}
rest = next(last);
last.cdr = .nil;
self.pushChunk(head);
}
}
/// Снять с учёта все невыданные ячейки: пачки и остатки пулов. После
/// этого `live` снова значит «выдано», а `used` равен числу выданных.
fn unreserve(self: *Shared) void {
const h = self.heap;
while (self.popChunk()) |head| self.unreserveChain(head);
for (self.mutators.items) |m| {
self.unreserveChain(m.local);
m.local = null;
}
h.used = h.live.count();
h.peak = @max(h.peak, h.used);
}
fn unreserveChain(self: *Shared, first: ?*Cell) void {
var cell = first;
while (cell) |c| {
self.heap.live.unset(self.heap.indexOf(c));
cell = next(c);
}
}
fn pushChunk(self: *Shared, head: *Cell) void {
head.car = if (self.chunks) |c| .fromCell(c) else .nil;
self.chunks = head;
}
fn popChunk(self: *Shared) ?*Cell {
const head = self.chunks orelse return null;
self.chunks = if (head.car.isCons()) head.car.asCell() else null;
head.car = .nil;
return head;
}
};
/// Итог параллельного вычисления.
pub const Report = struct {
/// Сколько раз за прогон останавливали мир.
collections: usize,
/// Сколько ячеек выделили все потоки и сколько раз брали блокировку за ними.
allocated: usize,
refills: usize,
/// Сколько всего стоял мир: ожидание остановки плюс сама сборка.
stopped_ns: u64,
/// От запуска первого потока до конца последнего, без нарезки пачек.
elapsed_ns: u64,
};
/// Вычислить каждую форму в своём потоке на общей куче машины и сложить
/// значения в `results`. Все `define` должны быть сделаны до вызова: таблицу
/// глобальных имён потоки только читают. Формы и результаты на время работы
/// лежат в явном корне, поэтому переживают любую сборку.
pub fn evalParallel(
vm: *Vm,
io: Io,
forms: []const Value,
results: []Value,
options: Options,
) (errors.Error || std.Thread.SpawnError)!Report {
std.debug.assert(forms.len == results.len);
const gpa = vm.gpa;
const n = forms.len;
// Первая половина корня это формы, вторая результаты.
var roots: heap_mod.RootList = .empty;
defer roots.deinit(gpa);
try roots.appendSlice(gpa, forms);
try roots.appendNTimes(gpa, .nil, n);
try vm.heap.pushRoots(&roots);
defer vm.heap.popRoots();
const failures = try gpa.alloc(?errors.Error, n);
defer gpa.free(failures);
@memset(failures, null);
const threads = try gpa.alloc(std.Thread, n);
defer gpa.free(threads);
var shared: Shared = .init(&vm.heap, io, options.batch);
defer shared.deinit();
const collections_before = vm.heap.collections;
vm.heap.shared = &shared;
defer vm.heap.shared = null;
const started = Io.Timestamp.now(io, .awake);
var spawned: usize = 0;
defer for (threads[0..spawned]) |t| t.join();
for (0..n) |i| {
threads[i] = try std.Thread.spawn(
.{ .stack_size = options.stack_size },
worker,
.{ &shared, vm, &roots.items[i], &roots.items[n + i], &failures[i] },
);
spawned += 1;
}
for (threads) |t| t.join();
spawned = 0;
const elapsed = started.durationTo(Io.Timestamp.now(io, .awake));
for (failures) |f| if (f) |err| return err;
@memcpy(results, roots.items[n..]);
return .{
.collections = vm.heap.collections - collections_before,
.allocated = shared.allocated,
.refills = shared.refills,
.stopped_ns = shared.stopped_ns,
.elapsed_ns = @intCast(elapsed.nanoseconds),
};
}
fn worker(shared: *Shared, vm: *Vm, form: *const Value, result: *Value, failure: *?errors.Error) void {
var m: Mutator = .{};
shared.attach(&m) catch |err| {
failure.* = err;
return;
};
defer shared.detach(&m);
// Результат ложится в корень до `detach`: пока поток подключён и не
// стоит, мир без него не остановится, так что сборка его не опередит.
result.* = eval.eval(vm, form.*, .nil) catch |err| {
failure.* = err;
return;
};
}
Две вещи, которые легко пропустить при чтении. current это ещё одна threadlocal: через неё take находит мутатор своего потока, не получая его параметром, потому что heap.cons о потоках ничего не знает. И worker кладёт результат в корень до detach: пока поток подключён и не стоит, мир без него не остановится, так что сборка его не опередит, а после detach его стек уже никто не читает. Главный поток к куче не подключён, его стек сборщик не видит, поэтому evalParallel держит формы и результаты в явном корне.
Правки кучи и сборщика
Куча учится отдавать выделение общему пулу. Счётчик выделенных ячеек переехал из cons в однопоточный take: в общем режиме каждый поток считает свои ячейки сам, в Mutator.allocated, иначе все потоки писали бы в одну линию кэша на каждом выделении.
const std = @import("std");
const gc = @import("gc.zig");
+const shared = @import("shared.zig");
const symbols = @import("symbols.zig");
const value = @import("value.zig");
@@
/// Рабочий стек пометки. Живёт между сборками, чтобы не выделяться заново.
mark_stack: std.ArrayList(*Cell) = .empty,
+ /// Общий режим шага 70: ячейки выдаёт `shared.Shared`, и ему же
+ /// принадлежит сборка. Пусто значит, что кучей пользуется один поток.
+ shared: ?*shared.Shared = null,
+
collections: usize = 0,
freed_total: usize = 0,
used: usize = 0,
peak: usize = 0,
- /// Сколько пар и замыканий выделено за жизнь кучи.
+ /// Сколько пар и замыканий выделено за жизнь кучи. В общем режиме эти
+ /// счётчики стоят: каждый поток считает свои ячейки сам (`Mutator`),
+ /// иначе потоки дрались бы за одну линию кэша на каждом выделении.
cells_allocated: usize = 0,
closures_allocated: usize = 0,
@@
/// собрать мусор. Сборка консервативная: нас позвали из глубины `eval`,
/// и живые значения лежат в кадрах Zig, про которые знает только стек.
fn take(self: *Heap) Allocator.Error!*Cell {
+ if (self.shared) |s| return s.take();
if (self.free_list == null) _ = try gc.collect(self, .conservative);
const cell = self.free_list orelse return error.OutOfMemory;
self.free_list = if (cell.cdr.isCons()) cell.cdr.asCell() else null;
self.live.set(self.indexOf(cell));
self.used += 1;
self.peak = @max(self.peak, self.used);
+ self.cells_allocated += 1;
return cell;
}
@@
pub fn cons(self: *Heap, a: Value, d: Value) Allocator.Error!Value {
const cell = try self.take();
cell.* = .{ .car = a, .cdr = d };
- self.cells_allocated += 1;
return .fromCell(cell);
}
@@
pub fn closure(self: *Heap, params: Value, body: Value, env: Value) Allocator.Error!Value {
const rest = try self.cons(body, env);
const head = try self.cons(params, rest);
- self.closures_allocated += 1;
+ if (self.shared == null) self.closures_allocated += 1;
return .fromClosure(head.asCell());
}
Сборщик отдаёт наружу пометку участка стека, пометку от глобальных имён и sweep, а цикл скана стека выносится в markRange, чтобы читать им и чужие стеки:
// --- Mark ---
-fn markValues(heap: *Heap) Allocator.Error!void {
+pub fn markValues(heap: *Heap) Allocator.Error!void {
var globals = heap.globals.valueIterator();
while (globals.next()) |v| try markFrom(heap, v.*);
for (heap.roots.items) |list| {
@@
/// может лежать единственная ссылка на ячейку: вызывающий положил её туда и
/// вправе рассчитывать, что она переживёт вызов. Остальные регистры
/// вызывающий перед вызовом сам сбросил на стек.
-const saved_count = switch (builtin.cpu.arch) {
+pub const saved_count = switch (builtin.cpu.arch) {
.x86_64 => 6,
.aarch64 => 11,
else => @compileError("консервативный скан: опиши сохраняемые регистры этой архитектуры"),
@@
/// Сложить сохраняемые регистры в массив. Это тот же приём, которым `setjmp`
/// снимает слепок регистров, только без libc.
-inline fn spillRegisters(regs: *[saved_count]usize) void {
+pub inline fn spillRegisters(regs: *[saved_count]usize) void {
switch (builtin.cpu.arch) {
.x86_64 => asm volatile (
\\ movq %%rbx, (%%rax)
@@
if (heap.stack_bottom == null) heap.stack_bottom = sys.stackBottom(top);
// Без дна стека скан невозможен, а сборка без скана освободила бы живое.
const bottom = heap.stack_bottom orelse return error.OutOfMemory;
+ try markRange(heap, top, bottom);
+}
- var at = top;
+/// Пометить всё, на что показывают слова участка стека `[top, bottom)`.
+/// С шага 70 так читаются и стеки чужих потоков, остановленных сборщиком.
+pub fn markRange(heap: *Heap, top: usize, bottom: usize) Allocator.Error!void {
+ var at = std.mem.alignForward(usize, top, @sizeOf(usize));
while (at < bottom) : (at += @sizeOf(usize)) {
const word = @as(*const usize, @ptrFromInt(at)).*;
try markWord(heap, word);
@@
/// Всё занятое и не помеченное освобождается, пометки гасятся, свободный
/// список собирается заново по возрастанию адресов.
-fn sweep(heap: *Heap) usize {
+pub fn sweep(heap: *Heap) usize {
var freed: usize = 0;
for (heap.cells, 0..) |*cell, i| {
if (!heap.live.isSet(i)) continue;
pub const programs = @import("programs.zig");
pub const reader = @import("reader.zig");
pub const repl = @import("repl.zig");
+pub const shared = @import("shared.zig");
pub const symbols = @import("symbols.zig");
pub const sys = @import("sys.zig");
pub const value = @import("value.zig");
const std = @import("std");
/// Шаги проекта в порядке появления.
-const steps = [_][]const u8{ "05", "06", "13", "14", "15", "19", "32", "35", "40", "56", "58", "59", "59_dsw" };
+const steps = [_][]const u8{ "05", "06", "13", "14", "15", "19", "32", "35", "40", "56", "58", "59", "59_dsw", "70" };
/// Программы на zl, которые вшиваются в модуль через `@embedFile`.
const zl_programs = [_][]const u8{
Комментарий в src/lower.zig, который говорил про поле vm.failure, поправлен на «в машине (Vm.fail)», код там не менялся.
Бенч: гонка и общая куча числом
Режим zig build bench -- threads делает две вещи. Сначала раздаёт свободный список на миллион ячеек нескольким потокам без блокировки и под мьютексом и считает, сколько выдач было сверх числа ячеек: каждая лишняя это ячейка, доставшаяся двум потокам. Потом гоняет (fib 23) обходом дерева в 1, 2, 4 и 8 потоках на общей куче, с пачкой в одну ячейку и в 256, на большом пуле, где сборок нет, и на маленьком, где они идут постоянно.
//! zig build bench все исполнители, которые идут на этой машине
//! zig build bench -- tree loop только перечисленные
//! zig build bench -- cmov развилка в JIT: переход против cmov (x86-64)
//! zig build bench -- cache обход списка при разных раскладках ячеек
+//! zig build bench -- threads гонка на свободном списке и общая куча в 1, 2, 4, 8 потоках
//!
//! Элемент нагрузки `fib` это один вызов `fib`, нагрузки `last` одна
//! cons-ячейка. CPE снимается измерителем из урока про CPE: минимум из серии
//! на каждом размере и прямая по наименьшим квадратам, наклон прямой это
@@
std.mem.doNotOptimizeAway(result);
}
};
-const Mode = enum { executors, cmov, cache };
+const Mode = enum { executors, cmov, cache, threads };
pub fn main(init: std.process.Init) !void {
var out_buf: [4096]u8 = undefined;
var stdout = std.Io.File.stdout().writerStreaming(init.io, &out_buf);
@@
if (std.mem.eql(u8, arg, "cmov")) {
mode = .cmov;
} else if (std.mem.eql(u8, arg, "cache")) {
mode = .cache;
+ } else if (std.mem.eql(u8, arg, "threads")) {
+ mode = .threads;
} else {
chosen[count] = std.meta.stringToEnum(Executor, arg) orelse {
try out.print("исполнитель: tree, loop, labeled или jit, а не {s}\n", .{arg});
try out.flush();
@@
const done = switch (mode) {
.executors => measureAll(io, out, chosen),
.cmov => measureSelect(io, out),
.cache => measureCache(io, out),
+ .threads => measureThreads(io, out),
};
done catch |err| std.debug.panic("бенч: {t}", .{err});
}
@@
try out.print("{t:<12} {d:>8.1}\n", .{ layout, series.items()[0].cycles / @as(f64, @floatFromInt(n)) });
try out.flush();
}
}
+
+// --- Потоки на одной куче ---
+
+const thread_counts = [_]usize{ 1, 2, 4, 8 };
+
+/// Список из n ячеек, связанных через `cdr`, как свободный список пула.
+fn freeChain(gpa: std.mem.Allocator, n: usize) ![]zl.value.Cell {
+ const cells = try gpa.alloc(zl.value.Cell, n);
+ for (cells, 0..) |*c, i| c.* = .{ .car = .nil, .cdr = if (i + 1 < n) .fromCell(&cells[i + 1]) else .nil };
+ return cells;
+}
+
+/// Раздача без блокировки или под мьютексом. Каждый поток снимает ячейки,
+/// пока список не кончится, и считает, сколько взял.
+const Race = struct {
+ list: zl.shared.RacyList,
+ io: std.Io,
+ locked: bool,
+ mutex: std.Io.Mutex = .init,
+
+ fn take(self: *Race) ?*zl.value.Cell {
+ if (!self.locked) return self.list.take();
+ self.mutex.lockUncancelable(self.io);
+ defer self.mutex.unlock(self.io);
+ return self.list.take();
+ }
+
+ fn drain(self: *Race, got: *usize) void {
+ while (self.take()) |_| got.* += 1;
+ }
+};
+
+/// Сколько ячеек выдано сверх их числа: каждая лишняя выдача это ячейка,
+/// которая досталась двум потокам.
+fn raceExtra(io: std.Io, cells: []zl.value.Cell, threads: usize, locked: bool) !usize {
+ var race: Race = .{ .list = .init(&cells[0]), .io = io, .locked = locked };
+ var got: [8]usize = @splat(0);
+ var handles: [8]std.Thread = undefined;
+ for (0..threads) |i| handles[i] = try std.Thread.spawn(.{}, Race.drain, .{ &race, &got[i] });
+ for (handles[0..threads]) |t| t.join();
+ var total: usize = 0;
+ for (got[0..threads]) |g| total += g;
+ return total - cells.len;
+}
+
+fn measureThreads(io: std.Io, out: *std.Io.Writer) !void {
+ const gpa = std.heap.smp_allocator;
+ const cpus = std.Thread.getCpuCount() catch 0;
+ try out.print("логических ядер: {d}\n\n", .{cpus});
+
+ const n = 1 << 20;
+ try out.print("гонка: {d} ячеек в свободном списке, выдано лишних (сумма пяти прогонов)\n", .{n});
+ try out.writeAll("потоков без блокировки под мьютексом\n");
+ for (thread_counts[1..]) |threads| {
+ var extra: [2]usize = .{ 0, 0 };
+ for ([_]bool{ false, true }, &extra) |locked, *e| {
+ for (0..5) |_| {
+ const cells = try freeChain(gpa, n);
+ defer gpa.free(cells);
+ e.* += try raceExtra(io, cells, threads, locked);
+ }
+ }
+ try out.print("{d:>7} {d:>14} {d:>13}\n", .{ threads, extra[0], extra[1] });
+ try out.flush();
+ }
+
+ const k = 23;
+ // Большой пул не успевает кончиться: видно чистую цену выделения.
+ // Маленький кончается много раз: видно цену остановки мира.
+ for ([_]u6{ 24, 18 }) |cells_log| {
+ try out.print("\nобщая куча на 2^{d} ячеек: каждый поток считает (fib {d}) обходом дерева, лучший из пяти прогонов\n", .{ cells_log, k });
+ try out.writeAll("пачка потоков мс ускорение сборок мир стоял, мс блокировок на 1000 ячеек\n");
+ for ([_]usize{ 1, 256 }) |batch| {
+ var base_ms: f64 = 0;
+ for (thread_counts) |threads| {
+ var best_ms = std.math.inf(f64);
+ var report: zl.shared.Report = undefined;
+ for (0..5) |_| {
+ const run = try parallelFib(io, gpa, @as(usize, 1) << cells_log, k, threads, batch);
+ if (run.ms < best_ms) {
+ best_ms = run.ms;
+ report = run.report;
+ }
+ }
+ if (threads == 1) base_ms = best_ms;
+ // Работа растёт вместе с числом потоков, поэтому ускорение это
+ // работа за единицу времени против одного потока.
+ const speedup = @as(f64, @floatFromInt(threads)) * base_ms / best_ms;
+ const per_1000 = 1000.0 * @as(f64, @floatFromInt(report.refills)) / @as(f64, @floatFromInt(report.allocated));
+ const stopped_ms = @as(f64, @floatFromInt(report.stopped_ns)) / 1e6;
+ try out.print("{d:>5} {d:>7} {d:>5.1} {d:>9.2} {d:>6} {d:>13.1} {d:>8.1}\n", .{ batch, threads, best_ms, speedup, report.collections, stopped_ms, per_1000 });
+ try out.flush();
+ }
+ }
+ }
+}
+
+const ParallelRun = struct { ms: f64, report: zl.shared.Report };
+
+fn parallelFib(io: std.Io, gpa: std.mem.Allocator, cells: usize, k: u32, threads: usize, batch: usize) !ParallelRun {
+ var vm: Vm = try .initWith(gpa, .{ .cells = cells });
+ defer vm.deinit();
+ _ = try zl.eval.evalSource(&vm, workload.definitions);
+ var text_buf: [32]u8 = undefined;
+ const form = try workload.form(&vm, try std.fmt.bufPrint(&text_buf, "(fib {d})", .{k}));
+ var forms: [8]Value = @splat(form);
+ var results: [8]Value = undefined;
+
+ const report = try zl.shared.evalParallel(&vm, io, forms[0..threads], results[0..threads], .{ .batch = batch });
+ for (results[0..threads]) |r| std.debug.assert(r.asFixnum() == 28657);
+ return .{ .ms = @as(f64, @floatFromInt(report.elapsed_ns)) / 1e6, .report = report };
+}
Тесты шага
//! Шаг 70: несколько потоков на одной куче.
//!
//! Гонка на свободном списке разыгрывается шагами: чередование задано
//! тестом, а не планировщиком, поэтому тест не хрупкий. Настоящие потоки
//! гоняются только на исправленной версии, где исход обязан быть одним при
//! любом чередовании: каждая ячейка выдана ровно одному потоку, результаты
//! верны, живое переживает сборку, мусор уходит.
const std = @import("std");
const zl = @import("zl");
const eval = zl.eval;
const gc = zl.gc;
const printer = zl.printer;
const reader = zl.reader;
const shared = zl.shared;
const Cell = zl.value.Cell;
const Heap = zl.heap.Heap;
const heap_mod = zl.heap;
const Value = zl.Value;
const Vm = zl.Vm;
const testing = std.testing;
fn expectPrinted(vm: *Vm, v: Value, expected: []const u8) !void {
var out: std.Io.Writer.Allocating = .init(testing.allocator);
defer out.deinit();
try printer.write(&out.writer, vm, v);
try testing.expectEqualStrings(expected, out.written());
}
// --- 1. Гонка ---
test "два потока без блокировки снимают одну голову: ячейка выдана дважды" {
var cells: [3]Cell = undefined;
cells[0] = .{ .car = .nil, .cdr = .fromCell(&cells[1]) };
cells[1] = .{ .car = .nil, .cdr = .fromCell(&cells[2]) };
cells[2] = .{ .car = .nil, .cdr = .nil };
var list: shared.RacyList = .init(&cells[0]);
// Чередование из урока: A читает, B читает, A пишет, B пишет.
const a_seen = list.load().?;
const b_seen = list.load().?;
const a_got = list.commit(a_seen);
const b_got = list.commit(b_seen);
try testing.expectEqual(&cells[0], a_got);
try testing.expectEqual(&cells[0], b_got);
// Голова сдвинулась на одну ячейку, хотя сняли две.
try testing.expectEqual(&cells[1], list.take().?);
// Последовательно, без чередования, тот же список раздаёт ячейки честно.
var fair: shared.RacyList = .init(&cells[0]);
try testing.expectEqual(&cells[0], fair.take().?);
try testing.expectEqual(&cells[1], fair.take().?);
try testing.expectEqual(&cells[2], fair.take().?);
try testing.expectEqual(@as(?*Cell, null), fair.take());
}
// --- 2. Лечение: общий пул под мьютексом и локальные пулы ---
const per_thread = 2000;
const thread_count = 4;
/// Поток берёт `per_thread` ячеек и подписывает каждую своим номером.
fn grabCells(s: *shared.Shared, id: usize, got: []*Cell, failure: *?anyerror) void {
var m: shared.Mutator = .{};
s.attach(&m) catch |err| {
failure.* = err;
return;
};
defer s.detach(&m);
for (got) |*slot| {
const v = s.heap.cons(.fromFixnum(@intCast(id)), .nil) catch |err| {
failure.* = err;
return;
};
slot.* = v.asCell();
}
}
test "под мьютексом и с локальными пулами каждая ячейка достаётся одному потоку" {
for ([_]usize{ 1, 64 }) |batch| {
// Пул больше, чем возьмут все потоки: сборка тут не нужна, и
// ячейки, про которые знает только тест, никто не освободит.
var heap: Heap = try .init(testing.allocator, .{ .cells = 4 * thread_count * per_thread });
defer heap.deinit();
var s: shared.Shared = .init(&heap, testing.io, batch);
heap.shared = &s;
var got: [thread_count][per_thread]*Cell = undefined;
var failures: [thread_count]?anyerror = @splat(null);
var threads: [thread_count]std.Thread = undefined;
for (&threads, 0..) |*t, id| t.* = try std.Thread.spawn(.{}, grabCells, .{ &s, id, &got[id], &failures[id] });
for (threads) |t| t.join();
heap.shared = null;
// Остатки пачек возвращаются куче.
s.deinit();
for (failures) |f| try testing.expectEqual(@as(?anyerror, null), f);
var seen: std.DynamicBitSetUnmanaged = try .initEmpty(testing.allocator, heap.cells.len);
defer seen.deinit(testing.allocator);
for (got, 0..) |cells, id| {
for (cells) |cell| {
const index = heap.indexOf(cell);
try testing.expect(!seen.isSet(index));
seen.set(index);
// Выданная дважды ячейка несла бы номер того, кто писал последним.
try testing.expectEqual(@as(i64, @intCast(id)), cell.car.asFixnum());
}
}
try testing.expectEqual(thread_count * per_thread, seen.count());
try testing.expectEqual(@as(usize, 0), heap.collections);
try testing.expectEqual(s.allocated, thread_count * per_thread);
// Пачки и пулы вернулись: заняты ровно выданные ячейки.
try testing.expectEqual(thread_count * per_thread, heap.stats().used);
// Пачка в одну ячейку значит блокировку на каждое выделение.
if (batch == 1) try testing.expectEqual(s.allocated, s.refills);
if (batch == 64) try testing.expect(s.refills <= thread_count * (per_thread / 64 + 1));
}
}
// --- 3. Остановка мира ---
/// `ring` строит кольцо из n ячеек и выбрасывает его, возвращая n. Кольцо
/// это циклический мусор: счётчик ссылок его не освободил бы. `churn`
/// строит k колец подряд и складывает их длины.
const definitions =
\\(define range (lambda (n) (cond ((= n 0) nil) (t (cons n (range (- n 1)))))))
\\(define last (lambda (xs) (cond ((atom (cdr xs)) xs) (t (last (cdr xs))))))
\\(define ring (lambda (n) ((lambda (xs) (cond ((set-cdr! (last xs) xs) (car xs)))) (range n))))
\\(define churn (lambda (k acc) (cond ((= k 0) acc) (t (churn (- k 1) (+ acc (ring 8)))))))
;
/// Форма потока `i`: список `keep` строится до того, как начнётся мусор, и
/// живёт только в окружении этого потока, то есть на его стеке.
fn churnForm(vm: *Vm, i: usize, buf: []u8) !Value {
const k = 100 + 25 * i;
return reader.readOne(vm, try std.fmt.bufPrint(buf, "((lambda (keep) (cons (churn {d} 0) keep)) (range 5))", .{k}));
}
test "потоки строят циклический мусор, мир останавливается, живое цело, мусор ушёл" {
for ([_]usize{ 1, 64 }) |batch| {
var vm: Vm = try .initWith(testing.allocator, .{ .cells = 1 << 14 });
defer vm.deinit();
_ = try eval.evalSource(&vm, definitions);
// Формы переживают точную сборку между прогонами только в корне:
// стек в точном режиме не читается.
var forms: heap_mod.RootList = .empty;
defer forms.deinit(testing.allocator);
try vm.heap.pushRoots(&forms);
defer vm.heap.popRoots();
var buf: [96]u8 = undefined;
for (0..thread_count) |i| try forms.append(testing.allocator, try churnForm(&vm, i, &buf));
var results: [thread_count]Value = undefined;
_ = try gc.collect(&vm.heap, .precise);
const baseline = vm.heap.stats().used;
// Два прогона подряд: если сборка что-то теряет или держит лишнее,
// второй прогон это покажет.
for (0..2) |_| {
const report = try shared.evalParallel(&vm, testing.io, forms.items, &results, .{ .batch = batch });
// Пул в 16 тысяч ячеек, а выделено в разы больше: без сборок
// прогон бы не уложился.
try testing.expect(report.allocated > 4 * vm.heap.cells.len);
try testing.expect(report.collections > 0);
for (results, 0..) |r, i| {
var expected: [64]u8 = undefined;
const text = try std.fmt.bufPrint(&expected, "({d} 5 4 3 2 1)", .{8 * (100 + 25 * i)});
try expectPrinted(&vm, r, text);
}
// Всё выделенное потоками мусор, результаты тоже: они проверены,
// и в корнях их больше нет.
_ = try gc.collect(&vm.heap, .precise);
try testing.expectEqual(baseline, vm.heap.stats().used);
}
}
}
test "ошибка в потоке всплывает из evalParallel, и машина остаётся годной" {
var vm: Vm = try .initWith(testing.allocator, .{ .cells = 1 << 14 });
defer vm.deinit();
_ = try eval.evalSource(&vm, definitions);
var buf: [96]u8 = undefined;
var forms: [thread_count]Value = undefined;
for (&forms, 0..) |*f, i| f.* = try churnForm(&vm, i, &buf);
forms[2] = try reader.readOne(&vm, "(car 5)");
var results: [thread_count]Value = undefined;
try testing.expectError(error.NotAPair, shared.evalParallel(&vm, testing.io, &forms, &results, .{}));
// Без падающей формы тот же набор идёт чисто.
forms[2] = try churnForm(&vm, 2, &buf);
_ = try shared.evalParallel(&vm, testing.io, &forms, &results, .{});
try expectPrinted(&vm, results[2], "(1200 5 4 3 2 1)");
}
- Первый тест разобран выше: гонка на свободном списке, разыгранная шагами.
- Второй запускает четыре потока по 2000 ячеек с пачками 1 и 64 на пуле, которого хватает без сборки. Каждая ячейка обязана встретиться ровно один раз и нести номер своего потока: ячейка, выданная дважды, несла бы номер того, кто писал последним. При пачке 1 блокировок ровно столько, сколько ячеек.
- Третий самый важный. Функция
ringстроит кольцо из восьми ячеек черезset-cdr!и выбрасывает его: это циклический мусор, который счётчик ссылок не освободил бы никогда. Четыре потока строят сотни колец на пуле в 16 384 ячейки, выделяя в разы больше, чем в нём есть, поэтому мир останавливается много раз. У каждого потока есть списокkeep, построенный до мусора, который живёт только в окружении этого потока, то есть на его стеке. Он обязан пережить все сборки, в том числе запрошенные другими потоками. После прогона точная сборка возвращает кучу ровно к исходному числу занятых ячеек, и так два прогона подряд. - Четвёртый проверяет, что ошибка примитива в одном потоке (
(car 5)) всплывает изevalParallel, а машина остаётся годной для следующего прогона.
Ловушка, на которой споткнулся автор эталона и которая годится в урок: формы, прочитанные в локальный массив теста, точная сборка между прогонами освобождала. В режиме .precise стек не читается, и массив на стеке для сборщика не существует. Поэтому формы лежат в явном корне, heap.pushRoots.
$ zig build test -Dstep=70 --summary all
Build Summary: 5/5 steps succeeded; 30/30 tests passed
$ zig build test --summary all
Build Summary: 31/31 steps succeeded; 117/133 tests passed (16 skipped)
Кроме четырёх тестов шага, -Dstep=70 гонит 26 тестов, которые лежат рядом с кодом модуля. Пропущенные в полном прогоне 16 тестов требуют x86-64 (JIT). В контейнере linux/amd64 проходят все 133, в linux/arm64 те же 117. Автор эталона прогнал шаг по двадцать раз в каждом режиме сборки на macOS и по восемь в каждом контейнере: ни одного падения.
Замер
zig build bench -- threads, сборка ReleaseFast, Apple M4 Max (12 производительных и 4 энергоэффективных ядра, 16 логических), 64 ГБ, macOS 26.6.2. Машину в это время грузили соседние процессы: load average 24 за минуту на старте (11 за пять минут) и 21 в конце, а два прогона подряд на такой машине расходятся на 10 до 20 процентов. Поэтому смотри на форму, а не на третий знак:
гонка: 1048576 ячеек в свободном списке, выдано лишних (сумма пяти прогонов)
потоков без блокировки под мьютексом
2 5383172 0
4 15544146 0
8 40525750 0
общая куча на 2^24 ячеек: каждый поток считает (fib 23) обходом дерева, лучший из пяти прогонов
пачка потоков мс ускорение сборок мир стоял, мс блокировок на 1000 ячеек
1 1 11.9 1.00 0 0.0 1000.0
1 2 116.7 0.20 0 0.0 1000.0
1 4 190.5 0.25 0 0.0 1000.0
1 8 341.2 0.28 0 0.0 1000.0
256 1 8.4 1.00 0 0.0 3.9
256 2 9.9 1.69 0 0.0 3.9
256 4 9.6 3.47 0 0.0 3.9
256 8 11.7 5.73 0 0.0 3.9
общая куча на 2^18 ячеек: каждый поток считает (fib 23) обходом дерева, лучший из пяти прогонов
пачка потоков мс ускорение сборок мир стоял, мс блокировок на 1000 ячеек
1 1 14.0 1.00 2 2.1 1000.0
1 2 117.5 0.24 5 5.4 1000.0
1 4 199.5 0.28 11 13.1 1000.0
1 8 384.9 0.29 22 26.2 1000.0
256 1 11.3 1.00 2 2.3 3.9
256 2 15.2 1.48 5 5.9 3.9
256 4 22.9 1.97 11 13.0 3.9
256 8 37.5 2.40 22 26.3 3.9
Как читать. Каждый поток считает свой (fib 23), работа растёт вместе с числом потоков, поэтому «ускорение» здесь это работа за единицу времени относительно одного потока: восемь потоков, сделавшие восьмикратную работу за то же время, дали бы 8.
- Гонка на настоящих потоках раздаёт миллионы лишних ячеек, и чем потоков больше, тем больше. Под мьютексом ни одной.
- Мьютекс на каждую ячейку (пачка 1) медленнее одного потока при любом числе потоков: 0,20 до 0,29. Блокировка на каждом
consи есть вся работа, потоки стоят в очереди к одному мьютексу, а линия с ним перебегает между ядрами. - Локальные пулы (пачка 256) на большом пуле, где сборок нет, дают около 6 на восьми потоках. Блокировка берётся 3,9 раза на тысячу ячеек, то есть раз на пачку.
- Маленький пул упирается в остановку мира. Сборок становится больше вместе с работой, каждая сборка последовательна (пометка, sweep всего пула, нарезка), и на восьми потоках мир стоит больше двух третей прогона: 26,3 из 37,5 мс. Ускорение застревает около 2,4.
Последний пункт это закон Амдала в чистом виде: последовательная часть программы ограничивает ускорение, сколько ни добавляй потоков. В следующем уроке мы выведем его формулой и измерим на psum. Настоящие сборщики лечат это параллельной пометкой, ленивым sweep и поколениями, это за рамками шага.
На macOS
Весь урок работает на macOS напрямую: std.Thread сделан поверх pthreads, а Io.Mutex и Io.Condition засыпают через __ulock_wait2 вместо futex. Живые выводы badcnt, тесты обоих шагов и замер zl сняты на macOS 26.6.2 (Apple M4 Max). Тесты шага 70 в our-tiny и три прогона badcnt повторены в контейнере ghcr.io/bondiano/runner-zig:dev (linux/arm64, Debian 12, ядро OrbStack), шаг 70 в zl проверен в linux/arm64 и linux/amd64.
Ассемблер x86-64 в уроке снят кросс-компиляцией: zig build asm -Dtarget=x86_64-linux -Doptimize=ReleaseFast в examples/our-zlisp и zig build-obj -target x86_64-linux для цикла badcnt, дизассемблер objdump из Xcode. Запустить такой бинарник на Mac можно только в контейнере --platform linux/amd64, а там его исполняет Rosetta, которая переводит x86-64 в команды arm64. Гонка badcnt проявится и там.
Переменная threadlocal на macOS устроена иначе, чем на Linux. Возьмём маленький файл с одной переменной потока:
//! Переменная потока и две функции, которые её пишут и читают. Собери
//! под macOS и под Linux и сравни, как выглядит доступ к ней.
threadlocal var failure: u16 = 0;
export fn fail(e: u16) void {
failure = e;
}
export fn take() u16 {
return failure;
}
Под Linux (zig build-obj -O ReleaseFast -target aarch64-linux tls.zig) запись это указатель потока плюс смещение, а релокации TLSLE (local exec) говорят компоновщику вписать смещение прямо в команды:
$ objdump -d -r --no-show-raw-insn tls.o
0000000000000020 <tls.fail>:
20: stp x29, x30, [sp, #-0x10]!
24: mov x29, sp
28: mrs x8, TPIDR_EL0
2c: add x8, x8, #0x0, lsl #12 // =0x0
000000000000002c: R_AARCH64_TLSLE_ADD_TPREL_HI12 tls.failure
30: add x8, x8, #0x0
0000000000000030: R_AARCH64_TLSLE_ADD_TPREL_LO12_NC tls.failure
34: strh w0, [x8]
38: ldp x29, x30, [sp], #0x10
3c: ret
Под macOS (zig build-obj -O ReleaseFast tls.zig) Mach-O не кладёт переменные потока по известному смещению от регистра. Код загружает дескриптор переменной, достаёт из него адрес функции доступа, которую подставил загрузчик, и вызывает её; функция возвращает адрес экземпляра этого потока:
$ objdump -d -r --no-show-raw-insn tls.o
0000000000000024 <_tls.fail>:
24: stp x29, x30, [sp, #-0x10]!
28: mov x29, sp
2c: mov x8, x0
30: adrp x0, 0x0 <ltmp0>
0000000000000030: ARM64_RELOC_TLVP_LOAD_PAGE21 _tls.failure
34: ldr x0, [x0]
0000000000000034: ARM64_RELOC_TLVP_LOAD_PAGEOFF12 _tls.failure
38: ldr x9, [x0]
3c: blr x9
40: strh w8, [x0]
44: ldp x29, x30, [sp], #0x10
48: ret
Поэтому ассемблер zl мы снимали для Linux: там доступ к TLS виден одной-двумя командами.
Дно стека потока zl узнаёт через pthread_get_stackaddr_np на macOS и разбором /proc/self/maps на Linux; /proc на Mac нет. Детектор гонок ThreadSanitizer из Zig 0.16.0 работает только на Linux, на macOS 26 он падает при старте. Им мы займёмся в уроке 72, в контейнере.
Практика
Ограниченный буфер своими руками. Кольцо, мьютекс и два семафора уже объявлены, insert и remove написаны как обычное кольцо без синхронизации: в одном потоке оно работает, а с несколькими затирает непрочитанные слоты и отдаёт мусор из пустого буфера. Перепиши обе функции так, как в Sbuf выше. Три отличия от эталона: мьютекс здесь настоящий Io.Mutex, индексы заворачиваются сразу (rear = (rear + 1) % n), а wait и lock отменяемые, так что ошибку error.Canceled надо пробросить через try.
Тесты проверяют порядок в одном потоке, оборот кольца, то, что insert на полном буфере и remove на пустом засыпают, пока другой поток их не освободит, и прогон трёх производителей с тремя потребителями через четыре слота: каждый из 6000 элементов приходит ровно один раз и в порядке вставки своего производителя.
Упражнения
Итоги
- Потоки делят всё адресное пространство, кроме регистров. Переменная разделяемая, если на какой-то её экземпляр ссылаются несколько потоков. Глобальная и переменная контейнера в функции существуют в одном экземпляре, локальная по экземпляру на вызов в своём стеке,
threadlocalпо экземпляру на поток. cnt += 1это загрузка, изменение и сохранение. На aarch64 это три инструкции, на x86-64 однаaddqс операндом в памяти, но и она не атомарна безlock.volatileзапрещает компилятору склеивать обращения, атомарности он не даёт.- Граф выполнения показывает перемешивание инструкций как траекторию. Критические секции двух потоков образуют небезопасную зону, и траектория через неё может потерять обновление. Гонка не обязана проявиться при каждом запуске.
- Семафор это неотрицательное целое с P и V. Двоичный семафор вокруг критической секции создаёт запретную зону, которая накрывает небезопасную. Мьютекс это двоичный семафор с хозяином.
- В Zig 0.16
Mutex,Semaphore,ConditionиRwLockживут вstd.Ioи принимаютio: ожиданием управляет реализацияIo. Ожидания это точки отмены,Uncancelableждут до конца.Io.Mutexсделан на futex, аIo.Semaphoreна мьютексе и условной переменной, и V в нём не передаёт разрешение из рук в руки. sbuf:slotsсчитает свободное,itemsзанятое, мьютекс защищает кольцо. P над счётным семафором всегда до мьютекса, иначе взаимоблокировка.- Условная переменная отпускает мьютекс и засыпает неделимо, условие проверяется в
while. Семафор удобен для счётных ресурсов, условная переменная для произвольных условий. - Первая задача о читателях и писателях морит голодом писателей, вторая читателей.
std.Io.RwLockпропускает писателей вперёд. - В общей куче
zlошибка примитива переехала вthreadlocal: на x86-64 это запись по%fs, на aarch64 черезTPIDR_EL0. Ячейки выдаются пачками под мьютексом из локальных пулов, сборщик останавливает мир в безопасной точкеrefill, а условная переменная ждёт, пока встанут все. Мьютекс на каждую ячейку медленнее одного потока, пачки дают около 6 на восьми, а последовательная сборка ограничивает ускорение по закону Амдала.
Дальше
Сегодня у нас появились все инструменты главы: семафор для исключения и для порядка, условная переменная, замок читателей и писателей, и мы увидели, чего стоит каждое решение, вплоть до инструкции. Куча zl стала общей, и первый же замер показал две вещи, которые станут темой следующего урока: блокировка на каждой операции делает программу медленнее последовательной, а последовательная часть, вроде нашей сборки, ограничивает ускорение, сколько ни добавляй потоков. Там sbuf станет очередью соединений для сервера с заранее созданным пулом потоков, psum в пяти версиях покажет ускорение и эффективность на реальном железе, закон Амдала получит формулу, а ложное разделение получит цену в миллисекундах. Раннер zbox получит пул воркеров и очередь задач, а zt prof нарисует flame graph сервера под нагрузкой, где видно, сколько времени уходит на блокировки.
домашка