190 lines
6.7 KiB
Markdown
190 lines
6.7 KiB
Markdown
# 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 (неявно) │
|
||
```
|
||
|
||
## Структура канала
|
||
|
||
```rust
|
||
#[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
|
||
|
||
```rust
|
||
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() — выделить канал
|
||
|
||
```rust
|
||
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) — запись результатов
|
||
|
||
```rust
|
||
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) — ожидание результата
|
||
|
||
```rust
|
||
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 — синхронная обёртка
|
||
|
||
```rust
|
||
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)
|
||
}
|
||
```
|
||
|
||
## Глобальный экземпляр
|
||
|
||
```rust
|
||
static ROUTER: GlobalRouter = GlobalRouter {
|
||
is_ready: AtomicBool::new(false),
|
||
inner: UnsafeCell::new(None),
|
||
};
|
||
```
|
||
|
||
Double-checked initialization:
|
||
1. `init()` — сетап всех каналов, `is_ready = true`.
|
||
2. `get_router()` — проверяет `is_ready`, затем `unwrap_unchecked()`.
|
||
|
||
**Безопасность**: `is_ready` устанавливается один раз и никогда не
|
||
очищается, поэтому TOCTOU между проверкой и unwrap — безопасен.
|
||
|
||
## dispatch() — утилита
|
||
|
||
```rust
|
||
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.
|