Skip to main content

strat9_kernel/ipc/
reply.rs

1//! IPC call/reply support : both synchronous (blocking) and async (ring-based).
2//!
3//! # Design  concept
4//!
5//! -  `ReplyTarget::Sync` is analogous to an seL4 endpoint call+reply
6//! -  `ReplyTarget::AsyncRing` combines both: the ring *is* the completion port
7//! -   Send/Receive/Reply : exactly our sync `wait_for_reply`/`deliver_reply`
8//! -   The CQE `user_data` correlates the reply to the original submission
9
10use super::message::IpcMessage;
11use crate::{
12    async_io::{complete::push_completion_for_ring, ring::find_ring},
13    memory::UserSliceWrite,
14    process::TaskId,
15    sync::{SpinLock, WaitQueue},
16};
17use alloc::{collections::BTreeMap, sync::Arc, vec::Vec};
18
19// ===========================================================================
20// ReplyTarget : where to route the reply
21// ===========================================================================
22
23/// Describes how to deliver a reply once the server responds.
24enum ReplyTarget {
25    Sync {
26        msg: Option<IpcMessage>,
27        waitq: Arc<WaitQueue>,
28    },
29
30    AsyncRing {
31        ring_id: u64,
32        user_data: u64,
33        reply_buf: u64,
34    },
35}
36
37struct ReplySlot {
38    target: ReplyTarget,
39
40    waiting_on: Option<TaskId>,
41}
42
43struct ReplyRegistry {
44    slots: BTreeMap<TaskId, ReplySlot>,
45}
46
47impl ReplyRegistry {
48    const fn new() -> Self {
49        ReplyRegistry {
50            slots: BTreeMap::new(),
51        }
52    }
53}
54
55static REPLIES: SpinLock<ReplyRegistry> = SpinLock::new(ReplyRegistry::new());
56
57// ===========================================================================
58// Helpers
59// ===========================================================================
60
61fn epipe_reply() -> IpcMessage {
62    let mut err = IpcMessage::new(0x80);
63    let epipe: u32 = 32;
64    err.payload[0..4].copy_from_slice(&epipe.to_le_bytes());
65    err
66}
67
68// ===========================================================================
69// Public API
70// ===========================================================================
71
72/// Block the current task waiting for a reply message (synchronous call path).
73///
74/// The caller (via `SYS_IPC_CALL`) blocks on a `WaitQueue` until the server
75/// calls `deliver_reply`.  Returns the reply message; returns an EPIPE error
76/// if the slot was removed while waiting (server died).
77pub fn wait_for_reply(task_id: TaskId, waiting_on: TaskId) -> IpcMessage {
78    let waitq = {
79        let mut registry = REPLIES.lock();
80        let slot = registry.slots.entry(task_id).or_insert_with(|| ReplySlot {
81            target: ReplyTarget::Sync {
82                msg: None,
83                waitq: Arc::new(WaitQueue::new()),
84            },
85            waiting_on: Some(waiting_on),
86        });
87        slot.waiting_on = Some(waiting_on); // Should never happen : a task cannot be both sync-waiting
88                                            // and have an async ring pending on the same slot.
89        match &slot.target {
90            ReplyTarget::Sync { waitq, .. } => waitq.clone(),
91            ReplyTarget::AsyncRing { .. } => {
92                return epipe_reply();
93            }
94        }
95    };
96
97    let msg = waitq.wait_until(|| {
98        // P2 fix: check for pending signals to avoid livelock.
99        if crate::process::signal::has_pending_signals() {
100            return Some(epipe_reply());
101        }
102        let mut registry = REPLIES.lock();
103        match registry.slots.get_mut(&task_id) {
104            Some(ReplySlot {
105                target: ReplyTarget::Sync { msg, .. },
106                ..
107            }) => msg.take(),
108            _ => Some(epipe_reply()),
109        }
110    });
111
112    let mut registry = REPLIES.lock();
113    registry.slots.remove(&task_id);
114    msg
115}
116
117/// Register a pending ring-based `IpcCall` for `caller`.
118///
119/// When the server calls `deliver_reply(caller, msg)`, the kernel will:
120///   - Copy `msg` into the caller's buffer at `reply_buf`.
121///   - Push a CQE to the caller's ring with the given `user_data`.
122///
123
124pub fn register_ring_call(
125    caller: TaskId,
126    waiting_on: TaskId,
127    ring_id: u64,
128    user_data: u64,
129    reply_buf: u64,
130) {
131    let mut registry = REPLIES.lock();
132    registry.slots.insert(
133        caller,
134        ReplySlot {
135            target: ReplyTarget::AsyncRing {
136                ring_id,
137                user_data,
138                reply_buf,
139            },
140            waiting_on: Some(waiting_on),
141        },
142    );
143}
144
145pub fn cancel_replies_waiting_on(dead_task: TaskId) {
146    let mut actions = alloc::vec::Vec::new();
147
148    {
149        let mut registry = REPLIES.lock();
150        let to_cancel: Vec<TaskId> = registry
151            .slots
152            .iter()
153            .filter(|(_, slot)| slot.waiting_on == Some(dead_task))
154            .map(|(id, _)| *id)
155            .collect();
156
157        for id in to_cancel {
158            let slot = registry.slots.remove(&id);
159            if let Some(slot) = slot {
160                actions.push((id, slot.target));
161            }
162        }
163    }
164
165    for (_id, target) in actions {
166        match target {
167            ReplyTarget::Sync { waitq, .. } => {
168                // Slot removed from registry; the blocked task's
169                // wait_until closure will return epipe_reply() when it
170                // can't find its slot.
171                waitq.wake_all();
172            }
173            ReplyTarget::AsyncRing {
174                ring_id, user_data, ..
175            } => {
176                push_completion_for_ring_by_id(ring_id, user_data, -32, 0);
177            }
178        }
179    }
180}
181
182/// Error returned by [`deliver_reply`].
183#[derive(Debug, Clone, Copy, PartialEq, Eq)]
184pub enum DeliverError {
185    /// The target has no pending call awaiting a reply.
186    NoPendingCall,
187    /// The replying task is not the server the target is waiting on.
188    /// Without this check, any process could forge replies into another
189    /// process's blocked `SYS_IPC_CALL` (reply spoofing).
190    NotResponder,
191}
192
193/// Check whether `responder` is authorized to reply to `target`.
194///
195/// Returns `Ok(())` if the authorization check passes, or the appropriate
196/// `DeliverError` otherwise.  This is used by `sys_ipc_reply` to validate
197/// authorization **before** performing handle transfer, preventing capability
198/// injection into arbitrary processes.
199pub fn check_authorization(responder: TaskId, target: TaskId) -> Result<(), DeliverError> {
200    let registry = REPLIES.lock();
201    let slot = registry
202        .slots
203        .get(&target)
204        .ok_or(DeliverError::NoPendingCall)?;
205    if slot.waiting_on != Some(responder) {
206        return Err(DeliverError::NotResponder);
207    }
208    Ok(())
209}
210
211/// Deliver a reply message to the given task.
212///
213/// # Authorization
214///
215/// `responder` must match the `waiting_on` server recorded when the target
216/// issued its call (`SYS_IPC_CALL` records the port owner). Replies from
217/// any other task are rejected with [`DeliverError::NotResponder`] —
218/// otherwise any process able to guess a blocked task id could inject
219/// forged responses into its RPCs.
220pub fn deliver_reply(
221    responder: TaskId,
222    target: TaskId,
223    msg: IpcMessage,
224) -> Result<(), DeliverError> {
225    let mut registry = REPLIES.lock();
226
227    // Fast path: look up the slot without inserting a new one.
228    let slot = registry
229        .slots
230        .get_mut(&target)
231        .ok_or(DeliverError::NoPendingCall)?;
232
233    // Only the awaited server may answer this call.
234    if slot.waiting_on != Some(responder) {
235        return Err(DeliverError::NotResponder);
236    }
237
238    match &mut slot.target {
239        ReplyTarget::Sync {
240            msg: slot_msg,
241            waitq,
242        } => {
243            slot_msg.replace(msg);
244            let wq = waitq.clone();
245            drop(registry);
246            wq.wake_one();
247            Ok(())
248        }
249        ReplyTarget::AsyncRing {
250            ring_id,
251            user_data,
252            reply_buf,
253        } => {
254            let ring_id = *ring_id;
255            let user_data = *user_data;
256            let reply_buf = *reply_buf;
257
258            registry.slots.remove(&target);
259
260            let msg_size = core::mem::size_of::<IpcMessage>();
261            let mut raw = [0u8; core::mem::size_of::<IpcMessage>()];
262            crate::ipc::message::ipc_message_to_raw(&msg, &mut raw);
263            if let Ok(user) = UserSliceWrite::new(reply_buf, msg_size) {
264                let _ = user.copy_from(&raw);
265            }
266
267            drop(registry);
268
269            // Push a CQE to the caller's async ring.
270            push_completion_for_ring_by_id(ring_id, user_data, 0, 0);
271            Ok(())
272        }
273    }
274}
275
276/// Thin wrapper : resolve `ring_id` to `&Ring` and then push CQE.
277fn push_completion_for_ring_by_id(ring_id: u64, user_data: u64, result: i32, flags: u32) {
278    if let Some(ring) = find_ring(ring_id) {
279        push_completion_for_ring(&ring, user_data, result, flags);
280    }
281}