dataplane_sdk/core/model/
data_flow.rs1use std::collections::HashMap;
14
15use bon::Builder;
16use chrono::Utc;
17use serde::{Deserialize, Serialize};
18use serde_json::Value;
19use thiserror::Error;
20
21use super::data_address::DataAddress;
22
23#[derive(Builder, Clone, Debug, PartialEq)]
24#[builder(on(String, into))]
25pub struct DataFlow {
26 pub id: String,
27 #[builder(default = DataFlowState::Initiating)]
28 pub state: DataFlowState,
29 pub profile: String,
30 pub kind: DataFlowType,
31 pub agreement_id: String,
32 pub dataset_id: String,
33 pub dataspace_context: String,
34 pub participant_id: String,
35 pub counter_party_id: String,
36 pub control_plane_id: String,
37 pub participant_context_id: String,
38 pub suspension_reason: Option<String>,
39 pub termination_reason: Option<String>,
40 #[builder(default)]
41 pub labels: Vec<String>,
42 #[builder(default)]
43 pub metadata: HashMap<String, Value>,
44 #[builder(into)]
45 pub data_address: Option<DataAddress>,
46 #[builder(default = Utc::now())]
47 pub created_at: chrono::DateTime<chrono::Utc>,
48 #[builder(default = Utc::now())]
49 pub updated_at: chrono::DateTime<chrono::Utc>,
50}
51
52#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
53#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
54pub enum DataFlowType {
55 Consumer,
56 Provider,
57}
58
59#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
60#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
61pub enum DataFlowState {
62 Initiating,
63 Preparing,
64 Prepared,
65 Starting,
66 Started,
67 Suspended,
68 Completed,
69 Terminated,
70}
71
72#[derive(Debug, Error)]
73pub enum TransitionError {
74 #[error("Invalid state transition: {0}")]
75 InvalidTransition(String),
76}
77
78impl DataFlow {
79 pub fn transition_to_preparing(&mut self) -> Result<(), TransitionError> {
80 match self.state {
81 DataFlowState::Initiating => {
82 self.state = DataFlowState::Preparing;
83 self.updated_at = chrono::Utc::now();
84 Ok(())
85 }
86 DataFlowState::Preparing => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
88 "Invalid state for dataflow {} transition from {:?} to PREPARING",
89 self.id, self.state
90 ))),
91 }
92 }
93
94 pub fn transition_to_prepared(&mut self) -> Result<(), TransitionError> {
95 match self.state {
96 DataFlowState::Initiating | DataFlowState::Preparing => {
97 self.state = DataFlowState::Prepared;
98 self.updated_at = chrono::Utc::now();
99 Ok(())
100 }
101 DataFlowState::Prepared => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
103 "Invalid state for dataflow {} transition from {:?} to PREPARED",
104 self.id, self.state
105 ))),
106 }
107 }
108
109 pub fn transition_to_starting(&mut self) -> Result<(), TransitionError> {
110 match self.state {
111 DataFlowState::Initiating | DataFlowState::Prepared => {
112 self.state = DataFlowState::Starting;
113 self.updated_at = chrono::Utc::now();
114 Ok(())
115 }
116 DataFlowState::Starting => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
118 "Invalid state for dataflow {} transition from {:?} to STARTING",
119 self.id, self.state
120 ))),
121 }
122 }
123
124 pub fn transition_to_started(&mut self) -> Result<(), TransitionError> {
125 match self.state {
126 DataFlowState::Initiating
127 | DataFlowState::Starting
128 | DataFlowState::Prepared
129 | DataFlowState::Suspended => {
130 self.state = DataFlowState::Started;
131 self.updated_at = chrono::Utc::now();
132 Ok(())
133 }
134 DataFlowState::Started => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
136 "Invalid state for dataflow {} transition from {:?} to STARTED",
137 self.id, self.state
138 ))),
139 }
140 }
141
142 pub fn transition_to_suspended(
143 &mut self,
144 reason: Option<String>,
145 ) -> Result<(), TransitionError> {
146 match self.state {
147 DataFlowState::Started => {
148 self.state = DataFlowState::Suspended;
149 self.suspension_reason = reason;
150 self.updated_at = chrono::Utc::now();
151 Ok(())
152 }
153 DataFlowState::Suspended => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
155 "Invalid state for dataflow {} transition from {:?} to SUSPENDED",
156 self.id, self.state
157 ))),
158 }
159 }
160
161 pub fn transition_to_completed(&mut self) -> Result<(), TransitionError> {
162 match self.state {
163 DataFlowState::Started => {
164 self.state = DataFlowState::Completed;
165 self.updated_at = chrono::Utc::now();
166 Ok(())
167 }
168 DataFlowState::Completed => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
170 "Invalid state for dataflow {} transition from {:?} to COMPLETED",
171 self.id, self.state
172 ))),
173 }
174 }
175
176 pub fn transition_to_terminated(
177 &mut self,
178 reason: Option<String>,
179 ) -> Result<(), TransitionError> {
180 match self.state {
181 DataFlowState::Completed => Err(TransitionError::InvalidTransition(format!(
182 "Invalid state for dataflow {} transition from {:?} to TERMINATED",
183 self.id, self.state
184 ))),
185 DataFlowState::Terminated => Ok(()), _ => {
187 self.state = DataFlowState::Terminated;
188 self.termination_reason = reason;
189 self.updated_at = chrono::Utc::now();
190 Ok(())
191 }
192 }
193 }
194}