(
mut tx: Tx,
mut rx: Rx,
state_ref: StateRef,
senders: &Senders,
stream_events: StreamEvents,
extra_message: Option<ToClientMessage>,
)
| 126 | } |
| 127 | |
| 128 | async 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 | } |
no test coverage detected