Skip to main content

strat9_kernel/vfs/
pipe.rs

1//! Kernel pipe implementation.
2//!
3//! A pipe is a unidirectional byte stream between two file descriptors.
4//! The read end blocks when empty; the write end returns EPIPE when
5//! the read end is closed.
6
7use crate::{
8    sync::{waitqueue::WaitQueue, SpinLock},
9    syscall::error::SyscallError,
10};
11use alloc::sync::Arc;
12
13const PIPE_BUF_SIZE: usize = 4096;
14
15/// Shared state for one pipe instance.
16struct PipeInner {
17    buf: [u8; PIPE_BUF_SIZE],
18    read_pos: usize,
19    write_pos: usize,
20    /// Number of bytes currently buffered.
21    len: usize,
22    read_closed: bool,
23    write_closed: bool,
24    /// Number of open file-descriptions referencing the read end.
25    /// The end is marked closed only when this reaches zero.
26    read_refs: usize,
27    /// Number of open file-descriptions referencing the write end.
28    write_refs: usize,
29}
30
31impl PipeInner {
32    /// Creates a new instance.
33    fn new() -> Self {
34        PipeInner {
35            buf: [0u8; PIPE_BUF_SIZE],
36            read_pos: 0,
37            write_pos: 0,
38            len: 0,
39            read_closed: false,
40            write_closed: false,
41            read_refs: 1,
42            write_refs: 1,
43        }
44    }
45
46    /// Returns whether empty.
47    #[allow(dead_code)]
48    fn is_empty(&self) -> bool {
49        self.len == 0
50    }
51
52    /// Returns whether full.
53    #[allow(dead_code)]
54    fn is_full(&self) -> bool {
55        self.len >= PIPE_BUF_SIZE
56    }
57
58    /// Performs the available read operation.
59    fn available_read(&self) -> usize {
60        self.len
61    }
62
63    /// Performs the available write operation.
64    fn available_write(&self) -> usize {
65        PIPE_BUF_SIZE - self.len
66    }
67}
68
69/// Shared pipe handle.
70pub struct Pipe {
71    inner: SpinLock<PipeInner>,
72    /// Woken when data becomes available (or write end is closed).
73    readers: WaitQueue,
74    /// Woken when space becomes available (or read end is closed).
75    writers: WaitQueue,
76}
77
78impl Pipe {
79    /// Creates a new instance.
80    pub fn new() -> Arc<Self> {
81        Arc::new(Pipe {
82            inner: SpinLock::new(PipeInner::new()),
83            readers: WaitQueue::new(),
84            writers: WaitQueue::new(),
85        })
86    }
87
88    /// Read from the pipe. Returns 0 on EOF (write end closed + empty).
89    /// If `non_block` is true and no data is available, returns EAGAIN.
90    pub fn read(&self, buf: &mut [u8], non_block: bool) -> Result<usize, SyscallError> {
91        loop {
92            // Phase 1: try to consume data without blocking.
93            {
94                let mut inner = self.inner.lock();
95                if inner.available_read() > 0 {
96                    let to_read = core::cmp::min(buf.len(), inner.available_read());
97                    for i in 0..to_read {
98                        buf[i] = inner.buf[inner.read_pos];
99                        inner.read_pos = (inner.read_pos + 1) % PIPE_BUF_SIZE;
100                    }
101                    inner.len -= to_read;
102                    self.writers.wake_one();
103                    return Ok(to_read);
104                }
105                if inner.write_closed {
106                    return Ok(0); // EOF
107                }
108                if non_block {
109                    return Err(SyscallError::Again);
110                }
111            }
112
113            // Phase 2: wait for data or close.
114            let result = self.readers.wait_until(|| {
115                let inner = self.inner.lock();
116                if inner.available_read() > 0 {
117                    return Some(true);
118                }
119                if inner.write_closed {
120                    return Some(false);
121                }
122                None
123            });
124
125            match result {
126                true => continue,
127                false => return Ok(0),
128            }
129        }
130    }
131
132    /// Write to the pipe. Returns EPIPE if read end is closed.
133    /// If `non_block` is true and buffer is full, returns EAGAIN.
134    pub fn write(&self, buf: &[u8], non_block: bool) -> Result<usize, SyscallError> {
135        if buf.is_empty() {
136            return Ok(0);
137        }
138
139        let mut total = 0;
140        while total < buf.len() {
141            let non_block = non_block;
142            let wrote_some = self.writers.wait_until(|| {
143                let mut inner = self.inner.lock();
144
145                if inner.read_closed {
146                    if total > 0 {
147                        return Some(Ok(0));
148                    }
149                    return Some(Err(SyscallError::Pipe));
150                }
151
152                if inner.available_write() > 0 {
153                    let to_write = core::cmp::min(buf.len() - total, inner.available_write());
154                    for i in 0..to_write {
155                        let wp = inner.write_pos;
156                        inner.buf[wp] = buf[total + i];
157                        inner.write_pos = (inner.write_pos + 1) % PIPE_BUF_SIZE;
158                    }
159                    inner.len += to_write;
160                    total += to_write;
161                    self.readers.wake_one();
162                    return Some(Ok(0));
163                }
164
165                if non_block {
166                    return Some(Err(SyscallError::Again));
167                }
168
169                None
170            });
171
172            match wrote_some {
173                Ok(_) => {
174                    if total >= buf.len() {
175                        break;
176                    }
177                }
178                // Transient: buffer full in non-blocking mode : surface EAGAIN.
179                Err(SyscallError::Again) if non_block => return Err(SyscallError::Again),
180                // F13 FIX (L2 host harness, testing-findings.md): EPIPE from
181                // a closed read end was silently swallowed by the previous
182                // catch-all arm, spinning the writing task forever at 100%
183                // CPU instead of failing with EPIPE. Propagate it.
184                Err(e) => return Err(e),
185            }
186        }
187        Ok(total)
188    }
189
190    /// Increment the read-end refcount (called on dup/fork).
191    pub fn dup_read(&self) {
192        self.inner.lock().read_refs += 1;
193    }
194
195    /// Increment the write-end refcount (called on dup/fork).
196    pub fn dup_write(&self) {
197        self.inner.lock().write_refs += 1;
198    }
199
200    /// Decrement the read-end refcount; marks the end closed only when it
201    /// reaches zero.  Returns true if the end was actually closed.
202    pub fn close_read(&self) -> bool {
203        let mut inner = self.inner.lock();
204        if inner.read_refs == 0 {
205            return false;
206        }
207        inner.read_refs -= 1;
208        if inner.read_refs == 0 {
209            inner.read_closed = true;
210            drop(inner);
211            self.writers.wake_all();
212            true
213        } else {
214            false
215        }
216    }
217
218    /// Decrement the write-end refcount; marks the end closed only when it
219    /// reaches zero.  Returns true if the end was actually closed.
220    pub fn close_write(&self) -> bool {
221        let mut inner = self.inner.lock();
222        if inner.write_refs == 0 {
223            return false;
224        }
225        inner.write_refs -= 1;
226        if inner.write_refs == 0 {
227            inner.write_closed = true;
228            drop(inner);
229            self.readers.wake_all();
230            true
231        } else {
232            false
233        }
234    }
235}
236
237// ============================================================================
238// Pipe as a VFS Scheme
239// ============================================================================
240
241use super::scheme::{
242    finalize_pseudo_stat, DirEntry, FileStat, OpenFlags, OpenResult, Scheme, DEV_PIPEFS,
243};
244use alloc::{collections::BTreeMap, vec::Vec};
245use core::sync::atomic::{AtomicU64, Ordering};
246
247/// A scheme that manages kernel pipes.
248///
249/// Each pipe gets two file_ids: even = read end, odd = write end.
250pub struct PipeScheme {
251    pipes: SpinLock<BTreeMap<u64, Arc<Pipe>>>,
252    /// Per-fd open flags (for O_NONBLOCK).
253    flags: SpinLock<BTreeMap<u64, OpenFlags>>,
254}
255
256static NEXT_PIPE_ID: AtomicU64 = AtomicU64::new(2); // Start at 2 (even numbers)
257
258impl PipeScheme {
259    /// Creates a new instance.
260    pub fn new() -> Self {
261        PipeScheme {
262            pipes: SpinLock::new(BTreeMap::new()),
263            flags: SpinLock::new(BTreeMap::new()),
264        }
265    }
266
267    /// Create a new pipe pair. Returns (read_file_id, write_file_id).
268    pub fn create_pipe(&self) -> (u64, Arc<Pipe>) {
269        let base_id = NEXT_PIPE_ID.fetch_add(2, Ordering::SeqCst);
270        let pipe = Pipe::new();
271        self.pipes.lock().insert(base_id, pipe.clone());
272        (base_id, pipe)
273    }
274
275    /// Store open flags for a pipe fd.
276    pub fn set_flags(&self, file_id: u64, flags: OpenFlags) {
277        self.flags.lock().insert(file_id, flags);
278    }
279
280    /// Returns pipe.
281    fn get_pipe(&self, file_id: u64) -> Result<Arc<Pipe>, SyscallError> {
282        let base = file_id & !1; // Even = base
283        self.pipes
284            .lock()
285            .get(&base)
286            .cloned()
287            .ok_or(SyscallError::BadHandle)
288    }
289
290    /// Returns whether read end.
291    fn is_read_end(file_id: u64) -> bool {
292        file_id & 1 == 0
293    }
294}
295
296impl Scheme for PipeScheme {
297    /// Performs the open operation.
298    fn open(&self, _path: &str, _flags: OpenFlags) -> Result<OpenResult, SyscallError> {
299        Err(SyscallError::NotSupported) // Pipes are created via sys_pipe, not open()
300    }
301
302    /// Performs the read operation.
303    fn read(&self, file_id: u64, _offset: u64, buf: &mut [u8]) -> Result<usize, SyscallError> {
304        if !Self::is_read_end(file_id) {
305            return Err(SyscallError::PermissionDenied);
306        }
307        let pipe = self.get_pipe(file_id)?;
308        let non_block = self
309            .flags
310            .lock()
311            .get(&file_id)
312            .map(|f| f.contains(OpenFlags::NONBLOCK))
313            .unwrap_or(false);
314        pipe.read(buf, non_block)
315    }
316
317    /// Performs the write operation.
318    fn write(&self, file_id: u64, _offset: u64, buf: &[u8]) -> Result<usize, SyscallError> {
319        if Self::is_read_end(file_id) {
320            return Err(SyscallError::PermissionDenied);
321        }
322        let pipe = self.get_pipe(file_id)?;
323        let non_block = self
324            .flags
325            .lock()
326            .get(&file_id)
327            .map(|f| f.contains(OpenFlags::NONBLOCK))
328            .unwrap_or(false);
329        pipe.write(buf, non_block)
330    }
331
332    /// Performs the close operation.
333    fn close(&self, file_id: u64) -> Result<(), SyscallError> {
334        let pipe = self.get_pipe(file_id)?;
335        if Self::is_read_end(file_id) {
336            pipe.close_read();
337        } else {
338            pipe.close_write();
339        }
340
341        // Remove the shared Pipe entry only when both ends are fully closed
342        // (both refcounts have reached zero).
343        let base = file_id & !1;
344        let inner = pipe.inner.lock();
345        if inner.read_closed && inner.write_closed {
346            drop(inner);
347            self.pipes.lock().remove(&base);
348        }
349        Ok(())
350    }
351
352    /// Performs the stat operation.
353    fn stat(&self, file_id: u64) -> Result<FileStat, SyscallError> {
354        let pipe = self.get_pipe(file_id)?;
355        let inner = pipe.inner.lock();
356        Ok(finalize_pseudo_stat(
357            FileStat {
358                st_ino: file_id,
359                st_mode: 0o010600, // S_IFIFO | rw-------
360                st_nlink: 1,
361                st_size: inner.len as u64,
362                st_blksize: PIPE_BUF_SIZE as u64,
363                st_blocks: 0,
364                ..FileStat::zeroed()
365            },
366            DEV_PIPEFS,
367            0,
368        ))
369    }
370
371    /// Performs the readdir operation.
372    fn readdir(&self, _file_id: u64) -> Result<Vec<DirEntry>, SyscallError> {
373        Err(SyscallError::InvalidArgument)
374    }
375
376    /// Performs the size operation.
377    fn size(&self, file_id: u64) -> Result<u64, SyscallError> {
378        let pipe = self.get_pipe(file_id)?;
379        let len = pipe.inner.lock().len;
380        Ok(len as u64)
381    }
382}