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 event_name(&self) -> &str {
135        &self.event_name
136    }
137
138    async fn handle(&self, event: &ClaimedOutboxEvent) -> AppResult<()> {
139        self.invoke(event).await
140    }
141}
142
143#[async_trait::async_trait]
144trait ProviderEventActionRunner: std::fmt::Debug + Send + Sync {
145    async fn run_actions(
146        &self,
147        event: &ClaimedOutboxEvent,
148        handler_name: &str,
149        actions: Vec<ProviderEventResultAction>,
150    ) -> AppResult<()>;
151}
152
153#[derive(Debug, Clone)]
154pub struct ProviderEventHostActionRunner {
155    runtime: RuntimeClient,
156    function_registry: Arc<FunctionRegistry>,
157    allowed_function_names: BTreeSet<String>,
158}
159
160impl ProviderEventHostActionRunner {
161    #[must_use]
162    pub fn new(
163        runtime: RuntimeClient,
164        function_registry: Arc<FunctionRegistry>,
165        allowed_function_names: impl IntoIterator<Item = String>,
166    ) -> Self {
167        Self {
168            runtime,
169            function_registry,
170            allowed_function_names: allowed_function_names.into_iter().collect(),
171        }
172    }
173}
174
175#[async_trait::async_trait]
176impl ProviderEventActionRunner for ProviderEventHostActionRunner {
177    async fn run_actions(
178        &self,
179        event: &ClaimedOutboxEvent,
180        handler_name: &str,
181        actions: Vec<ProviderEventResultAction>,
182    ) -> AppResult<()> {
183        if actions.len() > MAX_EVENT_HANDLER_RESULT_ACTIONS {
184            return Err(AppError::new(
185                ErrorCode::Validation,
186                format!(
187                    "provider event handler {handler_name} returned too many result actions: {}",
188                    actions.len()
189                ),
190            ));
191        }
192
193        for (index, action) in actions.into_iter().enumerate() {
194            match action {
195                ProviderEventResultAction::EnqueueFunction {
196                    function_name,
197                    input,
198                } => {
199                    self.enqueue_function(event, handler_name, index, function_name, input)
200                        .await?;
201                }
202            }
203        }
204
205        Ok(())
206    }
207}
208
209impl ProviderEventHostActionRunner {
210    async fn enqueue_function(
211        &self,
212        event: &ClaimedOutboxEvent,
213        handler_name: &str,
214        action_index: usize,
215        function_name: String,
216        input: serde_json::Value,
217    ) -> AppResult<()> {
218        if !self.allowed_function_names.contains(&function_name) {
219            return Err(AppError::new(
220                ErrorCode::Validation,
221                format!(
222                    "provider event handler {handler_name} requested runtime function {function_name} that is not declared by its module"
223                ),
224            ));
225        }
226
227        let definition = self.function_registry.get(&function_name).ok_or_else(|| {
228            AppError::new(
229                ErrorCode::Internal,
230                format!("provider event handler {handler_name} requested unregistered runtime function {function_name}"),
231            )
232        })?;
233        let run_id = self
234            .runtime
235            .enqueue_function(EnqueueFunctionRequest {
236                function_name: function_name.clone(),
237                input_json: input,
238                correlation_id: CorrelationId::new(event.correlation_id.clone()),
239                actor: actor_from_event(event),
240                tenant_id: tenant_from_event(event),
241                tenancy_mode: platform_runtime::FunctionTenancyMode::Optional,
242                trace: trace_context_from_headers(&event.headers),
243                causation_id: Some(format!(
244                    "provider_event_handler:{}:{handler_name}:{action_index}",
245                    event.id
246                )),
247                max_attempts: Some(runtime_max_attempts_for_enqueue(
248                    definition.retry_policy.max_attempts,
249                )),
250            })
251            .await?;
252
253        tracing::info!(
254            outbox_event_id = %event.id,
255            handler_name = %handler_name,
256            function_name = %function_name,
257            function_run_id = %run_id,
258            "provider event handler enqueued runtime function"
259        );
260
261        Ok(())
262    }
263}
264
265fn tenant_from_event(event: &ClaimedOutboxEvent) -> Option<platform_core::TenantId> {
266    event
267        .headers
268        .get("tenant_id")
269        .cloned()
270        .and_then(|value| serde_json::from_value(value).ok())
271}
272
273#[derive(Debug)]
274struct RejectingProviderEventActionRunner;
275
276#[async_trait::async_trait]
277impl ProviderEventActionRunner for RejectingProviderEventActionRunner {
278    async fn run_actions(
279        &self,
280        _event: &ClaimedOutboxEvent,
281        handler_name: &str,
282        actions: Vec<ProviderEventResultAction>,
283    ) -> AppResult<()> {
284        if actions.is_empty() {
285            return Ok(());
286        }
287
288        Err(AppError::new(
289            ErrorCode::Validation,
290            format!(
291                "provider event handler {handler_name} returned result actions but host actions are not configured"
292            ),
293        ))
294    }
295}
296
297fn actor_from_event(event: &ClaimedOutboxEvent) -> ActorContext {
298    event
299        .headers
300        .get("actor")
301        .cloned()
302        .and_then(|actor| serde_json::from_value(actor).ok())
303        .unwrap_or_default()
304}
305
306pub(crate) fn validate_event_handler_name(value: &str) -> AppResult<()> {
307    validate_path_segment(
308        value,
309        "provider event handler name must be a stable path segment",
310    )
311}
312
313pub(crate) fn validate_event_name(value: &str) -> AppResult<()> {
314    validate_path_segment(value, "provider event name must be a stable path segment")
315}
316
317fn runtime_max_attempts_for_enqueue(max_attempts: u32) -> i32 {
318    i32::try_from(max_attempts).unwrap_or(i32::MAX)
319}