| 53 | |
| 54 | impl Actor for VideoInput { |
| 55 | fn start( |
| 56 | mut self: Box<Self>, |
| 57 | _: Context, |
| 58 | resources: ResourceCollection, |
| 59 | ) -> rt::task::JoinHandle<Result<()>> { |
| 60 | rt::task::spawn(async move { |
| 61 | let mut recv_conns = vec![]; |
| 62 | |
| 63 | let mut id = 0u64; |
| 64 | for _ in 0..self.repeat { |
| 65 | for url in self.urls.iter() { |
| 66 | if Path::new(&url).is_file() { |
| 67 | id += 1; |
| 68 | // create multiple stream |
| 69 | self.out |
| 70 | .create(id.to_owned(), resources.clone(), Default::default()) |
| 71 | .await |
| 72 | .expect("broker has closed"); |
| 73 | let (_, port) = self.out.fetch().await.expect("broker has closed"); |
| 74 | let url_cloned = url.clone(); |
| 75 | let port_cloned = port.clone(); |
| 76 | flow_rs::rt::task::spawn_blocking( |
| 77 | move || -> Result<(), ffmpeg_next::Error> { |
| 78 | if let Err(err) = codec::decode_video(id, &url_cloned, &port_cloned) |
| 79 | { |
| 80 | error!("video[{}] {} decode fault: {:?}", id, url_cloned, err); |
| 81 | Err(err) |
| 82 | } else { |
| 83 | port_cloned.close(); |
| 84 | Ok(()) |
| 85 | } |
| 86 | }, |
| 87 | ); |
| 88 | |
| 89 | // recv streams |
| 90 | let (_, port) = self.inp.fetch().await.unwrap(); |
| 91 | recv_conns.push(async move { while port.recv_any().await.is_ok() {} }); |
| 92 | } |
| 93 | } |
| 94 | } |
| 95 | |
| 96 | join_all(recv_conns).await; |
| 97 | Ok(()) |
| 98 | }) |
| 99 | } |
| 100 | } |