1pub mod compute;
13mod differential;
14pub(super) mod initialize;
15mod prometheus;
16mod reachability;
17mod resource_usage;
18mod timely;
19
20use std::any::Any;
21use std::collections::BTreeMap;
22use std::marker::PhantomData;
23use std::rc::Rc;
24use std::time::{Duration, Instant};
25
26use ::prometheus::IntCounter;
27use ::timely::container::{CapacityContainerBuilder, PushInto};
28use ::timely::dataflow::Stream;
29use ::timely::dataflow::channels::pact::Pipeline;
30use ::timely::dataflow::operators::capture::{Event, EventLink, EventPusher};
31use ::timely::dataflow::operators::generic::Session;
32use ::timely::dataflow::operators::{Capability, CapabilityTrait, InputCapability, Operator};
33use ::timely::progress::Timestamp as TimelyTimestamp;
34use ::timely::scheduling::Activator;
35use ::timely::{Container, ContainerBuilder};
36use differential_dataflow::trace::Batcher;
37use mz_compute_client::logging::{ComputeLog, DifferentialLog, LogVariant, TimelyLog};
38use mz_expr::{MirScalarExpr, permutation_for_arrangement};
39use mz_repr::{Datum, Diff, Row, RowPacker, RowRef, Timestamp};
40use mz_timely_util::activator::RcActivator;
41use mz_timely_util::columnar::builder::ColumnBuilder;
42use mz_timely_util::operator::consolidate_pact;
43
44use crate::logging::compute::Logger as ComputeLogger;
45use crate::typedefs::RowRowAgent;
46
47pub use crate::logging::initialize::initialize;
48
49pub(super) type Update<D> = (D, Timestamp, Diff);
51pub(super) type OutputSession<'a, 'b, CB> =
54 Session<'a, 'b, Timestamp, CB, InputCapability<Timestamp>>;
55pub(super) type OutputSessionVec<'a, 'b, D> =
57 OutputSession<'a, 'b, CapacityContainerBuilder<Vec<D>>>;
58pub(super) type OutputSessionColumnar<'a, 'b, D> = OutputSession<'a, 'b, ColumnBuilder<D>>;
60
61struct BatchLogger<C, P>
63where
64 P: EventPusher<Timestamp, C>,
65{
66 time_ms: Timestamp,
68 event_pusher: P,
70 interval_ms: u128,
74 records_total: IntCounter,
76 _marker: PhantomData<C>,
77}
78
79impl<C, P> BatchLogger<C, P>
80where
81 P: EventPusher<Timestamp, C>,
82{
83 fn new(event_pusher: P, interval_ms: u128, records_total: IntCounter) -> Self {
85 BatchLogger {
86 time_ms: Timestamp::minimum(),
87 event_pusher,
88 interval_ms,
89 records_total,
90 _marker: PhantomData,
91 }
92 }
93}
94
95impl<C, P> BatchLogger<C, P>
96where
97 P: EventPusher<Timestamp, C>,
98 C: Container,
99{
100 fn publish_batch(&mut self, data: C) {
102 let records = u64::try_from(data.record_count()).expect("record count is non-negative");
103 self.records_total.inc_by(records);
104 self.event_pusher.push(Event::Messages(self.time_ms, data));
105 }
106
107 fn report_progress(&mut self, time: Duration) -> bool {
111 let time_ms = ((time.as_millis() / self.interval_ms) + 1) * self.interval_ms;
112 let new_time_ms: Timestamp = time_ms.try_into().expect("must fit");
113 if self.time_ms < new_time_ms {
114 self.event_pusher
115 .push(Event::Progress(vec![(new_time_ms, 1), (self.time_ms, -1)]));
116 self.time_ms = new_time_ms;
117 true
118 } else {
119 false
120 }
121 }
122}
123
124impl<C, P> Drop for BatchLogger<C, P>
125where
126 P: EventPusher<Timestamp, C>,
127{
128 fn drop(&mut self) {
129 self.event_pusher
130 .push(Event::Progress(vec![(self.time_ms, -1)]));
131 }
132}
133
134#[derive(Clone)]
144struct EventQueue<C, const N: usize = 1> {
145 links: [Rc<EventLink<Timestamp, C>>; N],
146 activator: RcActivator,
147}
148
149impl<C, const N: usize> EventQueue<C, N> {
150 fn new(name: &str) -> Self {
151 let activator_name = format!("{name}_activator");
152 let activate_after = 128;
153 Self {
154 links: [(); N].map(|_| Rc::new(EventLink::new())),
155 activator: RcActivator::new(activator_name, activate_after),
156 }
157 }
158}
159
160#[derive(Default)]
162struct SharedLoggingState {
163 arrangement_size_activators: BTreeMap<usize, Activator>,
165 compute_logger: Option<ComputeLogger>,
167}
168
169pub(crate) struct PermutedRowPacker {
171 key: Vec<usize>,
172 value: Vec<usize>,
173 key_row: Row,
174 value_row: Row,
175}
176
177impl PermutedRowPacker {
178 pub(crate) fn new<V: Into<LogVariant>>(variant: V) -> Self {
180 let variant = variant.into();
181 let key = variant.index_by();
182 let (_, value) = permutation_for_arrangement(
183 &key.iter()
184 .cloned()
185 .map(MirScalarExpr::column)
186 .collect::<Vec<_>>(),
187 variant.desc().arity(),
188 );
189 Self {
190 key,
191 value,
192 key_row: Row::default(),
193 value_row: Row::default(),
194 }
195 }
196
197 pub(crate) fn pack_slice(&mut self, datums: &[Datum]) -> (&RowRef, &RowRef) {
199 self.pack_by_index(|packer, index| packer.push(datums[index]))
200 }
201
202 pub(crate) fn pack_by_index<F: Fn(&mut RowPacker, usize)>(
204 &mut self,
205 logic: F,
206 ) -> (&RowRef, &RowRef) {
207 let mut packer = self.key_row.packer();
208 for index in &self.key {
209 logic(&mut packer, *index);
210 }
211
212 let mut packer = self.value_row.packer();
213 for index in &self.value {
214 logic(&mut packer, *index);
215 }
216
217 (&self.key_row, &self.value_row)
218 }
219}
220
221pub(super) fn downgrade_to_interval_boundary(
233 cap: &mut Capability<Timestamp>,
234 activator: &Activator,
235 now: Instant,
236 start_offset: Duration,
237 interval_ms: u128,
238) -> Timestamp {
239 let elapsed = now.elapsed().as_millis();
240 let time_ms: u128 = ((elapsed + start_offset.as_millis()) / interval_ms + 1) * interval_ms;
241 let ts: Timestamp = time_ms.try_into().expect("must fit");
242 cap.downgrade(&ts);
243
244 let next_boundary_ms = time_ms - start_offset.as_millis();
245 let next_activation =
246 now + Duration::from_millis(next_boundary_ms.try_into().expect("must fit"));
247 activator.activate_after(next_activation.saturating_duration_since(Instant::now()));
248
249 ts
250}
251
252pub(super) fn emit_snapshot_diff<K, V, CB, P, F>(
262 session: &mut Session<'_, '_, Timestamp, CB, P>,
263 packer: &mut PermutedRowPacker,
264 prev: &BTreeMap<K, V>,
265 current: &BTreeMap<K, V>,
266 ts: Timestamp,
267 pack: F,
268) where
269 K: Ord,
270 V: PartialEq,
271 CB: ContainerBuilder + for<'a> PushInto<((&'a RowRef, &'a RowRef), Timestamp, Diff)>,
272 P: CapabilityTrait<Timestamp>,
273 F: for<'a> Fn(&'a mut PermutedRowPacker, &K, &V) -> (&'a RowRef, &'a RowRef),
274{
275 for (key, value) in prev {
276 if current.get(key) != Some(value) {
277 let row = pack(packer, key, value);
278 session.give((row, ts, Diff::MINUS_ONE));
279 }
280 }
281 for (key, value) in current {
282 if prev.get(key) != Some(value) {
283 let row = pack(packer, key, value);
284 session.give((row, ts, Diff::ONE));
285 }
286 }
287}
288
289struct LogCollection {
291 trace: RowRowAgent<Timestamp, Diff>,
293 token: Rc<dyn Any>,
295}
296
297pub(super) fn consolidate_and_pack<'scope, Chu, B, CB, L, F, C>(
304 input: Stream<'scope, Timestamp, C>,
305 log: L,
306 mut logic: F,
307) -> Stream<'scope, Timestamp, CB::Container>
308where
309 B: Batcher<Time = Timestamp> + 'static,
310 Chu: ContainerBuilder<Container = B::Output> + for<'a> PushInto<&'a mut C> + 'static,
311 C: Container + Clone + 'static,
312 B::Output: Clone,
313 CB: ContainerBuilder,
314 L: Into<LogVariant>,
315 F: FnMut(B::Output, &mut PermutedRowPacker, &mut OutputSession<CB>) + 'static,
316{
317 let log = log.into();
318 let c_name = &format!("Consolidate {log:?}");
320 let u_name = &format!("ToRow {log:?}");
321 let mut packer = PermutedRowPacker::new(log);
322 let consolidated = consolidate_pact::<Chu, B, _, _>(input, Pipeline, c_name);
323 consolidated.unary::<CB, _, _, _>(Pipeline, u_name, |_, _| {
324 move |input, output| {
325 input.for_each_time(|time, data| {
326 let mut session = output.session_with_builder(&time);
327 for item in data.flatten().flat_map(|data| data.drain(..)) {
328 logic(item, &mut packer, &mut session);
329 }
330 });
331 }
332 })
333}