| 286 | |
| 287 | #[tokio::test] |
| 288 | async fn test_jsonl() -> Result<()> { |
| 289 | let testvalues = [ |
| 290 | Event::ProgressSteps { |
| 291 | task: "sometask".into(), |
| 292 | description: "somedesc".into(), |
| 293 | id: "someid".into(), |
| 294 | steps_cached: 0, |
| 295 | steps: 0, |
| 296 | steps_total: 3, |
| 297 | subtasks: Vec::new(), |
| 298 | }, |
| 299 | Event::ProgressBytes { |
| 300 | task: "sometask".into(), |
| 301 | description: "somedesc".into(), |
| 302 | id: "someid".into(), |
| 303 | bytes_cached: 0, |
| 304 | bytes: 11, |
| 305 | bytes_total: 42, |
| 306 | steps_cached: 0, |
| 307 | steps: 0, |
| 308 | steps_total: 3, |
| 309 | subtasks: Vec::new(), |
| 310 | }, |
| 311 | ]; |
| 312 | let (send, recv) = tokio::net::unix::pipe::pipe()?; |
| 313 | let testvalues_sender = testvalues.iter().cloned(); |
| 314 | let sender = async move { |
| 315 | let w = ProgressWriter::try_from(send)?; |
| 316 | for value in testvalues_sender { |
| 317 | w.send(value).await; |
| 318 | } |
| 319 | anyhow::Ok(()) |
| 320 | }; |
| 321 | let testvalues = &testvalues; |
| 322 | let receiver = async move { |
| 323 | let tf = BufReader::new(recv); |
| 324 | let mut expected = testvalues.iter(); |
| 325 | let mut lines = tf.lines(); |
| 326 | let mut got_first = false; |
| 327 | while let Some(line) = lines.next_line().await? { |
| 328 | let found: Event = serde_json::from_str(&line)?; |
| 329 | let expected_value = if !got_first { |
| 330 | got_first = true; |
| 331 | &Event::Start { |
| 332 | version: API_VERSION.into(), |
| 333 | } |
| 334 | } else { |
| 335 | expected.next().unwrap() |
| 336 | }; |
| 337 | assert_eq!(&found, expected_value); |
| 338 | } |
| 339 | anyhow::Ok(()) |
| 340 | }; |
| 341 | tokio::try_join!(sender, receiver)?; |
| 342 | Ok(()) |
| 343 | } |
| 344 | } |