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)]
1756#[derive(Clone, PartialEq, ::prost::Message)]
1757pub struct MatchRecognizeNode {
1758 #[prost(uint32, repeated, tag = "1")]
1759 pub partition_by: ::prost::alloc::vec::Vec<u32>,
1760 #[prost(message, repeated, tag = "2")]
1763 pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1764 #[prost(message, repeated, tag = "3")]
1765 pub measures: ::prost::alloc::vec::Vec<MatchRecognizeMeasure>,
1766 #[prost(message, repeated, tag = "5")]
1767 pub defines: ::prost::alloc::vec::Vec<MatchRecognizeDefine>,
1768 #[prost(message, optional, tag = "12")]
1772 pub pattern_node: ::core::option::Option<MatchRecognizePatternNode>,
1773 #[prost(message, optional, tag = "8")]
1777 pub state_table: ::core::option::Option<super::catalog::Table>,
1778 #[prost(message, optional, tag = "9")]
1780 pub after_match_skip: ::core::option::Option<MatchRecognizeAfterMatchSkip>,
1781 #[prost(message, optional, tag = "11")]
1784 pub within: ::core::option::Option<super::expr::ExprNode>,
1785 #[prost(message, optional, tag = "15")]
1789 pub within_deadline: ::core::option::Option<super::expr::ExprNode>,
1790 #[prost(enumeration = "MatchRecognizeInputMode", tag = "16")]
1794 pub input_mode: i32,
1795}
1796#[derive(prost_helpers::AnyPB)]
1800#[derive(Clone, PartialEq, ::prost::Message)]
1801pub struct MatchRecognizePatternNode {
1802 #[prost(oneof = "match_recognize_pattern_node::Node", tags = "1, 2, 3, 4, 5")]
1803 pub node: ::core::option::Option<match_recognize_pattern_node::Node>,
1804}
1805pub mod match_recognize_pattern_node {
1807 #[derive(prost_helpers::AnyPB)]
1808 #[derive(Clone, PartialEq, ::prost::Oneof)]
1809 pub enum Node {
1810 #[prost(string, tag = "1")]
1812 Var(::prost::alloc::string::String),
1813 #[prost(message, tag = "2")]
1815 Concat(super::MatchRecognizePatternSeq),
1816 #[prost(message, tag = "3")]
1818 Alternation(super::MatchRecognizePatternSeq),
1819 #[prost(message, tag = "4")]
1821 Quantified(::prost::alloc::boxed::Box<super::MatchRecognizeQuantifiedPattern>),
1822 #[prost(message, tag = "5")]
1824 Permute(super::MatchRecognizePermutePattern),
1825 }
1826}
1827#[derive(prost_helpers::AnyPB)]
1829#[derive(Clone, PartialEq, ::prost::Message)]
1830pub struct MatchRecognizePatternSeq {
1831 #[prost(message, repeated, tag = "1")]
1832 pub patterns: ::prost::alloc::vec::Vec<MatchRecognizePatternNode>,
1833}
1834#[derive(prost_helpers::AnyPB)]
1836#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1837pub struct MatchRecognizePermutePattern {
1838 #[prost(string, repeated, tag = "1")]
1839 pub vars: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
1840}
1841#[derive(prost_helpers::AnyPB)]
1843#[derive(Clone, PartialEq, ::prost::Message)]
1844pub struct MatchRecognizeQuantifiedPattern {
1845 #[prost(message, optional, boxed, tag = "1")]
1846 pub inner: ::core::option::Option<
1847 ::prost::alloc::boxed::Box<MatchRecognizePatternNode>,
1848 >,
1849 #[prost(message, optional, tag = "2")]
1850 pub quantifier: ::core::option::Option<MatchRecognizeQuantifier>,
1851 #[prost(bool, tag = "3")]
1853 pub reluctant: bool,
1854}
1855#[derive(prost_helpers::AnyPB)]
1857#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1858pub struct MatchRecognizeQuantifier {
1859 #[prost(enumeration = "match_recognize_quantifier::Kind", tag = "1")]
1860 pub kind: i32,
1861 #[prost(uint32, tag = "2")]
1863 pub min: u32,
1864 #[prost(uint32, optional, tag = "3")]
1866 pub max: ::core::option::Option<u32>,
1867}
1868pub mod match_recognize_quantifier {
1870 #[derive(prost_helpers::AnyPB)]
1871 #[derive(
1872 Clone,
1873 Copy,
1874 Debug,
1875 PartialEq,
1876 Eq,
1877 Hash,
1878 PartialOrd,
1879 Ord,
1880 ::prost::Enumeration
1881 )]
1882 #[repr(i32)]
1883 pub enum Kind {
1884 Unspecified = 0,
1885 Star = 1,
1887 Plus = 2,
1889 Question = 3,
1891 Range = 4,
1893 }
1894 impl Kind {
1895 pub fn as_str_name(&self) -> &'static str {
1900 match self {
1901 Self::Unspecified => "KIND_UNSPECIFIED",
1902 Self::Star => "KIND_STAR",
1903 Self::Plus => "KIND_PLUS",
1904 Self::Question => "KIND_QUESTION",
1905 Self::Range => "KIND_RANGE",
1906 }
1907 }
1908 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
1910 match value {
1911 "KIND_UNSPECIFIED" => Some(Self::Unspecified),
1912 "KIND_STAR" => Some(Self::Star),
1913 "KIND_PLUS" => Some(Self::Plus),
1914 "KIND_QUESTION" => Some(Self::Question),
1915 "KIND_RANGE" => Some(Self::Range),
1916 _ => None,
1917 }
1918 }
1919 }
1920}
1921#[derive(prost_helpers::AnyPB)]
1925#[derive(Clone, PartialEq, ::prost::Message)]
1926pub struct MatchRecognizeDefine {
1927 #[prost(string, tag = "1")]
1928 pub symbol: ::prost::alloc::string::String,
1929 #[prost(message, optional, tag = "2")]
1931 pub condition: ::core::option::Option<super::expr::ExprNode>,
1932 #[prost(message, repeated, tag = "3")]
1933 pub slots: ::prost::alloc::vec::Vec<MatchRecognizeDefineSlot>,
1934}
1935#[derive(prost_helpers::AnyPB)]
1937#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1938pub struct MatchRecognizeAfterMatchSkip {
1939 #[prost(enumeration = "match_recognize_after_match_skip::Mode", tag = "1")]
1940 pub mode: i32,
1941 #[prost(string, optional, tag = "2")]
1943 pub target: ::core::option::Option<::prost::alloc::string::String>,
1944}
1945pub mod match_recognize_after_match_skip {
1947 #[derive(prost_helpers::AnyPB)]
1948 #[derive(
1949 Clone,
1950 Copy,
1951 Debug,
1952 PartialEq,
1953 Eq,
1954 Hash,
1955 PartialOrd,
1956 Ord,
1957 ::prost::Enumeration
1958 )]
1959 #[repr(i32)]
1960 pub enum Mode {
1961 Unspecified = 0,
1962 PastLastRow = 1,
1964 ToNextRow = 2,
1966 ToFirst = 3,
1968 ToLast = 4,
1969 }
1970 impl Mode {
1971 pub fn as_str_name(&self) -> &'static str {
1976 match self {
1977 Self::Unspecified => "MODE_UNSPECIFIED",
1978 Self::PastLastRow => "MODE_PAST_LAST_ROW",
1979 Self::ToNextRow => "MODE_TO_NEXT_ROW",
1980 Self::ToFirst => "MODE_TO_FIRST",
1981 Self::ToLast => "MODE_TO_LAST",
1982 }
1983 }
1984 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
1986 match value {
1987 "MODE_UNSPECIFIED" => Some(Self::Unspecified),
1988 "MODE_PAST_LAST_ROW" => Some(Self::PastLastRow),
1989 "MODE_TO_NEXT_ROW" => Some(Self::ToNextRow),
1990 "MODE_TO_FIRST" => Some(Self::ToFirst),
1991 "MODE_TO_LAST" => Some(Self::ToLast),
1992 _ => None,
1993 }
1994 }
1995 }
1996}
1997#[derive(prost_helpers::AnyPB)]
1999#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2000pub struct MatchRecognizeDefineSlot {
2001 #[prost(enumeration = "match_recognize_define_slot::Kind", tag = "1")]
2002 pub kind: i32,
2003 #[prost(string, repeated, tag = "2")]
2004 pub vars: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
2005 #[prost(uint32, tag = "3")]
2006 pub col_idx: u32,
2007 #[prost(uint32, tag = "4")]
2008 pub offset: u32,
2009}
2010pub mod match_recognize_define_slot {
2012 #[derive(prost_helpers::AnyPB)]
2013 #[derive(
2014 Clone,
2015 Copy,
2016 Debug,
2017 PartialEq,
2018 Eq,
2019 Hash,
2020 PartialOrd,
2021 Ord,
2022 ::prost::Enumeration
2023 )]
2024 #[repr(i32)]
2025 pub enum Kind {
2026 Unspecified = 0,
2027 SelfCol = 1,
2029 Prev = 2,
2031 Next = 3,
2032 RunningFirst = 4,
2037 RunningLast = 5,
2038 }
2039 impl Kind {
2040 pub fn as_str_name(&self) -> &'static str {
2045 match self {
2046 Self::Unspecified => "KIND_UNSPECIFIED",
2047 Self::SelfCol => "KIND_SELF_COL",
2048 Self::Prev => "KIND_PREV",
2049 Self::Next => "KIND_NEXT",
2050 Self::RunningFirst => "KIND_RUNNING_FIRST",
2051 Self::RunningLast => "KIND_RUNNING_LAST",
2052 }
2053 }
2054 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2056 match value {
2057 "KIND_UNSPECIFIED" => Some(Self::Unspecified),
2058 "KIND_SELF_COL" => Some(Self::SelfCol),
2059 "KIND_PREV" => Some(Self::Prev),
2060 "KIND_NEXT" => Some(Self::Next),
2061 "KIND_RUNNING_FIRST" => Some(Self::RunningFirst),
2062 "KIND_RUNNING_LAST" => Some(Self::RunningLast),
2063 _ => None,
2064 }
2065 }
2066 }
2067}
2068#[derive(prost_helpers::AnyPB)]
2072#[derive(Clone, PartialEq, ::prost::Message)]
2073pub struct MatchRecognizeMeasure {
2074 #[prost(message, optional, tag = "1")]
2076 pub expr: ::core::option::Option<super::expr::ExprNode>,
2077 #[prost(string, tag = "2")]
2078 pub name: ::prost::alloc::string::String,
2079 #[prost(message, repeated, tag = "3")]
2080 pub slots: ::prost::alloc::vec::Vec<MatchRecognizeMeasureSlot>,
2081}
2082#[derive(prost_helpers::AnyPB)]
2084#[derive(Clone, PartialEq, ::prost::Message)]
2085pub struct MatchRecognizeMeasureSlot {
2086 #[prost(enumeration = "match_recognize_measure_slot::Kind", tag = "1")]
2087 pub kind: i32,
2088 #[prost(string, repeated, tag = "2")]
2091 pub vars: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
2092 #[prost(uint32, tag = "3")]
2093 pub col_idx: u32,
2094 #[prost(message, optional, tag = "4")]
2095 pub data_type: ::core::option::Option<super::data::DataType>,
2096 #[prost(message, optional, tag = "5")]
2099 pub agg_call: ::core::option::Option<super::expr::AggCall>,
2100}
2101pub mod match_recognize_measure_slot {
2103 #[derive(prost_helpers::AnyPB)]
2104 #[derive(
2105 Clone,
2106 Copy,
2107 Debug,
2108 PartialEq,
2109 Eq,
2110 Hash,
2111 PartialOrd,
2112 Ord,
2113 ::prost::Enumeration
2114 )]
2115 #[repr(i32)]
2116 pub enum Kind {
2117 Unspecified = 0,
2118 Last = 1,
2121 First = 2,
2123 Classifier = 3,
2125 CountStar = 4,
2127 Count = 5,
2129 Min = 6,
2131 Max = 7,
2132 Sum = 8,
2135 }
2136 impl Kind {
2137 pub fn as_str_name(&self) -> &'static str {
2142 match self {
2143 Self::Unspecified => "KIND_UNSPECIFIED",
2144 Self::Last => "KIND_LAST",
2145 Self::First => "KIND_FIRST",
2146 Self::Classifier => "KIND_CLASSIFIER",
2147 Self::CountStar => "KIND_COUNT_STAR",
2148 Self::Count => "KIND_COUNT",
2149 Self::Min => "KIND_MIN",
2150 Self::Max => "KIND_MAX",
2151 Self::Sum => "KIND_SUM",
2152 }
2153 }
2154 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2156 match value {
2157 "KIND_UNSPECIFIED" => Some(Self::Unspecified),
2158 "KIND_LAST" => Some(Self::Last),
2159 "KIND_FIRST" => Some(Self::First),
2160 "KIND_CLASSIFIER" => Some(Self::Classifier),
2161 "KIND_COUNT_STAR" => Some(Self::CountStar),
2162 "KIND_COUNT" => Some(Self::Count),
2163 "KIND_MIN" => Some(Self::Min),
2164 "KIND_MAX" => Some(Self::Max),
2165 "KIND_SUM" => Some(Self::Sum),
2166 _ => None,
2167 }
2168 }
2169 }
2170}
2171#[derive(prost_helpers::AnyPB)]
2172#[derive(Clone, PartialEq, ::prost::Message)]
2173pub struct StreamNode {
2174 #[prost(uint64, tag = "1", wrapper = "crate::id::StreamNodeLocalOperatorId")]
2177 pub operator_id: crate::id::StreamNodeLocalOperatorId,
2178 #[prost(message, repeated, tag = "3")]
2180 pub input: ::prost::alloc::vec::Vec<StreamNode>,
2181 #[prost(uint32, repeated, tag = "2")]
2182 pub stream_key: ::prost::alloc::vec::Vec<u32>,
2183 #[prost(enumeration = "stream_node::StreamKind", tag = "24")]
2184 pub stream_kind: i32,
2185 #[prost(string, tag = "18")]
2186 pub identity: ::prost::alloc::string::String,
2187 #[prost(message, repeated, tag = "19")]
2189 pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2190 #[prost(
2191 oneof = "stream_node::NodeBody",
2192 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"
2193 )]
2194 pub node_body: ::core::option::Option<stream_node::NodeBody>,
2195}
2196pub mod stream_node {
2198 #[derive(prost_helpers::AnyPB)]
2201 #[derive(
2202 Clone,
2203 Copy,
2204 Debug,
2205 PartialEq,
2206 Eq,
2207 Hash,
2208 PartialOrd,
2209 Ord,
2210 ::prost::Enumeration
2211 )]
2212 #[repr(i32)]
2213 pub enum StreamKind {
2214 Retract = 0,
2216 AppendOnly = 1,
2217 Upsert = 2,
2218 }
2219 impl StreamKind {
2220 pub fn as_str_name(&self) -> &'static str {
2225 match self {
2226 Self::Retract => "STREAM_KIND_RETRACT",
2227 Self::AppendOnly => "STREAM_KIND_APPEND_ONLY",
2228 Self::Upsert => "STREAM_KIND_UPSERT",
2229 }
2230 }
2231 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2233 match value {
2234 "STREAM_KIND_RETRACT" => Some(Self::Retract),
2235 "STREAM_KIND_APPEND_ONLY" => Some(Self::AppendOnly),
2236 "STREAM_KIND_UPSERT" => Some(Self::Upsert),
2237 _ => None,
2238 }
2239 }
2240 }
2241 #[derive(prost_helpers::AnyPB)]
2242 #[derive(::enum_as_inner::EnumAsInner, ::strum::Display, ::strum::EnumDiscriminants)]
2243 #[derive(::prost_helpers::StreamNodeBodyVariants)]
2244 #[strum_discriminants(derive(::strum::Display, Hash))]
2245 #[derive(Clone, PartialEq, ::prost::Oneof)]
2246 pub enum NodeBody {
2247 #[prost(message, tag = "100")]
2248 Source(::prost::alloc::boxed::Box<super::SourceNode>),
2249 #[prost(message, tag = "101")]
2250 Project(::prost::alloc::boxed::Box<super::ProjectNode>),
2251 #[prost(message, tag = "102")]
2252 Filter(::prost::alloc::boxed::Box<super::FilterNode>),
2253 #[prost(message, tag = "103")]
2254 Materialize(::prost::alloc::boxed::Box<super::MaterializeNode>),
2255 #[prost(message, tag = "104")]
2256 StatelessSimpleAgg(::prost::alloc::boxed::Box<super::SimpleAggNode>),
2257 #[prost(message, tag = "105")]
2258 SimpleAgg(::prost::alloc::boxed::Box<super::SimpleAggNode>),
2259 #[prost(message, tag = "106")]
2260 HashAgg(::prost::alloc::boxed::Box<super::HashAggNode>),
2261 #[prost(message, tag = "107")]
2262 AppendOnlyTopN(::prost::alloc::boxed::Box<super::TopNNode>),
2263 #[prost(message, tag = "108")]
2264 HashJoin(::prost::alloc::boxed::Box<super::HashJoinNode>),
2265 #[prost(message, tag = "109")]
2266 TopN(::prost::alloc::boxed::Box<super::TopNNode>),
2267 #[prost(message, tag = "110")]
2268 HopWindow(::prost::alloc::boxed::Box<super::HopWindowNode>),
2269 #[prost(message, tag = "111")]
2270 Merge(::prost::alloc::boxed::Box<super::MergeNode>),
2271 #[prost(message, tag = "112")]
2272 Exchange(::prost::alloc::boxed::Box<super::ExchangeNode>),
2273 #[prost(message, tag = "113")]
2274 StreamScan(::prost::alloc::boxed::Box<super::StreamScanNode>),
2275 #[prost(message, tag = "114")]
2276 BatchPlan(::prost::alloc::boxed::Box<super::BatchPlanNode>),
2277 #[prost(message, tag = "115")]
2278 Lookup(::prost::alloc::boxed::Box<super::LookupNode>),
2279 #[prost(message, tag = "116")]
2280 Arrange(::prost::alloc::boxed::Box<super::ArrangeNode>),
2281 #[prost(message, tag = "117")]
2282 LookupUnion(::prost::alloc::boxed::Box<super::LookupUnionNode>),
2283 #[prost(message, tag = "118")]
2284 Union(super::UnionNode),
2285 #[prost(message, tag = "119")]
2286 DeltaIndexJoin(::prost::alloc::boxed::Box<super::DeltaIndexJoinNode>),
2287 #[prost(message, tag = "120")]
2288 Sink(::prost::alloc::boxed::Box<super::SinkNode>),
2289 #[prost(message, tag = "121")]
2290 Expand(::prost::alloc::boxed::Box<super::ExpandNode>),
2291 #[prost(message, tag = "122")]
2292 DynamicFilter(::prost::alloc::boxed::Box<super::DynamicFilterNode>),
2293 #[prost(message, tag = "123")]
2294 ProjectSet(::prost::alloc::boxed::Box<super::ProjectSetNode>),
2295 #[prost(message, tag = "124")]
2296 GroupTopN(::prost::alloc::boxed::Box<super::GroupTopNNode>),
2297 #[prost(message, tag = "125")]
2298 Sort(::prost::alloc::boxed::Box<super::SortNode>),
2299 #[prost(message, tag = "126")]
2300 WatermarkFilter(::prost::alloc::boxed::Box<super::WatermarkFilterNode>),
2301 #[prost(message, tag = "127")]
2302 Dml(::prost::alloc::boxed::Box<super::DmlNode>),
2303 #[prost(message, tag = "128")]
2304 RowIdGen(::prost::alloc::boxed::Box<super::RowIdGenNode>),
2305 #[prost(message, tag = "129")]
2306 Now(::prost::alloc::boxed::Box<super::NowNode>),
2307 #[prost(message, tag = "130")]
2308 AppendOnlyGroupTopN(::prost::alloc::boxed::Box<super::GroupTopNNode>),
2309 #[prost(message, tag = "131")]
2310 TemporalJoin(::prost::alloc::boxed::Box<super::TemporalJoinNode>),
2311 #[prost(message, tag = "132")]
2312 BarrierRecv(::prost::alloc::boxed::Box<super::BarrierRecvNode>),
2313 #[prost(message, tag = "133")]
2314 Values(::prost::alloc::boxed::Box<super::ValuesNode>),
2315 #[prost(message, tag = "134")]
2316 AppendOnlyDedup(::prost::alloc::boxed::Box<super::DedupNode>),
2317 #[prost(message, tag = "135")]
2318 NoOp(super::NoOpNode),
2319 #[prost(message, tag = "136")]
2320 EowcOverWindow(::prost::alloc::boxed::Box<super::EowcOverWindowNode>),
2321 #[prost(message, tag = "137")]
2322 OverWindow(::prost::alloc::boxed::Box<super::OverWindowNode>),
2323 #[prost(message, tag = "138")]
2324 StreamFsFetch(::prost::alloc::boxed::Box<super::StreamFsFetchNode>),
2325 #[prost(message, tag = "139")]
2326 StreamCdcScan(::prost::alloc::boxed::Box<super::StreamCdcScanNode>),
2327 #[prost(message, tag = "140")]
2328 CdcFilter(::prost::alloc::boxed::Box<super::CdcFilterNode>),
2329 #[prost(message, tag = "142")]
2330 SourceBackfill(::prost::alloc::boxed::Box<super::SourceBackfillNode>),
2331 #[prost(message, tag = "143")]
2332 Changelog(::prost::alloc::boxed::Box<super::ChangeLogNode>),
2333 #[prost(message, tag = "144")]
2334 LocalApproxPercentile(
2335 ::prost::alloc::boxed::Box<super::LocalApproxPercentileNode>,
2336 ),
2337 #[prost(message, tag = "145")]
2338 GlobalApproxPercentile(
2339 ::prost::alloc::boxed::Box<super::GlobalApproxPercentileNode>,
2340 ),
2341 #[prost(message, tag = "146")]
2342 RowMerge(::prost::alloc::boxed::Box<super::RowMergeNode>),
2343 #[prost(message, tag = "147")]
2344 AsOfJoin(::prost::alloc::boxed::Box<super::AsOfJoinNode>),
2345 #[prost(message, tag = "148")]
2346 SyncLogStore(::prost::alloc::boxed::Box<super::SyncLogStoreNode>),
2347 #[prost(message, tag = "149")]
2348 MaterializedExprs(::prost::alloc::boxed::Box<super::MaterializedExprsNode>),
2349 #[prost(message, tag = "150")]
2350 VectorIndexWrite(::prost::alloc::boxed::Box<super::VectorIndexWriteNode>),
2351 #[prost(message, tag = "151")]
2352 UpstreamSinkUnion(::prost::alloc::boxed::Box<super::UpstreamSinkUnionNode>),
2353 #[prost(message, tag = "152")]
2354 LocalityProvider(::prost::alloc::boxed::Box<super::LocalityProviderNode>),
2355 #[prost(message, tag = "153")]
2356 EowcGapFill(::prost::alloc::boxed::Box<super::EowcGapFillNode>),
2357 #[prost(message, tag = "154")]
2358 GapFill(::prost::alloc::boxed::Box<super::GapFillNode>),
2359 #[prost(message, tag = "155")]
2360 VectorIndexLookupJoin(
2361 ::prost::alloc::boxed::Box<super::VectorIndexLookupJoinNode>,
2362 ),
2363 #[prost(message, tag = "156")]
2364 IcebergWithPkIndexWriter(
2365 ::prost::alloc::boxed::Box<super::IcebergWithPkIndexWriterNode>,
2366 ),
2367 #[prost(message, tag = "157")]
2368 IcebergWithPkIndexPositionDeleteMerger(
2369 ::prost::alloc::boxed::Box<super::IcebergWithPkIndexPositionDeleteMergerNode>,
2370 ),
2371 #[prost(message, tag = "158")]
2372 MatchRecognize(::prost::alloc::boxed::Box<super::MatchRecognizeNode>),
2373 #[prost(message, tag = "159")]
2374 CompactionResolver(::prost::alloc::boxed::Box<super::CompactionResolverNode>),
2375 }
2376}
2377#[derive(prost_helpers::AnyPB)]
2391#[derive(Clone, PartialEq, ::prost::Message)]
2392pub struct DispatchOutputMapping {
2393 #[prost(uint32, repeated, tag = "1")]
2395 pub indices: ::prost::alloc::vec::Vec<u32>,
2396 #[prost(message, repeated, tag = "2")]
2402 pub types: ::prost::alloc::vec::Vec<dispatch_output_mapping::TypePair>,
2403}
2404pub mod dispatch_output_mapping {
2406 #[derive(prost_helpers::AnyPB)]
2407 #[derive(Clone, PartialEq, ::prost::Message)]
2408 pub struct TypePair {
2409 #[prost(message, optional, tag = "1")]
2410 pub upstream: ::core::option::Option<super::super::data::DataType>,
2411 #[prost(message, optional, tag = "2")]
2412 pub downstream: ::core::option::Option<super::super::data::DataType>,
2413 }
2414}
2415#[derive(prost_helpers::AnyPB)]
2418#[derive(Clone, PartialEq, ::prost::Message)]
2419pub struct DispatchStrategy {
2420 #[prost(enumeration = "DispatcherType", tag = "1")]
2421 pub r#type: i32,
2422 #[prost(uint32, repeated, tag = "2")]
2423 pub dist_key_indices: ::prost::alloc::vec::Vec<u32>,
2424 #[prost(message, optional, tag = "3")]
2425 pub output_mapping: ::core::option::Option<DispatchOutputMapping>,
2426}
2427#[derive(prost_helpers::AnyPB)]
2430#[derive(Clone, PartialEq, ::prost::Message)]
2431pub struct Dispatcher {
2432 #[prost(enumeration = "DispatcherType", tag = "1")]
2433 pub r#type: i32,
2434 #[prost(uint32, repeated, tag = "2")]
2437 pub dist_key_indices: ::prost::alloc::vec::Vec<u32>,
2438 #[prost(message, optional, tag = "6")]
2440 pub output_mapping: ::core::option::Option<DispatchOutputMapping>,
2441 #[prost(message, optional, tag = "3")]
2444 pub hash_mapping: ::core::option::Option<ActorMapping>,
2445 #[prost(uint32, tag = "4", wrapper = "crate::id::FragmentId")]
2448 pub dispatcher_id: crate::id::FragmentId,
2449 #[prost(uint32, repeated, tag = "5", wrapper = "crate::id::ActorId")]
2451 pub downstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
2452}
2453#[derive(prost_helpers::AnyPB)]
2455#[derive(Clone, PartialEq, ::prost::Message)]
2456pub struct StreamActor {
2457 #[prost(uint32, tag = "1", wrapper = "crate::id::ActorId")]
2458 pub actor_id: crate::id::ActorId,
2459 #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
2460 pub fragment_id: crate::id::FragmentId,
2461 #[prost(message, repeated, tag = "4")]
2462 pub dispatcher: ::prost::alloc::vec::Vec<Dispatcher>,
2463 #[prost(message, optional, tag = "8")]
2466 pub vnode_bitmap: ::core::option::Option<super::common::Buffer>,
2467 #[prost(string, tag = "9")]
2469 pub mview_definition: ::prost::alloc::string::String,
2470 #[prost(message, optional, tag = "10")]
2472 pub expr_context: ::core::option::Option<super::plan_common::ExprContext>,
2473 #[prost(string, tag = "11")]
2475 pub config_override: ::prost::alloc::string::String,
2476}
2477#[derive(prost_helpers::AnyPB)]
2479#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2480pub struct StreamContext {
2481 #[prost(string, tag = "1")]
2483 pub timezone: ::prost::alloc::string::String,
2484 #[prost(string, tag = "2")]
2486 pub config_override: ::prost::alloc::string::String,
2487}
2488#[derive(prost_helpers::AnyPB)]
2489#[derive(Clone, PartialEq, ::prost::Message)]
2490pub struct BackfillOrder {
2491 #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::RelationId")]
2492 pub order: ::std::collections::HashMap<
2493 crate::id::RelationId,
2494 super::common::Uint32Vector,
2495 >,
2496}
2497#[derive(prost_helpers::AnyPB)]
2502#[derive(Clone, PartialEq, ::prost::Message)]
2503pub struct StreamFragmentGraph {
2504 #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::FragmentId")]
2506 pub fragments: ::std::collections::HashMap<
2507 crate::id::FragmentId,
2508 stream_fragment_graph::StreamFragment,
2509 >,
2510 #[prost(message, repeated, tag = "2")]
2512 pub edges: ::prost::alloc::vec::Vec<stream_fragment_graph::StreamFragmentEdge>,
2513 #[prost(uint32, repeated, tag = "3", wrapper = "crate::id::TableId")]
2514 pub dependent_table_ids: ::prost::alloc::vec::Vec<crate::id::TableId>,
2515 #[prost(uint32, tag = "4")]
2516 pub table_ids_cnt: u32,
2517 #[prost(message, optional, tag = "5")]
2518 pub ctx: ::core::option::Option<StreamContext>,
2519 #[prost(message, optional, tag = "6")]
2521 pub parallelism: ::core::option::Option<stream_fragment_graph::Parallelism>,
2522 #[prost(message, optional, tag = "9")]
2524 pub backfill_parallelism: ::core::option::Option<stream_fragment_graph::Parallelism>,
2525 #[prost(string, tag = "10")]
2527 pub adaptive_parallelism_strategy: ::prost::alloc::string::String,
2528 #[prost(string, tag = "11")]
2530 pub backfill_adaptive_parallelism_strategy: ::prost::alloc::string::String,
2531 #[prost(uint32, tag = "7")]
2541 pub max_parallelism: u32,
2542 #[prost(message, optional, tag = "8")]
2544 pub backfill_order: ::core::option::Option<BackfillOrder>,
2545}
2546pub mod stream_fragment_graph {
2548 #[derive(prost_helpers::AnyPB)]
2549 #[derive(Clone, PartialEq, ::prost::Message)]
2550 pub struct StreamFragment {
2551 #[prost(uint32, tag = "1", wrapper = "crate::id::FragmentId")]
2553 pub fragment_id: crate::id::FragmentId,
2554 #[prost(message, optional, tag = "2")]
2556 pub node: ::core::option::Option<super::StreamNode>,
2557 #[prost(uint32, tag = "3")]
2559 pub fragment_type_mask: u32,
2560 #[prost(bool, tag = "4")]
2564 pub requires_singleton: bool,
2565 }
2566 #[derive(prost_helpers::AnyPB)]
2567 #[derive(Clone, PartialEq, ::prost::Message)]
2568 pub struct StreamFragmentEdge {
2569 #[prost(message, optional, tag = "1")]
2571 pub dispatch_strategy: ::core::option::Option<super::DispatchStrategy>,
2572 #[prost(uint64, tag = "3")]
2576 pub link_id: u64,
2577 #[prost(uint32, tag = "4", wrapper = "crate::id::FragmentId")]
2578 pub upstream_id: crate::id::FragmentId,
2579 #[prost(uint32, tag = "5", wrapper = "crate::id::FragmentId")]
2580 pub downstream_id: crate::id::FragmentId,
2581 }
2582 #[derive(prost_helpers::AnyPB)]
2583 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
2584 pub struct Parallelism {
2585 #[prost(uint64, tag = "1")]
2586 pub parallelism: u64,
2587 }
2588}
2589#[derive(prost_helpers::AnyPB)]
2591#[derive(Clone, PartialEq, ::prost::Message)]
2592pub struct SinkSchemaChange {
2593 #[prost(message, repeated, tag = "1")]
2596 pub original_schema: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2597 #[prost(oneof = "sink_schema_change::Op", tags = "2, 3")]
2599 pub op: ::core::option::Option<sink_schema_change::Op>,
2600}
2601pub mod sink_schema_change {
2603 #[derive(prost_helpers::AnyPB)]
2605 #[derive(Clone, PartialEq, ::prost::Oneof)]
2606 pub enum Op {
2607 #[prost(message, tag = "2")]
2609 AddColumns(super::SinkAddColumnsOp),
2610 #[prost(message, tag = "3")]
2612 DropColumns(super::SinkDropColumnsOp),
2613 }
2614}
2615#[derive(prost_helpers::AnyPB)]
2617#[derive(Clone, PartialEq, ::prost::Message)]
2618pub struct SinkAddColumnsOp {
2619 #[prost(message, repeated, tag = "1")]
2621 pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2622}
2623#[derive(prost_helpers::AnyPB)]
2625#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2626pub struct SinkDropColumnsOp {
2627 #[prost(string, repeated, tag = "1")]
2629 pub column_names: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
2630}
2631#[derive(prost_helpers::AnyPB)]
2632#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2633#[repr(i32)]
2634pub enum SinkLogStoreType {
2635 Unspecified = 0,
2637 KvLogStore = 1,
2638 InMemoryLogStore = 2,
2639}
2640impl SinkLogStoreType {
2641 pub fn as_str_name(&self) -> &'static str {
2646 match self {
2647 Self::Unspecified => "SINK_LOG_STORE_TYPE_UNSPECIFIED",
2648 Self::KvLogStore => "SINK_LOG_STORE_TYPE_KV_LOG_STORE",
2649 Self::InMemoryLogStore => "SINK_LOG_STORE_TYPE_IN_MEMORY_LOG_STORE",
2650 }
2651 }
2652 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2654 match value {
2655 "SINK_LOG_STORE_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2656 "SINK_LOG_STORE_TYPE_KV_LOG_STORE" => Some(Self::KvLogStore),
2657 "SINK_LOG_STORE_TYPE_IN_MEMORY_LOG_STORE" => Some(Self::InMemoryLogStore),
2658 _ => None,
2659 }
2660 }
2661}
2662#[derive(prost_helpers::AnyPB)]
2663#[derive(prost_helpers::Version)]
2664#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2665#[repr(i32)]
2666pub enum AggNodeVersion {
2667 Unspecified = 0,
2668 Issue12140 = 1,
2670 Issue13465 = 2,
2672}
2673impl AggNodeVersion {
2674 pub fn as_str_name(&self) -> &'static str {
2679 match self {
2680 Self::Unspecified => "AGG_NODE_VERSION_UNSPECIFIED",
2681 Self::Issue12140 => "AGG_NODE_VERSION_ISSUE_12140",
2682 Self::Issue13465 => "AGG_NODE_VERSION_ISSUE_13465",
2683 }
2684 }
2685 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2687 match value {
2688 "AGG_NODE_VERSION_UNSPECIFIED" => Some(Self::Unspecified),
2689 "AGG_NODE_VERSION_ISSUE_12140" => Some(Self::Issue12140),
2690 "AGG_NODE_VERSION_ISSUE_13465" => Some(Self::Issue13465),
2691 _ => None,
2692 }
2693 }
2694}
2695#[derive(prost_helpers::AnyPB)]
2696#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2697#[repr(i32)]
2698pub enum InequalityType {
2699 Unspecified = 0,
2700 LessThan = 1,
2701 LessThanOrEqual = 2,
2702 GreaterThan = 3,
2703 GreaterThanOrEqual = 4,
2704}
2705impl InequalityType {
2706 pub fn as_str_name(&self) -> &'static str {
2711 match self {
2712 Self::Unspecified => "INEQUALITY_TYPE_UNSPECIFIED",
2713 Self::LessThan => "INEQUALITY_TYPE_LESS_THAN",
2714 Self::LessThanOrEqual => "INEQUALITY_TYPE_LESS_THAN_OR_EQUAL",
2715 Self::GreaterThan => "INEQUALITY_TYPE_GREATER_THAN",
2716 Self::GreaterThanOrEqual => "INEQUALITY_TYPE_GREATER_THAN_OR_EQUAL",
2717 }
2718 }
2719 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2721 match value {
2722 "INEQUALITY_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2723 "INEQUALITY_TYPE_LESS_THAN" => Some(Self::LessThan),
2724 "INEQUALITY_TYPE_LESS_THAN_OR_EQUAL" => Some(Self::LessThanOrEqual),
2725 "INEQUALITY_TYPE_GREATER_THAN" => Some(Self::GreaterThan),
2726 "INEQUALITY_TYPE_GREATER_THAN_OR_EQUAL" => Some(Self::GreaterThanOrEqual),
2727 _ => None,
2728 }
2729 }
2730}
2731#[derive(prost_helpers::AnyPB)]
2732#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2733#[repr(i32)]
2734pub enum JoinEncodingType {
2735 Unspecified = 0,
2736 MemoryOptimized = 1,
2737 CpuOptimized = 2,
2738}
2739impl JoinEncodingType {
2740 pub fn as_str_name(&self) -> &'static str {
2745 match self {
2746 Self::Unspecified => "UNSPECIFIED",
2747 Self::MemoryOptimized => "MEMORY_OPTIMIZED",
2748 Self::CpuOptimized => "CPU_OPTIMIZED",
2749 }
2750 }
2751 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2753 match value {
2754 "UNSPECIFIED" => Some(Self::Unspecified),
2755 "MEMORY_OPTIMIZED" => Some(Self::MemoryOptimized),
2756 "CPU_OPTIMIZED" => Some(Self::CpuOptimized),
2757 _ => None,
2758 }
2759 }
2760}
2761#[derive(prost_helpers::AnyPB)]
2763#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2764#[repr(i32)]
2765pub enum StreamScanType {
2766 Unspecified = 0,
2767 #[deprecated]
2769 Chain = 1,
2770 #[deprecated]
2772 Rearrange = 2,
2773 #[deprecated]
2775 Backfill = 3,
2776 UpstreamOnly = 4,
2778 ArrangementBackfill = 5,
2780 SnapshotBackfill = 6,
2782 CrossDbSnapshotBackfill = 7,
2784}
2785impl StreamScanType {
2786 pub fn as_str_name(&self) -> &'static str {
2791 match self {
2792 Self::Unspecified => "STREAM_SCAN_TYPE_UNSPECIFIED",
2793 #[allow(deprecated)]
2794 Self::Chain => "STREAM_SCAN_TYPE_CHAIN",
2795 #[allow(deprecated)]
2796 Self::Rearrange => "STREAM_SCAN_TYPE_REARRANGE",
2797 #[allow(deprecated)]
2798 Self::Backfill => "STREAM_SCAN_TYPE_BACKFILL",
2799 Self::UpstreamOnly => "STREAM_SCAN_TYPE_UPSTREAM_ONLY",
2800 Self::ArrangementBackfill => "STREAM_SCAN_TYPE_ARRANGEMENT_BACKFILL",
2801 Self::SnapshotBackfill => "STREAM_SCAN_TYPE_SNAPSHOT_BACKFILL",
2802 Self::CrossDbSnapshotBackfill => {
2803 "STREAM_SCAN_TYPE_CROSS_DB_SNAPSHOT_BACKFILL"
2804 }
2805 }
2806 }
2807 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2809 match value {
2810 "STREAM_SCAN_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2811 "STREAM_SCAN_TYPE_CHAIN" => Some(#[allow(deprecated)] Self::Chain),
2812 "STREAM_SCAN_TYPE_REARRANGE" => Some(#[allow(deprecated)] Self::Rearrange),
2813 "STREAM_SCAN_TYPE_BACKFILL" => Some(#[allow(deprecated)] Self::Backfill),
2814 "STREAM_SCAN_TYPE_UPSTREAM_ONLY" => Some(Self::UpstreamOnly),
2815 "STREAM_SCAN_TYPE_ARRANGEMENT_BACKFILL" => Some(Self::ArrangementBackfill),
2816 "STREAM_SCAN_TYPE_SNAPSHOT_BACKFILL" => Some(Self::SnapshotBackfill),
2817 "STREAM_SCAN_TYPE_CROSS_DB_SNAPSHOT_BACKFILL" => {
2818 Some(Self::CrossDbSnapshotBackfill)
2819 }
2820 _ => None,
2821 }
2822 }
2823}
2824#[derive(prost_helpers::AnyPB)]
2825#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2826#[repr(i32)]
2827pub enum OverWindowCachePolicy {
2828 Unspecified = 0,
2829 Full = 1,
2830 Recent = 2,
2831 RecentFirstN = 3,
2832 RecentLastN = 4,
2833}
2834impl OverWindowCachePolicy {
2835 pub fn as_str_name(&self) -> &'static str {
2840 match self {
2841 Self::Unspecified => "OVER_WINDOW_CACHE_POLICY_UNSPECIFIED",
2842 Self::Full => "OVER_WINDOW_CACHE_POLICY_FULL",
2843 Self::Recent => "OVER_WINDOW_CACHE_POLICY_RECENT",
2844 Self::RecentFirstN => "OVER_WINDOW_CACHE_POLICY_RECENT_FIRST_N",
2845 Self::RecentLastN => "OVER_WINDOW_CACHE_POLICY_RECENT_LAST_N",
2846 }
2847 }
2848 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2850 match value {
2851 "OVER_WINDOW_CACHE_POLICY_UNSPECIFIED" => Some(Self::Unspecified),
2852 "OVER_WINDOW_CACHE_POLICY_FULL" => Some(Self::Full),
2853 "OVER_WINDOW_CACHE_POLICY_RECENT" => Some(Self::Recent),
2854 "OVER_WINDOW_CACHE_POLICY_RECENT_FIRST_N" => Some(Self::RecentFirstN),
2855 "OVER_WINDOW_CACHE_POLICY_RECENT_LAST_N" => Some(Self::RecentLastN),
2856 _ => None,
2857 }
2858 }
2859}
2860#[derive(prost_helpers::AnyPB)]
2861#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2862#[repr(i32)]
2863pub enum MatchRecognizeInputMode {
2864 Unspecified = 0,
2865 EventTime = 1,
2866 ProcessingTime = 2,
2867}
2868impl MatchRecognizeInputMode {
2869 pub fn as_str_name(&self) -> &'static str {
2874 match self {
2875 Self::Unspecified => "MATCH_RECOGNIZE_INPUT_MODE_UNSPECIFIED",
2876 Self::EventTime => "MATCH_RECOGNIZE_INPUT_MODE_EVENT_TIME",
2877 Self::ProcessingTime => "MATCH_RECOGNIZE_INPUT_MODE_PROCESSING_TIME",
2878 }
2879 }
2880 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2882 match value {
2883 "MATCH_RECOGNIZE_INPUT_MODE_UNSPECIFIED" => Some(Self::Unspecified),
2884 "MATCH_RECOGNIZE_INPUT_MODE_EVENT_TIME" => Some(Self::EventTime),
2885 "MATCH_RECOGNIZE_INPUT_MODE_PROCESSING_TIME" => Some(Self::ProcessingTime),
2886 _ => None,
2887 }
2888 }
2889}
2890#[derive(prost_helpers::AnyPB)]
2891#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2892#[repr(i32)]
2893pub enum DispatcherType {
2894 Unspecified = 0,
2895 Hash = 1,
2897 Broadcast = 2,
2902 Simple = 3,
2904 NoShuffle = 4,
2908}
2909impl DispatcherType {
2910 pub fn as_str_name(&self) -> &'static str {
2915 match self {
2916 Self::Unspecified => "DISPATCHER_TYPE_UNSPECIFIED",
2917 Self::Hash => "DISPATCHER_TYPE_HASH",
2918 Self::Broadcast => "DISPATCHER_TYPE_BROADCAST",
2919 Self::Simple => "DISPATCHER_TYPE_SIMPLE",
2920 Self::NoShuffle => "DISPATCHER_TYPE_NO_SHUFFLE",
2921 }
2922 }
2923 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2925 match value {
2926 "DISPATCHER_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2927 "DISPATCHER_TYPE_HASH" => Some(Self::Hash),
2928 "DISPATCHER_TYPE_BROADCAST" => Some(Self::Broadcast),
2929 "DISPATCHER_TYPE_SIMPLE" => Some(Self::Simple),
2930 "DISPATCHER_TYPE_NO_SHUFFLE" => Some(Self::NoShuffle),
2931 _ => None,
2932 }
2933 }
2934}