(
&self,
timely_worker: &mut TimelyWorker,
client_rx: mpsc::UnboundedReceiver<(
Uuid,
mpsc::UnboundedReceiver<StorageCommand>,
mpsc::Unbound
| 121 | const NAME: &str = "storage"; |
| 122 | |
| 123 | fn run_worker( |
| 124 | &self, |
| 125 | timely_worker: &mut TimelyWorker, |
| 126 | client_rx: mpsc::UnboundedReceiver<( |
| 127 | Uuid, |
| 128 | mpsc::UnboundedReceiver<StorageCommand>, |
| 129 | mpsc::UnboundedSender<StorageResponse>, |
| 130 | )>, |
| 131 | ) { |
| 132 | // Register a timely logger that forwards events to the compute logging dataflow. |
| 133 | // Assign by local worker index so storage worker x matches compute worker x. |
| 134 | let local_index = timely_worker.index() % self.workers_per_process; |
| 135 | let writer = self.timely_log_writers.lock().unwrap()[local_index].take(); |
| 136 | if let Some(writer) = writer { |
| 137 | use timely::dataflow::operators::capture::{Event, EventPusher}; |
| 138 | use timely::logging::TimelyEventBuilder; |
| 139 | |
| 140 | // We use an approach similar to compute's logging: wrap the writer in |
| 141 | // a BatchLogger that translates Logger callbacks into Event pushes, |
| 142 | // then register the Logger with timely's log_register. |
| 143 | let interval_ms = 1000u128; // 1 second batching interval |
| 144 | let mut time_ms = mz_repr::Timestamp::from(0u64); |
| 145 | let mut event_pusher = writer; |
| 146 | let now = std::time::Instant::now(); |
| 147 | let start_offset = std::time::SystemTime::now() |
| 148 | .duration_since(std::time::SystemTime::UNIX_EPOCH) |
| 149 | .expect("Failed to get duration since Unix epoch"); |
| 150 | |
| 151 | let logger = timely::logging_core::Logger::<TimelyEventBuilder>::new( |
| 152 | now, |
| 153 | start_offset, |
| 154 | move |time: &std::time::Duration, |
| 155 | data: &mut Option<Vec<(std::time::Duration, TimelyEvent)>>| { |
| 156 | if let Some(mut data) = data.take() { |
| 157 | // Filter park/unpark events and remap IDs before handing events |
| 158 | // off to compute. Compute's park tracking assumes a single |
| 159 | // timely runtime; mixing in storage's park events would break |
| 160 | // it. Remapping ensures storage operator/channel IDs don't |
| 161 | // collide with compute's. |
| 162 | data.retain_mut(|(_, event)| { |
| 163 | if matches!(event, TimelyEvent::Park(_)) { |
| 164 | return false; |
| 165 | } |
| 166 | remap_timely_event_ids(event); |
| 167 | true |
| 168 | }); |
| 169 | event_pusher.push(Event::Messages(time_ms, data)); |
| 170 | } else { |
| 171 | // Advance progress. |
| 172 | let new_time_ms: u64 = (((time.as_millis() / interval_ms) + 1) |
| 173 | * interval_ms) |
| 174 | .try_into() |
| 175 | .expect("must fit"); |
| 176 | let new_time_ms = mz_repr::Timestamp::from(new_time_ms); |
| 177 | if time_ms < new_time_ms { |
| 178 | event_pusher |
| 179 | .push(Event::Progress(vec![(new_time_ms, 1), (time_ms, -1)])); |
| 180 | time_ms = new_time_ms; |
nothing calls this directly
no test coverage detected