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}