1mod metadata;
16
17use std::collections::{BTreeMap, HashMap};
18use std::fmt;
19
20use api::v1::ExpireAfter;
21use api::v1::flow::flow_request::Body as PbFlowRequest;
22use api::v1::flow::{CreateRequest, FlowRequest, FlowRequestHeader};
23use async_trait::async_trait;
24use chrono::{DateTime, Utc};
25use common_catalog::format_full_flow_name;
26use common_procedure::error::{FromJsonSnafu, ToJsonSnafu};
27use common_procedure::{
28 Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, ProcedureState,
29 Result as ProcedureResult, Status,
30};
31use common_telemetry::info;
32use common_telemetry::tracing_context::TracingContext;
33use futures::future::join_all;
34use itertools::Itertools;
35use serde::{Deserialize, Serialize};
36use snafu::{OptionExt, ResultExt, ensure};
37use strum::AsRefStr;
38use table::metadata::TableId;
39use table::table_name::TableName;
40
41use crate::cache_invalidator::Context;
42use crate::ddl::DdlContext;
43use crate::ddl::event::flow::{CREATE_FLOW_EVENT_TYPE, CreateFlowEventIntent, FlowDdlEvent};
44use crate::ddl::utils::{add_peer_context_if_needed, map_to_procedure_error};
45use crate::error::{self, Result, UnexpectedSnafu};
46use crate::instruction::{CacheIdent, CreateFlow, DropFlow};
47use crate::key::flow::flow_info::{FlowInfoValue, FlowScheduleConfig, FlowStatus};
48use crate::key::flow::flow_route::FlowRouteValue;
49use crate::key::table_name::TableNameKey;
50use crate::key::{DeserializedValueWithBytes, FlowId, FlowPartitionId};
51use crate::lock_key::{CatalogLock, FlowNameLock};
52use crate::metrics;
53use crate::peer::Peer;
54use crate::rpc::ddl::{CreateFlowTask, FlowQueryContext, QueryContext};
55
56pub struct CreateFlowProcedure {
58 pub context: DdlContext,
59 pub data: CreateFlowData,
60}
61
62impl CreateFlowProcedure {
63 pub const TYPE_NAME: &'static str = "metasrv-procedure::CreateFlow";
64
65 pub fn new(task: CreateFlowTask, query_context: QueryContext, context: DdlContext) -> Self {
67 Self {
68 context,
69 data: CreateFlowData {
70 task,
71 flow_id: None,
72 peers: vec![],
73 source_table_ids: vec![],
74 unresolved_source_table_names: vec![],
75 flow_context: without_scheduled_time_extension(query_context).into(),
76 state: CreateFlowState::Prepare,
77 prev_flow_info_value: None,
78 did_replace: false,
79 flow_type: None,
80 },
81 }
82 }
83
84 pub fn from_json(json: &str, context: DdlContext) -> ProcedureResult<Self> {
86 let data = serde_json::from_str(json).context(FromJsonSnafu)?;
87 Ok(CreateFlowProcedure { context, data })
88 }
89
90 pub(crate) async fn on_prepare(&mut self) -> Result<Status> {
91 let catalog_name = &self.data.task.catalog_name;
92 let flow_name = &self.data.task.flow_name;
93 let sink_table_name = &self.data.task.sink_table_name;
94 let create_if_not_exists = self.data.task.create_if_not_exists;
95 let or_replace = self.data.task.or_replace;
96
97 validate_flow_options(&self.data.task)?;
98
99 let flow_name_value = self
100 .context
101 .flow_metadata_manager
102 .flow_name_manager()
103 .get(catalog_name, flow_name)
104 .await?;
105
106 if create_if_not_exists && or_replace {
107 return error::UnsupportedSnafu {
109 operation: "Create flow with both `IF NOT EXISTS` and `OR REPLACE`",
110 }
111 .fail();
112 }
113
114 if let Some(value) = flow_name_value {
115 ensure!(
116 create_if_not_exists || or_replace,
117 error::FlowAlreadyExistsSnafu {
118 flow_name: format_full_flow_name(catalog_name, flow_name),
119 }
120 );
121
122 let flow_id = value.flow_id();
123 if create_if_not_exists {
124 info!("Flow already exists, flow_id: {}", flow_id);
125 return Ok(Status::done_with_output(flow_id));
126 }
127
128 let flow_id = value.flow_id();
129 let peers = self
130 .context
131 .flow_metadata_manager
132 .flow_route_manager()
133 .routes(flow_id)
134 .await?
135 .into_iter()
136 .map(|(_, value)| value.peer)
137 .collect::<Vec<_>>();
138 self.data.flow_id = Some(flow_id);
139 self.data.peers = peers;
140 info!("Replacing flow, flow_id: {}", flow_id);
141
142 let flow_info_value = self
143 .context
144 .flow_metadata_manager
145 .flow_info_manager()
146 .get_raw(flow_id)
147 .await?;
148
149 ensure!(
150 flow_info_value.is_some(),
151 error::FlowNotFoundSnafu {
152 flow_name: format_full_flow_name(catalog_name, flow_name),
153 }
154 );
155
156 self.data.prev_flow_info_value = flow_info_value;
157 }
158
159 let exists = self
161 .context
162 .table_metadata_manager
163 .table_name_manager()
164 .exists(TableNameKey::new(
165 &sink_table_name.catalog_name,
166 &sink_table_name.schema_name,
167 &sink_table_name.table_name,
168 ))
169 .await?;
170 if exists {
173 common_telemetry::warn!("Table already exists, table: {}", sink_table_name);
174 }
175
176 self.collect_source_tables().await?;
177 ensure!(
178 self.data.unresolved_source_table_names.is_empty()
179 || defer_on_missing_source(&self.data.task)?,
180 error::UnsupportedSnafu {
181 operation: format!(
182 "Create flow with missing source tables requires WITH ('{DEFER_ON_MISSING_SOURCE_KEY}'='true'): {}",
183 self.data
184 .unresolved_source_table_names
185 .iter()
186 .map(ToString::to_string)
187 .join(", ")
188 )
189 }
190 );
191 self.ensure_supported_replace_transition()?;
192
193 let sink_table_name = &self.data.task.sink_table_name;
195 if self
196 .data
197 .task
198 .source_table_names
199 .iter()
200 .any(|source| source == sink_table_name)
201 {
202 return error::UnsupportedSnafu {
203 operation: format!(
204 "Creating flow with source and sink table being the same: {}",
205 sink_table_name
206 ),
207 }
208 .fail();
209 }
210
211 if self.data.flow_id.is_none() {
212 self.allocate_flow_id().await?;
213 }
214 self.data.flow_type = Some(get_flow_type_from_options(&self.data.task)?);
215
216 resolve_schedule_defaults_into_task(
221 &mut self.data.task,
222 self.data
223 .prev_flow_info_value
224 .as_ref()
225 .map(|v| v.get_inner_ref()),
226 )?;
227
228 self.data.state = if self.data.is_pending() {
229 self.data.peers.clear();
230 CreateFlowState::CreateMetadata
231 } else {
232 CreateFlowState::CreateFlows
233 };
234
235 Ok(Status::executing(true))
236 }
237
238 fn ensure_supported_replace_transition(&self) -> Result<()> {
239 if !self.data.task.or_replace {
240 return Ok(());
241 }
242
243 let Some(prev_flow_info) = self.data.prev_flow_info_value.as_ref() else {
244 return Ok(());
245 };
246 let prev_pending = prev_flow_info.get_inner_ref().is_pending();
247 let new_pending = self.data.is_pending();
248 ensure!(
249 prev_pending == new_pending,
250 error::UnsupportedSnafu {
251 operation: "Replacing between pending and active flow states is not supported yet"
252 }
253 );
254
255 Ok(())
256 }
257
258 async fn on_flownode_create_flows(&mut self) -> Result<Status> {
259 let mut create_flow = Vec::with_capacity(self.data.peers.len());
261 for peer in &self.data.peers {
262 let requester = self.context.node_manager.flownode(peer).await;
263 let request = FlowRequest {
264 header: Some(FlowRequestHeader {
265 tracing_context: TracingContext::from_current_span().to_w3c(),
266 query_context: Some(
268 without_scheduled_time_extension(QueryContext::from(
269 self.data.flow_context.clone(),
270 ))
271 .into(),
272 ),
273 }),
274 body: Some(PbFlowRequest::Create((&self.data).into())),
275 };
276 create_flow.push(async move {
277 requester
278 .handle(request)
279 .await
280 .map_err(add_peer_context_if_needed(peer.clone()))
281 });
282 }
283 info!(
284 "Creating flow({:?}, type={:?}) on flownodes with peers={:?}",
285 self.data.flow_id, self.data.flow_type, self.data.peers
286 );
287 join_all(create_flow)
288 .await
289 .into_iter()
290 .collect::<Result<Vec<_>>>()?;
291
292 self.data.state = CreateFlowState::CreateMetadata;
293 Ok(Status::executing(true))
294 }
295
296 async fn on_create_metadata(&mut self) -> Result<Status> {
301 let flow_id = self.data.flow_id.unwrap();
303 let (flow_info, flow_routes) = (&self.data).into();
304 if let Some(prev_flow_value) = self.data.prev_flow_info_value.as_ref()
305 && self.data.task.or_replace
306 {
307 self.context
308 .flow_metadata_manager
309 .update_flow_metadata(flow_id, prev_flow_value, &flow_info, flow_routes)
310 .await?;
311 info!("Replaced flow metadata for flow {flow_id}");
312 self.data.did_replace = true;
313 } else {
314 self.context
315 .flow_metadata_manager
316 .create_flow_metadata(flow_id, flow_info, flow_routes)
317 .await?;
318 info!("Created flow metadata for flow {flow_id}");
319 }
320
321 self.data.state = CreateFlowState::InvalidateFlowCache;
322 Ok(Status::executing(true))
323 }
324
325 async fn on_broadcast(&mut self) -> Result<Status> {
326 debug_assert!(self.data.state == CreateFlowState::InvalidateFlowCache);
327 let flow_id = self.data.flow_id.unwrap();
329 let did_replace = self.data.did_replace;
330 let ctx = Context {
331 subject: Some("Invalidate flow cache by creating flow".to_string()),
332 };
333
334 let mut caches = vec![];
335
336 if did_replace {
338 let old_flow_info = self.data.prev_flow_info_value.as_ref().unwrap();
339
340 caches.extend([CacheIdent::DropFlow(DropFlow {
342 flow_id,
343 source_table_ids: old_flow_info.source_table_ids.clone(),
344 flow_part2node_id: old_flow_info.flownode_ids().clone().into_iter().collect(),
345 })]);
346 }
347
348 let (_flow_info, flow_routes) = (&self.data).into();
349 let flow_part2peers = flow_routes
350 .into_iter()
351 .map(|(part_id, route)| (part_id, route.peer))
352 .collect();
353
354 caches.extend([
355 CacheIdent::CreateFlow(CreateFlow {
356 flow_id,
357 source_table_ids: self.data.source_table_ids.clone(),
358 partition_to_peer_mapping: flow_part2peers,
359 }),
360 CacheIdent::FlowId(flow_id),
361 ]);
362
363 self.context
364 .cache_invalidator
365 .invalidate(&ctx, &caches)
366 .await?;
367
368 Ok(Status::done_with_output(flow_id))
369 }
370}
371
372#[async_trait]
373impl Procedure for CreateFlowProcedure {
374 fn type_name(&self) -> &str {
375 Self::TYPE_NAME
376 }
377
378 async fn execute(&mut self, _ctx: &ProcedureContext) -> ProcedureResult<Status> {
379 let state = &self.data.state;
380
381 let _timer = metrics::METRIC_META_PROCEDURE_CREATE_FLOW
382 .with_label_values(&[state.as_ref()])
383 .start_timer();
384
385 match state {
386 CreateFlowState::Prepare => self.on_prepare().await,
387 CreateFlowState::CreateFlows => self.on_flownode_create_flows().await,
388 CreateFlowState::CreateMetadata => self.on_create_metadata().await,
389 CreateFlowState::InvalidateFlowCache => self.on_broadcast().await,
390 }
391 .map_err(map_to_procedure_error)
392 }
393
394 fn dump(&self) -> ProcedureResult<String> {
395 serde_json::to_string(&self.data).context(ToJsonSnafu)
396 }
397
398 fn lock_key(&self) -> LockKey {
399 let catalog_name = &self.data.task.catalog_name;
400 let flow_name = &self.data.task.flow_name;
401
402 LockKey::new(vec![
403 CatalogLock::Read(catalog_name).into(),
404 FlowNameLock::new(catalog_name, flow_name).into(),
405 ])
406 }
407
408 fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn common_event_recorder::Event>> {
409 if !ctx.event_type_filter.allows(CREATE_FLOW_EVENT_TYPE) {
410 return None;
411 }
412
413 let event = match &ctx.trigger {
414 EventTrigger::Submitted => FlowDdlEvent::create_submitted(
415 &self.data.task.catalog_name,
416 &self.data.task.flow_name,
417 CreateFlowEventIntent {
418 or_replace: self.data.task.or_replace,
419 create_if_not_exists: self.data.task.create_if_not_exists,
420 expire_after: self.data.task.expire_after,
421 eval_interval_secs: self.data.task.eval_interval_secs,
422 },
423 ),
424 EventTrigger::Succeeded => {
425 let flow_id = match ctx.lifecycle_state {
426 ProcedureState::Done {
427 output: Some(output),
428 } => output
429 .downcast_ref::<FlowId>()
430 .copied()
431 .or(self.data.flow_id),
432 _ => self.data.flow_id,
433 };
434 FlowDdlEvent::create_succeeded(
435 &self.data.task.catalog_name,
436 &self.data.task.flow_name,
437 flow_id,
438 )
439 }
440 _ => FlowDdlEvent::create_lifecycle(
441 &self.data.task.catalog_name,
442 &self.data.task.flow_name,
443 ),
444 };
445
446 Some(Box::new(event))
447 }
448}
449
450pub fn get_flow_type_from_options(flow_task: &CreateFlowTask) -> Result<FlowType> {
451 let flow_type = flow_task
452 .flow_options
453 .get(FlowType::FLOW_TYPE_KEY)
454 .map(|s| s.as_str());
455 match flow_type {
456 Some(FlowType::BATCHING) => Ok(FlowType::Batching),
457 Some(FlowType::STREAMING) => Ok(FlowType::Streaming),
458 Some(unknown) => UnexpectedSnafu {
459 err_msg: format!("Unknown flow type: {}", unknown),
460 }
461 .fail(),
462 None => Ok(FlowType::Batching),
463 }
464}
465
466pub const DEFER_ON_MISSING_SOURCE_KEY: &str = "defer_on_missing_source";
468
469pub const INTERNAL_EVAL_OFFSET_KEY: &str = "__greptime_internal_eval_offset_secs";
476
477pub const INTERNAL_EVAL_SCHEDULE_KEY: &str = "__greptime_internal_eval_schedule";
484
485const FLOW_SCHEDULED_TIME_MILLIS_EXTENSION_KEY: &str = "flow.scheduled_time_millis";
486
487fn without_scheduled_time_extension(mut query_context: QueryContext) -> QueryContext {
488 query_context
489 .extensions
490 .remove(FLOW_SCHEDULED_TIME_MILLIS_EXTENSION_KEY);
491 query_context
492}
493
494pub fn defer_on_missing_source(flow_task: &CreateFlowTask) -> Result<bool> {
495 flow_task
496 .flow_options
497 .get(DEFER_ON_MISSING_SOURCE_KEY)
498 .map(|value| {
499 value
500 .trim()
501 .to_ascii_lowercase()
502 .parse::<bool>()
503 .map_err(|_| {
504 error::UnexpectedSnafu {
505 err_msg: format!(
506 "Invalid flow option '{DEFER_ON_MISSING_SOURCE_KEY}': {value}"
507 ),
508 }
509 .build()
510 })
511 })
512 .transpose()
513 .map(|value| value.unwrap_or(false))
514}
515
516pub fn validate_flow_options(flow_task: &CreateFlowTask) -> Result<()> {
517 if let Some(secs) = flow_task.eval_interval_secs
519 && secs <= 0
520 {
521 return UnexpectedSnafu {
522 err_msg: format!("EVAL INTERVAL must be positive, got {secs} seconds"),
523 }
524 .fail();
525 }
526
527 for key in [INTERNAL_EVAL_OFFSET_KEY, INTERNAL_EVAL_SCHEDULE_KEY] {
528 if flow_task.flow_options.contains_key(key) {
529 return UnexpectedSnafu {
530 err_msg: format!("flow option '{key}' is reserved for internal use"),
531 }
532 .fail();
533 }
534 }
535
536 if let Some(offset_secs) = flow_task.eval_offset_secs {
539 let Some(eval_interval_secs) = flow_task.eval_interval_secs else {
540 return UnexpectedSnafu {
541 err_msg: "EVAL OFFSET requires EVAL INTERVAL to be specified".to_string(),
542 }
543 .fail();
544 };
545 if !(0..eval_interval_secs).contains(&offset_secs) {
546 return UnexpectedSnafu {
547 err_msg: format!(
548 "EVAL OFFSET must be in range [0, EVAL INTERVAL), got {offset_secs} seconds with EVAL INTERVAL {eval_interval_secs} seconds"
549 ),
550 }
551 .fail();
552 }
553 }
554
555 for key in flow_task.flow_options.keys() {
556 match key.as_str() {
557 DEFER_ON_MISSING_SOURCE_KEY
558 | FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_KEY
559 | FlowType::FLOW_TYPE_KEY => {}
560 unknown => {
561 return UnexpectedSnafu {
562 err_msg: format!(
563 "Unknown flow option '{unknown}', supported user options: {DEFER_ON_MISSING_SOURCE_KEY}, {FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_KEY}"
564 ),
565 }
566 .fail();
567 }
568 }
569 }
570
571 if let Some(value) = flow_task
572 .flow_options
573 .get(FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_KEY)
574 && value != FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_SEQUENCE_RANGE
575 {
576 value.parse::<bool>().map_err(|_| {
577 UnexpectedSnafu {
578 err_msg: format!(
579 "Invalid flow option {FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_KEY}: {value}"
580 ),
581 }
582 .build()
583 })?;
584 }
585
586 defer_on_missing_source(flow_task)?;
587 get_flow_type_from_options(flow_task)?;
588 Ok(())
589}
590
591pub(crate) fn ceil_to_boundary(time: i64, anchor: i64, interval: i64) -> Result<i64> {
598 if interval <= 0 {
599 return Ok(time);
600 }
601 if time <= anchor {
602 return Ok(anchor);
603 }
604
605 let diff = i128::from(time) - i128::from(anchor);
606 let interval = i128::from(interval);
607 let k = (diff + interval - 1) / interval;
608 let boundary = i128::from(anchor) + k * interval;
609
610 i64::try_from(boundary).map_err(|_| {
611 UnexpectedSnafu {
612 err_msg: format!(
613 "Cannot align time {time} to the next `anchor + k * interval` boundary (anchor={anchor}, interval={interval}): result {boundary} does not fit in i64"
614 ),
615 }
616 .build()
617 })
618}
619
620pub(crate) fn ceil_to_whole_sec(now: DateTime<Utc>) -> Result<i64> {
630 ceil_whole_sec_from_parts(now.timestamp(), now.timestamp_subsec_nanos() != 0)
631}
632
633pub(crate) fn ceil_whole_sec_from_parts(secs: i64, has_fraction: bool) -> Result<i64> {
637 if !has_fraction {
638 return Ok(secs);
639 }
640 secs.checked_add(1).context(error::UnexpectedSnafu {
641 err_msg: format!(
642 "Cannot round instant at second {secs} up to the next whole second: timestamp overflow"
643 ),
644 })
645}
646
647pub fn effective_eval_schedule_from_flow_info(
653 flow_info: &FlowInfoValue,
654) -> Result<Option<FlowScheduleConfig>> {
655 if let Some(schedule) = &flow_info.eval_schedule {
656 return Ok(Some(schedule.clone()));
657 }
658
659 let Some(eval_interval_secs) = flow_info.eval_interval_secs else {
660 return Ok(None);
661 };
662 if eval_interval_secs <= 0 {
663 return Ok(None);
664 }
665
666 let created_ceil = ceil_to_whole_sec(flow_info.created_time)?;
669 let start_secs = ceil_to_boundary(
670 created_ceil,
671 FlowScheduleConfig::DEFAULT_ANCHOR_SECS,
672 eval_interval_secs,
673 )?;
674
675 Ok(Some(FlowScheduleConfig::default_with_start(
676 start_secs,
677 eval_interval_secs,
678 )))
679}
680
681pub(crate) fn resolve_schedule_defaults_into_task(
702 task: &mut CreateFlowTask,
703 prev_flow_info: Option<&FlowInfoValue>,
704) -> Result<()> {
705 if task.eval_schedule.is_some() {
707 return Ok(());
708 }
709
710 let Some(eval_interval_secs) = task.eval_interval_secs else {
711 return Ok(());
712 };
713 if eval_interval_secs <= 0 {
714 return Ok(());
715 }
716
717 let anchor_secs = task.eval_offset_secs.unwrap_or(0);
718
719 if !(0..eval_interval_secs).contains(&anchor_secs) {
723 return Ok(());
724 }
725
726 if task.or_replace
729 && let Some(prev) = prev_flow_info
730 && let Some(old_sched) = effective_eval_schedule_from_flow_info(prev)?
731 {
732 let old_interval = prev.eval_interval_secs.unwrap_or(0);
733 if old_interval == eval_interval_secs && old_sched.anchor_secs == anchor_secs {
734 task.eval_schedule = Some(old_sched);
735 return Ok(());
736 }
737 }
738
739 let prepare_secs = ceil_to_whole_sec(chrono::Utc::now())?;
745 let start_secs = ceil_to_boundary(prepare_secs, anchor_secs, eval_interval_secs)?;
746
747 task.eval_schedule = Some(FlowScheduleConfig::with_anchor(
748 anchor_secs,
749 start_secs,
750 eval_interval_secs,
751 ));
752 Ok(())
753}
754
755fn user_runtime_flow_options(options: &HashMap<String, String>) -> HashMap<String, String> {
756 let mut options = options.clone();
757 options.remove(DEFER_ON_MISSING_SOURCE_KEY);
758 options.remove(INTERNAL_EVAL_SCHEDULE_KEY);
759 options.remove(INTERNAL_EVAL_OFFSET_KEY);
760 options
761}
762
763#[derive(Debug, Clone, Serialize, Deserialize, AsRefStr, PartialEq)]
765pub enum CreateFlowState {
766 Prepare,
768 CreateFlows,
770 InvalidateFlowCache,
772 CreateMetadata,
774}
775
776#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, Default)]
778pub enum FlowType {
779 #[default]
781 Batching,
782 Streaming,
784}
785
786pub const FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_KEY: &str =
787 "experimental_enable_incremental_read";
788pub const FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_SEQUENCE_RANGE: &str =
790 "__greptime_internal_exact_sequence_range";
791
792impl FlowType {
793 pub const BATCHING: &str = "batching";
794 pub const STREAMING: &str = "streaming";
795 pub const FLOW_TYPE_KEY: &str = "flow_type";
796}
797
798impl fmt::Display for FlowType {
799 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
800 match self {
801 FlowType::Batching => write!(f, "{}", FlowType::BATCHING),
802 FlowType::Streaming => write!(f, "{}", FlowType::STREAMING),
803 }
804 }
805}
806
807#[derive(Debug, Serialize, Deserialize)]
809pub struct CreateFlowData {
810 pub(crate) state: CreateFlowState,
811 pub(crate) task: CreateFlowTask,
812 pub(crate) flow_id: Option<FlowId>,
813 pub(crate) peers: Vec<Peer>,
814 pub(crate) source_table_ids: Vec<TableId>,
815 #[serde(default)]
816 pub(crate) unresolved_source_table_names: Vec<TableName>,
817 #[serde(alias = "query_context")]
819 pub(crate) flow_context: FlowQueryContext,
820 pub(crate) prev_flow_info_value: Option<DeserializedValueWithBytes<FlowInfoValue>>,
823 #[serde(default)]
826 pub(crate) did_replace: bool,
827 pub(crate) flow_type: Option<FlowType>,
828}
829
830impl CreateFlowData {
831 pub(crate) fn is_pending(&self) -> bool {
832 !self.unresolved_source_table_names.is_empty()
833 }
834
835 pub(crate) fn is_active(&self) -> bool {
836 !self.is_pending()
837 }
838}
839
840impl From<&CreateFlowData> for CreateRequest {
841 fn from(value: &CreateFlowData) -> Self {
842 let flow_id = value.flow_id.unwrap();
843 let source_table_ids = &value.source_table_ids;
844
845 let mut req = CreateRequest {
846 flow_id: Some(api::v1::FlowId { id: flow_id }),
847 source_table_ids: source_table_ids
848 .iter()
849 .map(|table_id| api::v1::TableId { id: *table_id })
850 .collect_vec(),
851 sink_table_name: Some(value.task.sink_table_name.clone().into()),
852 create_if_not_exists: true,
854 or_replace: value.task.or_replace,
855 expire_after: value.task.expire_after.map(|value| ExpireAfter { value }),
856 eval_interval: value
857 .task
858 .eval_interval_secs
859 .map(|seconds| api::v1::EvalInterval { seconds }),
860 comment: value.task.comment.clone(),
861 sql: value.task.sql.clone(),
862 flow_options: user_runtime_flow_options(&value.task.flow_options),
863 };
864
865 let flow_type = value.flow_type.unwrap_or_default().to_string();
866 req.flow_options
867 .insert(FlowType::FLOW_TYPE_KEY.to_string(), flow_type);
868
869 if let Some(ref sched) = value.task.eval_schedule {
871 let json = serde_json::to_string(sched)
872 .expect("FlowScheduleConfig serialization should not fail");
873 req.flow_options
874 .insert(INTERNAL_EVAL_SCHEDULE_KEY.to_string(), json);
875 }
876
877 req
878 }
879}
880
881impl From<&CreateFlowData> for (FlowInfoValue, Vec<(FlowPartitionId, FlowRouteValue)>) {
882 fn from(value: &CreateFlowData) -> Self {
883 let catalog_name = value.task.catalog_name.clone();
884 let flow_name = value.task.flow_name.clone();
885 let sink_table_name = value.task.sink_table_name.clone();
886 let expire_after = value.task.expire_after;
887 let eval_interval = value.task.eval_interval_secs;
888 let comment = value.task.comment.clone();
889 let sql = value.task.sql.clone();
890 let eval_schedule = value.task.eval_schedule.clone();
891
892 let mut options: HashMap<String, String> = value
896 .task
897 .flow_options
898 .iter()
899 .filter(|(k, _)| {
900 k.as_str() != INTERNAL_EVAL_SCHEDULE_KEY && k.as_str() != INTERNAL_EVAL_OFFSET_KEY
901 })
902 .map(|(k, v)| (k.clone(), v.clone()))
903 .collect();
904
905 let flownode_ids = value
906 .peers
907 .iter()
908 .enumerate()
909 .map(|(idx, peer)| (idx as u32, peer.id))
910 .collect::<BTreeMap<_, _>>();
911 let flow_routes = value
912 .peers
913 .iter()
914 .enumerate()
915 .map(|(idx, peer)| (idx as u32, FlowRouteValue { peer: peer.clone() }))
916 .collect::<Vec<_>>();
917
918 let flow_type = value.flow_type.unwrap_or_default().to_string();
919 options.insert(FlowType::FLOW_TYPE_KEY.to_string(), flow_type);
920
921 let mut create_time = chrono::Utc::now();
922 if let Some(prev_flow_value) = value.prev_flow_info_value.as_ref()
923 && value.task.or_replace
924 {
925 create_time = prev_flow_value.get_inner_ref().created_time;
926 }
927
928 let flow_info: FlowInfoValue = FlowInfoValue {
933 source_table_ids: value.source_table_ids.clone(),
934 all_source_table_names: value.task.source_table_names.clone(),
935 unresolved_source_table_names: value.unresolved_source_table_names.clone(),
936 sink_table_name: sink_table_name.clone(),
937 flownode_ids,
938 catalog_name: catalog_name.clone(),
939 query_context: Some(without_scheduled_time_extension(QueryContext::from(
940 value.flow_context.clone(),
941 ))),
942 flow_name: flow_name.clone(),
943 raw_sql: sql.clone(),
944 expire_after,
945 eval_interval_secs: eval_interval,
946 comment: comment.clone(),
947 options,
948 status: if value.is_active() {
949 FlowStatus::Active
950 } else {
951 FlowStatus::PendingSources
952 },
953 created_time: create_time,
954 updated_time: chrono::Utc::now(),
955 eval_schedule: eval_schedule.clone(),
956 };
957
958 (flow_info, flow_routes)
959 }
960}