1use std::cell::RefCell;
13use std::collections::BTreeMap;
14use std::rc::Rc;
15use std::time::Duration;
16
17use columnar::{Columnar, Index};
18use columnation::{Columnation, CopyRegion};
19use mz_compute_client::logging::LoggingConfig;
20use mz_ore::cast::CastFrom;
21use mz_repr::{Datum, Diff, Timestamp};
22use mz_timely_util::columnar::batcher;
23use mz_timely_util::columnar::builder::ColumnBuilder;
24use mz_timely_util::columnar::{Col2ValBatcher, columnar_exchange};
25use mz_timely_util::columnation::ColumnationChunker;
26use mz_timely_util::replay::MzReplay;
27use timely::dataflow::Scope;
28use timely::dataflow::channels::pact::{ExchangeCore, Pipeline};
29use timely::dataflow::operators::Operator;
30use timely::dataflow::operators::generic::OutputBuilder;
31use timely::dataflow::operators::generic::builder_rc::OperatorBuilder;
32use timely::dataflow::operators::generic::operator::empty;
33use timely::logging::{
34 ChannelsEvent, MessagesEvent, OperatesEvent, ParkEvent, ScheduleEvent, ShutdownEvent,
35 TimelyEvent,
36};
37use tracing::error;
38
39use crate::extensions::arrange::MzArrangeCore;
40use crate::logging::compute::{ComputeEvent, DataflowShutdown};
41use crate::logging::{
42 EventQueue, LogVariant, OutputSessionColumnar, OutputSessionVec, PermutedRowPacker, TimelyLog,
43 Update,
44};
45use crate::logging::{LogCollection, SharedLoggingState, consolidate_and_pack};
46use crate::typedefs::{KeyBatcher, KeyValBatcher, RowRowSpine};
47use mz_row_spine::RowRowBuilder;
48
49pub(super) struct Return {
51 pub collections: BTreeMap<LogVariant, LogCollection>,
53}
54
55pub(super) fn construct(
62 scope: Scope<'_, Timestamp>,
63 config: &LoggingConfig,
64 event_queue: EventQueue<Vec<(Duration, TimelyEvent)>>,
65 shared_state: Rc<RefCell<SharedLoggingState>>,
66) -> Return {
67 scope.scoped("timely logging", move |scope| {
68 let enable_logging = config.enable_logging;
69 let (logs, token) = if enable_logging {
70 event_queue.links.mz_replay(
71 scope,
72 "timely logs",
73 config.interval,
74 event_queue.activator,
75 )
76 } else {
77 let token: Rc<dyn std::any::Any> = Rc::new(Box::new(()));
78 (empty(scope), token)
79 };
80
81 let mut demux = OperatorBuilder::new("Timely Logging Demux".to_string(), scope.clone());
84 let mut input = demux.new_input(logs, Pipeline);
85 let (operates_out, operates) = demux.new_output();
86 let mut operates_out = OutputBuilder::from(operates_out);
87 let (channels_out, channels) = demux.new_output();
88 let mut channels_out = OutputBuilder::from(channels_out);
89 let (addresses_out, addresses) = demux.new_output();
90 let mut addresses_out = OutputBuilder::from(addresses_out);
91 let (parks_out, parks) = demux.new_output();
92 let mut parks_out = OutputBuilder::from(parks_out);
93 let (messages_sent_out, messages_sent) = demux.new_output();
94 let mut messages_sent_out = OutputBuilder::from(messages_sent_out);
95 let (messages_received_out, messages_received) = demux.new_output();
96 let mut messages_received_out = OutputBuilder::from(messages_received_out);
97 let (schedules_duration_out, schedules_duration) = demux.new_output();
98 let mut schedules_duration_out = OutputBuilder::from(schedules_duration_out);
99 let (schedules_histogram_out, schedules_histogram) = demux.new_output();
100 let mut schedules_histogram_out = OutputBuilder::from(schedules_histogram_out);
101 let (batches_sent_out, batches_sent) = demux.new_output();
102 let mut batches_sent_out = OutputBuilder::from(batches_sent_out);
103 let (batches_received_out, batches_received) = demux.new_output();
104 let mut batches_received_out = OutputBuilder::from(batches_received_out);
105
106 let worker_id = scope.index();
107 let mut demux_state = DemuxState::default();
108 demux.build(|_capability| {
109 let peers = scope.peers();
110 let logging_interval_ms = std::cmp::max(1, config.interval.as_millis());
111 move |_frontiers| {
112 let mut operates = operates_out.activate();
113 let mut channels = channels_out.activate();
114 let mut addresses = addresses_out.activate();
115 let mut parks = parks_out.activate();
116 let mut messages_sent = messages_sent_out.activate();
117 let mut messages_received = messages_received_out.activate();
118 let mut batches_sent = batches_sent_out.activate();
119 let mut batches_received = batches_received_out.activate();
120 let mut schedules_duration = schedules_duration_out.activate();
121 let mut schedules_histogram = schedules_histogram_out.activate();
122
123 input.for_each(|cap, data| {
124 let mut output_buffers = DemuxOutput {
125 operates: operates.session_with_builder(&cap),
126 channels: channels.session_with_builder(&cap),
127 addresses: addresses.session_with_builder(&cap),
128 parks: parks.session_with_builder(&cap),
129 messages_sent: messages_sent.session_with_builder(&cap),
130 messages_received: messages_received.session_with_builder(&cap),
131 schedules_duration: schedules_duration.session_with_builder(&cap),
132 schedules_histogram: schedules_histogram.session_with_builder(&cap),
133 batches_sent: batches_sent.session_with_builder(&cap),
134 batches_received: batches_received.session_with_builder(&cap),
135 };
136
137 for (time, event) in data.drain(..) {
138 if let TimelyEvent::Messages(msg) = &event {
139 match msg.is_send {
140 true => assert_eq!(msg.source, worker_id),
141 false => assert_eq!(msg.target, worker_id),
142 }
143 }
144
145 DemuxHandler {
146 state: &mut demux_state,
147 shared_state: &mut shared_state.borrow_mut(),
148 output: &mut output_buffers,
149 logging_interval_ms,
150 peers,
151 time,
152 }
153 .handle(event);
154 }
155 });
156 }
157 });
158
159 let operates = consolidate_and_pack::<
164 ColumnationChunker<_>,
165 KeyValBatcher<_, _, _, _>,
166 ColumnBuilder<_>,
167 _,
168 _,
169 _,
170 >(
171 operates,
172 TimelyLog::Operates,
173 move |data, packer, session| {
174 for ((id, name), time, diff) in data.iter() {
175 let data = packer.pack_slice(&[
176 Datum::UInt64(u64::cast_from(*id)),
177 Datum::UInt64(u64::cast_from(worker_id)),
178 Datum::String(name),
179 ]);
180 session.give((data, time, diff));
181 }
182 },
183 );
184
185 let channels = channels.unary::<ColumnBuilder<_>, _, _, _>(
188 Pipeline,
189 "ToRow Channels",
190 |_cap, _info| {
191 let mut packer = PermutedRowPacker::new(TimelyLog::Channels);
192 move |input, output| {
193 input.for_each_time(|time, data| {
194 let mut session = output.session_with_builder(&time);
195 for d in data.flat_map(|c| c.borrow().into_index_iter()) {
196 let ((datum, ()), time, diff) = d;
197 let (source_node, source_port) = datum.source;
198 let (target_node, target_port) = datum.target;
199 let data = packer.pack_slice(&[
200 Datum::UInt64(u64::cast_from(datum.id)),
201 Datum::UInt64(u64::cast_from(worker_id)),
202 Datum::UInt64(u64::cast_from(source_node)),
203 Datum::UInt64(u64::cast_from(source_port)),
204 Datum::UInt64(u64::cast_from(target_node)),
205 Datum::UInt64(u64::cast_from(target_port)),
206 Datum::String(
207 std::str::from_utf8(datum.typ).expect("valid string"),
208 ),
209 ]);
210 session.give((data, time, diff));
211 }
212 });
213 }
214 },
215 );
216
217 type KVB<K, V, T, D> = KeyValBatcher<K, V, T, D>;
219 type KB<K, T, D> = KeyBatcher<K, T, D>;
220
221 let addresses = consolidate_and_pack::<
222 ColumnationChunker<_>,
223 KVB<_, _, _, _>,
224 ColumnBuilder<_>,
225 _,
226 _,
227 _,
228 >(
229 addresses,
230 TimelyLog::Addresses,
231 move |data, packer, session| {
232 for ((id, address), time, diff) in data.iter() {
233 let data = packer.pack_by_index(|packer, index| match index {
234 0 => packer.push(Datum::UInt64(u64::cast_from(*id))),
235 1 => packer.push(Datum::UInt64(u64::cast_from(worker_id))),
236 2 => {
237 let list = address.iter().map(|i| Datum::UInt64(u64::cast_from(*i)));
238 packer.push_list(list)
239 }
240 _ => unreachable!("Addresses relation has three columns"),
241 });
242 session.give((data, time, diff));
243 }
244 },
245 );
246
247 let parks =
248 consolidate_and_pack::<ColumnationChunker<_>, KB<_, _, _>, ColumnBuilder<_>, _, _, _>(
249 parks,
250 TimelyLog::Parks,
251 move |data, packer, session| {
252 for ((datum, ()), time, diff) in data.iter() {
253 let data = packer.pack_slice(&[
254 Datum::UInt64(u64::cast_from(worker_id)),
255 Datum::UInt64(datum.duration_pow),
256 datum
257 .requested_pow
258 .map(Datum::UInt64)
259 .unwrap_or(Datum::Null),
260 ]);
261 session.give((data, time, diff));
262 }
263 },
264 );
265
266 let batches_sent =
267 consolidate_and_pack::<ColumnationChunker<_>, KB<_, _, _>, ColumnBuilder<_>, _, _, _>(
268 batches_sent,
269 TimelyLog::BatchesSent,
270 move |data, packer, session| {
271 for ((datum, ()), time, diff) in data.iter() {
272 let data = packer.pack_slice(&[
273 Datum::UInt64(u64::cast_from(datum.channel)),
274 Datum::UInt64(u64::cast_from(worker_id)),
275 Datum::UInt64(u64::cast_from(datum.worker)),
276 ]);
277 session.give((data, time, diff));
278 }
279 },
280 );
281
282 let batches_received =
283 consolidate_and_pack::<ColumnationChunker<_>, KB<_, _, _>, ColumnBuilder<_>, _, _, _>(
284 batches_received,
285 TimelyLog::BatchesReceived,
286 move |data, packer, session| {
287 for ((datum, ()), time, diff) in data.iter() {
288 let data = packer.pack_slice(&[
289 Datum::UInt64(u64::cast_from(datum.channel)),
290 Datum::UInt64(u64::cast_from(datum.worker)),
291 Datum::UInt64(u64::cast_from(worker_id)),
292 ]);
293 session.give((data, time, diff));
294 }
295 },
296 );
297
298 let messages_sent =
299 consolidate_and_pack::<ColumnationChunker<_>, KB<_, _, _>, ColumnBuilder<_>, _, _, _>(
300 messages_sent,
301 TimelyLog::MessagesSent,
302 move |data, packer, session| {
303 for ((datum, ()), time, diff) in data.iter() {
304 let data = packer.pack_slice(&[
305 Datum::UInt64(u64::cast_from(datum.channel)),
306 Datum::UInt64(u64::cast_from(worker_id)),
307 Datum::UInt64(u64::cast_from(datum.worker)),
308 ]);
309 session.give((data, time, diff));
310 }
311 },
312 );
313
314 let messages_received =
315 consolidate_and_pack::<ColumnationChunker<_>, KB<_, _, _>, ColumnBuilder<_>, _, _, _>(
316 messages_received,
317 TimelyLog::MessagesReceived,
318 move |data, packer, session| {
319 for ((datum, ()), time, diff) in data.iter() {
320 let data = packer.pack_slice(&[
321 Datum::UInt64(u64::cast_from(datum.channel)),
322 Datum::UInt64(u64::cast_from(datum.worker)),
323 Datum::UInt64(u64::cast_from(worker_id)),
324 ]);
325 session.give((data, time, diff));
326 }
327 },
328 );
329
330 let elapsed =
331 consolidate_and_pack::<ColumnationChunker<_>, KB<_, _, _>, ColumnBuilder<_>, _, _, _>(
332 schedules_duration,
333 TimelyLog::Elapsed,
334 move |data, packer, session| {
335 for ((operator, ()), time, diff) in data.iter() {
336 let data = packer.pack_slice(&[
337 Datum::UInt64(u64::cast_from(*operator)),
338 Datum::UInt64(u64::cast_from(worker_id)),
339 ]);
340 session.give((data, time, diff));
341 }
342 },
343 );
344
345 let histogram =
346 consolidate_and_pack::<ColumnationChunker<_>, KB<_, _, _>, ColumnBuilder<_>, _, _, _>(
347 schedules_histogram,
348 TimelyLog::Histogram,
349 move |data, packer, session| {
350 for ((datum, ()), time, diff) in data.iter() {
351 let data = packer.pack_slice(&[
352 Datum::UInt64(u64::cast_from(datum.operator)),
353 Datum::UInt64(u64::cast_from(worker_id)),
354 Datum::UInt64(datum.duration_pow),
355 ]);
356 session.give((data, time, diff));
357 }
358 },
359 );
360
361 let logs = {
362 use TimelyLog::*;
363 [
364 (Operates, operates),
365 (Channels, channels),
366 (Elapsed, elapsed),
367 (Histogram, histogram),
368 (Addresses, addresses),
369 (Parks, parks),
370 (MessagesSent, messages_sent),
371 (MessagesReceived, messages_received),
372 (BatchesSent, batches_sent),
373 (BatchesReceived, batches_received),
374 ]
375 };
376
377 let mut collections = BTreeMap::new();
379 for (variant, collection) in logs {
380 let variant = LogVariant::Timely(variant);
381 if config.index_logs.contains_key(&variant) {
382 type Batcher<K, V, T, R> = Col2ValBatcher<K, V, T, R>;
384 type Builder<T, R> = RowRowBuilder<T, R>;
385 let trace = collection
386 .mz_arrange_core::<
387 _,
388 batcher::Chunker<_>,
389 Batcher<_, _, _, _>,
390 Builder<_, _>,
391 RowRowSpine<_, _>,
392 >(
393 ExchangeCore::<ColumnBuilder<_>, _>::new_core(
394 columnar_exchange::<mz_repr::Row, mz_repr::Row, Timestamp, Diff>,
395 ),
396 &format!("Arrange {variant:?}"),
397 )
398 .trace;
399 let collection = LogCollection {
400 trace,
401 token: Rc::clone(&token),
402 };
403 collections.insert(variant, collection);
404 }
405 }
406
407 Return { collections }
408 })
409}
410
411#[derive(Default)]
413struct DemuxState {
414 operators: BTreeMap<usize, OperatesEvent>,
416 dataflow_channels: BTreeMap<usize, Vec<ChannelsEvent>>,
418 last_park: Option<Park>,
420 messages_sent: BTreeMap<usize, Box<[MessageCount]>>,
422 messages_received: BTreeMap<usize, Box<[MessageCount]>>,
424 schedule_starts: Vec<(usize, Duration)>,
426 schedules_data: BTreeMap<usize, Vec<(isize, Diff)>>,
429}
430
431struct Park {
432 time: Duration,
434 requested: Option<Duration>,
436}
437
438#[derive(Default, Copy, Clone, Debug)]
440struct MessageCount {
441 batches: i64,
443 records: Diff,
445}
446
447struct DemuxOutput<'a, 'b> {
453 operates: OutputSessionVec<'a, 'b, Update<(usize, String)>>,
454 channels: OutputSessionColumnar<'a, 'b, Update<(ChannelDatum, ())>>,
455 addresses: OutputSessionVec<'a, 'b, Update<(usize, Vec<usize>)>>,
456 parks: OutputSessionVec<'a, 'b, Update<(ParkDatum, ())>>,
457 batches_sent: OutputSessionVec<'a, 'b, Update<(MessageDatum, ())>>,
458 batches_received: OutputSessionVec<'a, 'b, Update<(MessageDatum, ())>>,
459 messages_sent: OutputSessionVec<'a, 'b, Update<(MessageDatum, ())>>,
460 messages_received: OutputSessionVec<'a, 'b, Update<(MessageDatum, ())>>,
461 schedules_duration: OutputSessionVec<'a, 'b, Update<(usize, ())>>,
462 schedules_histogram: OutputSessionVec<'a, 'b, Update<(ScheduleHistogramDatum, ())>>,
463}
464
465#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Columnar)]
466struct ChannelDatum {
467 id: usize,
468 source: (usize, usize),
469 target: (usize, usize),
470 typ: String,
471}
472
473#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
474struct ParkDatum {
475 duration_pow: u64,
476 requested_pow: Option<u64>,
477}
478
479impl Columnation for ParkDatum {
480 type InnerRegion = CopyRegion<Self>;
481}
482
483#[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
484struct MessageDatum {
485 channel: usize,
486 worker: usize,
487}
488
489impl Columnation for MessageDatum {
490 type InnerRegion = CopyRegion<Self>;
491}
492
493#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
494struct ScheduleHistogramDatum {
495 operator: usize,
496 duration_pow: u64,
497}
498
499impl Columnation for ScheduleHistogramDatum {
500 type InnerRegion = CopyRegion<Self>;
501}
502
503struct DemuxHandler<'a, 'b, 'c> {
505 state: &'a mut DemuxState,
507 shared_state: &'a mut SharedLoggingState,
509 output: &'a mut DemuxOutput<'b, 'c>,
511 logging_interval_ms: u128,
513 peers: usize,
515 time: Duration,
517}
518
519impl DemuxHandler<'_, '_, '_> {
520 fn ts(&self) -> Timestamp {
523 let time_ms = self.time.as_millis();
524 let interval = self.logging_interval_ms;
525 let rounded = (time_ms / interval + 1) * interval;
526 rounded.try_into().expect("must fit")
527 }
528
529 fn handle(&mut self, event: TimelyEvent) {
531 use TimelyEvent::*;
532
533 match event {
534 Operates(e) => self.handle_operates(e),
535 Channels(e) => self.handle_channels(e),
536 Shutdown(e) => self.handle_shutdown(e),
537 Park(e) => self.handle_park(e),
538 Messages(e) => self.handle_messages(e),
539 Schedule(e) => self.handle_schedule(e),
540 _ => (),
541 }
542 }
543
544 fn handle_operates(&mut self, event: OperatesEvent) {
545 let ts = self.ts();
546 let datum = (event.id, event.name.clone());
547 self.output.operates.give((datum, ts, Diff::ONE));
548
549 let datum = (event.id, event.addr.clone());
550 self.output.addresses.give((datum, ts, Diff::ONE));
551
552 self.state.operators.insert(event.id, event);
553 }
554
555 fn handle_channels(&mut self, event: ChannelsEvent) {
556 let ts = self.ts();
557 let datum = ChannelDatumReference {
558 id: event.id,
559 source: event.source,
560 target: event.target,
561 typ: &event.typ,
562 };
563 self.output.channels.give(((datum, ()), ts, Diff::ONE));
564
565 let datum = (event.id, event.scope_addr.clone());
566 self.output.addresses.give((datum, ts, Diff::ONE));
567
568 let dataflow_index = event.scope_addr[0];
569 self.state
570 .dataflow_channels
571 .entry(dataflow_index)
572 .or_default()
573 .push(event);
574 }
575
576 fn handle_shutdown(&mut self, event: ShutdownEvent) {
577 let Some(operator) = self.state.operators.remove(&event.id) else {
583 error!(operator_id = ?event.id, "missing operator entry at time of shutdown");
584 return;
585 };
586
587 let ts = self.ts();
589 let datum = (operator.id, operator.name);
590 self.output.operates.give((datum, ts, Diff::MINUS_ONE));
591
592 if let Some(schedules) = self.state.schedules_data.remove(&event.id) {
594 for (bucket, (count, elapsed_ns)) in IntoIterator::into_iter(schedules)
595 .enumerate()
596 .filter(|(_, (count, _))| *count != 0)
597 {
598 self.output
599 .schedules_duration
600 .give(((event.id, ()), ts, Diff::from(-elapsed_ns)));
601
602 let datum = ScheduleHistogramDatum {
603 operator: event.id,
604 duration_pow: 1 << bucket,
605 };
606 let diff = Diff::cast_from(-count);
607 self.output
608 .schedules_histogram
609 .give(((datum, ()), ts, diff));
610 }
611 }
612
613 if operator.addr.len() == 1 {
614 let dataflow_index = operator.addr[0];
615 self.handle_dataflow_shutdown(dataflow_index);
616 }
617
618 let datum = (operator.id, operator.addr);
619 self.output.addresses.give((datum, ts, Diff::MINUS_ONE));
620 }
621
622 fn handle_dataflow_shutdown(&mut self, dataflow_index: usize) {
623 self.shared_state.compute_logger.as_ref().map(|logger| {
625 logger.log(&(ComputeEvent::DataflowShutdown(DataflowShutdown { dataflow_index })))
626 });
627
628 let Some(channels) = self.state.dataflow_channels.remove(&dataflow_index) else {
630 return;
631 };
632
633 let ts = self.ts();
634 for channel in channels {
635 let datum = ChannelDatumReference {
637 id: channel.id,
638 source: channel.source,
639 target: channel.target,
640 typ: &channel.typ,
641 };
642 self.output
643 .channels
644 .give(((datum, ()), ts, Diff::MINUS_ONE));
645
646 let datum = (channel.id, channel.scope_addr);
647 self.output.addresses.give((datum, ts, Diff::MINUS_ONE));
648
649 if let Some(sent) = self.state.messages_sent.remove(&channel.id) {
651 for (target_worker, count) in sent.iter().enumerate() {
652 let datum = MessageDatum {
653 channel: channel.id,
654 worker: target_worker,
655 };
656 self.output
657 .messages_sent
658 .give(((datum, ()), ts, Diff::from(-count.records)));
659 self.output
660 .batches_sent
661 .give(((datum, ()), ts, Diff::from(-count.batches)));
662 }
663 }
664 if let Some(received) = self.state.messages_received.remove(&channel.id) {
665 for (source_worker, count) in received.iter().enumerate() {
666 let datum = MessageDatum {
667 channel: channel.id,
668 worker: source_worker,
669 };
670 self.output.messages_received.give((
671 (datum, ()),
672 ts,
673 Diff::from(-count.records),
674 ));
675 self.output.batches_received.give((
676 (datum, ()),
677 ts,
678 Diff::from(-count.batches),
679 ));
680 }
681 }
682 }
683 }
684
685 fn handle_park(&mut self, event: ParkEvent) {
686 match event {
687 ParkEvent::Park(requested) => {
688 let park = Park {
689 time: self.time,
690 requested,
691 };
692 let existing = self.state.last_park.replace(park);
693 if existing.is_some() {
694 error!("park without a succeeding unpark");
695 }
696 }
697 ParkEvent::Unpark => {
698 let Some(park) = self.state.last_park.take() else {
699 error!("unpark without a preceding park");
700 return;
701 };
702
703 let duration_ns = self.time.saturating_sub(park.time).as_nanos();
704 let duration_pow =
705 u64::try_from(duration_ns.next_power_of_two()).expect("must fit");
706 let requested_pow = park
707 .requested
708 .map(|r| u64::try_from(r.as_nanos().next_power_of_two()).expect("must fit"));
709
710 let ts = self.ts();
711 let datum = ParkDatum {
712 duration_pow,
713 requested_pow,
714 };
715 self.output.parks.give(((datum, ()), ts, Diff::ONE));
716 }
717 }
718 }
719
720 fn handle_messages(&mut self, event: MessagesEvent) {
721 let ts = self.ts();
722 let count = Diff::from(event.record_count);
723
724 if event.is_send {
725 let datum = MessageDatum {
726 channel: event.channel,
727 worker: event.target,
728 };
729 self.output.messages_sent.give(((datum, ()), ts, count));
730 self.output.batches_sent.give(((datum, ()), ts, Diff::ONE));
731
732 let sent_counts = self
733 .state
734 .messages_sent
735 .entry(event.channel)
736 .or_insert_with(|| vec![Default::default(); self.peers].into_boxed_slice());
737 sent_counts[event.target].records += count;
738 sent_counts[event.target].batches += 1;
739 } else {
740 let datum = MessageDatum {
741 channel: event.channel,
742 worker: event.source,
743 };
744 self.output.messages_received.give(((datum, ()), ts, count));
745 self.output
746 .batches_received
747 .give(((datum, ()), ts, Diff::ONE));
748
749 let received_counts = self
750 .state
751 .messages_received
752 .entry(event.channel)
753 .or_insert_with(|| vec![Default::default(); self.peers].into_boxed_slice());
754 received_counts[event.source].records += count;
755 received_counts[event.source].batches += 1;
756 }
757 }
758
759 fn handle_schedule(&mut self, event: ScheduleEvent) {
760 match event.start_stop {
761 timely::logging::StartStop::Start => {
762 self.state.schedule_starts.push((event.id, self.time));
763 }
764 timely::logging::StartStop::Stop => {
765 let Some((old_id, start_time)) = self.state.schedule_starts.pop() else {
766 error!(operator_id = ?event.id, "schedule stop without preceding start");
767 return;
768 };
769
770 if old_id != event.id {
771 error!(start_id = ?old_id, stop_id = ?event.id, "schedule stop without preceding start");
772 return;
773 }
774
775 let elapsed_ns = self.time.saturating_sub(start_time).as_nanos();
776 let elapsed_i64 = i64::try_from(elapsed_ns).expect("must fit");
777 let elapsed_diff = Diff::from(elapsed_i64);
778 let elapsed_pow = u64::try_from(elapsed_ns.next_power_of_two()).expect("must fit");
779
780 let ts = self.ts();
781 let datum = event.id;
782 self.output
783 .schedules_duration
784 .give(((datum, ()), ts, elapsed_diff));
785
786 let datum = ScheduleHistogramDatum {
787 operator: event.id,
788 duration_pow: elapsed_pow,
789 };
790 self.output
791 .schedules_histogram
792 .give(((datum, ()), ts, Diff::ONE));
793
794 let index = usize::cast_from(elapsed_pow.trailing_zeros());
796 let data = self.state.schedules_data.entry(event.id).or_default();
797 grow_vec(data, index);
798 let (count, duration) = &mut data[index];
799 *count += 1;
800 *duration += elapsed_diff;
801 }
802 }
803 }
804}
805
806fn grow_vec<T>(vec: &mut Vec<T>, index: usize)
810where
811 T: Clone + Default,
812{
813 if vec.len() <= index {
814 vec.resize(index + 1, Default::default());
815 }
816}