Skip to main content

platform_provider/
admin_data.rs

1use crate::ProviderHostEffectCoordinator;
2use crate::config::ProviderConfig;
3use crate::invocation::{self, InvocationContext};
4use crate::protocol::{
5    ProviderAdminGetRequest, ProviderAdminListRequest, ProviderAdminQueryRequest,
6    ProviderGetResponse, ProviderInvocationMode, ProviderListResponse, ProviderOperationKind,
7    ProviderQueryResponse,
8};
9use platform_core::{ActorContext, AppError, AppResult, ErrorCode, TraceContext};
10use platform_module::{AdminDataSource, AdminListQuery, AdminPage, AdminQuerySource};
11use serde_json::Value;
12use std::time::Duration;
13
14#[derive(Debug, Clone)]
15pub struct ProviderAdminDataSource {
16    client: reqwest::Client,
17    config: ProviderConfig,
18    effects: ProviderHostEffectCoordinator,
19}
20
21impl ProviderAdminDataSource {
22    pub fn new(config: ProviderConfig) -> AppResult<Self> {
23        let client = reqwest::Client::builder()
24            .timeout(Duration::from_millis(config.timeout_ms))
25            .build()
26            .map_err(|error| {
27                AppError::new(
28                    ErrorCode::Internal,
29                    format!("failed to build Provider Service client: {error}"),
30                )
31            })?;
32        Ok(Self {
33            client,
34            config,
35            effects: ProviderHostEffectCoordinator::rejecting(),
36        })
37    }
38
39    #[must_use]
40    pub fn with_effect_coordinator(mut self, effects: ProviderHostEffectCoordinator) -> Self {
41        self.effects = effects;
42        self
43    }
44
45    async fn invoke<T: serde::de::DeserializeOwned>(
46        &self,
47        kind: ProviderOperationKind,
48        binding: &str,
49        operation: &str,
50        payload: Value,
51    ) -> AppResult<T> {
52        let invocation_id = uuid::Uuid::now_v7().to_string();
53        let invocation = invocation::build(
54            &self.config,
55            kind,
56            operation,
57            "1",
58            ProviderInvocationMode::ReadOnly,
59            InvocationContext {
60                request_id: invocation_id.clone(),
61                invocation_id,
62                attempt: 1,
63                actor: ActorContext::System,
64                correlation_id: uuid::Uuid::now_v7().to_string(),
65                causation_id: None,
66                trace: TraceContext::default(),
67            },
68            payload,
69        )?;
70        let outcome = invocation::send(
71            &self.client,
72            &self.config,
73            &self.effects,
74            binding,
75            &invocation,
76        )
77        .await?;
78        serde_json::from_value(invocation::result(&invocation, outcome)?).map_err(|error| {
79            AppError::new(
80                ErrorCode::ExternalDependency,
81                format!("Provider admin result violated its contract: {error}"),
82            )
83        })
84    }
85}
86
87#[async_trait::async_trait]
88impl AdminDataSource for ProviderAdminDataSource {
89    async fn list(&self, entity: &str, query: &AdminListQuery) -> AppResult<AdminPage> {
90        let response: ProviderListResponse = self
91            .invoke(
92                ProviderOperationKind::AdminList,
93                "admin:list",
94                entity,
95                serde_json::to_value(ProviderAdminListRequest {
96                    entity: entity.to_owned(),
97                    limit: query.limit,
98                    cursor: query.cursor.clone(),
99                })
100                .map_err(|error| AppError::new(ErrorCode::Internal, error.to_string()))?,
101            )
102            .await?;
103        Ok(response.into())
104    }
105
106    async fn get(&self, entity: &str, id: &str) -> AppResult<Option<Value>> {
107        let response: ProviderGetResponse = self
108            .invoke(
109                ProviderOperationKind::AdminGet,
110                "admin:get",
111                entity,
112                serde_json::to_value(ProviderAdminGetRequest {
113                    entity: entity.to_owned(),
114                    id: id.to_owned(),
115                })
116                .map_err(|error| AppError::new(ErrorCode::Internal, error.to_string()))?,
117            )
118            .await?;
119        Ok(response.record)
120    }
121}
122
123#[async_trait::async_trait]
124impl AdminQuerySource for ProviderAdminDataSource {
125    async fn query(&self, query: &str) -> AppResult<Value> {
126        validate_admin_query_name(query)?;
127        let response: ProviderQueryResponse = self
128            .invoke(
129                ProviderOperationKind::AdminQuery,
130                "admin:query",
131                query,
132                serde_json::to_value(ProviderAdminQueryRequest {
133                    query: query.to_owned(),
134                })
135                .map_err(|error| AppError::new(ErrorCode::Internal, error.to_string()))?,
136            )
137            .await?;
138        Ok(response.data)
139    }
140}
141
142fn validate_admin_query_name(query: &str) -> AppResult<()> {
143    let valid = !query.is_empty()
144        && query.chars().all(|character| {
145            character.is_ascii_alphanumeric()
146                || character == '.'
147                || character == '_'
148                || character == '-'
149        });
150    if valid {
151        return Ok(());
152    }
153
154    Err(AppError::new(
155        ErrorCode::Validation,
156        "provider admin query name must be a stable path segment",
157    ))
158}