os/vfs/impls/
pipe_file.rs

1//! 管道文件实现
2//!
3//! 管道是流式单向通信设备,读端和写端分别由两个 [`PipeFile`] 实例表示。
4
5use crate::sync::SpinLock;
6use crate::vfs::{File, FileMode, FsError, InodeMetadata, InodeType, OpenFlags, TimeSpec};
7use alloc::collections::VecDeque;
8use alloc::sync::Arc;
9
10/// 管道环形缓冲区
11///
12/// 容量默认 4KB(POSIX 最小 512 字节)。
13struct PipeRingBuffer {
14    /// 内部缓冲区
15    buffer: VecDeque<u8>,
16    /// 缓冲区容量
17    capacity: usize,
18    /// 写端引用计数 (用于检测写端关闭)
19    write_end_count: usize,
20    /// 读端引用计数 (用于检测读端关闭)
21    read_end_count: usize,
22}
23
24impl PipeRingBuffer {
25    const DEFAULT_CAPACITY: usize = 4096; // POSIX 规定最小 512 字节
26    const MIN_CAPACITY: usize = 4096; // Linux 最小管道大小
27    const MAX_CAPACITY: usize = 1048576; // Linux 最大管道大小 (1MB)
28
29    fn new() -> Self {
30        Self {
31            buffer: VecDeque::with_capacity(Self::DEFAULT_CAPACITY),
32            capacity: Self::DEFAULT_CAPACITY,
33            write_end_count: 0,
34            read_end_count: 0,
35        }
36    }
37
38    /// 获取管道容量
39    fn get_capacity(&self) -> usize {
40        self.capacity
41    }
42
43    /// 设置管道容量
44    fn set_capacity(&mut self, new_capacity: usize) -> Result<(), FsError> {
45        if new_capacity < Self::MIN_CAPACITY || new_capacity > Self::MAX_CAPACITY {
46            return Err(FsError::InvalidArgument);
47        }
48
49        // 如果新容量小于当前数据量,拒绝修改
50        if new_capacity < self.buffer.len() {
51            return Err(FsError::InvalidArgument);
52        }
53
54        self.capacity = new_capacity;
55        Ok(())
56    }
57
58    /// 读取数据 (非阻塞)
59    fn read(&mut self, buf: &mut [u8]) -> Result<usize, FsError> {
60        if self.buffer.is_empty() {
61            // 写端已关闭且缓冲区为空 -> EOF
62            if self.write_end_count == 0 {
63                return Ok(0);
64            }
65            // 写端未关闭但缓冲区为空 -> 暂时返回0 (TODO: 配合调度器实现阻塞)
66            return Ok(0);
67        }
68
69        let nread = buf.len().min(self.buffer.len());
70        for i in 0..nread {
71            buf[i] = self.buffer.pop_front().unwrap();
72        }
73
74        Ok(nread)
75    }
76
77    /// 写入数据 (非阻塞)
78    fn write(&mut self, buf: &[u8]) -> Result<usize, FsError> {
79        // 读端已关闭 -> EPIPE (应发送 SIGPIPE 信号)
80        if self.read_end_count == 0 {
81            return Err(FsError::BrokenPipe);
82        }
83
84        // 缓冲区已满 -> 暂时只写入可用空间 (TODO: 阻塞等待)
85        let available = self.capacity - self.buffer.len();
86        if available == 0 {
87            return Ok(0);
88        }
89
90        let nwrite = buf.len().min(available);
91        for &byte in &buf[..nwrite] {
92            self.buffer.push_back(byte);
93        }
94
95        Ok(nwrite)
96    }
97}
98
99/// 管道文件实现
100///
101/// 特点:
102/// - 单向数据流 (读端和写端分别创建两个 PipeFile 实例)
103/// - 流式设备 (无 offset 概念,不支持 lseek)
104/// - 不依赖 Inode (纯内存结构)
105pub struct PipeFile {
106    /// 共享的环形缓冲区
107    buffer: Arc<SpinLock<PipeRingBuffer>>,
108    /// 文件端点类型
109    end_type: PipeEnd,
110    /// 打开标志位 (支持 O_NONBLOCK 等)
111    flags: SpinLock<OpenFlags>,
112    /// 异步 I/O 所有者 PID
113    owner: SpinLock<Option<i32>>,
114}
115
116#[derive(Debug, Clone, Copy, PartialEq, Eq)]
117enum PipeEnd {
118    Read,
119    Write,
120}
121
122impl PipeFile {
123    /// 创建管道对 (返回 [读端, 写端])
124    ///
125    /// # 示例
126    /// ```rust
127    /// let (pipe_read, pipe_write) = PipeFile::create_pair();
128    /// fd_table.install_at(3, Arc::new(pipe_read) as Arc<dyn File>)?;
129    /// fd_table.install_at(4, Arc::new(pipe_write) as Arc<dyn File>)?;
130    /// ```
131    pub fn create_pair() -> (Self, Self) {
132        let buffer = Arc::new(SpinLock::new(PipeRingBuffer::new()));
133
134        // 初始化引用计数
135        {
136            let mut buf = buffer.lock();
137            buf.read_end_count = 1;
138            buf.write_end_count = 1;
139        }
140
141        let read_end = Self {
142            buffer: buffer.clone(),
143            end_type: PipeEnd::Read,
144            flags: SpinLock::new(OpenFlags::empty()),
145            owner: SpinLock::new(None),
146        };
147
148        let write_end = Self {
149            buffer,
150            end_type: PipeEnd::Write,
151            flags: SpinLock::new(OpenFlags::empty()),
152            owner: SpinLock::new(None),
153        };
154
155        (read_end, write_end)
156    }
157
158    /// 设置文件状态标志 (F_SETFL)
159    pub fn set_flags(&self, new_flags: OpenFlags) -> Result<(), FsError> {
160        let mut flags = self.flags.lock();
161        *flags = new_flags;
162        Ok(())
163    }
164
165    /// 获取管道大小 (F_GETPIPE_SZ)
166    pub fn get_pipe_size(&self) -> usize {
167        self.buffer.lock().get_capacity()
168    }
169
170    /// 设置管道大小 (F_SETPIPE_SZ)
171    pub fn set_pipe_size(&self, new_size: usize) -> Result<(), FsError> {
172        self.buffer.lock().set_capacity(new_size)
173    }
174}
175
176impl File for PipeFile {
177    fn readable(&self) -> bool {
178        if self.end_type != PipeEnd::Read {
179            return false;
180        }
181        let buf = self.buffer.lock();
182        // Readable if buffer has data OR write end is closed (EOF)
183        !buf.buffer.is_empty() || buf.write_end_count == 0
184    }
185
186    fn writable(&self) -> bool {
187        if self.end_type != PipeEnd::Write {
188            return false;
189        }
190        let buf = self.buffer.lock();
191        // Writable if buffer has space AND read end is open
192        buf.read_end_count > 0 && buf.buffer.len() < buf.capacity
193    }
194
195    fn read(&self, buf: &mut [u8]) -> Result<usize, FsError> {
196        if !self.readable() {
197            return Err(FsError::InvalidArgument);
198        }
199
200        let mut ring_buf = self.buffer.lock();
201        let result = ring_buf.read(buf);
202        // Only wake up writers if we actually freed buffer space
203        if let Ok(bytes_read) = result {
204            if bytes_read > 0 {
205                crate::kernel::syscall::io::wake_poll_waiters();
206            }
207        }
208        result
209    }
210
211    fn write(&self, buf: &[u8]) -> Result<usize, FsError> {
212        if !self.writable() {
213            return Err(FsError::InvalidArgument);
214        }
215
216        let mut ring_buf = self.buffer.lock();
217        let result = ring_buf.write(buf);
218        // Only wake up readers if we actually wrote data
219        if let Ok(bytes_written) = result {
220            if bytes_written > 0 {
221                crate::kernel::syscall::io::wake_poll_waiters();
222            }
223        }
224        result
225    }
226
227    fn metadata(&self) -> Result<InodeMetadata, FsError> {
228        // 管道没有真实的 inode,返回虚拟元数据
229        Ok(InodeMetadata {
230            inode_no: 0,
231            inode_type: InodeType::Fifo,
232            size: 0,
233            mode: FileMode::S_IFIFO | FileMode::S_IRUSR | FileMode::S_IWUSR,
234            uid: 0,
235            gid: 0,
236            atime: TimeSpec::zero(),
237            mtime: TimeSpec::zero(),
238            ctime: TimeSpec::zero(),
239            nlinks: 1,
240            blocks: 0,
241            rdev: 0,
242        })
243    }
244
245    fn flags(&self) -> OpenFlags {
246        *self.flags.lock()
247    }
248
249    fn set_status_flags(&self, new_flags: OpenFlags) -> Result<(), FsError> {
250        self.set_flags(new_flags)
251    }
252
253    fn get_pipe_size(&self) -> Result<usize, FsError> {
254        Ok(self.buffer.lock().get_capacity())
255    }
256
257    fn set_pipe_size(&self, size: usize) -> Result<(), FsError> {
258        self.buffer.lock().set_capacity(size)
259    }
260
261    fn get_owner(&self) -> Result<i32, FsError> {
262        Ok(self.owner.lock().unwrap_or(0))
263    }
264
265    fn set_owner(&self, pid: i32) -> Result<(), FsError> {
266        *self.owner.lock() = if pid == 0 { None } else { Some(pid) };
267        Ok(())
268    }
269
270    // lseek 使用默认实现 (返回 NotSupported)
271    fn as_any(&self) -> &dyn core::any::Any {
272        self
273    }
274}
275
276impl Drop for PipeFile {
277    fn drop(&mut self) {
278        // 减少引用计数
279        let mut buf = self.buffer.lock();
280        match self.end_type {
281            PipeEnd::Read => buf.read_end_count -= 1,
282            PipeEnd::Write => buf.write_end_count -= 1,
283        }
284    }
285}