Skip to main content

strat9_kernel/ipc/
transport.rs

1//! IPC transport layer: traits, dispatch enum, and manager.
2//!
3//! This module defines the three-level transport hierarchy:
4//! - **N1 (TypeSafe)**: kernel-internal mailbox, same address space, ~3-10 cycles.
5//! - **N2 (LockFree)**: shared-memory SPSC ring with futex notification, ~400-4000 cycles.
6//! - **N3 (Mmu)**: thread migration with PCID-preserving CR3 switch (research track).
7//!
8//! The [`TransportManager`] selects the appropriate level per silo-pair at
9//! connection creation time using a configurable decision matrix.
10use super::{
11    lockfree_ring::{LockFreeRing, RingError},
12    mailbox::IntrusiveMailbox,
13    n3::{MigrationState, N3Transport},
14};
15use crate::{process, silo::SiloId, sync::SpinLock};
16use alloc::{collections::BTreeMap, sync::Arc, vec::Vec};
17use core::sync::atomic::{AtomicU32, AtomicU64, Ordering};
18// ---------------------------------------------------------------------------
19// Transport level enumeration
20// ---------------------------------------------------------------------------
21
22/// Isolation level of an IPC transport.
23#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
24#[repr(u8)]
25pub enum TransportLevel {
26    /// Rust type-based isolation : kernel-internal, same address space.
27    TypeSafe = 1,
28    /// Shared-memory lock-free SPSC ring with futex notification.
29    LockFree = 2,
30    /// MMU-based thread migration with PCID preservation (research).
31    Mmu = 3,
32}
33
34// ---------------------------------------------------------------------------
35// IpcTransport trait hierarchy
36// ---------------------------------------------------------------------------
37
38/// Describes the capabilities of an IPC transport.
39#[derive(Debug, Clone, Copy)]
40pub struct TransportCapabilities {
41    /// Maximum message size in bytes.
42    pub max_message_size: usize,
43    /// Whether the transport supports blocking send/recv.
44    pub blocking: bool,
45    /// Whether the transport supports zero-copy (DMA into buffers).
46    pub zero_copy: bool,
47    /// Whether the transport supports vectored I/O (scatter-gather).
48    pub vectored: bool,
49    /// Number of directions (1 = simplex, 2 = duplex).
50    pub directions: u8,
51    /// Approximate round-trip cost in CPU cycles (for scheduler hints).
52    pub estimated_cost_cycles: u32,
53}
54
55/// Core trait for all IPC transports.
56pub trait IpcTransport: Send + Sync {
57    /// The isolation level of this transport.
58    fn level(&self) -> TransportLevel;
59    /// Static capabilities of this transport.
60    fn capabilities(&self) -> TransportCapabilities;
61    /// Human-readable name for debugging / profiling.
62    fn name(&self) -> &'static str;
63}
64
65/// Producer side : sends messages into the transport.
66pub trait IpcProducer: IpcTransport {
67    /// Send a message, blocking if the buffer is full.
68    fn send(&self, msg: &[u8]) -> Result<(), IpcError>;
69    /// Send without blocking.
70    fn try_send(&self, msg: &[u8]) -> Result<(), IpcError>;
71    /// Send using vectored I/O (scatter-gather).
72    ///
73    /// Default implementation concatenates buffers and calls `send()`.
74    /// Subtypes may override for zero-copy scatter-gather.
75    fn send_vectored(&self, bufs: &[&[u8]]) -> Result<(), IpcError> {
76        let total: usize = bufs.iter().map(|b| b.len()).sum();
77        if total > self.capabilities().max_message_size {
78            return Err(IpcError::MessageTooLarge);
79        }
80        let mut buf = alloc::vec![0u8; total];
81        let mut offset = 0;
82        for b in bufs {
83            buf[offset..offset + b.len()].copy_from_slice(b);
84            offset += b.len();
85        }
86        self.send(&buf)
87    }
88}
89
90/// Consumer side : receives messages from the transport.
91pub trait IpcConsumer: IpcTransport {
92    /// Receive a message, blocking if the buffer is empty.
93    fn recv(&self, buf: &mut [u8]) -> Result<usize, IpcError>;
94    /// Receive without blocking.
95    fn try_recv(&self, buf: &mut [u8]) -> Result<Option<usize>, IpcError>;
96}
97
98/// Notification side : futex-based wakeup for blocking transports.
99pub trait IpcNotification: IpcTransport {
100    /// Notify the consumer that data is available.
101    fn notify_consumer(&self);
102    /// Notify the producer that space is available.
103    fn notify_producer(&self);
104    /// Block until a notification arrives.
105    fn wait_notification(&self) -> Result<(), IpcError>;
106}
107
108// ---------------------------------------------------------------------------
109// IpcError
110// ---------------------------------------------------------------------------
111
112/// Errors from IPC transport operations.
113#[derive(Debug, Clone, Copy, PartialEq, Eq)]
114pub enum IpcError {
115    /// Transport buffer is full.
116    WouldBlock,
117    /// Transport has been closed or disconnected.
118    Disconnected,
119    /// Message exceeds transport capacity.
120    MessageTooLarge,
121    /// Provided buffer is too small for pending message.
122    BufferTooSmall,
123    /// No transport exists for the given pair.
124    TransportNotFound,
125    /// Insufficient permissions.
126    PermissionDenied,
127    /// Generic transport failure.
128    TransportFailed,
129    /// `validate_rip()` failed : invalid or non-executable instruction pointer.
130    InvalidRip,
131    /// Watchdog timeout : migration was not completed in time.
132    TimedOut,
133}
134
135// ---------------------------------------------------------------------------
136// TransportEndpoint : sum-type dispatch (no vtable)
137// ---------------------------------------------------------------------------
138
139/// A concrete IPC transport endpoint, dispatched via enum (no `dyn` vtable).
140#[derive(Debug, Clone)]
141pub enum TransportEndpoint {
142    /// N1 kernel-internal mailbox.
143    Mailbox(Arc<IntrusiveMailbox>),
144    /// N2 lock-free SPSC ring.
145    LockFree(Arc<LockFreeRing>),
146    /// N3 MMU thread migration.
147    Mmu(Arc<super::n3::N3Transport>),
148}
149
150impl IpcTransport for TransportEndpoint {
151    fn level(&self) -> TransportLevel {
152        match self {
153            Self::Mailbox(_) => TransportLevel::TypeSafe,
154            Self::LockFree(_) => TransportLevel::LockFree,
155            Self::Mmu(_) => TransportLevel::Mmu,
156        }
157    }
158
159    fn capabilities(&self) -> TransportCapabilities {
160        match self {
161            Self::Mailbox(m) => m.capabilities(),
162            Self::LockFree(r) => r.capabilities(),
163            Self::Mmu(n) => n.capabilities(),
164        }
165    }
166
167    fn name(&self) -> &'static str {
168        match self {
169            Self::Mailbox(_) => "mailbox",
170            Self::LockFree(_) => "lockfree",
171            Self::Mmu(n) => n.name(),
172        }
173    }
174}
175
176impl IpcProducer for TransportEndpoint {
177    fn send(&self, msg: &[u8]) -> Result<(), IpcError> {
178        match self {
179            Self::Mailbox(m) => m.send(msg),
180            Self::LockFree(r) => r.send(msg),
181            Self::Mmu(n) => n.send(msg),
182        }
183    }
184
185    fn try_send(&self, msg: &[u8]) -> Result<(), IpcError> {
186        match self {
187            Self::Mailbox(m) => m.try_send(msg),
188            Self::LockFree(r) => r.try_send(msg),
189            Self::Mmu(n) => n.try_send(msg),
190        }
191    }
192}
193
194impl IpcConsumer for TransportEndpoint {
195    fn recv(&self, buf: &mut [u8]) -> Result<usize, IpcError> {
196        match self {
197            Self::Mailbox(m) => m.recv(buf),
198            Self::LockFree(r) => r.recv(buf),
199            Self::Mmu(n) => n.recv(buf),
200        }
201    }
202
203    fn try_recv(&self, buf: &mut [u8]) -> Result<Option<usize>, IpcError> {
204        match self {
205            Self::Mailbox(m) => m.try_recv(buf),
206            Self::LockFree(r) => r.try_recv(buf),
207            Self::Mmu(n) => n.try_recv(buf),
208        }
209    }
210}
211
212// LockFreeRing implements IpcTransport, IpcProducer, IpcConsumer
213impl IpcTransport for LockFreeRing {
214    fn level(&self) -> TransportLevel {
215        TransportLevel::LockFree
216    }
217
218    fn capabilities(&self) -> TransportCapabilities {
219        TransportCapabilities {
220            max_message_size: 2048,
221            blocking: true,
222            // P1 fix: messages live in ArrayQueue<Box<[u8]>> (kernel heap),
223            // not in shared-memory frames.  Advertise honestly.
224            zero_copy: false,
225            vectored: true,
226            directions: 2,
227            estimated_cost_cycles: 400,
228        }
229    }
230
231    fn name(&self) -> &'static str {
232        "lockfree"
233    }
234}
235
236impl IpcProducer for LockFreeRing {
237    fn send(&self, msg: &[u8]) -> Result<(), IpcError> {
238        self.write(msg).map_err(ring_to_ipc_error)?;
239        self.notify_consumer_raw();
240        Ok(())
241    }
242
243    fn try_send(&self, msg: &[u8]) -> Result<(), IpcError> {
244        self.write(msg).map_err(ring_to_ipc_error)?;
245        self.notify_consumer_raw();
246        Ok(())
247    }
248}
249
250impl IpcConsumer for LockFreeRing {
251    fn recv(&self, buf: &mut [u8]) -> Result<usize, IpcError> {
252        self.read(buf).map_err(ring_to_ipc_error)
253    }
254
255    fn try_recv(&self, buf: &mut [u8]) -> Result<Option<usize>, IpcError> {
256        self.try_read(buf).map_err(ring_to_ipc_error)
257    }
258}
259
260fn ring_to_ipc_error(e: RingError) -> IpcError {
261    match e {
262        RingError::Full => IpcError::WouldBlock,
263        RingError::Empty => IpcError::WouldBlock,
264        RingError::MessageTooLarge => IpcError::MessageTooLarge,
265        RingError::BufferTooSmall => IpcError::BufferTooSmall,
266        _ => IpcError::TransportFailed,
267    }
268}
269
270// ---------------------------------------------------------------------------
271// TransportId / TransportCreateResult
272// ---------------------------------------------------------------------------
273
274/// Unique identifier for a transport connection.
275#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
276pub struct TransportId(u64);
277
278impl TransportId {
279    fn new() -> Self {
280        static NEXT: AtomicU64 = AtomicU64::new(1);
281        TransportId(NEXT.fetch_add(1, Ordering::Relaxed))
282    }
283
284    /// Create a TransportId from a raw u64.
285    pub fn from_u64(raw: u64) -> Self {
286        TransportId(raw)
287    }
288
289    /// Get the raw u64 value.
290    pub fn as_u64(self) -> u64 {
291        self.0
292    }
293}
294
295/// Result of a successful transport creation.
296#[derive(Debug, Clone)]
297pub struct TransportCreateResult {
298    pub id: TransportId,
299    pub local: TransportEndpoint,
300    pub remote: TransportEndpoint,
301    pub level: TransportLevel,
302    /// (src_silo, dst_silo) pair used for cache invalidation on close.
303    pub(crate) pair: (u32, u32),
304}
305
306// ---------------------------------------------------------------------------
307// TransportConfig
308// ---------------------------------------------------------------------------
309
310/// User-supplied configuration for transport creation.
311#[derive(Debug, Clone)]
312pub struct TransportConfig {
313    /// Minimum acceptable isolation level.
314    pub min_level: TransportLevel,
315    /// Requested ring capacity in slots (N2 only).
316    pub ring_capacity: Option<u32>,
317    /// Requested slot size in bytes (N2 only).
318    pub slot_size: Option<usize>,
319}
320
321// ---------------------------------------------------------------------------
322// Decision matrix
323// ---------------------------------------------------------------------------
324
325/// Policy entry for a single cell in the decision matrix.
326#[derive(Debug, Clone)]
327struct TransportPolicyEntry {
328    level: TransportLevel,
329    ring_capacity: u32,
330}
331
332/// Default decision matrix for Tier × Tier transport selection.
333const DECISION_MATRIX: [[TransportPolicyEntry; 3]; 3] = [
334    // src = Critical
335    [
336        TransportPolicyEntry {
337            level: TransportLevel::TypeSafe,
338            ring_capacity: 0,
339        },
340        TransportPolicyEntry {
341            level: TransportLevel::TypeSafe,
342            ring_capacity: 0,
343        },
344        TransportPolicyEntry {
345            level: TransportLevel::LockFree,
346            ring_capacity: 256,
347        },
348    ],
349    // src = System
350    [
351        TransportPolicyEntry {
352            level: TransportLevel::TypeSafe,
353            ring_capacity: 0,
354        },
355        TransportPolicyEntry {
356            level: TransportLevel::LockFree,
357            ring_capacity: 256,
358        },
359        TransportPolicyEntry {
360            level: TransportLevel::LockFree,
361            ring_capacity: 256,
362        },
363    ],
364    // src = User
365    [
366        TransportPolicyEntry {
367            level: TransportLevel::LockFree,
368            ring_capacity: 256,
369        },
370        TransportPolicyEntry {
371            level: TransportLevel::LockFree,
372            ring_capacity: 256,
373        },
374        TransportPolicyEntry {
375            // User=>User: maximum isolation via MMU thread migration.
376            // N3 provides hermetic Ring 3 isolation without shared memory.
377            level: TransportLevel::Mmu,
378            ring_capacity: 0,
379        },
380    ],
381];
382
383/// Per-transport performance counters.
384#[derive(Debug, Clone)]
385pub struct TransportStats {
386    /// Total transports created.
387    pub created: u64,
388    /// Total messages sent (incremented by syscall handler).
389    pub sent: u64,
390    /// Total messages received.
391    pub received: u64,
392    /// Total errors.
393    pub errors: u64,
394    /// Current transport level.
395    pub level: TransportLevel,
396}
397
398impl TransportStats {
399    const fn new(level: TransportLevel) -> Self {
400        TransportStats {
401            created: 0,
402            sent: 0,
403            received: 0,
404            errors: 0,
405            level,
406        }
407    }
408}
409
410// ---------------------------------------------------------------------------
411// TransportEndpoint helpers for polling
412// ---------------------------------------------------------------------------
413
414impl TransportEndpoint {
415    /// Whether the transport has data to read.
416    pub fn has_data(&self) -> bool {
417        match self {
418            Self::Mailbox(m) => !m.is_empty(),
419            Self::LockFree(r) => r.has_data(),
420            Self::Mmu(n) => {
421                let frame = n.frame();
422                frame.msg_len > 0 && frame.generation.load(Ordering::Acquire) > 0
423            }
424        }
425    }
426
427    /// Whether the transport has space to write.
428    pub fn has_space(&self) -> bool {
429        match self {
430            Self::Mailbox(_) => true,
431            Self::LockFree(r) => r.has_space(),
432            Self::Mmu(n) => {
433                let frame = n.frame();
434                frame.state.load(Ordering::Acquire) == MigrationState::Ready as u8
435            }
436        }
437    }
438}
439
440// ---------------------------------------------------------------------------
441// TransportManager
442// ---------------------------------------------------------------------------
443
444/// Central transport manager : selects and creates IPC transports per silo pair.
445///
446/// Uses a static decision matrix (tier × tier) with optional dynamic
447/// overrides.  A small FIFO cache avoids redundant creation for frequently
448/// used pairs.
449pub struct TransportManager {
450    /// Static decision matrix: [src_tier][dst_tier] -> policy.
451    decision_matrix: [[TransportPolicyEntry; 3]; 3],
452    /// Dynamic overrides (set by silo admin).
453    policy_overrides: SpinLock<BTreeMap<(u32, u32), TransportPolicyEntry>>,
454    /// Active transport registry.
455    active: SpinLock<BTreeMap<TransportId, TransportCreateResult>>,
456    /// Simple FIFO transport cache (no_std, no allocation).
457    cache: SpinLock<TransportCache>,
458    /// Per-transport performance statistics.
459    pub stats: SpinLock<TransportStats>,
460}
461
462impl TransportManager {
463    /// Create a new manager with the default decision matrix.
464    pub const fn new() -> Self {
465        TransportManager {
466            decision_matrix: DECISION_MATRIX,
467            policy_overrides: SpinLock::new(BTreeMap::new()),
468            active: SpinLock::new(BTreeMap::new()),
469            cache: SpinLock::new(TransportCache::new()),
470            stats: SpinLock::new(TransportStats::new(TransportLevel::LockFree)),
471        }
472    }
473
474    /// Establish a transport between `src` and `dst`.
475    ///
476    /// Selection order:
477    /// 1. Cache hit (fast path for repeated pairs).
478    /// 2. Dynamic override (admin-configured policy).
479    /// 3. Static decision matrix (safe default).
480    pub fn establish(
481        &self,
482        src: SiloId,
483        dst: SiloId,
484        config: TransportConfig,
485    ) -> Result<TransportCreateResult, IpcError> {
486        let pair = (src.sid, dst.sid);
487
488        // 1. Cache lookup
489        {
490            let mut cache = self.cache.lock();
491            if let Some(cached) = cache.get(pair) {
492                if cached.level as u8 >= config.min_level as u8 {
493                    return Ok(cached.clone());
494                }
495            }
496        }
497
498        // 2. Dynamic overrides
499        {
500            let overrides = self.policy_overrides.lock();
501            if let Some(entry) = overrides.get(&pair) {
502                return self.create(pair, entry.level, entry.ring_capacity);
503            }
504        }
505
506        // 3. Static matrix
507        let entry = &self.decision_matrix[src.tier as usize][dst.tier as usize];
508        let level = if entry.level < config.min_level {
509            config.min_level
510        } else {
511            entry.level
512        };
513        self.create(
514            pair,
515            level,
516            config.ring_capacity.unwrap_or(entry.ring_capacity),
517        )
518    }
519
520    /// Create a transport at the given level.
521    fn create(
522        &self,
523        _pair: (u32, u32),
524        level: TransportLevel,
525        capacity: u32,
526    ) -> Result<TransportCreateResult, IpcError> {
527        let id = TransportId::new();
528
529        let (local, remote) = match level {
530            TransportLevel::TypeSafe => {
531                let mb = IntrusiveMailbox::new();
532                let arc = Arc::new(mb);
533                (
534                    TransportEndpoint::Mailbox(arc.clone()),
535                    TransportEndpoint::Mailbox(arc),
536                )
537            }
538            TransportLevel::LockFree => {
539                let ring = LockFreeRing::new(capacity.max(4), 2048)
540                    .map_err(|_| IpcError::TransportFailed)?;
541                (
542                    TransportEndpoint::LockFree(ring.clone()),
543                    TransportEndpoint::LockFree(ring),
544                )
545            }
546            TransportLevel::Mmu => {
547                // N3 MMU thread migration.
548                //
549                // Creates a unidirectional N3Transport: sender = current task,
550                // receiver = first task found in the destination silo.
551                // For full-duplex N3, two transports must be created
552                // (one in each direction).
553                let sender = process::current_task_clone().ok_or(IpcError::TransportFailed)?;
554                let sender_id = sender.id;
555
556                // Find a task in the destination silo.
557                let all_tasks = process::get_all_tasks().ok_or(IpcError::TransportFailed)?;
558                let receiver_task = all_tasks
559                    .iter()
560                    .find(|t| {
561                        crate::silo::try_silo_id_for_task(t.id).map_or(false, |sid| sid == _pair.1)
562                    })
563                    .cloned()
564                    .ok_or(IpcError::Disconnected)?;
565
566                let transport = N3Transport::new(sender_id, receiver_task.id)?;
567                let arc = Arc::new(transport);
568                (
569                    TransportEndpoint::Mmu(arc.clone()),
570                    TransportEndpoint::Mmu(arc),
571                )
572            }
573        };
574
575        let result = TransportCreateResult {
576            id,
577            local,
578            remote: remote.clone(),
579            level,
580            pair: _pair,
581        };
582
583        // Increment global transport creation counter
584        self.stats.lock().created += 1;
585
586        {
587            let mut cache = self.cache.lock();
588            cache.put(_pair, result.clone());
589        }
590        {
591            let mut active = self.active.lock();
592            active.insert(id, result.clone());
593        }
594
595        Ok(result)
596    }
597
598    /// Look up a transport endpoint by ID.
599    pub fn get_endpoint(&self, id: TransportId) -> Option<TransportEndpoint> {
600        let active = self.active.lock();
601        active.get(&id).map(|r| r.local.clone())
602    }
603
604    /// Remove a transport from the registry (called when all handles close).
605    /// Also invalidates the cache entry for the pair to prevent stale reuse.
606    pub fn close(&self, id: TransportId) -> Result<(), IpcError> {
607        let mut active = self.active.lock();
608        let removed = active.remove(&id).ok_or(IpcError::TransportNotFound)?;
609        // Invalidate cache entry to prevent stale reuse.
610        let mut cache = self.cache.lock();
611        cache.invalidate(removed.pair);
612        Ok(())
613    }
614
615    /// Override the transport policy for a specific silo pair.
616    pub fn set_policy(&self, src: u32, dst: u32, level: TransportLevel, capacity: u32) {
617        let mut overrides = self.policy_overrides.lock();
618        overrides.insert(
619            (src, dst),
620            TransportPolicyEntry {
621                level,
622                ring_capacity: capacity,
623            },
624        );
625    }
626}
627
628// ---------------------------------------------------------------------------
629// TransportCache : no_std FIFO cache
630// ---------------------------------------------------------------------------
631
632const CACHE_SIZE: usize = 128;
633
634struct TransportCache {
635    entries: [((u32, u32), Option<TransportCreateResult>); CACHE_SIZE],
636    next: usize,
637}
638
639impl TransportCache {
640    const fn new() -> Self {
641        const NONE: ((u32, u32), Option<TransportCreateResult>) = ((0, 0), None);
642        TransportCache {
643            entries: [NONE; CACHE_SIZE],
644            next: 0,
645        }
646    }
647
648    fn get(&mut self, key: (u32, u32)) -> Option<&TransportCreateResult> {
649        self.entries
650            .iter()
651            .find_map(|(k, v)| if *k == key { v.as_ref() } else { None })
652    }
653
654    fn put(&mut self, key: (u32, u32), value: TransportCreateResult) {
655        let idx = self.next;
656        self.entries[idx] = (key, Some(value));
657        self.next = (self.next + 1) % CACHE_SIZE;
658    }
659
660    /// Invalidate any cached entry for the given key.
661    fn invalidate(&mut self, key: (u32, u32)) {
662        if let Some(entry) = self.entries.iter_mut().find(|(k, _)| *k == key) {
663            entry.1 = None;
664        }
665    }
666}
667
668// ---------------------------------------------------------------------------
669// TypedLockFreeRing : typestate pattern for SPSC role safety
670// ---------------------------------------------------------------------------
671
672use core::marker::PhantomData;
673
674/// Producer role marker: only this side can call `write()`.
675pub struct Producer;
676
677/// Consumer role marker: can only call `read()`.
678pub struct Consumer;
679
680/// SPSC ring with compile-time role enforcement.
681///
682/// `TypedLockFreeRing<Producer>` has both `write()` and `read()`.
683/// `TypedLockFreeRing<Consumer>` has only `read()`.
684/// The compiler rejects any misuse at compile time.
685pub struct TypedLockFreeRing<Role> {
686    inner: Arc<LockFreeRing>,
687    _role: PhantomData<Role>,
688}
689
690impl TypedLockFreeRing<Producer> {
691    /// Write into the ring (producer only).
692    pub fn write(&self, data: &[u8]) -> Result<(), RingError> {
693        self.inner.write(data)
694    }
695}
696
697impl<T> TypedLockFreeRing<T> {
698    /// Read from the ring (both roles).
699    pub fn read(&self, buf: &mut [u8]) -> Result<usize, RingError> {
700        self.inner.read(buf)
701    }
702}
703
704/// Create a typed SPSC pair returning separate producer and consumer ends.
705pub fn create_spsc_pair(
706    cap: u32,
707    slot_size: usize,
708) -> Result<(TypedLockFreeRing<Producer>, TypedLockFreeRing<Consumer>), RingError> {
709    let ring = LockFreeRing::new(cap, slot_size)?;
710    Ok((
711        TypedLockFreeRing {
712            inner: ring.clone(),
713            _role: PhantomData,
714        },
715        TypedLockFreeRing {
716            inner: ring,
717            _role: PhantomData,
718        },
719    ))
720}