From 74ad7e89cd402e9b18ee49141e2d5d3ae7ecf80e Mon Sep 17 00:00:00 2001 From: Enrico Risa Date: Wed, 29 Jul 2026 14:40:16 +0200 Subject: [PATCH] feat: add support for claims in the start/prepare messages --- .../migrations/20260218164640_create_data_flows.sql | 1 + crates/sdk-postgres/src/data_flow/model.rs | 3 +++ crates/sdk-postgres/src/data_flow/repo.rs | 5 +++-- crates/sdk/src/core/db/test_suite.rs | 5 +++++ crates/sdk/src/core/model/data_flow.rs | 2 ++ crates/sdk/src/core/model/messages.rs | 6 ++++++ crates/sdk/src/sdk/internal.rs | 2 ++ 7 files changed, 22 insertions(+), 2 deletions(-) diff --git a/crates/sdk-postgres/migrations/20260218164640_create_data_flows.sql b/crates/sdk-postgres/migrations/20260218164640_create_data_flows.sql index 8a140ce..dece723 100644 --- a/crates/sdk-postgres/migrations/20260218164640_create_data_flows.sql +++ b/crates/sdk-postgres/migrations/20260218164640_create_data_flows.sql @@ -18,6 +18,7 @@ CREATE TABLE IF NOT EXISTS data_flows ( suspension_reason TEXT, termination_reason TEXT, metadata JSONB NOT NULL DEFAULT '{}', + claims JSONB NOT NULL DEFAULT '{}', labels JSONB NOT NULL DEFAULT '[]', created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP diff --git a/crates/sdk-postgres/src/data_flow/model.rs b/crates/sdk-postgres/src/data_flow/model.rs index 99bf2cb..9122e6f 100644 --- a/crates/sdk-postgres/src/data_flow/model.rs +++ b/crates/sdk-postgres/src/data_flow/model.rs @@ -37,6 +37,8 @@ pub struct DataFlow { pub termination_reason: Option, #[builder(default)] pub metadata: Json>, + #[builder(default)] + pub claims: Json>, #[builder(into)] pub data_address: Option>, #[builder(into)] @@ -78,6 +80,7 @@ impl From for dataplane_sdk::core::model::data_flow::DataFlow { .labels(flow.labels.0) .agreement_id(flow.agreement_id) .metadata(flow.metadata.0) + .claims(flow.claims.0) .dataset_id(flow.dataset_id) .dataspace_context(flow.dataspace_context) .participant_id(flow.participant_id) diff --git a/crates/sdk-postgres/src/data_flow/repo.rs b/crates/sdk-postgres/src/data_flow/repo.rs index 0576541..a48dcae 100644 --- a/crates/sdk-postgres/src/data_flow/repo.rs +++ b/crates/sdk-postgres/src/data_flow/repo.rs @@ -33,8 +33,8 @@ impl DataFlowRepo for PgDataFlowRepo { async fn create(&self, tx: &mut Self::Transaction, flow: &DataFlow) -> DbResult<()> { let result = sqlx::query( r#" - INSERT INTO data_flows (id, participant_id, dataspace_context, participant_context_id, counter_party_id, dataset_id, agreement_id, state, profile, type, data_address, control_plane_id, labels, metadata, created_at, updated_at) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16) + INSERT INTO data_flows (id, participant_id, dataspace_context, participant_context_id, counter_party_id, dataset_id, agreement_id, state, profile, type, data_address, control_plane_id, labels, metadata, claims, created_at, updated_at) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17) "#, ) .bind(&flow.id) @@ -51,6 +51,7 @@ impl DataFlowRepo for PgDataFlowRepo { .bind(&flow.control_plane_id) .bind(Json(flow.labels.clone())) .bind(Json(flow.metadata.clone())) + .bind(Json(flow.claims.clone())) .bind(flow.created_at) .bind(flow.updated_at) .execute(&mut *tx.0) diff --git a/crates/sdk/src/core/db/test_suite.rs b/crates/sdk/src/core/db/test_suite.rs index 66cd2df..bf47b5c 100644 --- a/crates/sdk/src/core/db/test_suite.rs +++ b/crates/sdk/src/core/db/test_suite.rs @@ -50,6 +50,11 @@ pub fn create_data_flow(id: &str) -> DataFlow { .into_iter() .collect(), ) + .claims( + vec![("claim".to_string(), "value".into())] + .into_iter() + .collect(), + ) .dataset_id("dataset_id") .dataspace_context("dataspace_context") .participant_id("participant_id") diff --git a/crates/sdk/src/core/model/data_flow.rs b/crates/sdk/src/core/model/data_flow.rs index b21a83b..4e17896 100644 --- a/crates/sdk/src/core/model/data_flow.rs +++ b/crates/sdk/src/core/model/data_flow.rs @@ -41,6 +41,8 @@ pub struct DataFlow { pub labels: Vec, #[builder(default)] pub metadata: HashMap, + #[builder(default)] + pub claims: HashMap, #[builder(into)] pub data_address: Option, #[builder(default = Utc::now())] diff --git a/crates/sdk/src/core/model/messages.rs b/crates/sdk/src/core/model/messages.rs index 67e6579..f8e8786 100644 --- a/crates/sdk/src/core/model/messages.rs +++ b/crates/sdk/src/core/model/messages.rs @@ -37,6 +37,9 @@ pub struct DataFlowStartMessage { #[builder(default)] #[serde(default)] pub metadata: HashMap, + #[builder(default)] + #[serde(default)] + pub claims: HashMap, } #[derive(Debug, Serialize, Deserialize, Clone, Builder)] @@ -64,6 +67,9 @@ pub struct DataFlowPrepareMessage { #[builder(default)] #[serde(default)] pub metadata: HashMap, + #[builder(default)] + #[serde(default)] + pub claims: HashMap, } #[derive(Debug, Builder, Serialize, Deserialize, Clone)] diff --git a/crates/sdk/src/sdk/internal.rs b/crates/sdk/src/sdk/internal.rs index 40cb42a..79e70b1 100644 --- a/crates/sdk/src/sdk/internal.rs +++ b/crates/sdk/src/sdk/internal.rs @@ -60,6 +60,7 @@ where .participant_context_id(participant_context_id) .state(DataFlowState::Initiating) .metadata(req.metadata) + .claims(req.claims) .participant_id(req.participant_id) .dataspace_context(req.dataspace_context) .dataset_id(req.dataset_id) @@ -107,6 +108,7 @@ where .participant_context_id(participant_context_id) .state(DataFlowState::Initiating) .metadata(req.metadata) + .claims(req.claims) .participant_id(req.participant_id) .dataspace_context(req.dataspace_context) .dataset_id(req.dataset_id)