Skip to main content

dataplane_sdk/core/model/
data_flow.rs

1//  Copyright (c) 2026 Metaform Systems, Inc
2//
3//  This program and the accompanying materials are made available under the
4//  terms of the Apache License, Version 2.0 which is available at
5//  https://www.apache.org/licenses/LICENSE-2.0
6//
7//    SPDX-License-Identifier: Apache-2.0
8//
9//    Contributors:
10//         Metaform Systems, Inc. - initial API and implementation
11//
12
13use 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(()), // already in preparing, idempotent
87            _ => 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(()), // already in prepared, idempotent
102            _ => 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(()), // already in starting, idempotent
117            _ => 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(()), // already in started, idempotent
135            _ => 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(()), // already in suspended, idempotent
154            _ => 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(()), // already in completed, idempotent
169            _ => 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(()), // already in terminated, idempotent
186            _ => {
187                self.state = DataFlowState::Terminated;
188                self.termination_reason = reason;
189                self.updated_at = chrono::Utc::now();
190                Ok(())
191            }
192        }
193    }
194}