strat9_kernel/ipc/
reply.rs1use 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
19enum 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
57fn 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
68pub 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); 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 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
117pub 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 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
184pub enum DeliverError {
185 NoPendingCall,
187 NotResponder,
191}
192
193pub 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
211pub fn deliver_reply(
221 responder: TaskId,
222 target: TaskId,
223 msg: IpcMessage,
224) -> Result<(), DeliverError> {
225 let mut registry = REPLIES.lock();
226
227 let slot = registry
229 .slots
230 .get_mut(&target)
231 .ok_or(DeliverError::NoPendingCall)?;
232
233 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_completion_for_ring_by_id(ring_id, user_data, 0, 0);
271 Ok(())
272 }
273 }
274}
275
276fn 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}