6.7 KiB
6.7 KiB
PM Router: Lock-Free маршрутизация ответов
Концептуальная модель
PM Router — это центральный коммутатор, соединяющий исполнителей запросов (PMActors) с потребителями результатов (процессы/потоки).
- Он предоставляет 65536 каналов для асинхронной коммуникации.
- Каналы lock-free: alloc/release канала — CAS на атомарных счётчиках.
- Исключает блокировки между продюсером и консьюмером.
Процесс A PMActor #1
│ │
│ alloc_channel() │
│ ─────────► router ──────► │
│ │
│ submit_request(Alloc) │
│ ────────────────────────► │
│ │ process_messages()
│ │ ──────► buddy.alloc()
│ │
│ ◄─────────────── │
│ route_responses([Resp]) │
│ router │
│ │
│ wait_for_response(ch) │
│ ◄═══ RESULT ════ │
│ │
│ free_channel (неявно) │
Структура канала
#[repr(align(64))]
pub struct Channel {
state: AtomicU8, // FREE=0, PENDING=1, READY=2
next_free: AtomicU16, // указатель в стеке свободных каналов
result: UnsafeCell<Option<PMResult>>, // ячейка результата
}
align(64)— каждая кэш-линия содержит ровно один канал.state— конечный автомат: FREE → PENDING → READY → FREE.UnsafeCell— потому что write происходит из route_responses, read — из wait_for_response. Синхронизация через state.
Структура PMRouter
pub struct PMRouter {
channels: Box<[Channel]>, // 65536 каналов, Box<[T]> в куче
free_head: AtomicU16, // стек свободных каналов
}
Стек свободных каналов
Изначально все каналы свободны, free_head = 1 (канал 0 зарезервирован
как «discard» — ответы на канал 0 игнорируются).
free_head ──► channel[1].next_free = 2
channel[2].next_free = 3
channel[3].next_free = 4
...
channel[65535].next_free = 0 (NULL)
Операции
alloc_channel() — выделить канал
pub fn alloc_channel(&self) -> Option<u16> {
let mut head = self.free_head.load(Acquire);
loop {
if head == 0 { return None; } // нет свободных
let next = self.channels[head].next_free.load(Relaxed);
// CAS: free_head = head → next
match self.free_head.compare_exchange_weak(head, next, AcqRel, Acquire) {
Ok(_) => {
self.channels[head].state.store(STATE_PENDING, Release);
return Some(head);
}
Err(new) => head = new,
}
}
}
Lock-free: CAS на free_head позволяет нескольким продюсерам
конкурировать без блокировок.
route_responses(responses) — запись результатов
pub fn route_responses(&self, responses: Vec<PMResponse>) {
for resp in responses {
if resp.channel_id == 0 { continue; }
let channel = &self.channels[resp.channel_id as usize];
unsafe { *channel.result.get() = Some(resp.result); }
channel.state.store(STATE_READY, Release);
// TODO: Focus Mode — пробуждение ожидающего потока
}
}
wait_for_response(id) — ожидание результата
pub fn wait_for_response(&self, id: u16) -> PMResult {
let channel = &self.channels[id as usize];
// busy-wait с HLT
while channel.state.load(Acquire) != STATE_READY {
unsafe { asm!("hlt") }; // CPU остановка до прерывания
}
let result = unsafe { (*channel.result.get()).take().unwrap() };
channel.state.store(STATE_FREE, Release);
// Возвращаем канал в стек свободных
let mut head = self.free_head.load(Relaxed);
loop {
channel.next_free.store(head, Relaxed);
match self.free_head.compare_exchange_weak(head, id, Release, Relaxed) {
Ok(_) => break,
Err(new) => head = new,
}
}
result
}
request_and_wait — синхронная обёртка
pub fn request_and_wait<F>(actor: &PMActor, req_builder: F) -> PMResult
where
F: FnOnce(u16) -> PMRequest
{
let router = get_router();
let channel_id = router.alloc_channel().expect("OOM");
let req = req_builder(channel_id);
actor.submit_request(req).expect("Inbox full");
router.wait_for_response(channel_id)
}
Глобальный экземпляр
static ROUTER: GlobalRouter = GlobalRouter {
is_ready: AtomicBool::new(false),
inner: UnsafeCell::new(None),
};
Double-checked initialization:
init()— сетап всех каналов,is_ready = true.get_router()— проверяетis_ready, затемunwrap_unchecked().
Безопасность: is_ready устанавливается один раз и никогда не
очищается, поэтому TOCTOU между проверкой и unwrap — безопасен.
dispatch() — утилита
pub fn dispatch(responses: Vec<PMResponse>) {
get_router().route_responses(responses);
}
Ключевые свойства
| Свойство | Значение |
|---|---|
| Количество каналов | 65536 |
| Размер канала | 64 байта (1 cache line) |
| alloc_channel | Lock-free, O(1) |
| route_responses | O(N), N = количество ответов |
| wait_for_response | Lock-free + HLT |
| Состояния | FREE → PENDING → READY → FREE |
Где используется
- В
kmain(): тестирование PMActor через Router. - Будущие IPC и процесс-менеджер будут использовать Router для асинхронного взаимодействия с PMActors.