1use std::any::Any;
2use std::borrow::Cow;
3use std::cmp::max;
4use std::fmt::{Display, Error as FmtError, Formatter};
5use std::fs::File;
6use std::hash::Hasher;
7use std::io::{Error as IoError, Read};
8use std::ops::{Add, AddAssign, Range};
9use std::sync::Arc;
10
11use actix_web::HttpRequest;
12use anyhow::Result as AnyResult;
13use dbsp::operator::input::StagedBuffers;
14use erased_serde::Serialize as ErasedSerialize;
15use feldera_types::config::ConnectorConfig;
16use feldera_types::program_schema::Relation;
17use feldera_types::serde_with_context::FieldParseError;
18use serde::Serialize;
19use serde::de::StdError;
20
21use crate::ConnectorMetadata;
22use crate::catalog::{InputCollectionHandle, SerBatchReader};
23use crate::errors::controller::ControllerError;
24use crate::postprocess::Postprocessor;
25use crate::preprocess::Preprocessor;
26use crate::transport::{OutputBatchType, Step};
27
28pub trait InputFormat: Send + Sync {
32 fn name(&self) -> Cow<'static, str>;
34
35 fn config_from_http_request(
49 &self,
50 endpoint_name: &str,
51 request: &HttpRequest,
52 ) -> Result<Box<dyn ErasedSerialize>, ControllerError>;
53
54 fn new_parser(
62 &self,
63 endpoint_name: &str,
64 input_stream: &InputCollectionHandle,
65 config: &serde_json::Value,
66 ) -> Result<Box<dyn Parser>, ControllerError>;
67}
68
69pub trait InputBuffer: Any + Send {
95 fn flush(&mut self);
98
99 fn len(&self) -> BufferSize;
101
102 fn hash(&self, hasher: &mut dyn Hasher);
110
111 fn is_empty(&self) -> bool {
112 self.len().is_empty()
113 }
114
115 fn take_some(&mut self, n: usize) -> Option<Box<dyn InputBuffer>>;
135
136 fn take_all(&mut self) -> Option<Box<dyn InputBuffer>> {
137 self.take_some(usize::MAX)
138 }
139}
140
141#[derive(Copy, Clone, Debug, Default, PartialEq, Eq)]
143pub struct BufferSize {
144 pub records: usize,
146
147 pub bytes: usize,
152}
153
154impl BufferSize {
155 pub fn empty() -> Self {
157 Self::default()
158 }
159
160 pub fn is_empty(&self) -> bool {
162 self.records == 0 && self.bytes == 0
163 }
164}
165
166impl Add for BufferSize {
167 type Output = Self;
168
169 fn add(self, rhs: Self) -> Self::Output {
170 BufferSize {
171 records: self.records + rhs.records,
172 bytes: self.bytes + rhs.bytes,
173 }
174 }
175}
176
177impl AddAssign for BufferSize {
178 fn add_assign(&mut self, rhs: Self) {
179 self.records += rhs.records;
180 self.bytes += rhs.bytes;
181 }
182}
183
184impl InputBuffer for Option<Box<dyn InputBuffer>> {
185 fn len(&self) -> BufferSize {
186 self.as_ref()
187 .map_or(BufferSize::empty(), |buffer| buffer.len())
188 }
189
190 fn hash(&self, hasher: &mut dyn Hasher) {
191 if let Some(buffer) = self {
192 buffer.hash(hasher)
193 }
194 }
195
196 fn flush(&mut self) {
197 if let Some(buffer) = self.as_mut() {
198 buffer.flush()
199 }
200 }
201
202 fn take_some(&mut self, n: usize) -> Option<Box<dyn InputBuffer>> {
203 self.as_mut().and_then(|buffer| buffer.take_some(n))
204 }
205}
206
207impl InputBuffer for Box<dyn InputBuffer> {
208 fn len(&self) -> BufferSize {
209 self.as_ref().len()
210 }
211
212 fn hash(&self, hasher: &mut dyn Hasher) {
213 self.as_ref().hash(hasher)
214 }
215
216 fn flush(&mut self) {
217 self.as_mut().flush()
218 }
219
220 fn take_some(&mut self, n: usize) -> Option<Box<dyn InputBuffer>> {
221 self.as_mut().take_some(n)
222 }
223}
224
225impl InputBuffer for Vec<Box<dyn InputBuffer>> {
226 fn flush(&mut self) {
227 for v in self.iter_mut() {
228 v.flush();
229 }
230 }
231
232 fn len(&self) -> BufferSize {
233 let mut size = BufferSize::empty();
234 for v in self.iter() {
235 size += v.len();
236 }
237 size
238 }
239
240 fn hash(&self, hasher: &mut dyn Hasher) {
241 for v in self.iter() {
242 v.hash(hasher);
243 }
244 }
245
246 fn take_some(&mut self, n: usize) -> Option<Box<dyn InputBuffer>> {
247 let mut result = Vec::new();
248 let mut remaining = n;
249 let mut index = 0;
251 for v in self.iter_mut() {
252 if remaining == 0 {
253 break;
254 }
255 let buf = v.take_some(remaining);
256 if let Some(ib) = buf {
257 let len = ib.len().records;
258 if remaining >= len {
259 index += 1;
261 }
262 remaining = remaining.saturating_sub(len);
263 result.push(ib);
264 }
265 }
266 self.drain(0..index);
267 if result.is_empty() {
268 None
269 } else {
270 Some(Box::new(result))
271 }
272 }
273}
274
275pub fn flatten_nested<T>(buffers: Vec<Box<dyn InputBuffer>>) -> Vec<Box<T>>
278where
279 T: Any,
280{
281 fn inner<T>(input: Vec<Box<dyn InputBuffer>>, output: &mut Vec<Box<T>>)
282 where
283 T: Any,
284 {
285 for buffer in input {
286 let any = buffer as Box<dyn Any>;
287 match any.downcast::<Vec<Box<dyn InputBuffer>>>() {
288 Ok(vec) => inner(*vec, output),
289 Err(any) => output.push(any.downcast().unwrap()),
290 }
291 }
292 }
293
294 let mut output = Vec::new();
295 inner(buffers, &mut output);
296 output
297}
298
299pub struct StagedInputBuffer {
315 buffer: Box<dyn StagedBuffers>,
316 size: BufferSize,
317}
318
319impl StagedInputBuffer {
320 pub fn new(buffer: Box<dyn StagedBuffers>, size: BufferSize) -> Self {
321 Self { buffer, size }
322 }
323}
324
325impl InputBuffer for StagedInputBuffer {
326 fn flush(&mut self) {
327 self.buffer.flush()
328 }
329
330 fn len(&self) -> BufferSize {
331 self.size
332 }
333
334 fn hash(&self, _hasher: &mut dyn Hasher) {
335 unimplemented!()
336 }
337
338 fn take_some(&mut self, _n: usize) -> Option<Box<dyn InputBuffer>> {
339 unimplemented!()
340 }
341}
342
343pub trait Parser: Send + Sync {
345 fn parse(
351 &mut self,
352 data: &[u8],
353 metadata: Option<ConnectorMetadata>,
354 ) -> (Option<Box<dyn InputBuffer>>, Vec<ParseError>);
355
356 fn stage(&self, buffers: Vec<Box<dyn InputBuffer>>) -> Box<dyn StagedBuffers>;
361
362 fn splitter(&self) -> Box<dyn Splitter>;
365
366 fn fork(&self) -> Box<dyn Parser>;
371}
372
373pub struct StreamingPreprocessedParser {
375 preprocessor: Box<dyn Preprocessor>,
376 stream_splitter: StreamSplitter,
377 parser: Box<dyn Parser>,
378}
379
380impl StreamingPreprocessedParser {
381 pub fn new(preprocessor: Box<dyn Preprocessor>, parser: Box<dyn Parser>) -> Self {
382 Self {
383 preprocessor,
384 stream_splitter: StreamSplitter::new(parser.splitter()),
385 parser,
386 }
387 }
388}
389
390impl Parser for StreamingPreprocessedParser {
391 fn parse(
392 &mut self,
393 data: &[u8],
394 metadata: Option<ConnectorMetadata>,
395 ) -> (Option<Box<dyn InputBuffer>>, Vec<ParseError>) {
396 let (pre_data, mut pre_errors) = self
397 .preprocessor
398 .process_with_metadata(data, metadata.as_ref());
399 self.stream_splitter.append(&pre_data);
400 let mut parsed: Vec<Box<dyn InputBuffer>> = Vec::new();
401 while let Some(chunk) = self.stream_splitter.next(true) {
402 let (parsed_data, mut parse_errors) = self.parser.parse(chunk, metadata.clone());
403 pre_errors.append(&mut parse_errors);
404 if let Some(data) = parsed_data {
405 parsed.push(data);
406 }
407 }
408 if parsed.is_empty() {
409 (None, pre_errors)
410 } else {
411 (Some(Box::new(parsed)), pre_errors)
412 }
413 }
414
415 fn stage(&self, buffers: Vec<Box<dyn InputBuffer>>) -> Box<dyn StagedBuffers> {
416 self.parser.stage(buffers)
417 }
418
419 fn splitter(&self) -> Box<dyn Splitter> {
420 let pre_splitter = self.preprocessor.splitter();
421 if let Some(splitter) = pre_splitter {
422 return splitter;
423 }
424 self.parser.splitter()
425 }
426
427 fn fork(&self) -> Box<dyn Parser> {
428 Box::new(StreamingPreprocessedParser::new(
429 self.preprocessor.fork(),
430 self.parser.fork(),
431 ))
432 }
433}
434
435pub struct MessageOrientedPreprocessedParser {
437 preprocessor: Box<dyn Preprocessor>,
438 parser: Box<dyn Parser>,
439}
440
441impl MessageOrientedPreprocessedParser {
442 pub fn new(preprocessor: Box<dyn Preprocessor>, parser: Box<dyn Parser>) -> Self {
443 Self {
444 preprocessor,
445 parser,
446 }
447 }
448}
449
450impl Parser for MessageOrientedPreprocessedParser {
451 fn parse(
452 &mut self,
453 data: &[u8],
454 metadata: Option<ConnectorMetadata>,
455 ) -> (Option<Box<dyn InputBuffer>>, Vec<ParseError>) {
456 let (pre_data, mut pre_errors) = self
457 .preprocessor
458 .process_with_metadata(data, metadata.as_ref());
459 let mut parser_splitter = self.parser.splitter();
460 let mut parsed: Vec<Box<dyn InputBuffer>> = Vec::new();
461 let mut remaining = pre_data.as_slice();
462 while !remaining.is_empty() {
464 let chunk;
465 let split_offset = parser_splitter.input(remaining).unwrap_or(remaining.len());
466 (chunk, remaining) = remaining.split_at(split_offset);
467 let (parsed_data, mut parse_errors) = self.parser.parse(chunk, metadata.clone());
468 pre_errors.append(&mut parse_errors);
469 if let Some(data) = parsed_data {
470 parsed.push(data);
471 }
472 }
473 if parsed.is_empty() {
474 (None, pre_errors)
475 } else {
476 (Some(Box::new(parsed)), pre_errors)
477 }
478 }
479
480 fn stage(&self, buffers: Vec<Box<dyn InputBuffer>>) -> Box<dyn StagedBuffers> {
481 self.parser.stage(buffers)
482 }
483
484 fn splitter(&self) -> Box<dyn Splitter> {
485 let pre_splitter = self.preprocessor.splitter();
486 if let Some(splitter) = pre_splitter {
487 return splitter;
488 }
489 self.parser.splitter()
490 }
491
492 fn fork(&self) -> Box<dyn Parser> {
493 Box::new(MessageOrientedPreprocessedParser::new(
494 self.preprocessor.fork(),
495 self.parser.fork(),
496 ))
497 }
498}
499
500pub trait Splitter: Send + Sync {
506 fn input(&mut self, data: &[u8]) -> Option<usize>;
515
516 fn clear(&mut self);
519}
520
521pub struct StreamSplitter {
528 buffer: Vec<u8>,
529 start: u64,
530 fragment: Range<usize>,
531 fed: usize,
532 splitter: Box<dyn Splitter>,
533}
534
535impl StreamSplitter {
536 pub fn new(splitter: Box<dyn Splitter>) -> Self {
538 Self {
539 buffer: Vec::new(),
540 start: 0,
541 fragment: 0..0,
542 fed: 0,
543 splitter,
544 }
545 }
546
547 pub fn next(&mut self, eoi: bool) -> Option<&[u8]> {
551 match self
552 .splitter
553 .input(&self.buffer[self.fed..self.fragment.end])
554 {
555 Some(n) => {
556 let chunk = &self.buffer[self.fragment.start..self.fed + n];
557 self.fed += n;
558 self.fragment.start = self.fed;
559 Some(chunk)
560 }
561 None => {
562 self.fed = self.fragment.end;
563 if eoi && !self.fragment.is_empty() {
564 let chunk = &self.buffer[self.fragment.clone()];
565 self.fragment.start = self.fragment.end;
566 Some(chunk)
567 } else {
568 None
569 }
570 }
571 }
572 }
573
574 pub fn append(&mut self, data: &[u8]) {
576 let final_len = self.fragment.len() + data.len();
577 if final_len > self.buffer.len() {
578 self.buffer.reserve(final_len - self.buffer.len());
579 }
580 self.buffer.copy_within(self.fragment.clone(), 0);
581 self.buffer.resize(self.fragment.len(), 0);
582 self.buffer.extend(data);
583 self.fed -= self.fragment.start;
584 self.start += self.fragment.start as u64;
585 self.fragment = 0..self.buffer.len();
586 }
587
588 pub fn read(
592 &mut self,
593 file: &mut File,
594 buffer_size: usize,
595 limit: usize,
596 ) -> Result<usize, IoError> {
597 if self.fragment.start != 0 {
599 self.buffer.copy_within(self.fragment.clone(), 0);
600 self.fed -= self.fragment.start;
601 self.start += self.fragment.start as u64;
602 self.fragment = 0..self.fragment.len();
603 }
604
605 if self.fragment.len() == self.buffer.len() {
607 self.buffer
608 .resize(max(buffer_size, self.buffer.capacity() * 2), 0);
609 }
610
611 let mut space = &mut self.buffer[self.fragment.len()..];
613 if space.len() > limit {
614 space = &mut space[..limit];
615 }
616 let result = file.read(space);
617 if let Ok(n) = result {
618 self.fragment.end += n;
619 }
620 result
621 }
622
623 pub fn position(&self) -> u64 {
626 self.start + self.fragment.start as u64
627 }
628
629 pub fn seek(&mut self, offset: u64) {
632 self.start = offset;
633 self.fragment = 0..0;
634 self.fed = 0;
635 self.splitter.clear();
636 }
637
638 pub fn reset(&mut self) {
640 self.seek(0);
641 }
642}
643
644pub trait OutputFormat: Send + Sync {
645 fn name(&self) -> Cow<'static, str>;
647
648 fn config_from_http_request(
653 &self,
654 endpoint_name: &str,
655 request: &HttpRequest,
656 ) -> Result<Box<dyn ErasedSerialize>, ControllerError>;
657
658 fn new_encoder(
673 &self,
674 endpoint_name: &str,
675 config: &ConnectorConfig,
676 key_schema: &Option<Relation>,
677 value_schema: &Relation,
678 consumer: Box<dyn OutputConsumer>,
679 is_index: bool,
680 ) -> Result<Box<dyn Encoder>, ControllerError>;
681}
682
683pub trait Encoder: Send {
684 fn consumer(&mut self) -> &mut dyn OutputConsumer;
686
687 fn encode(&mut self, batch: Arc<dyn SerBatchReader>) -> AnyResult<()>;
696}
697
698#[doc(hidden)]
699pub trait OutputConsumer: Send {
700 fn max_buffer_size_bytes(&self) -> usize;
703
704 fn batch_start(&mut self, step: Step, batch_type: OutputBatchType);
705
706 fn push_buffer(&mut self, buffer: &[u8], num_records: usize);
708
709 fn push_key(
711 &mut self,
712 key: Option<&[u8]>,
713 val: Option<&[u8]>,
714 headers: &[(&str, Option<&[u8]>)],
715 num_records: usize,
716 );
717 fn batch_end(&mut self);
718
719 fn memory(&self) -> usize {
724 0
725 }
726}
727
728pub type PostprocessorErrorCallback = Box<dyn Fn(anyhow::Error) + Send + Sync>;
734
735pub struct PostprocessedConsumer {
742 inner: Box<dyn OutputConsumer>,
743 postprocessor: Box<dyn Postprocessor>,
744 error_cb: PostprocessorErrorCallback,
745}
746
747impl PostprocessedConsumer {
748 pub fn new(
749 inner: Box<dyn OutputConsumer>,
750 postprocessor: Box<dyn Postprocessor>,
751 error_cb: PostprocessorErrorCallback,
752 ) -> Self {
753 Self {
754 inner,
755 postprocessor,
756 error_cb,
757 }
758 }
759}
760
761impl OutputConsumer for PostprocessedConsumer {
762 fn max_buffer_size_bytes(&self) -> usize {
763 self.inner.max_buffer_size_bytes()
764 }
765
766 fn batch_start(&mut self, step: Step, batch_type: OutputBatchType) {
767 self.postprocessor.batch_start(step, batch_type);
768 self.inner.batch_start(step, batch_type);
769 }
770
771 fn push_buffer(&mut self, buffer: &[u8], num_records: usize) {
772 match self.postprocessor.push_buffer(buffer) {
773 Ok(transformed) => self.inner.push_buffer(&transformed, num_records),
774 Err(e) => (self.error_cb)(e),
775 }
776 }
777
778 fn push_key(
779 &mut self,
780 key: Option<&[u8]>,
781 val: Option<&[u8]>,
782 headers: &[(&str, Option<&[u8]>)],
783 num_records: usize,
784 ) {
785 match self.postprocessor.push_key(key, val, headers) {
786 Ok((k, v, h)) => {
787 let h_refs: Vec<(&str, Option<&[u8]>)> =
788 h.iter().map(|(k, v)| (k.as_str(), v.as_deref())).collect();
789 self.inner
790 .push_key(k.as_deref(), v.as_deref(), &h_refs, num_records);
791 }
792 Err(e) => (self.error_cb)(e),
793 }
794 }
795
796 fn batch_end(&mut self) {
797 self.postprocessor.batch_end();
798 self.inner.batch_end();
799 }
800
801 fn memory(&self) -> usize {
802 self.inner.memory() + self.postprocessor.memory()
803 }
804}
805
806pub const MAX_DUPLICATES: i64 = 1_000_000;
812
813pub const MAX_RECORD_LEN_IN_ERRMSG: usize = 4096;
816
817#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
819#[serde(transparent)]
820pub struct ParseError(Box<ParseErrorInner>);
823impl Display for ParseError {
824 fn fmt(&self, f: &mut Formatter<'_>) -> Result<(), FmtError> {
825 self.0.fmt(f)
826 }
827}
828
829impl StdError for ParseError {}
830
831impl ParseError {
832 pub fn new(
833 description: String,
834 event_number: Option<u64>,
835 field: Option<String>,
836 invalid_text: Option<&str>,
837 invalid_bytes: Option<&[u8]>,
838 suggestion: Option<Cow<'static, str>>,
839 error_tag: Option<String>,
840 ) -> Self {
841 Self(Box::new(ParseErrorInner::new(
842 description,
843 event_number,
844 field,
845 invalid_text,
846 invalid_bytes,
847 suggestion,
848 error_tag,
849 )))
850 }
851
852 pub fn text_event_error<E>(
853 msg: &str,
854 error: E,
855 event_number: u64,
856 invalid_text: Option<&str>,
857 suggestion: Option<Cow<'static, str>>,
858 ) -> Self
859 where
860 E: ToString,
861 {
862 Self(Box::new(ParseErrorInner::text_event_error(
863 msg,
864 error,
865 event_number,
866 invalid_text,
867 suggestion,
868 )))
869 }
870
871 pub fn text_envelope_error(
872 description: String,
873 invalid_text: &str,
874 suggestion: Option<Cow<'static, str>>,
875 ) -> Self {
876 Self(Box::new(ParseErrorInner::text_envelope_error(
877 description,
878 invalid_text,
879 suggestion,
880 )))
881 }
882
883 pub fn bin_event_error(
884 description: String,
885 event_number: u64,
886 invalid_bytes: &[u8],
887 suggestion: Option<Cow<'static, str>>,
888 ) -> Self {
889 Self(Box::new(ParseErrorInner::bin_event_error(
890 description,
891 event_number,
892 invalid_bytes,
893 suggestion,
894 )))
895 }
896
897 pub fn bin_envelope_error(
898 description: String,
899 invalid_bytes: &[u8],
900 suggestion: Option<Cow<'static, str>>,
901 ) -> Self {
902 Self(Box::new(ParseErrorInner::bin_envelope_error(
903 description,
904 invalid_bytes,
905 suggestion,
906 )))
907 }
908
909 pub fn map_description<F>(self, f: F) -> Self
913 where
914 F: FnOnce(&str) -> String,
915 {
916 let mut inner = self.0;
917 let description = f(&inner.description);
918 inner.description = description;
919 Self(inner)
920 }
921
922 pub fn get_error_tag(&self) -> Option<String> {
923 self.0.get_error_tag()
924 }
925}
926
927#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
928pub struct ParseErrorInner {
929 description: String,
931
932 event_number: Option<u64>,
940
941 field: Option<String>,
946
947 invalid_bytes: Option<Vec<u8>>,
952
953 invalid_text: Option<String>,
957
958 suggestion: Option<Cow<'static, str>>,
961
962 tag: Option<String>,
964}
965
966impl Display for ParseErrorInner {
967 fn fmt(&self, f: &mut Formatter<'_>) -> Result<(), FmtError> {
968 let event = if let Some(event_number) = self.event_number {
969 format!(" (event #{})", event_number)
970 } else {
971 String::new()
972 };
973
974 let invalid_fragment = if let Some(invalid_bytes) = &self.invalid_bytes {
975 format!("\nInvalid bytes: {invalid_bytes:?}")
976 } else if let Some(invalid_text) = &self.invalid_text {
977 format!("\nInvalid fragment: '{invalid_text}'")
978 } else {
979 String::new()
980 };
981
982 let suggestion = if let Some(suggestion) = &self.suggestion {
983 format!("\n{suggestion}")
984 } else {
985 String::new()
986 };
987
988 write!(
989 f,
990 "Parse error{event}: {}{invalid_fragment}{suggestion}",
991 self.description
992 )
993 }
994}
995
996impl ParseErrorInner {
997 pub fn new(
998 description: String,
999 event_number: Option<u64>,
1000 field: Option<String>,
1001 invalid_text: Option<&str>,
1002 invalid_bytes: Option<&[u8]>,
1003 suggestion: Option<Cow<'static, str>>,
1004 error_tag: Option<String>,
1005 ) -> Self {
1006 Self {
1007 description,
1008 event_number,
1009 field,
1010 invalid_text: invalid_text.map(str::to_string),
1011 invalid_bytes: invalid_bytes.map(ToOwned::to_owned),
1012 suggestion,
1013 tag: error_tag,
1014 }
1015 }
1016
1017 pub fn text_event_error<E>(
1020 msg: &str,
1021 error: E,
1022 event_number: u64,
1023 invalid_text: Option<&str>,
1024 suggestion: Option<Cow<'static, str>>,
1025 ) -> Self
1026 where
1027 E: ToString,
1028 {
1029 let err_str = error.to_string();
1030 let (descr, field) = if let Some(offset) = err_str.find("{\"field\":") {
1034 if let Some(Ok(err)) = serde_json::Deserializer::from_str(&err_str[offset..])
1035 .into_iter::<FieldParseError>()
1036 .next()
1037 {
1038 (err.description, Some(err.field))
1039 } else {
1040 (err_str, None)
1041 }
1042 } else {
1043 (err_str, None)
1044 };
1045 let column_name = if let Some(field) = &field {
1046 format!(": error parsing field '{field}'")
1047 } else {
1048 String::new()
1049 };
1050
1051 Self::new(
1052 format!("{msg}{column_name}: {descr}",),
1053 Some(event_number),
1054 field,
1055 invalid_text,
1056 None,
1057 suggestion,
1058 Some("text_event_err".to_string()),
1059 )
1060 }
1061
1062 pub fn text_envelope_error(
1066 description: String,
1067 invalid_text: &str,
1068 suggestion: Option<Cow<'static, str>>,
1069 ) -> Self {
1070 Self::new(
1071 description,
1072 None,
1073 None,
1074 Some(invalid_text),
1075 None,
1076 suggestion,
1077 Some("text_envelope_err".to_string()),
1078 )
1079 }
1080
1081 pub fn bin_event_error(
1084 description: String,
1085 event_number: u64,
1086 invalid_bytes: &[u8],
1087 suggestion: Option<Cow<'static, str>>,
1088 ) -> Self {
1089 Self::new(
1090 description,
1091 Some(event_number),
1092 None,
1093 None,
1094 Some(invalid_bytes),
1095 suggestion,
1096 Some("bin_event_err".to_string()),
1097 )
1098 }
1099
1100 pub fn bin_envelope_error(
1104 description: String,
1105 invalid_bytes: &[u8],
1106 suggestion: Option<Cow<'static, str>>,
1107 ) -> Self {
1108 Self::new(
1109 description,
1110 None,
1111 None,
1112 None,
1113 Some(invalid_bytes),
1114 suggestion,
1115 Some("bin_envelope_err".to_string()),
1116 )
1117 }
1118
1119 pub fn get_error_tag(&self) -> Option<String> {
1120 self.tag.clone()
1121 }
1122}