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}
124pub mod update_mutation {
126 #[derive(prost_helpers::AnyPB)]
127 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
128 pub struct DispatcherUpdate {
129 #[prost(uint32, tag = "1", wrapper = "crate::id::ActorId")]
131 pub actor_id: crate::id::ActorId,
132 #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
133 pub dispatcher_id: crate::id::FragmentId,
134 #[prost(message, optional, tag = "3")]
137 pub hash_mapping: ::core::option::Option<super::ActorMapping>,
138 #[prost(uint32, repeated, tag = "4", wrapper = "crate::id::ActorId")]
140 pub added_downstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
141 #[prost(uint32, repeated, tag = "5", wrapper = "crate::id::ActorId")]
143 pub removed_downstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
144 }
145 #[derive(prost_helpers::AnyPB)]
146 #[derive(Clone, PartialEq, ::prost::Message)]
147 pub struct MergeUpdate {
148 #[prost(uint32, tag = "1", wrapper = "crate::id::ActorId")]
150 pub actor_id: crate::id::ActorId,
151 #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
152 pub upstream_fragment_id: crate::id::FragmentId,
153 #[prost(uint32, optional, tag = "5", wrapper = "crate::id::FragmentId")]
157 pub new_upstream_fragment_id: ::core::option::Option<crate::id::FragmentId>,
158 #[prost(message, repeated, tag = "3")]
160 pub added_upstream_actors: ::prost::alloc::vec::Vec<
161 super::super::common::ActorInfo,
162 >,
163 #[prost(uint32, repeated, tag = "4", wrapper = "crate::id::ActorId")]
166 pub removed_upstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
167 }
168}
169#[derive(prost_helpers::AnyPB)]
170#[derive(Clone, PartialEq, ::prost::Message)]
171pub struct SourceChangeSplitMutation {
172 #[prost(map = "uint32, message", tag = "2", wrapper = "crate::id::ActorId")]
174 pub actor_splits: ::std::collections::HashMap<
175 crate::id::ActorId,
176 super::source::ConnectorSplits,
177 >,
178}
179#[derive(prost_helpers::AnyPB)]
180#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
181pub struct PauseMutation {}
182#[derive(prost_helpers::AnyPB)]
183#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
184pub struct ResumeMutation {}
185#[derive(prost_helpers::AnyPB)]
186#[derive(Clone, PartialEq, ::prost::Message)]
187pub struct ThrottleMutation {
188 #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::FragmentId")]
189 pub fragment_throttle: ::std::collections::HashMap<
190 crate::id::FragmentId,
191 throttle_mutation::ThrottleConfig,
192 >,
193}
194pub mod throttle_mutation {
196 #[derive(prost_helpers::AnyPB)]
197 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
198 pub struct ThrottleConfig {
199 #[prost(uint32, optional, tag = "1")]
200 pub rate_limit: ::core::option::Option<u32>,
201 #[prost(enumeration = "super::super::common::ThrottleType", tag = "2")]
202 pub throttle_type: i32,
203 }
204}
205#[derive(prost_helpers::AnyPB)]
206#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
207pub struct SubscriptionUpstreamInfo {
208 #[prost(uint32, tag = "1", wrapper = "crate::id::SubscriberId")]
210 pub subscriber_id: crate::id::SubscriberId,
211 #[prost(uint32, tag = "2", wrapper = "crate::id::TableId")]
212 pub upstream_mv_table_id: crate::id::TableId,
213}
214#[derive(prost_helpers::AnyPB)]
215#[derive(Clone, PartialEq, ::prost::Message)]
216pub struct DropSubscriptionsMutation {
217 #[prost(message, repeated, tag = "1")]
218 pub info: ::prost::alloc::vec::Vec<SubscriptionUpstreamInfo>,
219}
220#[derive(prost_helpers::AnyPB)]
221#[derive(Clone, PartialEq, ::prost::Message)]
222pub struct ConnectorPropsChangeMutation {
223 #[prost(map = "uint32, message", tag = "1")]
224 pub connector_props_infos: ::std::collections::HashMap<
225 u32,
226 connector_props_change_mutation::ConnectorPropsInfo,
227 >,
228}
229pub mod connector_props_change_mutation {
231 #[derive(prost_helpers::AnyPB)]
232 #[derive(Clone, PartialEq, ::prost::Message)]
233 pub struct ConnectorPropsInfo {
234 #[prost(map = "string, string", tag = "1")]
235 pub connector_props_info: ::std::collections::HashMap<
236 ::prost::alloc::string::String,
237 ::prost::alloc::string::String,
238 >,
239 }
240}
241#[derive(prost_helpers::AnyPB)]
242#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
243pub struct StartFragmentBackfillMutation {
244 #[prost(uint32, repeated, tag = "1", wrapper = "crate::id::FragmentId")]
245 pub fragment_ids: ::prost::alloc::vec::Vec<crate::id::FragmentId>,
246}
247#[derive(prost_helpers::AnyPB)]
248#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
249pub struct RefreshStartMutation {
250 #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
252 pub table_id: crate::id::TableId,
253 #[prost(uint32, tag = "2", wrapper = "crate::id::SourceId")]
255 pub associated_source_id: crate::id::SourceId,
256}
257#[derive(prost_helpers::AnyPB)]
258#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
259pub struct ListFinishMutation {
260 #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
262 pub associated_source_id: crate::id::SourceId,
263}
264#[derive(prost_helpers::AnyPB)]
265#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
266pub struct LoadFinishMutation {
267 #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
269 pub associated_source_id: crate::id::SourceId,
270}
271#[derive(prost_helpers::AnyPB)]
272#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
273pub struct ResetSourceMutation {
274 #[prost(uint32, tag = "1")]
276 pub source_id: u32,
277}
278#[derive(prost_helpers::AnyPB)]
281#[derive(Clone, PartialEq, ::prost::Message)]
282pub struct InjectSourceOffsetsMutation {
283 #[prost(uint32, tag = "1")]
285 pub source_id: u32,
286 #[prost(map = "string, string", tag = "2")]
288 pub split_offsets: ::std::collections::HashMap<
289 ::prost::alloc::string::String,
290 ::prost::alloc::string::String,
291 >,
292}
293#[derive(prost_helpers::AnyPB)]
295#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
296pub struct IcebergPkIndexCompactionContext {
297 #[prost(uint32, tag = "1", wrapper = "crate::id::SinkId")]
298 pub sink_id: crate::id::SinkId,
299 #[prost(uint64, tag = "2", wrapper = "crate::id::IcebergCompactionTaskId")]
300 pub task_id: crate::id::IcebergCompactionTaskId,
301 #[prost(enumeration = "iceberg_pk_index_compaction_context::Phase", tag = "3")]
302 pub phase: i32,
303}
304pub mod iceberg_pk_index_compaction_context {
306 #[derive(prost_helpers::AnyPB)]
307 #[derive(
308 Clone,
309 Copy,
310 Debug,
311 PartialEq,
312 Eq,
313 Hash,
314 PartialOrd,
315 Ord,
316 ::prost::Enumeration
317 )]
318 #[repr(i32)]
319 pub enum Phase {
320 Unspecified = 0,
321 Begin = 1,
322 End = 2,
323 }
324 impl Phase {
325 pub fn as_str_name(&self) -> &'static str {
330 match self {
331 Self::Unspecified => "PHASE_UNSPECIFIED",
332 Self::Begin => "PHASE_BEGIN",
333 Self::End => "PHASE_END",
334 }
335 }
336 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
338 match value {
339 "PHASE_UNSPECIFIED" => Some(Self::Unspecified),
340 "PHASE_BEGIN" => Some(Self::Begin),
341 "PHASE_END" => Some(Self::End),
342 _ => None,
343 }
344 }
345 }
346}
347#[derive(prost_helpers::AnyPB)]
348#[derive(Clone, PartialEq, ::prost::Message)]
349pub struct BarrierMutation {
350 #[prost(
351 oneof = "barrier_mutation::Mutation",
352 tags = "3, 4, 5, 6, 7, 8, 10, 12, 13, 14, 15, 16, 17, 18, 19"
353 )]
354 pub mutation: ::core::option::Option<barrier_mutation::Mutation>,
355}
356pub mod barrier_mutation {
358 #[derive(prost_helpers::AnyPB)]
359 #[derive(Clone, PartialEq, ::prost::Oneof)]
360 pub enum Mutation {
361 #[prost(message, tag = "3")]
363 Add(super::AddMutation),
364 #[prost(message, tag = "4")]
367 Stop(super::StopMutation),
368 #[prost(message, tag = "5")]
370 Update(super::UpdateMutation),
371 #[prost(message, tag = "6")]
373 Splits(super::SourceChangeSplitMutation),
374 #[prost(message, tag = "7")]
376 Pause(super::PauseMutation),
377 #[prost(message, tag = "8")]
379 Resume(super::ResumeMutation),
380 #[prost(message, tag = "10")]
382 Throttle(super::ThrottleMutation),
383 #[prost(message, tag = "12")]
385 DropSubscriptions(super::DropSubscriptionsMutation),
386 #[prost(message, tag = "13")]
388 ConnectorPropsChange(super::ConnectorPropsChangeMutation),
389 #[prost(message, tag = "14")]
395 StartFragmentBackfill(super::StartFragmentBackfillMutation),
396 #[prost(message, tag = "15")]
398 RefreshStart(super::RefreshStartMutation),
399 #[prost(message, tag = "16")]
401 LoadFinish(super::LoadFinishMutation),
402 #[prost(message, tag = "17")]
404 ListFinish(super::ListFinishMutation),
405 #[prost(message, tag = "18")]
407 ResetSource(super::ResetSourceMutation),
408 #[prost(message, tag = "19")]
410 InjectSourceOffsets(super::InjectSourceOffsetsMutation),
411 }
412}
413#[derive(prost_helpers::AnyPB)]
414#[derive(Clone, PartialEq, ::prost::Message)]
415pub struct Barrier {
416 #[prost(message, optional, tag = "1")]
417 pub epoch: ::core::option::Option<super::data::Epoch>,
418 #[prost(message, optional, tag = "3")]
419 pub mutation: ::core::option::Option<BarrierMutation>,
420 #[prost(map = "string, string", tag = "2")]
422 pub tracing_context: ::std::collections::HashMap<
423 ::prost::alloc::string::String,
424 ::prost::alloc::string::String,
425 >,
426 #[prost(enumeration = "barrier::BarrierKind", tag = "9")]
428 pub kind: i32,
429 #[prost(message, optional, tag = "10")]
430 pub iceberg_pk_index_compaction: ::core::option::Option<
431 IcebergPkIndexCompactionContext,
432 >,
433}
434pub mod barrier {
436 #[derive(prost_helpers::AnyPB)]
437 #[derive(::enum_as_inner::EnumAsInner)]
438 #[derive(
439 Clone,
440 Copy,
441 Debug,
442 PartialEq,
443 Eq,
444 Hash,
445 PartialOrd,
446 Ord,
447 ::prost::Enumeration
448 )]
449 #[repr(i32)]
450 pub enum BarrierKind {
451 Unspecified = 0,
452 Initial = 1,
455 Barrier = 2,
457 Checkpoint = 3,
459 }
460 impl BarrierKind {
461 pub fn as_str_name(&self) -> &'static str {
466 match self {
467 Self::Unspecified => "BARRIER_KIND_UNSPECIFIED",
468 Self::Initial => "BARRIER_KIND_INITIAL",
469 Self::Barrier => "BARRIER_KIND_BARRIER",
470 Self::Checkpoint => "BARRIER_KIND_CHECKPOINT",
471 }
472 }
473 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
475 match value {
476 "BARRIER_KIND_UNSPECIFIED" => Some(Self::Unspecified),
477 "BARRIER_KIND_INITIAL" => Some(Self::Initial),
478 "BARRIER_KIND_BARRIER" => Some(Self::Barrier),
479 "BARRIER_KIND_CHECKPOINT" => Some(Self::Checkpoint),
480 _ => None,
481 }
482 }
483 }
484}
485#[derive(prost_helpers::AnyPB)]
486#[derive(Clone, PartialEq, ::prost::Message)]
487pub struct Watermark {
488 #[prost(message, optional, tag = "1")]
490 pub column: ::core::option::Option<super::expr::InputRef>,
491 #[prost(message, optional, tag = "3")]
493 pub val: ::core::option::Option<super::data::Datum>,
494}
495#[derive(prost_helpers::AnyPB)]
496#[derive(Clone, PartialEq, ::prost::Message)]
497pub struct StreamMessage {
498 #[prost(oneof = "stream_message::StreamMessage", tags = "1, 2, 3")]
499 pub stream_message: ::core::option::Option<stream_message::StreamMessage>,
500}
501pub mod stream_message {
503 #[derive(prost_helpers::AnyPB)]
504 #[derive(Clone, PartialEq, ::prost::Oneof)]
505 pub enum StreamMessage {
506 #[prost(message, tag = "1")]
507 StreamChunk(super::super::data::StreamChunk),
508 #[prost(message, tag = "2")]
509 Barrier(super::Barrier),
510 #[prost(message, tag = "3")]
511 Watermark(super::Watermark),
512 }
513}
514#[derive(prost_helpers::AnyPB)]
515#[derive(Clone, PartialEq, ::prost::Message)]
516pub struct StreamMessageBatch {
517 #[prost(oneof = "stream_message_batch::StreamMessageBatch", tags = "1, 2, 3")]
518 pub stream_message_batch: ::core::option::Option<
519 stream_message_batch::StreamMessageBatch,
520 >,
521}
522pub mod stream_message_batch {
524 #[derive(prost_helpers::AnyPB)]
525 #[derive(Clone, PartialEq, ::prost::Message)]
526 pub struct BarrierBatch {
527 #[prost(message, repeated, tag = "1")]
528 pub barriers: ::prost::alloc::vec::Vec<super::Barrier>,
529 }
530 #[derive(prost_helpers::AnyPB)]
531 #[derive(Clone, PartialEq, ::prost::Oneof)]
532 pub enum StreamMessageBatch {
533 #[prost(message, tag = "1")]
534 StreamChunk(super::super::data::StreamChunk),
535 #[prost(message, tag = "2")]
536 BarrierBatch(BarrierBatch),
537 #[prost(message, tag = "3")]
538 Watermark(super::Watermark),
539 }
540}
541#[derive(prost_helpers::AnyPB)]
543#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
544pub struct ActorMapping {
545 #[prost(uint32, repeated, tag = "1")]
546 pub original_indices: ::prost::alloc::vec::Vec<u32>,
547 #[prost(uint32, repeated, tag = "2", wrapper = "crate::id::ActorId")]
548 pub data: ::prost::alloc::vec::Vec<crate::id::ActorId>,
549}
550#[derive(prost_helpers::AnyPB)]
551#[derive(Clone, PartialEq, ::prost::Message)]
552pub struct Columns {
553 #[prost(message, repeated, tag = "1")]
554 pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
555}
556#[derive(prost_helpers::AnyPB)]
557#[derive(Clone, PartialEq, ::prost::Message)]
558pub struct StreamSource {
559 #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
560 pub source_id: crate::id::SourceId,
561 #[prost(message, optional, tag = "2")]
562 pub state_table: ::core::option::Option<super::catalog::Table>,
563 #[prost(uint32, optional, tag = "3")]
564 pub row_id_index: ::core::option::Option<u32>,
565 #[prost(message, repeated, tag = "4")]
566 pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
567 #[prost(btree_map = "string, string", tag = "6")]
568 pub with_properties: ::prost::alloc::collections::BTreeMap<
569 ::prost::alloc::string::String,
570 ::prost::alloc::string::String,
571 >,
572 #[prost(message, optional, tag = "7")]
573 pub info: ::core::option::Option<super::catalog::StreamSourceInfo>,
574 #[prost(string, tag = "8")]
575 pub source_name: ::prost::alloc::string::String,
576 #[prost(uint32, optional, tag = "9")]
578 pub rate_limit: ::core::option::Option<u32>,
579 #[prost(btree_map = "string, message", tag = "10")]
580 pub secret_refs: ::prost::alloc::collections::BTreeMap<
581 ::prost::alloc::string::String,
582 super::secret::SecretRef,
583 >,
584 #[prost(message, optional, tag = "11")]
586 pub downstream_columns: ::core::option::Option<Columns>,
587 #[prost(message, optional, tag = "12")]
588 pub refresh_mode: ::core::option::Option<super::plan_common::SourceRefreshMode>,
589 #[prost(uint32, optional, tag = "13", wrapper = "crate::id::TableId")]
590 pub associated_table_id: ::core::option::Option<crate::id::TableId>,
591}
592#[derive(prost_helpers::AnyPB)]
594#[derive(Clone, PartialEq, ::prost::Message)]
595pub struct StreamFsFetch {
596 #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
597 pub source_id: crate::id::SourceId,
598 #[prost(message, optional, tag = "2")]
599 pub state_table: ::core::option::Option<super::catalog::Table>,
600 #[prost(uint32, optional, tag = "3")]
601 pub row_id_index: ::core::option::Option<u32>,
602 #[prost(message, repeated, tag = "4")]
603 pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
604 #[prost(btree_map = "string, string", tag = "6")]
605 pub with_properties: ::prost::alloc::collections::BTreeMap<
606 ::prost::alloc::string::String,
607 ::prost::alloc::string::String,
608 >,
609 #[prost(message, optional, tag = "7")]
610 pub info: ::core::option::Option<super::catalog::StreamSourceInfo>,
611 #[prost(string, tag = "8")]
612 pub source_name: ::prost::alloc::string::String,
613 #[prost(uint32, optional, tag = "9")]
615 pub rate_limit: ::core::option::Option<u32>,
616 #[prost(btree_map = "string, message", tag = "10")]
617 pub secret_refs: ::prost::alloc::collections::BTreeMap<
618 ::prost::alloc::string::String,
619 super::secret::SecretRef,
620 >,
621 #[prost(message, optional, tag = "11")]
622 pub refresh_mode: ::core::option::Option<super::plan_common::SourceRefreshMode>,
623 #[prost(uint32, optional, tag = "12", wrapper = "crate::id::TableId")]
624 pub associated_table_id: ::core::option::Option<crate::id::TableId>,
625}
626#[derive(prost_helpers::AnyPB)]
629#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
630pub struct BarrierRecvNode {}
631#[derive(prost_helpers::AnyPB)]
632#[derive(Clone, PartialEq, ::prost::Message)]
633pub struct SourceNode {
634 #[prost(message, optional, tag = "1")]
637 pub source_inner: ::core::option::Option<StreamSource>,
638}
639#[derive(prost_helpers::AnyPB)]
640#[derive(Clone, PartialEq, ::prost::Message)]
641pub struct StreamFsFetchNode {
642 #[prost(message, optional, tag = "1")]
643 pub node_inner: ::core::option::Option<StreamFsFetch>,
644}
645#[derive(prost_helpers::AnyPB)]
648#[derive(Clone, PartialEq, ::prost::Message)]
649pub struct SourceBackfillNode {
650 #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
651 pub upstream_source_id: crate::id::SourceId,
652 #[prost(uint32, optional, tag = "2")]
653 pub row_id_index: ::core::option::Option<u32>,
654 #[prost(message, repeated, tag = "3")]
655 pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
656 #[prost(message, optional, tag = "4")]
657 pub info: ::core::option::Option<super::catalog::StreamSourceInfo>,
658 #[prost(string, tag = "5")]
659 pub source_name: ::prost::alloc::string::String,
660 #[prost(btree_map = "string, string", tag = "6")]
661 pub with_properties: ::prost::alloc::collections::BTreeMap<
662 ::prost::alloc::string::String,
663 ::prost::alloc::string::String,
664 >,
665 #[prost(uint32, optional, tag = "7")]
667 pub rate_limit: ::core::option::Option<u32>,
668 #[prost(message, optional, tag = "8")]
670 pub state_table: ::core::option::Option<super::catalog::Table>,
671 #[prost(btree_map = "string, message", tag = "9")]
672 pub secret_refs: ::prost::alloc::collections::BTreeMap<
673 ::prost::alloc::string::String,
674 super::secret::SecretRef,
675 >,
676}
677#[derive(prost_helpers::AnyPB)]
678#[derive(Clone, PartialEq, ::prost::Message)]
679pub struct SinkDesc {
680 #[prost(uint32, tag = "1", wrapper = "crate::id::SinkId")]
681 pub id: crate::id::SinkId,
682 #[prost(string, tag = "2")]
683 pub name: ::prost::alloc::string::String,
684 #[prost(string, tag = "3")]
685 pub definition: ::prost::alloc::string::String,
686 #[prost(message, repeated, tag = "5")]
687 pub plan_pk: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
688 #[prost(uint32, repeated, tag = "6")]
689 pub downstream_pk: ::prost::alloc::vec::Vec<u32>,
690 #[prost(uint32, repeated, tag = "7")]
691 pub distribution_key: ::prost::alloc::vec::Vec<u32>,
692 #[prost(btree_map = "string, string", tag = "8")]
693 pub properties: ::prost::alloc::collections::BTreeMap<
694 ::prost::alloc::string::String,
695 ::prost::alloc::string::String,
696 >,
697 #[prost(enumeration = "super::catalog::SinkType", tag = "9")]
699 pub sink_type: i32,
700 #[prost(message, repeated, tag = "10")]
701 pub column_catalogs: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
702 #[prost(string, tag = "11")]
703 pub db_name: ::prost::alloc::string::String,
704 #[prost(string, tag = "12")]
707 pub sink_from_name: ::prost::alloc::string::String,
708 #[prost(message, optional, tag = "13")]
709 pub format_desc: ::core::option::Option<super::catalog::SinkFormatDesc>,
710 #[prost(uint32, optional, tag = "14")]
711 pub target_table: ::core::option::Option<u32>,
712 #[prost(uint64, optional, tag = "15")]
713 pub extra_partition_col_idx: ::core::option::Option<u64>,
714 #[prost(btree_map = "string, message", tag = "16")]
715 pub secret_refs: ::prost::alloc::collections::BTreeMap<
716 ::prost::alloc::string::String,
717 super::secret::SecretRef,
718 >,
719 #[prost(bool, tag = "17")]
723 pub raw_ignore_delete: bool,
724}
725#[derive(prost_helpers::AnyPB)]
726#[derive(Clone, PartialEq, ::prost::Message)]
727pub struct SinkNode {
728 #[prost(message, optional, tag = "1")]
729 pub sink_desc: ::core::option::Option<SinkDesc>,
730 #[prost(message, optional, tag = "2")]
732 pub table: ::core::option::Option<super::catalog::Table>,
733 #[prost(enumeration = "SinkLogStoreType", tag = "3")]
734 pub log_store_type: i32,
735 #[prost(uint32, optional, tag = "4")]
736 pub rate_limit: ::core::option::Option<u32>,
737}
738#[derive(prost_helpers::AnyPB)]
739#[derive(Clone, PartialEq, ::prost::Message)]
740pub struct IcebergWithPkIndexWriterNode {
741 #[prost(message, optional, tag = "1")]
742 pub sink_desc: ::core::option::Option<SinkDesc>,
743 #[prost(message, optional, tag = "2")]
744 pub pk_index_table: ::core::option::Option<super::catalog::Table>,
745}
746#[derive(prost_helpers::AnyPB)]
747#[derive(Clone, PartialEq, ::prost::Message)]
748pub struct IcebergWithPkIndexPositionDeleteMergerNode {
749 #[prost(message, optional, tag = "1")]
750 pub sink_desc: ::core::option::Option<SinkDesc>,
751}
752#[derive(prost_helpers::AnyPB)]
754#[derive(Clone, PartialEq, ::prost::Message)]
755pub struct CompactionResolverNode {
756 #[prost(message, optional, tag = "1")]
757 pub sink_desc: ::core::option::Option<SinkDesc>,
758 #[prost(message, optional, tag = "2")]
759 pub pk_index_table: ::core::option::Option<super::catalog::Table>,
760 #[prost(string, repeated, tag = "3")]
761 pub output_data_file_paths: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
762 #[prost(string, repeated, tag = "4")]
763 pub input_data_file_paths: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
764 #[prost(int64, tag = "5")]
765 pub read_snapshot_id: i64,
766 #[prost(uint64, tag = "6", wrapper = "crate::id::IcebergCompactionTaskId")]
767 pub compaction_task_id: crate::id::IcebergCompactionTaskId,
768}
769#[derive(prost_helpers::AnyPB)]
770#[derive(Clone, PartialEq, ::prost::Message)]
771pub struct ProjectNode {
772 #[prost(message, repeated, tag = "1")]
773 pub select_list: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
774 #[prost(uint32, repeated, tag = "2")]
778 pub watermark_input_cols: ::prost::alloc::vec::Vec<u32>,
779 #[prost(uint32, repeated, tag = "3")]
780 pub watermark_output_cols: ::prost::alloc::vec::Vec<u32>,
781 #[prost(uint32, repeated, tag = "4")]
782 pub nondecreasing_exprs: ::prost::alloc::vec::Vec<u32>,
783 #[prost(bool, tag = "5")]
786 pub noop_update_hint: bool,
787}
788#[derive(prost_helpers::AnyPB)]
789#[derive(Clone, PartialEq, ::prost::Message)]
790pub struct FilterNode {
791 #[prost(message, optional, tag = "1")]
792 pub search_condition: ::core::option::Option<super::expr::ExprNode>,
793}
794#[derive(prost_helpers::AnyPB)]
795#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
796pub struct ChangeLogNode {
797 #[prost(bool, tag = "1")]
799 pub need_op: bool,
800 #[prost(uint32, repeated, tag = "2")]
801 pub distribution_keys: ::prost::alloc::vec::Vec<u32>,
802}
803#[derive(prost_helpers::AnyPB)]
804#[derive(Clone, PartialEq, ::prost::Message)]
805pub struct CdcFilterNode {
806 #[prost(message, optional, tag = "1")]
807 pub search_condition: ::core::option::Option<super::expr::ExprNode>,
808 #[prost(uint32, tag = "2", wrapper = "crate::id::SourceId")]
809 pub upstream_source_id: crate::id::SourceId,
810}
811#[derive(prost_helpers::AnyPB)]
820#[derive(Clone, PartialEq, ::prost::Message)]
821pub struct MaterializeNode {
822 #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
823 pub table_id: crate::id::TableId,
824 #[prost(message, repeated, tag = "2")]
826 pub column_orders: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
827 #[prost(message, optional, tag = "3")]
834 pub table: ::core::option::Option<super::catalog::Table>,
835 #[prost(message, optional, tag = "5")]
844 pub staging_table: ::core::option::Option<super::catalog::Table>,
845 #[prost(message, optional, tag = "6")]
858 pub refresh_progress_table: ::core::option::Option<super::catalog::Table>,
859 #[prost(bool, tag = "7")]
861 pub cleaned_by_ttl_watermark: bool,
862}
863#[derive(prost_helpers::AnyPB)]
864#[derive(Clone, PartialEq, ::prost::Message)]
865pub struct AggCallState {
866 #[prost(oneof = "agg_call_state::Inner", tags = "1, 3")]
867 pub inner: ::core::option::Option<agg_call_state::Inner>,
868}
869pub mod agg_call_state {
871 #[derive(prost_helpers::AnyPB)]
873 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
874 pub struct ValueState {}
875 #[derive(prost_helpers::AnyPB)]
877 #[derive(Clone, PartialEq, ::prost::Message)]
878 pub struct MaterializedInputState {
879 #[prost(message, optional, tag = "1")]
880 pub table: ::core::option::Option<super::super::catalog::Table>,
881 #[prost(uint32, repeated, tag = "2")]
883 pub included_upstream_indices: ::prost::alloc::vec::Vec<u32>,
884 #[prost(uint32, repeated, tag = "3")]
885 pub table_value_indices: ::prost::alloc::vec::Vec<u32>,
886 #[prost(message, repeated, tag = "4")]
887 pub order_columns: ::prost::alloc::vec::Vec<super::super::common::ColumnOrder>,
888 }
889 #[derive(prost_helpers::AnyPB)]
890 #[derive(Clone, PartialEq, ::prost::Oneof)]
891 pub enum Inner {
892 #[prost(message, tag = "1")]
893 ValueState(ValueState),
894 #[prost(message, tag = "3")]
895 MaterializedInputState(MaterializedInputState),
896 }
897}
898#[derive(prost_helpers::AnyPB)]
899#[derive(Clone, PartialEq, ::prost::Message)]
900pub struct SimpleAggNode {
901 #[prost(message, repeated, tag = "1")]
902 pub agg_calls: ::prost::alloc::vec::Vec<super::expr::AggCall>,
903 #[prost(message, repeated, tag = "3")]
904 pub agg_call_states: ::prost::alloc::vec::Vec<AggCallState>,
905 #[prost(message, optional, tag = "4")]
906 pub intermediate_state_table: ::core::option::Option<super::catalog::Table>,
907 #[prost(bool, tag = "5")]
910 pub is_append_only: bool,
911 #[prost(map = "uint32, message", tag = "6")]
912 pub distinct_dedup_tables: ::std::collections::HashMap<u32, super::catalog::Table>,
913 #[prost(uint32, tag = "7")]
914 pub row_count_index: u32,
915 #[prost(enumeration = "AggNodeVersion", tag = "8")]
916 pub version: i32,
917 #[prost(bool, tag = "9")]
920 pub must_output_per_barrier: bool,
921}
922#[derive(prost_helpers::AnyPB)]
923#[derive(Clone, PartialEq, ::prost::Message)]
924pub struct HashAggNode {
925 #[prost(uint32, repeated, tag = "1")]
926 pub group_key: ::prost::alloc::vec::Vec<u32>,
927 #[prost(message, repeated, tag = "2")]
928 pub agg_calls: ::prost::alloc::vec::Vec<super::expr::AggCall>,
929 #[prost(message, repeated, tag = "3")]
930 pub agg_call_states: ::prost::alloc::vec::Vec<AggCallState>,
931 #[prost(message, optional, tag = "4")]
932 pub intermediate_state_table: ::core::option::Option<super::catalog::Table>,
933 #[prost(bool, tag = "5")]
936 pub is_append_only: bool,
937 #[prost(map = "uint32, message", tag = "6")]
938 pub distinct_dedup_tables: ::std::collections::HashMap<u32, super::catalog::Table>,
939 #[prost(uint32, tag = "7")]
940 pub row_count_index: u32,
941 #[prost(bool, tag = "8")]
942 pub emit_on_window_close: bool,
943 #[prost(enumeration = "AggNodeVersion", tag = "9")]
944 pub version: i32,
945}
946#[derive(prost_helpers::AnyPB)]
947#[derive(Clone, PartialEq, ::prost::Message)]
948pub struct TopNNode {
949 #[prost(uint64, tag = "1")]
951 pub limit: u64,
952 #[prost(uint64, tag = "2")]
953 pub offset: u64,
954 #[prost(message, optional, tag = "3")]
955 pub table: ::core::option::Option<super::catalog::Table>,
956 #[prost(message, repeated, tag = "4")]
957 pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
958 #[prost(bool, tag = "5")]
959 pub with_ties: bool,
960}
961#[derive(prost_helpers::AnyPB)]
962#[derive(Clone, PartialEq, ::prost::Message)]
963pub struct GroupTopNNode {
964 #[prost(uint64, tag = "1")]
966 pub limit: u64,
967 #[prost(uint64, tag = "2")]
968 pub offset: u64,
969 #[prost(uint32, repeated, tag = "3")]
970 pub group_key: ::prost::alloc::vec::Vec<u32>,
971 #[prost(message, optional, tag = "4")]
972 pub table: ::core::option::Option<super::catalog::Table>,
973 #[prost(message, repeated, tag = "5")]
974 pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
975 #[prost(bool, tag = "6")]
976 pub with_ties: bool,
977}
978#[derive(prost_helpers::AnyPB)]
979#[derive(Clone, PartialEq, ::prost::Message)]
980pub struct DeltaExpression {
981 #[prost(enumeration = "super::expr::expr_node::Type", tag = "1")]
982 pub delta_type: i32,
983 #[prost(message, optional, tag = "2")]
984 pub delta: ::core::option::Option<super::expr::ExprNode>,
985}
986#[derive(prost_helpers::AnyPB)]
988#[derive(Clone, PartialEq, ::prost::Message)]
989pub struct InequalityPair {
990 #[prost(uint32, tag = "1")]
992 pub key_required_larger: u32,
993 #[prost(uint32, tag = "2")]
995 pub key_required_smaller: u32,
996 #[prost(bool, tag = "3")]
998 pub clean_state: bool,
999 #[prost(message, optional, tag = "4")]
1001 pub delta_expression: ::core::option::Option<DeltaExpression>,
1002}
1003#[derive(prost_helpers::AnyPB)]
1004#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1005pub struct InequalityPairV2 {
1006 #[prost(uint32, tag = "1")]
1008 pub left_idx: u32,
1009 #[prost(uint32, tag = "2")]
1011 pub right_idx: u32,
1012 #[prost(bool, tag = "3")]
1014 pub clean_left_state: bool,
1015 #[prost(bool, tag = "4")]
1017 pub clean_right_state: bool,
1018 #[prost(enumeration = "InequalityType", tag = "5")]
1020 pub op: i32,
1021}
1022#[derive(prost_helpers::AnyPB)]
1023#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1024pub struct JoinKeyWatermarkIndex {
1025 #[prost(uint32, tag = "1")]
1027 pub index: u32,
1028 #[prost(bool, tag = "2")]
1030 pub do_state_cleaning: bool,
1031}
1032#[derive(prost_helpers::AnyPB)]
1033#[derive(Clone, PartialEq, ::prost::Message)]
1034pub struct HashJoinWatermarkHandleDesc {
1035 #[prost(message, repeated, tag = "1")]
1037 pub watermark_indices_in_jk: ::prost::alloc::vec::Vec<JoinKeyWatermarkIndex>,
1038 #[prost(message, repeated, tag = "2")]
1040 pub inequality_pairs: ::prost::alloc::vec::Vec<InequalityPairV2>,
1041}
1042#[derive(prost_helpers::AnyPB)]
1043#[derive(Clone, PartialEq, ::prost::Message)]
1044pub struct HashJoinNode {
1045 #[prost(enumeration = "super::plan_common::JoinType", tag = "1")]
1046 pub join_type: i32,
1047 #[prost(int32, repeated, tag = "2")]
1048 pub left_key: ::prost::alloc::vec::Vec<i32>,
1049 #[prost(int32, repeated, tag = "3")]
1050 pub right_key: ::prost::alloc::vec::Vec<i32>,
1051 #[prost(message, optional, tag = "4")]
1052 pub condition: ::core::option::Option<super::expr::ExprNode>,
1053 #[prost(message, optional, tag = "6")]
1055 pub left_table: ::core::option::Option<super::catalog::Table>,
1056 #[prost(message, optional, tag = "7")]
1058 pub right_table: ::core::option::Option<super::catalog::Table>,
1059 #[prost(message, optional, tag = "8")]
1061 pub left_degree_table: ::core::option::Option<super::catalog::Table>,
1062 #[prost(message, optional, tag = "9")]
1064 pub right_degree_table: ::core::option::Option<super::catalog::Table>,
1065 #[prost(uint32, repeated, tag = "10")]
1067 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1068 #[prost(uint32, repeated, tag = "11")]
1073 pub left_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1074 #[prost(uint32, repeated, tag = "12")]
1079 pub right_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1080 #[prost(bool, repeated, tag = "13")]
1081 pub null_safe: ::prost::alloc::vec::Vec<bool>,
1082 #[prost(bool, tag = "14")]
1085 pub is_append_only: bool,
1086 #[deprecated]
1089 #[prost(enumeration = "JoinEncodingType", tag = "15")]
1090 pub join_encoding_type: i32,
1091 #[prost(message, optional, tag = "17")]
1093 pub watermark_handle_desc: ::core::option::Option<HashJoinWatermarkHandleDesc>,
1094}
1095#[derive(prost_helpers::AnyPB)]
1096#[derive(Clone, PartialEq, ::prost::Message)]
1097pub struct AsOfJoinNode {
1098 #[prost(enumeration = "super::plan_common::AsOfJoinType", tag = "1")]
1099 pub join_type: i32,
1100 #[prost(int32, repeated, tag = "2")]
1101 pub left_key: ::prost::alloc::vec::Vec<i32>,
1102 #[prost(int32, repeated, tag = "3")]
1103 pub right_key: ::prost::alloc::vec::Vec<i32>,
1104 #[prost(message, optional, tag = "4")]
1106 pub left_table: ::core::option::Option<super::catalog::Table>,
1107 #[prost(message, optional, tag = "5")]
1109 pub right_table: ::core::option::Option<super::catalog::Table>,
1110 #[prost(uint32, repeated, tag = "6")]
1112 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1113 #[prost(uint32, repeated, tag = "7")]
1117 pub left_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1118 #[prost(uint32, repeated, tag = "8")]
1122 pub right_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1123 #[prost(bool, repeated, tag = "9")]
1124 pub null_safe: ::prost::alloc::vec::Vec<bool>,
1125 #[prost(message, optional, tag = "10")]
1126 pub asof_desc: ::core::option::Option<super::plan_common::AsOfJoinDesc>,
1127 #[deprecated]
1130 #[prost(enumeration = "JoinEncodingType", tag = "11")]
1131 pub join_encoding_type: i32,
1132 #[prost(bool, optional, tag = "12")]
1136 pub use_cache: ::core::option::Option<bool>,
1137}
1138#[derive(prost_helpers::AnyPB)]
1139#[derive(Clone, PartialEq, ::prost::Message)]
1140pub struct TemporalJoinNode {
1141 #[prost(enumeration = "super::plan_common::JoinType", tag = "1")]
1142 pub join_type: i32,
1143 #[prost(int32, repeated, tag = "2")]
1144 pub left_key: ::prost::alloc::vec::Vec<i32>,
1145 #[prost(int32, repeated, tag = "3")]
1146 pub right_key: ::prost::alloc::vec::Vec<i32>,
1147 #[prost(bool, repeated, tag = "4")]
1148 pub null_safe: ::prost::alloc::vec::Vec<bool>,
1149 #[prost(message, optional, tag = "5")]
1150 pub condition: ::core::option::Option<super::expr::ExprNode>,
1151 #[prost(uint32, repeated, tag = "6")]
1153 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1154 #[prost(message, optional, tag = "7")]
1156 pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1157 #[prost(uint32, repeated, tag = "8")]
1159 pub table_output_indices: ::prost::alloc::vec::Vec<u32>,
1160 #[prost(message, optional, tag = "9")]
1162 pub memo_table: ::core::option::Option<super::catalog::Table>,
1163 #[prost(bool, tag = "10")]
1165 pub is_nested_loop: bool,
1166}
1167#[derive(prost_helpers::AnyPB)]
1168#[derive(Clone, PartialEq, ::prost::Message)]
1169pub struct DynamicFilterNode {
1170 #[prost(uint32, tag = "1")]
1171 pub left_key: u32,
1172 #[prost(message, optional, tag = "2")]
1174 pub condition: ::core::option::Option<super::expr::ExprNode>,
1175 #[prost(message, optional, tag = "3")]
1177 pub left_table: ::core::option::Option<super::catalog::Table>,
1178 #[prost(message, optional, tag = "4")]
1180 pub right_table: ::core::option::Option<super::catalog::Table>,
1181 #[deprecated]
1188 #[prost(bool, tag = "5")]
1189 pub condition_always_relax: bool,
1190 #[prost(bool, tag = "6")]
1192 pub cleaned_by_watermark: bool,
1193}
1194#[derive(prost_helpers::AnyPB)]
1197#[derive(Clone, PartialEq, ::prost::Message)]
1198pub struct DeltaIndexJoinNode {
1199 #[prost(enumeration = "super::plan_common::JoinType", tag = "1")]
1200 pub join_type: i32,
1201 #[prost(int32, repeated, tag = "2")]
1202 pub left_key: ::prost::alloc::vec::Vec<i32>,
1203 #[prost(int32, repeated, tag = "3")]
1204 pub right_key: ::prost::alloc::vec::Vec<i32>,
1205 #[prost(message, optional, tag = "4")]
1206 pub condition: ::core::option::Option<super::expr::ExprNode>,
1207 #[prost(uint32, tag = "7", wrapper = "crate::id::TableId")]
1209 pub left_table_id: crate::id::TableId,
1210 #[prost(uint32, tag = "8", wrapper = "crate::id::TableId")]
1212 pub right_table_id: crate::id::TableId,
1213 #[prost(message, optional, tag = "9")]
1215 pub left_info: ::core::option::Option<ArrangementInfo>,
1216 #[prost(message, optional, tag = "10")]
1218 pub right_info: ::core::option::Option<ArrangementInfo>,
1219 #[prost(uint32, repeated, tag = "11")]
1221 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1222}
1223#[derive(prost_helpers::AnyPB)]
1224#[derive(Clone, PartialEq, ::prost::Message)]
1225pub struct HopWindowNode {
1226 #[prost(uint32, tag = "1")]
1227 pub time_col: u32,
1228 #[prost(message, optional, tag = "2")]
1229 pub window_slide: ::core::option::Option<super::data::Interval>,
1230 #[prost(message, optional, tag = "3")]
1231 pub window_size: ::core::option::Option<super::data::Interval>,
1232 #[prost(uint32, repeated, tag = "4")]
1233 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1234 #[prost(message, repeated, tag = "5")]
1235 pub window_start_exprs: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
1236 #[prost(message, repeated, tag = "6")]
1237 pub window_end_exprs: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
1238}
1239#[derive(prost_helpers::AnyPB)]
1240#[derive(Clone, PartialEq, ::prost::Message)]
1241pub struct MergeNode {
1242 #[deprecated]
1249 #[prost(uint32, repeated, packed = "false", tag = "1")]
1250 pub upstream_actor_id: ::prost::alloc::vec::Vec<u32>,
1251 #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
1252 pub upstream_fragment_id: crate::id::FragmentId,
1253 #[prost(enumeration = "DispatcherType", tag = "3")]
1256 pub upstream_dispatcher_type: i32,
1257 #[deprecated]
1259 #[prost(message, repeated, tag = "4")]
1260 pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
1261 #[prost(bool, tag = "5")]
1263 pub allow_no_initial_upstream: bool,
1264}
1265#[derive(prost_helpers::AnyPB)]
1268#[derive(Clone, PartialEq, ::prost::Message)]
1269pub struct ExchangeNode {
1270 #[prost(message, optional, tag = "1")]
1271 pub strategy: ::core::option::Option<DispatchStrategy>,
1272}
1273#[derive(prost_helpers::AnyPB)]
1279#[derive(Clone, PartialEq, ::prost::Message)]
1280pub struct StreamScanNode {
1281 #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
1282 pub table_id: crate::id::TableId,
1283 #[prost(int32, repeated, tag = "2")]
1288 pub upstream_column_ids: ::prost::alloc::vec::Vec<i32>,
1289 #[prost(uint32, repeated, tag = "3")]
1294 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1295 #[prost(enumeration = "StreamScanType", tag = "4")]
1300 pub stream_scan_type: i32,
1301 #[prost(message, optional, tag = "5")]
1303 pub state_table: ::core::option::Option<super::catalog::Table>,
1304 #[prost(message, optional, tag = "7")]
1307 pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1308 #[prost(uint32, optional, tag = "8")]
1310 pub rate_limit: ::core::option::Option<u32>,
1311 #[deprecated]
1313 #[prost(uint32, tag = "9")]
1314 pub snapshot_read_barrier_interval: u32,
1315 #[prost(message, optional, tag = "10")]
1318 pub arrangement_table: ::core::option::Option<super::catalog::Table>,
1319 #[prost(uint64, optional, tag = "11")]
1320 pub snapshot_backfill_epoch: ::core::option::Option<u64>,
1321 #[prost(message, optional, tag = "13")]
1323 pub pk_scan_range: ::core::option::Option<super::batch_plan::ScanRange>,
1324}
1325#[derive(prost_helpers::AnyPB)]
1327#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1328pub struct StreamCdcScanOptions {
1329 #[prost(bool, tag = "1")]
1331 pub disable_backfill: bool,
1332 #[prost(uint32, tag = "2")]
1333 pub snapshot_barrier_interval: u32,
1334 #[prost(uint32, tag = "3")]
1335 pub snapshot_batch_size: u32,
1336 #[prost(uint32, tag = "4")]
1337 pub backfill_parallelism: u32,
1338 #[prost(uint64, tag = "5")]
1339 pub backfill_num_rows_per_split: u64,
1340 #[prost(bool, tag = "6")]
1341 pub backfill_as_even_splits: bool,
1342 #[prost(uint32, tag = "7")]
1343 pub backfill_split_pk_column_index: u32,
1344}
1345#[derive(prost_helpers::AnyPB)]
1346#[derive(Clone, PartialEq, ::prost::Message)]
1347pub struct StreamCdcScanNode {
1348 #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
1349 pub table_id: crate::id::TableId,
1350 #[prost(int32, repeated, tag = "2")]
1353 pub upstream_column_ids: ::prost::alloc::vec::Vec<i32>,
1354 #[prost(uint32, repeated, tag = "3")]
1356 pub output_indices: ::prost::alloc::vec::Vec<u32>,
1357 #[prost(message, optional, tag = "4")]
1359 pub state_table: ::core::option::Option<super::catalog::Table>,
1360 #[prost(message, optional, tag = "5")]
1362 pub cdc_table_desc: ::core::option::Option<super::plan_common::ExternalTableDesc>,
1363 #[prost(uint32, optional, tag = "6")]
1365 pub rate_limit: ::core::option::Option<u32>,
1366 #[prost(bool, tag = "7")]
1369 pub disable_backfill: bool,
1370 #[prost(message, optional, tag = "8")]
1371 pub options: ::core::option::Option<StreamCdcScanOptions>,
1372}
1373#[derive(prost_helpers::AnyPB)]
1377#[derive(Clone, PartialEq, ::prost::Message)]
1378pub struct BatchPlanNode {
1379 #[prost(message, optional, tag = "1")]
1380 pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1381 #[prost(int32, repeated, tag = "2")]
1382 pub column_ids: ::prost::alloc::vec::Vec<i32>,
1383}
1384#[derive(prost_helpers::AnyPB)]
1385#[derive(Clone, PartialEq, ::prost::Message)]
1386pub struct ArrangementInfo {
1387 #[prost(message, repeated, tag = "1")]
1390 pub arrange_key_orders: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1391 #[prost(message, repeated, tag = "2")]
1393 pub column_descs: ::prost::alloc::vec::Vec<super::plan_common::ColumnDesc>,
1394 #[prost(message, optional, tag = "4")]
1396 pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1397 #[prost(uint32, repeated, tag = "5")]
1399 pub output_col_idx: ::prost::alloc::vec::Vec<u32>,
1400}
1401#[derive(prost_helpers::AnyPB)]
1404#[derive(Clone, PartialEq, ::prost::Message)]
1405pub struct ArrangeNode {
1406 #[prost(message, optional, tag = "1")]
1408 pub table_info: ::core::option::Option<ArrangementInfo>,
1409 #[prost(uint32, repeated, tag = "2")]
1411 pub distribution_key: ::prost::alloc::vec::Vec<u32>,
1412 #[prost(message, optional, tag = "3")]
1414 pub table: ::core::option::Option<super::catalog::Table>,
1415}
1416#[derive(prost_helpers::AnyPB)]
1418#[derive(Clone, PartialEq, ::prost::Message)]
1419pub struct LookupNode {
1420 #[prost(int32, repeated, tag = "1")]
1422 pub arrange_key: ::prost::alloc::vec::Vec<i32>,
1423 #[prost(int32, repeated, tag = "2")]
1425 pub stream_key: ::prost::alloc::vec::Vec<i32>,
1426 #[prost(bool, tag = "3")]
1428 pub use_current_epoch: bool,
1429 #[prost(int32, repeated, tag = "4")]
1433 pub column_mapping: ::prost::alloc::vec::Vec<i32>,
1434 #[prost(message, optional, tag = "7")]
1436 pub arrangement_table_info: ::core::option::Option<ArrangementInfo>,
1437 #[prost(oneof = "lookup_node::ArrangementTableId", tags = "5, 6")]
1438 pub arrangement_table_id: ::core::option::Option<lookup_node::ArrangementTableId>,
1439}
1440pub mod lookup_node {
1442 #[derive(prost_helpers::AnyPB)]
1443 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Oneof)]
1444 pub enum ArrangementTableId {
1445 #[prost(uint32, tag = "5", wrapper = "crate::id::TableId")]
1447 TableId(crate::id::TableId),
1448 #[prost(uint32, tag = "6", wrapper = "crate::id::TableId")]
1450 IndexId(crate::id::TableId),
1451 }
1452}
1453#[derive(prost_helpers::AnyPB)]
1455#[derive(Clone, PartialEq, ::prost::Message)]
1456pub struct WatermarkFilterNode {
1457 #[prost(message, repeated, tag = "1")]
1459 pub watermark_descs: ::prost::alloc::vec::Vec<super::catalog::WatermarkDesc>,
1460 #[prost(message, repeated, tag = "2")]
1462 pub tables: ::prost::alloc::vec::Vec<super::catalog::Table>,
1463}
1464#[derive(prost_helpers::AnyPB)]
1466#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1467pub struct UnionNode {}
1468#[derive(prost_helpers::AnyPB)]
1470#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1471pub struct LookupUnionNode {
1472 #[prost(uint32, repeated, tag = "1")]
1473 pub order: ::prost::alloc::vec::Vec<u32>,
1474}
1475#[derive(prost_helpers::AnyPB)]
1476#[derive(Clone, PartialEq, ::prost::Message)]
1477pub struct ExpandNode {
1478 #[prost(message, repeated, tag = "1")]
1479 pub column_subsets: ::prost::alloc::vec::Vec<expand_node::Subset>,
1480}
1481pub mod expand_node {
1483 #[derive(prost_helpers::AnyPB)]
1484 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1485 pub struct Subset {
1486 #[prost(uint32, repeated, tag = "1")]
1487 pub column_indices: ::prost::alloc::vec::Vec<u32>,
1488 }
1489}
1490#[derive(prost_helpers::AnyPB)]
1491#[derive(Clone, PartialEq, ::prost::Message)]
1492pub struct ProjectSetNode {
1493 #[prost(message, repeated, tag = "1")]
1494 pub select_list: ::prost::alloc::vec::Vec<super::expr::ProjectSetSelectItem>,
1495 #[prost(uint32, repeated, tag = "2")]
1499 pub watermark_input_cols: ::prost::alloc::vec::Vec<u32>,
1500 #[prost(uint32, repeated, tag = "3")]
1501 pub watermark_expr_indices: ::prost::alloc::vec::Vec<u32>,
1502 #[prost(uint32, repeated, tag = "4")]
1503 pub nondecreasing_exprs: ::prost::alloc::vec::Vec<u32>,
1504}
1505#[derive(prost_helpers::AnyPB)]
1507#[derive(Clone, PartialEq, ::prost::Message)]
1508pub struct SortNode {
1509 #[prost(message, optional, tag = "1")]
1511 pub state_table: ::core::option::Option<super::catalog::Table>,
1512 #[prost(uint32, tag = "2")]
1514 pub sort_column_index: u32,
1515}
1516#[derive(prost_helpers::AnyPB)]
1518#[derive(Clone, PartialEq, ::prost::Message)]
1519pub struct DmlNode {
1520 #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
1522 pub table_id: crate::id::TableId,
1523 #[prost(uint64, tag = "3")]
1525 pub table_version_id: u64,
1526 #[prost(message, repeated, tag = "2")]
1528 pub column_descs: ::prost::alloc::vec::Vec<super::plan_common::ColumnDesc>,
1529 #[prost(uint32, optional, tag = "4")]
1530 pub rate_limit: ::core::option::Option<u32>,
1531}
1532#[derive(prost_helpers::AnyPB)]
1533#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1534pub struct RowIdGenNode {
1535 #[prost(uint64, tag = "1")]
1536 pub row_id_index: u64,
1537}
1538#[derive(prost_helpers::AnyPB)]
1539#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1540pub struct NowModeUpdateCurrent {}
1541#[derive(prost_helpers::AnyPB)]
1542#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1543pub struct NowModeGenerateSeries {
1544 #[prost(message, optional, tag = "1")]
1545 pub start_timestamp: ::core::option::Option<super::data::Datum>,
1546 #[prost(message, optional, tag = "2")]
1547 pub interval: ::core::option::Option<super::data::Datum>,
1548}
1549#[derive(prost_helpers::AnyPB)]
1550#[derive(Clone, PartialEq, ::prost::Message)]
1551pub struct NowNode {
1552 #[prost(message, optional, tag = "1")]
1554 pub state_table: ::core::option::Option<super::catalog::Table>,
1555 #[prost(oneof = "now_node::Mode", tags = "101, 102")]
1556 pub mode: ::core::option::Option<now_node::Mode>,
1557}
1558pub mod now_node {
1560 #[derive(prost_helpers::AnyPB)]
1561 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Oneof)]
1562 pub enum Mode {
1563 #[prost(message, tag = "101")]
1564 UpdateCurrent(super::NowModeUpdateCurrent),
1565 #[prost(message, tag = "102")]
1566 GenerateSeries(super::NowModeGenerateSeries),
1567 }
1568}
1569#[derive(prost_helpers::AnyPB)]
1570#[derive(Clone, PartialEq, ::prost::Message)]
1571pub struct ValuesNode {
1572 #[prost(message, repeated, tag = "1")]
1573 pub tuples: ::prost::alloc::vec::Vec<values_node::ExprTuple>,
1574 #[prost(message, repeated, tag = "2")]
1575 pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
1576}
1577pub mod values_node {
1579 #[derive(prost_helpers::AnyPB)]
1580 #[derive(Clone, PartialEq, ::prost::Message)]
1581 pub struct ExprTuple {
1582 #[prost(message, repeated, tag = "1")]
1583 pub cells: ::prost::alloc::vec::Vec<super::super::expr::ExprNode>,
1584 }
1585}
1586#[derive(prost_helpers::AnyPB)]
1587#[derive(Clone, PartialEq, ::prost::Message)]
1588pub struct DedupNode {
1589 #[prost(message, optional, tag = "1")]
1590 pub state_table: ::core::option::Option<super::catalog::Table>,
1591 #[prost(uint32, repeated, tag = "2")]
1592 pub dedup_column_indices: ::prost::alloc::vec::Vec<u32>,
1593}
1594#[derive(prost_helpers::AnyPB)]
1595#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1596pub struct NoOpNode {}
1597#[derive(prost_helpers::AnyPB)]
1598#[derive(Clone, PartialEq, ::prost::Message)]
1599pub struct EowcOverWindowNode {
1600 #[prost(message, repeated, tag = "1")]
1601 pub calls: ::prost::alloc::vec::Vec<super::expr::WindowFunction>,
1602 #[prost(uint32, repeated, tag = "2")]
1603 pub partition_by: ::prost::alloc::vec::Vec<u32>,
1604 #[prost(message, repeated, tag = "3")]
1606 pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1607 #[prost(message, optional, tag = "4")]
1608 pub state_table: ::core::option::Option<super::catalog::Table>,
1609 #[prost(message, optional, tag = "5")]
1612 pub intermediate_state_table: ::core::option::Option<super::catalog::Table>,
1613}
1614#[derive(prost_helpers::AnyPB)]
1615#[derive(Clone, PartialEq, ::prost::Message)]
1616pub struct OverWindowNode {
1617 #[prost(message, repeated, tag = "1")]
1618 pub calls: ::prost::alloc::vec::Vec<super::expr::WindowFunction>,
1619 #[prost(uint32, repeated, tag = "2")]
1620 pub partition_by: ::prost::alloc::vec::Vec<u32>,
1621 #[prost(message, repeated, tag = "3")]
1622 pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1623 #[prost(message, optional, tag = "4")]
1624 pub state_table: ::core::option::Option<super::catalog::Table>,
1625 #[deprecated]
1627 #[prost(enumeration = "OverWindowCachePolicy", tag = "5")]
1628 pub cache_policy: i32,
1629}
1630#[derive(prost_helpers::AnyPB)]
1631#[derive(Clone, Copy, PartialEq, ::prost::Message)]
1632pub struct LocalApproxPercentileNode {
1633 #[prost(double, tag = "1")]
1634 pub base: f64,
1635 #[prost(uint32, tag = "2")]
1636 pub percentile_index: u32,
1637}
1638#[derive(prost_helpers::AnyPB)]
1639#[derive(Clone, PartialEq, ::prost::Message)]
1640pub struct GlobalApproxPercentileNode {
1641 #[prost(double, tag = "1")]
1642 pub base: f64,
1643 #[prost(double, tag = "2")]
1644 pub quantile: f64,
1645 #[prost(message, optional, tag = "3")]
1646 pub bucket_state_table: ::core::option::Option<super::catalog::Table>,
1647 #[prost(message, optional, tag = "4")]
1648 pub count_state_table: ::core::option::Option<super::catalog::Table>,
1649}
1650#[derive(prost_helpers::AnyPB)]
1651#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1652pub struct RowMergeNode {
1653 #[prost(message, optional, tag = "1")]
1654 pub lhs_mapping: ::core::option::Option<super::catalog::ColIndexMapping>,
1655 #[prost(message, optional, tag = "2")]
1656 pub rhs_mapping: ::core::option::Option<super::catalog::ColIndexMapping>,
1657}
1658#[derive(prost_helpers::AnyPB)]
1659#[derive(Clone, PartialEq, ::prost::Message)]
1660pub struct SyncLogStoreNode {
1661 #[prost(message, optional, tag = "1")]
1662 pub log_store_table: ::core::option::Option<super::catalog::Table>,
1663 #[deprecated]
1665 #[prost(uint32, optional, tag = "2")]
1666 pub pause_duration_ms: ::core::option::Option<u32>,
1667 #[deprecated]
1669 #[prost(uint32, optional, tag = "3")]
1670 pub buffer_size: ::core::option::Option<u32>,
1671 #[prost(bool, tag = "4")]
1672 pub aligned: bool,
1673}
1674#[derive(prost_helpers::AnyPB)]
1675#[derive(Clone, PartialEq, ::prost::Message)]
1676pub struct MaterializedExprsNode {
1677 #[prost(message, repeated, tag = "1")]
1678 pub exprs: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
1679 #[prost(message, optional, tag = "2")]
1680 pub state_table: ::core::option::Option<super::catalog::Table>,
1681 #[prost(uint32, optional, tag = "3")]
1682 pub state_clean_col_idx: ::core::option::Option<u32>,
1683}
1684#[derive(prost_helpers::AnyPB)]
1685#[derive(Clone, PartialEq, ::prost::Message)]
1686pub struct VectorIndexWriteNode {
1687 #[prost(message, optional, tag = "1")]
1688 pub table: ::core::option::Option<super::catalog::Table>,
1689}
1690#[derive(prost_helpers::AnyPB)]
1691#[derive(Clone, PartialEq, ::prost::Message)]
1692pub struct VectorIndexLookupJoinNode {
1693 #[prost(message, optional, tag = "1")]
1694 pub reader_desc: ::core::option::Option<super::plan_common::VectorIndexReaderDesc>,
1695 #[prost(uint32, tag = "2")]
1696 pub vector_column_idx: u32,
1697}
1698#[derive(prost_helpers::AnyPB)]
1699#[derive(Clone, PartialEq, ::prost::Message)]
1700pub struct UpstreamSinkUnionNode {
1701 #[prost(message, repeated, tag = "1")]
1704 pub init_upstreams: ::prost::alloc::vec::Vec<UpstreamSinkInfo>,
1705}
1706#[derive(prost_helpers::AnyPB)]
1707#[derive(Clone, PartialEq, ::prost::Message)]
1708pub struct LocalityProviderNode {
1709 #[prost(uint32, repeated, tag = "1")]
1711 pub locality_columns: ::prost::alloc::vec::Vec<u32>,
1712 #[prost(message, optional, tag = "2")]
1714 pub state_table: ::core::option::Option<super::catalog::Table>,
1715 #[prost(message, optional, tag = "3")]
1717 pub progress_table: ::core::option::Option<super::catalog::Table>,
1718}
1719#[derive(prost_helpers::AnyPB)]
1720#[derive(Clone, PartialEq, ::prost::Message)]
1721pub struct EowcGapFillNode {
1722 #[prost(uint32, tag = "1")]
1723 pub time_column_index: u32,
1724 #[prost(message, optional, tag = "2")]
1725 pub interval: ::core::option::Option<super::expr::ExprNode>,
1726 #[prost(uint32, repeated, tag = "3")]
1727 pub fill_columns: ::prost::alloc::vec::Vec<u32>,
1728 #[prost(string, repeated, tag = "4")]
1729 pub fill_strategies: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
1730 #[prost(message, optional, tag = "6")]
1731 pub prev_row_table: ::core::option::Option<super::catalog::Table>,
1732 #[prost(uint32, repeated, tag = "7")]
1733 pub partition_by_indices: ::prost::alloc::vec::Vec<u32>,
1734}
1735#[derive(prost_helpers::AnyPB)]
1736#[derive(Clone, PartialEq, ::prost::Message)]
1737pub struct GapFillNode {
1738 #[prost(uint32, tag = "1")]
1739 pub time_column_index: u32,
1740 #[prost(message, optional, tag = "2")]
1741 pub interval: ::core::option::Option<super::expr::ExprNode>,
1742 #[prost(uint32, repeated, tag = "3")]
1743 pub fill_columns: ::prost::alloc::vec::Vec<u32>,
1744 #[prost(string, repeated, tag = "4")]
1745 pub fill_strategies: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
1746 #[prost(message, optional, tag = "5")]
1747 pub state_table: ::core::option::Option<super::catalog::Table>,
1748 #[prost(uint32, repeated, tag = "6")]
1749 pub partition_by_indices: ::prost::alloc::vec::Vec<u32>,
1750 #[prost(uint32, repeated, tag = "7")]
1751 pub pointer_key_indices: ::prost::alloc::vec::Vec<u32>,
1752}
1753#[derive(prost_helpers::AnyPB)]
1754#[derive(Clone, PartialEq, ::prost::Message)]
1755pub struct StreamNode {
1756 #[prost(uint64, tag = "1", wrapper = "crate::id::StreamNodeLocalOperatorId")]
1759 pub operator_id: crate::id::StreamNodeLocalOperatorId,
1760 #[prost(message, repeated, tag = "3")]
1762 pub input: ::prost::alloc::vec::Vec<StreamNode>,
1763 #[prost(uint32, repeated, tag = "2")]
1764 pub stream_key: ::prost::alloc::vec::Vec<u32>,
1765 #[prost(enumeration = "stream_node::StreamKind", tag = "24")]
1766 pub stream_kind: i32,
1767 #[prost(string, tag = "18")]
1768 pub identity: ::prost::alloc::string::String,
1769 #[prost(message, repeated, tag = "19")]
1771 pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
1772 #[prost(
1773 oneof = "stream_node::NodeBody",
1774 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, 159"
1775 )]
1776 pub node_body: ::core::option::Option<stream_node::NodeBody>,
1777}
1778pub mod stream_node {
1780 #[derive(prost_helpers::AnyPB)]
1783 #[derive(
1784 Clone,
1785 Copy,
1786 Debug,
1787 PartialEq,
1788 Eq,
1789 Hash,
1790 PartialOrd,
1791 Ord,
1792 ::prost::Enumeration
1793 )]
1794 #[repr(i32)]
1795 pub enum StreamKind {
1796 Retract = 0,
1798 AppendOnly = 1,
1799 Upsert = 2,
1800 }
1801 impl StreamKind {
1802 pub fn as_str_name(&self) -> &'static str {
1807 match self {
1808 Self::Retract => "STREAM_KIND_RETRACT",
1809 Self::AppendOnly => "STREAM_KIND_APPEND_ONLY",
1810 Self::Upsert => "STREAM_KIND_UPSERT",
1811 }
1812 }
1813 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
1815 match value {
1816 "STREAM_KIND_RETRACT" => Some(Self::Retract),
1817 "STREAM_KIND_APPEND_ONLY" => Some(Self::AppendOnly),
1818 "STREAM_KIND_UPSERT" => Some(Self::Upsert),
1819 _ => None,
1820 }
1821 }
1822 }
1823 #[derive(prost_helpers::AnyPB)]
1824 #[derive(::enum_as_inner::EnumAsInner, ::strum::Display, ::strum::EnumDiscriminants)]
1825 #[derive(::prost_helpers::StreamNodeBodyVariants)]
1826 #[strum_discriminants(derive(::strum::Display, Hash))]
1827 #[derive(Clone, PartialEq, ::prost::Oneof)]
1828 pub enum NodeBody {
1829 #[prost(message, tag = "100")]
1830 Source(::prost::alloc::boxed::Box<super::SourceNode>),
1831 #[prost(message, tag = "101")]
1832 Project(::prost::alloc::boxed::Box<super::ProjectNode>),
1833 #[prost(message, tag = "102")]
1834 Filter(::prost::alloc::boxed::Box<super::FilterNode>),
1835 #[prost(message, tag = "103")]
1836 Materialize(::prost::alloc::boxed::Box<super::MaterializeNode>),
1837 #[prost(message, tag = "104")]
1838 StatelessSimpleAgg(::prost::alloc::boxed::Box<super::SimpleAggNode>),
1839 #[prost(message, tag = "105")]
1840 SimpleAgg(::prost::alloc::boxed::Box<super::SimpleAggNode>),
1841 #[prost(message, tag = "106")]
1842 HashAgg(::prost::alloc::boxed::Box<super::HashAggNode>),
1843 #[prost(message, tag = "107")]
1844 AppendOnlyTopN(::prost::alloc::boxed::Box<super::TopNNode>),
1845 #[prost(message, tag = "108")]
1846 HashJoin(::prost::alloc::boxed::Box<super::HashJoinNode>),
1847 #[prost(message, tag = "109")]
1848 TopN(::prost::alloc::boxed::Box<super::TopNNode>),
1849 #[prost(message, tag = "110")]
1850 HopWindow(::prost::alloc::boxed::Box<super::HopWindowNode>),
1851 #[prost(message, tag = "111")]
1852 Merge(::prost::alloc::boxed::Box<super::MergeNode>),
1853 #[prost(message, tag = "112")]
1854 Exchange(::prost::alloc::boxed::Box<super::ExchangeNode>),
1855 #[prost(message, tag = "113")]
1856 StreamScan(::prost::alloc::boxed::Box<super::StreamScanNode>),
1857 #[prost(message, tag = "114")]
1858 BatchPlan(::prost::alloc::boxed::Box<super::BatchPlanNode>),
1859 #[prost(message, tag = "115")]
1860 Lookup(::prost::alloc::boxed::Box<super::LookupNode>),
1861 #[prost(message, tag = "116")]
1862 Arrange(::prost::alloc::boxed::Box<super::ArrangeNode>),
1863 #[prost(message, tag = "117")]
1864 LookupUnion(::prost::alloc::boxed::Box<super::LookupUnionNode>),
1865 #[prost(message, tag = "118")]
1866 Union(super::UnionNode),
1867 #[prost(message, tag = "119")]
1868 DeltaIndexJoin(::prost::alloc::boxed::Box<super::DeltaIndexJoinNode>),
1869 #[prost(message, tag = "120")]
1870 Sink(::prost::alloc::boxed::Box<super::SinkNode>),
1871 #[prost(message, tag = "121")]
1872 Expand(::prost::alloc::boxed::Box<super::ExpandNode>),
1873 #[prost(message, tag = "122")]
1874 DynamicFilter(::prost::alloc::boxed::Box<super::DynamicFilterNode>),
1875 #[prost(message, tag = "123")]
1876 ProjectSet(::prost::alloc::boxed::Box<super::ProjectSetNode>),
1877 #[prost(message, tag = "124")]
1878 GroupTopN(::prost::alloc::boxed::Box<super::GroupTopNNode>),
1879 #[prost(message, tag = "125")]
1880 Sort(::prost::alloc::boxed::Box<super::SortNode>),
1881 #[prost(message, tag = "126")]
1882 WatermarkFilter(::prost::alloc::boxed::Box<super::WatermarkFilterNode>),
1883 #[prost(message, tag = "127")]
1884 Dml(::prost::alloc::boxed::Box<super::DmlNode>),
1885 #[prost(message, tag = "128")]
1886 RowIdGen(::prost::alloc::boxed::Box<super::RowIdGenNode>),
1887 #[prost(message, tag = "129")]
1888 Now(::prost::alloc::boxed::Box<super::NowNode>),
1889 #[prost(message, tag = "130")]
1890 AppendOnlyGroupTopN(::prost::alloc::boxed::Box<super::GroupTopNNode>),
1891 #[prost(message, tag = "131")]
1892 TemporalJoin(::prost::alloc::boxed::Box<super::TemporalJoinNode>),
1893 #[prost(message, tag = "132")]
1894 BarrierRecv(::prost::alloc::boxed::Box<super::BarrierRecvNode>),
1895 #[prost(message, tag = "133")]
1896 Values(::prost::alloc::boxed::Box<super::ValuesNode>),
1897 #[prost(message, tag = "134")]
1898 AppendOnlyDedup(::prost::alloc::boxed::Box<super::DedupNode>),
1899 #[prost(message, tag = "135")]
1900 NoOp(super::NoOpNode),
1901 #[prost(message, tag = "136")]
1902 EowcOverWindow(::prost::alloc::boxed::Box<super::EowcOverWindowNode>),
1903 #[prost(message, tag = "137")]
1904 OverWindow(::prost::alloc::boxed::Box<super::OverWindowNode>),
1905 #[prost(message, tag = "138")]
1906 StreamFsFetch(::prost::alloc::boxed::Box<super::StreamFsFetchNode>),
1907 #[prost(message, tag = "139")]
1908 StreamCdcScan(::prost::alloc::boxed::Box<super::StreamCdcScanNode>),
1909 #[prost(message, tag = "140")]
1910 CdcFilter(::prost::alloc::boxed::Box<super::CdcFilterNode>),
1911 #[prost(message, tag = "142")]
1912 SourceBackfill(::prost::alloc::boxed::Box<super::SourceBackfillNode>),
1913 #[prost(message, tag = "143")]
1914 Changelog(::prost::alloc::boxed::Box<super::ChangeLogNode>),
1915 #[prost(message, tag = "144")]
1916 LocalApproxPercentile(
1917 ::prost::alloc::boxed::Box<super::LocalApproxPercentileNode>,
1918 ),
1919 #[prost(message, tag = "145")]
1920 GlobalApproxPercentile(
1921 ::prost::alloc::boxed::Box<super::GlobalApproxPercentileNode>,
1922 ),
1923 #[prost(message, tag = "146")]
1924 RowMerge(::prost::alloc::boxed::Box<super::RowMergeNode>),
1925 #[prost(message, tag = "147")]
1926 AsOfJoin(::prost::alloc::boxed::Box<super::AsOfJoinNode>),
1927 #[prost(message, tag = "148")]
1928 SyncLogStore(::prost::alloc::boxed::Box<super::SyncLogStoreNode>),
1929 #[prost(message, tag = "149")]
1930 MaterializedExprs(::prost::alloc::boxed::Box<super::MaterializedExprsNode>),
1931 #[prost(message, tag = "150")]
1932 VectorIndexWrite(::prost::alloc::boxed::Box<super::VectorIndexWriteNode>),
1933 #[prost(message, tag = "151")]
1934 UpstreamSinkUnion(::prost::alloc::boxed::Box<super::UpstreamSinkUnionNode>),
1935 #[prost(message, tag = "152")]
1936 LocalityProvider(::prost::alloc::boxed::Box<super::LocalityProviderNode>),
1937 #[prost(message, tag = "153")]
1938 EowcGapFill(::prost::alloc::boxed::Box<super::EowcGapFillNode>),
1939 #[prost(message, tag = "154")]
1940 GapFill(::prost::alloc::boxed::Box<super::GapFillNode>),
1941 #[prost(message, tag = "155")]
1942 VectorIndexLookupJoin(
1943 ::prost::alloc::boxed::Box<super::VectorIndexLookupJoinNode>,
1944 ),
1945 #[prost(message, tag = "156")]
1946 IcebergWithPkIndexWriter(
1947 ::prost::alloc::boxed::Box<super::IcebergWithPkIndexWriterNode>,
1948 ),
1949 #[prost(message, tag = "157")]
1950 IcebergWithPkIndexPositionDeleteMerger(
1951 ::prost::alloc::boxed::Box<super::IcebergWithPkIndexPositionDeleteMergerNode>,
1952 ),
1953 #[prost(message, tag = "159")]
1954 CompactionResolver(::prost::alloc::boxed::Box<super::CompactionResolverNode>),
1955 }
1956}
1957#[derive(prost_helpers::AnyPB)]
1971#[derive(Clone, PartialEq, ::prost::Message)]
1972pub struct DispatchOutputMapping {
1973 #[prost(uint32, repeated, tag = "1")]
1975 pub indices: ::prost::alloc::vec::Vec<u32>,
1976 #[prost(message, repeated, tag = "2")]
1982 pub types: ::prost::alloc::vec::Vec<dispatch_output_mapping::TypePair>,
1983}
1984pub mod dispatch_output_mapping {
1986 #[derive(prost_helpers::AnyPB)]
1987 #[derive(Clone, PartialEq, ::prost::Message)]
1988 pub struct TypePair {
1989 #[prost(message, optional, tag = "1")]
1990 pub upstream: ::core::option::Option<super::super::data::DataType>,
1991 #[prost(message, optional, tag = "2")]
1992 pub downstream: ::core::option::Option<super::super::data::DataType>,
1993 }
1994}
1995#[derive(prost_helpers::AnyPB)]
1998#[derive(Clone, PartialEq, ::prost::Message)]
1999pub struct DispatchStrategy {
2000 #[prost(enumeration = "DispatcherType", tag = "1")]
2001 pub r#type: i32,
2002 #[prost(uint32, repeated, tag = "2")]
2003 pub dist_key_indices: ::prost::alloc::vec::Vec<u32>,
2004 #[prost(message, optional, tag = "3")]
2005 pub output_mapping: ::core::option::Option<DispatchOutputMapping>,
2006}
2007#[derive(prost_helpers::AnyPB)]
2010#[derive(Clone, PartialEq, ::prost::Message)]
2011pub struct Dispatcher {
2012 #[prost(enumeration = "DispatcherType", tag = "1")]
2013 pub r#type: i32,
2014 #[prost(uint32, repeated, tag = "2")]
2017 pub dist_key_indices: ::prost::alloc::vec::Vec<u32>,
2018 #[prost(message, optional, tag = "6")]
2020 pub output_mapping: ::core::option::Option<DispatchOutputMapping>,
2021 #[prost(message, optional, tag = "3")]
2024 pub hash_mapping: ::core::option::Option<ActorMapping>,
2025 #[prost(uint32, tag = "4", wrapper = "crate::id::FragmentId")]
2028 pub dispatcher_id: crate::id::FragmentId,
2029 #[prost(uint32, repeated, tag = "5", wrapper = "crate::id::ActorId")]
2031 pub downstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
2032}
2033#[derive(prost_helpers::AnyPB)]
2035#[derive(Clone, PartialEq, ::prost::Message)]
2036pub struct StreamActor {
2037 #[prost(uint32, tag = "1", wrapper = "crate::id::ActorId")]
2038 pub actor_id: crate::id::ActorId,
2039 #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
2040 pub fragment_id: crate::id::FragmentId,
2041 #[prost(message, repeated, tag = "4")]
2042 pub dispatcher: ::prost::alloc::vec::Vec<Dispatcher>,
2043 #[prost(message, optional, tag = "8")]
2046 pub vnode_bitmap: ::core::option::Option<super::common::Buffer>,
2047 #[prost(string, tag = "9")]
2049 pub mview_definition: ::prost::alloc::string::String,
2050 #[prost(message, optional, tag = "10")]
2052 pub expr_context: ::core::option::Option<super::plan_common::ExprContext>,
2053 #[prost(string, tag = "11")]
2055 pub config_override: ::prost::alloc::string::String,
2056}
2057#[derive(prost_helpers::AnyPB)]
2059#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2060pub struct StreamContext {
2061 #[prost(string, tag = "1")]
2063 pub timezone: ::prost::alloc::string::String,
2064 #[prost(string, tag = "2")]
2066 pub config_override: ::prost::alloc::string::String,
2067}
2068#[derive(prost_helpers::AnyPB)]
2069#[derive(Clone, PartialEq, ::prost::Message)]
2070pub struct BackfillOrder {
2071 #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::RelationId")]
2072 pub order: ::std::collections::HashMap<
2073 crate::id::RelationId,
2074 super::common::Uint32Vector,
2075 >,
2076}
2077#[derive(prost_helpers::AnyPB)]
2082#[derive(Clone, PartialEq, ::prost::Message)]
2083pub struct StreamFragmentGraph {
2084 #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::FragmentId")]
2086 pub fragments: ::std::collections::HashMap<
2087 crate::id::FragmentId,
2088 stream_fragment_graph::StreamFragment,
2089 >,
2090 #[prost(message, repeated, tag = "2")]
2092 pub edges: ::prost::alloc::vec::Vec<stream_fragment_graph::StreamFragmentEdge>,
2093 #[prost(uint32, repeated, tag = "3", wrapper = "crate::id::TableId")]
2094 pub dependent_table_ids: ::prost::alloc::vec::Vec<crate::id::TableId>,
2095 #[prost(uint32, tag = "4")]
2096 pub table_ids_cnt: u32,
2097 #[prost(message, optional, tag = "5")]
2098 pub ctx: ::core::option::Option<StreamContext>,
2099 #[prost(message, optional, tag = "6")]
2101 pub parallelism: ::core::option::Option<stream_fragment_graph::Parallelism>,
2102 #[prost(message, optional, tag = "9")]
2104 pub backfill_parallelism: ::core::option::Option<stream_fragment_graph::Parallelism>,
2105 #[prost(string, tag = "10")]
2107 pub adaptive_parallelism_strategy: ::prost::alloc::string::String,
2108 #[prost(string, tag = "11")]
2110 pub backfill_adaptive_parallelism_strategy: ::prost::alloc::string::String,
2111 #[prost(uint32, tag = "7")]
2121 pub max_parallelism: u32,
2122 #[prost(message, optional, tag = "8")]
2124 pub backfill_order: ::core::option::Option<BackfillOrder>,
2125}
2126pub mod stream_fragment_graph {
2128 #[derive(prost_helpers::AnyPB)]
2129 #[derive(Clone, PartialEq, ::prost::Message)]
2130 pub struct StreamFragment {
2131 #[prost(uint32, tag = "1", wrapper = "crate::id::FragmentId")]
2133 pub fragment_id: crate::id::FragmentId,
2134 #[prost(message, optional, tag = "2")]
2136 pub node: ::core::option::Option<super::StreamNode>,
2137 #[prost(uint32, tag = "3")]
2139 pub fragment_type_mask: u32,
2140 #[prost(bool, tag = "4")]
2144 pub requires_singleton: bool,
2145 }
2146 #[derive(prost_helpers::AnyPB)]
2147 #[derive(Clone, PartialEq, ::prost::Message)]
2148 pub struct StreamFragmentEdge {
2149 #[prost(message, optional, tag = "1")]
2151 pub dispatch_strategy: ::core::option::Option<super::DispatchStrategy>,
2152 #[prost(uint64, tag = "3")]
2156 pub link_id: u64,
2157 #[prost(uint32, tag = "4", wrapper = "crate::id::FragmentId")]
2158 pub upstream_id: crate::id::FragmentId,
2159 #[prost(uint32, tag = "5", wrapper = "crate::id::FragmentId")]
2160 pub downstream_id: crate::id::FragmentId,
2161 }
2162 #[derive(prost_helpers::AnyPB)]
2163 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
2164 pub struct Parallelism {
2165 #[prost(uint64, tag = "1")]
2166 pub parallelism: u64,
2167 }
2168}
2169#[derive(prost_helpers::AnyPB)]
2171#[derive(Clone, PartialEq, ::prost::Message)]
2172pub struct SinkSchemaChange {
2173 #[prost(message, repeated, tag = "1")]
2176 pub original_schema: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2177 #[prost(oneof = "sink_schema_change::Op", tags = "2, 3")]
2179 pub op: ::core::option::Option<sink_schema_change::Op>,
2180}
2181pub mod sink_schema_change {
2183 #[derive(prost_helpers::AnyPB)]
2185 #[derive(Clone, PartialEq, ::prost::Oneof)]
2186 pub enum Op {
2187 #[prost(message, tag = "2")]
2189 AddColumns(super::SinkAddColumnsOp),
2190 #[prost(message, tag = "3")]
2192 DropColumns(super::SinkDropColumnsOp),
2193 }
2194}
2195#[derive(prost_helpers::AnyPB)]
2197#[derive(Clone, PartialEq, ::prost::Message)]
2198pub struct SinkAddColumnsOp {
2199 #[prost(message, repeated, tag = "1")]
2201 pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2202}
2203#[derive(prost_helpers::AnyPB)]
2205#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2206pub struct SinkDropColumnsOp {
2207 #[prost(string, repeated, tag = "1")]
2209 pub column_names: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
2210}
2211#[derive(prost_helpers::AnyPB)]
2212#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2213#[repr(i32)]
2214pub enum SinkLogStoreType {
2215 Unspecified = 0,
2217 KvLogStore = 1,
2218 InMemoryLogStore = 2,
2219}
2220impl SinkLogStoreType {
2221 pub fn as_str_name(&self) -> &'static str {
2226 match self {
2227 Self::Unspecified => "SINK_LOG_STORE_TYPE_UNSPECIFIED",
2228 Self::KvLogStore => "SINK_LOG_STORE_TYPE_KV_LOG_STORE",
2229 Self::InMemoryLogStore => "SINK_LOG_STORE_TYPE_IN_MEMORY_LOG_STORE",
2230 }
2231 }
2232 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2234 match value {
2235 "SINK_LOG_STORE_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2236 "SINK_LOG_STORE_TYPE_KV_LOG_STORE" => Some(Self::KvLogStore),
2237 "SINK_LOG_STORE_TYPE_IN_MEMORY_LOG_STORE" => Some(Self::InMemoryLogStore),
2238 _ => None,
2239 }
2240 }
2241}
2242#[derive(prost_helpers::AnyPB)]
2243#[derive(prost_helpers::Version)]
2244#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2245#[repr(i32)]
2246pub enum AggNodeVersion {
2247 Unspecified = 0,
2248 Issue12140 = 1,
2250 Issue13465 = 2,
2252}
2253impl AggNodeVersion {
2254 pub fn as_str_name(&self) -> &'static str {
2259 match self {
2260 Self::Unspecified => "AGG_NODE_VERSION_UNSPECIFIED",
2261 Self::Issue12140 => "AGG_NODE_VERSION_ISSUE_12140",
2262 Self::Issue13465 => "AGG_NODE_VERSION_ISSUE_13465",
2263 }
2264 }
2265 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2267 match value {
2268 "AGG_NODE_VERSION_UNSPECIFIED" => Some(Self::Unspecified),
2269 "AGG_NODE_VERSION_ISSUE_12140" => Some(Self::Issue12140),
2270 "AGG_NODE_VERSION_ISSUE_13465" => Some(Self::Issue13465),
2271 _ => None,
2272 }
2273 }
2274}
2275#[derive(prost_helpers::AnyPB)]
2276#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2277#[repr(i32)]
2278pub enum InequalityType {
2279 Unspecified = 0,
2280 LessThan = 1,
2281 LessThanOrEqual = 2,
2282 GreaterThan = 3,
2283 GreaterThanOrEqual = 4,
2284}
2285impl InequalityType {
2286 pub fn as_str_name(&self) -> &'static str {
2291 match self {
2292 Self::Unspecified => "INEQUALITY_TYPE_UNSPECIFIED",
2293 Self::LessThan => "INEQUALITY_TYPE_LESS_THAN",
2294 Self::LessThanOrEqual => "INEQUALITY_TYPE_LESS_THAN_OR_EQUAL",
2295 Self::GreaterThan => "INEQUALITY_TYPE_GREATER_THAN",
2296 Self::GreaterThanOrEqual => "INEQUALITY_TYPE_GREATER_THAN_OR_EQUAL",
2297 }
2298 }
2299 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2301 match value {
2302 "INEQUALITY_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2303 "INEQUALITY_TYPE_LESS_THAN" => Some(Self::LessThan),
2304 "INEQUALITY_TYPE_LESS_THAN_OR_EQUAL" => Some(Self::LessThanOrEqual),
2305 "INEQUALITY_TYPE_GREATER_THAN" => Some(Self::GreaterThan),
2306 "INEQUALITY_TYPE_GREATER_THAN_OR_EQUAL" => Some(Self::GreaterThanOrEqual),
2307 _ => None,
2308 }
2309 }
2310}
2311#[derive(prost_helpers::AnyPB)]
2312#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2313#[repr(i32)]
2314pub enum JoinEncodingType {
2315 Unspecified = 0,
2316 MemoryOptimized = 1,
2317 CpuOptimized = 2,
2318}
2319impl JoinEncodingType {
2320 pub fn as_str_name(&self) -> &'static str {
2325 match self {
2326 Self::Unspecified => "UNSPECIFIED",
2327 Self::MemoryOptimized => "MEMORY_OPTIMIZED",
2328 Self::CpuOptimized => "CPU_OPTIMIZED",
2329 }
2330 }
2331 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2333 match value {
2334 "UNSPECIFIED" => Some(Self::Unspecified),
2335 "MEMORY_OPTIMIZED" => Some(Self::MemoryOptimized),
2336 "CPU_OPTIMIZED" => Some(Self::CpuOptimized),
2337 _ => None,
2338 }
2339 }
2340}
2341#[derive(prost_helpers::AnyPB)]
2343#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2344#[repr(i32)]
2345pub enum StreamScanType {
2346 Unspecified = 0,
2347 #[deprecated]
2349 Chain = 1,
2350 #[deprecated]
2352 Rearrange = 2,
2353 #[deprecated]
2355 Backfill = 3,
2356 UpstreamOnly = 4,
2358 ArrangementBackfill = 5,
2360 SnapshotBackfill = 6,
2362 CrossDbSnapshotBackfill = 7,
2364}
2365impl StreamScanType {
2366 pub fn as_str_name(&self) -> &'static str {
2371 match self {
2372 Self::Unspecified => "STREAM_SCAN_TYPE_UNSPECIFIED",
2373 #[allow(deprecated)]
2374 Self::Chain => "STREAM_SCAN_TYPE_CHAIN",
2375 #[allow(deprecated)]
2376 Self::Rearrange => "STREAM_SCAN_TYPE_REARRANGE",
2377 #[allow(deprecated)]
2378 Self::Backfill => "STREAM_SCAN_TYPE_BACKFILL",
2379 Self::UpstreamOnly => "STREAM_SCAN_TYPE_UPSTREAM_ONLY",
2380 Self::ArrangementBackfill => "STREAM_SCAN_TYPE_ARRANGEMENT_BACKFILL",
2381 Self::SnapshotBackfill => "STREAM_SCAN_TYPE_SNAPSHOT_BACKFILL",
2382 Self::CrossDbSnapshotBackfill => {
2383 "STREAM_SCAN_TYPE_CROSS_DB_SNAPSHOT_BACKFILL"
2384 }
2385 }
2386 }
2387 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2389 match value {
2390 "STREAM_SCAN_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2391 "STREAM_SCAN_TYPE_CHAIN" => Some(#[allow(deprecated)] Self::Chain),
2392 "STREAM_SCAN_TYPE_REARRANGE" => Some(#[allow(deprecated)] Self::Rearrange),
2393 "STREAM_SCAN_TYPE_BACKFILL" => Some(#[allow(deprecated)] Self::Backfill),
2394 "STREAM_SCAN_TYPE_UPSTREAM_ONLY" => Some(Self::UpstreamOnly),
2395 "STREAM_SCAN_TYPE_ARRANGEMENT_BACKFILL" => Some(Self::ArrangementBackfill),
2396 "STREAM_SCAN_TYPE_SNAPSHOT_BACKFILL" => Some(Self::SnapshotBackfill),
2397 "STREAM_SCAN_TYPE_CROSS_DB_SNAPSHOT_BACKFILL" => {
2398 Some(Self::CrossDbSnapshotBackfill)
2399 }
2400 _ => None,
2401 }
2402 }
2403}
2404#[derive(prost_helpers::AnyPB)]
2405#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2406#[repr(i32)]
2407pub enum OverWindowCachePolicy {
2408 Unspecified = 0,
2409 Full = 1,
2410 Recent = 2,
2411 RecentFirstN = 3,
2412 RecentLastN = 4,
2413}
2414impl OverWindowCachePolicy {
2415 pub fn as_str_name(&self) -> &'static str {
2420 match self {
2421 Self::Unspecified => "OVER_WINDOW_CACHE_POLICY_UNSPECIFIED",
2422 Self::Full => "OVER_WINDOW_CACHE_POLICY_FULL",
2423 Self::Recent => "OVER_WINDOW_CACHE_POLICY_RECENT",
2424 Self::RecentFirstN => "OVER_WINDOW_CACHE_POLICY_RECENT_FIRST_N",
2425 Self::RecentLastN => "OVER_WINDOW_CACHE_POLICY_RECENT_LAST_N",
2426 }
2427 }
2428 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2430 match value {
2431 "OVER_WINDOW_CACHE_POLICY_UNSPECIFIED" => Some(Self::Unspecified),
2432 "OVER_WINDOW_CACHE_POLICY_FULL" => Some(Self::Full),
2433 "OVER_WINDOW_CACHE_POLICY_RECENT" => Some(Self::Recent),
2434 "OVER_WINDOW_CACHE_POLICY_RECENT_FIRST_N" => Some(Self::RecentFirstN),
2435 "OVER_WINDOW_CACHE_POLICY_RECENT_LAST_N" => Some(Self::RecentLastN),
2436 _ => None,
2437 }
2438 }
2439}
2440#[derive(prost_helpers::AnyPB)]
2441#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2442#[repr(i32)]
2443pub enum DispatcherType {
2444 Unspecified = 0,
2445 Hash = 1,
2447 Broadcast = 2,
2452 Simple = 3,
2454 NoShuffle = 4,
2458}
2459impl DispatcherType {
2460 pub fn as_str_name(&self) -> &'static str {
2465 match self {
2466 Self::Unspecified => "DISPATCHER_TYPE_UNSPECIFIED",
2467 Self::Hash => "DISPATCHER_TYPE_HASH",
2468 Self::Broadcast => "DISPATCHER_TYPE_BROADCAST",
2469 Self::Simple => "DISPATCHER_TYPE_SIMPLE",
2470 Self::NoShuffle => "DISPATCHER_TYPE_NO_SHUFFLE",
2471 }
2472 }
2473 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2475 match value {
2476 "DISPATCHER_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2477 "DISPATCHER_TYPE_HASH" => Some(Self::Hash),
2478 "DISPATCHER_TYPE_BROADCAST" => Some(Self::Broadcast),
2479 "DISPATCHER_TYPE_SIMPLE" => Some(Self::Simple),
2480 "DISPATCHER_TYPE_NO_SHUFFLE" => Some(Self::NoShuffle),
2481 _ => None,
2482 }
2483 }
2484}