mz_compute/render/
errors.rs1use columnar::Columnar;
28use columnation::{Columnation, Region};
29use mz_expr::EvalError;
30use mz_proto::{ProtoType, RustType};
31use mz_repr::Row;
32use mz_storage_types::errors::{DataflowError, ProtoDataflowError};
33use prost::Message;
34use serde::{Deserialize, Serialize};
35use std::fmt;
36
37#[derive(
44 Clone,
45 Eq,
46 PartialEq,
47 Ord,
48 PartialOrd,
49 Hash,
50 Serialize,
51 Deserialize,
52 Columnar
53)]
54#[columnar(derive(Eq, PartialEq, Ord, PartialOrd))]
55pub struct DataflowErrorSer(Vec<u8>);
56
57impl DataflowErrorSer {
58 pub fn deserialize(&self) -> DataflowError {
64 let proto = ProtoDataflowError::decode(self.0.as_slice())
65 .expect("DataflowErrorSer: invalid proto bytes");
66 proto
67 .into_rust()
68 .expect("DataflowErrorSer: failed to convert proto to DataflowError")
69 }
70}
71
72impl From<DataflowError> for DataflowErrorSer {
73 fn from(err: DataflowError) -> Self {
74 DataflowErrorSer(err.into_proto().encode_to_vec())
75 }
76}
77
78impl From<EvalError> for DataflowErrorSer {
79 fn from(err: EvalError) -> Self {
80 DataflowErrorSer::from(DataflowError::from(err))
83 }
84}
85
86impl fmt::Display for DataflowErrorSer {
87 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
88 self.deserialize().fmt(f)
89 }
90}
91
92impl fmt::Debug for DataflowErrorSer {
93 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
94 f.debug_tuple("DataflowErrorSer")
95 .field(&format_args!("{} bytes", self.0.len()))
96 .finish()
97 }
98}
99
100impl Columnation for DataflowErrorSer {
101 type InnerRegion = DataflowErrorSerRegion;
102}
103
104impl crate::sink::correction_v2::DataBytes for DataflowErrorSer {
105 fn data_bytes(&self) -> usize {
106 std::mem::size_of::<Self>() + self.0.len()
107 }
108}
109
110#[derive(Default)]
112pub struct DataflowErrorSerRegion {
113 inner: <Vec<u8> as Columnation>::InnerRegion,
114}
115
116impl Region for DataflowErrorSerRegion {
117 type Item = DataflowErrorSer;
118
119 unsafe fn copy(&mut self, item: &Self::Item) -> Self::Item {
120 DataflowErrorSer(unsafe { self.inner.copy(&item.0) })
122 }
123
124 fn clear(&mut self) {
125 self.inner.clear();
126 }
127
128 fn reserve_items<'a, I>(&mut self, items: I)
129 where
130 I: Iterator<Item = &'a Self::Item> + Clone,
131 {
132 self.inner.reserve_items(items.map(|item| &item.0));
133 }
134
135 fn reserve_regions<'a, I>(&mut self, regions: I)
136 where
137 I: Iterator<Item = &'a Self> + Clone,
138 {
139 self.inner.reserve_regions(regions.map(|r| &r.inner));
140 }
141
142 fn heap_size(&self, callback: impl FnMut(usize, usize)) {
143 self.inner.heap_size(callback);
144 }
145}
146
147pub(super) trait MaybeValidatingRow<T, E> {
151 fn ok(t: T) -> Self;
152 fn into_error() -> Option<fn(E) -> Self>;
153}
154
155impl<E> MaybeValidatingRow<Row, E> for Row {
156 fn ok(t: Row) -> Self {
157 t
158 }
159
160 fn into_error() -> Option<fn(E) -> Self> {
161 None
162 }
163}
164
165impl<E> MaybeValidatingRow<(), E> for () {
166 fn ok(t: ()) -> Self {
167 t
168 }
169
170 fn into_error() -> Option<fn(E) -> Self> {
171 None
172 }
173}
174
175impl<E, R> MaybeValidatingRow<Vec<R>, E> for Vec<R> {
176 fn ok(t: Vec<R>) -> Self {
177 t
178 }
179
180 fn into_error() -> Option<fn(E) -> Self> {
181 None
182 }
183}
184
185impl<T, E> MaybeValidatingRow<T, E> for Result<T, E> {
186 fn ok(row: T) -> Self {
187 Ok(row)
188 }
189
190 fn into_error() -> Option<fn(E) -> Self> {
191 Some(Err)
192 }
193}
194
195#[derive(Clone)]
198pub(super) struct ErrorLogger {
199 dataflow_name: String,
200}
201
202impl ErrorLogger {
203 pub fn new(dataflow_name: String) -> Self {
204 Self { dataflow_name }
205 }
206
207 pub fn log(&self, message: &'static str, details: &str) {
224 tracing::warn!(
225 dataflow = self.dataflow_name,
226 "[customer-data] {message} ({details})"
227 );
228 tracing::error!(message);
229 }
230
231 pub fn soft_panic_or_log(&self, message: &'static str, details: &str) {
235 tracing::warn!(
236 dataflow = self.dataflow_name,
237 "[customer-data] {message} ({details})"
238 );
239 mz_ore::soft_panic_or_log!("{}", message);
240 }
241}
242
243#[cfg(test)]
244mod tests {
245 use super::*;
246 use mz_storage_types::errors::DataflowError;
247 use proptest::prelude::*;
248
249 #[mz_ore::test]
250 #[cfg_attr(miri, ignore)]
251 fn proptest_roundtrip_canonical() {
252 proptest!(|(err in any::<DataflowError>())| {
253 let ser = DataflowErrorSer::from(err.clone());
254
255 let deserialized = ser.deserialize();
257 let re_serialized = DataflowErrorSer::from(deserialized);
258 prop_assert_eq!(&ser, &re_serialized,
259 "Canonicality violation: round-trip produced different bytes");
260
261 let ser2 = DataflowErrorSer::from(err);
263 prop_assert_eq!(&ser, &ser2,
264 "Canonicality violation: same error produced different bytes");
265 });
266 }
267
268 #[mz_ore::test]
269 fn display_roundtrip() {
270 let eval_err = EvalError::DivisionByZero;
271 let dfe = DataflowError::from(eval_err.clone());
272 let ser = DataflowErrorSer::from(eval_err);
273
274 assert_eq!(dfe.to_string(), ser.to_string());
275 }
276}