1#[derive(prost_helpers::AnyPB)]
3#[derive(Clone, PartialEq, ::prost::Message)]
4pub struct Dispatchers {
5 #[prost(message, repeated, tag = "1")]
6 pub dispatchers: ::prost::alloc::vec::Vec<Dispatcher>,
7}
8#[derive(prost_helpers::AnyPB)]
9#[derive(Clone, PartialEq, ::prost::Message)]
10pub struct UpstreamSinkInfo {
11 #[prost(uint32, tag = "1", wrapper = "crate::id::FragmentId")]
12 pub upstream_fragment_id: crate::id::FragmentId,
13 #[prost(message, repeated, tag = "2")]
14 pub sink_output_schema: ::prost::alloc::vec::Vec<super::plan_common::Field>,
15 #[prost(message, repeated, tag = "3")]
16 pub project_exprs: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
17}
18#[derive(prost_helpers::AnyPB)]
19#[derive(Clone, PartialEq, ::prost::Message)]
20pub struct AddMutation {
21 #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::ActorId")]
23 pub actor_dispatchers: ::std::collections::HashMap<crate::id::ActorId, Dispatchers>,
24 #[prost(uint32, repeated, tag = "3", wrapper = "crate::id::ActorId")]
26 pub added_actors: ::prost::alloc::vec::Vec<crate::id::ActorId>,
27 #[prost(map = "uint32, message", tag = "2", wrapper = "crate::id::ActorId")]
31 pub actor_splits: ::std::collections::HashMap<
32 crate::id::ActorId,
33 super::source::ConnectorSplits,
34 >,
35 #[prost(bool, tag = "4")]
38 pub pause: bool,
39 #[prost(message, repeated, tag = "5")]
40 pub subscriptions_to_add: ::prost::alloc::vec::Vec<SubscriptionUpstreamInfo>,
41 #[prost(uint32, repeated, tag = "6", wrapper = "crate::id::FragmentId")]
43 pub backfill_nodes_to_pause: ::prost::alloc::vec::Vec<crate::id::FragmentId>,
44 #[prost(message, optional, tag = "7")]
46 pub actor_cdc_table_snapshot_splits: ::core::option::Option<
47 super::source::CdcTableSnapshotSplitsWithGeneration,
48 >,
49 #[prost(map = "uint32, message", tag = "8", wrapper = "crate::id::FragmentId")]
51 pub new_upstream_sinks: ::std::collections::HashMap<
52 crate::id::FragmentId,
53 add_mutation::NewUpstreamSink,
54 >,
55 #[prost(uint32, repeated, tag = "9", wrapper = "crate::id::ActorId")]
57 pub dropped_actors: ::prost::alloc::vec::Vec<crate::id::ActorId>,
58 #[prost(uint32, repeated, tag = "10", wrapper = "crate::id::SinkId")]
60 pub sink_log_store_flush: ::prost::alloc::vec::Vec<crate::id::SinkId>,
61}
62pub mod add_mutation {
64 #[derive(prost_helpers::AnyPB)]
65 #[derive(Clone, PartialEq, ::prost::Message)]
66 pub struct NewUpstreamSink {
67 #[prost(message, optional, tag = "1")]
68 pub info: ::core::option::Option<super::UpstreamSinkInfo>,
69 #[prost(message, repeated, tag = "2")]
70 pub upstream_actors: ::prost::alloc::vec::Vec<super::super::common::ActorInfo>,
71 }
72}
73#[derive(prost_helpers::AnyPB)]
74#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
75pub struct StopMutation {
76 #[prost(uint32, repeated, tag = "1", wrapper = "crate::id::ActorId")]
77 pub actors: ::prost::alloc::vec::Vec<crate::id::ActorId>,
78 #[prost(uint32, repeated, tag = "2", wrapper = "crate::id::FragmentId")]
80 pub dropped_sink_fragments: ::prost::alloc::vec::Vec<crate::id::FragmentId>,
81}
82#[derive(prost_helpers::AnyPB)]
83#[derive(Clone, PartialEq, ::prost::Message)]
84pub struct UpdateMutation {
85 #[prost(message, repeated, tag = "1")]
87 pub dispatcher_update: ::prost::alloc::vec::Vec<update_mutation::DispatcherUpdate>,
88 #[prost(message, repeated, tag = "2")]
90 pub merge_update: ::prost::alloc::vec::Vec<update_mutation::MergeUpdate>,
91 #[prost(map = "uint32, message", tag = "3", wrapper = "crate::id::ActorId")]
93 pub actor_vnode_bitmap_update: ::std::collections::HashMap<
94 crate::id::ActorId,
95 super::common::Buffer,
96 >,
97 #[prost(uint32, repeated, tag = "4", wrapper = "crate::id::ActorId")]
99 pub dropped_actors: ::prost::alloc::vec::Vec<crate::id::ActorId>,
100 #[prost(map = "uint32, message", tag = "5", wrapper = "crate::id::ActorId")]
103 pub actor_splits: ::std::collections::HashMap<
104 crate::id::ActorId,
105 super::source::ConnectorSplits,
106 >,
107 #[prost(map = "uint32, message", tag = "6", wrapper = "crate::id::ActorId")]
110 pub actor_new_dispatchers: ::std::collections::HashMap<
111 crate::id::ActorId,
112 Dispatchers,
113 >,
114 #[prost(message, optional, tag = "7")]
116 pub actor_cdc_table_snapshot_splits: ::core::option::Option<
117 super::source::CdcTableSnapshotSplitsWithGeneration,
118 >,
119 #[prost(map = "uint32, message", tag = "8")]
120 pub sink_schema_change: ::std::collections::HashMap<u32, SinkSchemaChange>,
121 #[prost(message, repeated, tag = "9")]
122 pub subscriptions_to_drop: ::prost::alloc::vec::Vec<SubscriptionUpstreamInfo>,
123 #[prost(message, optional, tag = "10")]
124 pub iceberg_pk_index_compaction: ::core::option::Option<
125 IcebergPkIndexCompactionContext,
126 >,
127}
128pub mod update_mutation {
130 #[derive(prost_helpers::AnyPB)]
131 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
132 pub struct DispatcherUpdate {
133 #[prost(uint32, tag = "1", wrapper = "crate::id::ActorId")]
135 pub actor_id: crate::id::ActorId,
136 #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
137 pub dispatcher_id: crate::id::FragmentId,
138 #[prost(message, optional, tag = "3")]
141 pub hash_mapping: ::core::option::Option<super::ActorMapping>,
142 #[prost(uint32, repeated, tag = "4", wrapper = "crate::id::ActorId")]
144 pub added_downstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
145 #[prost(uint32, repeated, tag = "5", wrapper = "crate::id::ActorId")]
147 pub removed_downstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
148 }
149 #[derive(prost_helpers::AnyPB)]
150 #[derive(Clone, PartialEq, ::prost::Message)]
151 pub struct MergeUpdate {
152 #[prost(uint32, tag = "1", wrapper = "crate::id::ActorId")]
154 pub actor_id: crate::id::ActorId,
155 #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
156 pub upstream_fragment_id: crate::id::FragmentId,
157 #[prost(uint32, optional, tag = "5", wrapper = "crate::id::FragmentId")]
161 pub new_upstream_fragment_id: ::core::option::Option<crate::id::FragmentId>,
162 #[prost(message, repeated, tag = "3")]
164 pub added_upstream_actors: ::prost::alloc::vec::Vec<
165 super::super::common::ActorInfo,
166 >,
167 #[prost(uint32, repeated, tag = "4", wrapper = "crate::id::ActorId")]
170 pub removed_upstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
171 }
172}
173#[derive(prost_helpers::AnyPB)]
174#[derive(Clone, PartialEq, ::prost::Message)]
175pub struct SourceChangeSplitMutation {
176 #[prost(map = "uint32, message", tag = "2", wrapper = "crate::id::ActorId")]
178 pub actor_splits: ::std::collections::HashMap<
179 crate::id::ActorId,
180 super::source::ConnectorSplits,
181 >,
182}
183#[derive(prost_helpers::AnyPB)]
184#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
185pub struct PauseMutation {}
186#[derive(prost_helpers::AnyPB)]
187#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
188pub struct ResumeMutation {}
189#[derive(prost_helpers::AnyPB)]
190#[derive(Clone, PartialEq, ::prost::Message)]
191pub struct ThrottleMutation {
192 #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::FragmentId")]
193 pub fragment_throttle: ::std::collections::HashMap<
194 crate::id::FragmentId,
195 throttle_mutation::ThrottleConfig,
196 >,
197}
198pub mod throttle_mutation {
200 #[derive(prost_helpers::AnyPB)]
201 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
202 pub struct ThrottleConfig {
203 #[prost(uint32, optional, tag = "1")]
204 pub rate_limit: ::core::option::Option<u32>,
205 #[prost(enumeration = "super::super::common::ThrottleType", tag = "2")]
206 pub throttle_type: i32,
207 }
208}
209#[derive(prost_helpers::AnyPB)]
210#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
211pub struct SubscriptionUpstreamInfo {
212 #[prost(uint32, tag = "1", wrapper = "crate::id::SubscriberId")]
214 pub subscriber_id: crate::id::SubscriberId,
215 #[prost(uint32, tag = "2", wrapper = "crate::id::TableId")]
216 pub upstream_mv_table_id: crate::id::TableId,
217}
218#[derive(prost_helpers::AnyPB)]
219#[derive(Clone, PartialEq, ::prost::Message)]
220pub struct DropSubscriptionsMutation {
221 #[prost(message, repeated, tag = "1")]
222 pub info: ::prost::alloc::vec::Vec<SubscriptionUpstreamInfo>,
223}
224#[derive(prost_helpers::AnyPB)]
225#[derive(Clone, PartialEq, ::prost::Message)]
226pub struct ConnectorPropsChangeMutation {
227 #[prost(map = "uint32, message", tag = "1")]
228 pub connector_props_infos: ::std::collections::HashMap<
229 u32,
230 connector_props_change_mutation::ConnectorPropsInfo,
231 >,
232}
233pub mod connector_props_change_mutation {
235 #[derive(prost_helpers::AnyPB)]
236 #[derive(Clone, PartialEq, ::prost::Message)]
237 pub struct ConnectorPropsInfo {
238 #[prost(map = "string, string", tag = "1")]
239 pub connector_props_info: ::std::collections::HashMap<
240 ::prost::alloc::string::String,
241 ::prost::alloc::string::String,
242 >,
243 }
244}
245#[derive(prost_helpers::AnyPB)]
246#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
247pub struct StartFragmentBackfillMutation {
248 #[prost(uint32, repeated, tag = "1", wrapper = "crate::id::FragmentId")]
249 pub fragment_ids: ::prost::alloc::vec::Vec<crate::id::FragmentId>,
250}
251#[derive(prost_helpers::AnyPB)]
252#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
253pub struct RefreshStartMutation {
254 #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
256 pub table_id: crate::id::TableId,
257 #[prost(uint32, tag = "2", wrapper = "crate::id::SourceId")]
259 pub associated_source_id: crate::id::SourceId,
260}
261#[derive(prost_helpers::AnyPB)]
262#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
263pub struct ListFinishMutation {
264 #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
266 pub associated_source_id: crate::id::SourceId,
267}
268#[derive(prost_helpers::AnyPB)]
269#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
270pub struct LoadFinishMutation {
271 #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
273 pub associated_source_id: crate::id::SourceId,
274}
275#[derive(prost_helpers::AnyPB)]
276#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
277pub struct ResetSourceMutation {
278 #[prost(uint32, tag = "1")]
280 pub source_id: u32,
281}
282#[derive(prost_helpers::AnyPB)]
285#[derive(Clone, PartialEq, ::prost::Message)]
286pub struct InjectSourceOffsetsMutation {
287 #[prost(uint32, tag = "1")]
289 pub source_id: u32,
290 #[prost(map = "string, string", tag = "2")]
292 pub split_offsets: ::std::collections::HashMap<
293 ::prost::alloc::string::String,
294 ::prost::alloc::string::String,
295 >,
296}
297#[derive(prost_helpers::AnyPB)]
299#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
300pub struct IcebergPkIndexCompactionContext {
301 #[prost(uint32, tag = "1", wrapper = "crate::id::SinkId")]
302 pub sink_id: crate::id::SinkId,
303 #[prost(uint64, tag = "2", wrapper = "crate::id::IcebergCompactionTaskId")]
304 pub task_id: crate::id::IcebergCompactionTaskId,
305 #[prost(enumeration = "iceberg_pk_index_compaction_context::Phase", tag = "3")]
306 pub phase: i32,
307 #[prost(message, optional, tag = "4")]
309 pub resolver_task_input: ::core::option::Option<
310 iceberg_pk_index_compaction_context::ResolverTaskInput,
311 >,
312}
313pub mod iceberg_pk_index_compaction_context {
315 #[derive(prost_helpers::AnyPB)]
316 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
317 pub struct ResolverTaskInput {
318 #[prost(string, repeated, tag = "1")]
319 pub output_data_file_paths: ::prost::alloc::vec::Vec<
320 ::prost::alloc::string::String,
321 >,
322 #[prost(string, repeated, tag = "2")]
323 pub input_data_file_paths: ::prost::alloc::vec::Vec<
324 ::prost::alloc::string::String,
325 >,
326 #[prost(int64, tag = "3")]
327 pub read_snapshot_id: i64,
328 }
329 #[derive(prost_helpers::AnyPB)]
330 #[derive(
331 Clone,
332 Copy,
333 Debug,
334 PartialEq,
335 Eq,
336 Hash,
337 PartialOrd,
338 Ord,
339 ::prost::Enumeration
340 )]
341 #[repr(i32)]
342 pub enum Phase {
343 Unspecified = 0,
344 Begin = 1,
345 End = 2,
346 }
347 impl Phase {
348 pub fn as_str_name(&self) -> &'static str {
353 match self {
354 Self::Unspecified => "PHASE_UNSPECIFIED",
355 Self::Begin => "PHASE_BEGIN",
356 Self::End => "PHASE_END",
357 }
358 }
359 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
361 match value {
362 "PHASE_UNSPECIFIED" => Some(Self::Unspecified),
363 "PHASE_BEGIN" => Some(Self::Begin),
364 "PHASE_END" => Some(Self::End),
365 _ => None,
366 }
367 }
368 }
369}
370#[derive(prost_helpers::AnyPB)]
371#[derive(Clone, PartialEq, ::prost::Message)]
372pub struct BarrierMutation {
373 #[prost(
374 oneof = "barrier_mutation::Mutation",
375 tags = "3, 4, 5, 6, 7, 8, 10, 12, 13, 14, 15, 16, 17, 18, 19"
376 )]
377 pub mutation: ::core::option::Option<barrier_mutation::Mutation>,
378}
379pub mod barrier_mutation {
381 #[derive(prost_helpers::AnyPB)]
382 #[derive(Clone, PartialEq, ::prost::Oneof)]
383 pub enum Mutation {
384 #[prost(message, tag = "3")]
386 Add(super::AddMutation),
387 #[prost(message, tag = "4")]
390 Stop(super::StopMutation),
391 #[prost(message, tag = "5")]
393 Update(super::UpdateMutation),
394 #[prost(message, tag = "6")]
396 Splits(super::SourceChangeSplitMutation),
397 #[prost(message, tag = "7")]
399 Pause(super::PauseMutation),
400 #[prost(message, tag = "8")]
402 Resume(super::ResumeMutation),
403 #[prost(message, tag = "10")]
405 Throttle(super::ThrottleMutation),
406 #[prost(message, tag = "12")]
408 DropSubscriptions(super::DropSubscriptionsMutation),
409 #[prost(message, tag = "13")]
411 ConnectorPropsChange(super::ConnectorPropsChangeMutation),
412 #[prost(message, tag = "14")]
418 StartFragmentBackfill(super::StartFragmentBackfillMutation),
419 #[prost(message, tag = "15")]
421 RefreshStart(super::RefreshStartMutation),
422 #[prost(message, tag = "16")]
424 LoadFinish(super::LoadFinishMutation),
425 #[prost(message, tag = "17")]
427 ListFinish(super::ListFinishMutation),
428 #[prost(message, tag = "18")]
430 ResetSource(super::ResetSourceMutation),
431 #[prost(message, tag = "19")]
433 InjectSourceOffsets(super::InjectSourceOffsetsMutation),
434 }
435}
436#[derive(prost_helpers::AnyPB)]
437#[derive(Clone, PartialEq, ::prost::Message)]
438pub struct Barrier {
439 #[prost(message, optional, tag = "1")]
440 pub epoch: ::core::option::Option<super::data::Epoch>,
441 #[prost(message, optional, tag = "3")]
442 pub mutation: ::core::option::Option<BarrierMutation>,
443 #[prost(map = "string, string", tag = "2")]
445 pub tracing_context: ::std::collections::HashMap<
446 ::prost::alloc::string::String,
447 ::prost::alloc::string::String,
448 >,
449 #[prost(enumeration = "barrier::BarrierKind", tag = "9")]
451 pub kind: i32,
452}
453pub mod barrier {
455 #[derive(prost_helpers::AnyPB)]
456 #[derive(::enum_as_inner::EnumAsInner)]
457 #[derive(
458 Clone,
459 Copy,
460 Debug,
461 PartialEq,
462 Eq,
463 Hash,
464 PartialOrd,
465 Ord,
466 ::prost::Enumeration
467 )]
468 #[repr(i32)]
469 pub enum BarrierKind {
470 Unspecified = 0,
471 Initial = 1,
474 Barrier = 2,
476 Checkpoint = 3,
478 }
479 impl BarrierKind {
480 pub fn as_str_name(&self) -> &'static str {
485 match self {
486 Self::Unspecified => "BARRIER_KIND_UNSPECIFIED",
487 Self::Initial => "BARRIER_KIND_INITIAL",
488 Self::Barrier => "BARRIER_KIND_BARRIER",
489 Self::Checkpoint => "BARRIER_KIND_CHECKPOINT",
490 }
491 }
492 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
494 match value {
495 "BARRIER_KIND_UNSPECIFIED" => Some(Self::Unspecified),
496 "BARRIER_KIND_INITIAL" => Some(Self::Initial),
497 "BARRIER_KIND_BARRIER" => Some(Self::Barrier),
498 "BARRIER_KIND_CHECKPOINT" => Some(Self::Checkpoint),
499 _ => None,
500 }
501 }
502 }
503}
504#[derive(prost_helpers::AnyPB)]
505#[derive(Clone, PartialEq, ::prost::Message)]
506pub struct Watermark {
507 #[prost(message, optional, tag = "1")]
509 pub column: ::core::option::Option<super::expr::InputRef>,
510 #[prost(message, optional, tag = "3")]
512 pub val: ::core::option::Option<super::data::Datum>,
513}
514#[derive(prost_helpers::AnyPB)]
515#[derive(Clone, PartialEq, ::prost::Message)]
516pub struct StreamMessage {
517 #[prost(oneof = "stream_message::StreamMessage", tags = "1, 2, 3")]
518 pub stream_message: ::core::option::Option<stream_message::StreamMessage>,
519}
520pub mod stream_message {
522 #[derive(prost_helpers::AnyPB)]
523 #[derive(Clone, PartialEq, ::prost::Oneof)]
524 pub enum StreamMessage {
525 #[prost(message, tag = "1")]
526 StreamChunk(super::super::data::StreamChunk),
527 #[prost(message, tag = "2")]
528 Barrier(super::Barrier),
529 #[prost(message, tag = "3")]
530 Watermark(super::Watermark),
531 }
532}
533#[derive(prost_helpers::AnyPB)]
534#[derive(Clone, PartialEq, ::prost::Message)]
535pub struct StreamMessageBatch {
536 #[prost(oneof = "stream_message_batch::StreamMessageBatch", tags = "1, 2, 3")]
537 pub stream_message_batch: ::core::option::Option<
538 stream_message_batch::StreamMessageBatch,
539 >,
540}
541pub mod stream_message_batch {
543 #[derive(prost_helpers::AnyPB)]
544 #[derive(Clone, PartialEq, ::prost::Message)]
545 pub struct BarrierBatch {
546 #[prost(message, repeated, tag = "1")]
547 pub barriers: ::prost::alloc::vec::Vec<super::Barrier>,
548 }
549 #[derive(prost_helpers::AnyPB)]
550 #[derive(Clone, PartialEq, ::prost::Oneof)]
551 pub enum StreamMessageBatch {
552 #[prost(message, tag = "1")]
553 StreamChunk(super::super::data::StreamChunk),
554 #[prost(message, tag = "2")]
555 BarrierBatch(BarrierBatch),
556 #[prost(message, tag = "3")]
557 Watermark(super::Watermark),
558 }
559}
560#[derive(prost_helpers::AnyPB)]
562#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
563pub struct ActorMapping {
564 #[prost(uint32, repeated, tag = "1")]
565 pub original_indices: ::prost::alloc::vec::Vec<u32>,
566 #[prost(uint32, repeated, tag = "2", wrapper = "crate::id::ActorId")]
567 pub data: ::prost::alloc::vec::Vec<crate::id::ActorId>,
568}
569#[derive(prost_helpers::AnyPB)]
570#[derive(Clone, PartialEq, ::prost::Message)]
571pub struct Columns {
572 #[prost(message, repeated, tag = "1")]
573 pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
574}
575#[derive(prost_helpers::AnyPB)]
576#[derive(Clone, PartialEq, ::prost::Message)]
577pub struct StreamSource {
578 #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
579 pub source_id: crate::id::SourceId,
580 #[prost(message, optional, tag = "2")]
581 pub state_table: ::core::option::Option<super::catalog::Table>,
582 #[prost(uint32, optional, tag = "3")]
583 pub row_id_index: ::core::option::Option<u32>,
584 #[prost(message, repeated, tag = "4")]
585 pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
586 #[prost(btree_map = "string, string", tag = "6")]
587 pub with_properties: ::prost::alloc::collections::BTreeMap<
588 ::prost::alloc::string::String,
589 ::prost::alloc::string::String,
590 >,
591 #[prost(message, optional, tag = "7")]
592 pub info: ::core::option::Option<super::catalog::StreamSourceInfo>,
593 #[prost(string, tag = "8")]
594 pub source_name: ::prost::alloc::string::String,
595 #[prost(uint32, optional, tag = "9")]
597 pub rate_limit: ::core::option::Option<u32>,
598 #[prost(btree_map = "string, message", tag = "10")]
599 pub secret_refs: ::prost::alloc::collections::BTreeMap<
600 ::prost::alloc::string::String,
601 super::secret::SecretRef,
602 >,
603 #[prost(message, optional, tag = "11")]
605 pub downstream_columns: ::core::option::Option<Columns>,
606 #[prost(message, optional, tag = "12")]
607 pub refresh_mode: ::core::option::Option<super::plan_common::SourceRefreshMode>,
608 #[prost(uint32, optional, tag = "13", wrapper = "crate::id::TableId")]
609 pub associated_table_id: ::core::option::Option<crate::id::TableId>,
610}
611#[derive(prost_helpers::AnyPB)]
613#[derive(Clone, PartialEq, ::prost::Message)]
614pub struct StreamFsFetch {
615 #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
616 pub source_id: crate::id::SourceId,
617 #[prost(message, optional, tag = "2")]
618 pub state_table: ::core::option::Option<super::catalog::Table>,
619 #[prost(uint32, optional, tag = "3")]
620 pub row_id_index: ::core::option::Option<u32>,
621 #[prost(message, repeated, tag = "4")]
622 pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
623 #[prost(btree_map = "string, string", tag = "6")]
624 pub with_properties: ::prost::alloc::collections::BTreeMap<
625 ::prost::alloc::string::String,
626 ::prost::alloc::string::String,
627 >,
628 #[prost(message, optional, tag = "7")]
629 pub info: ::core::option::Option<super::catalog::StreamSourceInfo>,
630 #[prost(string, tag = "8")]
631 pub source_name: ::prost::alloc::string::String,
632 #[prost(uint32, optional, tag = "9")]
634 pub rate_limit: ::core::option::Option<u32>,
635 #[prost(btree_map = "string, message", tag = "10")]
636 pub secret_refs: ::prost::alloc::collections::BTreeMap<
637 ::prost::alloc::string::String,
638 super::secret::SecretRef,
639 >,
640 #[prost(message, optional, tag = "11")]
641 pub refresh_mode: ::core::option::Option<super::plan_common::SourceRefreshMode>,
642 #[prost(uint32, optional, tag = "12", wrapper = "crate::id::TableId")]
643 pub associated_table_id: ::core::option::Option<crate::id::TableId>,
644}
645#[derive(prost_helpers::AnyPB)]
648#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
649pub struct BarrierRecvNode {}
650#[derive(prost_helpers::AnyPB)]
651#[derive(Clone, PartialEq, ::prost::Message)]
652pub struct SourceNode {
653 #[prost(message, optional, tag = "1")]
656 pub source_inner: ::core::option::Option<StreamSource>,
657}
658#[derive(prost_helpers::AnyPB)]
659#[derive(Clone, PartialEq, ::prost::Message)]
660pub struct StreamFsFetchNode {
661 #[prost(message, optional, tag = "1")]
662 pub node_inner: ::core::option::Option<StreamFsFetch>,
663}
664#[derive(prost_helpers::AnyPB)]
667#[derive(Clone, PartialEq, ::prost::Message)]
668pub struct SourceBackfillNode {
669 #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
670 pub upstream_source_id: crate::id::SourceId,
671 #[prost(uint32, optional, tag = "2")]
672 pub row_id_index: ::core::option::Option<u32>,
673 #[prost(message, repeated, tag = "3")]
674 pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
675 #[prost(message, optional, tag = "4")]
676 pub info: ::core::option::Option<super::catalog::StreamSourceInfo>,
677 #[prost(string, tag = "5")]
678 pub source_name: ::prost::alloc::string::String,
679 #[prost(btree_map = "string, string", tag = "6")]
680 pub with_properties: ::prost::alloc::collections::BTreeMap<
681 ::prost::alloc::string::String,
682 ::prost::alloc::string::String,
683 >,
684 #[prost(uint32, optional, tag = "7")]
686 pub rate_limit: ::core::option::Option<u32>,
687 #[prost(message, optional, tag = "8")]
689 pub state_table: ::core::option::Option<super::catalog::Table>,
690 #[prost(btree_map = "string, message", tag = "9")]
691 pub secret_refs: ::prost::alloc::collections::BTreeMap<
692 ::prost::alloc::string::String,
693 super::secret::SecretRef,
694 >,
695}
696#[derive(prost_helpers::AnyPB)]
697#[derive(Clone, PartialEq, ::prost::Message)]
698pub struct SinkDesc {
699 #[prost(uint32, tag = "1", wrapper = "crate::id::SinkId")]
700 pub id: crate::id::SinkId,
701 #[prost(string, tag = "2")]
702 pub name: ::prost::alloc::string::String,
703 #[prost(string, tag = "3")]
704 pub definition: ::prost::alloc::string::String,
705 #[prost(message, repeated, tag = "5")]
706 pub plan_pk: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
707 #[prost(uint32, repeated, tag = "6")]
708 pub downstream_pk: ::prost::alloc::vec::Vec<u32>,
709 #[prost(uint32, repeated, tag = "7")]
710 pub distribution_key: ::prost::alloc::vec::Vec<u32>,
711 #[prost(btree_map = "string, string", tag = "8")]
712 pub properties: ::prost::alloc::collections::BTreeMap<
713 ::prost::alloc::string::String,
714 ::prost::alloc::string::String,
715 >,
716 #[prost(enumeration = "super::catalog::SinkType", tag = "9")]
718 pub sink_type: i32,
719 #[prost(message, repeated, tag = "10")]
720 pub column_catalogs: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
721 #[prost(string, tag = "11")]
722 pub db_name: ::prost::alloc::string::String,
723 #[prost(string, tag = "12")]
726 pub sink_from_name: ::prost::alloc::string::String,
727 #[prost(message, optional, tag = "13")]
728 pub format_desc: ::core::option::Option<super::catalog::SinkFormatDesc>,
729 #[prost(uint32, optional, tag = "14")]
730 pub target_table: ::core::option::Option<u32>,
731 #[prost(uint64, optional, tag = "15")]
732 pub extra_partition_col_idx: ::core::option::Option<u64>,
733 #[prost(btree_map = "string, message", tag = "16")]
734 pub secret_refs: ::prost::alloc::collections::BTreeMap<
735 ::prost::alloc::string::String,
736 super::secret::SecretRef,
737 >,
738 #[prost(bool, tag = "17")]
742 pub raw_ignore_delete: bool,
743}
744#[derive(prost_helpers::AnyPB)]
745#[derive(Clone, PartialEq, ::prost::Message)]
746pub struct SinkNode {
747 #[prost(message, optional, tag = "1")]
748 pub sink_desc: ::core::option::Option<SinkDesc>,
749 #[prost(message, optional, tag = "2")]
751 pub table: ::core::option::Option<super::catalog::Table>,
752 #[prost(enumeration = "SinkLogStoreType", tag = "3")]
753 pub log_store_type: i32,
754 #[prost(uint32, optional, tag = "4")]
755 pub rate_limit: ::core::option::Option<u32>,
756}
757#[derive(prost_helpers::AnyPB)]
758#[derive(Clone, PartialEq, ::prost::Message)]
759pub struct IcebergWithPkIndexWriterNode {
760 #[prost(message, optional, tag = "1")]
761 pub sink_desc: ::core::option::Option<SinkDesc>,
762 #[prost(message, optional, tag = "2")]
763 pub pk_index_table: ::core::option::Option<super::catalog::Table>,
764}
765#[derive(prost_helpers::AnyPB)]
766#[derive(Clone, PartialEq, ::prost::Message)]
767pub struct IcebergWithPkIndexPositionDeleteMergerNode {
768 #[prost(message, optional, tag = "1")]
769 pub sink_desc: ::core::option::Option<SinkDesc>,
770}
771#[derive(prost_helpers::AnyPB)]
776#[derive(Clone, PartialEq, ::prost::Message)]
777pub struct CompactionResolverNode {
778 #[prost(uint32, tag = "1", wrapper = "crate::id::SinkId")]
780 pub sink_id: crate::id::SinkId,
781 #[prost(btree_map = "string, string", tag = "2")]
782 pub properties: ::prost::alloc::collections::BTreeMap<
783 ::prost::alloc::string::String,
784 ::prost::alloc::string::String,
785 >,
786 #[prost(btree_map = "string, message", tag = "3")]
787 pub secret_refs: ::prost::alloc::collections::BTreeMap<
788 ::prost::alloc::string::String,
789 super::secret::SecretRef,
790 >,
791 #[prost(message, repeated, tag = "4")]
793 pub pk_columns: ::prost::alloc::vec::Vec<compaction_resolver_node::PkColumn>,
794}
795pub mod compaction_resolver_node {
797 #[derive(prost_helpers::AnyPB)]
798 #[derive(Clone, PartialEq, ::prost::Message)]
799 pub struct PkColumn {
800 #[prost(uint32, tag = "1")]
802 pub data_file_index: u32,
803 #[prost(message, optional, tag = "2")]
805 pub column_desc: ::core::option::Option<super::super::plan_common::ColumnDesc>,
806 }
807}
808#[derive(prost_helpers::AnyPB)]
809#[derive(Clone, PartialEq, ::prost::Message)]
810pub struct ProjectNode {
811 #[prost(message, repeated, tag = "1")]
812 pub select_list: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
813 #[prost(uint32, repeated, tag = "2")]
817 pub watermark_input_cols: ::prost::alloc::vec::Vec<u32>,
818 #[prost(uint32, repeated, tag = "3")]
819 pub watermark_output_cols: ::prost::alloc::vec::Vec<u32>,
820 #[prost(uint32, repeated, tag = "4")]
821 pub nondecreasing_exprs: ::prost::alloc::vec::Vec<u32>,
822 #[prost(bool, tag = "5")]
825 pub noop_update_hint: bool,
826}
827#[derive(prost_helpers::AnyPB)]
828#[derive(Clone, PartialEq, ::prost::Message)]
829pub struct FilterNode {
830 #[prost(message, optional, tag = "1")]
831 pub search_condition: ::core::option::Option<super::expr::ExprNode>,
832}
833#[derive(prost_helpers::AnyPB)]
834#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
835pub struct ChangeLogNode {
836 #[prost(bool, tag = "1")]
838 pub need_op: bool,
839 #[prost(uint32, repeated, tag = "2")]
840 pub distribution_keys: ::prost::alloc::vec::Vec<u32>,
841 #[prost(uint32, repeated, tag = "4")]
842 pub stream_keys: ::prost::alloc::vec::Vec<u32>,
843 #[prost(oneof = "change_log_node::Mode", tags = "3")]
844 pub mode: ::core::option::Option<change_log_node::Mode>,
845}
846pub mod change_log_node {
848 #[derive(prost_helpers::AnyPB)]
849 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
850 pub struct Keyed {}
851 #[derive(prost_helpers::AnyPB)]
852 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Oneof)]
853 pub enum Mode {
854 #[prost(message, tag = "3")]
855 Keyed(Keyed),
856 }
857}
858#[derive(prost_helpers::AnyPB)]
859#[derive(Clone, PartialEq, ::prost::Message)]
860pub struct CdcFilterNode {
861 #[prost(message, optional, tag = "1")]
862 pub search_condition: ::core::option::Option<super::expr::ExprNode>,
863 #[prost(uint32, tag = "2", wrapper = "crate::id::SourceId")]
864 pub upstream_source_id: crate::id::SourceId,
865}
866#[derive(prost_helpers::AnyPB)]
875#[derive(Clone, PartialEq, ::prost::Message)]
876pub struct MaterializeNode {
877 #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
878 pub table_id: crate::id::TableId,
879 #[prost(message, repeated, tag = "2")]
881 pub column_orders: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
882 #[prost(message, optional, tag = "3")]
889 pub table: ::core::option::Option<super::catalog::Table>,
890 #[prost(message, optional, tag = "5")]
899 pub staging_table: ::core::option::Option<super::catalog::Table>,
900 #[prost(message, optional, tag = "6")]
913 pub refresh_progress_table: ::core::option::Option<super::catalog::Table>,
914 #[prost(bool, tag = "7")]
916 pub cleaned_by_ttl_watermark: bool,
917}
918#[derive(prost_helpers::AnyPB)]
919#[derive(Clone, PartialEq, ::prost::Message)]
920pub struct AggCallState {
921 #[prost(oneof = "agg_call_state::Inner", tags = "1, 3")]
922 pub inner: ::core::option::Option<agg_call_state::Inner>,
923}
924pub mod agg_call_state {
926 #[derive(prost_helpers::AnyPB)]
928 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
929 pub struct ValueState {}
930 #[derive(prost_helpers::AnyPB)]
932 #[derive(Clone, PartialEq, ::prost::Message)]
933 pub struct MaterializedInputState {
934 #[prost(message, optional, tag = "1")]
935 pub table: ::core::option::Option<super::super::catalog::Table>,
936 #[prost(uint32, repeated, tag = "2")]
938 pub included_upstream_indices: ::prost::alloc::vec::Vec<u32>,
939 #[prost(uint32, repeated, tag = "3")]
940 pub table_value_indices: ::prost::alloc::vec::Vec<u32>,
941 #[prost(message, repeated, tag = "4")]
942 pub order_columns: ::prost::alloc::vec::Vec<super::super::common::ColumnOrder>,
943 }
944 #[derive(prost_helpers::AnyPB)]
945 #[derive(Clone, PartialEq, ::prost::Oneof)]
946 pub enum Inner {
947 #[prost(message, tag = "1")]
948 ValueState(ValueState),
949 #[prost(message, tag = "3")]
950 MaterializedInputState(MaterializedInputState),
951 }
952}
953#[derive(prost_helpers::AnyPB)]
954#[derive(Clone, PartialEq, ::prost::Message)]
955pub struct SimpleAggNode {
956 #[prost(message, repeated, tag = "1")]
957 pub agg_calls: ::prost::alloc::vec::Vec<super::expr::AggCall>,
958 #[prost(message, repeated, tag = "3")]
959 pub agg_call_states: ::prost::alloc::vec::Vec<AggCallState>,
960 #[prost(message, optional, tag = "4")]
961 pub intermediate_state_table: ::core::option::Option<super::catalog::Table>,
962 #[prost(bool, tag = "5")]
965 pub is_append_only: bool,
966 #[prost(map = "uint32, message", tag = "6")]
967 pub distinct_dedup_tables: ::std::collections::HashMap<u32, super::catalog::Table>,
968 #[prost(uint32, tag = "7")]
969 pub row_count_index: u32,
970 #[prost(enumeration = "AggNodeVersion", tag = "8")]
971 pub version: i32,
972 #[prost(bool, tag = "9")]
975 pub must_output_per_barrier: bool,
976}
977#[derive(prost_helpers::AnyPB)]
978#[derive(Clone, PartialEq, ::prost::Message)]
979pub struct HashAggNode {
980 #[prost(uint32, repeated, tag = "1")]
981 pub group_key: ::prost::alloc::vec::Vec<u32>,
982 #[prost(message, repeated, tag = "2")]
983 pub agg_calls: ::prost::alloc::vec::Vec<super::expr::AggCall>,
984 #[prost(message, repeated, tag = "3")]
985 pub agg_call_states: ::prost::alloc::vec::Vec<AggCallState>,
986 #[prost(message, optional, tag = "4")]
987 pub intermediate_state_table: ::core::option::Option<super::catalog::Table>,
988 #[prost(bool, tag = "5")]
991 pub is_append_only: bool,
992 #[prost(map = "uint32, message", tag = "6")]
993 pub distinct_dedup_tables: ::std::collections::HashMap<u32, super::catalog::Table>,
994 #[prost(uint32, tag = "7")]
995 pub row_count_index: u32,
996 #[prost(bool, tag = "8")]
997 pub emit_on_window_close: bool,
998 #[prost(enumeration = "AggNodeVersion", tag = "9")]
999 pub version: i32,
1000}
1001#[derive(prost_helpers::AnyPB)]
1002#[derive(Clone, PartialEq, ::prost::Message)]
1003pub struct TopNNode {
1004 #[prost(uint64, tag = "1")]
1006 pub limit: u64,
1007 #[prost(uint64, tag = "2")]
1008 pub offset: u64,
1009 #[prost(message, optional, tag = "3")]
1010 pub table: ::core::option::Option<super::catalog::Table>,
1011 #[prost(message, repeated, tag = "4")]
1012 pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1013 #[prost(bool, tag = "5")]
1014 pub with_ties: bool,
1015}
1016#[derive(prost_helpers::AnyPB)]
1017#[derive(Clone, PartialEq, ::prost::Message)]
1018pub struct GroupTopNNode {
1019 #[prost(uint64, tag = "1")]
1021 pub limit: u64,
1022 #[prost(uint64, tag = "2")]
1023 pub offset: u64,
1024 #[prost(uint32, repeated, tag = "3")]
1025 pub group_key: ::prost::alloc::vec::Vec<u32>,
1026 #[prost(message, optional, tag = "4")]
1027 pub table: ::core::option::Option<super::catalog::Table>,
1028 #[prost(message, repeated, tag = "5")]
1029 pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1030 #[prost(bool, tag = "6")]
1031 pub with_ties: bool,
1032}
1033#[derive(prost_helpers::AnyPB)]
1034#[derive(Clone, PartialEq, ::prost::Message)]
1035pub struct DeltaExpression {
1036 #[prost(enumeration = "super::expr::expr_node::Type", tag = "1")]
1037 pub delta_type: i32,
1038 #[prost(message, optional, tag = "2")]
1039 pub delta: ::core::option::Option<super::expr::ExprNode>,
1040}
1041#[derive(prost_helpers::AnyPB)]
1043#[derive(Clone, PartialEq, ::prost::Message)]
1044pub struct InequalityPair {
1045 #[prost(uint32, tag = "1")]
1047 pub key_required_larger: u32,
1048 #[prost(uint32, tag = "2")]
1050 pub key_required_smaller: u32,
1051 #[prost(bool, tag = "3")]
1053 pub clean_state: bool,
1054 #[prost(message, optional, tag = "4")]
1056 pub delta_expression: ::core::option::Option<DeltaExpression>,
1057}
1058#[derive(prost_helpers::AnyPB)]
1059#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1060pub struct InequalityPairV2 {
1061 #[prost(uint32, tag = "1")]
1063 pub left_idx: u32,
1064 #[prost(uint32, tag = "2")]
1066 pub right_idx: u32,
1067 #[prost(bool, tag = "3")]
1069 pub clean_left_state: bool,
1070 #[prost(bool, tag = "4")]
1072 pub clean_right_state: bool,
1073 #[prost(enumeration = "InequalityType", tag = "5")]
1075 pub op: i32,
1076}
1077#[derive(prost_helpers::AnyPB)]
1078#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1079pub struct JoinKeyWatermarkIndex {
1080 #[prost(uint32, tag = "1")]
1082 pub index: u32,
1083 #[prost(bool, tag = "2")]
1085 pub do_state_cleaning: bool,
1086}
1087#[derive(prost_helpers::AnyPB)]
1088#[derive(Clone, PartialEq, ::prost::Message)]
1089pub struct HashJoinWatermarkHandleDesc {
1090 #[prost(message, repeated, tag = "1")]
1092 pub watermark_indices_in_jk: ::prost::alloc::vec::Vec<JoinKeyWatermarkIndex>,
1093 #[prost(message, repeated, tag = "2")]
1095 pub inequality_pairs: ::prost::alloc::vec::Vec<InequalityPairV2>,
1096}
1097#[derive(prost_helpers::AnyPB)]
1098#[derive(Clone, PartialEq, ::prost::Message)]
1099pub struct HashJoinNode {
1100 #[prost(enumeration = "super::plan_common::JoinType", tag = "1")]
1101 pub join_type: i32,
1102 #[prost(int32, repeated, tag = "2")]
1103 pub left_key: ::prost::alloc::vec::Vec<i32>,
1104 #[prost(int32, repeated, tag = "3")]
1105 pub right_key: ::prost::alloc::vec::Vec<i32>,
1106 #[prost(message, optional, tag = "4")]
1107 pub condition: ::core::option::Option<super::expr::ExprNode>,
1108 #[prost(message, optional, tag = "6")]
1110 pub left_table: ::core::option::Option<super::catalog::Table>,
1111 #[prost(message, optional, tag = "7")]
1113 pub right_table: ::core::option::Option<super::catalog::Table>,
1114 #[prost(message, optional, tag = "8")]
1116 pub left_degree_table: ::core::option::Option<super::catalog::Table>,
1117 #[prost(message, optional, tag = "9")]
1119 pub right_degree_table: ::core::option::Option<super::catalog::Table>,
1120 #[prost(uint32, repeated, tag = "10")]
1122 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1123 #[prost(uint32, repeated, tag = "11")]
1128 pub left_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1129 #[prost(uint32, repeated, tag = "12")]
1134 pub right_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1135 #[prost(bool, repeated, tag = "13")]
1136 pub null_safe: ::prost::alloc::vec::Vec<bool>,
1137 #[prost(bool, tag = "14")]
1140 pub is_append_only: bool,
1141 #[deprecated]
1144 #[prost(enumeration = "JoinEncodingType", tag = "15")]
1145 pub join_encoding_type: i32,
1146 #[prost(message, optional, tag = "17")]
1148 pub watermark_handle_desc: ::core::option::Option<HashJoinWatermarkHandleDesc>,
1149}
1150#[derive(prost_helpers::AnyPB)]
1151#[derive(Clone, PartialEq, ::prost::Message)]
1152pub struct AsOfJoinNode {
1153 #[prost(enumeration = "super::plan_common::AsOfJoinType", tag = "1")]
1154 pub join_type: i32,
1155 #[prost(int32, repeated, tag = "2")]
1156 pub left_key: ::prost::alloc::vec::Vec<i32>,
1157 #[prost(int32, repeated, tag = "3")]
1158 pub right_key: ::prost::alloc::vec::Vec<i32>,
1159 #[prost(message, optional, tag = "4")]
1161 pub left_table: ::core::option::Option<super::catalog::Table>,
1162 #[prost(message, optional, tag = "5")]
1164 pub right_table: ::core::option::Option<super::catalog::Table>,
1165 #[prost(uint32, repeated, tag = "6")]
1167 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1168 #[prost(uint32, repeated, tag = "7")]
1172 pub left_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1173 #[prost(uint32, repeated, tag = "8")]
1177 pub right_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1178 #[prost(bool, repeated, tag = "9")]
1179 pub null_safe: ::prost::alloc::vec::Vec<bool>,
1180 #[prost(message, optional, tag = "10")]
1181 pub asof_desc: ::core::option::Option<super::plan_common::AsOfJoinDesc>,
1182 #[deprecated]
1185 #[prost(enumeration = "JoinEncodingType", tag = "11")]
1186 pub join_encoding_type: i32,
1187 #[prost(bool, optional, tag = "12")]
1191 pub use_cache: ::core::option::Option<bool>,
1192}
1193#[derive(prost_helpers::AnyPB)]
1194#[derive(Clone, PartialEq, ::prost::Message)]
1195pub struct TemporalJoinNode {
1196 #[prost(enumeration = "super::plan_common::JoinType", tag = "1")]
1197 pub join_type: i32,
1198 #[prost(int32, repeated, tag = "2")]
1199 pub left_key: ::prost::alloc::vec::Vec<i32>,
1200 #[prost(int32, repeated, tag = "3")]
1201 pub right_key: ::prost::alloc::vec::Vec<i32>,
1202 #[prost(bool, repeated, tag = "4")]
1203 pub null_safe: ::prost::alloc::vec::Vec<bool>,
1204 #[prost(message, optional, tag = "5")]
1205 pub condition: ::core::option::Option<super::expr::ExprNode>,
1206 #[prost(uint32, repeated, tag = "6")]
1208 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1209 #[prost(message, optional, tag = "7")]
1211 pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1212 #[prost(uint32, repeated, tag = "8")]
1214 pub table_output_indices: ::prost::alloc::vec::Vec<u32>,
1215 #[prost(message, optional, tag = "9")]
1217 pub memo_table: ::core::option::Option<super::catalog::Table>,
1218 #[prost(bool, tag = "10")]
1220 pub is_nested_loop: bool,
1221 #[prost(bool, tag = "11")]
1223 pub is_broadcast: bool,
1224}
1225#[derive(prost_helpers::AnyPB)]
1226#[derive(Clone, PartialEq, ::prost::Message)]
1227pub struct DynamicFilterNode {
1228 #[prost(uint32, tag = "1")]
1229 pub left_key: u32,
1230 #[prost(message, optional, tag = "2")]
1232 pub condition: ::core::option::Option<super::expr::ExprNode>,
1233 #[prost(message, optional, tag = "3")]
1235 pub left_table: ::core::option::Option<super::catalog::Table>,
1236 #[prost(message, optional, tag = "4")]
1238 pub right_table: ::core::option::Option<super::catalog::Table>,
1239 #[deprecated]
1246 #[prost(bool, tag = "5")]
1247 pub condition_always_relax: bool,
1248 #[prost(bool, tag = "6")]
1250 pub cleaned_by_watermark: bool,
1251}
1252#[derive(prost_helpers::AnyPB)]
1255#[derive(Clone, PartialEq, ::prost::Message)]
1256pub struct DeltaIndexJoinNode {
1257 #[prost(enumeration = "super::plan_common::JoinType", tag = "1")]
1258 pub join_type: i32,
1259 #[prost(int32, repeated, tag = "2")]
1260 pub left_key: ::prost::alloc::vec::Vec<i32>,
1261 #[prost(int32, repeated, tag = "3")]
1262 pub right_key: ::prost::alloc::vec::Vec<i32>,
1263 #[prost(message, optional, tag = "4")]
1264 pub condition: ::core::option::Option<super::expr::ExprNode>,
1265 #[prost(uint32, tag = "7", wrapper = "crate::id::TableId")]
1267 pub left_table_id: crate::id::TableId,
1268 #[prost(uint32, tag = "8", wrapper = "crate::id::TableId")]
1270 pub right_table_id: crate::id::TableId,
1271 #[prost(message, optional, tag = "9")]
1273 pub left_info: ::core::option::Option<ArrangementInfo>,
1274 #[prost(message, optional, tag = "10")]
1276 pub right_info: ::core::option::Option<ArrangementInfo>,
1277 #[prost(uint32, repeated, tag = "11")]
1279 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1280}
1281#[derive(prost_helpers::AnyPB)]
1282#[derive(Clone, PartialEq, ::prost::Message)]
1283pub struct HopWindowNode {
1284 #[prost(uint32, tag = "1")]
1285 pub time_col: u32,
1286 #[prost(message, optional, tag = "2")]
1287 pub window_slide: ::core::option::Option<super::data::Interval>,
1288 #[prost(message, optional, tag = "3")]
1289 pub window_size: ::core::option::Option<super::data::Interval>,
1290 #[prost(uint32, repeated, tag = "4")]
1291 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1292 #[prost(message, repeated, tag = "5")]
1293 pub window_start_exprs: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
1294 #[prost(message, repeated, tag = "6")]
1295 pub window_end_exprs: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
1296}
1297#[derive(prost_helpers::AnyPB)]
1298#[derive(Clone, PartialEq, ::prost::Message)]
1299pub struct MergeNode {
1300 #[deprecated]
1307 #[prost(uint32, repeated, packed = "false", tag = "1")]
1308 pub upstream_actor_id: ::prost::alloc::vec::Vec<u32>,
1309 #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
1310 pub upstream_fragment_id: crate::id::FragmentId,
1311 #[prost(enumeration = "DispatcherType", tag = "3")]
1314 pub upstream_dispatcher_type: i32,
1315 #[deprecated]
1317 #[prost(message, repeated, tag = "4")]
1318 pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
1319 #[prost(bool, tag = "5")]
1323 pub allow_empty_upstream: bool,
1324}
1325#[derive(prost_helpers::AnyPB)]
1328#[derive(Clone, PartialEq, ::prost::Message)]
1329pub struct ExchangeNode {
1330 #[prost(message, optional, tag = "1")]
1331 pub strategy: ::core::option::Option<DispatchStrategy>,
1332}
1333#[derive(prost_helpers::AnyPB)]
1339#[derive(Clone, PartialEq, ::prost::Message)]
1340pub struct StreamScanNode {
1341 #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
1342 pub table_id: crate::id::TableId,
1343 #[prost(int32, repeated, tag = "2")]
1348 pub upstream_column_ids: ::prost::alloc::vec::Vec<i32>,
1349 #[prost(uint32, repeated, tag = "3")]
1354 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1355 #[prost(enumeration = "StreamScanType", tag = "4")]
1360 pub stream_scan_type: i32,
1361 #[prost(message, optional, tag = "5")]
1363 pub state_table: ::core::option::Option<super::catalog::Table>,
1364 #[prost(message, optional, tag = "7")]
1367 pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1368 #[prost(uint32, optional, tag = "8")]
1370 pub rate_limit: ::core::option::Option<u32>,
1371 #[deprecated]
1373 #[prost(uint32, tag = "9")]
1374 pub snapshot_read_barrier_interval: u32,
1375 #[prost(message, optional, tag = "10")]
1378 pub arrangement_table: ::core::option::Option<super::catalog::Table>,
1379 #[prost(uint64, optional, tag = "11")]
1380 pub snapshot_backfill_epoch: ::core::option::Option<u64>,
1381 #[prost(message, optional, tag = "13")]
1383 pub pk_scan_range: ::core::option::Option<super::batch_plan::ScanRange>,
1384}
1385#[derive(prost_helpers::AnyPB)]
1387#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1388pub struct StreamCdcScanOptions {
1389 #[prost(bool, tag = "1")]
1391 pub disable_backfill: bool,
1392 #[prost(uint32, tag = "2")]
1393 pub snapshot_barrier_interval: u32,
1394 #[prost(uint32, tag = "3")]
1395 pub snapshot_batch_size: u32,
1396 #[prost(uint32, tag = "4")]
1397 pub backfill_parallelism: u32,
1398 #[prost(uint64, tag = "5")]
1399 pub backfill_num_rows_per_split: u64,
1400 #[prost(bool, tag = "6")]
1401 pub backfill_as_even_splits: bool,
1402 #[prost(uint32, tag = "7")]
1403 pub backfill_split_pk_column_index: u32,
1404}
1405#[derive(prost_helpers::AnyPB)]
1406#[derive(Clone, PartialEq, ::prost::Message)]
1407pub struct StreamCdcScanNode {
1408 #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
1409 pub table_id: crate::id::TableId,
1410 #[prost(int32, repeated, tag = "2")]
1413 pub upstream_column_ids: ::prost::alloc::vec::Vec<i32>,
1414 #[prost(uint32, repeated, tag = "3")]
1416 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1417 #[prost(message, optional, tag = "4")]
1419 pub state_table: ::core::option::Option<super::catalog::Table>,
1420 #[prost(message, optional, tag = "5")]
1422 pub cdc_table_desc: ::core::option::Option<super::plan_common::ExternalTableDesc>,
1423 #[prost(uint32, optional, tag = "6")]
1425 pub rate_limit: ::core::option::Option<u32>,
1426 #[prost(bool, tag = "7")]
1429 pub disable_backfill: bool,
1430 #[prost(message, optional, tag = "8")]
1431 pub options: ::core::option::Option<StreamCdcScanOptions>,
1432}
1433#[derive(prost_helpers::AnyPB)]
1437#[derive(Clone, PartialEq, ::prost::Message)]
1438pub struct BatchPlanNode {
1439 #[prost(message, optional, tag = "1")]
1440 pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1441 #[prost(int32, repeated, tag = "2")]
1442 pub column_ids: ::prost::alloc::vec::Vec<i32>,
1443}
1444#[derive(prost_helpers::AnyPB)]
1445#[derive(Clone, PartialEq, ::prost::Message)]
1446pub struct ArrangementInfo {
1447 #[prost(message, repeated, tag = "1")]
1450 pub arrange_key_orders: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1451 #[prost(message, repeated, tag = "2")]
1453 pub column_descs: ::prost::alloc::vec::Vec<super::plan_common::ColumnDesc>,
1454 #[prost(message, optional, tag = "4")]
1456 pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1457 #[prost(uint32, repeated, tag = "5")]
1459 pub output_col_idx: ::prost::alloc::vec::Vec<u32>,
1460}
1461#[derive(prost_helpers::AnyPB)]
1464#[derive(Clone, PartialEq, ::prost::Message)]
1465pub struct ArrangeNode {
1466 #[prost(message, optional, tag = "1")]
1468 pub table_info: ::core::option::Option<ArrangementInfo>,
1469 #[prost(uint32, repeated, tag = "2")]
1471 pub distribution_key: ::prost::alloc::vec::Vec<u32>,
1472 #[prost(message, optional, tag = "3")]
1474 pub table: ::core::option::Option<super::catalog::Table>,
1475}
1476#[derive(prost_helpers::AnyPB)]
1478#[derive(Clone, PartialEq, ::prost::Message)]
1479pub struct LookupNode {
1480 #[prost(int32, repeated, tag = "1")]
1482 pub arrange_key: ::prost::alloc::vec::Vec<i32>,
1483 #[prost(int32, repeated, tag = "2")]
1485 pub stream_key: ::prost::alloc::vec::Vec<i32>,
1486 #[prost(bool, tag = "3")]
1488 pub use_current_epoch: bool,
1489 #[prost(int32, repeated, tag = "4")]
1493 pub column_mapping: ::prost::alloc::vec::Vec<i32>,
1494 #[prost(message, optional, tag = "7")]
1496 pub arrangement_table_info: ::core::option::Option<ArrangementInfo>,
1497 #[prost(oneof = "lookup_node::ArrangementTableId", tags = "5, 6")]
1498 pub arrangement_table_id: ::core::option::Option<lookup_node::ArrangementTableId>,
1499}
1500pub mod lookup_node {
1502 #[derive(prost_helpers::AnyPB)]
1503 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Oneof)]
1504 pub enum ArrangementTableId {
1505 #[prost(uint32, tag = "5", wrapper = "crate::id::TableId")]
1507 TableId(crate::id::TableId),
1508 #[prost(uint32, tag = "6", wrapper = "crate::id::TableId")]
1510 IndexId(crate::id::TableId),
1511 }
1512}
1513#[derive(prost_helpers::AnyPB)]
1515#[derive(Clone, PartialEq, ::prost::Message)]
1516pub struct WatermarkFilterNode {
1517 #[prost(message, repeated, tag = "1")]
1519 pub watermark_descs: ::prost::alloc::vec::Vec<super::catalog::WatermarkDesc>,
1520 #[prost(message, repeated, tag = "2")]
1522 pub tables: ::prost::alloc::vec::Vec<super::catalog::Table>,
1523}
1524#[derive(prost_helpers::AnyPB)]
1526#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1527pub struct UnionNode {}
1528#[derive(prost_helpers::AnyPB)]
1530#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1531pub struct LookupUnionNode {
1532 #[prost(uint32, repeated, tag = "1")]
1533 pub order: ::prost::alloc::vec::Vec<u32>,
1534}
1535#[derive(prost_helpers::AnyPB)]
1536#[derive(Clone, PartialEq, ::prost::Message)]
1537pub struct ExpandNode {
1538 #[prost(message, repeated, tag = "1")]
1539 pub column_subsets: ::prost::alloc::vec::Vec<expand_node::Subset>,
1540}
1541pub mod expand_node {
1543 #[derive(prost_helpers::AnyPB)]
1544 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1545 pub struct Subset {
1546 #[prost(uint32, repeated, tag = "1")]
1547 pub column_indices: ::prost::alloc::vec::Vec<u32>,
1548 }
1549}
1550#[derive(prost_helpers::AnyPB)]
1551#[derive(Clone, PartialEq, ::prost::Message)]
1552pub struct ProjectSetNode {
1553 #[prost(message, repeated, tag = "1")]
1554 pub select_list: ::prost::alloc::vec::Vec<super::expr::ProjectSetSelectItem>,
1555 #[prost(uint32, repeated, tag = "2")]
1559 pub watermark_input_cols: ::prost::alloc::vec::Vec<u32>,
1560 #[prost(uint32, repeated, tag = "3")]
1561 pub watermark_expr_indices: ::prost::alloc::vec::Vec<u32>,
1562 #[prost(uint32, repeated, tag = "4")]
1563 pub nondecreasing_exprs: ::prost::alloc::vec::Vec<u32>,
1564}
1565#[derive(prost_helpers::AnyPB)]
1567#[derive(Clone, PartialEq, ::prost::Message)]
1568pub struct SortNode {
1569 #[prost(message, optional, tag = "1")]
1571 pub state_table: ::core::option::Option<super::catalog::Table>,
1572 #[prost(uint32, tag = "2")]
1574 pub sort_column_index: u32,
1575}
1576#[derive(prost_helpers::AnyPB)]
1578#[derive(Clone, PartialEq, ::prost::Message)]
1579pub struct DmlNode {
1580 #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
1582 pub table_id: crate::id::TableId,
1583 #[prost(uint64, tag = "3")]
1585 pub table_version_id: u64,
1586 #[prost(message, repeated, tag = "2")]
1588 pub column_descs: ::prost::alloc::vec::Vec<super::plan_common::ColumnDesc>,
1589 #[prost(uint32, optional, tag = "4")]
1590 pub rate_limit: ::core::option::Option<u32>,
1591}
1592#[derive(prost_helpers::AnyPB)]
1593#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1594pub struct RowIdGenNode {
1595 #[prost(uint64, tag = "1")]
1596 pub row_id_index: u64,
1597}
1598#[derive(prost_helpers::AnyPB)]
1599#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1600pub struct NowModeUpdateCurrent {}
1601#[derive(prost_helpers::AnyPB)]
1602#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1603pub struct NowModeGenerateSeries {
1604 #[prost(message, optional, tag = "1")]
1605 pub start_timestamp: ::core::option::Option<super::data::Datum>,
1606 #[prost(message, optional, tag = "2")]
1607 pub interval: ::core::option::Option<super::data::Datum>,
1608}
1609#[derive(prost_helpers::AnyPB)]
1610#[derive(Clone, PartialEq, ::prost::Message)]
1611pub struct NowNode {
1612 #[prost(message, optional, tag = "1")]
1614 pub state_table: ::core::option::Option<super::catalog::Table>,
1615 #[prost(oneof = "now_node::Mode", tags = "101, 102")]
1616 pub mode: ::core::option::Option<now_node::Mode>,
1617}
1618pub mod now_node {
1620 #[derive(prost_helpers::AnyPB)]
1621 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Oneof)]
1622 pub enum Mode {
1623 #[prost(message, tag = "101")]
1624 UpdateCurrent(super::NowModeUpdateCurrent),
1625 #[prost(message, tag = "102")]
1626 GenerateSeries(super::NowModeGenerateSeries),
1627 }
1628}
1629#[derive(prost_helpers::AnyPB)]
1630#[derive(Clone, PartialEq, ::prost::Message)]
1631pub struct ValuesNode {
1632 #[prost(message, repeated, tag = "1")]
1633 pub tuples: ::prost::alloc::vec::Vec<values_node::ExprTuple>,
1634 #[prost(message, repeated, tag = "2")]
1635 pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
1636}
1637pub mod values_node {
1639 #[derive(prost_helpers::AnyPB)]
1640 #[derive(Clone, PartialEq, ::prost::Message)]
1641 pub struct ExprTuple {
1642 #[prost(message, repeated, tag = "1")]
1643 pub cells: ::prost::alloc::vec::Vec<super::super::expr::ExprNode>,
1644 }
1645}
1646#[derive(prost_helpers::AnyPB)]
1647#[derive(Clone, PartialEq, ::prost::Message)]
1648pub struct DedupNode {
1649 #[prost(message, optional, tag = "1")]
1650 pub state_table: ::core::option::Option<super::catalog::Table>,
1651 #[prost(uint32, repeated, tag = "2")]
1652 pub dedup_column_indices: ::prost::alloc::vec::Vec<u32>,
1653}
1654#[derive(prost_helpers::AnyPB)]
1655#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1656pub struct NoOpNode {}
1657#[derive(prost_helpers::AnyPB)]
1658#[derive(Clone, PartialEq, ::prost::Message)]
1659pub struct EowcOverWindowNode {
1660 #[prost(message, repeated, tag = "1")]
1661 pub calls: ::prost::alloc::vec::Vec<super::expr::WindowFunction>,
1662 #[prost(uint32, repeated, tag = "2")]
1663 pub partition_by: ::prost::alloc::vec::Vec<u32>,
1664 #[prost(message, repeated, tag = "3")]
1666 pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1667 #[prost(message, optional, tag = "4")]
1668 pub state_table: ::core::option::Option<super::catalog::Table>,
1669 #[prost(message, optional, tag = "5")]
1672 pub intermediate_state_table: ::core::option::Option<super::catalog::Table>,
1673}
1674#[derive(prost_helpers::AnyPB)]
1675#[derive(Clone, PartialEq, ::prost::Message)]
1676pub struct OverWindowNode {
1677 #[prost(message, repeated, tag = "1")]
1678 pub calls: ::prost::alloc::vec::Vec<super::expr::WindowFunction>,
1679 #[prost(uint32, repeated, tag = "2")]
1680 pub partition_by: ::prost::alloc::vec::Vec<u32>,
1681 #[prost(message, repeated, tag = "3")]
1682 pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1683 #[prost(message, optional, tag = "4")]
1684 pub state_table: ::core::option::Option<super::catalog::Table>,
1685 #[deprecated]
1687 #[prost(enumeration = "OverWindowCachePolicy", tag = "5")]
1688 pub cache_policy: i32,
1689 #[prost(bool, tag = "6")]
1694 pub enable_state_cleaning: bool,
1695}
1696#[derive(prost_helpers::AnyPB)]
1697#[derive(Clone, Copy, PartialEq, ::prost::Message)]
1698pub struct LocalApproxPercentileNode {
1699 #[prost(double, tag = "1")]
1700 pub base: f64,
1701 #[prost(uint32, tag = "2")]
1702 pub percentile_index: u32,
1703}
1704#[derive(prost_helpers::AnyPB)]
1705#[derive(Clone, PartialEq, ::prost::Message)]
1706pub struct GlobalApproxPercentileNode {
1707 #[prost(double, tag = "1")]
1708 pub base: f64,
1709 #[prost(double, tag = "2")]
1710 pub quantile: f64,
1711 #[prost(message, optional, tag = "3")]
1712 pub bucket_state_table: ::core::option::Option<super::catalog::Table>,
1713 #[prost(message, optional, tag = "4")]
1714 pub count_state_table: ::core::option::Option<super::catalog::Table>,
1715}
1716#[derive(prost_helpers::AnyPB)]
1717#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1718pub struct RowMergeNode {
1719 #[prost(message, optional, tag = "1")]
1720 pub lhs_mapping: ::core::option::Option<super::catalog::ColIndexMapping>,
1721 #[prost(message, optional, tag = "2")]
1722 pub rhs_mapping: ::core::option::Option<super::catalog::ColIndexMapping>,
1723}
1724#[derive(prost_helpers::AnyPB)]
1725#[derive(Clone, PartialEq, ::prost::Message)]
1726pub struct SyncLogStoreNode {
1727 #[prost(message, optional, tag = "1")]
1728 pub log_store_table: ::core::option::Option<super::catalog::Table>,
1729 #[deprecated]
1731 #[prost(uint32, optional, tag = "2")]
1732 pub pause_duration_ms: ::core::option::Option<u32>,
1733 #[deprecated]
1735 #[prost(uint32, optional, tag = "3")]
1736 pub buffer_size: ::core::option::Option<u32>,
1737 #[prost(bool, tag = "4")]
1738 pub aligned: bool,
1739}
1740#[derive(prost_helpers::AnyPB)]
1741#[derive(Clone, PartialEq, ::prost::Message)]
1742pub struct MaterializedExprsNode {
1743 #[prost(message, repeated, tag = "1")]
1744 pub exprs: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
1745 #[prost(message, optional, tag = "2")]
1746 pub state_table: ::core::option::Option<super::catalog::Table>,
1747 #[prost(uint32, optional, tag = "3")]
1748 pub state_clean_col_idx: ::core::option::Option<u32>,
1749}
1750#[derive(prost_helpers::AnyPB)]
1751#[derive(Clone, PartialEq, ::prost::Message)]
1752pub struct VectorIndexWriteNode {
1753 #[prost(message, optional, tag = "1")]
1754 pub table: ::core::option::Option<super::catalog::Table>,
1755}
1756#[derive(prost_helpers::AnyPB)]
1757#[derive(Clone, PartialEq, ::prost::Message)]
1758pub struct VectorIndexLookupJoinNode {
1759 #[prost(message, optional, tag = "1")]
1760 pub reader_desc: ::core::option::Option<super::plan_common::VectorIndexReaderDesc>,
1761 #[prost(uint32, tag = "2")]
1762 pub vector_column_idx: u32,
1763}
1764#[derive(prost_helpers::AnyPB)]
1765#[derive(Clone, PartialEq, ::prost::Message)]
1766pub struct UpstreamSinkUnionNode {
1767 #[prost(message, repeated, tag = "1")]
1770 pub init_upstreams: ::prost::alloc::vec::Vec<UpstreamSinkInfo>,
1771}
1772#[derive(prost_helpers::AnyPB)]
1773#[derive(Clone, PartialEq, ::prost::Message)]
1774pub struct LocalityProviderNode {
1775 #[prost(uint32, repeated, tag = "1")]
1777 pub locality_columns: ::prost::alloc::vec::Vec<u32>,
1778 #[prost(message, optional, tag = "2")]
1780 pub state_table: ::core::option::Option<super::catalog::Table>,
1781 #[prost(message, optional, tag = "3")]
1783 pub progress_table: ::core::option::Option<super::catalog::Table>,
1784 #[prost(uint32, optional, tag = "4")]
1786 pub rate_limit: ::core::option::Option<u32>,
1787}
1788#[derive(prost_helpers::AnyPB)]
1789#[derive(Clone, PartialEq, ::prost::Message)]
1790pub struct EowcGapFillNode {
1791 #[prost(uint32, tag = "1")]
1792 pub time_column_index: u32,
1793 #[prost(message, optional, tag = "2")]
1794 pub interval: ::core::option::Option<super::expr::ExprNode>,
1795 #[prost(uint32, repeated, tag = "3")]
1796 pub fill_columns: ::prost::alloc::vec::Vec<u32>,
1797 #[prost(string, repeated, tag = "4")]
1798 pub fill_strategies: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
1799 #[prost(message, optional, tag = "6")]
1800 pub prev_row_table: ::core::option::Option<super::catalog::Table>,
1801 #[prost(uint32, repeated, tag = "7")]
1802 pub partition_by_indices: ::prost::alloc::vec::Vec<u32>,
1803}
1804#[derive(prost_helpers::AnyPB)]
1805#[derive(Clone, PartialEq, ::prost::Message)]
1806pub struct GapFillNode {
1807 #[prost(uint32, tag = "1")]
1808 pub time_column_index: u32,
1809 #[prost(message, optional, tag = "2")]
1810 pub interval: ::core::option::Option<super::expr::ExprNode>,
1811 #[prost(uint32, repeated, tag = "3")]
1812 pub fill_columns: ::prost::alloc::vec::Vec<u32>,
1813 #[prost(string, repeated, tag = "4")]
1814 pub fill_strategies: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
1815 #[prost(message, optional, tag = "5")]
1816 pub state_table: ::core::option::Option<super::catalog::Table>,
1817 #[prost(uint32, repeated, tag = "6")]
1818 pub partition_by_indices: ::prost::alloc::vec::Vec<u32>,
1819 #[prost(uint32, repeated, tag = "7")]
1820 pub pointer_key_indices: ::prost::alloc::vec::Vec<u32>,
1821}
1822#[derive(prost_helpers::AnyPB)]
1825#[derive(Clone, PartialEq, ::prost::Message)]
1826pub struct MatchRecognizeNode {
1827 #[prost(uint32, repeated, tag = "1")]
1828 pub partition_by: ::prost::alloc::vec::Vec<u32>,
1829 #[prost(message, repeated, tag = "2")]
1832 pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1833 #[prost(message, repeated, tag = "3")]
1834 pub measures: ::prost::alloc::vec::Vec<MatchRecognizeMeasure>,
1835 #[prost(message, repeated, tag = "5")]
1836 pub defines: ::prost::alloc::vec::Vec<MatchRecognizeDefine>,
1837 #[prost(message, optional, tag = "12")]
1841 pub pattern_node: ::core::option::Option<MatchRecognizePatternNode>,
1842 #[prost(message, optional, tag = "8")]
1846 pub state_table: ::core::option::Option<super::catalog::Table>,
1847 #[prost(message, optional, tag = "9")]
1849 pub after_match_skip: ::core::option::Option<MatchRecognizeAfterMatchSkip>,
1850 #[prost(message, optional, tag = "11")]
1853 pub within: ::core::option::Option<super::expr::ExprNode>,
1854 #[prost(message, optional, tag = "15")]
1858 pub within_deadline: ::core::option::Option<super::expr::ExprNode>,
1859 #[prost(enumeration = "MatchRecognizeInputMode", tag = "16")]
1863 pub input_mode: i32,
1864}
1865#[derive(prost_helpers::AnyPB)]
1869#[derive(Clone, PartialEq, ::prost::Message)]
1870pub struct MatchRecognizePatternNode {
1871 #[prost(oneof = "match_recognize_pattern_node::Node", tags = "1, 2, 3, 4, 5")]
1872 pub node: ::core::option::Option<match_recognize_pattern_node::Node>,
1873}
1874pub mod match_recognize_pattern_node {
1876 #[derive(prost_helpers::AnyPB)]
1877 #[derive(Clone, PartialEq, ::prost::Oneof)]
1878 pub enum Node {
1879 #[prost(string, tag = "1")]
1881 Var(::prost::alloc::string::String),
1882 #[prost(message, tag = "2")]
1884 Concat(super::MatchRecognizePatternSeq),
1885 #[prost(message, tag = "3")]
1887 Alternation(super::MatchRecognizePatternSeq),
1888 #[prost(message, tag = "4")]
1890 Quantified(::prost::alloc::boxed::Box<super::MatchRecognizeQuantifiedPattern>),
1891 #[prost(message, tag = "5")]
1893 Permute(super::MatchRecognizePermutePattern),
1894 }
1895}
1896#[derive(prost_helpers::AnyPB)]
1898#[derive(Clone, PartialEq, ::prost::Message)]
1899pub struct MatchRecognizePatternSeq {
1900 #[prost(message, repeated, tag = "1")]
1901 pub patterns: ::prost::alloc::vec::Vec<MatchRecognizePatternNode>,
1902}
1903#[derive(prost_helpers::AnyPB)]
1905#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1906pub struct MatchRecognizePermutePattern {
1907 #[prost(string, repeated, tag = "1")]
1908 pub vars: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
1909}
1910#[derive(prost_helpers::AnyPB)]
1912#[derive(Clone, PartialEq, ::prost::Message)]
1913pub struct MatchRecognizeQuantifiedPattern {
1914 #[prost(message, optional, boxed, tag = "1")]
1915 pub inner: ::core::option::Option<
1916 ::prost::alloc::boxed::Box<MatchRecognizePatternNode>,
1917 >,
1918 #[prost(message, optional, tag = "2")]
1919 pub quantifier: ::core::option::Option<MatchRecognizeQuantifier>,
1920 #[prost(bool, tag = "3")]
1922 pub reluctant: bool,
1923}
1924#[derive(prost_helpers::AnyPB)]
1926#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1927pub struct MatchRecognizeQuantifier {
1928 #[prost(enumeration = "match_recognize_quantifier::Kind", tag = "1")]
1929 pub kind: i32,
1930 #[prost(uint32, tag = "2")]
1932 pub min: u32,
1933 #[prost(uint32, optional, tag = "3")]
1935 pub max: ::core::option::Option<u32>,
1936}
1937pub mod match_recognize_quantifier {
1939 #[derive(prost_helpers::AnyPB)]
1940 #[derive(
1941 Clone,
1942 Copy,
1943 Debug,
1944 PartialEq,
1945 Eq,
1946 Hash,
1947 PartialOrd,
1948 Ord,
1949 ::prost::Enumeration
1950 )]
1951 #[repr(i32)]
1952 pub enum Kind {
1953 Unspecified = 0,
1954 Star = 1,
1956 Plus = 2,
1958 Question = 3,
1960 Range = 4,
1962 }
1963 impl Kind {
1964 pub fn as_str_name(&self) -> &'static str {
1969 match self {
1970 Self::Unspecified => "KIND_UNSPECIFIED",
1971 Self::Star => "KIND_STAR",
1972 Self::Plus => "KIND_PLUS",
1973 Self::Question => "KIND_QUESTION",
1974 Self::Range => "KIND_RANGE",
1975 }
1976 }
1977 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
1979 match value {
1980 "KIND_UNSPECIFIED" => Some(Self::Unspecified),
1981 "KIND_STAR" => Some(Self::Star),
1982 "KIND_PLUS" => Some(Self::Plus),
1983 "KIND_QUESTION" => Some(Self::Question),
1984 "KIND_RANGE" => Some(Self::Range),
1985 _ => None,
1986 }
1987 }
1988 }
1989}
1990#[derive(prost_helpers::AnyPB)]
1994#[derive(Clone, PartialEq, ::prost::Message)]
1995pub struct MatchRecognizeDefine {
1996 #[prost(string, tag = "1")]
1997 pub symbol: ::prost::alloc::string::String,
1998 #[prost(message, optional, tag = "2")]
2000 pub condition: ::core::option::Option<super::expr::ExprNode>,
2001 #[prost(message, repeated, tag = "3")]
2002 pub slots: ::prost::alloc::vec::Vec<MatchRecognizeDefineSlot>,
2003}
2004#[derive(prost_helpers::AnyPB)]
2006#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2007pub struct MatchRecognizeAfterMatchSkip {
2008 #[prost(enumeration = "match_recognize_after_match_skip::Mode", tag = "1")]
2009 pub mode: i32,
2010 #[prost(string, optional, tag = "2")]
2012 pub target: ::core::option::Option<::prost::alloc::string::String>,
2013}
2014pub mod match_recognize_after_match_skip {
2016 #[derive(prost_helpers::AnyPB)]
2017 #[derive(
2018 Clone,
2019 Copy,
2020 Debug,
2021 PartialEq,
2022 Eq,
2023 Hash,
2024 PartialOrd,
2025 Ord,
2026 ::prost::Enumeration
2027 )]
2028 #[repr(i32)]
2029 pub enum Mode {
2030 Unspecified = 0,
2031 PastLastRow = 1,
2033 ToNextRow = 2,
2035 ToFirst = 3,
2037 ToLast = 4,
2038 }
2039 impl Mode {
2040 pub fn as_str_name(&self) -> &'static str {
2045 match self {
2046 Self::Unspecified => "MODE_UNSPECIFIED",
2047 Self::PastLastRow => "MODE_PAST_LAST_ROW",
2048 Self::ToNextRow => "MODE_TO_NEXT_ROW",
2049 Self::ToFirst => "MODE_TO_FIRST",
2050 Self::ToLast => "MODE_TO_LAST",
2051 }
2052 }
2053 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2055 match value {
2056 "MODE_UNSPECIFIED" => Some(Self::Unspecified),
2057 "MODE_PAST_LAST_ROW" => Some(Self::PastLastRow),
2058 "MODE_TO_NEXT_ROW" => Some(Self::ToNextRow),
2059 "MODE_TO_FIRST" => Some(Self::ToFirst),
2060 "MODE_TO_LAST" => Some(Self::ToLast),
2061 _ => None,
2062 }
2063 }
2064 }
2065}
2066#[derive(prost_helpers::AnyPB)]
2068#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2069pub struct MatchRecognizeDefineSlot {
2070 #[prost(enumeration = "match_recognize_define_slot::Kind", tag = "1")]
2071 pub kind: i32,
2072 #[prost(string, repeated, tag = "2")]
2073 pub vars: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
2074 #[prost(uint32, tag = "3")]
2075 pub col_idx: u32,
2076 #[prost(uint32, tag = "4")]
2077 pub offset: u32,
2078}
2079pub mod match_recognize_define_slot {
2081 #[derive(prost_helpers::AnyPB)]
2082 #[derive(
2083 Clone,
2084 Copy,
2085 Debug,
2086 PartialEq,
2087 Eq,
2088 Hash,
2089 PartialOrd,
2090 Ord,
2091 ::prost::Enumeration
2092 )]
2093 #[repr(i32)]
2094 pub enum Kind {
2095 Unspecified = 0,
2096 SelfCol = 1,
2098 Prev = 2,
2100 Next = 3,
2101 RunningFirst = 4,
2106 RunningLast = 5,
2107 }
2108 impl Kind {
2109 pub fn as_str_name(&self) -> &'static str {
2114 match self {
2115 Self::Unspecified => "KIND_UNSPECIFIED",
2116 Self::SelfCol => "KIND_SELF_COL",
2117 Self::Prev => "KIND_PREV",
2118 Self::Next => "KIND_NEXT",
2119 Self::RunningFirst => "KIND_RUNNING_FIRST",
2120 Self::RunningLast => "KIND_RUNNING_LAST",
2121 }
2122 }
2123 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2125 match value {
2126 "KIND_UNSPECIFIED" => Some(Self::Unspecified),
2127 "KIND_SELF_COL" => Some(Self::SelfCol),
2128 "KIND_PREV" => Some(Self::Prev),
2129 "KIND_NEXT" => Some(Self::Next),
2130 "KIND_RUNNING_FIRST" => Some(Self::RunningFirst),
2131 "KIND_RUNNING_LAST" => Some(Self::RunningLast),
2132 _ => None,
2133 }
2134 }
2135 }
2136}
2137#[derive(prost_helpers::AnyPB)]
2141#[derive(Clone, PartialEq, ::prost::Message)]
2142pub struct MatchRecognizeMeasure {
2143 #[prost(message, optional, tag = "1")]
2145 pub expr: ::core::option::Option<super::expr::ExprNode>,
2146 #[prost(string, tag = "2")]
2147 pub name: ::prost::alloc::string::String,
2148 #[prost(message, repeated, tag = "3")]
2149 pub slots: ::prost::alloc::vec::Vec<MatchRecognizeMeasureSlot>,
2150}
2151#[derive(prost_helpers::AnyPB)]
2153#[derive(Clone, PartialEq, ::prost::Message)]
2154pub struct MatchRecognizeMeasureSlot {
2155 #[prost(enumeration = "match_recognize_measure_slot::Kind", tag = "1")]
2156 pub kind: i32,
2157 #[prost(string, repeated, tag = "2")]
2160 pub vars: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
2161 #[prost(uint32, tag = "3")]
2162 pub col_idx: u32,
2163 #[prost(message, optional, tag = "4")]
2164 pub data_type: ::core::option::Option<super::data::DataType>,
2165 #[prost(message, optional, tag = "5")]
2168 pub agg_call: ::core::option::Option<super::expr::AggCall>,
2169}
2170pub mod match_recognize_measure_slot {
2172 #[derive(prost_helpers::AnyPB)]
2173 #[derive(
2174 Clone,
2175 Copy,
2176 Debug,
2177 PartialEq,
2178 Eq,
2179 Hash,
2180 PartialOrd,
2181 Ord,
2182 ::prost::Enumeration
2183 )]
2184 #[repr(i32)]
2185 pub enum Kind {
2186 Unspecified = 0,
2187 Last = 1,
2190 First = 2,
2192 Classifier = 3,
2194 CountStar = 4,
2196 Count = 5,
2198 Min = 6,
2200 Max = 7,
2201 Sum = 8,
2204 }
2205 impl Kind {
2206 pub fn as_str_name(&self) -> &'static str {
2211 match self {
2212 Self::Unspecified => "KIND_UNSPECIFIED",
2213 Self::Last => "KIND_LAST",
2214 Self::First => "KIND_FIRST",
2215 Self::Classifier => "KIND_CLASSIFIER",
2216 Self::CountStar => "KIND_COUNT_STAR",
2217 Self::Count => "KIND_COUNT",
2218 Self::Min => "KIND_MIN",
2219 Self::Max => "KIND_MAX",
2220 Self::Sum => "KIND_SUM",
2221 }
2222 }
2223 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2225 match value {
2226 "KIND_UNSPECIFIED" => Some(Self::Unspecified),
2227 "KIND_LAST" => Some(Self::Last),
2228 "KIND_FIRST" => Some(Self::First),
2229 "KIND_CLASSIFIER" => Some(Self::Classifier),
2230 "KIND_COUNT_STAR" => Some(Self::CountStar),
2231 "KIND_COUNT" => Some(Self::Count),
2232 "KIND_MIN" => Some(Self::Min),
2233 "KIND_MAX" => Some(Self::Max),
2234 "KIND_SUM" => Some(Self::Sum),
2235 _ => None,
2236 }
2237 }
2238 }
2239}
2240#[derive(prost_helpers::AnyPB)]
2241#[derive(Clone, PartialEq, ::prost::Message)]
2242pub struct StreamNode {
2243 #[prost(uint64, tag = "1", wrapper = "crate::id::StreamNodeLocalOperatorId")]
2246 pub operator_id: crate::id::StreamNodeLocalOperatorId,
2247 #[prost(message, repeated, tag = "3")]
2249 pub input: ::prost::alloc::vec::Vec<StreamNode>,
2250 #[prost(uint32, repeated, tag = "2")]
2251 pub stream_key: ::prost::alloc::vec::Vec<u32>,
2252 #[prost(enumeration = "stream_node::StreamKind", tag = "24")]
2253 pub stream_kind: i32,
2254 #[prost(string, tag = "18")]
2255 pub identity: ::prost::alloc::string::String,
2256 #[prost(message, repeated, tag = "19")]
2258 pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2259 #[prost(
2260 oneof = "stream_node::NodeBody",
2261 tags = "100, 101, 102, 103, 104, 105, 106, 107, 108, 109, 110, 111, 112, 113, 114, 115, 116, 117, 118, 119, 120, 121, 122, 123, 124, 125, 126, 127, 128, 129, 130, 131, 132, 133, 134, 135, 136, 137, 138, 139, 140, 142, 143, 144, 145, 146, 147, 148, 149, 150, 151, 152, 153, 154, 155, 156, 157, 158, 159"
2262 )]
2263 pub node_body: ::core::option::Option<stream_node::NodeBody>,
2264}
2265pub mod stream_node {
2267 #[derive(prost_helpers::AnyPB)]
2270 #[derive(
2271 Clone,
2272 Copy,
2273 Debug,
2274 PartialEq,
2275 Eq,
2276 Hash,
2277 PartialOrd,
2278 Ord,
2279 ::prost::Enumeration
2280 )]
2281 #[repr(i32)]
2282 pub enum StreamKind {
2283 Retract = 0,
2285 AppendOnly = 1,
2286 Upsert = 2,
2287 }
2288 impl StreamKind {
2289 pub fn as_str_name(&self) -> &'static str {
2294 match self {
2295 Self::Retract => "STREAM_KIND_RETRACT",
2296 Self::AppendOnly => "STREAM_KIND_APPEND_ONLY",
2297 Self::Upsert => "STREAM_KIND_UPSERT",
2298 }
2299 }
2300 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2302 match value {
2303 "STREAM_KIND_RETRACT" => Some(Self::Retract),
2304 "STREAM_KIND_APPEND_ONLY" => Some(Self::AppendOnly),
2305 "STREAM_KIND_UPSERT" => Some(Self::Upsert),
2306 _ => None,
2307 }
2308 }
2309 }
2310 #[derive(prost_helpers::AnyPB)]
2311 #[derive(::enum_as_inner::EnumAsInner, ::strum::Display, ::strum::EnumDiscriminants)]
2312 #[derive(::prost_helpers::StreamNodeBodyVariants)]
2313 #[strum_discriminants(derive(::strum::Display, Hash))]
2314 #[derive(Clone, PartialEq, ::prost::Oneof)]
2315 pub enum NodeBody {
2316 #[prost(message, tag = "100")]
2317 Source(::prost::alloc::boxed::Box<super::SourceNode>),
2318 #[prost(message, tag = "101")]
2319 Project(::prost::alloc::boxed::Box<super::ProjectNode>),
2320 #[prost(message, tag = "102")]
2321 Filter(::prost::alloc::boxed::Box<super::FilterNode>),
2322 #[prost(message, tag = "103")]
2323 Materialize(::prost::alloc::boxed::Box<super::MaterializeNode>),
2324 #[prost(message, tag = "104")]
2325 StatelessSimpleAgg(::prost::alloc::boxed::Box<super::SimpleAggNode>),
2326 #[prost(message, tag = "105")]
2327 SimpleAgg(::prost::alloc::boxed::Box<super::SimpleAggNode>),
2328 #[prost(message, tag = "106")]
2329 HashAgg(::prost::alloc::boxed::Box<super::HashAggNode>),
2330 #[prost(message, tag = "107")]
2331 AppendOnlyTopN(::prost::alloc::boxed::Box<super::TopNNode>),
2332 #[prost(message, tag = "108")]
2333 HashJoin(::prost::alloc::boxed::Box<super::HashJoinNode>),
2334 #[prost(message, tag = "109")]
2335 TopN(::prost::alloc::boxed::Box<super::TopNNode>),
2336 #[prost(message, tag = "110")]
2337 HopWindow(::prost::alloc::boxed::Box<super::HopWindowNode>),
2338 #[prost(message, tag = "111")]
2339 Merge(::prost::alloc::boxed::Box<super::MergeNode>),
2340 #[prost(message, tag = "112")]
2341 Exchange(::prost::alloc::boxed::Box<super::ExchangeNode>),
2342 #[prost(message, tag = "113")]
2343 StreamScan(::prost::alloc::boxed::Box<super::StreamScanNode>),
2344 #[prost(message, tag = "114")]
2345 BatchPlan(::prost::alloc::boxed::Box<super::BatchPlanNode>),
2346 #[prost(message, tag = "115")]
2347 Lookup(::prost::alloc::boxed::Box<super::LookupNode>),
2348 #[prost(message, tag = "116")]
2349 Arrange(::prost::alloc::boxed::Box<super::ArrangeNode>),
2350 #[prost(message, tag = "117")]
2351 LookupUnion(::prost::alloc::boxed::Box<super::LookupUnionNode>),
2352 #[prost(message, tag = "118")]
2353 Union(super::UnionNode),
2354 #[prost(message, tag = "119")]
2355 DeltaIndexJoin(::prost::alloc::boxed::Box<super::DeltaIndexJoinNode>),
2356 #[prost(message, tag = "120")]
2357 Sink(::prost::alloc::boxed::Box<super::SinkNode>),
2358 #[prost(message, tag = "121")]
2359 Expand(::prost::alloc::boxed::Box<super::ExpandNode>),
2360 #[prost(message, tag = "122")]
2361 DynamicFilter(::prost::alloc::boxed::Box<super::DynamicFilterNode>),
2362 #[prost(message, tag = "123")]
2363 ProjectSet(::prost::alloc::boxed::Box<super::ProjectSetNode>),
2364 #[prost(message, tag = "124")]
2365 GroupTopN(::prost::alloc::boxed::Box<super::GroupTopNNode>),
2366 #[prost(message, tag = "125")]
2367 Sort(::prost::alloc::boxed::Box<super::SortNode>),
2368 #[prost(message, tag = "126")]
2369 WatermarkFilter(::prost::alloc::boxed::Box<super::WatermarkFilterNode>),
2370 #[prost(message, tag = "127")]
2371 Dml(::prost::alloc::boxed::Box<super::DmlNode>),
2372 #[prost(message, tag = "128")]
2373 RowIdGen(::prost::alloc::boxed::Box<super::RowIdGenNode>),
2374 #[prost(message, tag = "129")]
2375 Now(::prost::alloc::boxed::Box<super::NowNode>),
2376 #[prost(message, tag = "130")]
2377 AppendOnlyGroupTopN(::prost::alloc::boxed::Box<super::GroupTopNNode>),
2378 #[prost(message, tag = "131")]
2379 TemporalJoin(::prost::alloc::boxed::Box<super::TemporalJoinNode>),
2380 #[prost(message, tag = "132")]
2381 BarrierRecv(::prost::alloc::boxed::Box<super::BarrierRecvNode>),
2382 #[prost(message, tag = "133")]
2383 Values(::prost::alloc::boxed::Box<super::ValuesNode>),
2384 #[prost(message, tag = "134")]
2385 AppendOnlyDedup(::prost::alloc::boxed::Box<super::DedupNode>),
2386 #[prost(message, tag = "135")]
2387 NoOp(super::NoOpNode),
2388 #[prost(message, tag = "136")]
2389 EowcOverWindow(::prost::alloc::boxed::Box<super::EowcOverWindowNode>),
2390 #[prost(message, tag = "137")]
2391 OverWindow(::prost::alloc::boxed::Box<super::OverWindowNode>),
2392 #[prost(message, tag = "138")]
2393 StreamFsFetch(::prost::alloc::boxed::Box<super::StreamFsFetchNode>),
2394 #[prost(message, tag = "139")]
2395 StreamCdcScan(::prost::alloc::boxed::Box<super::StreamCdcScanNode>),
2396 #[prost(message, tag = "140")]
2397 CdcFilter(::prost::alloc::boxed::Box<super::CdcFilterNode>),
2398 #[prost(message, tag = "142")]
2399 SourceBackfill(::prost::alloc::boxed::Box<super::SourceBackfillNode>),
2400 #[prost(message, tag = "143")]
2401 Changelog(::prost::alloc::boxed::Box<super::ChangeLogNode>),
2402 #[prost(message, tag = "144")]
2403 LocalApproxPercentile(
2404 ::prost::alloc::boxed::Box<super::LocalApproxPercentileNode>,
2405 ),
2406 #[prost(message, tag = "145")]
2407 GlobalApproxPercentile(
2408 ::prost::alloc::boxed::Box<super::GlobalApproxPercentileNode>,
2409 ),
2410 #[prost(message, tag = "146")]
2411 RowMerge(::prost::alloc::boxed::Box<super::RowMergeNode>),
2412 #[prost(message, tag = "147")]
2413 AsOfJoin(::prost::alloc::boxed::Box<super::AsOfJoinNode>),
2414 #[prost(message, tag = "148")]
2415 SyncLogStore(::prost::alloc::boxed::Box<super::SyncLogStoreNode>),
2416 #[prost(message, tag = "149")]
2417 MaterializedExprs(::prost::alloc::boxed::Box<super::MaterializedExprsNode>),
2418 #[prost(message, tag = "150")]
2419 VectorIndexWrite(::prost::alloc::boxed::Box<super::VectorIndexWriteNode>),
2420 #[prost(message, tag = "151")]
2421 UpstreamSinkUnion(::prost::alloc::boxed::Box<super::UpstreamSinkUnionNode>),
2422 #[prost(message, tag = "152")]
2423 LocalityProvider(::prost::alloc::boxed::Box<super::LocalityProviderNode>),
2424 #[prost(message, tag = "153")]
2425 EowcGapFill(::prost::alloc::boxed::Box<super::EowcGapFillNode>),
2426 #[prost(message, tag = "154")]
2427 GapFill(::prost::alloc::boxed::Box<super::GapFillNode>),
2428 #[prost(message, tag = "155")]
2429 VectorIndexLookupJoin(
2430 ::prost::alloc::boxed::Box<super::VectorIndexLookupJoinNode>,
2431 ),
2432 #[prost(message, tag = "156")]
2433 IcebergWithPkIndexWriter(
2434 ::prost::alloc::boxed::Box<super::IcebergWithPkIndexWriterNode>,
2435 ),
2436 #[prost(message, tag = "157")]
2437 IcebergWithPkIndexPositionDeleteMerger(
2438 ::prost::alloc::boxed::Box<super::IcebergWithPkIndexPositionDeleteMergerNode>,
2439 ),
2440 #[prost(message, tag = "158")]
2441 MatchRecognize(::prost::alloc::boxed::Box<super::MatchRecognizeNode>),
2442 #[prost(message, tag = "159")]
2443 CompactionResolver(::prost::alloc::boxed::Box<super::CompactionResolverNode>),
2444 }
2445}
2446#[derive(prost_helpers::AnyPB)]
2460#[derive(Clone, PartialEq, ::prost::Message)]
2461pub struct DispatchOutputMapping {
2462 #[prost(uint32, repeated, tag = "1")]
2464 pub indices: ::prost::alloc::vec::Vec<u32>,
2465 #[prost(message, repeated, tag = "2")]
2471 pub types: ::prost::alloc::vec::Vec<dispatch_output_mapping::TypePair>,
2472}
2473pub mod dispatch_output_mapping {
2475 #[derive(prost_helpers::AnyPB)]
2476 #[derive(Clone, PartialEq, ::prost::Message)]
2477 pub struct TypePair {
2478 #[prost(message, optional, tag = "1")]
2479 pub upstream: ::core::option::Option<super::super::data::DataType>,
2480 #[prost(message, optional, tag = "2")]
2481 pub downstream: ::core::option::Option<super::super::data::DataType>,
2482 }
2483}
2484#[derive(prost_helpers::AnyPB)]
2487#[derive(Clone, PartialEq, ::prost::Message)]
2488pub struct DispatchStrategy {
2489 #[prost(enumeration = "DispatcherType", tag = "1")]
2490 pub r#type: i32,
2491 #[prost(uint32, repeated, tag = "2")]
2492 pub dist_key_indices: ::prost::alloc::vec::Vec<u32>,
2493 #[prost(message, optional, tag = "3")]
2494 pub output_mapping: ::core::option::Option<DispatchOutputMapping>,
2495}
2496#[derive(prost_helpers::AnyPB)]
2499#[derive(Clone, PartialEq, ::prost::Message)]
2500pub struct Dispatcher {
2501 #[prost(enumeration = "DispatcherType", tag = "1")]
2502 pub r#type: i32,
2503 #[prost(uint32, repeated, tag = "2")]
2506 pub dist_key_indices: ::prost::alloc::vec::Vec<u32>,
2507 #[prost(message, optional, tag = "6")]
2509 pub output_mapping: ::core::option::Option<DispatchOutputMapping>,
2510 #[prost(message, optional, tag = "3")]
2513 pub hash_mapping: ::core::option::Option<ActorMapping>,
2514 #[prost(uint32, tag = "4", wrapper = "crate::id::FragmentId")]
2517 pub dispatcher_id: crate::id::FragmentId,
2518 #[prost(uint32, repeated, tag = "5", wrapper = "crate::id::ActorId")]
2520 pub downstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
2521}
2522#[derive(prost_helpers::AnyPB)]
2524#[derive(Clone, PartialEq, ::prost::Message)]
2525pub struct StreamActor {
2526 #[prost(uint32, tag = "1", wrapper = "crate::id::ActorId")]
2527 pub actor_id: crate::id::ActorId,
2528 #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
2529 pub fragment_id: crate::id::FragmentId,
2530 #[prost(message, repeated, tag = "4")]
2531 pub dispatcher: ::prost::alloc::vec::Vec<Dispatcher>,
2532 #[prost(message, optional, tag = "8")]
2535 pub vnode_bitmap: ::core::option::Option<super::common::Buffer>,
2536 #[prost(string, tag = "9")]
2538 pub mview_definition: ::prost::alloc::string::String,
2539 #[prost(message, optional, tag = "10")]
2541 pub expr_context: ::core::option::Option<super::plan_common::ExprContext>,
2542 #[prost(string, tag = "11")]
2544 pub config_override: ::prost::alloc::string::String,
2545}
2546#[derive(prost_helpers::AnyPB)]
2548#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2549pub struct StreamContext {
2550 #[prost(string, tag = "1")]
2552 pub timezone: ::prost::alloc::string::String,
2553 #[prost(string, tag = "2")]
2555 pub config_override: ::prost::alloc::string::String,
2556}
2557#[derive(prost_helpers::AnyPB)]
2558#[derive(Clone, PartialEq, ::prost::Message)]
2559pub struct BackfillOrder {
2560 #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::RelationId")]
2561 pub order: ::std::collections::HashMap<
2562 crate::id::RelationId,
2563 super::common::Uint32Vector,
2564 >,
2565}
2566#[derive(prost_helpers::AnyPB)]
2571#[derive(Clone, PartialEq, ::prost::Message)]
2572pub struct StreamFragmentGraph {
2573 #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::FragmentId")]
2575 pub fragments: ::std::collections::HashMap<
2576 crate::id::FragmentId,
2577 stream_fragment_graph::StreamFragment,
2578 >,
2579 #[prost(message, repeated, tag = "2")]
2581 pub edges: ::prost::alloc::vec::Vec<stream_fragment_graph::StreamFragmentEdge>,
2582 #[prost(uint32, repeated, tag = "3", wrapper = "crate::id::TableId")]
2583 pub dependent_table_ids: ::prost::alloc::vec::Vec<crate::id::TableId>,
2584 #[prost(uint32, tag = "4")]
2585 pub table_ids_cnt: u32,
2586 #[prost(message, optional, tag = "5")]
2587 pub ctx: ::core::option::Option<StreamContext>,
2588 #[prost(message, optional, tag = "6")]
2590 pub parallelism: ::core::option::Option<stream_fragment_graph::Parallelism>,
2591 #[prost(message, optional, tag = "9")]
2593 pub backfill_parallelism: ::core::option::Option<stream_fragment_graph::Parallelism>,
2594 #[prost(string, tag = "10")]
2596 pub adaptive_parallelism_strategy: ::prost::alloc::string::String,
2597 #[prost(string, tag = "11")]
2599 pub backfill_adaptive_parallelism_strategy: ::prost::alloc::string::String,
2600 #[prost(uint32, tag = "7")]
2610 pub max_parallelism: u32,
2611 #[prost(message, optional, tag = "8")]
2613 pub backfill_order: ::core::option::Option<BackfillOrder>,
2614}
2615pub mod stream_fragment_graph {
2617 #[derive(prost_helpers::AnyPB)]
2618 #[derive(Clone, PartialEq, ::prost::Message)]
2619 pub struct StreamFragment {
2620 #[prost(uint32, tag = "1", wrapper = "crate::id::FragmentId")]
2622 pub fragment_id: crate::id::FragmentId,
2623 #[prost(message, optional, tag = "2")]
2625 pub node: ::core::option::Option<super::StreamNode>,
2626 #[prost(uint32, tag = "3")]
2628 pub fragment_type_mask: u32,
2629 #[prost(bool, tag = "4")]
2633 pub requires_singleton: bool,
2634 }
2635 #[derive(prost_helpers::AnyPB)]
2636 #[derive(Clone, PartialEq, ::prost::Message)]
2637 pub struct StreamFragmentEdge {
2638 #[prost(message, optional, tag = "1")]
2640 pub dispatch_strategy: ::core::option::Option<super::DispatchStrategy>,
2641 #[prost(uint64, tag = "3")]
2645 pub link_id: u64,
2646 #[prost(uint32, tag = "4", wrapper = "crate::id::FragmentId")]
2647 pub upstream_id: crate::id::FragmentId,
2648 #[prost(uint32, tag = "5", wrapper = "crate::id::FragmentId")]
2649 pub downstream_id: crate::id::FragmentId,
2650 }
2651 #[derive(prost_helpers::AnyPB)]
2652 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
2653 pub struct Parallelism {
2654 #[prost(uint64, tag = "1")]
2655 pub parallelism: u64,
2656 }
2657}
2658#[derive(prost_helpers::AnyPB)]
2660#[derive(Clone, PartialEq, ::prost::Message)]
2661pub struct SinkSchemaChange {
2662 #[prost(message, repeated, tag = "1")]
2665 pub original_schema: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2666 #[prost(oneof = "sink_schema_change::Op", tags = "2, 3")]
2668 pub op: ::core::option::Option<sink_schema_change::Op>,
2669}
2670pub mod sink_schema_change {
2672 #[derive(prost_helpers::AnyPB)]
2674 #[derive(Clone, PartialEq, ::prost::Oneof)]
2675 pub enum Op {
2676 #[prost(message, tag = "2")]
2678 AddColumns(super::SinkAddColumnsOp),
2679 #[prost(message, tag = "3")]
2681 DropColumns(super::SinkDropColumnsOp),
2682 }
2683}
2684#[derive(prost_helpers::AnyPB)]
2686#[derive(Clone, PartialEq, ::prost::Message)]
2687pub struct SinkAddColumnsOp {
2688 #[prost(message, repeated, tag = "1")]
2690 pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2691}
2692#[derive(prost_helpers::AnyPB)]
2694#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2695pub struct SinkDropColumnsOp {
2696 #[prost(string, repeated, tag = "1")]
2698 pub column_names: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
2699}
2700#[derive(prost_helpers::AnyPB)]
2701#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2702#[repr(i32)]
2703pub enum SinkLogStoreType {
2704 Unspecified = 0,
2706 KvLogStore = 1,
2707 InMemoryLogStore = 2,
2708}
2709impl SinkLogStoreType {
2710 pub fn as_str_name(&self) -> &'static str {
2715 match self {
2716 Self::Unspecified => "SINK_LOG_STORE_TYPE_UNSPECIFIED",
2717 Self::KvLogStore => "SINK_LOG_STORE_TYPE_KV_LOG_STORE",
2718 Self::InMemoryLogStore => "SINK_LOG_STORE_TYPE_IN_MEMORY_LOG_STORE",
2719 }
2720 }
2721 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2723 match value {
2724 "SINK_LOG_STORE_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2725 "SINK_LOG_STORE_TYPE_KV_LOG_STORE" => Some(Self::KvLogStore),
2726 "SINK_LOG_STORE_TYPE_IN_MEMORY_LOG_STORE" => Some(Self::InMemoryLogStore),
2727 _ => None,
2728 }
2729 }
2730}
2731#[derive(prost_helpers::AnyPB)]
2732#[derive(prost_helpers::Version)]
2733#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2734#[repr(i32)]
2735pub enum AggNodeVersion {
2736 Unspecified = 0,
2737 Issue12140 = 1,
2739 Issue13465 = 2,
2741}
2742impl AggNodeVersion {
2743 pub fn as_str_name(&self) -> &'static str {
2748 match self {
2749 Self::Unspecified => "AGG_NODE_VERSION_UNSPECIFIED",
2750 Self::Issue12140 => "AGG_NODE_VERSION_ISSUE_12140",
2751 Self::Issue13465 => "AGG_NODE_VERSION_ISSUE_13465",
2752 }
2753 }
2754 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2756 match value {
2757 "AGG_NODE_VERSION_UNSPECIFIED" => Some(Self::Unspecified),
2758 "AGG_NODE_VERSION_ISSUE_12140" => Some(Self::Issue12140),
2759 "AGG_NODE_VERSION_ISSUE_13465" => Some(Self::Issue13465),
2760 _ => None,
2761 }
2762 }
2763}
2764#[derive(prost_helpers::AnyPB)]
2765#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2766#[repr(i32)]
2767pub enum InequalityType {
2768 Unspecified = 0,
2769 LessThan = 1,
2770 LessThanOrEqual = 2,
2771 GreaterThan = 3,
2772 GreaterThanOrEqual = 4,
2773}
2774impl InequalityType {
2775 pub fn as_str_name(&self) -> &'static str {
2780 match self {
2781 Self::Unspecified => "INEQUALITY_TYPE_UNSPECIFIED",
2782 Self::LessThan => "INEQUALITY_TYPE_LESS_THAN",
2783 Self::LessThanOrEqual => "INEQUALITY_TYPE_LESS_THAN_OR_EQUAL",
2784 Self::GreaterThan => "INEQUALITY_TYPE_GREATER_THAN",
2785 Self::GreaterThanOrEqual => "INEQUALITY_TYPE_GREATER_THAN_OR_EQUAL",
2786 }
2787 }
2788 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2790 match value {
2791 "INEQUALITY_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2792 "INEQUALITY_TYPE_LESS_THAN" => Some(Self::LessThan),
2793 "INEQUALITY_TYPE_LESS_THAN_OR_EQUAL" => Some(Self::LessThanOrEqual),
2794 "INEQUALITY_TYPE_GREATER_THAN" => Some(Self::GreaterThan),
2795 "INEQUALITY_TYPE_GREATER_THAN_OR_EQUAL" => Some(Self::GreaterThanOrEqual),
2796 _ => None,
2797 }
2798 }
2799}
2800#[derive(prost_helpers::AnyPB)]
2801#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2802#[repr(i32)]
2803pub enum JoinEncodingType {
2804 Unspecified = 0,
2805 MemoryOptimized = 1,
2806 CpuOptimized = 2,
2807}
2808impl JoinEncodingType {
2809 pub fn as_str_name(&self) -> &'static str {
2814 match self {
2815 Self::Unspecified => "UNSPECIFIED",
2816 Self::MemoryOptimized => "MEMORY_OPTIMIZED",
2817 Self::CpuOptimized => "CPU_OPTIMIZED",
2818 }
2819 }
2820 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2822 match value {
2823 "UNSPECIFIED" => Some(Self::Unspecified),
2824 "MEMORY_OPTIMIZED" => Some(Self::MemoryOptimized),
2825 "CPU_OPTIMIZED" => Some(Self::CpuOptimized),
2826 _ => None,
2827 }
2828 }
2829}
2830#[derive(prost_helpers::AnyPB)]
2832#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2833#[repr(i32)]
2834pub enum StreamScanType {
2835 Unspecified = 0,
2836 #[deprecated]
2838 Chain = 1,
2839 #[deprecated]
2841 Rearrange = 2,
2842 #[deprecated]
2844 Backfill = 3,
2845 UpstreamOnly = 4,
2847 ArrangementBackfill = 5,
2849 SnapshotBackfill = 6,
2851 CrossDbSnapshotBackfill = 7,
2853}
2854impl StreamScanType {
2855 pub fn as_str_name(&self) -> &'static str {
2860 match self {
2861 Self::Unspecified => "STREAM_SCAN_TYPE_UNSPECIFIED",
2862 #[allow(deprecated)]
2863 Self::Chain => "STREAM_SCAN_TYPE_CHAIN",
2864 #[allow(deprecated)]
2865 Self::Rearrange => "STREAM_SCAN_TYPE_REARRANGE",
2866 #[allow(deprecated)]
2867 Self::Backfill => "STREAM_SCAN_TYPE_BACKFILL",
2868 Self::UpstreamOnly => "STREAM_SCAN_TYPE_UPSTREAM_ONLY",
2869 Self::ArrangementBackfill => "STREAM_SCAN_TYPE_ARRANGEMENT_BACKFILL",
2870 Self::SnapshotBackfill => "STREAM_SCAN_TYPE_SNAPSHOT_BACKFILL",
2871 Self::CrossDbSnapshotBackfill => {
2872 "STREAM_SCAN_TYPE_CROSS_DB_SNAPSHOT_BACKFILL"
2873 }
2874 }
2875 }
2876 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2878 match value {
2879 "STREAM_SCAN_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2880 "STREAM_SCAN_TYPE_CHAIN" => Some(#[allow(deprecated)] Self::Chain),
2881 "STREAM_SCAN_TYPE_REARRANGE" => Some(#[allow(deprecated)] Self::Rearrange),
2882 "STREAM_SCAN_TYPE_BACKFILL" => Some(#[allow(deprecated)] Self::Backfill),
2883 "STREAM_SCAN_TYPE_UPSTREAM_ONLY" => Some(Self::UpstreamOnly),
2884 "STREAM_SCAN_TYPE_ARRANGEMENT_BACKFILL" => Some(Self::ArrangementBackfill),
2885 "STREAM_SCAN_TYPE_SNAPSHOT_BACKFILL" => Some(Self::SnapshotBackfill),
2886 "STREAM_SCAN_TYPE_CROSS_DB_SNAPSHOT_BACKFILL" => {
2887 Some(Self::CrossDbSnapshotBackfill)
2888 }
2889 _ => None,
2890 }
2891 }
2892}
2893#[derive(prost_helpers::AnyPB)]
2894#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2895#[repr(i32)]
2896pub enum OverWindowCachePolicy {
2897 Unspecified = 0,
2898 Full = 1,
2899 Recent = 2,
2900 RecentFirstN = 3,
2901 RecentLastN = 4,
2902}
2903impl OverWindowCachePolicy {
2904 pub fn as_str_name(&self) -> &'static str {
2909 match self {
2910 Self::Unspecified => "OVER_WINDOW_CACHE_POLICY_UNSPECIFIED",
2911 Self::Full => "OVER_WINDOW_CACHE_POLICY_FULL",
2912 Self::Recent => "OVER_WINDOW_CACHE_POLICY_RECENT",
2913 Self::RecentFirstN => "OVER_WINDOW_CACHE_POLICY_RECENT_FIRST_N",
2914 Self::RecentLastN => "OVER_WINDOW_CACHE_POLICY_RECENT_LAST_N",
2915 }
2916 }
2917 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2919 match value {
2920 "OVER_WINDOW_CACHE_POLICY_UNSPECIFIED" => Some(Self::Unspecified),
2921 "OVER_WINDOW_CACHE_POLICY_FULL" => Some(Self::Full),
2922 "OVER_WINDOW_CACHE_POLICY_RECENT" => Some(Self::Recent),
2923 "OVER_WINDOW_CACHE_POLICY_RECENT_FIRST_N" => Some(Self::RecentFirstN),
2924 "OVER_WINDOW_CACHE_POLICY_RECENT_LAST_N" => Some(Self::RecentLastN),
2925 _ => None,
2926 }
2927 }
2928}
2929#[derive(prost_helpers::AnyPB)]
2930#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2931#[repr(i32)]
2932pub enum MatchRecognizeInputMode {
2933 Unspecified = 0,
2934 EventTime = 1,
2935 ProcessingTime = 2,
2936}
2937impl MatchRecognizeInputMode {
2938 pub fn as_str_name(&self) -> &'static str {
2943 match self {
2944 Self::Unspecified => "MATCH_RECOGNIZE_INPUT_MODE_UNSPECIFIED",
2945 Self::EventTime => "MATCH_RECOGNIZE_INPUT_MODE_EVENT_TIME",
2946 Self::ProcessingTime => "MATCH_RECOGNIZE_INPUT_MODE_PROCESSING_TIME",
2947 }
2948 }
2949 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2951 match value {
2952 "MATCH_RECOGNIZE_INPUT_MODE_UNSPECIFIED" => Some(Self::Unspecified),
2953 "MATCH_RECOGNIZE_INPUT_MODE_EVENT_TIME" => Some(Self::EventTime),
2954 "MATCH_RECOGNIZE_INPUT_MODE_PROCESSING_TIME" => Some(Self::ProcessingTime),
2955 _ => None,
2956 }
2957 }
2958}
2959#[derive(prost_helpers::AnyPB)]
2960#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2961#[repr(i32)]
2962pub enum DispatcherType {
2963 Unspecified = 0,
2964 Hash = 1,
2966 Broadcast = 2,
2971 Simple = 3,
2973 NoShuffle = 4,
2977}
2978impl DispatcherType {
2979 pub fn as_str_name(&self) -> &'static str {
2984 match self {
2985 Self::Unspecified => "DISPATCHER_TYPE_UNSPECIFIED",
2986 Self::Hash => "DISPATCHER_TYPE_HASH",
2987 Self::Broadcast => "DISPATCHER_TYPE_BROADCAST",
2988 Self::Simple => "DISPATCHER_TYPE_SIMPLE",
2989 Self::NoShuffle => "DISPATCHER_TYPE_NO_SHUFFLE",
2990 }
2991 }
2992 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2994 match value {
2995 "DISPATCHER_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2996 "DISPATCHER_TYPE_HASH" => Some(Self::Hash),
2997 "DISPATCHER_TYPE_BROADCAST" => Some(Self::Broadcast),
2998 "DISPATCHER_TYPE_SIMPLE" => Some(Self::Simple),
2999 "DISPATCHER_TYPE_NO_SHUFFLE" => Some(Self::NoShuffle),
3000 _ => None,
3001 }
3002 }
3003}