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