shadow_rs/host/descriptor/
pipe.rs1use std::sync::Arc;
2
3use atomic_refcell::AtomicRefCell;
4use linux_api::errno::Errno;
5use linux_api::ioctls::IoctlRequest;
6use linux_api::stat::SFlag;
7use shadow_shim_helper_rs::syscall_types::ForeignPtr;
8
9use crate::cshadow as c;
10use crate::host::descriptor::listener::{StateEventSource, StateListenHandle, StateListenerFilter};
11use crate::host::descriptor::shared_buf::{
12 BufferHandle, BufferSignals, BufferState, ReaderHandle, SharedBuf, WriterHandle,
13};
14use crate::host::descriptor::{FileMode, FileSignals, FileState, FileStatus};
15use crate::host::memory_manager::MemoryManager;
16use crate::host::syscall::io::{IoVec, IoVecReader, IoVecWriter};
17use crate::host::syscall::types::{SyscallError, SyscallResult};
18use crate::utility::HostTreePointer;
19use crate::utility::callback_queue::CallbackQueue;
20
21pub struct Pipe {
22 buffer: Option<Arc<AtomicRefCell<SharedBuf>>>,
23 event_source: StateEventSource,
24 state: FileState,
25 mode: FileMode,
26 status: FileStatus,
27 write_mode: WriteMode,
28 buffer_event_handle: Option<BufferHandle>,
29 reader_handle: Option<ReaderHandle>,
30 writer_handle: Option<WriterHandle>,
31 has_open_file: bool,
34}
35
36impl Pipe {
37 pub fn new(mode: FileMode, status: FileStatus) -> Self {
40 Self {
41 buffer: None,
42 event_source: StateEventSource::new(),
43 state: FileState::ACTIVE,
44 mode,
45 status,
46 write_mode: WriteMode::Stream,
47 buffer_event_handle: None,
48 reader_handle: None,
49 writer_handle: None,
50 has_open_file: false,
51 }
52 }
53
54 pub fn status(&self) -> FileStatus {
55 self.status
56 }
57
58 pub fn set_status(&mut self, status: FileStatus) {
59 self.status = status;
60 }
61
62 pub fn mode(&self) -> FileMode {
63 self.mode
64 }
65
66 pub fn has_open_file(&self) -> bool {
67 self.has_open_file
68 }
69
70 pub fn supports_sa_restart(&self) -> bool {
71 true
72 }
73
74 pub fn set_has_open_file(&mut self, val: bool) {
75 self.has_open_file = val;
76 }
77
78 pub fn max_size(&self) -> usize {
79 self.buffer.as_ref().unwrap().borrow().max_len()
80 }
81
82 pub fn close(&mut self, cb_queue: &mut CallbackQueue) -> Result<(), SyscallError> {
83 if self.state.contains(FileState::CLOSED) {
84 log::warn!("Attempting to close an already-closed pipe");
85 }
86
87 if let Some(h) = self.buffer_event_handle.take() {
89 h.stop_listening()
90 }
91
92 if let Some(writer_handle) = self.writer_handle.take() {
94 self.buffer
95 .as_ref()
96 .unwrap()
97 .borrow_mut()
98 .remove_writer(writer_handle, cb_queue);
99 }
100
101 if let Some(reader_handle) = self.reader_handle.take() {
103 self.buffer
104 .as_ref()
105 .unwrap()
106 .borrow_mut()
107 .remove_reader(reader_handle, cb_queue);
108 }
109
110 self.buffer = None;
112
113 self.update_state(
115 FileState::CLOSED | FileState::ACTIVE | FileState::READABLE | FileState::WRITABLE,
116 FileState::CLOSED,
117 FileSignals::empty(),
118 cb_queue,
119 );
120
121 Ok(())
122 }
123
124 pub fn readv(
125 &mut self,
126 iovs: &[IoVec],
127 offset: Option<libc::off_t>,
128 _flags: libc::c_int,
129 mem: &mut MemoryManager,
130 cb_queue: &mut CallbackQueue,
131 ) -> Result<libc::ssize_t, SyscallError> {
132 if offset.is_some() {
134 return Err(linux_api::errno::Errno::ESPIPE.into());
135 }
136
137 if !self.mode.contains(FileMode::READ) {
139 return Err(linux_api::errno::Errno::EBADF.into());
140 }
141
142 let num_bytes_to_read: libc::size_t = iovs.iter().map(|x| x.len).sum();
143
144 let mut writer = IoVecWriter::new(iovs, mem);
145
146 let (num_copied, _num_removed_from_buf) = self
147 .buffer
148 .as_ref()
149 .unwrap()
150 .borrow_mut()
151 .read(&mut writer, cb_queue)?;
152
153 if num_copied == 0
158 && num_bytes_to_read != 0
159 && self.buffer.as_ref().unwrap().borrow().num_writers() > 0
160 {
161 Err(Errno::EWOULDBLOCK.into())
162 } else {
163 Ok(num_copied.try_into().unwrap())
164 }
165 }
166
167 pub fn writev(
168 &mut self,
169 iovs: &[IoVec],
170 offset: Option<libc::off_t>,
171 _flags: libc::c_int,
172 mem: &mut MemoryManager,
173 cb_queue: &mut CallbackQueue,
174 ) -> Result<libc::ssize_t, SyscallError> {
175 if offset.is_some() {
177 return Err(linux_api::errno::Errno::ESPIPE.into());
178 }
179
180 if !self.mode.contains(FileMode::WRITE) {
182 return Err(linux_api::errno::Errno::EBADF.into());
183 }
184
185 let mut buffer = self.buffer.as_ref().unwrap().borrow_mut();
186
187 if buffer.num_readers() == 0 {
188 return Err(linux_api::errno::Errno::EPIPE.into());
189 }
190
191 if self.write_mode == WriteMode::Packet && !self.status.contains(FileStatus::O_DIRECT) {
192 self.write_mode = WriteMode::Stream;
194 } else if self.write_mode == WriteMode::Stream && self.status.contains(FileStatus::O_DIRECT)
195 {
196 if !buffer.has_data() {
200 self.write_mode = WriteMode::Packet;
201 }
202 }
203
204 let len: libc::size_t = iovs.iter().map(|x| x.len).sum();
205
206 let mut reader = IoVecReader::new(iovs, mem);
207
208 let num_copied = match self.write_mode {
209 WriteMode::Stream => buffer.write_stream(&mut reader, len, cb_queue)?,
210 WriteMode::Packet => {
211 let mut num_written = 0;
212
213 loop {
214 let bytes_remaining = len - num_written;
216
217 if bytes_remaining == 0 {
219 break num_written;
220 }
221
222 let bytes_to_write = std::cmp::min(bytes_remaining, libc::PIPE_BUF);
224
225 if let Err(e) = buffer.write_packet(&mut reader, bytes_to_write, cb_queue) {
226 if num_written > 0 {
228 break num_written;
229 }
230 return Err(e.into());
231 }
232
233 num_written += bytes_to_write;
234 }
235 }
236 };
237
238 Ok(num_copied.try_into().unwrap())
239 }
240
241 pub fn ioctl(
242 &mut self,
243 request: IoctlRequest,
244 _arg_ptr: ForeignPtr<()>,
245 _memory_manager: &mut MemoryManager,
246 ) -> SyscallResult {
247 log::warn!("We do not yet handle ioctl request {request:?} on pipes");
248 Err(Errno::EINVAL.into())
249 }
250
251 pub fn lseek(
252 &mut self,
253 _off: linux_api::posix_types::kernel_off_t,
254 _whence: linux_api::unistd::LSeekWhence,
255 ) -> Result<linux_api::posix_types::kernel_off_t, SyscallError> {
256 Err(Errno::ESPIPE.into())
257 }
258
259 pub fn stat(&self) -> Result<linux_api::stat::stat, SyscallError> {
260 warn_once_then_debug!("Not all fields of 'struct stat' are implemented for pipes");
261
262 Ok(linux_api::stat::stat {
263 lst_dev: 0,
266 lst_ino: 0,
267 lst_nlink: 1,
269 lst_mode: (SFlag::S_IFIFO | SFlag::S_IRUSR | SFlag::S_IWUSR).bits(),
273 lst_uid: 0,
276 lst_gid: 0,
277 lst_rdev: 0,
278 lst_size: 0,
281 lst_blksize: 0,
283 lst_blocks: 0,
284 lst_atime: 0,
285 lst_atime_nsec: 0,
286 lst_mtime: 0,
287 lst_mtime_nsec: 0,
288 lst_ctime: 0,
289 lst_ctime_nsec: 0,
290 ..shadow_pod::zeroed()
293 })
294 }
295
296 pub fn connect_to_buffer(
297 arc: &Arc<AtomicRefCell<Self>>,
298 buffer: Arc<AtomicRefCell<SharedBuf>>,
299 cb_queue: &mut CallbackQueue,
300 ) {
301 let weak = Arc::downgrade(arc);
302 let pipe = &mut *arc.borrow_mut();
303
304 pipe.buffer = Some(buffer);
305
306 if pipe.mode.contains(FileMode::WRITE) {
307 pipe.writer_handle = Some(
308 pipe.buffer
309 .as_ref()
310 .unwrap()
311 .borrow_mut()
312 .add_writer(cb_queue),
313 );
314 }
315
316 if pipe.mode.contains(FileMode::READ) {
317 pipe.reader_handle = Some(
318 pipe.buffer
319 .as_ref()
320 .unwrap()
321 .borrow_mut()
322 .add_reader(cb_queue),
323 );
324 }
325
326 let mut monitoring_state = BufferState::empty();
328
329 if pipe.mode.contains(FileMode::READ) {
332 monitoring_state.insert(BufferState::READABLE);
333 monitoring_state.insert(BufferState::NO_WRITERS);
334 }
335
336 if pipe.mode.contains(FileMode::WRITE) {
339 monitoring_state.insert(BufferState::WRITABLE);
340 monitoring_state.insert(BufferState::NO_READERS);
341 }
342
343 let monitoring_signals = BufferSignals::BUFFER_GREW;
355
356 let handle = pipe.buffer.as_ref().unwrap().borrow_mut().add_listener(
357 monitoring_state,
358 monitoring_signals,
359 move |buffer_state, buffer_signals, cb_queue| {
360 if let Some(pipe) = weak.upgrade() {
362 let mut pipe = pipe.borrow_mut();
363
364 pipe.align_state_to_buffer(buffer_state, buffer_signals, cb_queue);
366 }
367 },
368 );
369
370 pipe.buffer_event_handle = Some(handle);
371
372 let buffer_state = pipe.buffer.as_ref().unwrap().borrow().state();
374 pipe.align_state_to_buffer(buffer_state, BufferSignals::empty(), cb_queue);
375 }
376
377 pub fn add_listener(
378 &mut self,
379 monitoring_state: FileState,
380 monitoring_signals: FileSignals,
381 filter: StateListenerFilter,
382 notify_fn: impl Fn(FileState, FileState, FileSignals, &mut CallbackQueue)
383 + Send
384 + Sync
385 + 'static,
386 ) -> StateListenHandle {
387 self.event_source
388 .add_listener(monitoring_state, monitoring_signals, filter, notify_fn)
389 }
390
391 pub fn add_legacy_listener(&mut self, ptr: HostTreePointer<c::StatusListener>) {
392 self.event_source.add_legacy_listener(ptr);
393 }
394
395 pub fn remove_legacy_listener(&mut self, ptr: *mut c::StatusListener) {
396 self.event_source.remove_legacy_listener(ptr);
397 }
398
399 pub fn state(&self) -> FileState {
400 self.state
401 }
402
403 fn align_state_to_buffer(
408 &mut self,
409 buffer_state: BufferState,
410 buffer_signals: BufferSignals,
411 cb_queue: &mut CallbackQueue,
412 ) {
413 let mut mask = FileState::empty();
414 let mut file_state = FileState::empty();
415 let mut file_signals = FileSignals::empty();
416
417 if self.state.contains(FileState::CLOSED) {
419 return;
420 }
421
422 if self.mode.contains(FileMode::READ) {
424 mask.insert(FileState::READABLE);
425 if buffer_state.intersects(BufferState::READABLE | BufferState::NO_WRITERS) {
427 file_state.insert(FileState::READABLE);
428 }
429 if buffer_signals.intersects(BufferSignals::BUFFER_GREW) {
430 file_signals.insert(FileSignals::READ_BUFFER_GREW);
431 }
432 }
433
434 if self.mode.contains(FileMode::WRITE) {
436 mask.insert(FileState::WRITABLE);
437 if buffer_state.intersects(BufferState::WRITABLE | BufferState::NO_READERS) {
439 file_state.insert(FileState::WRITABLE);
440 }
441 }
442
443 self.update_state(mask, file_state, file_signals, cb_queue);
445 }
446
447 fn update_state(
448 &mut self,
449 mask: FileState,
450 state: FileState,
451 signals: FileSignals,
452 cb_queue: &mut CallbackQueue,
453 ) {
454 let old_state = self.state;
455
456 self.state.remove(mask);
458 self.state.insert(state & mask);
459
460 self.handle_state_change(old_state, signals, cb_queue);
461 }
462
463 fn handle_state_change(
464 &mut self,
465 old_state: FileState,
466 signals: FileSignals,
467 cb_queue: &mut CallbackQueue,
468 ) {
469 let states_changed = self.state ^ old_state;
470
471 if states_changed.is_empty() && signals.is_empty() {
473 return;
474 }
475
476 self.event_source
477 .notify_listeners(self.state, states_changed, signals, cb_queue);
478 }
479}
480
481#[derive(Debug, PartialEq, Eq)]
482enum WriteMode {
483 Stream,
484 Packet,
485}