MCPcopy Create free account
hub / github.com/It4innovations/hyperqueue / get_stream

Method get_stream

crates/hyperqueue/src/worker/streamer.rs:54–98  ·  view source on GitHub ↗

Get tasm stream input end, if a connection to a stream server is not established, the new one is created */

(
        &mut self,
        streamer_ref: &StreamerRef,
        stream_path: &Path,
        task_id: TaskId,
        instance_id: InstanceId,
    )

Source from the content-addressed store, hash-verified

52 the new one is created
53 */
54 pub fn get_stream(
55 &mut self,
56 streamer_ref: &StreamerRef,
57 stream_path: &Path,
58 task_id: TaskId,
59 instance_id: InstanceId,
60 ) -> crate::Result<StreamSender> {
61 log::debug!("New stream for {task_id} ({})", stream_path.display());
62 let sender = if let Some(ref mut info) = self.streams.get_mut(stream_path) {
63 info.sender.clone()
64 } else {
65 log::debug!(
66 "Starting a new stream instance, stream_path = {}",
67 stream_path.display()
68 );
69 if !stream_path.is_dir() {
70 std::fs::create_dir_all(stream_path).map_err(|error| {
71 HqError::GenericError(format!(
72 "Cannot create stream directory {}: {error:?}",
73 stream_path.display()
74 ))
75 })?;
76 }
77 let (queue_sender, queue_receiver) = channel(STREAMER_BUFFER_SIZE);
78 let stream = StreamDescriptor {
79 sender: queue_sender.clone(),
80 };
81 self.streams.insert(stream_path.to_path_buf(), stream);
82 let stream_path = stream_path.to_path_buf();
83 let streamer_ref = streamer_ref.clone();
84 spawn_local(async move {
85 if let Err(e) = stream_writer(&streamer_ref, &stream_path, queue_receiver).await {
86 log::error!("Stream failed: {e}")
87 }
88 let mut streamer = streamer_ref.get_mut();
89 assert!(streamer.streams.remove(&stream_path).is_some());
90 });
91 queue_sender
92 };
93 Ok(StreamSender {
94 task_id,
95 instance_id,
96 sender,
97 })
98 }
99}
100
101define_wrapped_type!(StreamerRef, Streamer, pub);

Callers 1

create_task_futureFunction · 0.80

Calls 4

stream_writerFunction · 0.85
get_mutMethod · 0.45
cloneMethod · 0.45
insertMethod · 0.45

Tested by

no test coverage detected