feat(ipc): implement synchronous Rendezvous IPC and multi-thread cooperative scheduler

This commit is contained in:
RarDog
2026-09-02 16:52:03 +03:00
parent 6cbd13b4aa
commit 448b0e710e
11 changed files with 712 additions and 97 deletions
+58 -16
View File
@@ -1,5 +1,8 @@
//! x86_64 Fast Syscall / Sysret Setup and Fast-Path Dispatch
//! x86_64 Fast Syscall / Sysret Setup and Fast-Path & Rendezvous Dispatch
use crate::ipc::endpoint::EP_MANAGER;
use crate::ipc::fastpath::{CAP_KERNEL_CONTROL, handle_fastpath_call};
use crate::sched::SCHEDULER;
use crate::kprintln;
const IA32_EFER: u32 = 0xC0000080;
@@ -39,20 +42,13 @@ extern "C" {
}
pub unsafe fn init() {
// 1. Enable System Call Extension (SCE) in IA32_EFER
let efer = rdmsr(IA32_EFER);
wrmsr(IA32_EFER, efer | 1);
// 2. Set STAR MSR:
// Bits [47:32] = Kernel CS/SS selectors (0x0008)
// Bits [63:48] = User CS/SS selectors base (0x0010 -> SS=0x18, CS=0x20)
let star = (0x0010u64 << 48) | (0x0008u64 << 32);
wrmsr(IA32_STAR, star);
// 3. Set LSTAR to our assembly entry trampoline
wrmsr(IA32_LSTAR, syscall_entry as *const () as usize as u64);
// 4. Set FMASK to clear IF (bit 9, 0x200), DF (bit 10), TF (bit 8)
wrmsr(IA32_FMASK, 0x00000200);
}
@@ -70,8 +66,13 @@ pub struct SyscallRegisters {
/// Dispatcher called directly from syscall_entry assembly trampoline
#[no_mangle]
pub extern "C" fn kernel_syscall_dispatcher(regs: &mut SyscallRegisters) {
let current_tid = unsafe {
let sched = &mut *core::ptr::addr_of_mut!(SCHEDULER);
sched.get_current_mut().id
};
match regs.rax {
// SYS_IPC_CALL = 1
// SYS_IPC_CALL = 1 (Client: Send & Wait for Reply)
1 => {
let dest_cap = regs.rdi;
let opcode = regs.rsi;
@@ -80,16 +81,57 @@ pub extern "C" fn kernel_syscall_dispatcher(regs: &mut SyscallRegisters) {
let arg2 = regs.r9;
let arg3 = regs.r10;
let (ret, out0, out1) = crate::ipc::fastpath::handle_fastpath_call(
dest_cap, opcode, arg0, arg1, arg2, arg3,
);
if dest_cap == CAP_KERNEL_CONTROL {
let (ret, out0, out1) = handle_fastpath_call(dest_cap, opcode, arg0, arg1, arg2, arg3);
regs.rax = ret as u64;
regs.rdx = out0;
regs.r8 = out1;
} else {
// Inter-Process Rendezvous IPC Call
unsafe {
let ep_mgr = &mut *core::ptr::addr_of_mut!(EP_MANAGER);
match ep_mgr.ipc_call(dest_cap, opcode, [arg0, arg1, arg2, arg3], current_tid) {
Ok((out0, out1)) => {
regs.rax = 0;
regs.rdx = out0;
regs.r8 = out1;
}
Err(code) => {
regs.rax = code as u64;
}
}
}
}
}
// SYS_IPC_REPLY_RECV = 2 (Server: Reply to Caller & Wait for Request)
2 => {
let listen_ep = regs.rdi;
let reply_to_tid = regs.rsi;
let reply_val0 = regs.rdx;
let reply_val1 = regs.r8;
regs.rax = ret as u64;
regs.rdx = out0;
regs.r8 = out1;
unsafe {
let ep_mgr = &mut *core::ptr::addr_of_mut!(EP_MANAGER);
match ep_mgr.ipc_reply_recv(listen_ep, reply_to_tid, reply_val0, reply_val1, current_tid) {
Ok((sender_tid, opcode, a0, a1)) => {
regs.rax = 0;
regs.rdi = sender_tid;
regs.rsi = opcode;
regs.rdx = a0;
regs.r8 = a1;
}
Err(code) => {
regs.rax = code as u64;
}
}
}
}
// SYS_THREAD_YIELD = 7
7 => {
unsafe {
let sched = &mut *core::ptr::addr_of_mut!(SCHEDULER);
sched.yield_current();
}
regs.rax = 0;
}
// SYS_LOG_DEBUG = 9
@@ -99,7 +141,7 @@ pub extern "C" fn kernel_syscall_dispatcher(regs: &mut SyscallRegisters) {
if !msg_ptr.is_null() && len > 0 && len < 4096 {
let slice = unsafe { core::slice::from_raw_parts(msg_ptr, len) };
if let Ok(s) = core::str::from_utf8(slice) {
kprintln!("[USER LOG] {}", s);
kprintln!("[USER TID {}] {}", current_tid, s);
regs.rax = 0;
return;
}
+195
View File
@@ -0,0 +1,195 @@
//! IPC Endpoints and Synchronous Rendezvous Engine
use crate::ipc::fastpath::FastPathMessage;
use crate::sched::SCHEDULER;
use crate::sched::thread::ThreadState;
use crate::kprintln;
pub const MAX_ENDPOINTS: usize = 32;
pub const MAX_QUEUE_PER_EP: usize = 8;
#[derive(Clone, Copy)]
pub struct Endpoint {
pub id: u64,
pub receiver_tid: u64,
pub send_queue: [u64; MAX_QUEUE_PER_EP],
pub send_count: usize,
}
impl Endpoint {
pub const fn empty() -> Self {
Self {
id: 0,
receiver_tid: 0,
send_queue: [0; MAX_QUEUE_PER_EP],
send_count: 0,
}
}
}
pub struct EndpointManager {
pub endpoints: [Endpoint; MAX_ENDPOINTS],
}
pub static mut EP_MANAGER: EndpointManager = EndpointManager {
endpoints: [const { Endpoint::empty() }; MAX_ENDPOINTS],
};
impl EndpointManager {
pub fn init(&mut self) {
for i in 0..MAX_ENDPOINTS {
self.endpoints[i].id = i as u64;
}
}
pub fn get_endpoint_mut(&mut self, ep_id: u64) -> Option<&mut Endpoint> {
let idx = ep_id as usize;
if idx < MAX_ENDPOINTS {
Some(&mut self.endpoints[idx])
} else {
None
}
}
/// Client performs synchronous Call (Send + Wait for Reply)
pub unsafe fn ipc_call(
&mut self,
dest_ep: u64,
opcode: u64,
args: [u64; 4],
caller_tid: u64,
) -> Result<(u64, u64), i64> {
let ep = self.get_endpoint_mut(dest_ep).ok_or(-2i64)?; // Invalid Cap/Endpoint
let sched = &mut *core::ptr::addr_of_mut!(SCHEDULER);
if ep.receiver_tid != 0 {
// Rendezvous: Server is already blocked waiting for a request!
let server_tid = ep.receiver_tid;
ep.receiver_tid = 0; // Clear receiver wait state
// Transfer message directly into server's context
if let Some(server) = sched.get_thread_mut(server_tid) {
server.caller_tid = caller_tid;
server.ipc_buffer = FastPathMessage {
dest_cap: dest_ep,
opcode,
args,
};
server.state = ThreadState::Ready;
}
// Block caller on reply from server
if let Some(caller) = sched.get_thread_mut(caller_tid) {
caller.state = ThreadState::BlockedOnReply(server_tid);
}
// Direct context handoff to server
sched.schedule();
// When unblocked after reply, return results stored in caller's buffer
if let Some(caller) = sched.get_thread_mut(caller_tid) {
return Ok((caller.ipc_buffer.args[0], caller.ipc_buffer.args[1]));
}
Ok((0, 0))
} else {
// Receiver not waiting: Enqueue sender and block
if ep.send_count >= MAX_QUEUE_PER_EP {
return Err(-9); // SYS_ERR_BUSY
}
ep.send_queue[ep.send_count] = caller_tid;
ep.send_count += 1;
if let Some(caller) = sched.get_thread_mut(caller_tid) {
caller.ipc_buffer = FastPathMessage {
dest_cap: dest_ep,
opcode,
args,
};
caller.state = ThreadState::BlockedOnSend(dest_ep);
}
sched.schedule();
if let Some(caller) = sched.get_thread_mut(caller_tid) {
return Ok((caller.ipc_buffer.args[0], caller.ipc_buffer.args[1]));
}
Ok((0, 0))
}
}
/// Server handles Reply to current caller and waits for Next Request
pub unsafe fn ipc_reply_recv(
&mut self,
listen_ep: u64,
reply_to_tid: u64,
reply_val0: u64,
reply_val1: u64,
server_tid: u64,
) -> Result<(u64, u64, u64, u64), i64> {
let ep = self.get_endpoint_mut(listen_ep).ok_or(-2i64)?;
let sched = &mut *core::ptr::addr_of_mut!(SCHEDULER);
// 1. Deliver reply to previous caller if requested
if reply_to_tid != 0 {
if let Some(caller) = sched.get_thread_mut(reply_to_tid) {
caller.ipc_buffer.args[0] = reply_val0;
caller.ipc_buffer.args[1] = reply_val1;
caller.state = ThreadState::Ready;
}
}
// 2. Check if a sender is already queued
if ep.send_count > 0 {
let next_sender_tid = ep.send_queue[0];
// Shift queue
for i in 0..(ep.send_count - 1) {
ep.send_queue[i] = ep.send_queue[i + 1];
}
ep.send_count -= 1;
let mut ret_opcode = 0;
let mut ret_arg0 = 0;
let mut ret_arg1 = 0;
if let Some(sender) = sched.get_thread_mut(next_sender_tid) {
ret_opcode = sender.ipc_buffer.opcode;
ret_arg0 = sender.ipc_buffer.args[0];
ret_arg1 = sender.ipc_buffer.args[1];
sender.state = ThreadState::BlockedOnReply(server_tid);
}
if let Some(srv) = sched.get_thread_mut(server_tid) {
srv.caller_tid = next_sender_tid;
}
return Ok((next_sender_tid, ret_opcode, ret_arg0, ret_arg1));
}
// 3. No sender queued: block server on endpoint
ep.receiver_tid = server_tid;
if let Some(srv) = sched.get_thread_mut(server_tid) {
srv.state = ThreadState::BlockedOnRecv(listen_ep);
}
sched.schedule();
// Server woke up with incoming message delivered
if let Some(srv) = sched.get_thread_mut(server_tid) {
let sender_tid = srv.caller_tid;
let opcode = srv.ipc_buffer.opcode;
let a0 = srv.ipc_buffer.args[0];
let a1 = srv.ipc_buffer.args[1];
return Ok((sender_tid, opcode, a0, a1));
}
Ok((0, 0, 0, 0))
}
}
pub fn init() {
unsafe {
(*core::ptr::addr_of_mut!(EP_MANAGER)).init();
kprintln!("[IPC] Synchronous Rendezvous Engine initialized ({} Endpoints).", MAX_ENDPOINTS);
}
}
+3 -1
View File
@@ -1,5 +1,7 @@
pub mod fastpath;
pub mod endpoint;
pub fn init() {
crate::kprintln!("[IPC] Capability-based Fast-Path IPC Ready.");
endpoint::init();
crate::kprintln!("[IPC] Capability-based Fast-Path & Rendezvous IPC Ready.");
}
+33 -27
View File
@@ -11,26 +11,17 @@ pub mod loader;
use core::panic::PanicInfo;
use limine_requests::*;
use sched::SCHEDULER;
// Include assembly trampolines directly into the binary
core::arch::global_asm!(include_str!("../asm/context.S"));
core::arch::global_asm!(include_str!("../asm/syscall_entry.S"));
core::arch::global_asm!(include_str!("../asm/ring3_enter.S"));
extern "C" {
fn enter_user_mode(
entry_rip: u64,
user_rsp: u64,
pml4_paddr: u64,
kernel_stack_top: u64,
) -> !;
}
#[no_mangle]
pub extern "C" fn _start() -> ! {
// 1. Validate Limine Base Revision
if BASE_REVISION.revision != 3 {
// Revision mismatch
loop { core::hint::spin_loop(); }
}
@@ -84,29 +75,44 @@ pub extern "C" fn _start() -> ! {
);
kprintln!("[TEST] IPC Fast-Path Ping Result: status={}, opcode={:#x}, val={}", status, resp_op, val);
// 9. Inspect and Load Initial Userspace Process (Ring 3 Transition)
// 9. Inspect and Load Initial Userspace Processes
unsafe {
let mod_resp = MODULE_REQUEST.response;
if !mod_resp.is_null() && (*mod_resp).module_count > 0 {
kprintln!("[BOOT] Loading root userspace module [0] into Ring 3...");
let mod_file = *(*mod_resp).modules;
let elf_slice = core::slice::from_raw_parts((*mod_file).address, (*mod_file).size as usize);
let count = (*mod_resp).module_count;
kprintln!("[BOOT] Detected {} userspace boot modules. Creating tasks...", count);
match loader::elf::load_elf(elf_slice, hhdm_offset) {
Ok(proc) => {
kprintln!("[BOOT] Transitioning CPU to Ring 3 (User Mode) via IRETQ...");
kprintln!("-------------------------------------------------------");
enter_user_mode(
proc.entry_rip,
proc.user_rsp,
proc.pml4_paddr,
proc.kernel_stack_top,
);
}
Err(err) => {
kprintln!("[ERROR] Failed to load ELF process: {}", err);
let sched_ptr = core::ptr::addr_of_mut!(SCHEDULER);
for i in 0..count {
let mod_file = *(*mod_resp).modules.add(i as usize);
let elf_slice = core::slice::from_raw_parts((*mod_file).address, (*mod_file).size as usize);
let name = match i {
0 => "init_server",
1 => "uart_driver",
_ => "user_server",
};
match loader::elf::load_elf(elf_slice, hhdm_offset) {
Ok(proc) => {
let tid = (*sched_ptr).create_user_thread(
name,
proc.entry_rip,
proc.user_rsp,
proc.pml4_paddr,
proc.kernel_stack_top,
);
kprintln!("[BOOT] Registered Thread '{}' with TID {:?}", name, tid);
}
Err(err) => {
kprintln!("[ERROR] Failed to load module [{}]: {}", i, err);
}
}
}
kprintln!("[BOOT] Starting Multi-Tasking & IPC Rendezvous...");
kprintln!("-------------------------------------------------------");
(*sched_ptr).schedule();
} else {
kprintln!("[BOOT] No userspace modules detected. Staying in Ring 0 idle loop.");
}
+203 -26
View File
@@ -1,35 +1,212 @@
//! Minimal Thread & Task Scheduler Subsystem
//! Preemptive & Cooperative Microkernel Scheduler
#[repr(C)]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ThreadState {
Ready,
Running,
BlockedOnIpc,
Dead,
pub mod thread;
use thread::{Thread, ThreadState, MAX_THREADS};
use crate::arch::gdt::set_kernel_stack;
use crate::mm::vmm::VMM;
use crate::kprintln;
extern "C" {
fn switch_to(prev_rsp: *mut u64, next_rsp: u64);
}
#[repr(C)]
pub struct ContextFrame {
pub r15: u64,
pub r14: u64,
pub r13: u64,
pub r12: u64,
pub rbx: u64,
pub rbp: u64,
pub rsp: u64,
pub rip: u64,
pub rflags: u64,
pub struct Scheduler {
pub threads: [Thread; MAX_THREADS],
pub current_tid: usize,
pub next_tid: u64,
}
pub struct Thread {
pub id: u64,
pub state: ThreadState,
pub rsp: u64,
pub cr3: u64,
pub kernel_stack_top: u64,
pub static mut SCHEDULER: Scheduler = Scheduler {
threads: [const { Thread::empty() }; MAX_THREADS],
current_tid: 0,
next_tid: 1,
};
impl Scheduler {
pub fn init(&mut self) {
// Slot 0 is reserved for Kernel Idle/Bootstrap thread
self.threads[0].id = 0;
self.threads[0].name = "kernel_idle";
self.threads[0].state = ThreadState::Running;
self.current_tid = 0;
self.next_tid = 1;
kprintln!("[SCHED] Scheduler initialized (Multi-Thread TCB table ready).");
}
pub fn create_user_thread(
&mut self,
name: &'static str,
entry_rip: u64,
user_rsp: u64,
pml4_paddr: u64,
kernel_stack_top: u64,
) -> Option<u64> {
for i in 1..MAX_THREADS {
if self.threads[i].state == ThreadState::Unused {
let tid = self.next_tid;
self.next_tid += 1;
self.threads[i].id = tid;
self.threads[i].name = name;
self.threads[i].state = ThreadState::Ready;
self.threads[i].pml4_paddr = pml4_paddr;
self.threads[i].kernel_stack_top = kernel_stack_top;
self.threads[i].user_rsp = user_rsp;
self.threads[i].user_rip = entry_rip;
// Setup initial kernel stack for IRETQ return
unsafe {
let kstack = kernel_stack_top as *mut u64;
let stack_ptr = kstack.sub(16);
*stack_ptr.add(0) = 0; // r15
*stack_ptr.add(1) = 0; // r14
*stack_ptr.add(2) = 0; // r13
*stack_ptr.add(3) = 0; // r12
*stack_ptr.add(4) = 0; // rbx
*stack_ptr.add(5) = 0; // rbp
*stack_ptr.add(6) = initial_user_trampoline as *const () as usize as u64; // RIP for switch_to
self.threads[i].context.rsp = stack_ptr as u64;
}
kprintln!("[SCHED] Created User Thread [TID {}] '{}' (RIP: {:#x}, RSP: {:#x})",
tid, name, entry_rip, user_rsp);
return Some(tid);
}
}
None
}
pub fn get_current_mut(&mut self) -> &mut Thread {
&mut self.threads[self.current_tid]
}
pub fn get_thread_mut(&mut self, tid: u64) -> Option<&mut Thread> {
for t in self.threads.iter_mut() {
if t.id == tid && t.state != ThreadState::Unused {
return Some(t);
}
}
None
}
pub fn block_current(&mut self, new_state: ThreadState) {
self.threads[self.current_tid].state = new_state;
self.schedule();
}
pub fn unblock(&mut self, tid: u64) {
if let Some(t) = self.get_thread_mut(tid) {
if t.state != ThreadState::Unused && t.state != ThreadState::Dead {
t.state = ThreadState::Ready;
}
}
}
pub fn yield_current(&mut self) {
if self.threads[self.current_tid].state == ThreadState::Running {
self.threads[self.current_tid].state = ThreadState::Ready;
}
self.schedule();
}
pub fn schedule(&mut self) {
let prev_idx = self.current_tid;
let mut next_idx = (prev_idx + 1) % MAX_THREADS;
// Find next Ready thread (Round-Robin)
let mut found = false;
for _ in 0..MAX_THREADS {
if self.threads[next_idx].state == ThreadState::Ready {
found = true;
break;
}
next_idx = (next_idx + 1) % MAX_THREADS;
}
if !found {
if self.threads[prev_idx].state == ThreadState::Running {
return;
}
next_idx = 0; // idle thread
}
if prev_idx == next_idx && self.threads[prev_idx].state == ThreadState::Running {
return;
}
if self.threads[prev_idx].state == ThreadState::Running {
self.threads[prev_idx].state = ThreadState::Ready;
}
self.current_tid = next_idx;
self.threads[next_idx].state = ThreadState::Running;
let next_pml4 = self.threads[next_idx].pml4_paddr;
let next_kstack = self.threads[next_idx].kernel_stack_top;
let next_rsp = self.threads[next_idx].context.rsp;
let prev_rsp_ptr = core::ptr::addr_of_mut!(self.threads[prev_idx].context.rsp);
unsafe {
if next_pml4 != 0 {
(*core::ptr::addr_of_mut!(VMM)).load_cr3(next_pml4);
}
if next_kstack != 0 {
set_kernel_stack(next_kstack);
}
switch_to(prev_rsp_ptr, next_rsp);
}
}
}
pub fn init() {
crate::kprintln!("[SCHED] Scheduler initialized (Round-Robin preemptive stub).");
unsafe {
(*core::ptr::addr_of_mut!(SCHEDULER)).init();
}
}
#[no_mangle]
pub extern "C" fn initial_user_trampoline() {
let (entry_rip, user_rsp) = unsafe {
let sched = &mut *core::ptr::addr_of_mut!(SCHEDULER);
let curr = sched.get_current_mut();
(curr.user_rip, curr.user_rsp)
};
unsafe {
// IRETQ to Ring 3
core::arch::asm!(
"mov ax, 0x1B",
"mov ds, ax",
"mov es, ax",
"mov fs, ax",
"mov gs, ax",
"push 0x1B", // User SS
"push {user_rsp}", // User RSP
"push 0x0202", // RFLAGS (IF enabled)
"push 0x23", // User CS
"push {entry_rip}", // User RIP
"xor rax, rax",
"xor rbx, rbx",
"xor rcx, rcx",
"xor rdx, rdx",
"xor rsi, rsi",
"xor rdi, rdi",
"xor rbp, rbp",
"xor r8, r8",
"xor r9, r9",
"xor r10, r10",
"xor r11, r11",
"xor r12, r12",
"xor r13, r13",
"xor r14, r14",
"xor r15, r15",
"iretq",
user_rsp = in(reg) user_rsp,
entry_rip = in(reg) entry_rip,
options(noreturn)
);
}
}
+79
View File
@@ -0,0 +1,79 @@
//! Thread Control Block (TCB) & Context Definition
use crate::ipc::fastpath::FastPathMessage;
pub const MAX_THREADS: usize = 16;
pub const KERNEL_STACK_SIZE: usize = 16384; // 16 KiB
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ThreadState {
Unused,
Ready,
Running,
BlockedOnSend(u64), // Waiting on destination endpoint ID
BlockedOnRecv(u64), // Waiting for messages on endpoint ID
BlockedOnReply(u64), // Client waiting for specific server TID reply
Dead,
}
#[repr(C)]
#[derive(Clone, Copy)]
pub struct SavedContext {
pub r15: u64,
pub r14: u64,
pub r13: u64,
pub r12: u64,
pub rbx: u64,
pub rbp: u64,
pub rsp: u64, // Saved kernel RSP
}
impl SavedContext {
pub const fn empty() -> Self {
Self {
r15: 0,
r14: 0,
r13: 0,
r12: 0,
rbx: 0,
rbp: 0,
rsp: 0,
}
}
}
pub struct Thread {
pub id: u64,
pub name: &'static str,
pub state: ThreadState,
pub context: SavedContext,
pub pml4_paddr: u64,
pub kernel_stack_base: u64,
pub kernel_stack_top: u64,
pub user_rsp: u64,
pub user_rip: u64,
pub ipc_buffer: FastPathMessage,
pub caller_tid: u64, // TID of client waiting for reply
}
impl Thread {
pub const fn empty() -> Self {
Self {
id: 0,
name: "",
state: ThreadState::Unused,
context: SavedContext::empty(),
pml4_paddr: 0,
kernel_stack_base: 0,
kernel_stack_top: 0,
user_rsp: 0,
user_rip: 0,
ipc_buffer: FastPathMessage {
dest_cap: 0,
opcode: 0,
args: [0; 4],
},
caller_tid: 0,
}
}
}