Skip to main content

platform_provider/
event.rs

1use crate::ProviderHostEffectCoordinator;
2use crate::config::ProviderConfig;
3use crate::invocation::{self, InvocationContext};
4use crate::protocol::{
5    ProviderEventHandleRequest, ProviderEventHandleResponse, ProviderEventResultAction,
6    ProviderInvocationMode, ProviderOperationKind,
7};
8use crate::validation::validate_path_segment;
9use platform_core::{
10    ActorContext, AppError, AppResult, ClaimedOutboxEvent, CorrelationId, ErrorCode, EventHandler,
11    trace_context_from_headers,
12};
13use platform_runtime::{EnqueueFunctionRequest, FunctionRegistry, RuntimeClient};
14use std::collections::BTreeSet;
15use std::sync::Arc;
16use std::time::Duration;
17
18const MAX_EVENT_HANDLER_RESULT_ACTIONS: usize = 1;
19
20#[derive(Debug, Clone)]
21pub struct ProviderEventHandler {
22    client: reqwest::Client,
23    config: ProviderConfig,
24    handler_name: String,
25    event_name: String,
26    action_runner: Arc<dyn ProviderEventActionRunner>,
27    effects: ProviderHostEffectCoordinator,
28}
29
30impl ProviderEventHandler {
31    pub fn new(
32        config: ProviderConfig,
33        handler_name: impl Into<String>,
34        event_name: impl Into<String>,
35        effects: ProviderHostEffectCoordinator,
36    ) -> AppResult<Self> {
37        let handler_name = handler_name.into();
38        let event_name = event_name.into();
39        validate_event_handler_name(&handler_name)?;
40        validate_event_name(&event_name)?;
41        let client = reqwest::Client::builder()
42            .timeout(Duration::from_millis(config.timeout_ms))
43            .build()
44            .map_err(|error| {
45                AppError::new(
46                    ErrorCode::Internal,
47                    format!("failed to build provider event handler client: {error}"),
48                )
49            })?;
50        Ok(Self {
51            client,
52            config,
53            handler_name,
54            event_name,
55            action_runner: Arc::new(RejectingProviderEventActionRunner),
56            effects,
57        })
58    }
59
60    #[must_use]
61    pub fn with_host_action_runner(mut self, action_runner: ProviderEventHostActionRunner) -> Self {
62        self.action_runner = Arc::new(action_runner);
63        self
64    }
65
66    pub async fn invoke(&self, event: &ClaimedOutboxEvent) -> AppResult<()> {
67        let request_body = ProviderEventHandleRequest {
68            request_id: format!("{}:{}", event.id, self.handler_name),
69            outbox_event_id: event.id.clone(),
70            handler_name: self.handler_name.clone(),
71            event_name: event.event_name.clone(),
72            event_version: event.event_version,
73            source_module: event.source_module.clone(),
74            aggregate_type: event.aggregate_type.clone(),
75            aggregate_id: event.aggregate_id.clone(),
76            correlation_id: event.correlation_id.clone(),
77            causation_id: event.causation_id.clone(),
78            occurred_at: event.occurred_at.to_rfc3339(),
79            actor: actor_from_event(event),
80            trace: trace_context_from_headers(&event.headers),
81            payload: event.payload.clone(),
82            headers: event.headers.clone(),
83        };
84        let invocation = invocation::build(
85            &self.config,
86            ProviderOperationKind::EventHandler,
87            &self.handler_name,
88            event.event_version.to_string(),
89            ProviderInvocationMode::Durable,
90            InvocationContext {
91                invocation_id: request_body.request_id.clone(),
92                request_id: request_body.request_id.clone(),
93                attempt: 1,
94                actor: request_body.actor.clone(),
95                correlation_id: request_body.correlation_id.clone(),
96                causation_id: request_body.causation_id.clone(),
97                trace: request_body.trace.clone(),
98            },
99            serde_json::to_value(&request_body).map_err(|error| {
100                AppError::new(
101                    ErrorCode::Internal,
102                    format!("encode Provider payload: {error}"),
103                )
104            })?,
105        )?;
106        let outcome = invocation::send(
107            &self.client,
108            &self.config,
109            &self.effects,
110            "events:handle",
111            &invocation,
112        )
113        .await?;
114        let value = invocation::result(&invocation, outcome)?;
115        let response: ProviderEventHandleResponse = if value.is_null() {
116            ProviderEventHandleResponse::default()
117        } else {
118            serde_json::from_value(value).map_err(|error| {
119                AppError::new(
120                    ErrorCode::ExternalDependency,
121                    format!("Provider event result violated its contract: {error}"),
122                )
123            })?
124        };
125        self.action_runner
126            .run_actions(event, &self.handler_name, response.actions)
127            .await?;
128        Ok(())
129    }
130}
131
132#[async_trait::async_trait]
133impl EventHandler for ProviderEventHandler {
134    fn handler_name(&self) -> &str {
135        &self.handler_name
136    }
137
138    fn event_name(&self) -> &str {
139        &self.event_name
140    }
141
142    async fn handle(&self, event: &ClaimedOutboxEvent) -> AppResult<()> {
143        self.invoke(event).await
144    }
145}
146
147#[async_trait::async_trait]
148trait ProviderEventActionRunner: std::fmt::Debug + Send + Sync {
149    async fn run_actions(
150        &self,
151        event: &ClaimedOutboxEvent,
152        handler_name: &str,
153        actions: Vec<ProviderEventResultAction>,
154    ) -> AppResult<()>;
155}
156
157#[derive(Debug, Clone)]
158pub struct ProviderEventHostActionRunner {
159    runtime: RuntimeClient,
160    function_registry: Arc<FunctionRegistry>,
161    allowed_function_names: BTreeSet<String>,
162}
163
164impl ProviderEventHostActionRunner {
165    #[must_use]
166    pub fn new(
167        runtime: RuntimeClient,
168        function_registry: Arc<FunctionRegistry>,
169        allowed_function_names: impl IntoIterator<Item = String>,
170    ) -> Self {
171        Self {
172            runtime,
173            function_registry,
174            allowed_function_names: allowed_function_names.into_iter().collect(),
175        }
176    }
177}
178
179#[async_trait::async_trait]
180impl ProviderEventActionRunner for ProviderEventHostActionRunner {
181    async fn run_actions(
182        &self,
183        event: &ClaimedOutboxEvent,
184        handler_name: &str,
185        actions: Vec<ProviderEventResultAction>,
186    ) -> AppResult<()> {
187        if actions.len() > MAX_EVENT_HANDLER_RESULT_ACTIONS {
188            return Err(AppError::new(
189                ErrorCode::Validation,
190                format!(
191                    "provider event handler {handler_name} returned too many result actions: {}",
192                    actions.len()
193                ),
194            ));
195        }
196
197        for (index, action) in actions.into_iter().enumerate() {
198            match action {
199                ProviderEventResultAction::EnqueueFunction {
200                    function_name,
201                    input,
202                } => {
203                    self.enqueue_function(event, handler_name, index, function_name, input)
204                        .await?;
205                }
206            }
207        }
208
209        Ok(())
210    }
211}
212
213impl ProviderEventHostActionRunner {
214    async fn enqueue_function(
215        &self,
216        event: &ClaimedOutboxEvent,
217        handler_name: &str,
218        action_index: usize,
219        function_name: String,
220        input: serde_json::Value,
221    ) -> AppResult<()> {
222        if !self.allowed_function_names.contains(&function_name) {
223            return Err(AppError::new(
224                ErrorCode::Validation,
225                format!(
226                    "provider event handler {handler_name} requested runtime function {function_name} that is not declared by its module"
227                ),
228            ));
229        }
230
231        let definition = self.function_registry.get(&function_name).ok_or_else(|| {
232            AppError::new(
233                ErrorCode::Internal,
234                format!("provider event handler {handler_name} requested unregistered runtime function {function_name}"),
235            )
236        })?;
237        let run_id = self
238            .runtime
239            .enqueue_function(EnqueueFunctionRequest {
240                function_name: function_name.clone(),
241                input_json: input,
242                correlation_id: CorrelationId::new(event.correlation_id.clone()),
243                actor: actor_from_event(event),
244                tenant_id: tenant_from_event(event),
245                tenancy_mode: platform_runtime::FunctionTenancyMode::Optional,
246                trace: trace_context_from_headers(&event.headers),
247                causation_id: Some(format!(
248                    "provider_event_handler:{}:{handler_name}:{action_index}",
249                    event.id
250                )),
251                max_attempts: Some(runtime_max_attempts_for_enqueue(
252                    definition.retry_policy.max_attempts,
253                )),
254            })
255            .await?;
256
257        tracing::info!(
258            outbox_event_id = %event.id,
259            handler_name = %handler_name,
260            function_name = %function_name,
261            function_run_id = %run_id,
262            "provider event handler enqueued runtime function"
263        );
264
265        Ok(())
266    }
267}
268
269fn tenant_from_event(event: &ClaimedOutboxEvent) -> Option<platform_core::TenantId> {
270    event
271        .headers
272        .get("tenant_id")
273        .cloned()
274        .and_then(|value| serde_json::from_value(value).ok())
275}
276
277#[derive(Debug)]
278struct RejectingProviderEventActionRunner;
279
280#[async_trait::async_trait]
281impl ProviderEventActionRunner for RejectingProviderEventActionRunner {
282    async fn run_actions(
283        &self,
284        _event: &ClaimedOutboxEvent,
285        handler_name: &str,
286        actions: Vec<ProviderEventResultAction>,
287    ) -> AppResult<()> {
288        if actions.is_empty() {
289            return Ok(());
290        }
291
292        Err(AppError::new(
293            ErrorCode::Validation,
294            format!(
295                "provider event handler {handler_name} returned result actions but host actions are not configured"
296            ),
297        ))
298    }
299}
300
301fn actor_from_event(event: &ClaimedOutboxEvent) -> ActorContext {
302    event
303        .headers
304        .get("actor")
305        .cloned()
306        .and_then(|actor| serde_json::from_value(actor).ok())
307        .unwrap_or_default()
308}
309
310pub(crate) fn validate_event_handler_name(value: &str) -> AppResult<()> {
311    validate_path_segment(
312        value,
313        "provider event handler name must be a stable path segment",
314    )
315}
316
317pub(crate) fn validate_event_name(value: &str) -> AppResult<()> {
318    validate_path_segment(value, "provider event name must be a stable path segment")
319}
320
321fn runtime_max_attempts_for_enqueue(max_attempts: u32) -> i32 {
322    i32::try_from(max_attempts).unwrap_or(i32::MAX)
323}