Skip to main content

mz_storage/render/
sources.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//! Logic related to the creation of dataflow sources.
11//!
12//! See [`render_source`] for more details.
13
14use std::collections::BTreeMap;
15use std::iter;
16use std::sync::Arc;
17
18use differential_dataflow::{AsCollection, VecCollection};
19use mz_ore::cast::CastLossy;
20use mz_persist_client::operators::shard_source::SnapshotMode;
21use mz_repr::{Datum, Diff, GlobalId, Row, RowPacker};
22use mz_storage_operators::persist_source;
23use mz_storage_operators::persist_source::Subtime;
24use mz_storage_types::controller::CollectionMetadata;
25use mz_storage_types::dyncfgs;
26use mz_storage_types::errors::{
27    DataflowError, DecodeError, EnvelopeError, UpsertError, UpsertNullKeyError, UpsertValueError,
28};
29use mz_storage_types::parameters::StorageMaxInflightBytesConfig;
30use mz_storage_types::sources::envelope::{KeyEnvelope, NoneEnvelope, UpsertEnvelope, UpsertStyle};
31use mz_storage_types::sources::*;
32use mz_timely_util::builder_async::PressOnDropButton;
33use mz_timely_util::operator::CollectionExt;
34use mz_timely_util::order::refine_antichain;
35use serde::{Deserialize, Serialize};
36use timely::container::CapacityContainerBuilder;
37use timely::dataflow::StreamVec;
38use timely::dataflow::operators::vec::Map;
39use timely::dataflow::operators::{ConnectLoop, Feedback, Leave, OkErr};
40use timely::dataflow::scope::Scope;
41use timely::progress::{Antichain, Timestamp};
42
43use crate::decode::{render_decode_cdcv2, render_decode_delimited};
44use crate::healthcheck::{HealthStatusMessage, StatusNamespace};
45use crate::source::types::{DecodeResult, SourceOutput, SourceRender};
46use crate::source::{self, RawSourceCreationConfig, SourceExportCreationConfig};
47use crate::upsert::{UpsertKey, UpsertSourceTime, UpsertValue};
48
49/// _Renders_ complete _differential_ collections
50/// that represent the final source and its errors
51/// as requested by the original `CREATE SOURCE` statement,
52/// encapsulated in the passed `SourceInstanceDesc`.
53///
54/// The first element in the returned tuple is the pair of Collections,
55/// the second is a type-erased token that will keep the source
56/// alive as long as it is not dropped.
57///
58/// This function is intended to implement the recipe described here:
59/// <https://github.com/MaterializeInc/materialize/blob/main/doc/developer/platform/architecture-storage.md#source-ingestion>
60pub fn render_source<'scope, 'root, C>(
61    scope: Scope<'scope, mz_repr::Timestamp>,
62    root_scope: Scope<'root, ()>,
63    dataflow_debug_name: &String,
64    connection: C,
65    description: IngestionDescription<CollectionMetadata>,
66    resume_stream: StreamVec<'scope, mz_repr::Timestamp, ()>,
67    storage_state: &crate::storage_state::StorageState,
68    base_source_config: RawSourceCreationConfig,
69) -> (
70    BTreeMap<
71        GlobalId,
72        (
73            VecCollection<'scope, mz_repr::Timestamp, Row, Diff>,
74            VecCollection<'scope, mz_repr::Timestamp, DataflowError, Diff>,
75        ),
76    >,
77    Vec<StreamVec<'root, (), HealthStatusMessage>>,
78    StreamVec<'scope, mz_repr::Timestamp, ()>,
79    Vec<PressOnDropButton>,
80)
81where
82    C: SourceConnection + SourceRender + 'static,
83    C::Time: UpsertSourceTime,
84{
85    // Tokens that we should return from the method.
86    let mut needed_tokens = Vec::new();
87
88    // Note that this `render_source` attaches a single _instance_ of a source
89    // to the passed `Scope`, and this instance may be disabled if the
90    // source type does not support multiple instances. `render_source`
91    // is called on each timely worker as part of
92    // [`super::build_storage_dataflow`].
93
94    // A set of channels (1 per worker) used to signal rehydration being finished
95    // to raw sources. These are channels and not timely streams because they
96    // have to cross a scope boundary.
97    //
98    // Note that these will be entirely subsumed by full `hydration` backpressure,
99    // once that is implemented.
100    let (starter, mut start_signal) = tokio::sync::mpsc::channel::<()>(1);
101    let start_signal = async move {
102        let _ = start_signal.recv().await;
103    };
104
105    // Build the _raw_ ok and error sources using `create_raw_source` and the
106    // correct `SourceReader` implementations
107    let (exports, health, remap_upper, source_tokens) = source::create_raw_source(
108        scope,
109        root_scope,
110        storage_state,
111        resume_stream,
112        &base_source_config,
113        connection,
114        start_signal,
115    );
116
117    needed_tokens.extend(source_tokens);
118
119    let mut health_streams = Vec::with_capacity(exports.len() + 1);
120    health_streams.push(health);
121
122    let mut outputs = BTreeMap::new();
123    for (export_id, export) in exports {
124        type CB<C> = CapacityContainerBuilder<C>;
125        let (ok_stream, err_stream) =
126            export.map_fallible::<CB<_>, CB<_>, _, _, _>("export-demux-ok-err", |r| r);
127
128        // All sources should push their various error streams into this vector,
129        // whose contents will be concatenated and inserted along the collection.
130        // All subsources include the non-definite errors of the ingestion
131        let mut error_collections = Vec::new();
132
133        let data_config = base_source_config.source_exports[&export_id]
134            .data_config
135            .clone();
136        let (ok, extra_tokens, health_stream) = render_source_stream(
137            scope,
138            dataflow_debug_name,
139            export_id,
140            ok_stream,
141            data_config,
142            &description,
143            &mut error_collections,
144            storage_state,
145            &base_source_config,
146            starter.clone(),
147        );
148        needed_tokens.extend(extra_tokens);
149
150        // Flatten the error collections.
151        let err_collection = match error_collections.len() {
152            0 => err_stream,
153            _ => err_stream.concatenate(error_collections),
154        };
155
156        outputs.insert(export_id, (ok, err_collection));
157
158        health_streams.extend(health_stream.into_iter().map(|s| s.leave(root_scope)));
159    }
160    (outputs, health_streams, remap_upper, needed_tokens)
161}
162
163/// Completes the rendering of a particular source stream by applying decoding and envelope
164/// processing as necessary
165fn render_source_stream<'scope, FromTime>(
166    scope: Scope<'scope, mz_repr::Timestamp>,
167    dataflow_debug_name: &String,
168    export_id: GlobalId,
169    ok_source: VecCollection<'scope, mz_repr::Timestamp, SourceOutput<FromTime>, Diff>,
170    data_config: SourceExportDataConfig,
171    description: &IngestionDescription<CollectionMetadata>,
172    error_collections: &mut Vec<VecCollection<'scope, mz_repr::Timestamp, DataflowError, Diff>>,
173    storage_state: &crate::storage_state::StorageState,
174    base_source_config: &RawSourceCreationConfig,
175    rehydrated_token: impl std::any::Any + 'static,
176) -> (
177    VecCollection<'scope, mz_repr::Timestamp, Row, Diff>,
178    Vec<PressOnDropButton>,
179    Vec<StreamVec<'scope, mz_repr::Timestamp, HealthStatusMessage>>,
180)
181where
182    FromTime: Timestamp + Sync,
183    FromTime: UpsertSourceTime,
184{
185    let mut needed_tokens = vec![];
186
187    // Use the envelope and encoding configs for this particular source export
188    let SourceExportDataConfig { encoding, envelope } = data_config;
189
190    let SourceDesc {
191        connection: _,
192        timestamp_interval: _,
193    } = description.desc;
194
195    let (decoded_stream, decode_health) = match encoding {
196        None => (
197            ok_source.map(|r| DecodeResult {
198                // This is safe because the current set of sources produce
199                // either:
200                // 1. Non-nullable keys
201                // 2. No keys at all.
202                //
203                // Please see the comment on `key_envelope_no_encoding` in
204                // `mz_sql::plan::statement::ddl` for more details.
205                key: Some(Ok(r.key)),
206                value: Some(Ok(r.value)),
207                metadata: r.metadata,
208                from_time: r.from_time,
209            }),
210            None,
211        ),
212        Some(encoding) => {
213            let (decoded_stream, decode_health) = render_decode_delimited(
214                ok_source,
215                encoding.key,
216                encoding.value,
217                dataflow_debug_name.clone(),
218                storage_state.metrics.decode_defs.clone(),
219                storage_state.storage_configuration.clone(),
220            );
221            (decoded_stream, Some(decode_health))
222        }
223    };
224
225    // render envelopes
226    let (envelope_ok, envelope_health) = match &envelope {
227        SourceEnvelope::Upsert(upsert_envelope) => {
228            let upsert_input = upsert_commands(decoded_stream, upsert_envelope.clone());
229
230            let persist_clients = Arc::clone(&storage_state.persist_clients);
231            // TODO: Get this to work with the as_of.
232            let resume_upper = base_source_config.resume_uppers[&export_id].clone();
233
234            let upper_ts = resume_upper
235                .as_option()
236                .expect("resuming an already finished ingestion")
237                .clone();
238            let outer_mz_scope = scope.clone();
239            let (upsert, health_update) = scope.scoped(
240                &format!("upsert_rehydration_backpressure({})", export_id),
241                |scope| {
242                    let (
243                        previous_ok,
244                        previous_err,
245                        previous_token,
246                        feedback_handle,
247                        backpressure_metrics,
248                    ) = {
249                        let as_of = Antichain::from_elem(upper_ts.saturating_sub(1));
250
251                        let backpressure_max_inflight_bytes = get_backpressure_max_inflight_bytes(
252                            &storage_state
253                                .storage_configuration
254                                .parameters
255                                .storage_dataflow_max_inflight_bytes_config,
256                            &storage_state.instance_context.cluster_memory_limit,
257                        );
258
259                        let (feedback_handle, flow_control, backpressure_metrics) =
260                            if let Some(storage_dataflow_max_inflight_bytes) =
261                                backpressure_max_inflight_bytes
262                            {
263                                tracing::info!(
264                                    ?backpressure_max_inflight_bytes,
265                                    "timely-{} using backpressure in upsert for source {}",
266                                    base_source_config.worker_id,
267                                    export_id
268                                );
269                                if !storage_state
270                                    .storage_configuration
271                                    .parameters
272                                    .storage_dataflow_max_inflight_bytes_config
273                                    .disk_only
274                                    || storage_state.instance_context.scratch_directory.is_some()
275                                {
276                                    let (feedback_handle, feedback_data) =
277                                        scope.feedback(Default::default());
278
279                                    // TODO(guswynn): cleanup
280                                    let backpressure_metrics = Some(
281                                        base_source_config
282                                            .metrics
283                                            .get_backpressure_metrics(export_id, scope.index()),
284                                    );
285
286                                    (
287                                        Some(feedback_handle),
288                                        Some(persist_source::FlowControl {
289                                            progress_stream: feedback_data,
290                                            max_inflight_bytes: storage_dataflow_max_inflight_bytes,
291                                            summary: (Default::default(), Subtime::least_summary()),
292                                            metrics: backpressure_metrics
293                                                .as_ref()
294                                                .map(|m| m.operator_metrics()),
295                                        }),
296                                        backpressure_metrics,
297                                    )
298                                } else {
299                                    (None, None, None)
300                                }
301                            } else {
302                                (None, None, None)
303                            };
304
305                        let storage_metadata = description.source_exports[&export_id]
306                            .storage_metadata
307                            .clone();
308
309                        let error_handler =
310                            storage_state.error_handler("upsert_rehydration", export_id);
311
312                        let (ok_stream, err_stream, tok) = persist_source::persist_source_core(
313                            outer_mz_scope,
314                            scope,
315                            export_id,
316                            persist_clients,
317                            storage_metadata,
318                            None,
319                            Some(as_of),
320                            SnapshotMode::Include,
321                            Antichain::new(),
322                            None,
323                            flow_control,
324                            false.then_some(|| unreachable!()),
325                            async {},
326                            error_handler,
327                        );
328                        (
329                            ok_stream.as_collection(),
330                            err_stream.as_collection(),
331                            Some(tok),
332                            feedback_handle,
333                            backpressure_metrics,
334                        )
335                    };
336
337                    let export_statistics = storage_state
338                        .aggregated_statistics
339                        .get_source(&export_id)
340                        .expect("statistics initialized")
341                        .clone();
342                    let export_config = SourceExportCreationConfig {
343                        id: export_id,
344                        worker_id: base_source_config.worker_id,
345                        metrics: base_source_config.metrics.clone(),
346                        source_statistics: export_statistics,
347                    };
348                    let (upsert, health_update, snapshot_progress, upsert_token) =
349                        if dyncfgs::ENABLE_UPSERT_V2
350                            .get(storage_state.storage_configuration.config_set())
351                        {
352                            // Resolved here, at operator construction, so the
353                            // dataflow keeps one stash flavor for its whole
354                            // life even if the flag flips underneath it.
355                            let stash_flavor =
356                                crate::upsert_continual_feedback_v2::UpsertStashFlavor::from_config(
357                                    storage_state.storage_configuration.config_set(),
358                                );
359                            crate::upsert::upsert_v2(
360                                upsert_input.enter(scope),
361                                upsert_envelope.clone(),
362                                refine_antichain(&resume_upper),
363                                previous_ok,
364                                previous_err,
365                                previous_token,
366                                export_config,
367                                backpressure_metrics,
368                                stash_flavor,
369                            )
370                        } else {
371                            crate::upsert::upsert(
372                                upsert_input.enter(scope),
373                                upsert_envelope.clone(),
374                                refine_antichain(&resume_upper),
375                                previous_ok,
376                                previous_err,
377                                previous_token,
378                                export_config,
379                                &storage_state.instance_context,
380                                &storage_state.storage_configuration,
381                                &storage_state.dataflow_parameters,
382                                backpressure_metrics,
383                            )
384                        };
385
386                    // Even though we register the `persist_sink` token at a top-level,
387                    // which will stop any data from being committed, we also register
388                    // a token for the `upsert` operator which may be in the middle of
389                    // rehydration processing the `persist_source` input above.
390                    needed_tokens.push(upsert_token);
391
392                    // If configured, delay raw sources until we rehydrate the upsert
393                    // source. Otherwise, drop the token, unblocking the sources at the
394                    // end rendering.
395                    if dyncfgs::DELAY_SOURCES_PAST_REHYDRATION
396                        .get(storage_state.storage_configuration.config_set())
397                    {
398                        crate::upsert::rehydration_finished(
399                            scope.clone(),
400                            base_source_config,
401                            rehydrated_token,
402                            refine_antichain(&resume_upper),
403                            snapshot_progress.clone(),
404                        );
405                    } else {
406                        drop(rehydrated_token)
407                    };
408
409                    // If backpressure from persist is enabled, we connect the upsert operator's
410                    // snapshot progress to the persist source feedback handle.
411                    if let Some(feedback_handle) = feedback_handle {
412                        snapshot_progress.connect_loop(feedback_handle);
413                    }
414
415                    (
416                        upsert.leave(outer_mz_scope),
417                        health_update
418                            .map(|(id, update)| HealthStatusMessage {
419                                id,
420                                namespace: StatusNamespace::Upsert,
421                                update,
422                            })
423                            .leave(outer_mz_scope),
424                    )
425                },
426            );
427
428            let (upsert_ok, upsert_err) = upsert.inner.ok_err(split_ok_err);
429            error_collections.push(upsert_err.as_collection());
430
431            (upsert_ok.as_collection(), Some(health_update))
432        }
433        SourceEnvelope::None(none_envelope) => {
434            let results = append_metadata_to_value(decoded_stream);
435
436            let flattened_stream = flatten_results_prepend_keys(none_envelope, results);
437
438            let (stream, errors) = flattened_stream.inner.ok_err(split_ok_err);
439
440            error_collections.push(errors.as_collection());
441            (stream.as_collection(), None)
442        }
443        SourceEnvelope::CdcV2 => {
444            let (oks, token) = render_decode_cdcv2(&decoded_stream);
445            needed_tokens.push(token);
446            (oks, None)
447        }
448    };
449
450    // Return the collections and any needed tokens.
451    let health = decode_health.into_iter().chain(envelope_health).collect();
452    (envelope_ok, needed_tokens, health)
453}
454
455// Returns the maximum limit of inflight bytes for backpressure based on given config
456// and the current cluster size
457fn get_backpressure_max_inflight_bytes(
458    inflight_bytes_config: &StorageMaxInflightBytesConfig,
459    cluster_memory_limit: &Option<usize>,
460) -> Option<usize> {
461    let StorageMaxInflightBytesConfig {
462        max_inflight_bytes_default,
463        max_inflight_bytes_cluster_size_fraction,
464        disk_only: _,
465    } = inflight_bytes_config;
466
467    // Will use backpressure only if the default inflight value is provided
468    if max_inflight_bytes_default.is_some() {
469        let current_cluster_max_bytes_limit =
470            cluster_memory_limit.as_ref().and_then(|cluster_memory| {
471                max_inflight_bytes_cluster_size_fraction.map(|fraction| {
472                    // We just need close the correct % of bytes here, so we just use lossy casts.
473                    usize::cast_lossy(f64::cast_lossy(*cluster_memory) * fraction)
474                })
475            });
476        current_cluster_max_bytes_limit.or(*max_inflight_bytes_default)
477    } else {
478        None
479    }
480}
481
482// TODO: Maybe we should finally move this to some central place and re-use. There seem to be
483// enough instances of this by now.
484fn split_ok_err<O, E, T, D>(x: (Result<O, E>, T, D)) -> Result<(O, T, D), (E, T, D)> {
485    match x {
486        (Ok(ok), ts, diff) => Ok((ok, ts, diff)),
487        (Err(err), ts, diff) => Err((err, ts, diff)),
488    }
489}
490
491/// After handling metadata insertion, we split streams into key/value parts for convenience
492#[derive(
493    Debug,
494    Clone,
495    Hash,
496    PartialEq,
497    Eq,
498    Ord,
499    PartialOrd,
500    Serialize,
501    Deserialize
502)]
503struct KV {
504    key: Option<Result<Row, DecodeError>>,
505    val: Option<Result<Row, DecodeError>>,
506}
507
508fn append_metadata_to_value<'scope, T: Timestamp, FromTime: Timestamp>(
509    results: VecCollection<'scope, T, DecodeResult<FromTime>, Diff>,
510) -> VecCollection<'scope, T, KV, Diff> {
511    results.map(move |res| {
512        let val = res.value.map(|val_result| {
513            val_result.map(|mut val| {
514                if !res.metadata.is_empty() {
515                    RowPacker::for_existing_row(&mut val).extend_by_row(&res.metadata);
516                }
517                val
518            })
519        });
520
521        KV { val, key: res.key }
522    })
523}
524
525/// Convert from streams of [`DecodeResult`] to UpsertCommands, inserting the Key according to [`KeyEnvelope`]
526fn upsert_commands<'scope, T: Timestamp, FromTime: Timestamp>(
527    input: VecCollection<'scope, T, DecodeResult<FromTime>, Diff>,
528    upsert_envelope: UpsertEnvelope,
529) -> VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff> {
530    let mut row_buf = Row::default();
531    input.map(move |result| {
532        let from_time = result.from_time;
533
534        let key = match result.key {
535            Some(Ok(key)) => Ok(key),
536            None => Err(UpsertError::NullKey(UpsertNullKeyError)),
537            Some(Err(err)) => Err(UpsertError::KeyDecode(err)),
538        };
539
540        // If we have a well-formed key we can continue, otherwise we're upserting an error
541        let key = match key {
542            Ok(key) => key,
543            Err(err) => match result.value {
544                Some(_) => {
545                    return (
546                        UpsertKey::from_key(Err(&err)),
547                        Some(Err(Box::new(err))),
548                        from_time,
549                    );
550                }
551                None => return (UpsertKey::from_key(Err(&err)), None, from_time),
552            },
553        };
554
555        // We can now apply the key envelope
556        let key_row = match upsert_envelope.style {
557            // flattened or debezium
558            UpsertStyle::Debezium { .. }
559            | UpsertStyle::Default(KeyEnvelope::Flattened)
560            | UpsertStyle::ValueErrInline {
561                key_envelope: KeyEnvelope::Flattened,
562                error_column: _,
563            } => key,
564            // named
565            UpsertStyle::Default(KeyEnvelope::Named(_))
566            | UpsertStyle::ValueErrInline {
567                key_envelope: KeyEnvelope::Named(_),
568                error_column: _,
569            } => {
570                if key.iter().nth(1).is_none() {
571                    key
572                } else {
573                    row_buf.packer().push_list(key.iter());
574                    row_buf.clone()
575                }
576            }
577            UpsertStyle::Default(KeyEnvelope::None)
578            | UpsertStyle::ValueErrInline {
579                key_envelope: KeyEnvelope::None,
580                error_column: _,
581            } => unreachable!(),
582        };
583
584        let key = UpsertKey::from_key(Ok(&key_row));
585
586        let metadata = result.metadata;
587
588        let value = match result.value {
589            Some(Ok(ref row)) => match upsert_envelope.style {
590                UpsertStyle::Debezium { after_idx } => match row.iter().nth(after_idx).unwrap() {
591                    Datum::List(after) => {
592                        let mut packer = row_buf.packer();
593                        packer.extend(after.iter());
594                        packer.extend_by_row(&metadata);
595                        Some(Ok(row_buf.clone()))
596                    }
597                    Datum::Null => None,
598                    d => panic!("type error: expected record, found {:?}", d),
599                },
600                UpsertStyle::Default(_) => {
601                    let mut packer = row_buf.packer();
602                    packer.extend_by_row(&key_row);
603                    packer.extend_by_row(row);
604                    packer.extend_by_row(&metadata);
605                    Some(Ok(row_buf.clone()))
606                }
607                UpsertStyle::ValueErrInline { .. } => {
608                    let mut packer = row_buf.packer();
609                    packer.extend_by_row(&key_row);
610                    // The 'error' column is null
611                    packer.push(Datum::Null);
612                    packer.extend_by_row(row);
613                    packer.extend_by_row(&metadata);
614                    Some(Ok(row_buf.clone()))
615                }
616            },
617            Some(Err(inner)) => {
618                match upsert_envelope.style {
619                    UpsertStyle::ValueErrInline { .. } => {
620                        let mut count = 0;
621                        // inline the error in the data output
622                        let err_string = inner.to_string();
623                        let mut packer = row_buf.packer();
624                        for datum in key_row.iter() {
625                            packer.push(datum);
626                            count += 1;
627                        }
628                        // The 'error' column is a record with a 'description' column
629                        packer.push_list(iter::once(Datum::String(&err_string)));
630                        count += 1;
631                        let metadata_len = metadata.as_row_ref().iter().count();
632                        // push nulls for all value columns
633                        packer.extend(
634                            iter::repeat(Datum::Null)
635                                .take(upsert_envelope.source_arity - count - metadata_len),
636                        );
637                        packer.extend_by_row(&metadata);
638                        Some(Ok(row_buf.clone()))
639                    }
640                    _ => Some(Err(Box::new(UpsertError::Value(UpsertValueError {
641                        for_key: key_row,
642                        inner,
643                    })))),
644                }
645            }
646            None => None,
647        };
648
649        (key, value, from_time)
650    })
651}
652
653/// Convert from streams of [`DecodeResult`] to Rows, inserting the Key according to [`KeyEnvelope`]
654fn flatten_results_prepend_keys<'scope, T: Timestamp>(
655    none_envelope: &NoneEnvelope,
656    results: VecCollection<'scope, T, KV, Diff>,
657) -> VecCollection<'scope, T, Result<Row, DataflowError>, Diff> {
658    let NoneEnvelope {
659        key_envelope,
660        key_arity,
661    } = none_envelope;
662
663    let null_key_columns = Row::pack_slice(&vec![Datum::Null; *key_arity]);
664
665    match key_envelope {
666        KeyEnvelope::None => {
667            results.flat_map(|KV { val, .. }| val.map(|result| result.map_err(Into::into)))
668        }
669        KeyEnvelope::Flattened => results
670            .flat_map(raise_key_value_errors)
671            .map(move |maybe_kv| {
672                maybe_kv.map(|(key, value)| {
673                    let mut key = key.unwrap_or_else(|| null_key_columns.clone());
674                    RowPacker::for_existing_row(&mut key).extend_by_row(&value);
675                    key
676                })
677            }),
678        KeyEnvelope::Named(_) => {
679            results
680                .flat_map(raise_key_value_errors)
681                .map(move |maybe_kv| {
682                    maybe_kv.map(|(key, value)| {
683                        let mut key = key.unwrap_or_else(|| null_key_columns.clone());
684                        // Named semantics rename a key that is a single column, and encode a
685                        // multi-column field as a struct with that name
686                        let row = if key.iter().nth(1).is_none() {
687                            RowPacker::for_existing_row(&mut key).extend_by_row(&value);
688                            key
689                        } else {
690                            let mut new_row = Row::default();
691                            let mut packer = new_row.packer();
692                            packer.push_list(key.iter());
693                            packer.extend_by_row(&value);
694                            new_row
695                        };
696                        row
697                    })
698                })
699        }
700    }
701}
702
703/// Handle possibly missing key or value portions of messages
704fn raise_key_value_errors(
705    KV { key, val }: KV,
706) -> Option<Result<(Option<Row>, Row), DataflowError>> {
707    match (key, val) {
708        (Some(Ok(key)), Some(Ok(value))) => Some(Ok((Some(key), value))),
709        (None, Some(Ok(value))) => Some(Ok((None, value))),
710        // always prioritize the value error if either or both have an error
711        (_, Some(Err(e))) => Some(Err(e.into())),
712        (Some(Err(e)), _) => Some(Err(e.into())),
713        (None, None) => None,
714        // TODO(petrosagg): these errors would be better grouped under an EnvelopeError enum
715        _ => Some(Err(DataflowError::from(EnvelopeError::Flat(
716            "Value not present for message".into(),
717        )))),
718    }
719}
720
721#[cfg(test)]
722mod test {
723    use super::*;
724
725    #[mz_ore::test]
726    fn test_no_default() {
727        let config = StorageMaxInflightBytesConfig {
728            max_inflight_bytes_default: None,
729            max_inflight_bytes_cluster_size_fraction: Some(0.5),
730            disk_only: false,
731        };
732        let memory_limit = Some(1000);
733
734        let backpressure_inflight_bytes_limit =
735            get_backpressure_max_inflight_bytes(&config, &memory_limit);
736
737        assert_eq!(backpressure_inflight_bytes_limit, None)
738    }
739
740    #[mz_ore::test]
741    fn test_no_matching_size() {
742        let config = StorageMaxInflightBytesConfig {
743            max_inflight_bytes_default: Some(10000),
744            max_inflight_bytes_cluster_size_fraction: Some(0.5),
745            disk_only: false,
746        };
747
748        let backpressure_inflight_bytes_limit = get_backpressure_max_inflight_bytes(&config, &None);
749
750        assert_eq!(
751            backpressure_inflight_bytes_limit,
752            config.max_inflight_bytes_default
753        )
754    }
755
756    #[mz_ore::test]
757    fn test_calculated_cluster_limit() {
758        let config = StorageMaxInflightBytesConfig {
759            max_inflight_bytes_default: Some(10000),
760            max_inflight_bytes_cluster_size_fraction: Some(0.5),
761            disk_only: false,
762        };
763        let memory_limit = Some(2000);
764
765        let backpressure_inflight_bytes_limit =
766            get_backpressure_max_inflight_bytes(&config, &memory_limit);
767
768        // the limit should be 50% of 2000 i.e. 1000
769        assert_eq!(backpressure_inflight_bytes_limit, Some(1000));
770    }
771}