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}