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(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(()), // already in preparing, idempotent
89            _ => 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(()), // already in prepared, idempotent
104            _ => 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(()), // already in starting, idempotent
119            _ => 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(()), // already in started, idempotent
137            _ => 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(()), // already in suspended, idempotent
156            _ => 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(()), // already in completed, idempotent
171            _ => 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(()), // already in terminated, idempotent
188            _ => {
189                self.state = DataFlowState::Terminated;
190                self.termination_reason = reason;
191                self.updated_at = chrono::Utc::now();
192                Ok(())
193            }
194        }
195    }
196}