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;
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 ensure_runtime_owner()?;
300 let scope_stack = current_scope_stack();
301 let (scope, event, subscribers, emission_scope_stack) = {
302 let scope_guard = scope_stack
303 .read()
304 .map_err(|error| scope_stack_lock_error(error, "pop"))?;
305 let top = scope_guard.top();
306 if top.uuid != *params.handle_uuid {
307 if scope_guard.find(params.handle_uuid).is_some() {
308 return Err(FlowError::InvalidArgument(
309 "scope handle is not at the top of the stack".into(),
310 ));
311 }
312 return Err(FlowError::NotFound("scope handle not found".into()));
313 }
314 let scope_subscribers = scope_guard.collect_scope_local_subscribers();
315 let subscribers = snapshot_event_subscribers(scope_subscribers)?;
316 let scope = top.clone();
317 let context = global_context();
318 let state = context
319 .read()
320 .map_err(|error| FlowError::Internal(error.to_string()))?;
321 let event = state.build_scope_end_event(
322 EndScopeHandleParams::builder()
323 .handle(&scope)
324 .data_opt(params.output)
325 .timestamp_opt(params.timestamp)
326 .metadata_opt(params.metadata)
327 .build(),
328 );
329 (scope, event, subscribers, scope_stack.clone())
330 };
331 let sanitizers = snapshot_event_sanitizers(&event, &emission_scope_stack).unwrap_or_default();
335 let publication_scope_stack = snapshot_scope_stack(&emission_scope_stack)?;
336 let removed = task_scope_remove(params.handle_uuid)?;
337 debug_assert_eq!(removed.uuid, scope.uuid);
338 let _ = subscriber_dispatcher::dispatch_sanitized_event(
339 event,
340 sanitizers,
341 &subscribers,
342 publication_scope_stack,
343 );
344 Ok(())
345}
346
347pub fn event(params: EmitMarkEventParams<'_>) -> Result<()> {
373 ensure_runtime_owner()?;
374 let parent_uuid = resolve_parent_uuid(params.parent);
375 let scope_stack = current_scope_stack();
376 let (event, subscribers, emission_scope_stack) = {
377 let subscribers = if params.name == COMPACTION_EVENT_NAME {
378 let mut scope_guard = scope_stack
379 .write()
380 .map_err(|error| scope_stack_lock_error(error, "mark"))?;
381 let subscribers =
382 snapshot_event_subscribers(scope_guard.collect_scope_local_subscribers())?;
383 scope_guard.mark_agent_fresh(parent_uuid);
384 subscribers
385 } else {
386 let scope_guard = scope_stack
387 .read()
388 .map_err(|error| scope_stack_lock_error(error, "mark"))?;
389 snapshot_event_subscribers(scope_guard.collect_scope_local_subscribers())?
390 };
391 let context = global_context();
392 let state = context
393 .read()
394 .map_err(|error| FlowError::Internal(error.to_string()))?;
395 let event = state.create_event(MarkEvent::new(
396 BaseEvent::builder()
397 .name(params.name)
398 .parent_uuid_opt(parent_uuid)
399 .timestamp(params.timestamp.unwrap_or_else(Utc::now))
400 .data_opt(params.data)
401 .data_schema_opt(params.data_schema)
402 .metadata_opt(params.metadata)
403 .build(),
404 params.category,
405 params.category_profile,
406 ));
407 (event, subscribers, scope_stack.clone())
408 };
409 let sanitizers = snapshot_event_sanitizers(&event, &emission_scope_stack).unwrap_or_default();
410 let _ = subscriber_dispatcher::dispatch_sanitized_event(
411 event,
412 sanitizers,
413 &subscribers,
414 emission_scope_stack,
415 );
416 Ok(())
417}