1use crate::api::event::{BaseEvent, CategoryProfile, DataSchema, EventCategory, MarkEvent};
5use crate::api::runtime::global_context;
6use crate::api::runtime::scope_stack::snapshot_scope_stack;
7use crate::api::runtime::subscriber_dispatcher::{self, SubscriberDelivery};
8use crate::api::runtime::{
9 current_scope_stack, task_scope_push, task_scope_remove, task_scope_top,
10};
11use crate::api::shared::{
12 ensure_runtime_owner, resolve_parent_uuid, snapshot_event_sanitizers,
13 snapshot_event_subscribers,
14};
15use crate::error::{FlowError, Result};
16use crate::json::Json;
17use chrono::{DateTime, Utc};
18use serde::{Deserialize, Serialize};
19use typed_builder::TypedBuilder;
20use uuid::Uuid;
21
22pub use nemo_relay_types::api::scope::{HandleAttributes, ScopeAttributes, ScopeType};
23
24pub const COMPACTION_EVENT_NAME: &str = "compaction";
26
27#[derive(Debug, Clone, Serialize, Deserialize, TypedBuilder)]
29#[builder(field_defaults(setter(strip_option(ignore_invalid, fallback_suffix = "_opt"))))]
30pub struct ScopeHandle {
31 #[builder(default = Uuid::now_v7())]
33 pub uuid: Uuid,
34 #[builder(default = Utc::now())]
36 pub started_at: DateTime<Utc>,
37 pub scope_type: ScopeType,
39 #[builder(setter(into))]
41 pub name: String,
42 #[builder(default)]
44 pub data: Option<Json>,
45 #[builder(default)]
47 pub metadata: Option<Json>,
48 #[builder(default = ScopeAttributes::empty())]
50 pub attributes: ScopeAttributes,
51 #[builder(default)]
53 pub parent_uuid: Option<Uuid>,
54}
55
56fn scope_stack_lock_error(error: impl std::fmt::Display, operation: &'static str) -> FlowError {
57 log::error!(
58 target: "nemo_relay.runtime",
59 event = "scope_stack_unavailable",
60 operation = operation;
61 "Scope operation failed because the scope stack lock is poisoned: {error}"
62 );
63 FlowError::Internal(error.to_string())
64}
65
66#[derive(TypedBuilder)]
68#[builder(field_defaults(setter(strip_option(ignore_invalid, fallback_suffix = "_opt"))))]
69pub struct PushScopeParams<'a> {
70 pub name: &'a str,
72 pub scope_type: ScopeType,
74 #[builder(default)]
76 pub parent: Option<&'a ScopeHandle>,
77 #[builder(default = ScopeAttributes::empty())]
79 pub attributes: ScopeAttributes,
80 #[builder(default)]
82 pub data: Option<Json>,
83 #[builder(default)]
85 pub metadata: Option<Json>,
86 #[builder(default)]
88 pub input: Option<Json>,
89 #[builder(default)]
91 pub timestamp: Option<DateTime<Utc>>,
92}
93
94#[derive(Debug, Clone, TypedBuilder)]
96#[builder(field_defaults(setter(strip_option(ignore_invalid, fallback_suffix = "_opt"))))]
97pub struct CreateScopeHandleParams<'a> {
98 pub name: &'a str,
100 #[builder(default)]
102 pub parent_uuid: Option<Uuid>,
103 pub scope_type: ScopeType,
105 #[builder(default = ScopeAttributes::empty())]
107 pub attributes: ScopeAttributes,
108 #[builder(default)]
110 pub data: Option<Json>,
111 #[builder(default)]
113 pub metadata: Option<Json>,
114 #[builder(default)]
117 pub timestamp: Option<DateTime<Utc>>,
118}
119
120#[derive(Debug, Clone, TypedBuilder)]
122#[builder(field_defaults(setter(strip_option(ignore_invalid, fallback_suffix = "_opt"))))]
123pub struct EndScopeHandleParams<'a> {
124 pub handle: &'a ScopeHandle,
126 #[builder(default)]
128 pub data: Option<Json>,
129 #[builder(default)]
131 pub metadata: Option<Json>,
132 #[builder(default)]
136 pub timestamp: Option<DateTime<Utc>>,
137}
138
139#[derive(TypedBuilder)]
141#[builder(field_defaults(setter(strip_option(ignore_invalid, fallback_suffix = "_opt"))))]
142pub struct PopScopeParams<'a> {
143 pub handle_uuid: &'a Uuid,
145 #[builder(default)]
147 pub output: Option<Json>,
148 #[builder(default)]
150 pub metadata: Option<Json>,
151 #[builder(default)]
155 pub timestamp: Option<DateTime<Utc>>,
156}
157
158#[derive(TypedBuilder)]
160#[builder(field_defaults(setter(strip_option(ignore_invalid, fallback_suffix = "_opt"))))]
161pub struct EmitMarkEventParams<'a> {
162 pub name: &'a str,
164 #[builder(default)]
166 pub parent: Option<&'a ScopeHandle>,
167 #[builder(default)]
169 pub data: Option<Json>,
170 #[builder(default)]
172 pub data_schema: Option<DataSchema>,
173 #[builder(default)]
175 pub metadata: Option<Json>,
176 #[builder(default)]
178 pub category: Option<EventCategory>,
179 #[builder(default)]
181 pub category_profile: Option<CategoryProfile>,
182 #[builder(default)]
185 pub timestamp: Option<DateTime<Utc>>,
186}
187
188pub fn get_handle() -> Result<ScopeHandle> {
201 ensure_runtime_owner()?;
202 Ok(task_scope_top())
203}
204
205pub fn push_scope(params: PushScopeParams<'_>) -> Result<ScopeHandle> {
234 ensure_runtime_owner()?;
235 let parent_uuid = resolve_parent_uuid(params.parent);
236 let (handle, event, subscribers, emission_scope_stack) = {
237 let scope_stack = current_scope_stack();
238 let scope_guard = scope_stack
239 .read()
240 .map_err(|error| scope_stack_lock_error(error, "push"))?;
241 let scope_subscribers = scope_guard.collect_scope_local_subscribers();
242 let subscribers = snapshot_event_subscribers(scope_subscribers)?;
243 let context = global_context();
244 let state = context
245 .read()
246 .map_err(|error| FlowError::Internal(error.to_string()))?;
247 let handle_params = CreateScopeHandleParams::builder()
248 .name(params.name)
249 .parent_uuid_opt(parent_uuid)
250 .scope_type(params.scope_type)
251 .attributes(params.attributes)
252 .data_opt(params.data)
253 .metadata_opt(params.metadata)
254 .timestamp_opt(params.timestamp)
255 .build();
256 let handle = state.create_scope_handle(handle_params);
257 let event = state.build_scope_start_event(&handle, params.input);
258 (handle, event, subscribers, scope_stack.clone())
259 };
260 task_scope_push(handle.clone());
261 let sanitizers = snapshot_event_sanitizers(&event, &emission_scope_stack).unwrap_or_default();
262 let _ = subscriber_dispatcher::dispatch_sanitized_event(
263 event,
264 sanitizers,
265 &subscribers,
266 emission_scope_stack,
267 );
268 Ok(handle)
269}
270
271pub fn pop_scope(params: PopScopeParams<'_>) -> Result<()> {
299 pop_scope_inner(params, false).map(|_| ())
300}
301
302#[doc(hidden)]
307pub fn pop_scope_with_subscriber_delivery(
308 params: PopScopeParams<'_>,
309) -> Result<SubscriberDelivery> {
310 pop_scope_inner(params, true)?.ok_or_else(|| {
311 FlowError::Internal("tracked scope pop did not create a subscriber delivery receipt".into())
312 })
313}
314
315fn pop_scope_inner(
316 params: PopScopeParams<'_>,
317 track_delivery: bool,
318) -> Result<Option<SubscriberDelivery>> {
319 ensure_runtime_owner()?;
320 let scope_stack = current_scope_stack();
321 let (scope, event, subscribers, emission_scope_stack) = {
322 let scope_guard = scope_stack
323 .read()
324 .map_err(|error| scope_stack_lock_error(error, "pop"))?;
325 let top = scope_guard.top();
326 if top.uuid != *params.handle_uuid {
327 if scope_guard.find(params.handle_uuid).is_some() {
328 return Err(FlowError::InvalidArgument(
329 "scope handle is not at the top of the stack".into(),
330 ));
331 }
332 return Err(FlowError::NotFound("scope handle not found".into()));
333 }
334 let scope_subscribers = scope_guard.collect_scope_local_subscribers();
335 let subscribers = snapshot_event_subscribers(scope_subscribers)?;
336 let scope = top.clone();
337 let context = global_context();
338 let state = context
339 .read()
340 .map_err(|error| FlowError::Internal(error.to_string()))?;
341 let event = state.build_scope_end_event(
342 EndScopeHandleParams::builder()
343 .handle(&scope)
344 .data_opt(params.output)
345 .timestamp_opt(params.timestamp)
346 .metadata_opt(params.metadata)
347 .build(),
348 );
349 (scope, event, subscribers, scope_stack.clone())
350 };
351 let sanitizers = snapshot_event_sanitizers(&event, &emission_scope_stack).unwrap_or_default();
355 let publication_scope_stack = snapshot_scope_stack(&emission_scope_stack)?;
356 let removed = task_scope_remove(params.handle_uuid)?;
357 debug_assert_eq!(removed.uuid, scope.uuid);
358 if track_delivery {
359 subscriber_dispatcher::dispatch_sanitized_event_with_delivery(
360 event,
361 sanitizers,
362 &subscribers,
363 publication_scope_stack,
364 )
365 .map(Some)
366 } else {
367 let _ = subscriber_dispatcher::dispatch_sanitized_event(
368 event,
369 sanitizers,
370 &subscribers,
371 publication_scope_stack,
372 );
373 Ok(None)
374 }
375}
376
377pub fn event(params: EmitMarkEventParams<'_>) -> Result<()> {
403 ensure_runtime_owner()?;
404 let parent_uuid = resolve_parent_uuid(params.parent);
405 let scope_stack = current_scope_stack();
406 let (event, subscribers, emission_scope_stack) = {
407 let subscribers = if params.name == COMPACTION_EVENT_NAME {
408 let mut scope_guard = scope_stack
409 .write()
410 .map_err(|error| scope_stack_lock_error(error, "mark"))?;
411 let subscribers =
412 snapshot_event_subscribers(scope_guard.collect_scope_local_subscribers())?;
413 scope_guard.mark_agent_fresh(parent_uuid);
414 subscribers
415 } else {
416 let scope_guard = scope_stack
417 .read()
418 .map_err(|error| scope_stack_lock_error(error, "mark"))?;
419 snapshot_event_subscribers(scope_guard.collect_scope_local_subscribers())?
420 };
421 let context = global_context();
422 let state = context
423 .read()
424 .map_err(|error| FlowError::Internal(error.to_string()))?;
425 let event = state.create_event(MarkEvent::new(
426 BaseEvent::builder()
427 .name(params.name)
428 .parent_uuid_opt(parent_uuid)
429 .timestamp(params.timestamp.unwrap_or_else(Utc::now))
430 .data_opt(params.data)
431 .data_schema_opt(params.data_schema)
432 .metadata_opt(params.metadata)
433 .build(),
434 params.category,
435 params.category_profile,
436 ));
437 (event, subscribers, scope_stack.clone())
438 };
439 let sanitizers = snapshot_event_sanitizers(&event, &emission_scope_stack).unwrap_or_default();
440 let _ = subscriber_dispatcher::dispatch_sanitized_event(
441 event,
442 sanitizers,
443 &subscribers,
444 emission_scope_stack,
445 );
446 Ok(())
447}