| 148 | #[async_trait::async_trait] |
| 149 | impl InputStream for WasiStdin { |
| 150 | fn read(&mut self, size: usize) -> Result<Bytes, StreamError> { |
| 151 | if size == 0 { |
| 152 | return Ok(Bytes::new()); |
| 153 | } |
| 154 | let g = GlobalStdin::get(); |
| 155 | let mut locked = g.state.lock().unwrap(); |
| 156 | match mem::replace(&mut *locked, StdinState::ReadRequested(size)) { |
| 157 | StdinState::ReadNotRequested => { |
| 158 | g.read_requested.notify_one(); |
| 159 | Ok(Bytes::new()) |
| 160 | } |
| 161 | StdinState::ReadRequested(prev_size) => { |
| 162 | // Preserve the larger of the two requested sizes |
| 163 | // so the worker thread allocates an adequate buffer. |
| 164 | *locked = StdinState::ReadRequested(prev_size.max(size)); |
| 165 | Ok(Bytes::new()) |
| 166 | } |
| 167 | StdinState::Data(mut data) => { |
| 168 | let size = data.len().min(size); |
| 169 | let bytes = data.split_to(size); |
| 170 | *locked = if data.is_empty() { |
| 171 | StdinState::ReadNotRequested |
| 172 | } else { |
| 173 | StdinState::Data(data) |
| 174 | }; |
| 175 | Ok(bytes.freeze()) |
| 176 | } |
| 177 | StdinState::Error(e) => { |
| 178 | *locked = StdinState::Closed; |
| 179 | Err(StreamError::LastOperationFailed(e.into())) |
| 180 | } |
| 181 | StdinState::Closed => { |
| 182 | *locked = StdinState::Closed; |
| 183 | Err(StreamError::Closed) |
| 184 | } |
| 185 | } |
| 186 | } |
| 187 | } |
| 188 | |
| 189 | #[async_trait::async_trait] |