| 158 | } |
| 159 | |
| 160 | async fn stream_writer( |
| 161 | streamer_ref: &StreamerRef, |
| 162 | path: &Path, |
| 163 | mut receiver: Receiver<StreamerMessage>, |
| 164 | ) -> crate::Result<()> { |
| 165 | let uid: String = rand::rng() |
| 166 | .sample_iter(&Alphanumeric) |
| 167 | .take(12) |
| 168 | .map(char::from) |
| 169 | .collect::<String>(); |
| 170 | let mut path = path.to_path_buf(); |
| 171 | path.push(format!("{uid}.{STREAM_FILE_SUFFIX}")); |
| 172 | log::debug!("Opening stream file {}", path.display()); |
| 173 | let mut file = BufWriter::new(File::create(path).await?); |
| 174 | file.write_all(STREAM_FILE_HEADER).await?; |
| 175 | let mut buffer: Vec<u8> = Vec::new(); |
| 176 | { |
| 177 | let streamer = streamer_ref.get(); |
| 178 | let header = StreamFileHeader { |
| 179 | server_uid: Cow::Borrowed(&streamer.server_uid), |
| 180 | worker_id: streamer.worker_id, |
| 181 | }; |
| 182 | StreamSerializationConfig::config().serialize_into(&mut buffer, &header)?; |
| 183 | }; |
| 184 | file.write_all(&buffer).await?; |
| 185 | while let Some(message) = receiver.recv().await { |
| 186 | match message { |
| 187 | StreamerMessage::Write { header, data } => { |
| 188 | log::debug!("Waiting data chunk into stream file"); |
| 189 | buffer.clear(); |
| 190 | StreamSerializationConfig::config().serialize_into(&mut buffer, &header)?; |
| 191 | file.write_all(&buffer).await?; |
| 192 | if !data.is_empty() { |
| 193 | file.write_all(&data).await? |
| 194 | } |
| 195 | } |
| 196 | StreamerMessage::Flush(callback) => { |
| 197 | log::debug!("Waiting for flush of stream file"); |
| 198 | file.flush().await?; |
| 199 | if callback.send(()).is_err() { |
| 200 | log::debug!("Flush callback failed") |
| 201 | } |
| 202 | } |
| 203 | } |
| 204 | } |
| 205 | log::debug!("Stream receiver is closed"); |
| 206 | Ok(()) |
| 207 | } |