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}