Skip to main content

mz_compute/logging/
timely.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Logging dataflows for events generated by timely dataflow.
11
12use 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
49/// The return type of [`construct`].
50pub(super) struct Return {
51    /// Collections to export.
52    pub collections: BTreeMap<LogVariant, LogCollection>,
53}
54
55/// Constructs the logging dataflow fragment for timely logs.
56///
57/// Params
58/// * `scope`: The Timely scope hosting the log analysis dataflow.
59/// * `config`: Logging configuration
60/// * `event_queue`: The source to read log events from.
61pub(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        // Build a demux operator that splits the replayed event stream up into the separate
82        // logging streams.
83        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        // Encode the contents of each logging stream into its expected `Row` format.
160        // We pre-arrange the logging streams to force a consolidation and reduce the amount of
161        // updates that reach `Row` encoding.
162
163        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        // TODO: `consolidate_and_pack` requires columnation, which `ChannelDatum` does not
186        // implement. Consider consolidating here once we support columnar.
187        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        // Types to make rustfmt happy.
218        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        // Build the output arrangements.
378        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                // Extract types to make rustfmt happy.
383                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/// State maintained by the demux operator.
412#[derive(Default)]
413struct DemuxState {
414    /// Information about live operators, indexed by operator ID.
415    operators: BTreeMap<usize, OperatesEvent>,
416    /// Maps dataflow IDs to channels in the dataflow.
417    dataflow_channels: BTreeMap<usize, Vec<ChannelsEvent>>,
418    /// Information about the last requested park.
419    last_park: Option<Park>,
420    /// Maps channel IDs to boxed slices counting the messages sent to each target worker.
421    messages_sent: BTreeMap<usize, Box<[MessageCount]>>,
422    /// Maps channel IDs to boxed slices counting the messages received from each source worker.
423    messages_received: BTreeMap<usize, Box<[MessageCount]>>,
424    /// Stores for scheduled operators the time when they were scheduled.
425    schedule_starts: Vec<(usize, Duration)>,
426    /// Maps operator IDs to a vector recording the (count, elapsed_ns) values in each histogram
427    /// bucket.
428    schedules_data: BTreeMap<usize, Vec<(isize, Diff)>>,
429}
430
431struct Park {
432    /// Time when the park occurred.
433    time: Duration,
434    /// Requested park time.
435    requested: Option<Duration>,
436}
437
438/// Organize message counts into number of batches and records.
439#[derive(Default, Copy, Clone, Debug)]
440struct MessageCount {
441    /// The number of batches sent across a channel.
442    batches: i64,
443    /// The number of records sent across a channel.
444    records: Diff,
445}
446
447/// Bundled output buffers used by the demux operator.
448//
449// We use tuples rather than dedicated `*Datum` structs for `operates` and `addresses` to avoid
450// having to manually implement `Columnation`. If `Columnation` could be `#[derive]`ed, that
451// wouldn't be an issue.
452struct 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
503/// Event handler of the demux operator.
504struct DemuxHandler<'a, 'b, 'c> {
505    /// State kept by the demux operator.
506    state: &'a mut DemuxState,
507    /// State shared across log receivers.
508    shared_state: &'a mut SharedLoggingState,
509    /// Demux output buffers.
510    output: &'a mut DemuxOutput<'b, 'c>,
511    /// The logging interval specifying the time granularity for the updates.
512    logging_interval_ms: u128,
513    /// The number of timely workers.
514    peers: usize,
515    /// The current event time.
516    time: Duration,
517}
518
519impl DemuxHandler<'_, '_, '_> {
520    /// Return the timestamp associated with the current event, based on the event time and the
521    /// logging interval.
522    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    /// Handle the given timely event.
530    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        // Dropped operators should result in a negative record for
578        // the `operates` collection, cancelling out the initial
579        // operator announcement.
580        // Remove logging for this operator.
581
582        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        // Retract operator information.
588        let ts = self.ts();
589        let datum = (operator.id, operator.name);
590        self.output.operates.give((datum, ts, Diff::MINUS_ONE));
591
592        // Retract schedules information for the operator
593        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        // Notify compute logging about the shutdown.
624        self.shared_state.compute_logger.as_ref().map(|logger| {
625            logger.log(&(ComputeEvent::DataflowShutdown(DataflowShutdown { dataflow_index })))
626        });
627
628        // When a dataflow shuts down, we need to retract all its channels.
629        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            // Retract channel description.
636            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            // Retract messages logged for this channel.
650            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                // Record count and elapsed time for later retraction.
795                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
806/// Grow the given vector so it fits the given index.
807///
808/// This does nothing if the vector is already large enough.
809fn 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}