MCPcopy Create free account
hub / github.com/MegEngine/MegFlow / start

Method start

flow-plugins/src/video_input.rs:55–99  ·  view source on GitHub ↗
(
        mut self: Box<Self>,
        _: Context,
        resources: ResourceCollection,
    )

Source from the content-addressed store, hash-verified

53
54impl 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}

Callers 2

newMethod · 0.45
__init__Method · 0.45

Calls 7

decode_videoFunction · 0.85
createMethod · 0.80
pushMethod · 0.80
recv_anyMethod · 0.80
cloneMethod · 0.45
fetchMethod · 0.45
closeMethod · 0.45

Tested by

no test coverage detected