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

Function start_streaming

crates/hyperqueue/src/server/client/mod.rs:128–190  ·  view source on GitHub ↗
(
    mut tx: Tx,
    mut rx: Rx,
    state_ref: StateRef,
    senders: &Senders,
    stream_events: StreamEvents,
    extra_message: Option<ToClientMessage>,
)

Source from the content-addressed store, hash-verified

126}
127
128async fn start_streaming<
129 Tx: Sink<ToClientMessage, Error = tako::Error> + Unpin + 'static,
130 Rx: Stream<Item = tako::Result<FromClientMessage>> + Unpin,
131>(
132 mut tx: Tx,
133 mut rx: Rx,
134 state_ref: StateRef,
135 senders: &Senders,
136 stream_events: StreamEvents,
137 extra_message: Option<ToClientMessage>,
138) where
139 Tx::Error: Debug,
140{
141 let StreamEvents {
142 mode,
143 enable_worker_overviews,
144 filter,
145 } = stream_events;
146 if enable_worker_overviews {
147 senders.server_control.add_worker_overview_listener();
148 }
149 log::debug!("Start streaming events to client");
150
151 /* We create two event queues, one for historic events and one for live events
152 So while historic events are loaded from the file and streamed, live events are already
153 collected and sent immediately once the historic events are sent */
154 let live = if mode.is_live_events_enabled() {
155 let (tx2, rx2) = mpsc::unbounded_channel::<Event>();
156 let listener_id = senders.events.register_listener(filter, tx2);
157 Some((rx2, listener_id))
158 } else {
159 None
160 };
161
162 if let Some(msg) = extra_message {
163 let _ = tx.send(msg).await;
164 }
165
166 // If we use a journal, we can replay historical events from it.
167 // If not, we can at least try to reconstruct a few basic events
168 // based on the current state.
169 if mode.is_past_events_enabled() {
170 let (tx1, rx1) = mpsc::unbounded_channel::<Event>();
171 if senders.events.is_journal_enabled() {
172 senders.events.start_journal_replay(tx1);
173 } else if let Err(e) = reconstruct_historical_events(state_ref, tx1) {
174 log::error!("Cannot reconstruct historical state: {e:?}");
175 }
176 stream_history_events(&mut tx, rx1).await;
177 }
178
179 if let Some((rx2, listener_id)) = live {
180 if mode.is_past_events_enabled() {
181 let _ = tx.send(ToClientMessage::EventLiveBoundary).await;
182 }
183 crate::server::client::stream_events(&mut tx, &mut rx, rx2).await;
184 senders.events.unregister_listener(listener_id);
185 }

Callers 1

client_rpc_loopFunction · 0.85

Calls 12

stream_history_eventsFunction · 0.85
stream_eventsFunction · 0.85
register_listenerMethod · 0.80
is_journal_enabledMethod · 0.80
start_journal_replayMethod · 0.80
unregister_listenerMethod · 0.80
sendMethod · 0.45

Tested by

no test coverage detected