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(default)]
45 pub claims: HashMap<String, Value>,
46 #[builder(into)]
47 pub data_address: Option<DataAddress>,
48 #[builder(default = Utc::now())]
49 pub created_at: chrono::DateTime<chrono::Utc>,
50 #[builder(default = Utc::now())]
51 pub updated_at: chrono::DateTime<chrono::Utc>,
52}
53
54#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
55#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
56pub enum DataFlowType {
57 Consumer,
58 Provider,
59}
60
61#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
62#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
63pub enum DataFlowState {
64 Initiating,
65 Preparing,
66 Prepared,
67 Starting,
68 Started,
69 Suspended,
70 Completed,
71 Terminated,
72}
73
74#[derive(Debug, Error)]
75pub enum TransitionError {
76 #[error("Invalid state transition: {0}")]
77 InvalidTransition(String),
78}
79
80impl DataFlow {
81 pub fn transition_to_preparing(&mut self) -> Result<(), TransitionError> {
82 match self.state {
83 DataFlowState::Initiating => {
84 self.state = DataFlowState::Preparing;
85 self.updated_at = chrono::Utc::now();
86 Ok(())
87 }
88 DataFlowState::Preparing => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
90 "Invalid state for dataflow {} transition from {:?} to PREPARING",
91 self.id, self.state
92 ))),
93 }
94 }
95
96 pub fn transition_to_prepared(&mut self) -> Result<(), TransitionError> {
97 match self.state {
98 DataFlowState::Initiating | DataFlowState::Preparing => {
99 self.state = DataFlowState::Prepared;
100 self.updated_at = chrono::Utc::now();
101 Ok(())
102 }
103 DataFlowState::Prepared => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
105 "Invalid state for dataflow {} transition from {:?} to PREPARED",
106 self.id, self.state
107 ))),
108 }
109 }
110
111 pub fn transition_to_starting(&mut self) -> Result<(), TransitionError> {
112 match self.state {
113 DataFlowState::Initiating | DataFlowState::Prepared => {
114 self.state = DataFlowState::Starting;
115 self.updated_at = chrono::Utc::now();
116 Ok(())
117 }
118 DataFlowState::Starting => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
120 "Invalid state for dataflow {} transition from {:?} to STARTING",
121 self.id, self.state
122 ))),
123 }
124 }
125
126 pub fn transition_to_started(&mut self) -> Result<(), TransitionError> {
127 match self.state {
128 DataFlowState::Initiating
129 | DataFlowState::Starting
130 | DataFlowState::Prepared
131 | DataFlowState::Suspended => {
132 self.state = DataFlowState::Started;
133 self.updated_at = chrono::Utc::now();
134 Ok(())
135 }
136 DataFlowState::Started => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
138 "Invalid state for dataflow {} transition from {:?} to STARTED",
139 self.id, self.state
140 ))),
141 }
142 }
143
144 pub fn transition_to_suspended(
145 &mut self,
146 reason: Option<String>,
147 ) -> Result<(), TransitionError> {
148 match self.state {
149 DataFlowState::Started => {
150 self.state = DataFlowState::Suspended;
151 self.suspension_reason = reason;
152 self.updated_at = chrono::Utc::now();
153 Ok(())
154 }
155 DataFlowState::Suspended => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
157 "Invalid state for dataflow {} transition from {:?} to SUSPENDED",
158 self.id, self.state
159 ))),
160 }
161 }
162
163 pub fn transition_to_completed(&mut self) -> Result<(), TransitionError> {
164 match self.state {
165 DataFlowState::Started => {
166 self.state = DataFlowState::Completed;
167 self.updated_at = chrono::Utc::now();
168 Ok(())
169 }
170 DataFlowState::Completed => Ok(()), _ => Err(TransitionError::InvalidTransition(format!(
172 "Invalid state for dataflow {} transition from {:?} to COMPLETED",
173 self.id, self.state
174 ))),
175 }
176 }
177
178 pub fn transition_to_terminated(
179 &mut self,
180 reason: Option<String>,
181 ) -> Result<(), TransitionError> {
182 match self.state {
183 DataFlowState::Completed => Err(TransitionError::InvalidTransition(format!(
184 "Invalid state for dataflow {} transition from {:?} to TERMINATED",
185 self.id, self.state
186 ))),
187 DataFlowState::Terminated => Ok(()), _ => {
189 self.state = DataFlowState::Terminated;
190 self.termination_reason = reason;
191 self.updated_at = chrono::Utc::now();
192 Ok(())
193 }
194 }
195 }
196}