Skip to main content

risingwave_pb/
stream_plan.rs

1// This file is @generated by prost-build.
2#[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    /// New dispatchers for each actor.
22    #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::ActorId")]
23    pub actor_dispatchers: ::std::collections::HashMap<crate::id::ActorId, Dispatchers>,
24    /// All actors to be added (to the main connected component of the graph) in this update.
25    #[prost(uint32, repeated, tag = "3", wrapper = "crate::id::ActorId")]
26    pub added_actors: ::prost::alloc::vec::Vec<crate::id::ActorId>,
27    /// We may embed a source change split mutation here.
28    /// `Source` and `SourceBackfill` are handled together here.
29    /// TODO: we may allow multiple mutations in a single barrier.
30    #[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    /// We may embed a pause mutation here.
36    /// TODO: we may allow multiple mutations in a single barrier.
37    #[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    /// nodes which should be paused initially.
42    #[prost(uint32, repeated, tag = "6", wrapper = "crate::id::FragmentId")]
43    pub backfill_nodes_to_pause: ::prost::alloc::vec::Vec<crate::id::FragmentId>,
44    /// CDC table snapshot splits
45    #[prost(message, optional, tag = "7")]
46    pub actor_cdc_table_snapshot_splits: ::core::option::Option<
47        super::source::CdcTableSnapshotSplitsWithGeneration,
48    >,
49    /// Use downstream_fragment_id as keys.
50    #[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    /// All actors to be dropped while adding the new streaming job.
56    #[prost(uint32, repeated, tag = "9", wrapper = "crate::id::ActorId")]
57    pub dropped_actors: ::prost::alloc::vec::Vec<crate::id::ActorId>,
58    /// Sink ids whose log store must be drained before this add barrier can be collected.
59    #[prost(uint32, repeated, tag = "10", wrapper = "crate::id::SinkId")]
60    pub sink_log_store_flush: ::prost::alloc::vec::Vec<crate::id::SinkId>,
61}
62/// Nested message and enum types in `AddMutation`.
63pub 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    /// Only sink-fragment in the sink job is recorded.
79    #[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    /// Dispatcher updates.
86    #[prost(message, repeated, tag = "1")]
87    pub dispatcher_update: ::prost::alloc::vec::Vec<update_mutation::DispatcherUpdate>,
88    /// Merge updates.
89    #[prost(message, repeated, tag = "2")]
90    pub merge_update: ::prost::alloc::vec::Vec<update_mutation::MergeUpdate>,
91    /// Vnode bitmap updates for each actor.
92    #[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    /// All actors to be dropped in this update.
98    #[prost(uint32, repeated, tag = "4", wrapper = "crate::id::ActorId")]
99    pub dropped_actors: ::prost::alloc::vec::Vec<crate::id::ActorId>,
100    /// Source updates.
101    /// `Source` and `SourceBackfill` are handled together here.
102    #[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    /// When modifying the Materialized View, we need to recreate the Dispatcher from the old upstream to the new TableFragment.
108    /// Consistent with the semantics in AddMutation.
109    #[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    /// CDC table snapshot splits
115    #[prost(message, optional, tag = "7")]
116    pub actor_cdc_table_snapshot_splits: ::core::option::Option<
117        super::source::CdcTableSnapshotSplitsWithGeneration,
118    >,
119    #[prost(map = "uint32, message", tag = "8")]
120    pub sink_schema_change: ::std::collections::HashMap<u32, SinkSchemaChange>,
121    #[prost(message, repeated, tag = "9")]
122    pub subscriptions_to_drop: ::prost::alloc::vec::Vec<SubscriptionUpstreamInfo>,
123    #[prost(message, optional, tag = "10")]
124    pub iceberg_pk_index_compaction: ::core::option::Option<
125        IcebergPkIndexCompactionContext,
126    >,
127}
128/// Nested message and enum types in `UpdateMutation`.
129pub mod update_mutation {
130    #[derive(prost_helpers::AnyPB)]
131    #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
132    pub struct DispatcherUpdate {
133        /// Dispatcher can be uniquely identified by a combination of actor id and dispatcher id.
134        #[prost(uint32, tag = "1", wrapper = "crate::id::ActorId")]
135        pub actor_id: crate::id::ActorId,
136        #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
137        pub dispatcher_id: crate::id::FragmentId,
138        /// The hash mapping for consistent hash.
139        /// For dispatcher types other than HASH, this is ignored.
140        #[prost(message, optional, tag = "3")]
141        pub hash_mapping: ::core::option::Option<super::ActorMapping>,
142        /// Added downstream actors.
143        #[prost(uint32, repeated, tag = "4", wrapper = "crate::id::ActorId")]
144        pub added_downstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
145        /// Removed downstream actors.
146        #[prost(uint32, repeated, tag = "5", wrapper = "crate::id::ActorId")]
147        pub removed_downstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
148    }
149    #[derive(prost_helpers::AnyPB)]
150    #[derive(Clone, PartialEq, ::prost::Message)]
151    pub struct MergeUpdate {
152        /// Merge executor can be uniquely identified by a combination of actor id and upstream fragment id.
153        #[prost(uint32, tag = "1", wrapper = "crate::id::ActorId")]
154        pub actor_id: crate::id::ActorId,
155        #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
156        pub upstream_fragment_id: crate::id::FragmentId,
157        /// * For scaling, this is always `None`.
158        /// * For plan change, the upstream fragment will be changed to a new one, and this will be `Some`.
159        ///   In this case, all the upstream actors should be removed and replaced by the `new` ones.
160        #[prost(uint32, optional, tag = "5", wrapper = "crate::id::FragmentId")]
161        pub new_upstream_fragment_id: ::core::option::Option<crate::id::FragmentId>,
162        /// Added upstream actors.
163        #[prost(message, repeated, tag = "3")]
164        pub added_upstream_actors: ::prost::alloc::vec::Vec<
165            super::super::common::ActorInfo,
166        >,
167        /// Removed upstream actors.
168        /// Note: this is empty for replace job.
169        #[prost(uint32, repeated, tag = "4", wrapper = "crate::id::ActorId")]
170        pub removed_upstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
171    }
172}
173#[derive(prost_helpers::AnyPB)]
174#[derive(Clone, PartialEq, ::prost::Message)]
175pub struct SourceChangeSplitMutation {
176    /// `Source` and `SourceBackfill` are handled together here.
177    #[prost(map = "uint32, message", tag = "2", wrapper = "crate::id::ActorId")]
178    pub actor_splits: ::std::collections::HashMap<
179        crate::id::ActorId,
180        super::source::ConnectorSplits,
181    >,
182}
183#[derive(prost_helpers::AnyPB)]
184#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
185pub struct PauseMutation {}
186#[derive(prost_helpers::AnyPB)]
187#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
188pub struct ResumeMutation {}
189#[derive(prost_helpers::AnyPB)]
190#[derive(Clone, PartialEq, ::prost::Message)]
191pub struct ThrottleMutation {
192    #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::FragmentId")]
193    pub fragment_throttle: ::std::collections::HashMap<
194        crate::id::FragmentId,
195        throttle_mutation::ThrottleConfig,
196    >,
197}
198/// Nested message and enum types in `ThrottleMutation`.
199pub mod throttle_mutation {
200    #[derive(prost_helpers::AnyPB)]
201    #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
202    pub struct ThrottleConfig {
203        #[prost(uint32, optional, tag = "1")]
204        pub rate_limit: ::core::option::Option<u32>,
205        #[prost(enumeration = "super::super::common::ThrottleType", tag = "2")]
206        pub throttle_type: i32,
207    }
208}
209#[derive(prost_helpers::AnyPB)]
210#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
211pub struct SubscriptionUpstreamInfo {
212    /// can either be subscription_id or table_id of creating TableFragments
213    #[prost(uint32, tag = "1", wrapper = "crate::id::SubscriberId")]
214    pub subscriber_id: crate::id::SubscriberId,
215    #[prost(uint32, tag = "2", wrapper = "crate::id::TableId")]
216    pub upstream_mv_table_id: crate::id::TableId,
217}
218#[derive(prost_helpers::AnyPB)]
219#[derive(Clone, PartialEq, ::prost::Message)]
220pub struct DropSubscriptionsMutation {
221    #[prost(message, repeated, tag = "1")]
222    pub info: ::prost::alloc::vec::Vec<SubscriptionUpstreamInfo>,
223}
224#[derive(prost_helpers::AnyPB)]
225#[derive(Clone, PartialEq, ::prost::Message)]
226pub struct ConnectorPropsChangeMutation {
227    #[prost(map = "uint32, message", tag = "1")]
228    pub connector_props_infos: ::std::collections::HashMap<
229        u32,
230        connector_props_change_mutation::ConnectorPropsInfo,
231    >,
232}
233/// Nested message and enum types in `ConnectorPropsChangeMutation`.
234pub mod connector_props_change_mutation {
235    #[derive(prost_helpers::AnyPB)]
236    #[derive(Clone, PartialEq, ::prost::Message)]
237    pub struct ConnectorPropsInfo {
238        #[prost(map = "string, string", tag = "1")]
239        pub connector_props_info: ::std::collections::HashMap<
240            ::prost::alloc::string::String,
241            ::prost::alloc::string::String,
242        >,
243    }
244}
245#[derive(prost_helpers::AnyPB)]
246#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
247pub struct StartFragmentBackfillMutation {
248    #[prost(uint32, repeated, tag = "1", wrapper = "crate::id::FragmentId")]
249    pub fragment_ids: ::prost::alloc::vec::Vec<crate::id::FragmentId>,
250}
251#[derive(prost_helpers::AnyPB)]
252#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
253pub struct RefreshStartMutation {
254    /// Table ID to start refresh operation.
255    #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
256    pub table_id: crate::id::TableId,
257    /// Associated source ID for this refresh operation.
258    #[prost(uint32, tag = "2", wrapper = "crate::id::SourceId")]
259    pub associated_source_id: crate::id::SourceId,
260}
261#[derive(prost_helpers::AnyPB)]
262#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
263pub struct ListFinishMutation {
264    /// Associated source ID for this list operation.
265    #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
266    pub associated_source_id: crate::id::SourceId,
267}
268#[derive(prost_helpers::AnyPB)]
269#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
270pub struct LoadFinishMutation {
271    /// Associated source ID for this load operation.
272    #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
273    pub associated_source_id: crate::id::SourceId,
274}
275#[derive(prost_helpers::AnyPB)]
276#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
277pub struct ResetSourceMutation {
278    /// Source ID to reset
279    #[prost(uint32, tag = "1")]
280    pub source_id: u32,
281}
282/// Inject specific offsets into source splits (UNSAFE - admin only).
283/// This overwrites the stored offsets and can cause data duplication or loss.
284#[derive(prost_helpers::AnyPB)]
285#[derive(Clone, PartialEq, ::prost::Message)]
286pub struct InjectSourceOffsetsMutation {
287    /// Source ID to inject offsets for
288    #[prost(uint32, tag = "1")]
289    pub source_id: u32,
290    /// Split ID -> offset (JSON-encoded based on connector type)
291    #[prost(map = "string, string", tag = "2")]
292    pub split_offsets: ::std::collections::HashMap<
293        ::prost::alloc::string::String,
294        ::prost::alloc::string::String,
295    >,
296}
297/// Context for an ongoing iceberg pk-index compaction.
298#[derive(prost_helpers::AnyPB)]
299#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
300pub struct IcebergPkIndexCompactionContext {
301    #[prost(uint32, tag = "1", wrapper = "crate::id::SinkId")]
302    pub sink_id: crate::id::SinkId,
303    #[prost(uint64, tag = "2", wrapper = "crate::id::IcebergCompactionTaskId")]
304    pub task_id: crate::id::IcebergCompactionTaskId,
305    #[prost(enumeration = "iceberg_pk_index_compaction_context::Phase", tag = "3")]
306    pub phase: i32,
307    /// Required for PHASE_BEGIN and absent for PHASE_END.
308    #[prost(message, optional, tag = "4")]
309    pub resolver_task_input: ::core::option::Option<
310        iceberg_pk_index_compaction_context::ResolverTaskInput,
311    >,
312}
313/// Nested message and enum types in `IcebergPkIndexCompactionContext`.
314pub mod iceberg_pk_index_compaction_context {
315    #[derive(prost_helpers::AnyPB)]
316    #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
317    pub struct ResolverTaskInput {
318        #[prost(string, repeated, tag = "1")]
319        pub output_data_file_paths: ::prost::alloc::vec::Vec<
320            ::prost::alloc::string::String,
321        >,
322        #[prost(string, repeated, tag = "2")]
323        pub input_data_file_paths: ::prost::alloc::vec::Vec<
324            ::prost::alloc::string::String,
325        >,
326        #[prost(int64, tag = "3")]
327        pub read_snapshot_id: i64,
328    }
329    #[derive(prost_helpers::AnyPB)]
330    #[derive(
331        Clone,
332        Copy,
333        Debug,
334        PartialEq,
335        Eq,
336        Hash,
337        PartialOrd,
338        Ord,
339        ::prost::Enumeration
340    )]
341    #[repr(i32)]
342    pub enum Phase {
343        Unspecified = 0,
344        Begin = 1,
345        End = 2,
346    }
347    impl Phase {
348        /// String value of the enum field names used in the ProtoBuf definition.
349        ///
350        /// The values are not transformed in any way and thus are considered stable
351        /// (if the ProtoBuf definition does not change) and safe for programmatic use.
352        pub fn as_str_name(&self) -> &'static str {
353            match self {
354                Self::Unspecified => "PHASE_UNSPECIFIED",
355                Self::Begin => "PHASE_BEGIN",
356                Self::End => "PHASE_END",
357            }
358        }
359        /// Creates an enum from field names used in the ProtoBuf definition.
360        pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
361            match value {
362                "PHASE_UNSPECIFIED" => Some(Self::Unspecified),
363                "PHASE_BEGIN" => Some(Self::Begin),
364                "PHASE_END" => Some(Self::End),
365                _ => None,
366            }
367        }
368    }
369}
370#[derive(prost_helpers::AnyPB)]
371#[derive(Clone, PartialEq, ::prost::Message)]
372pub struct BarrierMutation {
373    #[prost(
374        oneof = "barrier_mutation::Mutation",
375        tags = "3, 4, 5, 6, 7, 8, 10, 12, 13, 14, 15, 16, 17, 18, 19"
376    )]
377    pub mutation: ::core::option::Option<barrier_mutation::Mutation>,
378}
379/// Nested message and enum types in `BarrierMutation`.
380pub mod barrier_mutation {
381    #[derive(prost_helpers::AnyPB)]
382    #[derive(Clone, PartialEq, ::prost::Oneof)]
383    pub enum Mutation {
384        /// Add new dispatchers to some actors, used for creating materialized views.
385        #[prost(message, tag = "3")]
386        Add(super::AddMutation),
387        /// Stop a set of actors, used for dropping materialized views. Empty dispatchers will be
388        /// automatically removed.
389        #[prost(message, tag = "4")]
390        Stop(super::StopMutation),
391        /// Update outputs and hash mappings for some dispatchers, used for scaling and replace table.
392        #[prost(message, tag = "5")]
393        Update(super::UpdateMutation),
394        /// Change the split of some sources.
395        #[prost(message, tag = "6")]
396        Splits(super::SourceChangeSplitMutation),
397        /// Pause the dataflow of the whole streaming graph, only used for scaling.
398        #[prost(message, tag = "7")]
399        Pause(super::PauseMutation),
400        /// Resume the dataflow of the whole streaming graph, only used for scaling.
401        #[prost(message, tag = "8")]
402        Resume(super::ResumeMutation),
403        /// Alter the RateLimit of some specific executors
404        #[prost(message, tag = "10")]
405        Throttle(super::ThrottleMutation),
406        /// Drop subscription on mv
407        #[prost(message, tag = "12")]
408        DropSubscriptions(super::DropSubscriptionsMutation),
409        /// Alter sink/connector/source props
410        #[prost(message, tag = "13")]
411        ConnectorPropsChange(super::ConnectorPropsChangeMutation),
412        /// Start backfilling for specific fragments
413        /// This is separated from `ThrottleMutation::NO_RATE_LIMIT`.
414        /// This is because user may concurrently rate limit + use backfill order control.
415        /// If we use rate limit to pause / resume backfill fragments, if user manually
416        /// resumes some fragments, this will overwrite the backfill order configuration.
417        #[prost(message, tag = "14")]
418        StartFragmentBackfill(super::StartFragmentBackfillMutation),
419        /// Start refresh signal for refreshing tables
420        #[prost(message, tag = "15")]
421        RefreshStart(super::RefreshStartMutation),
422        /// Load finish signal for refreshing tables
423        #[prost(message, tag = "16")]
424        LoadFinish(super::LoadFinishMutation),
425        /// List finish signal for refreshing tables
426        #[prost(message, tag = "17")]
427        ListFinish(super::ListFinishMutation),
428        /// Reset CDC source offset to latest
429        #[prost(message, tag = "18")]
430        ResetSource(super::ResetSourceMutation),
431        /// Inject specific offsets into source splits (UNSAFE)
432        #[prost(message, tag = "19")]
433        InjectSourceOffsets(super::InjectSourceOffsetsMutation),
434    }
435}
436#[derive(prost_helpers::AnyPB)]
437#[derive(Clone, PartialEq, ::prost::Message)]
438pub struct Barrier {
439    #[prost(message, optional, tag = "1")]
440    pub epoch: ::core::option::Option<super::data::Epoch>,
441    #[prost(message, optional, tag = "3")]
442    pub mutation: ::core::option::Option<BarrierMutation>,
443    /// Used for tracing.
444    #[prost(map = "string, string", tag = "2")]
445    pub tracing_context: ::std::collections::HashMap<
446        ::prost::alloc::string::String,
447        ::prost::alloc::string::String,
448    >,
449    /// The kind of the barrier.
450    #[prost(enumeration = "barrier::BarrierKind", tag = "9")]
451    pub kind: i32,
452}
453/// Nested message and enum types in `Barrier`.
454pub mod barrier {
455    #[derive(prost_helpers::AnyPB)]
456    #[derive(::enum_as_inner::EnumAsInner)]
457    #[derive(
458        Clone,
459        Copy,
460        Debug,
461        PartialEq,
462        Eq,
463        Hash,
464        PartialOrd,
465        Ord,
466        ::prost::Enumeration
467    )]
468    #[repr(i32)]
469    pub enum BarrierKind {
470        Unspecified = 0,
471        /// The first barrier after a fresh start or recovery.
472        /// There will be no data associated with the previous epoch of the barrier.
473        Initial = 1,
474        /// A normal barrier. Data should be flushed locally.
475        Barrier = 2,
476        /// A checkpoint barrier. Data should be synchorized to the shared storage.
477        Checkpoint = 3,
478    }
479    impl BarrierKind {
480        /// String value of the enum field names used in the ProtoBuf definition.
481        ///
482        /// The values are not transformed in any way and thus are considered stable
483        /// (if the ProtoBuf definition does not change) and safe for programmatic use.
484        pub fn as_str_name(&self) -> &'static str {
485            match self {
486                Self::Unspecified => "BARRIER_KIND_UNSPECIFIED",
487                Self::Initial => "BARRIER_KIND_INITIAL",
488                Self::Barrier => "BARRIER_KIND_BARRIER",
489                Self::Checkpoint => "BARRIER_KIND_CHECKPOINT",
490            }
491        }
492        /// Creates an enum from field names used in the ProtoBuf definition.
493        pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
494            match value {
495                "BARRIER_KIND_UNSPECIFIED" => Some(Self::Unspecified),
496                "BARRIER_KIND_INITIAL" => Some(Self::Initial),
497                "BARRIER_KIND_BARRIER" => Some(Self::Barrier),
498                "BARRIER_KIND_CHECKPOINT" => Some(Self::Checkpoint),
499                _ => None,
500            }
501        }
502    }
503}
504#[derive(prost_helpers::AnyPB)]
505#[derive(Clone, PartialEq, ::prost::Message)]
506pub struct Watermark {
507    /// The reference to the watermark column in the stream's schema.
508    #[prost(message, optional, tag = "1")]
509    pub column: ::core::option::Option<super::expr::InputRef>,
510    /// The watermark value, there will be no record having a greater value in the watermark column.
511    #[prost(message, optional, tag = "3")]
512    pub val: ::core::option::Option<super::data::Datum>,
513}
514#[derive(prost_helpers::AnyPB)]
515#[derive(Clone, PartialEq, ::prost::Message)]
516pub struct StreamMessage {
517    #[prost(oneof = "stream_message::StreamMessage", tags = "1, 2, 3")]
518    pub stream_message: ::core::option::Option<stream_message::StreamMessage>,
519}
520/// Nested message and enum types in `StreamMessage`.
521pub mod stream_message {
522    #[derive(prost_helpers::AnyPB)]
523    #[derive(Clone, PartialEq, ::prost::Oneof)]
524    pub enum StreamMessage {
525        #[prost(message, tag = "1")]
526        StreamChunk(super::super::data::StreamChunk),
527        #[prost(message, tag = "2")]
528        Barrier(super::Barrier),
529        #[prost(message, tag = "3")]
530        Watermark(super::Watermark),
531    }
532}
533#[derive(prost_helpers::AnyPB)]
534#[derive(Clone, PartialEq, ::prost::Message)]
535pub struct StreamMessageBatch {
536    #[prost(oneof = "stream_message_batch::StreamMessageBatch", tags = "1, 2, 3")]
537    pub stream_message_batch: ::core::option::Option<
538        stream_message_batch::StreamMessageBatch,
539    >,
540}
541/// Nested message and enum types in `StreamMessageBatch`.
542pub mod stream_message_batch {
543    #[derive(prost_helpers::AnyPB)]
544    #[derive(Clone, PartialEq, ::prost::Message)]
545    pub struct BarrierBatch {
546        #[prost(message, repeated, tag = "1")]
547        pub barriers: ::prost::alloc::vec::Vec<super::Barrier>,
548    }
549    #[derive(prost_helpers::AnyPB)]
550    #[derive(Clone, PartialEq, ::prost::Oneof)]
551    pub enum StreamMessageBatch {
552        #[prost(message, tag = "1")]
553        StreamChunk(super::super::data::StreamChunk),
554        #[prost(message, tag = "2")]
555        BarrierBatch(BarrierBatch),
556        #[prost(message, tag = "3")]
557        Watermark(super::Watermark),
558    }
559}
560/// Hash mapping for compute node. Stores mapping from virtual node to actor id.
561#[derive(prost_helpers::AnyPB)]
562#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
563pub struct ActorMapping {
564    #[prost(uint32, repeated, tag = "1")]
565    pub original_indices: ::prost::alloc::vec::Vec<u32>,
566    #[prost(uint32, repeated, tag = "2", wrapper = "crate::id::ActorId")]
567    pub data: ::prost::alloc::vec::Vec<crate::id::ActorId>,
568}
569#[derive(prost_helpers::AnyPB)]
570#[derive(Clone, PartialEq, ::prost::Message)]
571pub struct Columns {
572    #[prost(message, repeated, tag = "1")]
573    pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
574}
575#[derive(prost_helpers::AnyPB)]
576#[derive(Clone, PartialEq, ::prost::Message)]
577pub struct StreamSource {
578    #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
579    pub source_id: crate::id::SourceId,
580    #[prost(message, optional, tag = "2")]
581    pub state_table: ::core::option::Option<super::catalog::Table>,
582    #[prost(uint32, optional, tag = "3")]
583    pub row_id_index: ::core::option::Option<u32>,
584    #[prost(message, repeated, tag = "4")]
585    pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
586    #[prost(btree_map = "string, string", tag = "6")]
587    pub with_properties: ::prost::alloc::collections::BTreeMap<
588        ::prost::alloc::string::String,
589        ::prost::alloc::string::String,
590    >,
591    #[prost(message, optional, tag = "7")]
592    pub info: ::core::option::Option<super::catalog::StreamSourceInfo>,
593    #[prost(string, tag = "8")]
594    pub source_name: ::prost::alloc::string::String,
595    /// Source rate limit
596    #[prost(uint32, optional, tag = "9")]
597    pub rate_limit: ::core::option::Option<u32>,
598    #[prost(btree_map = "string, message", tag = "10")]
599    pub secret_refs: ::prost::alloc::collections::BTreeMap<
600        ::prost::alloc::string::String,
601        super::secret::SecretRef,
602    >,
603    /// Downstream columns are used by list node to know which columns are needed.
604    #[prost(message, optional, tag = "11")]
605    pub downstream_columns: ::core::option::Option<Columns>,
606    #[prost(message, optional, tag = "12")]
607    pub refresh_mode: ::core::option::Option<super::plan_common::SourceRefreshMode>,
608    #[prost(uint32, optional, tag = "13", wrapper = "crate::id::TableId")]
609    pub associated_table_id: ::core::option::Option<crate::id::TableId>,
610}
611/// copy contents from StreamSource to prevent compatibility issues in the future
612#[derive(prost_helpers::AnyPB)]
613#[derive(Clone, PartialEq, ::prost::Message)]
614pub struct StreamFsFetch {
615    #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
616    pub source_id: crate::id::SourceId,
617    #[prost(message, optional, tag = "2")]
618    pub state_table: ::core::option::Option<super::catalog::Table>,
619    #[prost(uint32, optional, tag = "3")]
620    pub row_id_index: ::core::option::Option<u32>,
621    #[prost(message, repeated, tag = "4")]
622    pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
623    #[prost(btree_map = "string, string", tag = "6")]
624    pub with_properties: ::prost::alloc::collections::BTreeMap<
625        ::prost::alloc::string::String,
626        ::prost::alloc::string::String,
627    >,
628    #[prost(message, optional, tag = "7")]
629    pub info: ::core::option::Option<super::catalog::StreamSourceInfo>,
630    #[prost(string, tag = "8")]
631    pub source_name: ::prost::alloc::string::String,
632    /// Source rate limit
633    #[prost(uint32, optional, tag = "9")]
634    pub rate_limit: ::core::option::Option<u32>,
635    #[prost(btree_map = "string, message", tag = "10")]
636    pub secret_refs: ::prost::alloc::collections::BTreeMap<
637        ::prost::alloc::string::String,
638        super::secret::SecretRef,
639    >,
640    #[prost(message, optional, tag = "11")]
641    pub refresh_mode: ::core::option::Option<super::plan_common::SourceRefreshMode>,
642    #[prost(uint32, optional, tag = "12", wrapper = "crate::id::TableId")]
643    pub associated_table_id: ::core::option::Option<crate::id::TableId>,
644}
645/// The executor only for receiving barrier from the meta service. It always resides in the leaves
646/// of the streaming graph.
647#[derive(prost_helpers::AnyPB)]
648#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
649pub struct BarrierRecvNode {}
650#[derive(prost_helpers::AnyPB)]
651#[derive(Clone, PartialEq, ::prost::Message)]
652pub struct SourceNode {
653    /// The source node can contain either a stream source or nothing. So here we extract all
654    /// information about stream source to a message, and here it will be an `Option` in Rust.
655    #[prost(message, optional, tag = "1")]
656    pub source_inner: ::core::option::Option<StreamSource>,
657}
658#[derive(prost_helpers::AnyPB)]
659#[derive(Clone, PartialEq, ::prost::Message)]
660pub struct StreamFsFetchNode {
661    #[prost(message, optional, tag = "1")]
662    pub node_inner: ::core::option::Option<StreamFsFetch>,
663}
664/// / It's input must be a `MergeNode`, which connects to the upstream source job.
665/// / See `StreamSourceScan::adhoc_to_stream_prost` for the plan.
666#[derive(prost_helpers::AnyPB)]
667#[derive(Clone, PartialEq, ::prost::Message)]
668pub struct SourceBackfillNode {
669    #[prost(uint32, tag = "1", wrapper = "crate::id::SourceId")]
670    pub upstream_source_id: crate::id::SourceId,
671    #[prost(uint32, optional, tag = "2")]
672    pub row_id_index: ::core::option::Option<u32>,
673    #[prost(message, repeated, tag = "3")]
674    pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
675    #[prost(message, optional, tag = "4")]
676    pub info: ::core::option::Option<super::catalog::StreamSourceInfo>,
677    #[prost(string, tag = "5")]
678    pub source_name: ::prost::alloc::string::String,
679    #[prost(btree_map = "string, string", tag = "6")]
680    pub with_properties: ::prost::alloc::collections::BTreeMap<
681        ::prost::alloc::string::String,
682        ::prost::alloc::string::String,
683    >,
684    /// Backfill rate limit
685    #[prost(uint32, optional, tag = "7")]
686    pub rate_limit: ::core::option::Option<u32>,
687    /// `| partition_id | backfill_progress |`
688    #[prost(message, optional, tag = "8")]
689    pub state_table: ::core::option::Option<super::catalog::Table>,
690    #[prost(btree_map = "string, message", tag = "9")]
691    pub secret_refs: ::prost::alloc::collections::BTreeMap<
692        ::prost::alloc::string::String,
693        super::secret::SecretRef,
694    >,
695}
696#[derive(prost_helpers::AnyPB)]
697#[derive(Clone, PartialEq, ::prost::Message)]
698pub struct SinkDesc {
699    #[prost(uint32, tag = "1", wrapper = "crate::id::SinkId")]
700    pub id: crate::id::SinkId,
701    #[prost(string, tag = "2")]
702    pub name: ::prost::alloc::string::String,
703    #[prost(string, tag = "3")]
704    pub definition: ::prost::alloc::string::String,
705    #[prost(message, repeated, tag = "5")]
706    pub plan_pk: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
707    #[prost(uint32, repeated, tag = "6")]
708    pub downstream_pk: ::prost::alloc::vec::Vec<u32>,
709    #[prost(uint32, repeated, tag = "7")]
710    pub distribution_key: ::prost::alloc::vec::Vec<u32>,
711    #[prost(btree_map = "string, string", tag = "8")]
712    pub properties: ::prost::alloc::collections::BTreeMap<
713        ::prost::alloc::string::String,
714        ::prost::alloc::string::String,
715    >,
716    /// to be deprecated
717    #[prost(enumeration = "super::catalog::SinkType", tag = "9")]
718    pub sink_type: i32,
719    #[prost(message, repeated, tag = "10")]
720    pub column_catalogs: ::prost::alloc::vec::Vec<super::plan_common::ColumnCatalog>,
721    #[prost(string, tag = "11")]
722    pub db_name: ::prost::alloc::string::String,
723    /// If the sink is from table or mv, this is name of the table/mv. Otherwise
724    /// it is the name of the sink itself.
725    #[prost(string, tag = "12")]
726    pub sink_from_name: ::prost::alloc::string::String,
727    #[prost(message, optional, tag = "13")]
728    pub format_desc: ::core::option::Option<super::catalog::SinkFormatDesc>,
729    #[prost(uint32, optional, tag = "14")]
730    pub target_table: ::core::option::Option<u32>,
731    #[prost(uint64, optional, tag = "15")]
732    pub extra_partition_col_idx: ::core::option::Option<u64>,
733    #[prost(btree_map = "string, message", tag = "16")]
734    pub secret_refs: ::prost::alloc::collections::BTreeMap<
735        ::prost::alloc::string::String,
736        super::secret::SecretRef,
737    >,
738    /// Backward-compatible replacement for `SINK_TYPE_FORCE_APPEND_ONLY`.
739    ///
740    /// Should not directly access this field. Use method `ignore_delete()` instead.
741    #[prost(bool, tag = "17")]
742    pub raw_ignore_delete: bool,
743}
744#[derive(prost_helpers::AnyPB)]
745#[derive(Clone, PartialEq, ::prost::Message)]
746pub struct SinkNode {
747    #[prost(message, optional, tag = "1")]
748    pub sink_desc: ::core::option::Option<SinkDesc>,
749    /// A sink with a kv log store should have a table.
750    #[prost(message, optional, tag = "2")]
751    pub table: ::core::option::Option<super::catalog::Table>,
752    #[prost(enumeration = "SinkLogStoreType", tag = "3")]
753    pub log_store_type: i32,
754    #[prost(uint32, optional, tag = "4")]
755    pub rate_limit: ::core::option::Option<u32>,
756}
757#[derive(prost_helpers::AnyPB)]
758#[derive(Clone, PartialEq, ::prost::Message)]
759pub struct IcebergWithPkIndexWriterNode {
760    #[prost(message, optional, tag = "1")]
761    pub sink_desc: ::core::option::Option<SinkDesc>,
762    #[prost(message, optional, tag = "2")]
763    pub pk_index_table: ::core::option::Option<super::catalog::Table>,
764}
765#[derive(prost_helpers::AnyPB)]
766#[derive(Clone, PartialEq, ::prost::Message)]
767pub struct IcebergWithPkIndexPositionDeleteMergerNode {
768    #[prost(message, optional, tag = "1")]
769    pub sink_desc: ::core::option::Option<SinkDesc>,
770}
771/// Leaf of the inactive pk-index compaction pipeline rendered with the sink job.
772///
773/// This message existed in release-3.1, but no producer emitted it into a stream graph, so the old
774/// fields were never persisted as part of an actor definition.
775#[derive(prost_helpers::AnyPB)]
776#[derive(Clone, PartialEq, ::prost::Message)]
777pub struct CompactionResolverNode {
778    /// Placeholder from frontend; replaced with the actual job ID by Meta.
779    #[prost(uint32, tag = "1", wrapper = "crate::id::SinkId")]
780    pub sink_id: crate::id::SinkId,
781    #[prost(btree_map = "string, string", tag = "2")]
782    pub properties: ::prost::alloc::collections::BTreeMap<
783        ::prost::alloc::string::String,
784        ::prost::alloc::string::String,
785    >,
786    #[prost(btree_map = "string, message", tag = "3")]
787    pub secret_refs: ::prost::alloc::collections::BTreeMap<
788        ::prost::alloc::string::String,
789        super::secret::SecretRef,
790    >,
791    /// The order defines the primary-key order in resolver output rows.
792    #[prost(message, repeated, tag = "4")]
793    pub pk_columns: ::prost::alloc::vec::Vec<compaction_resolver_node::PkColumn>,
794}
795/// Nested message and enum types in `CompactionResolverNode`.
796pub mod compaction_resolver_node {
797    #[derive(prost_helpers::AnyPB)]
798    #[derive(Clone, PartialEq, ::prost::Message)]
799    pub struct PkColumn {
800        /// Column position in the Iceberg data-file row.
801        #[prost(uint32, tag = "1")]
802        pub data_file_index: u32,
803        /// Type and identity of the corresponding primary-key column.
804        #[prost(message, optional, tag = "2")]
805        pub column_desc: ::core::option::Option<super::super::plan_common::ColumnDesc>,
806    }
807}
808#[derive(prost_helpers::AnyPB)]
809#[derive(Clone, PartialEq, ::prost::Message)]
810pub struct ProjectNode {
811    #[prost(message, repeated, tag = "1")]
812    pub select_list: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
813    /// this two field is expressing a list of usize pair, which means when project receives a
814    /// watermark with `watermark_input_cols\[i\]` column index, it should derive a new watermark
815    /// with `watermark_output_cols\[i\]`th expression
816    #[prost(uint32, repeated, tag = "2")]
817    pub watermark_input_cols: ::prost::alloc::vec::Vec<u32>,
818    #[prost(uint32, repeated, tag = "3")]
819    pub watermark_output_cols: ::prost::alloc::vec::Vec<u32>,
820    #[prost(uint32, repeated, tag = "4")]
821    pub nondecreasing_exprs: ::prost::alloc::vec::Vec<u32>,
822    /// Whether there are likely no-op updates in the output chunks, so that eliminating them with
823    /// `StreamChunk::eliminate_adjacent_noop_update` could be beneficial.
824    #[prost(bool, tag = "5")]
825    pub noop_update_hint: bool,
826}
827#[derive(prost_helpers::AnyPB)]
828#[derive(Clone, PartialEq, ::prost::Message)]
829pub struct FilterNode {
830    #[prost(message, optional, tag = "1")]
831    pub search_condition: ::core::option::Option<super::expr::ExprNode>,
832}
833#[derive(prost_helpers::AnyPB)]
834#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
835pub struct ChangeLogNode {
836    /// Whether or not there is an op in the final output.
837    #[prost(bool, tag = "1")]
838    pub need_op: bool,
839    #[prost(uint32, repeated, tag = "2")]
840    pub distribution_keys: ::prost::alloc::vec::Vec<u32>,
841    #[prost(uint32, repeated, tag = "4")]
842    pub stream_keys: ::prost::alloc::vec::Vec<u32>,
843    #[prost(oneof = "change_log_node::Mode", tags = "3")]
844    pub mode: ::core::option::Option<change_log_node::Mode>,
845}
846/// Nested message and enum types in `ChangeLogNode`.
847pub mod change_log_node {
848    #[derive(prost_helpers::AnyPB)]
849    #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
850    pub struct Keyed {}
851    #[derive(prost_helpers::AnyPB)]
852    #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Oneof)]
853    pub enum Mode {
854        #[prost(message, tag = "3")]
855        Keyed(Keyed),
856    }
857}
858#[derive(prost_helpers::AnyPB)]
859#[derive(Clone, PartialEq, ::prost::Message)]
860pub struct CdcFilterNode {
861    #[prost(message, optional, tag = "1")]
862    pub search_condition: ::core::option::Option<super::expr::ExprNode>,
863    #[prost(uint32, tag = "2", wrapper = "crate::id::SourceId")]
864    pub upstream_source_id: crate::id::SourceId,
865}
866/// A materialized view is regarded as a table.
867/// In addition, we also specify primary key to MV for efficient point lookup during update and deletion.
868///
869/// The node will be used for both create mv and create index.
870///
871/// * When creating mv, `pk == distribution_key == column_orders`.
872/// * When creating index, `column_orders` will contain both
873///   arrange columns and pk columns, while distribution key will be arrange columns.
874#[derive(prost_helpers::AnyPB)]
875#[derive(Clone, PartialEq, ::prost::Message)]
876pub struct MaterializeNode {
877    #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
878    pub table_id: crate::id::TableId,
879    /// Column indexes and orders of primary key.
880    #[prost(message, repeated, tag = "2")]
881    pub column_orders: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
882    /// Primary materialized table that stores the final result data.
883    /// Purpose: This is the main table that users query against. It contains the
884    /// complete materialized view results with all columns from the SELECT clause.
885    /// Schema: Matches the output schema of the materialized view query
886    /// PK: Determined by the stream key of the materialized view
887    /// Distribution: Hash distributed by primary key columns for parallel processing
888    #[prost(message, optional, tag = "3")]
889    pub table: ::core::option::Option<super::catalog::Table>,
890    /// Staging table for refreshable materialized views during refresh operations.
891    /// Purpose: Temporary storage for collecting new/updated data during refresh.
892    /// This allows atomic replacement of old data with refreshed data.
893    /// Schema: Contains only the primary key columns from the main table (pk-only table)
894    /// - All PK columns with same data types as main table
895    /// - Same primary key definition as main table
896    /// Usage: Active only during refresh operations, empty otherwise
897    /// Lifecycle: Created -> populated during refresh -> merged with main table -> cleared
898    #[prost(message, optional, tag = "5")]
899    pub staging_table: ::core::option::Option<super::catalog::Table>,
900    /// Progress tracking table for refreshable materialized views.
901    /// Purpose: Tracks refresh operation progress per VirtualNode to enable fault-tolerant
902    /// resumable refresh operations. Stores checkpoint information for recovery.
903    /// Schema: Simplified variable-length schema (following backfill pattern):
904    /// - vnode (i32): VirtualNode identifier (PK)
905    /// - current_pos...: Current processing position (variable PK fields from upstream)
906    /// - is_completed (bool): Whether this vnode has completed processing
907    /// - processed_rows (i64): Number of rows processed so far in this vnode
908    /// PK: vnode (allows efficient per-vnode progress lookup and updates)
909    /// Distribution: Hash distributed by vnode for parallel progress tracking
910    /// Usage: Persists across refresh operations for resumability
911    /// Note: Stage info is now tracked in MaterializeExecutor memory for simplicity
912    #[prost(message, optional, tag = "6")]
913    pub refresh_progress_table: ::core::option::Option<super::catalog::Table>,
914    /// Whether the table can clean itself by TTL watermark, i.e., is defined with `WATERMARK ... WITH TTL`.
915    #[prost(bool, tag = "7")]
916    pub cleaned_by_ttl_watermark: bool,
917}
918#[derive(prost_helpers::AnyPB)]
919#[derive(Clone, PartialEq, ::prost::Message)]
920pub struct AggCallState {
921    #[prost(oneof = "agg_call_state::Inner", tags = "1, 3")]
922    pub inner: ::core::option::Option<agg_call_state::Inner>,
923}
924/// Nested message and enum types in `AggCallState`.
925pub mod agg_call_state {
926    /// the state is stored in the intermediate state table. used for count/sum/append-only extreme.
927    #[derive(prost_helpers::AnyPB)]
928    #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
929    pub struct ValueState {}
930    /// use the some column of the Upstream's materialization as the AggCall's state, used for extreme/string_agg/array_agg.
931    #[derive(prost_helpers::AnyPB)]
932    #[derive(Clone, PartialEq, ::prost::Message)]
933    pub struct MaterializedInputState {
934        #[prost(message, optional, tag = "1")]
935        pub table: ::core::option::Option<super::super::catalog::Table>,
936        /// for constructing state table column mapping
937        #[prost(uint32, repeated, tag = "2")]
938        pub included_upstream_indices: ::prost::alloc::vec::Vec<u32>,
939        #[prost(uint32, repeated, tag = "3")]
940        pub table_value_indices: ::prost::alloc::vec::Vec<u32>,
941        #[prost(message, repeated, tag = "4")]
942        pub order_columns: ::prost::alloc::vec::Vec<super::super::common::ColumnOrder>,
943    }
944    #[derive(prost_helpers::AnyPB)]
945    #[derive(Clone, PartialEq, ::prost::Oneof)]
946    pub enum Inner {
947        #[prost(message, tag = "1")]
948        ValueState(ValueState),
949        #[prost(message, tag = "3")]
950        MaterializedInputState(MaterializedInputState),
951    }
952}
953#[derive(prost_helpers::AnyPB)]
954#[derive(Clone, PartialEq, ::prost::Message)]
955pub struct SimpleAggNode {
956    #[prost(message, repeated, tag = "1")]
957    pub agg_calls: ::prost::alloc::vec::Vec<super::expr::AggCall>,
958    #[prost(message, repeated, tag = "3")]
959    pub agg_call_states: ::prost::alloc::vec::Vec<AggCallState>,
960    #[prost(message, optional, tag = "4")]
961    pub intermediate_state_table: ::core::option::Option<super::catalog::Table>,
962    /// Whether to optimize for append only stream.
963    /// It is true when the input is append-only
964    #[prost(bool, tag = "5")]
965    pub is_append_only: bool,
966    #[prost(map = "uint32, message", tag = "6")]
967    pub distinct_dedup_tables: ::std::collections::HashMap<u32, super::catalog::Table>,
968    #[prost(uint32, tag = "7")]
969    pub row_count_index: u32,
970    #[prost(enumeration = "AggNodeVersion", tag = "8")]
971    pub version: i32,
972    /// Required by the downstream `RowMergeNode`,
973    /// currently only used by the `approx_percentile`'s two phase plan
974    #[prost(bool, tag = "9")]
975    pub must_output_per_barrier: bool,
976}
977#[derive(prost_helpers::AnyPB)]
978#[derive(Clone, PartialEq, ::prost::Message)]
979pub struct HashAggNode {
980    #[prost(uint32, repeated, tag = "1")]
981    pub group_key: ::prost::alloc::vec::Vec<u32>,
982    #[prost(message, repeated, tag = "2")]
983    pub agg_calls: ::prost::alloc::vec::Vec<super::expr::AggCall>,
984    #[prost(message, repeated, tag = "3")]
985    pub agg_call_states: ::prost::alloc::vec::Vec<AggCallState>,
986    #[prost(message, optional, tag = "4")]
987    pub intermediate_state_table: ::core::option::Option<super::catalog::Table>,
988    /// Whether to optimize for append only stream.
989    /// It is true when the input is append-only
990    #[prost(bool, tag = "5")]
991    pub is_append_only: bool,
992    #[prost(map = "uint32, message", tag = "6")]
993    pub distinct_dedup_tables: ::std::collections::HashMap<u32, super::catalog::Table>,
994    #[prost(uint32, tag = "7")]
995    pub row_count_index: u32,
996    #[prost(bool, tag = "8")]
997    pub emit_on_window_close: bool,
998    #[prost(enumeration = "AggNodeVersion", tag = "9")]
999    pub version: i32,
1000}
1001#[derive(prost_helpers::AnyPB)]
1002#[derive(Clone, PartialEq, ::prost::Message)]
1003pub struct TopNNode {
1004    /// 0 means no limit as limit of 0 means this node should be optimized away
1005    #[prost(uint64, tag = "1")]
1006    pub limit: u64,
1007    #[prost(uint64, tag = "2")]
1008    pub offset: u64,
1009    #[prost(message, optional, tag = "3")]
1010    pub table: ::core::option::Option<super::catalog::Table>,
1011    #[prost(message, repeated, tag = "4")]
1012    pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1013    #[prost(bool, tag = "5")]
1014    pub with_ties: bool,
1015}
1016#[derive(prost_helpers::AnyPB)]
1017#[derive(Clone, PartialEq, ::prost::Message)]
1018pub struct GroupTopNNode {
1019    /// 0 means no limit as limit of 0 means this node should be optimized away
1020    #[prost(uint64, tag = "1")]
1021    pub limit: u64,
1022    #[prost(uint64, tag = "2")]
1023    pub offset: u64,
1024    #[prost(uint32, repeated, tag = "3")]
1025    pub group_key: ::prost::alloc::vec::Vec<u32>,
1026    #[prost(message, optional, tag = "4")]
1027    pub table: ::core::option::Option<super::catalog::Table>,
1028    #[prost(message, repeated, tag = "5")]
1029    pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1030    #[prost(bool, tag = "6")]
1031    pub with_ties: bool,
1032}
1033#[derive(prost_helpers::AnyPB)]
1034#[derive(Clone, PartialEq, ::prost::Message)]
1035pub struct DeltaExpression {
1036    #[prost(enumeration = "super::expr::expr_node::Type", tag = "1")]
1037    pub delta_type: i32,
1038    #[prost(message, optional, tag = "2")]
1039    pub delta: ::core::option::Option<super::expr::ExprNode>,
1040}
1041/// Deprecated: Use InequalityPairV2 instead.
1042#[derive(prost_helpers::AnyPB)]
1043#[derive(Clone, PartialEq, ::prost::Message)]
1044pub struct InequalityPair {
1045    /// Input index of greater side of inequality.
1046    #[prost(uint32, tag = "1")]
1047    pub key_required_larger: u32,
1048    /// Input index of less side of inequality.
1049    #[prost(uint32, tag = "2")]
1050    pub key_required_smaller: u32,
1051    /// Whether this condition is used to clean state table of `HashJoinExecutor`.
1052    #[prost(bool, tag = "3")]
1053    pub clean_state: bool,
1054    /// greater >= less + delta_expression, if `None`, it represents that greater >= less
1055    #[prost(message, optional, tag = "4")]
1056    pub delta_expression: ::core::option::Option<DeltaExpression>,
1057}
1058#[derive(prost_helpers::AnyPB)]
1059#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1060pub struct InequalityPairV2 {
1061    /// Index of the left side of the inequality (from left input).
1062    #[prost(uint32, tag = "1")]
1063    pub left_idx: u32,
1064    /// Index of the right side of the inequality (from right input, NOT offset by left_cols_num).
1065    #[prost(uint32, tag = "2")]
1066    pub right_idx: u32,
1067    /// Whether this condition is used to clean left state table of `HashJoinExecutor`.
1068    #[prost(bool, tag = "3")]
1069    pub clean_left_state: bool,
1070    /// Whether this condition is used to clean right state table of `HashJoinExecutor`.
1071    #[prost(bool, tag = "4")]
1072    pub clean_right_state: bool,
1073    /// Comparison operator: left_col `<op>` right_col (e.g., \<, \<=, >, >=).
1074    #[prost(enumeration = "InequalityType", tag = "5")]
1075    pub op: i32,
1076}
1077#[derive(prost_helpers::AnyPB)]
1078#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1079pub struct JoinKeyWatermarkIndex {
1080    /// Index in `left_key`/`right_key`.
1081    #[prost(uint32, tag = "1")]
1082    pub index: u32,
1083    /// Whether watermark on this join key triggers state cleaning of the counterpart.
1084    #[prost(bool, tag = "2")]
1085    pub do_state_cleaning: bool,
1086}
1087#[derive(prost_helpers::AnyPB)]
1088#[derive(Clone, PartialEq, ::prost::Message)]
1089pub struct HashJoinWatermarkHandleDesc {
1090    /// Join-key positions whose watermarks should be aligned across both sides.
1091    #[prost(message, repeated, tag = "1")]
1092    pub watermark_indices_in_jk: ::prost::alloc::vec::Vec<JoinKeyWatermarkIndex>,
1093    /// Inequality pairs used for watermark generation / state cleaning.
1094    #[prost(message, repeated, tag = "2")]
1095    pub inequality_pairs: ::prost::alloc::vec::Vec<InequalityPairV2>,
1096}
1097#[derive(prost_helpers::AnyPB)]
1098#[derive(Clone, PartialEq, ::prost::Message)]
1099pub struct HashJoinNode {
1100    #[prost(enumeration = "super::plan_common::JoinType", tag = "1")]
1101    pub join_type: i32,
1102    #[prost(int32, repeated, tag = "2")]
1103    pub left_key: ::prost::alloc::vec::Vec<i32>,
1104    #[prost(int32, repeated, tag = "3")]
1105    pub right_key: ::prost::alloc::vec::Vec<i32>,
1106    #[prost(message, optional, tag = "4")]
1107    pub condition: ::core::option::Option<super::expr::ExprNode>,
1108    /// Used for internal table states.
1109    #[prost(message, optional, tag = "6")]
1110    pub left_table: ::core::option::Option<super::catalog::Table>,
1111    /// Used for internal table states.
1112    #[prost(message, optional, tag = "7")]
1113    pub right_table: ::core::option::Option<super::catalog::Table>,
1114    /// Used for internal table states.
1115    #[prost(message, optional, tag = "8")]
1116    pub left_degree_table: ::core::option::Option<super::catalog::Table>,
1117    /// Used for internal table states.
1118    #[prost(message, optional, tag = "9")]
1119    pub right_degree_table: ::core::option::Option<super::catalog::Table>,
1120    /// The output indices of current node
1121    #[prost(uint32, repeated, tag = "10")]
1122    pub output_indices: ::prost::alloc::vec::Vec<u32>,
1123    /// Left deduped input pk indices. The pk of the left_table and
1124    /// left_degree_table is  \[left_join_key | left_deduped_input_pk_indices\]
1125    /// and is expected to be the shortest key which starts with
1126    /// the join key and satisfies unique constrain.
1127    #[prost(uint32, repeated, tag = "11")]
1128    pub left_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1129    /// Right deduped input pk indices. The pk of the right_table and
1130    /// right_degree_table is  \[right_join_key | right_deduped_input_pk_indices\]
1131    /// and is expected to be the shortest key which starts with
1132    /// the join key and satisfies unique constrain.
1133    #[prost(uint32, repeated, tag = "12")]
1134    pub right_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1135    #[prost(bool, repeated, tag = "13")]
1136    pub null_safe: ::prost::alloc::vec::Vec<bool>,
1137    /// Whether to optimize for append only stream.
1138    /// It is true when the input is append-only
1139    #[prost(bool, tag = "14")]
1140    pub is_append_only: bool,
1141    /// Which encoding will be used to encode join rows in operator cache.
1142    /// Deprecated. Use the one from `StreamingDeveloperConfig` instead.
1143    #[deprecated]
1144    #[prost(enumeration = "JoinEncodingType", tag = "15")]
1145    pub join_encoding_type: i32,
1146    /// Description of watermark-based and other state cleaning strategies.
1147    #[prost(message, optional, tag = "17")]
1148    pub watermark_handle_desc: ::core::option::Option<HashJoinWatermarkHandleDesc>,
1149}
1150#[derive(prost_helpers::AnyPB)]
1151#[derive(Clone, PartialEq, ::prost::Message)]
1152pub struct AsOfJoinNode {
1153    #[prost(enumeration = "super::plan_common::AsOfJoinType", tag = "1")]
1154    pub join_type: i32,
1155    #[prost(int32, repeated, tag = "2")]
1156    pub left_key: ::prost::alloc::vec::Vec<i32>,
1157    #[prost(int32, repeated, tag = "3")]
1158    pub right_key: ::prost::alloc::vec::Vec<i32>,
1159    /// Used for internal table states.
1160    #[prost(message, optional, tag = "4")]
1161    pub left_table: ::core::option::Option<super::catalog::Table>,
1162    /// Used for internal table states.
1163    #[prost(message, optional, tag = "5")]
1164    pub right_table: ::core::option::Option<super::catalog::Table>,
1165    /// The output indices of current node
1166    #[prost(uint32, repeated, tag = "6")]
1167    pub output_indices: ::prost::alloc::vec::Vec<u32>,
1168    /// Left deduped input pk indices. The pk of the left_table and
1169    /// The pk of the left_table is  \[left_join_key | left_inequality_key | left_deduped_input_pk_indices\]
1170    /// left_inequality_key is not used but for forward compatibility.
1171    #[prost(uint32, repeated, tag = "7")]
1172    pub left_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1173    /// Right deduped input pk indices.
1174    /// The pk of the right_table is  \[right_join_key | right_inequality_key | right_deduped_input_pk_indices\]
1175    /// right_inequality_key is not used but for forward compatibility.
1176    #[prost(uint32, repeated, tag = "8")]
1177    pub right_deduped_input_pk_indices: ::prost::alloc::vec::Vec<u32>,
1178    #[prost(bool, repeated, tag = "9")]
1179    pub null_safe: ::prost::alloc::vec::Vec<bool>,
1180    #[prost(message, optional, tag = "10")]
1181    pub asof_desc: ::core::option::Option<super::plan_common::AsOfJoinDesc>,
1182    /// Which encoding will be used to encode join rows in operator cache.
1183    /// Deprecated. Use the one from `StreamingDeveloperConfig` instead.
1184    #[deprecated]
1185    #[prost(enumeration = "JoinEncodingType", tag = "11")]
1186    pub join_encoding_type: i32,
1187    /// Whether to use the cache-based implementation (true) or the no-cache implementation (false).
1188    /// Controlled by session variable `streaming_asof_join_use_cache`.
1189    /// Optional so that legacy plans (without this field) default to true (cache-based).
1190    #[prost(bool, optional, tag = "12")]
1191    pub use_cache: ::core::option::Option<bool>,
1192}
1193#[derive(prost_helpers::AnyPB)]
1194#[derive(Clone, PartialEq, ::prost::Message)]
1195pub struct TemporalJoinNode {
1196    #[prost(enumeration = "super::plan_common::JoinType", tag = "1")]
1197    pub join_type: i32,
1198    #[prost(int32, repeated, tag = "2")]
1199    pub left_key: ::prost::alloc::vec::Vec<i32>,
1200    #[prost(int32, repeated, tag = "3")]
1201    pub right_key: ::prost::alloc::vec::Vec<i32>,
1202    #[prost(bool, repeated, tag = "4")]
1203    pub null_safe: ::prost::alloc::vec::Vec<bool>,
1204    #[prost(message, optional, tag = "5")]
1205    pub condition: ::core::option::Option<super::expr::ExprNode>,
1206    /// The output indices of current node
1207    #[prost(uint32, repeated, tag = "6")]
1208    pub output_indices: ::prost::alloc::vec::Vec<u32>,
1209    /// The table desc of the lookup side table.
1210    #[prost(message, optional, tag = "7")]
1211    pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1212    /// The output indices of the lookup side table
1213    #[prost(uint32, repeated, tag = "8")]
1214    pub table_output_indices: ::prost::alloc::vec::Vec<u32>,
1215    /// The state table used for non-append-only temporal join.
1216    #[prost(message, optional, tag = "9")]
1217    pub memo_table: ::core::option::Option<super::catalog::Table>,
1218    /// If it is a nested lool temporal join
1219    #[prost(bool, tag = "10")]
1220    pub is_nested_loop: bool,
1221    /// Whether lookup-side changes are broadcast to every temporal join actor.
1222    #[prost(bool, tag = "11")]
1223    pub is_broadcast: bool,
1224}
1225#[derive(prost_helpers::AnyPB)]
1226#[derive(Clone, PartialEq, ::prost::Message)]
1227pub struct DynamicFilterNode {
1228    #[prost(uint32, tag = "1")]
1229    pub left_key: u32,
1230    /// Must be one of \<, \<=, >, >=
1231    #[prost(message, optional, tag = "2")]
1232    pub condition: ::core::option::Option<super::expr::ExprNode>,
1233    /// Left table stores all states with predicate possibly not NULL.
1234    #[prost(message, optional, tag = "3")]
1235    pub left_table: ::core::option::Option<super::catalog::Table>,
1236    /// Right table stores single value from RHS of predicate.
1237    #[prost(message, optional, tag = "4")]
1238    pub right_table: ::core::option::Option<super::catalog::Table>,
1239    /// If the right side's change always make the condition more relaxed.
1240    /// In other words, make more record in the left side satisfy the condition.
1241    /// If this is true, we need to store LHS records which do not match the condition in the internal table.
1242    /// When the condition changes, we will tell downstream to insert the LHS records which now match the condition.
1243    /// If this is false, we need to store RHS records which match the condition in the internal table.
1244    /// When the condition changes, we will tell downstream to delete the LHS records which now no longer match the condition.
1245    #[deprecated]
1246    #[prost(bool, tag = "5")]
1247    pub condition_always_relax: bool,
1248    /// Whether the dynamic filter can clean its left state table by watermark.
1249    #[prost(bool, tag = "6")]
1250    pub cleaned_by_watermark: bool,
1251}
1252/// Delta join with two indexes. This is a pseudo plan node generated on frontend. On meta
1253/// service, it will be rewritten into lookup joins.
1254#[derive(prost_helpers::AnyPB)]
1255#[derive(Clone, PartialEq, ::prost::Message)]
1256pub struct DeltaIndexJoinNode {
1257    #[prost(enumeration = "super::plan_common::JoinType", tag = "1")]
1258    pub join_type: i32,
1259    #[prost(int32, repeated, tag = "2")]
1260    pub left_key: ::prost::alloc::vec::Vec<i32>,
1261    #[prost(int32, repeated, tag = "3")]
1262    pub right_key: ::prost::alloc::vec::Vec<i32>,
1263    #[prost(message, optional, tag = "4")]
1264    pub condition: ::core::option::Option<super::expr::ExprNode>,
1265    /// Table id of the left index.
1266    #[prost(uint32, tag = "7", wrapper = "crate::id::TableId")]
1267    pub left_table_id: crate::id::TableId,
1268    /// Table id of the right index.
1269    #[prost(uint32, tag = "8", wrapper = "crate::id::TableId")]
1270    pub right_table_id: crate::id::TableId,
1271    /// Info about the left index
1272    #[prost(message, optional, tag = "9")]
1273    pub left_info: ::core::option::Option<ArrangementInfo>,
1274    /// Info about the right index
1275    #[prost(message, optional, tag = "10")]
1276    pub right_info: ::core::option::Option<ArrangementInfo>,
1277    /// the output indices of current node
1278    #[prost(uint32, repeated, tag = "11")]
1279    pub output_indices: ::prost::alloc::vec::Vec<u32>,
1280}
1281#[derive(prost_helpers::AnyPB)]
1282#[derive(Clone, PartialEq, ::prost::Message)]
1283pub struct HopWindowNode {
1284    #[prost(uint32, tag = "1")]
1285    pub time_col: u32,
1286    #[prost(message, optional, tag = "2")]
1287    pub window_slide: ::core::option::Option<super::data::Interval>,
1288    #[prost(message, optional, tag = "3")]
1289    pub window_size: ::core::option::Option<super::data::Interval>,
1290    #[prost(uint32, repeated, tag = "4")]
1291    pub output_indices: ::prost::alloc::vec::Vec<u32>,
1292    #[prost(message, repeated, tag = "5")]
1293    pub window_start_exprs: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
1294    #[prost(message, repeated, tag = "6")]
1295    pub window_end_exprs: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
1296}
1297#[derive(prost_helpers::AnyPB)]
1298#[derive(Clone, PartialEq, ::prost::Message)]
1299pub struct MergeNode {
1300    /// **WARNING**: Use this field with caution.
1301    ///
1302    /// `upstream_actor_id` stored in the plan node in `Fragment` meta model cannot be directly used.
1303    /// See `compose_fragment`.
1304    /// The field is deprecated because the upstream actor info is provided separately instead of
1305    /// injected here in the node.
1306    #[deprecated]
1307    #[prost(uint32, repeated, packed = "false", tag = "1")]
1308    pub upstream_actor_id: ::prost::alloc::vec::Vec<u32>,
1309    #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
1310    pub upstream_fragment_id: crate::id::FragmentId,
1311    /// Type of the upstream dispatcher. If there's always one upstream according to this
1312    /// type, the compute node may use the `ReceiverExecutor` as an optimization.
1313    #[prost(enumeration = "DispatcherType", tag = "3")]
1314    pub upstream_dispatcher_type: i32,
1315    /// The schema of input columns. Already deprecated.
1316    #[deprecated]
1317    #[prost(message, repeated, tag = "4")]
1318    pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
1319    /// Allows this merge to have no upstream actors, either initially or after a MergeUpdate.
1320    /// When true, the merge must not use the singleton receiver optimization because a singleton
1321    /// receiver cannot transition to empty.
1322    #[prost(bool, tag = "5")]
1323    pub allow_empty_upstream: bool,
1324}
1325/// passed from frontend to meta, used by fragmenter to generate `MergeNode`
1326/// and maybe `DispatcherNode` later.
1327#[derive(prost_helpers::AnyPB)]
1328#[derive(Clone, PartialEq, ::prost::Message)]
1329pub struct ExchangeNode {
1330    #[prost(message, optional, tag = "1")]
1331    pub strategy: ::core::option::Option<DispatchStrategy>,
1332}
1333/// StreamScanNode reads data from upstream table first, and then pass all events to downstream.
1334/// It always these 2 inputs in the following order:
1335///
1336/// 1. A MergeNode (as a placeholder) of upstream.
1337/// 1. A BatchPlanNode for the snapshot read.
1338#[derive(prost_helpers::AnyPB)]
1339#[derive(Clone, PartialEq, ::prost::Message)]
1340pub struct StreamScanNode {
1341    #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
1342    pub table_id: crate::id::TableId,
1343    /// The columns from the upstream table that'll be internally required by this stream scan node.
1344    ///
1345    /// * For non-backfill stream scan node, it's the same as the output columns.
1346    /// * For backfill stream scan node, there're additionally primary key columns.
1347    #[prost(int32, repeated, tag = "2")]
1348    pub upstream_column_ids: ::prost::alloc::vec::Vec<i32>,
1349    /// The columns to be output by this stream scan node. The index is based on the internal required columns.
1350    ///
1351    /// * For non-backfill stream scan node, it's simply all the columns.
1352    /// * For backfill stream scan node, this strips the primary key columns if they're unnecessary.
1353    #[prost(uint32, repeated, tag = "3")]
1354    pub output_indices: ::prost::alloc::vec::Vec<u32>,
1355    /// Generally, the barrier needs to be rearranged during the MV creation process, so that data can
1356    /// be flushed to shared buffer periodically, instead of making the first epoch from batch query extra
1357    /// large. However, in some cases, e.g., shared state, the barrier cannot be rearranged in StreamScanNode.
1358    /// StreamScanType is used to decide which implementation for the StreamScanNode.
1359    #[prost(enumeration = "StreamScanType", tag = "4")]
1360    pub stream_scan_type: i32,
1361    /// / The state table used by Backfill operator for persisting internal state
1362    #[prost(message, optional, tag = "5")]
1363    pub state_table: ::core::option::Option<super::catalog::Table>,
1364    /// The upstream materialized view info used by backfill.
1365    /// Used iff `ChainType::Backfill`.
1366    #[prost(message, optional, tag = "7")]
1367    pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1368    /// The backfill rate limit for the stream scan node.
1369    #[prost(uint32, optional, tag = "8")]
1370    pub rate_limit: ::core::option::Option<u32>,
1371    /// Snapshot read every N barriers
1372    #[deprecated]
1373    #[prost(uint32, tag = "9")]
1374    pub snapshot_read_barrier_interval: u32,
1375    /// The state table used by ArrangementBackfill to replicate upstream mview's state table.
1376    /// Used iff `ChainType::ArrangementBackfill`.
1377    #[prost(message, optional, tag = "10")]
1378    pub arrangement_table: ::core::option::Option<super::catalog::Table>,
1379    #[prost(uint64, optional, tag = "11")]
1380    pub snapshot_backfill_epoch: ::core::option::Option<u64>,
1381    /// Full scan range pushdown (eq_conds + range bounds on PK).
1382    #[prost(message, optional, tag = "13")]
1383    pub pk_scan_range: ::core::option::Option<super::batch_plan::ScanRange>,
1384}
1385/// Config options for CDC backfill
1386#[derive(prost_helpers::AnyPB)]
1387#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1388pub struct StreamCdcScanOptions {
1389    /// Whether skip the backfill and only consume from upstream.
1390    #[prost(bool, tag = "1")]
1391    pub disable_backfill: bool,
1392    #[prost(uint32, tag = "2")]
1393    pub snapshot_barrier_interval: u32,
1394    #[prost(uint32, tag = "3")]
1395    pub snapshot_batch_size: u32,
1396    #[prost(uint32, tag = "4")]
1397    pub backfill_parallelism: u32,
1398    #[prost(uint64, tag = "5")]
1399    pub backfill_num_rows_per_split: u64,
1400    #[prost(bool, tag = "6")]
1401    pub backfill_as_even_splits: bool,
1402    #[prost(uint32, tag = "7")]
1403    pub backfill_split_pk_column_index: u32,
1404}
1405#[derive(prost_helpers::AnyPB)]
1406#[derive(Clone, PartialEq, ::prost::Message)]
1407pub struct StreamCdcScanNode {
1408    #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
1409    pub table_id: crate::id::TableId,
1410    /// The columns from the upstream table that'll be internally required by this stream scan node.
1411    /// Contains Primary Keys and Output columns.
1412    #[prost(int32, repeated, tag = "2")]
1413    pub upstream_column_ids: ::prost::alloc::vec::Vec<i32>,
1414    /// Strips the primary key columns if they're unnecessary.
1415    #[prost(uint32, repeated, tag = "3")]
1416    pub output_indices: ::prost::alloc::vec::Vec<u32>,
1417    /// The state table used by CdcBackfill operator for persisting internal state
1418    #[prost(message, optional, tag = "4")]
1419    pub state_table: ::core::option::Option<super::catalog::Table>,
1420    /// The external table that will be backfilled for CDC.
1421    #[prost(message, optional, tag = "5")]
1422    pub cdc_table_desc: ::core::option::Option<super::plan_common::ExternalTableDesc>,
1423    /// The backfill rate limit for the stream cdc scan node.
1424    #[prost(uint32, optional, tag = "6")]
1425    pub rate_limit: ::core::option::Option<u32>,
1426    /// Whether skip the backfill and only consume from upstream.
1427    /// keep it for backward compatibility, new stream plan will use `options.disable_backfill`
1428    #[prost(bool, tag = "7")]
1429    pub disable_backfill: bool,
1430    #[prost(message, optional, tag = "8")]
1431    pub options: ::core::option::Option<StreamCdcScanOptions>,
1432}
1433/// BatchPlanNode is used for mv on mv snapshot read.
1434/// BatchPlanNode is supposed to carry a batch plan that can be optimized with the streaming plan_common.
1435/// Currently, streaming to batch push down is not yet supported, BatchPlanNode is simply a table scan.
1436#[derive(prost_helpers::AnyPB)]
1437#[derive(Clone, PartialEq, ::prost::Message)]
1438pub struct BatchPlanNode {
1439    #[prost(message, optional, tag = "1")]
1440    pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1441    #[prost(int32, repeated, tag = "2")]
1442    pub column_ids: ::prost::alloc::vec::Vec<i32>,
1443}
1444#[derive(prost_helpers::AnyPB)]
1445#[derive(Clone, PartialEq, ::prost::Message)]
1446pub struct ArrangementInfo {
1447    /// Order key of the arrangement, including order by columns and pk from the materialize
1448    /// executor.
1449    #[prost(message, repeated, tag = "1")]
1450    pub arrange_key_orders: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1451    /// Column descs of the arrangement
1452    #[prost(message, repeated, tag = "2")]
1453    pub column_descs: ::prost::alloc::vec::Vec<super::plan_common::ColumnDesc>,
1454    /// Used to build storage table by stream lookup join of delta join.
1455    #[prost(message, optional, tag = "4")]
1456    pub table_desc: ::core::option::Option<super::plan_common::StorageTableDesc>,
1457    /// Output index columns
1458    #[prost(uint32, repeated, tag = "5")]
1459    pub output_col_idx: ::prost::alloc::vec::Vec<u32>,
1460}
1461/// Special node for shared state, which will only be produced in fragmenter. ArrangeNode will
1462/// produce a special Materialize executor, which materializes data for downstream to query.
1463#[derive(prost_helpers::AnyPB)]
1464#[derive(Clone, PartialEq, ::prost::Message)]
1465pub struct ArrangeNode {
1466    /// Info about the arrangement
1467    #[prost(message, optional, tag = "1")]
1468    pub table_info: ::core::option::Option<ArrangementInfo>,
1469    /// Hash key of the materialize node, which is a subset of pk.
1470    #[prost(uint32, repeated, tag = "2")]
1471    pub distribution_key: ::prost::alloc::vec::Vec<u32>,
1472    /// Used for internal table states.
1473    #[prost(message, optional, tag = "3")]
1474    pub table: ::core::option::Option<super::catalog::Table>,
1475}
1476/// Special node for shared state. LookupNode will join an arrangement with a stream.
1477#[derive(prost_helpers::AnyPB)]
1478#[derive(Clone, PartialEq, ::prost::Message)]
1479pub struct LookupNode {
1480    /// Join key of the arrangement side
1481    #[prost(int32, repeated, tag = "1")]
1482    pub arrange_key: ::prost::alloc::vec::Vec<i32>,
1483    /// Join key of the stream side
1484    #[prost(int32, repeated, tag = "2")]
1485    pub stream_key: ::prost::alloc::vec::Vec<i32>,
1486    /// Whether to join the current epoch of arrangement
1487    #[prost(bool, tag = "3")]
1488    pub use_current_epoch: bool,
1489    /// Sometimes we need to re-order the output data to meet the requirement of schema.
1490    /// By default, lookup executor will produce `<arrangement side, stream side>`. We
1491    /// will then apply the column mapping to the combined result.
1492    #[prost(int32, repeated, tag = "4")]
1493    pub column_mapping: ::prost::alloc::vec::Vec<i32>,
1494    /// Info about the arrangement
1495    #[prost(message, optional, tag = "7")]
1496    pub arrangement_table_info: ::core::option::Option<ArrangementInfo>,
1497    #[prost(oneof = "lookup_node::ArrangementTableId", tags = "5, 6")]
1498    pub arrangement_table_id: ::core::option::Option<lookup_node::ArrangementTableId>,
1499}
1500/// Nested message and enum types in `LookupNode`.
1501pub mod lookup_node {
1502    #[derive(prost_helpers::AnyPB)]
1503    #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Oneof)]
1504    pub enum ArrangementTableId {
1505        /// Table Id of the arrangement (when created along with join plan)
1506        #[prost(uint32, tag = "5", wrapper = "crate::id::TableId")]
1507        TableId(crate::id::TableId),
1508        /// Table Id of the arrangement (when using index)
1509        #[prost(uint32, tag = "6", wrapper = "crate::id::TableId")]
1510        IndexId(crate::id::TableId),
1511    }
1512}
1513/// WatermarkFilter needs to filter the upstream data by the water mark.
1514#[derive(prost_helpers::AnyPB)]
1515#[derive(Clone, PartialEq, ::prost::Message)]
1516pub struct WatermarkFilterNode {
1517    /// The watermark descs
1518    #[prost(message, repeated, tag = "1")]
1519    pub watermark_descs: ::prost::alloc::vec::Vec<super::catalog::WatermarkDesc>,
1520    /// The tables used to persist watermarks, the key is vnode.
1521    #[prost(message, repeated, tag = "2")]
1522    pub tables: ::prost::alloc::vec::Vec<super::catalog::Table>,
1523}
1524/// Acts like a merger, but on different inputs.
1525#[derive(prost_helpers::AnyPB)]
1526#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1527pub struct UnionNode {}
1528/// Special node for shared state. Merge and align barrier from upstreams. Pipe inputs in order.
1529#[derive(prost_helpers::AnyPB)]
1530#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1531pub struct LookupUnionNode {
1532    #[prost(uint32, repeated, tag = "1")]
1533    pub order: ::prost::alloc::vec::Vec<u32>,
1534}
1535#[derive(prost_helpers::AnyPB)]
1536#[derive(Clone, PartialEq, ::prost::Message)]
1537pub struct ExpandNode {
1538    #[prost(message, repeated, tag = "1")]
1539    pub column_subsets: ::prost::alloc::vec::Vec<expand_node::Subset>,
1540}
1541/// Nested message and enum types in `ExpandNode`.
1542pub mod expand_node {
1543    #[derive(prost_helpers::AnyPB)]
1544    #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1545    pub struct Subset {
1546        #[prost(uint32, repeated, tag = "1")]
1547        pub column_indices: ::prost::alloc::vec::Vec<u32>,
1548    }
1549}
1550#[derive(prost_helpers::AnyPB)]
1551#[derive(Clone, PartialEq, ::prost::Message)]
1552pub struct ProjectSetNode {
1553    #[prost(message, repeated, tag = "1")]
1554    pub select_list: ::prost::alloc::vec::Vec<super::expr::ProjectSetSelectItem>,
1555    /// this two field is expressing a list of usize pair, which means when project receives a
1556    /// watermark with `watermark_input_cols\[i\]` column index, it should derive a new watermark
1557    /// with `watermark_output_cols\[i\]`th expression
1558    #[prost(uint32, repeated, tag = "2")]
1559    pub watermark_input_cols: ::prost::alloc::vec::Vec<u32>,
1560    #[prost(uint32, repeated, tag = "3")]
1561    pub watermark_expr_indices: ::prost::alloc::vec::Vec<u32>,
1562    #[prost(uint32, repeated, tag = "4")]
1563    pub nondecreasing_exprs: ::prost::alloc::vec::Vec<u32>,
1564}
1565/// Sorts inputs and outputs ordered data based on watermark.
1566#[derive(prost_helpers::AnyPB)]
1567#[derive(Clone, PartialEq, ::prost::Message)]
1568pub struct SortNode {
1569    /// Persists data above watermark.
1570    #[prost(message, optional, tag = "1")]
1571    pub state_table: ::core::option::Option<super::catalog::Table>,
1572    /// Column index of watermark to perform sorting.
1573    #[prost(uint32, tag = "2")]
1574    pub sort_column_index: u32,
1575}
1576/// Merges two streams from streaming and batch for data manipulation.
1577#[derive(prost_helpers::AnyPB)]
1578#[derive(Clone, PartialEq, ::prost::Message)]
1579pub struct DmlNode {
1580    /// Id of the table on which DML performs.
1581    #[prost(uint32, tag = "1", wrapper = "crate::id::TableId")]
1582    pub table_id: crate::id::TableId,
1583    /// Version of the table.
1584    #[prost(uint64, tag = "3")]
1585    pub table_version_id: u64,
1586    /// Column descriptions of the table.
1587    #[prost(message, repeated, tag = "2")]
1588    pub column_descs: ::prost::alloc::vec::Vec<super::plan_common::ColumnDesc>,
1589    #[prost(uint32, optional, tag = "4")]
1590    pub rate_limit: ::core::option::Option<u32>,
1591}
1592#[derive(prost_helpers::AnyPB)]
1593#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1594pub struct RowIdGenNode {
1595    #[prost(uint64, tag = "1")]
1596    pub row_id_index: u64,
1597}
1598#[derive(prost_helpers::AnyPB)]
1599#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1600pub struct NowModeUpdateCurrent {}
1601#[derive(prost_helpers::AnyPB)]
1602#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1603pub struct NowModeGenerateSeries {
1604    #[prost(message, optional, tag = "1")]
1605    pub start_timestamp: ::core::option::Option<super::data::Datum>,
1606    #[prost(message, optional, tag = "2")]
1607    pub interval: ::core::option::Option<super::data::Datum>,
1608}
1609#[derive(prost_helpers::AnyPB)]
1610#[derive(Clone, PartialEq, ::prost::Message)]
1611pub struct NowNode {
1612    /// Persists emitted 'now'.
1613    #[prost(message, optional, tag = "1")]
1614    pub state_table: ::core::option::Option<super::catalog::Table>,
1615    #[prost(oneof = "now_node::Mode", tags = "101, 102")]
1616    pub mode: ::core::option::Option<now_node::Mode>,
1617}
1618/// Nested message and enum types in `NowNode`.
1619pub mod now_node {
1620    #[derive(prost_helpers::AnyPB)]
1621    #[derive(Clone, PartialEq, Eq, Hash, ::prost::Oneof)]
1622    pub enum Mode {
1623        #[prost(message, tag = "101")]
1624        UpdateCurrent(super::NowModeUpdateCurrent),
1625        #[prost(message, tag = "102")]
1626        GenerateSeries(super::NowModeGenerateSeries),
1627    }
1628}
1629#[derive(prost_helpers::AnyPB)]
1630#[derive(Clone, PartialEq, ::prost::Message)]
1631pub struct ValuesNode {
1632    #[prost(message, repeated, tag = "1")]
1633    pub tuples: ::prost::alloc::vec::Vec<values_node::ExprTuple>,
1634    #[prost(message, repeated, tag = "2")]
1635    pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
1636}
1637/// Nested message and enum types in `ValuesNode`.
1638pub mod values_node {
1639    #[derive(prost_helpers::AnyPB)]
1640    #[derive(Clone, PartialEq, ::prost::Message)]
1641    pub struct ExprTuple {
1642        #[prost(message, repeated, tag = "1")]
1643        pub cells: ::prost::alloc::vec::Vec<super::super::expr::ExprNode>,
1644    }
1645}
1646#[derive(prost_helpers::AnyPB)]
1647#[derive(Clone, PartialEq, ::prost::Message)]
1648pub struct DedupNode {
1649    #[prost(message, optional, tag = "1")]
1650    pub state_table: ::core::option::Option<super::catalog::Table>,
1651    #[prost(uint32, repeated, tag = "2")]
1652    pub dedup_column_indices: ::prost::alloc::vec::Vec<u32>,
1653}
1654#[derive(prost_helpers::AnyPB)]
1655#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1656pub struct NoOpNode {}
1657#[derive(prost_helpers::AnyPB)]
1658#[derive(Clone, PartialEq, ::prost::Message)]
1659pub struct EowcOverWindowNode {
1660    #[prost(message, repeated, tag = "1")]
1661    pub calls: ::prost::alloc::vec::Vec<super::expr::WindowFunction>,
1662    #[prost(uint32, repeated, tag = "2")]
1663    pub partition_by: ::prost::alloc::vec::Vec<u32>,
1664    /// use `repeated` in case of future extension, now only one column is allowed
1665    #[prost(message, repeated, tag = "3")]
1666    pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1667    #[prost(message, optional, tag = "4")]
1668    pub state_table: ::core::option::Option<super::catalog::Table>,
1669    /// Optional state table for persisting window function intermediate states.
1670    /// Currently used for numbering functions (row_number/rank/dense_rank).
1671    #[prost(message, optional, tag = "5")]
1672    pub intermediate_state_table: ::core::option::Option<super::catalog::Table>,
1673}
1674#[derive(prost_helpers::AnyPB)]
1675#[derive(Clone, PartialEq, ::prost::Message)]
1676pub struct OverWindowNode {
1677    #[prost(message, repeated, tag = "1")]
1678    pub calls: ::prost::alloc::vec::Vec<super::expr::WindowFunction>,
1679    #[prost(uint32, repeated, tag = "2")]
1680    pub partition_by: ::prost::alloc::vec::Vec<u32>,
1681    #[prost(message, repeated, tag = "3")]
1682    pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1683    #[prost(message, optional, tag = "4")]
1684    pub state_table: ::core::option::Option<super::catalog::Table>,
1685    /// Deprecated. Use the one from `StreamingDeveloperConfig` instead.
1686    #[deprecated]
1687    #[prost(enumeration = "OverWindowCachePolicy", tag = "5")]
1688    pub cache_policy: i32,
1689    /// Whether the executor is allowed to clean up state rows in recently touched partitions that
1690    /// can never be accessed again, driven by the watermark on the first `order_by` column. The
1691    /// optimizer only enables this when the input is append-only, all window frames are bounded
1692    /// `ROWS` frames, and the first `order_by` column is a watermark column.
1693    #[prost(bool, tag = "6")]
1694    pub enable_state_cleaning: bool,
1695}
1696#[derive(prost_helpers::AnyPB)]
1697#[derive(Clone, Copy, PartialEq, ::prost::Message)]
1698pub struct LocalApproxPercentileNode {
1699    #[prost(double, tag = "1")]
1700    pub base: f64,
1701    #[prost(uint32, tag = "2")]
1702    pub percentile_index: u32,
1703}
1704#[derive(prost_helpers::AnyPB)]
1705#[derive(Clone, PartialEq, ::prost::Message)]
1706pub struct GlobalApproxPercentileNode {
1707    #[prost(double, tag = "1")]
1708    pub base: f64,
1709    #[prost(double, tag = "2")]
1710    pub quantile: f64,
1711    #[prost(message, optional, tag = "3")]
1712    pub bucket_state_table: ::core::option::Option<super::catalog::Table>,
1713    #[prost(message, optional, tag = "4")]
1714    pub count_state_table: ::core::option::Option<super::catalog::Table>,
1715}
1716#[derive(prost_helpers::AnyPB)]
1717#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1718pub struct RowMergeNode {
1719    #[prost(message, optional, tag = "1")]
1720    pub lhs_mapping: ::core::option::Option<super::catalog::ColIndexMapping>,
1721    #[prost(message, optional, tag = "2")]
1722    pub rhs_mapping: ::core::option::Option<super::catalog::ColIndexMapping>,
1723}
1724#[derive(prost_helpers::AnyPB)]
1725#[derive(Clone, PartialEq, ::prost::Message)]
1726pub struct SyncLogStoreNode {
1727    #[prost(message, optional, tag = "1")]
1728    pub log_store_table: ::core::option::Option<super::catalog::Table>,
1729    /// Deprecated. Use the one from `StreamingDeveloperConfig` instead.
1730    #[deprecated]
1731    #[prost(uint32, optional, tag = "2")]
1732    pub pause_duration_ms: ::core::option::Option<u32>,
1733    /// Deprecated. Use the one from `StreamingDeveloperConfig` instead.
1734    #[deprecated]
1735    #[prost(uint32, optional, tag = "3")]
1736    pub buffer_size: ::core::option::Option<u32>,
1737    #[prost(bool, tag = "4")]
1738    pub aligned: bool,
1739}
1740#[derive(prost_helpers::AnyPB)]
1741#[derive(Clone, PartialEq, ::prost::Message)]
1742pub struct MaterializedExprsNode {
1743    #[prost(message, repeated, tag = "1")]
1744    pub exprs: ::prost::alloc::vec::Vec<super::expr::ExprNode>,
1745    #[prost(message, optional, tag = "2")]
1746    pub state_table: ::core::option::Option<super::catalog::Table>,
1747    #[prost(uint32, optional, tag = "3")]
1748    pub state_clean_col_idx: ::core::option::Option<u32>,
1749}
1750#[derive(prost_helpers::AnyPB)]
1751#[derive(Clone, PartialEq, ::prost::Message)]
1752pub struct VectorIndexWriteNode {
1753    #[prost(message, optional, tag = "1")]
1754    pub table: ::core::option::Option<super::catalog::Table>,
1755}
1756#[derive(prost_helpers::AnyPB)]
1757#[derive(Clone, PartialEq, ::prost::Message)]
1758pub struct VectorIndexLookupJoinNode {
1759    #[prost(message, optional, tag = "1")]
1760    pub reader_desc: ::core::option::Option<super::plan_common::VectorIndexReaderDesc>,
1761    #[prost(uint32, tag = "2")]
1762    pub vector_column_idx: u32,
1763}
1764#[derive(prost_helpers::AnyPB)]
1765#[derive(Clone, PartialEq, ::prost::Message)]
1766pub struct UpstreamSinkUnionNode {
1767    /// It is always empty in the persisted metadata, and get filled before we spawn the actors.
1768    /// The actual upstream info may be added and removed dynamically at runtime.
1769    #[prost(message, repeated, tag = "1")]
1770    pub init_upstreams: ::prost::alloc::vec::Vec<UpstreamSinkInfo>,
1771}
1772#[derive(prost_helpers::AnyPB)]
1773#[derive(Clone, PartialEq, ::prost::Message)]
1774pub struct LocalityProviderNode {
1775    /// Column indices that define locality
1776    #[prost(uint32, repeated, tag = "1")]
1777    pub locality_columns: ::prost::alloc::vec::Vec<u32>,
1778    /// State table for buffering input data
1779    #[prost(message, optional, tag = "2")]
1780    pub state_table: ::core::option::Option<super::catalog::Table>,
1781    /// Progress table for tracking backfill progress
1782    #[prost(message, optional, tag = "3")]
1783    pub progress_table: ::core::option::Option<super::catalog::Table>,
1784    /// Rate limit of the drain phase, set by `BACKFILL_RATE_LIMIT`.
1785    #[prost(uint32, optional, tag = "4")]
1786    pub rate_limit: ::core::option::Option<u32>,
1787}
1788#[derive(prost_helpers::AnyPB)]
1789#[derive(Clone, PartialEq, ::prost::Message)]
1790pub struct EowcGapFillNode {
1791    #[prost(uint32, tag = "1")]
1792    pub time_column_index: u32,
1793    #[prost(message, optional, tag = "2")]
1794    pub interval: ::core::option::Option<super::expr::ExprNode>,
1795    #[prost(uint32, repeated, tag = "3")]
1796    pub fill_columns: ::prost::alloc::vec::Vec<u32>,
1797    #[prost(string, repeated, tag = "4")]
1798    pub fill_strategies: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
1799    #[prost(message, optional, tag = "6")]
1800    pub prev_row_table: ::core::option::Option<super::catalog::Table>,
1801    #[prost(uint32, repeated, tag = "7")]
1802    pub partition_by_indices: ::prost::alloc::vec::Vec<u32>,
1803}
1804#[derive(prost_helpers::AnyPB)]
1805#[derive(Clone, PartialEq, ::prost::Message)]
1806pub struct GapFillNode {
1807    #[prost(uint32, tag = "1")]
1808    pub time_column_index: u32,
1809    #[prost(message, optional, tag = "2")]
1810    pub interval: ::core::option::Option<super::expr::ExprNode>,
1811    #[prost(uint32, repeated, tag = "3")]
1812    pub fill_columns: ::prost::alloc::vec::Vec<u32>,
1813    #[prost(string, repeated, tag = "4")]
1814    pub fill_strategies: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
1815    #[prost(message, optional, tag = "5")]
1816    pub state_table: ::core::option::Option<super::catalog::Table>,
1817    #[prost(uint32, repeated, tag = "6")]
1818    pub partition_by_indices: ::prost::alloc::vec::Vec<u32>,
1819    #[prost(uint32, repeated, tag = "7")]
1820    pub pointer_key_indices: ::prost::alloc::vec::Vec<u32>,
1821}
1822/// SQL:2016 MATCH_RECOGNIZE (row pattern recognition). v1: append-only input,
1823/// ONE ROW PER MATCH. PARTITION BY / ORDER BY are input column indices.
1824#[derive(prost_helpers::AnyPB)]
1825#[derive(Clone, PartialEq, ::prost::Message)]
1826pub struct MatchRecognizeNode {
1827    #[prost(uint32, repeated, tag = "1")]
1828    pub partition_by: ::prost::alloc::vec::Vec<u32>,
1829    /// ORDER BY keys, carried as ColumnOrder like every other ordered streaming node. v1 only supports
1830    /// the default ascending order (the binder rejects DESC / explicit NULLS); the executor asserts it.
1831    #[prost(message, repeated, tag = "2")]
1832    pub order_by: ::prost::alloc::vec::Vec<super::common::ColumnOrder>,
1833    #[prost(message, repeated, tag = "3")]
1834    pub measures: ::prost::alloc::vec::Vec<MatchRecognizeMeasure>,
1835    #[prost(message, repeated, tag = "5")]
1836    pub defines: ::prost::alloc::vec::Vec<MatchRecognizeDefine>,
1837    /// The PATTERN clause as a structured tree (the v1 subset: variables, concatenation, alternation,
1838    /// quantification and PERMUTE; anchors and exclusions are rejected at parse time). Field 7
1839    /// previously carried the pattern as text; it is now reserved.
1840    #[prost(message, optional, tag = "12")]
1841    pub pattern_node: ::core::option::Option<MatchRecognizePatternNode>,
1842    /// Rows referenced by live partial matches, keyed (partition..., order..., seq): the retained
1843    /// window the matcher re-feeds on recovery. Input arrives already ordered (see `input_mode`), so
1844    /// this holds only what live matches still need -- there is no out-of-order buffer here.
1845    #[prost(message, optional, tag = "8")]
1846    pub state_table: ::core::option::Option<super::catalog::Table>,
1847    /// AFTER MATCH SKIP strategy.
1848    #[prost(message, optional, tag = "9")]
1849    pub after_match_skip: ::core::option::Option<MatchRecognizeAfterMatchSkip>,
1850    /// WITHIN span check: a predicate over a synthetic `\[last_order_key, first_order_key\]` row
1851    /// (`InputRef(0) - InputRef(1) <= interval`). Absent when there is no WITHIN clause.
1852    #[prost(message, optional, tag = "11")]
1853    pub within: ::core::option::Option<super::expr::ExprNode>,
1854    /// WITHIN deadline: `first_order_key + interval` over a synthetic `\[first_order_key\]` row
1855    /// (`InputRef(0) + interval`). The watermark at which a partial starting at that row expires;
1856    /// used to wake an idle partition to evict it. Absent when there is no WITHIN clause.
1857    #[prost(message, optional, tag = "15")]
1858    pub within_deadline: ::core::option::Option<super::expr::ExprNode>,
1859    /// How the ordered-input requirement is satisfied. EVENT_TIME: an EowcSort upstream (same
1860    /// fragment) emits rows in full ORDER BY order, strictly below each forwarded watermark.
1861    /// PROCESSING_TIME (not yet planned): rows are matched in observed arrival order.
1862    #[prost(enumeration = "MatchRecognizeInputMode", tag = "16")]
1863    pub input_mode: i32,
1864}
1865/// A node in the structured row pattern tree, mirroring the executor-side pattern AST. Groups and
1866/// anchors are not represented: parenthesized groups are flattened during lowering, and anchors /
1867/// exclusions are rejected at planning time (not supported in v1).
1868#[derive(prost_helpers::AnyPB)]
1869#[derive(Clone, PartialEq, ::prost::Message)]
1870pub struct MatchRecognizePatternNode {
1871    #[prost(oneof = "match_recognize_pattern_node::Node", tags = "1, 2, 3, 4, 5")]
1872    pub node: ::core::option::Option<match_recognize_pattern_node::Node>,
1873}
1874/// Nested message and enum types in `MatchRecognizePatternNode`.
1875pub mod match_recognize_pattern_node {
1876    #[derive(prost_helpers::AnyPB)]
1877    #[derive(Clone, PartialEq, ::prost::Oneof)]
1878    pub enum Node {
1879        /// A pattern variable, e.g. `A`.
1880        #[prost(string, tag = "1")]
1881        Var(::prost::alloc::string::String),
1882        /// Concatenation, e.g. `A B C`.
1883        #[prost(message, tag = "2")]
1884        Concat(super::MatchRecognizePatternSeq),
1885        /// Alternation, e.g. `A | B | C`.
1886        #[prost(message, tag = "3")]
1887        Alternation(super::MatchRecognizePatternSeq),
1888        /// A quantified sub-pattern, e.g. `A+`.
1889        #[prost(message, tag = "4")]
1890        Quantified(::prost::alloc::boxed::Box<super::MatchRecognizeQuantifiedPattern>),
1891        /// `PERMUTE(a, b, ...)`.
1892        #[prost(message, tag = "5")]
1893        Permute(super::MatchRecognizePermutePattern),
1894    }
1895}
1896/// An ordered list of sub-patterns (the operands of a concatenation or alternation).
1897#[derive(prost_helpers::AnyPB)]
1898#[derive(Clone, PartialEq, ::prost::Message)]
1899pub struct MatchRecognizePatternSeq {
1900    #[prost(message, repeated, tag = "1")]
1901    pub patterns: ::prost::alloc::vec::Vec<MatchRecognizePatternNode>,
1902}
1903/// `PERMUTE(a, b, ...)` over the listed pattern variables.
1904#[derive(prost_helpers::AnyPB)]
1905#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
1906pub struct MatchRecognizePermutePattern {
1907    #[prost(string, repeated, tag = "1")]
1908    pub vars: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
1909}
1910/// A sub-pattern with a quantifier, e.g. `A+`, `A*?`, `A{1,3}`.
1911#[derive(prost_helpers::AnyPB)]
1912#[derive(Clone, PartialEq, ::prost::Message)]
1913pub struct MatchRecognizeQuantifiedPattern {
1914    #[prost(message, optional, boxed, tag = "1")]
1915    pub inner: ::core::option::Option<
1916        ::prost::alloc::boxed::Box<MatchRecognizePatternNode>,
1917    >,
1918    #[prost(message, optional, tag = "2")]
1919    pub quantifier: ::core::option::Option<MatchRecognizeQuantifier>,
1920    /// A trailing `?` (e.g. `A*?`) — prefer the fewest repetitions.
1921    #[prost(bool, tag = "3")]
1922    pub reluctant: bool,
1923}
1924/// A row-pattern quantifier.
1925#[derive(prost_helpers::AnyPB)]
1926#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
1927pub struct MatchRecognizeQuantifier {
1928    #[prost(enumeration = "match_recognize_quantifier::Kind", tag = "1")]
1929    pub kind: i32,
1930    /// For RANGE: minimum repetitions (defaults to 0).
1931    #[prost(uint32, tag = "2")]
1932    pub min: u32,
1933    /// For RANGE: maximum repetitions; absent means unbounded.
1934    #[prost(uint32, optional, tag = "3")]
1935    pub max: ::core::option::Option<u32>,
1936}
1937/// Nested message and enum types in `MatchRecognizeQuantifier`.
1938pub mod match_recognize_quantifier {
1939    #[derive(prost_helpers::AnyPB)]
1940    #[derive(
1941        Clone,
1942        Copy,
1943        Debug,
1944        PartialEq,
1945        Eq,
1946        Hash,
1947        PartialOrd,
1948        Ord,
1949        ::prost::Enumeration
1950    )]
1951    #[repr(i32)]
1952    pub enum Kind {
1953        Unspecified = 0,
1954        /// `*`
1955        Star = 1,
1956        /// `+`
1957        Plus = 2,
1958        /// `?`
1959        Question = 3,
1960        /// `{min,max}` (and its `{n}`, `{n,}`, `{,m}` forms).
1961        Range = 4,
1962    }
1963    impl Kind {
1964        /// String value of the enum field names used in the ProtoBuf definition.
1965        ///
1966        /// The values are not transformed in any way and thus are considered stable
1967        /// (if the ProtoBuf definition does not change) and safe for programmatic use.
1968        pub fn as_str_name(&self) -> &'static str {
1969            match self {
1970                Self::Unspecified => "KIND_UNSPECIFIED",
1971                Self::Star => "KIND_STAR",
1972                Self::Plus => "KIND_PLUS",
1973                Self::Question => "KIND_QUESTION",
1974                Self::Range => "KIND_RANGE",
1975            }
1976        }
1977        /// Creates an enum from field names used in the ProtoBuf definition.
1978        pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
1979            match value {
1980                "KIND_UNSPECIFIED" => Some(Self::Unspecified),
1981                "KIND_STAR" => Some(Self::Star),
1982                "KIND_PLUS" => Some(Self::Plus),
1983                "KIND_QUESTION" => Some(Self::Question),
1984                "KIND_RANGE" => Some(Self::Range),
1985                _ => None,
1986            }
1987        }
1988    }
1989}
1990/// A DEFINE predicate. The condition is evaluated, per candidate row during matching, over a
1991/// synthetic row whose i-th column is produced by slots\[i\] from the candidate row, its physical
1992/// neighbours, and the in-progress match's labels.
1993#[derive(prost_helpers::AnyPB)]
1994#[derive(Clone, PartialEq, ::prost::Message)]
1995pub struct MatchRecognizeDefine {
1996    #[prost(string, tag = "1")]
1997    pub symbol: ::prost::alloc::string::String,
1998    /// Predicate over the synthetic slot row: an InputRef(i) reads slots\[i\].
1999    #[prost(message, optional, tag = "2")]
2000    pub condition: ::core::option::Option<super::expr::ExprNode>,
2001    #[prost(message, repeated, tag = "3")]
2002    pub slots: ::prost::alloc::vec::Vec<MatchRecognizeDefineSlot>,
2003}
2004/// The AFTER MATCH SKIP clause: where the scan resumes after an emitted match.
2005#[derive(prost_helpers::AnyPB)]
2006#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2007pub struct MatchRecognizeAfterMatchSkip {
2008    #[prost(enumeration = "match_recognize_after_match_skip::Mode", tag = "1")]
2009    pub mode: i32,
2010    /// Target pattern variable, for TO_FIRST / TO_LAST only.
2011    #[prost(string, optional, tag = "2")]
2012    pub target: ::core::option::Option<::prost::alloc::string::String>,
2013}
2014/// Nested message and enum types in `MatchRecognizeAfterMatchSkip`.
2015pub mod match_recognize_after_match_skip {
2016    #[derive(prost_helpers::AnyPB)]
2017    #[derive(
2018        Clone,
2019        Copy,
2020        Debug,
2021        PartialEq,
2022        Eq,
2023        Hash,
2024        PartialOrd,
2025        Ord,
2026        ::prost::Enumeration
2027    )]
2028    #[repr(i32)]
2029    pub enum Mode {
2030        Unspecified = 0,
2031        /// Resume past the match's last row (the default).
2032        PastLastRow = 1,
2033        /// Resume at the row after the match's first row.
2034        ToNextRow = 2,
2035        /// Resume at the first / last row bound to `target` within the match.
2036        ToFirst = 3,
2037        ToLast = 4,
2038    }
2039    impl Mode {
2040        /// String value of the enum field names used in the ProtoBuf definition.
2041        ///
2042        /// The values are not transformed in any way and thus are considered stable
2043        /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2044        pub fn as_str_name(&self) -> &'static str {
2045            match self {
2046                Self::Unspecified => "MODE_UNSPECIFIED",
2047                Self::PastLastRow => "MODE_PAST_LAST_ROW",
2048                Self::ToNextRow => "MODE_TO_NEXT_ROW",
2049                Self::ToFirst => "MODE_TO_FIRST",
2050                Self::ToLast => "MODE_TO_LAST",
2051            }
2052        }
2053        /// Creates an enum from field names used in the ProtoBuf definition.
2054        pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2055            match value {
2056                "MODE_UNSPECIFIED" => Some(Self::Unspecified),
2057                "MODE_PAST_LAST_ROW" => Some(Self::PastLastRow),
2058                "MODE_TO_NEXT_ROW" => Some(Self::ToNextRow),
2059                "MODE_TO_FIRST" => Some(Self::ToFirst),
2060                "MODE_TO_LAST" => Some(Self::ToLast),
2061                _ => None,
2062            }
2063        }
2064    }
2065}
2066/// One input a DEFINE predicate reads.
2067#[derive(prost_helpers::AnyPB)]
2068#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2069pub struct MatchRecognizeDefineSlot {
2070    #[prost(enumeration = "match_recognize_define_slot::Kind", tag = "1")]
2071    pub kind: i32,
2072    #[prost(string, repeated, tag = "2")]
2073    pub vars: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
2074    #[prost(uint32, tag = "3")]
2075    pub col_idx: u32,
2076    #[prost(uint32, tag = "4")]
2077    pub offset: u32,
2078}
2079/// Nested message and enum types in `MatchRecognizeDefineSlot`.
2080pub mod match_recognize_define_slot {
2081    #[derive(prost_helpers::AnyPB)]
2082    #[derive(
2083        Clone,
2084        Copy,
2085        Debug,
2086        PartialEq,
2087        Eq,
2088        Hash,
2089        PartialOrd,
2090        Ord,
2091        ::prost::Enumeration
2092    )]
2093    #[repr(i32)]
2094    pub enum Kind {
2095        Unspecified = 0,
2096        /// The candidate row's own column.
2097        SelfCol = 1,
2098        /// PREV(col, offset) / NEXT(col, offset): a physical-offset row's column.
2099        Prev = 2,
2100        Next = 3,
2101        /// FIRST(var.col) / LAST(var.col): RUNNING navigation — the first / last row of the match so
2102        /// far that belongs to `vars`. The candidate row belongs to `vars` whenever `vars` contains the
2103        /// variable being defined: its label is still tentative, so it is not yet among the match's
2104        /// labels, but the standard counts it (a `var.col` reference *is* RUNNING LAST of that column).
2105        RunningFirst = 4,
2106        RunningLast = 5,
2107    }
2108    impl Kind {
2109        /// String value of the enum field names used in the ProtoBuf definition.
2110        ///
2111        /// The values are not transformed in any way and thus are considered stable
2112        /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2113        pub fn as_str_name(&self) -> &'static str {
2114            match self {
2115                Self::Unspecified => "KIND_UNSPECIFIED",
2116                Self::SelfCol => "KIND_SELF_COL",
2117                Self::Prev => "KIND_PREV",
2118                Self::Next => "KIND_NEXT",
2119                Self::RunningFirst => "KIND_RUNNING_FIRST",
2120                Self::RunningLast => "KIND_RUNNING_LAST",
2121            }
2122        }
2123        /// Creates an enum from field names used in the ProtoBuf definition.
2124        pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2125            match value {
2126                "KIND_UNSPECIFIED" => Some(Self::Unspecified),
2127                "KIND_SELF_COL" => Some(Self::SelfCol),
2128                "KIND_PREV" => Some(Self::Prev),
2129                "KIND_NEXT" => Some(Self::Next),
2130                "KIND_RUNNING_FIRST" => Some(Self::RunningFirst),
2131                "KIND_RUNNING_LAST" => Some(Self::RunningLast),
2132                _ => None,
2133            }
2134        }
2135    }
2136}
2137/// One MEASURES item. The expression is evaluated over a synthetic per-match row whose i-th column is
2138/// produced by `slots\[i\]`; the executor materializes that row once the match and its per-row pattern
2139/// variable labels are known.
2140#[derive(prost_helpers::AnyPB)]
2141#[derive(Clone, PartialEq, ::prost::Message)]
2142pub struct MatchRecognizeMeasure {
2143    /// Expression over the synthetic row: an InputRef(i) reads slots\[i\].
2144    #[prost(message, optional, tag = "1")]
2145    pub expr: ::core::option::Option<super::expr::ExprNode>,
2146    #[prost(string, tag = "2")]
2147    pub name: ::prost::alloc::string::String,
2148    #[prost(message, repeated, tag = "3")]
2149    pub slots: ::prost::alloc::vec::Vec<MatchRecognizeMeasureSlot>,
2150}
2151/// One navigation input of a measure expression.
2152#[derive(prost_helpers::AnyPB)]
2153#[derive(Clone, PartialEq, ::prost::Message)]
2154pub struct MatchRecognizeMeasureSlot {
2155    #[prost(enumeration = "match_recognize_measure_slot::Kind", tag = "1")]
2156    pub kind: i32,
2157    /// The pattern variables this slot navigates over: one for a plain variable, several for a SUBSET
2158    /// union variable. Empty for CLASSIFIER. A row matches if its label is any of these.
2159    #[prost(string, repeated, tag = "2")]
2160    pub vars: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
2161    #[prost(uint32, tag = "3")]
2162    pub col_idx: u32,
2163    #[prost(message, optional, tag = "4")]
2164    pub data_type: ::core::option::Option<super::data::DataType>,
2165    /// The aggregate to run for KIND_SUM, over a single input column (the projected `col_idx`) of
2166    /// the rows labeled by `vars`. Its single arg is an InputRef to column 0.
2167    #[prost(message, optional, tag = "5")]
2168    pub agg_call: ::core::option::Option<super::expr::AggCall>,
2169}
2170/// Nested message and enum types in `MatchRecognizeMeasureSlot`.
2171pub mod match_recognize_measure_slot {
2172    #[derive(prost_helpers::AnyPB)]
2173    #[derive(
2174        Clone,
2175        Copy,
2176        Debug,
2177        PartialEq,
2178        Eq,
2179        Hash,
2180        PartialOrd,
2181        Ord,
2182        ::prost::Enumeration
2183    )]
2184    #[repr(i32)]
2185    pub enum Kind {
2186        Unspecified = 0,
2187        /// LAST(var.col): column value of the last row labeled `var` in the match (also the meaning of
2188        /// a bare `var.col` under ONE ROW PER MATCH FINAL semantics).
2189        Last = 1,
2190        /// FIRST(var.col): column value of the first such row.
2191        First = 2,
2192        /// CLASSIFIER(): the pattern variable bound to the match's last row (`var`/`col_idx` unused).
2193        Classifier = 3,
2194        /// COUNT(\*): number of rows in the match (`var`/`col_idx` unused).
2195        CountStar = 4,
2196        /// COUNT(var.col): number of rows labeled `var` whose `col` is non-null.
2197        Count = 5,
2198        /// MIN(var.col) / MAX(var.col): min / max of `col` over rows labeled `var`.
2199        Min = 6,
2200        Max = 7,
2201        /// SUM(var.col): evaluated by `agg_call` over the rows labeled `var`. (AVG is lowered by the
2202        /// planner into a SUM slot plus a COUNT slot and a division expression.)
2203        Sum = 8,
2204    }
2205    impl Kind {
2206        /// String value of the enum field names used in the ProtoBuf definition.
2207        ///
2208        /// The values are not transformed in any way and thus are considered stable
2209        /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2210        pub fn as_str_name(&self) -> &'static str {
2211            match self {
2212                Self::Unspecified => "KIND_UNSPECIFIED",
2213                Self::Last => "KIND_LAST",
2214                Self::First => "KIND_FIRST",
2215                Self::Classifier => "KIND_CLASSIFIER",
2216                Self::CountStar => "KIND_COUNT_STAR",
2217                Self::Count => "KIND_COUNT",
2218                Self::Min => "KIND_MIN",
2219                Self::Max => "KIND_MAX",
2220                Self::Sum => "KIND_SUM",
2221            }
2222        }
2223        /// Creates an enum from field names used in the ProtoBuf definition.
2224        pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2225            match value {
2226                "KIND_UNSPECIFIED" => Some(Self::Unspecified),
2227                "KIND_LAST" => Some(Self::Last),
2228                "KIND_FIRST" => Some(Self::First),
2229                "KIND_CLASSIFIER" => Some(Self::Classifier),
2230                "KIND_COUNT_STAR" => Some(Self::CountStar),
2231                "KIND_COUNT" => Some(Self::Count),
2232                "KIND_MIN" => Some(Self::Min),
2233                "KIND_MAX" => Some(Self::Max),
2234                "KIND_SUM" => Some(Self::Sum),
2235                _ => None,
2236            }
2237        }
2238    }
2239}
2240#[derive(prost_helpers::AnyPB)]
2241#[derive(Clone, PartialEq, ::prost::Message)]
2242pub struct StreamNode {
2243    /// The id for the operator. This is local per mview.
2244    /// TODO: should better be a uint32.
2245    #[prost(uint64, tag = "1", wrapper = "crate::id::StreamNodeLocalOperatorId")]
2246    pub operator_id: crate::id::StreamNodeLocalOperatorId,
2247    /// Child node in plan aka. upstream nodes in the streaming DAG
2248    #[prost(message, repeated, tag = "3")]
2249    pub input: ::prost::alloc::vec::Vec<StreamNode>,
2250    #[prost(uint32, repeated, tag = "2")]
2251    pub stream_key: ::prost::alloc::vec::Vec<u32>,
2252    #[prost(enumeration = "stream_node::StreamKind", tag = "24")]
2253    pub stream_kind: i32,
2254    #[prost(string, tag = "18")]
2255    pub identity: ::prost::alloc::string::String,
2256    /// The schema of the plan node
2257    #[prost(message, repeated, tag = "19")]
2258    pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2259    #[prost(
2260        oneof = "stream_node::NodeBody",
2261        tags = "100, 101, 102, 103, 104, 105, 106, 107, 108, 109, 110, 111, 112, 113, 114, 115, 116, 117, 118, 119, 120, 121, 122, 123, 124, 125, 126, 127, 128, 129, 130, 131, 132, 133, 134, 135, 136, 137, 138, 139, 140, 142, 143, 144, 145, 146, 147, 148, 149, 150, 151, 152, 153, 154, 155, 156, 157, 158, 159"
2262    )]
2263    pub node_body: ::core::option::Option<stream_node::NodeBody>,
2264}
2265/// Nested message and enum types in `StreamNode`.
2266pub mod stream_node {
2267    /// This field used to be a `bool append_only`.
2268    /// Enum variants are ordered for backwards compatibility.
2269    #[derive(prost_helpers::AnyPB)]
2270    #[derive(
2271        Clone,
2272        Copy,
2273        Debug,
2274        PartialEq,
2275        Eq,
2276        Hash,
2277        PartialOrd,
2278        Ord,
2279        ::prost::Enumeration
2280    )]
2281    #[repr(i32)]
2282    pub enum StreamKind {
2283        /// buf:lint:ignore ENUM_ZERO_VALUE_SUFFIX
2284        Retract = 0,
2285        AppendOnly = 1,
2286        Upsert = 2,
2287    }
2288    impl StreamKind {
2289        /// String value of the enum field names used in the ProtoBuf definition.
2290        ///
2291        /// The values are not transformed in any way and thus are considered stable
2292        /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2293        pub fn as_str_name(&self) -> &'static str {
2294            match self {
2295                Self::Retract => "STREAM_KIND_RETRACT",
2296                Self::AppendOnly => "STREAM_KIND_APPEND_ONLY",
2297                Self::Upsert => "STREAM_KIND_UPSERT",
2298            }
2299        }
2300        /// Creates an enum from field names used in the ProtoBuf definition.
2301        pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2302            match value {
2303                "STREAM_KIND_RETRACT" => Some(Self::Retract),
2304                "STREAM_KIND_APPEND_ONLY" => Some(Self::AppendOnly),
2305                "STREAM_KIND_UPSERT" => Some(Self::Upsert),
2306                _ => None,
2307            }
2308        }
2309    }
2310    #[derive(prost_helpers::AnyPB)]
2311    #[derive(::enum_as_inner::EnumAsInner, ::strum::Display, ::strum::EnumDiscriminants)]
2312    #[derive(::prost_helpers::StreamNodeBodyVariants)]
2313    #[strum_discriminants(derive(::strum::Display, Hash))]
2314    #[derive(Clone, PartialEq, ::prost::Oneof)]
2315    pub enum NodeBody {
2316        #[prost(message, tag = "100")]
2317        Source(::prost::alloc::boxed::Box<super::SourceNode>),
2318        #[prost(message, tag = "101")]
2319        Project(::prost::alloc::boxed::Box<super::ProjectNode>),
2320        #[prost(message, tag = "102")]
2321        Filter(::prost::alloc::boxed::Box<super::FilterNode>),
2322        #[prost(message, tag = "103")]
2323        Materialize(::prost::alloc::boxed::Box<super::MaterializeNode>),
2324        #[prost(message, tag = "104")]
2325        StatelessSimpleAgg(::prost::alloc::boxed::Box<super::SimpleAggNode>),
2326        #[prost(message, tag = "105")]
2327        SimpleAgg(::prost::alloc::boxed::Box<super::SimpleAggNode>),
2328        #[prost(message, tag = "106")]
2329        HashAgg(::prost::alloc::boxed::Box<super::HashAggNode>),
2330        #[prost(message, tag = "107")]
2331        AppendOnlyTopN(::prost::alloc::boxed::Box<super::TopNNode>),
2332        #[prost(message, tag = "108")]
2333        HashJoin(::prost::alloc::boxed::Box<super::HashJoinNode>),
2334        #[prost(message, tag = "109")]
2335        TopN(::prost::alloc::boxed::Box<super::TopNNode>),
2336        #[prost(message, tag = "110")]
2337        HopWindow(::prost::alloc::boxed::Box<super::HopWindowNode>),
2338        #[prost(message, tag = "111")]
2339        Merge(::prost::alloc::boxed::Box<super::MergeNode>),
2340        #[prost(message, tag = "112")]
2341        Exchange(::prost::alloc::boxed::Box<super::ExchangeNode>),
2342        #[prost(message, tag = "113")]
2343        StreamScan(::prost::alloc::boxed::Box<super::StreamScanNode>),
2344        #[prost(message, tag = "114")]
2345        BatchPlan(::prost::alloc::boxed::Box<super::BatchPlanNode>),
2346        #[prost(message, tag = "115")]
2347        Lookup(::prost::alloc::boxed::Box<super::LookupNode>),
2348        #[prost(message, tag = "116")]
2349        Arrange(::prost::alloc::boxed::Box<super::ArrangeNode>),
2350        #[prost(message, tag = "117")]
2351        LookupUnion(::prost::alloc::boxed::Box<super::LookupUnionNode>),
2352        #[prost(message, tag = "118")]
2353        Union(super::UnionNode),
2354        #[prost(message, tag = "119")]
2355        DeltaIndexJoin(::prost::alloc::boxed::Box<super::DeltaIndexJoinNode>),
2356        #[prost(message, tag = "120")]
2357        Sink(::prost::alloc::boxed::Box<super::SinkNode>),
2358        #[prost(message, tag = "121")]
2359        Expand(::prost::alloc::boxed::Box<super::ExpandNode>),
2360        #[prost(message, tag = "122")]
2361        DynamicFilter(::prost::alloc::boxed::Box<super::DynamicFilterNode>),
2362        #[prost(message, tag = "123")]
2363        ProjectSet(::prost::alloc::boxed::Box<super::ProjectSetNode>),
2364        #[prost(message, tag = "124")]
2365        GroupTopN(::prost::alloc::boxed::Box<super::GroupTopNNode>),
2366        #[prost(message, tag = "125")]
2367        Sort(::prost::alloc::boxed::Box<super::SortNode>),
2368        #[prost(message, tag = "126")]
2369        WatermarkFilter(::prost::alloc::boxed::Box<super::WatermarkFilterNode>),
2370        #[prost(message, tag = "127")]
2371        Dml(::prost::alloc::boxed::Box<super::DmlNode>),
2372        #[prost(message, tag = "128")]
2373        RowIdGen(::prost::alloc::boxed::Box<super::RowIdGenNode>),
2374        #[prost(message, tag = "129")]
2375        Now(::prost::alloc::boxed::Box<super::NowNode>),
2376        #[prost(message, tag = "130")]
2377        AppendOnlyGroupTopN(::prost::alloc::boxed::Box<super::GroupTopNNode>),
2378        #[prost(message, tag = "131")]
2379        TemporalJoin(::prost::alloc::boxed::Box<super::TemporalJoinNode>),
2380        #[prost(message, tag = "132")]
2381        BarrierRecv(::prost::alloc::boxed::Box<super::BarrierRecvNode>),
2382        #[prost(message, tag = "133")]
2383        Values(::prost::alloc::boxed::Box<super::ValuesNode>),
2384        #[prost(message, tag = "134")]
2385        AppendOnlyDedup(::prost::alloc::boxed::Box<super::DedupNode>),
2386        #[prost(message, tag = "135")]
2387        NoOp(super::NoOpNode),
2388        #[prost(message, tag = "136")]
2389        EowcOverWindow(::prost::alloc::boxed::Box<super::EowcOverWindowNode>),
2390        #[prost(message, tag = "137")]
2391        OverWindow(::prost::alloc::boxed::Box<super::OverWindowNode>),
2392        #[prost(message, tag = "138")]
2393        StreamFsFetch(::prost::alloc::boxed::Box<super::StreamFsFetchNode>),
2394        #[prost(message, tag = "139")]
2395        StreamCdcScan(::prost::alloc::boxed::Box<super::StreamCdcScanNode>),
2396        #[prost(message, tag = "140")]
2397        CdcFilter(::prost::alloc::boxed::Box<super::CdcFilterNode>),
2398        #[prost(message, tag = "142")]
2399        SourceBackfill(::prost::alloc::boxed::Box<super::SourceBackfillNode>),
2400        #[prost(message, tag = "143")]
2401        Changelog(::prost::alloc::boxed::Box<super::ChangeLogNode>),
2402        #[prost(message, tag = "144")]
2403        LocalApproxPercentile(
2404            ::prost::alloc::boxed::Box<super::LocalApproxPercentileNode>,
2405        ),
2406        #[prost(message, tag = "145")]
2407        GlobalApproxPercentile(
2408            ::prost::alloc::boxed::Box<super::GlobalApproxPercentileNode>,
2409        ),
2410        #[prost(message, tag = "146")]
2411        RowMerge(::prost::alloc::boxed::Box<super::RowMergeNode>),
2412        #[prost(message, tag = "147")]
2413        AsOfJoin(::prost::alloc::boxed::Box<super::AsOfJoinNode>),
2414        #[prost(message, tag = "148")]
2415        SyncLogStore(::prost::alloc::boxed::Box<super::SyncLogStoreNode>),
2416        #[prost(message, tag = "149")]
2417        MaterializedExprs(::prost::alloc::boxed::Box<super::MaterializedExprsNode>),
2418        #[prost(message, tag = "150")]
2419        VectorIndexWrite(::prost::alloc::boxed::Box<super::VectorIndexWriteNode>),
2420        #[prost(message, tag = "151")]
2421        UpstreamSinkUnion(::prost::alloc::boxed::Box<super::UpstreamSinkUnionNode>),
2422        #[prost(message, tag = "152")]
2423        LocalityProvider(::prost::alloc::boxed::Box<super::LocalityProviderNode>),
2424        #[prost(message, tag = "153")]
2425        EowcGapFill(::prost::alloc::boxed::Box<super::EowcGapFillNode>),
2426        #[prost(message, tag = "154")]
2427        GapFill(::prost::alloc::boxed::Box<super::GapFillNode>),
2428        #[prost(message, tag = "155")]
2429        VectorIndexLookupJoin(
2430            ::prost::alloc::boxed::Box<super::VectorIndexLookupJoinNode>,
2431        ),
2432        #[prost(message, tag = "156")]
2433        IcebergWithPkIndexWriter(
2434            ::prost::alloc::boxed::Box<super::IcebergWithPkIndexWriterNode>,
2435        ),
2436        #[prost(message, tag = "157")]
2437        IcebergWithPkIndexPositionDeleteMerger(
2438            ::prost::alloc::boxed::Box<super::IcebergWithPkIndexPositionDeleteMergerNode>,
2439        ),
2440        #[prost(message, tag = "158")]
2441        MatchRecognize(::prost::alloc::boxed::Box<super::MatchRecognizeNode>),
2442        #[prost(message, tag = "159")]
2443        CompactionResolver(::prost::alloc::boxed::Box<super::CompactionResolverNode>),
2444    }
2445}
2446/// The method to map the upstream columns in the dispatcher before dispatching.
2447///
2448/// * For intra-job exchange, typically the upstream and downstream columns are the same. `indices`
2449///   will be `0..len` and `types` will be empty.
2450///
2451/// * For inter-job exchange,
2452///
2453///   * if the downstream only requires a subset of the upstream columns, `indices` will be the
2454///     indices of the required columns in the upstream columns.
2455///   * if some columns are added to the upstream, `indices` will help to maintain the same schema
2456///     from the view of the downstream.
2457///   * if some columns are altered to different (composite) types, `types` will be used to convert
2458///     the upstream columns to the downstream columns to maintain the same schema.
2459#[derive(prost_helpers::AnyPB)]
2460#[derive(Clone, PartialEq, ::prost::Message)]
2461pub struct DispatchOutputMapping {
2462    /// Indices of the columns to output.
2463    #[prost(uint32, repeated, tag = "1")]
2464    pub indices: ::prost::alloc::vec::Vec<u32>,
2465    /// Besides the indices, we may also need to convert the types of some columns.
2466    ///
2467    /// * If no type conversion is needed, this field will be empty.
2468    /// * If type conversion is needed, this will have the same length as `indices`. Only columns with
2469    ///   type conversion will have `upstream` and `downstream` field set.
2470    #[prost(message, repeated, tag = "2")]
2471    pub types: ::prost::alloc::vec::Vec<dispatch_output_mapping::TypePair>,
2472}
2473/// Nested message and enum types in `DispatchOutputMapping`.
2474pub mod dispatch_output_mapping {
2475    #[derive(prost_helpers::AnyPB)]
2476    #[derive(Clone, PartialEq, ::prost::Message)]
2477    pub struct TypePair {
2478        #[prost(message, optional, tag = "1")]
2479        pub upstream: ::core::option::Option<super::super::data::DataType>,
2480        #[prost(message, optional, tag = "2")]
2481        pub downstream: ::core::option::Option<super::super::data::DataType>,
2482    }
2483}
2484/// The property of an edge in the fragment graph.
2485/// This is essientially a "logical" version of `Dispatcher`. See the doc of `Dispatcher` for more details.
2486#[derive(prost_helpers::AnyPB)]
2487#[derive(Clone, PartialEq, ::prost::Message)]
2488pub struct DispatchStrategy {
2489    #[prost(enumeration = "DispatcherType", tag = "1")]
2490    pub r#type: i32,
2491    #[prost(uint32, repeated, tag = "2")]
2492    pub dist_key_indices: ::prost::alloc::vec::Vec<u32>,
2493    #[prost(message, optional, tag = "3")]
2494    pub output_mapping: ::core::option::Option<DispatchOutputMapping>,
2495}
2496/// A dispatcher redistribute messages.
2497/// We encode both the type and other usage information in the proto.
2498#[derive(prost_helpers::AnyPB)]
2499#[derive(Clone, PartialEq, ::prost::Message)]
2500pub struct Dispatcher {
2501    #[prost(enumeration = "DispatcherType", tag = "1")]
2502    pub r#type: i32,
2503    /// Indices of the columns to be used for hashing.
2504    /// For dispatcher types other than HASH, this is ignored.
2505    #[prost(uint32, repeated, tag = "2")]
2506    pub dist_key_indices: ::prost::alloc::vec::Vec<u32>,
2507    /// The method to map the upstream columns in the dispatcher before dispatching.
2508    #[prost(message, optional, tag = "6")]
2509    pub output_mapping: ::core::option::Option<DispatchOutputMapping>,
2510    /// The hash mapping for consistent hash.
2511    /// For dispatcher types other than HASH, this is ignored.
2512    #[prost(message, optional, tag = "3")]
2513    pub hash_mapping: ::core::option::Option<ActorMapping>,
2514    /// Dispatcher can be uniquely identified by a combination of actor id and dispatcher id.
2515    /// This is exactly the same as its downstream fragment id.
2516    #[prost(uint32, tag = "4", wrapper = "crate::id::FragmentId")]
2517    pub dispatcher_id: crate::id::FragmentId,
2518    /// Number of downstreams decides how many endpoints a dispatcher should dispatch.
2519    #[prost(uint32, repeated, tag = "5", wrapper = "crate::id::ActorId")]
2520    pub downstream_actor_id: ::prost::alloc::vec::Vec<crate::id::ActorId>,
2521}
2522/// A StreamActor is a running fragment of the overall stream graph,
2523#[derive(prost_helpers::AnyPB)]
2524#[derive(Clone, PartialEq, ::prost::Message)]
2525pub struct StreamActor {
2526    #[prost(uint32, tag = "1", wrapper = "crate::id::ActorId")]
2527    pub actor_id: crate::id::ActorId,
2528    #[prost(uint32, tag = "2", wrapper = "crate::id::FragmentId")]
2529    pub fragment_id: crate::id::FragmentId,
2530    #[prost(message, repeated, tag = "4")]
2531    pub dispatcher: ::prost::alloc::vec::Vec<Dispatcher>,
2532    /// Vnodes that the executors in this actor own.
2533    /// If the fragment is a singleton, this field will not be set and leave a `None`.
2534    #[prost(message, optional, tag = "8")]
2535    pub vnode_bitmap: ::core::option::Option<super::common::Buffer>,
2536    /// The SQL definition of this materialized view. Used for debugging only.
2537    #[prost(string, tag = "9")]
2538    pub mview_definition: ::prost::alloc::string::String,
2539    /// Provide the necessary context, e.g. session info like time zone, for the actor.
2540    #[prost(message, optional, tag = "10")]
2541    pub expr_context: ::core::option::Option<super::plan_common::ExprContext>,
2542    /// The config override for this actor.
2543    #[prost(string, tag = "11")]
2544    pub config_override: ::prost::alloc::string::String,
2545}
2546/// The streaming context associated with a stream plan
2547#[derive(prost_helpers::AnyPB)]
2548#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2549pub struct StreamContext {
2550    /// The timezone associated with the streaming plan. Only applies to MV for now.
2551    #[prost(string, tag = "1")]
2552    pub timezone: ::prost::alloc::string::String,
2553    /// The partial config of this job to override the global config.
2554    #[prost(string, tag = "2")]
2555    pub config_override: ::prost::alloc::string::String,
2556}
2557#[derive(prost_helpers::AnyPB)]
2558#[derive(Clone, PartialEq, ::prost::Message)]
2559pub struct BackfillOrder {
2560    #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::RelationId")]
2561    pub order: ::std::collections::HashMap<
2562        crate::id::RelationId,
2563        super::common::Uint32Vector,
2564    >,
2565}
2566/// Representation of a graph of stream fragments.
2567/// Generated by the fragmenter in the frontend, only used in DDL requests and never persisted.
2568///
2569/// For the persisted form, see `TableFragments`.
2570#[derive(prost_helpers::AnyPB)]
2571#[derive(Clone, PartialEq, ::prost::Message)]
2572pub struct StreamFragmentGraph {
2573    /// all the fragments in the graph.
2574    #[prost(map = "uint32, message", tag = "1", wrapper = "crate::id::FragmentId")]
2575    pub fragments: ::std::collections::HashMap<
2576        crate::id::FragmentId,
2577        stream_fragment_graph::StreamFragment,
2578    >,
2579    /// edges between fragments.
2580    #[prost(message, repeated, tag = "2")]
2581    pub edges: ::prost::alloc::vec::Vec<stream_fragment_graph::StreamFragmentEdge>,
2582    #[prost(uint32, repeated, tag = "3", wrapper = "crate::id::TableId")]
2583    pub dependent_table_ids: ::prost::alloc::vec::Vec<crate::id::TableId>,
2584    #[prost(uint32, tag = "4")]
2585    pub table_ids_cnt: u32,
2586    #[prost(message, optional, tag = "5")]
2587    pub ctx: ::core::option::Option<StreamContext>,
2588    /// If none, default parallelism will be applied.
2589    #[prost(message, optional, tag = "6")]
2590    pub parallelism: ::core::option::Option<stream_fragment_graph::Parallelism>,
2591    /// Parallelism to use during backfill. Falls back to `parallelism` if unset.
2592    #[prost(message, optional, tag = "9")]
2593    pub backfill_parallelism: ::core::option::Option<stream_fragment_graph::Parallelism>,
2594    /// The adaptive parallelism strategy for this streaming job, if explicitly set.
2595    #[prost(string, tag = "10")]
2596    pub adaptive_parallelism_strategy: ::prost::alloc::string::String,
2597    /// The adaptive parallelism strategy for the backfill override, if explicitly set.
2598    #[prost(string, tag = "11")]
2599    pub backfill_adaptive_parallelism_strategy: ::prost::alloc::string::String,
2600    /// Specified max parallelism, i.e., expected vnode count for the graph.
2601    ///
2602    /// The scheduler on the meta service will use this as a hint to decide the vnode count
2603    /// for each fragment.
2604    ///
2605    /// Note that the actual vnode count may be different from this value.
2606    /// For example, a no-shuffle exchange between current fragment graph and an existing
2607    /// upstream fragment graph requires two fragments to be in the same distribution,
2608    /// thus the same vnode count.
2609    #[prost(uint32, tag = "7")]
2610    pub max_parallelism: u32,
2611    /// The backfill order strategy for the fragments.
2612    #[prost(message, optional, tag = "8")]
2613    pub backfill_order: ::core::option::Option<BackfillOrder>,
2614}
2615/// Nested message and enum types in `StreamFragmentGraph`.
2616pub mod stream_fragment_graph {
2617    #[derive(prost_helpers::AnyPB)]
2618    #[derive(Clone, PartialEq, ::prost::Message)]
2619    pub struct StreamFragment {
2620        /// 0-based on frontend, and will be rewritten to global id on meta.
2621        #[prost(uint32, tag = "1", wrapper = "crate::id::FragmentId")]
2622        pub fragment_id: crate::id::FragmentId,
2623        /// root stream node in this fragment.
2624        #[prost(message, optional, tag = "2")]
2625        pub node: ::core::option::Option<super::StreamNode>,
2626        /// Bitwise-OR of `FragmentTypeFlag`s
2627        #[prost(uint32, tag = "3")]
2628        pub fragment_type_mask: u32,
2629        /// Mark whether this fragment requires exactly one actor.
2630        /// Note: if this is `false`, the fragment may still be a singleton according to the scheduler.
2631        /// One should check `meta.Fragment.distribution_type` for the final result.
2632        #[prost(bool, tag = "4")]
2633        pub requires_singleton: bool,
2634    }
2635    #[derive(prost_helpers::AnyPB)]
2636    #[derive(Clone, PartialEq, ::prost::Message)]
2637    pub struct StreamFragmentEdge {
2638        /// Dispatch strategy for the fragment.
2639        #[prost(message, optional, tag = "1")]
2640        pub dispatch_strategy: ::core::option::Option<super::DispatchStrategy>,
2641        /// A unique identifier of this edge. Generally it should be exchange node's operator id. When
2642        /// rewriting fragments into delta joins or when inserting 1-to-1 exchange, there will be
2643        /// virtual links generated.
2644        #[prost(uint64, tag = "3")]
2645        pub link_id: u64,
2646        #[prost(uint32, tag = "4", wrapper = "crate::id::FragmentId")]
2647        pub upstream_id: crate::id::FragmentId,
2648        #[prost(uint32, tag = "5", wrapper = "crate::id::FragmentId")]
2649        pub downstream_id: crate::id::FragmentId,
2650    }
2651    #[derive(prost_helpers::AnyPB)]
2652    #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
2653    pub struct Parallelism {
2654        #[prost(uint64, tag = "1")]
2655        pub parallelism: u64,
2656    }
2657}
2658/// Schema change operation for sink
2659#[derive(prost_helpers::AnyPB)]
2660#[derive(Clone, PartialEq, ::prost::Message)]
2661pub struct SinkSchemaChange {
2662    /// Original schema before this change.
2663    /// Used for validation and conflict detection.
2664    #[prost(message, repeated, tag = "1")]
2665    pub original_schema: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2666    /// Schema change operation (mutually exclusive)
2667    #[prost(oneof = "sink_schema_change::Op", tags = "2, 3")]
2668    pub op: ::core::option::Option<sink_schema_change::Op>,
2669}
2670/// Nested message and enum types in `SinkSchemaChange`.
2671pub mod sink_schema_change {
2672    /// Schema change operation (mutually exclusive)
2673    #[derive(prost_helpers::AnyPB)]
2674    #[derive(Clone, PartialEq, ::prost::Oneof)]
2675    pub enum Op {
2676        /// Add new columns to the schema
2677        #[prost(message, tag = "2")]
2678        AddColumns(super::SinkAddColumnsOp),
2679        /// Drop columns from the schema
2680        #[prost(message, tag = "3")]
2681        DropColumns(super::SinkDropColumnsOp),
2682    }
2683}
2684/// Add columns operation
2685#[derive(prost_helpers::AnyPB)]
2686#[derive(Clone, PartialEq, ::prost::Message)]
2687pub struct SinkAddColumnsOp {
2688    /// Columns to add
2689    #[prost(message, repeated, tag = "1")]
2690    pub fields: ::prost::alloc::vec::Vec<super::plan_common::Field>,
2691}
2692/// Drop columns operation
2693#[derive(prost_helpers::AnyPB)]
2694#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
2695pub struct SinkDropColumnsOp {
2696    /// Column names to drop
2697    #[prost(string, repeated, tag = "1")]
2698    pub column_names: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
2699}
2700#[derive(prost_helpers::AnyPB)]
2701#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2702#[repr(i32)]
2703pub enum SinkLogStoreType {
2704    /// / Default value is the normal in memory log store to be backward compatible with the previously unset value
2705    Unspecified = 0,
2706    KvLogStore = 1,
2707    InMemoryLogStore = 2,
2708}
2709impl SinkLogStoreType {
2710    /// String value of the enum field names used in the ProtoBuf definition.
2711    ///
2712    /// The values are not transformed in any way and thus are considered stable
2713    /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2714    pub fn as_str_name(&self) -> &'static str {
2715        match self {
2716            Self::Unspecified => "SINK_LOG_STORE_TYPE_UNSPECIFIED",
2717            Self::KvLogStore => "SINK_LOG_STORE_TYPE_KV_LOG_STORE",
2718            Self::InMemoryLogStore => "SINK_LOG_STORE_TYPE_IN_MEMORY_LOG_STORE",
2719        }
2720    }
2721    /// Creates an enum from field names used in the ProtoBuf definition.
2722    pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2723        match value {
2724            "SINK_LOG_STORE_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2725            "SINK_LOG_STORE_TYPE_KV_LOG_STORE" => Some(Self::KvLogStore),
2726            "SINK_LOG_STORE_TYPE_IN_MEMORY_LOG_STORE" => Some(Self::InMemoryLogStore),
2727            _ => None,
2728        }
2729    }
2730}
2731#[derive(prost_helpers::AnyPB)]
2732#[derive(prost_helpers::Version)]
2733#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2734#[repr(i32)]
2735pub enum AggNodeVersion {
2736    Unspecified = 0,
2737    /// <https://github.com/risingwavelabs/risingwave/issues/12140#issuecomment-1776289808>
2738    Issue12140 = 1,
2739    /// <https://github.com/risingwavelabs/risingwave/issues/13465#issuecomment-1821016508>
2740    Issue13465 = 2,
2741}
2742impl AggNodeVersion {
2743    /// String value of the enum field names used in the ProtoBuf definition.
2744    ///
2745    /// The values are not transformed in any way and thus are considered stable
2746    /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2747    pub fn as_str_name(&self) -> &'static str {
2748        match self {
2749            Self::Unspecified => "AGG_NODE_VERSION_UNSPECIFIED",
2750            Self::Issue12140 => "AGG_NODE_VERSION_ISSUE_12140",
2751            Self::Issue13465 => "AGG_NODE_VERSION_ISSUE_13465",
2752        }
2753    }
2754    /// Creates an enum from field names used in the ProtoBuf definition.
2755    pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2756        match value {
2757            "AGG_NODE_VERSION_UNSPECIFIED" => Some(Self::Unspecified),
2758            "AGG_NODE_VERSION_ISSUE_12140" => Some(Self::Issue12140),
2759            "AGG_NODE_VERSION_ISSUE_13465" => Some(Self::Issue13465),
2760            _ => None,
2761        }
2762    }
2763}
2764#[derive(prost_helpers::AnyPB)]
2765#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2766#[repr(i32)]
2767pub enum InequalityType {
2768    Unspecified = 0,
2769    LessThan = 1,
2770    LessThanOrEqual = 2,
2771    GreaterThan = 3,
2772    GreaterThanOrEqual = 4,
2773}
2774impl InequalityType {
2775    /// String value of the enum field names used in the ProtoBuf definition.
2776    ///
2777    /// The values are not transformed in any way and thus are considered stable
2778    /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2779    pub fn as_str_name(&self) -> &'static str {
2780        match self {
2781            Self::Unspecified => "INEQUALITY_TYPE_UNSPECIFIED",
2782            Self::LessThan => "INEQUALITY_TYPE_LESS_THAN",
2783            Self::LessThanOrEqual => "INEQUALITY_TYPE_LESS_THAN_OR_EQUAL",
2784            Self::GreaterThan => "INEQUALITY_TYPE_GREATER_THAN",
2785            Self::GreaterThanOrEqual => "INEQUALITY_TYPE_GREATER_THAN_OR_EQUAL",
2786        }
2787    }
2788    /// Creates an enum from field names used in the ProtoBuf definition.
2789    pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2790        match value {
2791            "INEQUALITY_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2792            "INEQUALITY_TYPE_LESS_THAN" => Some(Self::LessThan),
2793            "INEQUALITY_TYPE_LESS_THAN_OR_EQUAL" => Some(Self::LessThanOrEqual),
2794            "INEQUALITY_TYPE_GREATER_THAN" => Some(Self::GreaterThan),
2795            "INEQUALITY_TYPE_GREATER_THAN_OR_EQUAL" => Some(Self::GreaterThanOrEqual),
2796            _ => None,
2797        }
2798    }
2799}
2800#[derive(prost_helpers::AnyPB)]
2801#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2802#[repr(i32)]
2803pub enum JoinEncodingType {
2804    Unspecified = 0,
2805    MemoryOptimized = 1,
2806    CpuOptimized = 2,
2807}
2808impl JoinEncodingType {
2809    /// String value of the enum field names used in the ProtoBuf definition.
2810    ///
2811    /// The values are not transformed in any way and thus are considered stable
2812    /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2813    pub fn as_str_name(&self) -> &'static str {
2814        match self {
2815            Self::Unspecified => "UNSPECIFIED",
2816            Self::MemoryOptimized => "MEMORY_OPTIMIZED",
2817            Self::CpuOptimized => "CPU_OPTIMIZED",
2818        }
2819    }
2820    /// Creates an enum from field names used in the ProtoBuf definition.
2821    pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2822        match value {
2823            "UNSPECIFIED" => Some(Self::Unspecified),
2824            "MEMORY_OPTIMIZED" => Some(Self::MemoryOptimized),
2825            "CPU_OPTIMIZED" => Some(Self::CpuOptimized),
2826            _ => None,
2827        }
2828    }
2829}
2830/// Decides which kind of Executor will be used
2831#[derive(prost_helpers::AnyPB)]
2832#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2833#[repr(i32)]
2834pub enum StreamScanType {
2835    Unspecified = 0,
2836    /// ChainExecutor. Deprecated: only for existing persisted legacy jobs.
2837    #[deprecated]
2838    Chain = 1,
2839    /// RearrangedChainExecutor. Deprecated: only for existing persisted legacy jobs.
2840    #[deprecated]
2841    Rearrange = 2,
2842    /// BackfillExecutor. Deprecated: only for existing persisted legacy jobs.
2843    #[deprecated]
2844    Backfill = 3,
2845    /// ChainExecutor with upstream_only = true
2846    UpstreamOnly = 4,
2847    /// ArrangementBackfillExecutor
2848    ArrangementBackfill = 5,
2849    /// SnapshotBackfillExecutor
2850    SnapshotBackfill = 6,
2851    /// SnapshotBackfillExecutor
2852    CrossDbSnapshotBackfill = 7,
2853}
2854impl StreamScanType {
2855    /// String value of the enum field names used in the ProtoBuf definition.
2856    ///
2857    /// The values are not transformed in any way and thus are considered stable
2858    /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2859    pub fn as_str_name(&self) -> &'static str {
2860        match self {
2861            Self::Unspecified => "STREAM_SCAN_TYPE_UNSPECIFIED",
2862            #[allow(deprecated)]
2863            Self::Chain => "STREAM_SCAN_TYPE_CHAIN",
2864            #[allow(deprecated)]
2865            Self::Rearrange => "STREAM_SCAN_TYPE_REARRANGE",
2866            #[allow(deprecated)]
2867            Self::Backfill => "STREAM_SCAN_TYPE_BACKFILL",
2868            Self::UpstreamOnly => "STREAM_SCAN_TYPE_UPSTREAM_ONLY",
2869            Self::ArrangementBackfill => "STREAM_SCAN_TYPE_ARRANGEMENT_BACKFILL",
2870            Self::SnapshotBackfill => "STREAM_SCAN_TYPE_SNAPSHOT_BACKFILL",
2871            Self::CrossDbSnapshotBackfill => {
2872                "STREAM_SCAN_TYPE_CROSS_DB_SNAPSHOT_BACKFILL"
2873            }
2874        }
2875    }
2876    /// Creates an enum from field names used in the ProtoBuf definition.
2877    pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2878        match value {
2879            "STREAM_SCAN_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2880            "STREAM_SCAN_TYPE_CHAIN" => Some(#[allow(deprecated)] Self::Chain),
2881            "STREAM_SCAN_TYPE_REARRANGE" => Some(#[allow(deprecated)] Self::Rearrange),
2882            "STREAM_SCAN_TYPE_BACKFILL" => Some(#[allow(deprecated)] Self::Backfill),
2883            "STREAM_SCAN_TYPE_UPSTREAM_ONLY" => Some(Self::UpstreamOnly),
2884            "STREAM_SCAN_TYPE_ARRANGEMENT_BACKFILL" => Some(Self::ArrangementBackfill),
2885            "STREAM_SCAN_TYPE_SNAPSHOT_BACKFILL" => Some(Self::SnapshotBackfill),
2886            "STREAM_SCAN_TYPE_CROSS_DB_SNAPSHOT_BACKFILL" => {
2887                Some(Self::CrossDbSnapshotBackfill)
2888            }
2889            _ => None,
2890        }
2891    }
2892}
2893#[derive(prost_helpers::AnyPB)]
2894#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2895#[repr(i32)]
2896pub enum OverWindowCachePolicy {
2897    Unspecified = 0,
2898    Full = 1,
2899    Recent = 2,
2900    RecentFirstN = 3,
2901    RecentLastN = 4,
2902}
2903impl OverWindowCachePolicy {
2904    /// String value of the enum field names used in the ProtoBuf definition.
2905    ///
2906    /// The values are not transformed in any way and thus are considered stable
2907    /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2908    pub fn as_str_name(&self) -> &'static str {
2909        match self {
2910            Self::Unspecified => "OVER_WINDOW_CACHE_POLICY_UNSPECIFIED",
2911            Self::Full => "OVER_WINDOW_CACHE_POLICY_FULL",
2912            Self::Recent => "OVER_WINDOW_CACHE_POLICY_RECENT",
2913            Self::RecentFirstN => "OVER_WINDOW_CACHE_POLICY_RECENT_FIRST_N",
2914            Self::RecentLastN => "OVER_WINDOW_CACHE_POLICY_RECENT_LAST_N",
2915        }
2916    }
2917    /// Creates an enum from field names used in the ProtoBuf definition.
2918    pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2919        match value {
2920            "OVER_WINDOW_CACHE_POLICY_UNSPECIFIED" => Some(Self::Unspecified),
2921            "OVER_WINDOW_CACHE_POLICY_FULL" => Some(Self::Full),
2922            "OVER_WINDOW_CACHE_POLICY_RECENT" => Some(Self::Recent),
2923            "OVER_WINDOW_CACHE_POLICY_RECENT_FIRST_N" => Some(Self::RecentFirstN),
2924            "OVER_WINDOW_CACHE_POLICY_RECENT_LAST_N" => Some(Self::RecentLastN),
2925            _ => None,
2926        }
2927    }
2928}
2929#[derive(prost_helpers::AnyPB)]
2930#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2931#[repr(i32)]
2932pub enum MatchRecognizeInputMode {
2933    Unspecified = 0,
2934    EventTime = 1,
2935    ProcessingTime = 2,
2936}
2937impl MatchRecognizeInputMode {
2938    /// String value of the enum field names used in the ProtoBuf definition.
2939    ///
2940    /// The values are not transformed in any way and thus are considered stable
2941    /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2942    pub fn as_str_name(&self) -> &'static str {
2943        match self {
2944            Self::Unspecified => "MATCH_RECOGNIZE_INPUT_MODE_UNSPECIFIED",
2945            Self::EventTime => "MATCH_RECOGNIZE_INPUT_MODE_EVENT_TIME",
2946            Self::ProcessingTime => "MATCH_RECOGNIZE_INPUT_MODE_PROCESSING_TIME",
2947        }
2948    }
2949    /// Creates an enum from field names used in the ProtoBuf definition.
2950    pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2951        match value {
2952            "MATCH_RECOGNIZE_INPUT_MODE_UNSPECIFIED" => Some(Self::Unspecified),
2953            "MATCH_RECOGNIZE_INPUT_MODE_EVENT_TIME" => Some(Self::EventTime),
2954            "MATCH_RECOGNIZE_INPUT_MODE_PROCESSING_TIME" => Some(Self::ProcessingTime),
2955            _ => None,
2956        }
2957    }
2958}
2959#[derive(prost_helpers::AnyPB)]
2960#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
2961#[repr(i32)]
2962pub enum DispatcherType {
2963    Unspecified = 0,
2964    /// Dispatch by hash key, hashed by consistent hash.
2965    Hash = 1,
2966    /// Broadcast to all downstreams.
2967    ///
2968    /// Note a broadcast cannot be represented as multiple simple dispatchers, since they are
2969    /// different when we update dispatchers during scaling.
2970    Broadcast = 2,
2971    /// Only one downstream.
2972    Simple = 3,
2973    /// A special kind of exchange that doesn't involve shuffle. The upstream actor will be directly
2974    /// piped into the downstream actor, if there are the same number of actors. If number of actors
2975    /// are not the same, should use hash instead. Should be only used when distribution is the same.
2976    NoShuffle = 4,
2977}
2978impl DispatcherType {
2979    /// String value of the enum field names used in the ProtoBuf definition.
2980    ///
2981    /// The values are not transformed in any way and thus are considered stable
2982    /// (if the ProtoBuf definition does not change) and safe for programmatic use.
2983    pub fn as_str_name(&self) -> &'static str {
2984        match self {
2985            Self::Unspecified => "DISPATCHER_TYPE_UNSPECIFIED",
2986            Self::Hash => "DISPATCHER_TYPE_HASH",
2987            Self::Broadcast => "DISPATCHER_TYPE_BROADCAST",
2988            Self::Simple => "DISPATCHER_TYPE_SIMPLE",
2989            Self::NoShuffle => "DISPATCHER_TYPE_NO_SHUFFLE",
2990        }
2991    }
2992    /// Creates an enum from field names used in the ProtoBuf definition.
2993    pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
2994        match value {
2995            "DISPATCHER_TYPE_UNSPECIFIED" => Some(Self::Unspecified),
2996            "DISPATCHER_TYPE_HASH" => Some(Self::Hash),
2997            "DISPATCHER_TYPE_BROADCAST" => Some(Self::Broadcast),
2998            "DISPATCHER_TYPE_SIMPLE" => Some(Self::Simple),
2999            "DISPATCHER_TYPE_NO_SHUFFLE" => Some(Self::NoShuffle),
3000            _ => None,
3001        }
3002    }
3003}