os/vfs/impls/
pipe_file.rs1use crate::sync::SpinLock;
6use crate::vfs::{File, FileMode, FsError, InodeMetadata, InodeType, OpenFlags, TimeSpec};
7use alloc::collections::VecDeque;
8use alloc::sync::Arc;
9
10struct PipeRingBuffer {
14 buffer: VecDeque<u8>,
16 capacity: usize,
18 write_end_count: usize,
20 read_end_count: usize,
22}
23
24impl PipeRingBuffer {
25 const DEFAULT_CAPACITY: usize = 4096; const MIN_CAPACITY: usize = 4096; const MAX_CAPACITY: usize = 1048576; 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 fn get_capacity(&self) -> usize {
40 self.capacity
41 }
42
43 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 if new_capacity < self.buffer.len() {
51 return Err(FsError::InvalidArgument);
52 }
53
54 self.capacity = new_capacity;
55 Ok(())
56 }
57
58 fn read(&mut self, buf: &mut [u8]) -> Result<usize, FsError> {
60 if self.buffer.is_empty() {
61 if self.write_end_count == 0 {
63 return Ok(0);
64 }
65 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 fn write(&mut self, buf: &[u8]) -> Result<usize, FsError> {
79 if self.read_end_count == 0 {
81 return Err(FsError::BrokenPipe);
82 }
83
84 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
99pub struct PipeFile {
106 buffer: Arc<SpinLock<PipeRingBuffer>>,
108 end_type: PipeEnd,
110 flags: SpinLock<OpenFlags>,
112 owner: SpinLock<Option<i32>>,
114}
115
116#[derive(Debug, Clone, Copy, PartialEq, Eq)]
117enum PipeEnd {
118 Read,
119 Write,
120}
121
122impl PipeFile {
123 pub fn create_pair() -> (Self, Self) {
132 let buffer = Arc::new(SpinLock::new(PipeRingBuffer::new()));
133
134 {
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 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 pub fn get_pipe_size(&self) -> usize {
167 self.buffer.lock().get_capacity()
168 }
169
170 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 !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 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 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 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 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 fn as_any(&self) -> &dyn core::any::Any {
272 self
273 }
274}
275
276impl Drop for PipeFile {
277 fn drop(&mut self) {
278 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}