From 7bbdd9066f0b48ff29c5f2808cdfcdac72e03bb6 Mon Sep 17 00:00:00 2001 From: Alex Fournier Date: Tue, 29 Sep 2026 11:18:29 -0700 Subject: [PATCH 1/2] feat(plugin): add timestamped scope guards Signed-off-by: Alex Fournier --- crates/plugin/src/lib.rs | 103 +++++++++++++++++- crates/plugin/tests/typed_callbacks.rs | 140 ++++++++++++++++++++++++- 2 files changed, 237 insertions(+), 6 deletions(-) diff --git a/crates/plugin/src/lib.rs b/crates/plugin/src/lib.rs index 369e60f75..627383b55 100644 --- a/crates/plugin/src/lib.rs +++ b/crates/plugin/src/lib.rs @@ -18,6 +18,7 @@ use std::marker::{PhantomData, PhantomPinned}; use std::panic::{AssertUnwindSafe, catch_unwind}; use std::ptr; use std::sync::{Arc, Mutex}; +use std::time::{SystemTime, UNIX_EPOCH}; pub use nemo_relay_types::Json; pub use nemo_relay_types::api::event::{ @@ -1790,6 +1791,36 @@ impl PluginRuntime { }) } + /// Opens a scope and records `started_at` on its start event. + /// + /// This is typed SDK access to the timestamp slot already carried by the + /// native host's `scope_push` function. It does not introduce a distinct + /// scope event or change the native ABI. + pub fn scope_at( + &self, + name: &str, + scope_type: ScopeType, + data: Option<&Json>, + metadata: Option<&Json>, + input: Option<&Json>, + started_at: SystemTime, + ) -> Result> { + let timestamp = unix_micros(started_at)?; + let handle = push_scope_with_timestamp( + &self.host, + name, + scope_type.into(), + data, + metadata, + input, + Some(timestamp), + )?; + Ok(ScopeGuard { + runtime: self, + handle: Some(handle), + }) + } + /// Emits a mark event under the current scope. pub fn emit_mark( &self, @@ -1911,7 +1942,8 @@ impl From for NemoRelayNativeScopeType { } } -/// RAII guard for a host scope opened by [`PluginRuntime::scope`]. +/// RAII guard for a host scope opened by [`PluginRuntime::scope`] or +/// [`PluginRuntime::scope_at`]. /// /// A guard may move between threads only while its scope stack is bound on the /// destination thread. Async middleware restores that binding around each poll; @@ -1937,6 +1969,33 @@ impl<'a> ScopeGuard<'a> { self.handle.take(); Ok(()) } + + /// Pops the scope and records `ended_at` on its end event. + /// + /// This is typed SDK access to the timestamp slot already carried by the + /// native host's `scope_pop` function. The handle remains owned by this + /// guard if the host rejects the close, so [`Drop`] can still attempt the + /// ordinary cleanup path. + pub fn close_at( + &mut self, + output: Option<&Json>, + metadata: Option<&Json>, + ended_at: SystemTime, + ) -> Result<()> { + let Some(handle) = self.handle.as_ref() else { + return Ok(()); + }; + let timestamp = unix_micros(ended_at)?; + pop_scope_with_timestamp( + &self.runtime.host, + handle, + output, + metadata, + Some(timestamp), + )?; + self.handle.take(); + Ok(()) + } } impl Drop for ScopeGuard<'_> { @@ -2255,6 +2314,18 @@ pub fn push_scope<'a>( data: Option<&Json>, metadata: Option<&Json>, input: Option<&Json>, +) -> Result> { + push_scope_with_timestamp(host, name, scope_type, data, metadata, input, None) +} + +fn push_scope_with_timestamp<'a>( + host: &'a NemoRelayNativeHostApiV1, + name: &str, + scope_type: NemoRelayNativeScopeType, + data: Option<&Json>, + metadata: Option<&Json>, + input: Option<&Json>, + timestamp: Option, ) -> Result> { let name = HostString::new(host, name).ok_or_else(|| "failed to allocate scope name".to_string())?; @@ -2271,7 +2342,7 @@ pub fn push_scope<'a>( data.as_ptr(), metadata.as_ptr(), input.as_ptr(), - ptr::null(), + timestamp.as_ref().map_or(ptr::null(), ptr::from_ref), &mut out, ) }; @@ -2288,6 +2359,16 @@ pub fn pop_scope( handle: &ScopeHandle<'_>, output: Option<&Json>, metadata: Option<&Json>, +) -> Result<()> { + pop_scope_with_timestamp(host, handle, output, metadata, None) +} + +fn pop_scope_with_timestamp( + host: &NemoRelayNativeHostApiV1, + handle: &ScopeHandle<'_>, + output: Option<&Json>, + metadata: Option<&Json>, + timestamp: Option, ) -> Result<()> { let output = OptionalHostJson::new(host, output)?; let metadata = OptionalHostJson::new(host, metadata)?; @@ -2296,7 +2377,7 @@ pub fn pop_scope( handle.as_ptr(), output.as_ptr(), metadata.as_ptr(), - ptr::null(), + timestamp.as_ref().map_or(ptr::null(), ptr::from_ref), ) }; if status == NemoRelayStatus::Ok { @@ -2306,6 +2387,22 @@ pub fn pop_scope( } } +fn unix_micros(timestamp: SystemTime) -> Result { + let micros = match timestamp.duration_since(UNIX_EPOCH) { + Ok(duration) => i128::try_from(duration.as_micros()), + Err(error) => { + let duration = error.duration(); + i128::try_from(duration.as_micros()).map(|micros| { + // `Duration::as_micros` truncates toward zero, but a signed + // Unix timestamp must floor pre-epoch sub-microsecond values. + -micros - i128::from(duration.subsec_nanos() % 1_000 != 0) + }) + } + } + .map_err(|_| "scope timestamp exceeds the supported range".to_string())?; + i64::try_from(micros).map_err(|_| "scope timestamp exceeds the supported range".to_string()) +} + /// Emits a mark event under the current scope. pub fn emit_mark( host: &NemoRelayNativeHostApiV1, diff --git a/crates/plugin/tests/typed_callbacks.rs b/crates/plugin/tests/typed_callbacks.rs index ae6d44a78..0a9680ada 100644 --- a/crates/plugin/tests/typed_callbacks.rs +++ b/crates/plugin/tests/typed_callbacks.rs @@ -17,7 +17,7 @@ use std::sync::{ atomic::{AtomicBool, AtomicUsize, Ordering}, }; use std::task::Poll; -use std::time::{Duration, Instant}; +use std::time::{Duration, Instant, UNIX_EPOCH}; use futures::StreamExt; use nemo_relay_plugin::{ @@ -424,6 +424,7 @@ static SCOPE_GET_CURRENT_RETURNS_NULL: Mutex = Mutex::new(false); static SCOPE_PUSH_STATUS: Mutex = Mutex::new(NemoRelayStatus::Ok); static SCOPE_PUSH_RETURNS_NULL: Mutex = Mutex::new(false); static SCOPE_POP_STATUS: Mutex = Mutex::new(NemoRelayStatus::Ok); +static SCOPE_TIMESTAMPS: Mutex<(Option, Option)> = Mutex::new((None, None)); static EMIT_MARK_STATUS: Mutex = Mutex::new(NemoRelayStatus::Ok); static SCOPE_STACK_CREATE_STATUS: Mutex = Mutex::new(NemoRelayStatus::Ok); static SCOPE_STACK_CREATE_RETURNS_NULL: Mutex = Mutex::new(false); @@ -1406,7 +1407,7 @@ unsafe extern "C" fn capture_scope_push( data_json: *const NemoRelayNativeString, metadata_json: *const NemoRelayNativeString, input_json: *const NemoRelayNativeString, - _timestamp_unix_micros: *const i64, + timestamp_unix_micros: *const i64, out: *mut *mut NemoRelayNativeScopeHandle, ) -> NemoRelayStatus { if out.is_null() { @@ -1437,6 +1438,9 @@ unsafe extern "C" fn capture_scope_push( "push:{name}:{scope_type:?}:{attributes}:parent={}:data={data}:metadata={metadata}:input={input}", !parent.is_null() )); + if !timestamp_unix_micros.is_null() { + SCOPE_TIMESTAMPS.lock().unwrap().0 = Some(unsafe { *timestamp_unix_micros }); + } if *SCOPE_PUSH_RETURNS_NULL.lock().unwrap() { unsafe { *out = ptr::null_mut() }; } else { @@ -1449,7 +1453,7 @@ unsafe extern "C" fn capture_scope_pop( handle: *const NemoRelayNativeScopeHandle, output_json: *const NemoRelayNativeString, metadata_json: *const NemoRelayNativeString, - _timestamp_unix_micros: *const i64, + timestamp_unix_micros: *const i64, ) -> NemoRelayStatus { if handle.is_null() { return NemoRelayStatus::NullPointer; @@ -1471,6 +1475,9 @@ unsafe extern "C" fn capture_scope_pop( .lock() .unwrap() .push(format!("pop:output={output}:metadata={metadata}")); + if !timestamp_unix_micros.is_null() { + SCOPE_TIMESTAMPS.lock().unwrap().1 = Some(unsafe { *timestamp_unix_micros }); + } NemoRelayStatus::Ok } @@ -2930,6 +2937,7 @@ fn reset_state() { *SCOPE_PUSH_STATUS.lock().unwrap() = NemoRelayStatus::Ok; *SCOPE_PUSH_RETURNS_NULL.lock().unwrap() = false; *SCOPE_POP_STATUS.lock().unwrap() = NemoRelayStatus::Ok; + *SCOPE_TIMESTAMPS.lock().unwrap() = (None, None); *EMIT_MARK_STATUS.lock().unwrap() = NemoRelayStatus::Ok; *SCOPE_STACK_CREATE_STATUS.lock().unwrap() = NemoRelayStatus::Ok; *SCOPE_STACK_CREATE_RETURNS_NULL.lock().unwrap() = false; @@ -3293,6 +3301,132 @@ fn plugin_runtime_scope_mark_and_stack_helpers_call_host() { assert_eq!(SCOPE_STACK_BINDING_FREES.load(Ordering::SeqCst), 0); } +#[test] +fn plugin_runtime_forwards_historical_scope_timestamps() { + let _guard = begin_test(); + let host = test_host(); + let runtime = PluginRuntime::new(&host); + let cases = [ + ( + UNIX_EPOCH + Duration::from_micros(1_000_000), + UNIX_EPOCH + Duration::from_micros(1_250_000), + (Some(1_000_000), Some(1_250_000)), + ), + ( + UNIX_EPOCH - Duration::from_micros(1_250_000), + UNIX_EPOCH - Duration::from_micros(1_000_000), + (Some(-1_250_000), Some(-1_000_000)), + ), + ( + UNIX_EPOCH - Duration::from_nanos(1_500), + UNIX_EPOCH - Duration::from_nanos(500), + (Some(-2), Some(-1)), + ), + ]; + + for (started_at, ended_at, expected) in cases { + let mut scope = runtime + .scope_at( + "historical", + ScopeType::Custom, + None, + None, + None, + started_at, + ) + .unwrap(); + scope.close_at(None, None, ended_at).unwrap(); + assert_eq!(*SCOPE_TIMESTAMPS.lock().unwrap(), expected); + } +} + +#[test] +fn ordinary_scope_helpers_leave_native_timestamps_unset() { + let _guard = begin_test(); + let host = test_host(); + let runtime = PluginRuntime::new(&host); + + let mut scope = runtime + .scope("current-time", ScopeType::Custom, None, None, None) + .unwrap(); + scope.close(None, None).unwrap(); + + assert_eq!(*SCOPE_TIMESTAMPS.lock().unwrap(), (None, None)); +} + +#[test] +fn historical_scope_close_retains_ownership_after_failure_and_is_idempotent() { + let _guard = begin_test(); + let host = test_host(); + let runtime = PluginRuntime::new(&host); + let mut scope = runtime + .scope_at( + "historical", + ScopeType::Custom, + None, + None, + None, + UNIX_EPOCH + Duration::from_micros(10), + ) + .unwrap(); + + *SCOPE_POP_STATUS.lock().unwrap() = NemoRelayStatus::Internal; + assert_eq!( + scope + .close_at(None, None, UNIX_EPOCH + Duration::from_micros(20)) + .unwrap_err(), + "scope_pop failed: Internal" + ); + assert!(scope.handle().is_some()); + + *SCOPE_POP_STATUS.lock().unwrap() = NemoRelayStatus::Ok; + scope + .close_at(None, None, UNIX_EPOCH + Duration::from_micros(30)) + .unwrap(); + assert!(scope.handle().is_none()); + scope + .close_at(None, None, UNIX_EPOCH + Duration::from_micros(40)) + .unwrap(); + + assert_eq!( + RUNTIME_CALLS + .lock() + .unwrap() + .iter() + .filter(|call| call.starts_with("pop:")) + .count(), + 1 + ); + assert_eq!(*SCOPE_TIMESTAMPS.lock().unwrap(), (Some(10), Some(30))); +} + +#[test] +fn historical_scope_rejects_timestamps_outside_native_range() { + let _guard = begin_test(); + let host = test_host(); + let runtime = PluginRuntime::new(&host); + let Some(outside_native_range) = + UNIX_EPOCH.checked_add(Duration::from_micros(i64::MAX as u64 + 1)) + else { + // Windows FILETIME cannot represent a SystemTime this far after the + // epoch, so the public API cannot receive this overflow case there. + return; + }; + + assert_eq!( + expect_string_err(runtime.scope_at( + "historical", + ScopeType::Custom, + None, + None, + None, + outside_native_range, + )), + "scope timestamp exceeds the supported range" + ); + assert!(RUNTIME_CALLS.lock().unwrap().is_empty()); +} + #[test] fn plugin_runtime_logs_every_level_and_requires_native_abi_v6() { let _guard = begin_test(); From 7fe4b8a26f01b2116a395a6f4b41e416f31c50f4 Mon Sep 17 00:00:00 2001 From: Alex Fournier Date: Thu, 1 Oct 2026 11:18:34 -0700 Subject: [PATCH 2/2] feat(plugin): align timestamped worker scopes Signed-off-by: Alex Fournier --- crates/core/src/plugin/dynamic/worker.rs | 19 +++ .../core/tests/unit/dynamic_worker_tests.rs | 120 +++++++++++++++++- .../nemo/relay/worker/v1/plugin_worker.proto | 2 + crates/worker-proto/tests/proto_tests.rs | 38 +++++- crates/worker/src/lib.rs | 99 +++++++++++++++ crates/worker/tests/unit/timestamp_tests.rs | 49 +++++++ crates/worker/tests/worker_sdk_tests.rs | 57 ++++++++- .../native/runtime-events-and-scopes.mdx | 3 + .../workers/grpc-v1-protocol.mdx | 14 +- .../workers/runtime-events-and-scopes.mdx | 3 +- python/plugin/src/nemo_relay_plugin/_api.py | 72 +++++++---- python/tests/plugin/test_worker_sdk.py | 43 +++++++ 12 files changed, 486 insertions(+), 33 deletions(-) create mode 100644 crates/worker/tests/unit/timestamp_tests.rs diff --git a/crates/core/src/plugin/dynamic/worker.rs b/crates/core/src/plugin/dynamic/worker.rs index 622a379d9..a3cb641b1 100644 --- a/crates/core/src/plugin/dynamic/worker.rs +++ b/crates/core/src/plugin/dynamic/worker.rs @@ -13,6 +13,7 @@ use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Condvar, Mutex}; use std::time::Duration; +use chrono::{DateTime, Utc}; use futures_util::FutureExt; use nemo_relay_worker_proto::v1::plugin_worker_client::PluginWorkerClient; use nemo_relay_worker_proto::v1::relay_host_runtime_server::{ @@ -3148,6 +3149,7 @@ impl RelayHostRuntime for WorkerHostRuntimeService { .data_opt(optional_envelope_to_json(request.data)?) .metadata_opt(optional_envelope_to_json(request.metadata)?) .input_opt(optional_envelope_to_json(request.input)?) + .timestamp_opt(optional_worker_timestamp(request.timestamp_unix_micros)?) .build(), ) }); @@ -3189,6 +3191,10 @@ impl RelayHostRuntime for WorkerHostRuntimeService { let request = request.into_inner(); self.state .authorize(&request.activation_id, &request.auth_token)?; + let timestamp = match optional_worker_timestamp(request.timestamp_unix_micros) { + Ok(timestamp) => timestamp, + Err(err) => return Ok(Response::new(host_ack(Err(err)))), + }; let handle = self .state .scope_handles @@ -3204,6 +3210,7 @@ impl RelayHostRuntime for WorkerHostRuntimeService { .handle_uuid(&handle.handle.uuid) .output_opt(output) .metadata_opt(metadata) + .timestamp_opt(timestamp) .build(), ) }; @@ -3725,6 +3732,18 @@ fn optional_envelope_to_json(value: Option) -> FlowResult) -> FlowResult>> { + value + .map(|timestamp| { + DateTime::::from_timestamp_micros(timestamp).ok_or_else(|| { + FlowError::InvalidArgument( + "timestamp unix microseconds are outside supported range".into(), + ) + }) + }) + .transpose() +} + fn optional_typed_envelope( value: Option, field: &str, diff --git a/crates/core/tests/unit/dynamic_worker_tests.rs b/crates/core/tests/unit/dynamic_worker_tests.rs index a2e9b99d3..3c11ae4d7 100644 --- a/crates/core/tests/unit/dynamic_worker_tests.rs +++ b/crates/core/tests/unit/dynamic_worker_tests.rs @@ -7,7 +7,7 @@ use std::sync::{Arc, Mutex}; #[cfg(unix)] use std::os::unix::fs::PermissionsExt; -use crate::api::event::{BaseEvent, MarkEvent}; +use crate::api::event::{BaseEvent, Event, MarkEvent, ScopeCategory}; use crate::api::optimization::{ LlmOptimizationRecorder, record_llm_optimization_contribution, scope_llm_optimization_recorder, }; @@ -1480,6 +1480,7 @@ async fn dropping_callback_future_cancels_worker_and_cleans_host_state() { data: None, metadata: None, input: None, + timestamp_unix_micros: None, })) .await .expect("worker scope should push") @@ -2188,6 +2189,7 @@ async fn host_runtime_service_covers_auth_scope_and_ack_errors() { }), metadata: None, input: None, + timestamp_unix_micros: None, })) .await .expect("invalid JSON should be structured") @@ -2199,6 +2201,119 @@ async fn host_runtime_service_covers_auth_scope_and_ack_errors() { .contains("invalid JSON") ); + let historical_events = Arc::new(Mutex::new(Vec::::new())); + let captured_historical_events = Arc::clone(&historical_events); + crate::api::subscriber::register_subscriber( + "worker-historical-scope-timestamps", + Arc::new(move |event| { + if event.name() == "historical-scope" { + captured_historical_events + .lock() + .expect("historical events lock") + .push(event.clone()); + } + }), + ) + .expect("historical timestamp subscriber should register"); + + let historical_push = service + .push_scope(Request::new(PushScopeRequest { + activation_id: ACTIVATION_ID.into(), + auth_token: AUTH_TOKEN.into(), + scope: None, + name: "historical-scope".into(), + scope_type: ProtoScopeType::Custom as i32, + data: None, + metadata: None, + input: None, + timestamp_unix_micros: Some(-2), + })) + .await + .expect("historical scope should push") + .into_inner(); + assert!(historical_push.error.is_none()); + let historical_handle_id = historical_push.scope_handle_id; + assert_eq!( + state + .scope_handles + .lock() + .expect("scope handles lock") + .get(&historical_handle_id) + .expect("historical scope handle") + .handle + .started_at + .timestamp_micros(), + -2 + ); + + let invalid_pop = service + .pop_scope(Request::new(PopScopeRequest { + activation_id: ACTIVATION_ID.into(), + auth_token: AUTH_TOKEN.into(), + scope_handle_id: historical_handle_id.clone(), + output: None, + metadata: None, + timestamp_unix_micros: Some(i64::MAX), + })) + .await + .expect("invalid timestamp should return a host ack") + .into_inner(); + assert!(!invalid_pop.ok); + assert!( + invalid_pop + .error + .expect("invalid timestamp error") + .message + .contains("outside supported range") + ); + assert!( + state + .scope_handles + .lock() + .expect("scope handles lock") + .contains_key(&historical_handle_id), + "an invalid timestamp must not consume the pop handle" + ); + + let historical_pop = service + .pop_scope(Request::new(PopScopeRequest { + activation_id: ACTIVATION_ID.into(), + auth_token: AUTH_TOKEN.into(), + scope_handle_id: historical_handle_id.clone(), + output: None, + metadata: None, + timestamp_unix_micros: Some(0), + })) + .await + .expect("epoch timestamp should pop") + .into_inner(); + assert!(historical_pop.ok, "{:?}", historical_pop.error); + assert!( + !state + .scope_handles + .lock() + .expect("scope handles lock") + .contains_key(&historical_handle_id) + ); + crate::api::subscriber::flush_subscribers().expect("historical timestamp events should flush"); + assert!( + crate::api::subscriber::deregister_subscriber("worker-historical-scope-timestamps") + .expect("historical timestamp subscriber should deregister") + ); + { + let historical_events = historical_events.lock().expect("historical events lock"); + let historical_start = historical_events + .iter() + .find(|event| event.scope_category() == Some(ScopeCategory::Start)) + .expect("historical start event"); + let historical_end = historical_events + .iter() + .find(|event| event.scope_category() == Some(ScopeCategory::End)) + .expect("historical end event"); + assert_eq!(historical_start.timestamp().timestamp_micros(), -2); + assert_eq!(historical_end.timestamp().timestamp_micros(), 0); + } + let pop_error = service .pop_scope(Request::new(PopScopeRequest { activation_id: ACTIVATION_ID.into(), @@ -2206,6 +2321,7 @@ async fn host_runtime_service_covers_auth_scope_and_ack_errors() { scope_handle_id: "missing-scope".into(), output: None, metadata: None, + timestamp_unix_micros: None, })) .await .expect_err("missing scope handle should fail"); @@ -2443,6 +2559,7 @@ async fn host_runtime_service_reports_poisoned_internal_locks() { data: None, metadata: None, input: None, + timestamp_unix_micros: None, })) .await .expect_err("poisoned scope handle lock should fail"); @@ -2455,6 +2572,7 @@ async fn host_runtime_service_reports_poisoned_internal_locks() { scope_handle_id: "missing".into(), output: None, metadata: None, + timestamp_unix_micros: None, })) .await .expect_err("poisoned scope handle lock should fail"); diff --git a/crates/worker-proto/proto/nemo/relay/worker/v1/plugin_worker.proto b/crates/worker-proto/proto/nemo/relay/worker/v1/plugin_worker.proto index 7aa990e31..ec75f3b12 100644 --- a/crates/worker-proto/proto/nemo/relay/worker/v1/plugin_worker.proto +++ b/crates/worker-proto/proto/nemo/relay/worker/v1/plugin_worker.proto @@ -436,6 +436,7 @@ message PushScopeRequest { JsonEnvelope data = 6; JsonEnvelope metadata = 7; JsonEnvelope input = 8; + optional int64 timestamp_unix_micros = 9; } message PushScopeResponse { @@ -449,6 +450,7 @@ message PopScopeRequest { string scope_handle_id = 3; JsonEnvelope output = 4; JsonEnvelope metadata = 5; + optional int64 timestamp_unix_micros = 6; } message CreateScopeStackRequest { diff --git a/crates/worker-proto/tests/proto_tests.rs b/crates/worker-proto/tests/proto_tests.rs index 087abfeff..c4832af9c 100644 --- a/crates/worker-proto/tests/proto_tests.rs +++ b/crates/worker-proto/tests/proto_tests.rs @@ -6,9 +6,9 @@ use nemo_relay_worker_proto::v1::{ ConditionalMiddlewareGuardrailRegistration, ConditionalMiddlewareInvocation, EmitMarkRequest, GetRuntimeDiagnosticsRequest, GetRuntimeDiagnosticsResponse, HandshakeRequest, HealthRequest, - InvokeRequest, JsonEnvelope, JsonValue, RegisterConditionalMiddlewareGuardrailRequest, - RegistrationSurface, RuntimeDiagnostic, ScopeType, - ToolExecutionResult as ProtoToolExecutionResult, invoke_request, + InvokeRequest, JsonEnvelope, JsonValue, PopScopeRequest, PushScopeRequest, + RegisterConditionalMiddlewareGuardrailRequest, RegistrationSurface, RuntimeDiagnostic, + ScopeType, ToolExecutionResult as ProtoToolExecutionResult, invoke_request, }; use nemo_relay_worker_proto::{ WORKER_PROTOCOL_GRPC_V1, decode_json_envelope, decode_json_value, json_envelope, json_value, @@ -243,6 +243,38 @@ fn tool_execution_result_tolerates_unknown_protobuf_fields() { ); } +#[test] +fn scope_timestamps_are_additive_and_presence_aware() { + let legacy_push = PushScopeRequest::decode([].as_slice()).expect("decode legacy push"); + let legacy_pop = PopScopeRequest::decode([].as_slice()).expect("decode legacy pop"); + assert_eq!(legacy_push.timestamp_unix_micros, None); + assert_eq!(legacy_pop.timestamp_unix_micros, None); + + let epoch_push = PushScopeRequest { + timestamp_unix_micros: Some(0), + ..PushScopeRequest::default() + }; + assert_eq!(epoch_push.encode_to_vec(), vec![0x48, 0x00]); + assert_eq!( + PushScopeRequest::decode(epoch_push.encode_to_vec().as_slice()) + .expect("decode epoch push") + .timestamp_unix_micros, + Some(0) + ); + + let historical_pop = PopScopeRequest { + timestamp_unix_micros: Some(-1_250_000), + ..PopScopeRequest::default() + }; + assert_eq!(historical_pop.encode_to_vec()[0], 0x30); + assert_eq!( + PopScopeRequest::decode(historical_pop.encode_to_vec().as_slice()) + .expect("decode historical pop") + .timestamp_unix_micros, + Some(-1_250_000) + ); +} + #[test] fn emit_mark_additive_fields_preserve_legacy_wire_compatibility() { let legacy = EmitMarkRequest::decode(b"\x22\x04mark".as_slice()).expect("decode legacy mark"); diff --git a/crates/worker/src/lib.rs b/crates/worker/src/lib.rs index 953f384c3..f960dd7a1 100644 --- a/crates/worker/src/lib.rs +++ b/crates/worker/src/lib.rs @@ -29,6 +29,7 @@ use std::pin::Pin; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; use std::task::{Context, Poll}; +use std::time::{SystemTime, UNIX_EPOCH}; use futures_util::{Stream, StreamExt}; #[cfg(unix)] @@ -1439,6 +1440,54 @@ impl PluginRuntime { data: Option, metadata: Option, input: Option, + ) -> Result { + self.push_scope_with_timestamp( + scope_stack_id, + name, + scope_type, + data, + metadata, + input, + None, + ) + .await + } + + /// Pushes a scope through the host runtime using an explicit start time. + #[allow(clippy::too_many_arguments)] // Mirrors `push_scope` with one additive timestamp. + pub async fn push_scope_at( + &self, + scope_stack_id: Option<&str>, + name: &str, + scope_type: ScopeType, + data: Option, + metadata: Option, + input: Option, + started_at: SystemTime, + ) -> Result { + let timestamp = unix_micros(started_at)?; + self.push_scope_with_timestamp( + scope_stack_id, + name, + scope_type, + data, + metadata, + input, + Some(timestamp), + ) + .await + } + + #[allow(clippy::too_many_arguments)] // Shared implementation for the parallel public methods. + async fn push_scope_with_timestamp( + &self, + scope_stack_id: Option<&str>, + name: &str, + scope_type: ScopeType, + data: Option, + metadata: Option, + input: Option, + timestamp_unix_micros: Option, ) -> Result { let scope = scope_stack_id .map(scope_context) @@ -1454,6 +1503,7 @@ impl PluginRuntime { data: optional_json_envelope(data)?, metadata: optional_json_envelope(metadata)?, input: optional_json_envelope(input)?, + timestamp_unix_micros, })) .await .map_err(|err| WorkerSdkError::Transport(err.to_string()))? @@ -1470,6 +1520,30 @@ impl PluginRuntime { scope_handle_id: &str, output: Option, metadata: Option, + ) -> Result<()> { + self.pop_scope_with_timestamp(scope_handle_id, output, metadata, None) + .await + } + + /// Pops a scope through the host runtime using an explicit end time. + pub async fn pop_scope_at( + &self, + scope_handle_id: &str, + output: Option, + metadata: Option, + ended_at: SystemTime, + ) -> Result<()> { + let timestamp = unix_micros(ended_at)?; + self.pop_scope_with_timestamp(scope_handle_id, output, metadata, Some(timestamp)) + .await + } + + async fn pop_scope_with_timestamp( + &self, + scope_handle_id: &str, + output: Option, + metadata: Option, + timestamp_unix_micros: Option, ) -> Result<()> { let mut client = self.host_client().await?; let response = client @@ -1479,6 +1553,7 @@ impl PluginRuntime { scope_handle_id: scope_handle_id.into(), output: optional_json_envelope(output)?, metadata: optional_json_envelope(metadata)?, + timestamp_unix_micros, })) .await .map_err(|err| WorkerSdkError::Transport(err.to_string()))? @@ -1499,6 +1574,26 @@ impl PluginRuntime { } } +fn unix_micros(timestamp: SystemTime) -> Result { + let micros = match timestamp.duration_since(UNIX_EPOCH) { + Ok(duration) => i128::try_from(duration.as_micros()), + Err(error) => { + let duration = error.duration(); + i128::try_from(duration.as_micros()).map(|micros| { + // `Duration::as_micros` truncates toward zero, but a signed + // Unix timestamp must floor pre-epoch sub-microsecond values. + -micros - i128::from(duration.subsec_nanos() % 1_000 != 0) + }) + } + } + .map_err(|_| { + WorkerSdkError::InvalidInput("scope timestamp exceeds the supported range".into()) + })?; + i64::try_from(micros).map_err(|_| { + WorkerSdkError::InvalidInput("scope timestamp exceeds the supported range".into()) + }) +} + /// Explicit worker server configuration for tests and custom launchers. #[derive(Debug, Clone)] pub struct WorkerServerConfig { @@ -3581,3 +3676,7 @@ fn rustc_version_runtime() -> String { #[cfg(test)] #[path = "../tests/unit/codec_identity_tests.rs"] mod codec_identity_tests; + +#[cfg(test)] +#[path = "../tests/unit/timestamp_tests.rs"] +mod timestamp_tests; diff --git a/crates/worker/tests/unit/timestamp_tests.rs b/crates/worker/tests/unit/timestamp_tests.rs new file mode 100644 index 000000000..2a0a66cd3 --- /dev/null +++ b/crates/worker/tests/unit/timestamp_tests.rs @@ -0,0 +1,49 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +use std::time::Duration; + +use super::*; + +#[tokio::test] +async fn historical_scope_rejects_timestamp_outside_wire_range_before_rpc() { + let Some(outside_wire_range) = + UNIX_EPOCH.checked_add(Duration::from_micros(i64::MAX as u64 + 1)) + else { + // Windows FILETIME cannot represent a SystemTime this far after the + // epoch, so the public API cannot receive this overflow case there. + return; + }; + let runtime = PluginRuntime { + activation_id: "activation".into(), + auth_token: "token".into(), + host_endpoint: "unsupported://host".into(), + host_channel: Arc::new(OnceCell::new()), + conditional_middleware_callbacks: Arc::new(Mutex::new(HashMap::new())), + }; + + let push_error = runtime + .push_scope_at( + None, + "historical", + ScopeType::Custom, + None, + None, + None, + outside_wire_range, + ) + .await + .expect_err("out-of-range scope start should fail before the host call"); + let pop_error = runtime + .pop_scope_at("scope-handle", None, None, outside_wire_range) + .await + .expect_err("out-of-range scope end should fail before the host call"); + + for error in [push_error, pop_error] { + assert!(matches!(error, WorkerSdkError::InvalidInput(_))); + assert_eq!( + error.to_string(), + "invalid input: scope timestamp exceeds the supported range" + ); + } +} diff --git a/crates/worker/tests/worker_sdk_tests.rs b/crates/worker/tests/worker_sdk_tests.rs index ad6ee60a4..a6bb6718c 100644 --- a/crates/worker/tests/worker_sdk_tests.rs +++ b/crates/worker/tests/worker_sdk_tests.rs @@ -1083,6 +1083,25 @@ async fn worker_service_invokes_every_registration_surface() { assert!(calls.contains(&"mark:stream-poll:stack-1:parent-1".into())); assert!(calls.contains(&"push:scope-agent:explicit-stack:".into())); assert!(calls.contains(&"push:scope-unknown:explicit-stack:".into())); + let scope_requests = host.scope_requests(); + let historical_push = scope_requests + .push + .iter() + .find(|request| request.name == "worker-scope") + .expect("historical scope push"); + assert_eq!(historical_push.timestamp_unix_micros, Some(-2)); + let historical_pop = scope_requests + .pop + .iter() + .find(|request| request.scope_handle_id == "scope-handle-1") + .expect("historical scope pop"); + assert_eq!(historical_pop.timestamp_unix_micros, Some(0)); + let ordinary_push = scope_requests + .push + .iter() + .find(|request| request.name == "scope-agent") + .expect("ordinary scope push"); + assert_eq!(ordinary_push.timestamp_unix_micros, None); let telemetry_mark = host .marks() .into_iter() @@ -2381,9 +2400,19 @@ impl WorkerPlugin for SurfacePlugin { .await?; runtime.emit_mark("tool-exec-restored", None, None).await?; let handle = runtime - .push_scope(None, "worker-scope", ScopeType::Function, None, None, None) + .push_scope_at( + None, + "worker-scope", + ScopeType::Function, + None, + None, + None, + UNIX_EPOCH - Duration::from_nanos(1_500), + ) + .await?; + runtime + .pop_scope_at(&handle, None, None, UNIX_EPOCH) .await?; - runtime.pop_scope(&handle, None, None).await?; runtime.drop_scope_stack(&stack_id).await?; let mut next_value = next.call(value).await?; next_value.result = set_json_field(next_value.result, "phase", "tool_exec"); @@ -2613,6 +2642,12 @@ struct RuntimeRegistrationRequests { deregister: Vec, } +#[derive(Clone, Default)] +struct ScopeRequests { + push: Vec, + pop: Vec, +} + #[derive(Clone, Default)] struct MockHost { calls: Arc>>, @@ -2620,6 +2655,7 @@ struct MockHost { marks: Arc>>, failures: Arc>, runtime_registration_requests: Arc>, + scope_requests: Arc>, } impl MockHost { @@ -2653,6 +2689,13 @@ impl MockHost { .expect("runtime registration requests lock") .clone() } + + fn scope_requests(&self) -> ScopeRequests { + self.scope_requests + .lock() + .expect("scope requests lock") + .clone() + } } #[tonic::async_trait] @@ -2793,6 +2836,11 @@ impl RelayHostRuntime for MockHost { ) -> std::result::Result, Status> { let request = request.into_inner(); authorize_host(&request.activation_id, &request.auth_token)?; + self.scope_requests + .lock() + .expect("scope requests lock") + .push + .push(request.clone()); let scope = request.scope.expect("scope context"); self.record(format!( "push:{}:{}:{}", @@ -2816,6 +2864,11 @@ impl RelayHostRuntime for MockHost { ) -> std::result::Result, Status> { let request = request.into_inner(); authorize_host(&request.activation_id, &request.auth_token)?; + self.scope_requests + .lock() + .expect("scope requests lock") + .pop + .push(request.clone()); self.record(format!("pop:{}", request.scope_handle_id)); match self.failures().pop_scope { HostFailure::None => {} diff --git a/docs/build-plugins/native/runtime-events-and-scopes.mdx b/docs/build-plugins/native/runtime-events-and-scopes.mdx index da4fec04e..49a548582 100644 --- a/docs/build-plugins/native/runtime-events-and-scopes.mdx +++ b/docs/build-plugins/native/runtime-events-and-scopes.mdx @@ -79,6 +79,9 @@ pub(crate) fn emit_configured_runtime_events( `runtime.scope` returns an owning guard. Calling `close` supplies the end-event output; dropping an unclosed guard still follows the SDK cleanup path. `with_current` restores the prior stack after either `Ok` or `Err`, and dropping `isolated` releases the stack. +When a plugin replays work captured earlier, `runtime.scope_at` and +`ScopeGuard::close_at` accept the original `SystemTime` values while preserving the same +guard ownership and cleanup behavior. ## Emit Typed Telemetry and Read Diagnostics diff --git a/docs/build-plugins/workers/grpc-v1-protocol.mdx b/docs/build-plugins/workers/grpc-v1-protocol.mdx index 0230e7208..f63004d01 100644 --- a/docs/build-plugins/workers/grpc-v1-protocol.mdx +++ b/docs/build-plugins/workers/grpc-v1-protocol.mdx @@ -259,8 +259,8 @@ invocation ends. | `Log` | Emits an authenticated operational log record through Relay's configured sinks. | | `EmitMark` | Emits mark data and metadata under the supplied scope context. | | `GetRuntimeDiagnostics` | Returns the current bounded host-level diagnostic snapshot. | -| `PushScope` | Opens a typed scope with name, data, metadata, and input, returning the handle required for pop. | -| `PopScope` | Closes the owned scope handle with output and metadata. | +| `PushScope` | Opens a typed scope with name, data, metadata, input, and an optional historical start timestamp, returning the handle required for pop. | +| `PopScope` | Closes the owned scope handle with output, metadata, and an optional historical end timestamp. | | `CreateScopeStack` | Allocates an isolated stack and returns its opaque ID. | | `DropScopeStack` | Releases an isolated stack owned by the activation. | | `ToolNext` | Executes a tool continuation with JSON arguments and captured scope, returning a structural tool result or worker error. | @@ -277,8 +277,8 @@ The host-runtime request and response fields are complete in the following table | `Log` | `activation_id`, `auth_token`, required `level`, optional `target`, `message`, and JSON-object `fields` | `HostAck.ok` or `HostAck.error` | | `EmitMark` | `activation_id`, `auth_token`, captured `scope`, `name`, optional `data`, `metadata`, `data_schema`, `severity`, and `category` | `HostAck.ok` or `HostAck.error` | | `GetRuntimeDiagnostics` | `activation_id` and `auth_token` | Ordered `RuntimeDiagnostic` entries with `code`, `message`, and `count` | -| `PushScope` | `activation_id`, `auth_token`, captured `scope`, `name`, `scope_type`, and optional `data`, `metadata`, and `input` | `scope_handle_id` or `error` | -| `PopScope` | `activation_id`, `auth_token`, owned `scope_handle_id`, and optional `output` and `metadata` | `HostAck.ok` or `HostAck.error` | +| `PushScope` | `activation_id`, `auth_token`, captured `scope`, `name`, `scope_type`, and optional `data`, `metadata`, `input`, and `timestamp_unix_micros` | `scope_handle_id` or `error` | +| `PopScope` | `activation_id`, `auth_token`, owned `scope_handle_id`, and optional `output`, `metadata`, and `timestamp_unix_micros` | `HostAck.ok` or `HostAck.error` | | `CreateScopeStack` | `activation_id` and `auth_token` | `scope_stack_id` or `error` | | `DropScopeStack` | `activation_id`, `auth_token`, and owned `scope_stack_id` | `HostAck.ok` or `HostAck.error` | | `ToolNext` | `activation_id`, `auth_token`, `continuation_id`, JSON `value`, and captured `scope` | `ToolExecutionResultResponse.value` or `error` | @@ -292,6 +292,12 @@ The host-runtime request and response fields are complete in the following table agent, function, tool, LLM, retriever, embedder, reranker, guardrail, evaluator, custom, and unknown scopes. The unspecified wire value is invalid for a pushed scope. +Scope timestamps are signed Unix microseconds. Omitting the field preserves the normal +host-clock behavior; an explicitly present zero records the Unix epoch. The fields are an +additive `grpc-v1` extension. Workers that require historical timestamp preservation must +declare a Relay compatibility range beginning at 0.10, because an older protobuf host can +ignore an unknown field. + `LogRequest.level` is one of `error`, `warn`, `info`, `debug`, or `trace`. Relay prefixes the optional target with its worker-plugin namespace and applies its configured log filter. `fields` is a lossless JSON object: its keys become structured log attributes, so JSONL sinks diff --git a/docs/build-plugins/workers/runtime-events-and-scopes.mdx b/docs/build-plugins/workers/runtime-events-and-scopes.mdx index 37d330a74..37cb847dd 100644 --- a/docs/build-plugins/workers/runtime-events-and-scopes.mdx +++ b/docs/build-plugins/workers/runtime-events-and-scopes.mdx @@ -18,7 +18,7 @@ scope context while the `RelayHostRuntime` service performs the real work. | Emit a mark | Await `emit_mark` with name, data, and metadata. | Await `emit_mark` with the same JSON fields. | | Create or drop an isolated stack | Await `create_scope_stack` and `drop_scope_stack`. | Await the same operations and retain the returned stack ID only for its owned lifetime. | | Bind and restore a stack | Run a future through `with_scope_stack`. | Enter `bind_scope_stack`; use `clear_scope_stack` when code must deliberately run without a bound stack. | -| Push and pop a scope | Await `push_scope`, retain its handle, and await `pop_scope`. | Await the same pair; the handle, not a guessed scope ID, owns the pop. | +| Push and pop a scope | Await `push_scope`, retain its handle, and await `pop_scope`. Use `push_scope_at` and `pop_scope_at` when replaying captured work at its original `SystemTime`. | Await the same pair; pass a timezone-aware `timestamp` when replaying captured work. The handle, not a guessed scope ID, owns the pop. | | Inspect inherited context | The runtime proxy carries the task scope snapshot. | `current_scope_stack_id` and `current_parent_scope_id` expose the current binding. | | Emit typed telemetry | Use `emit_mark_with_options` or `emit_metric`. | Use the matching async runtime methods. | | Read runtime diagnostics | Await `runtime_diagnostics`. | Await `runtime_diagnostics`. | @@ -28,6 +28,7 @@ scope context while the `RelayHostRuntime` service performs the real work. Every host request carries the activation ID, token, and relevant scope context. A worker must not expose, log, or persist the activation token or codec capability IDs. Those values identify a live local capability and expire with the activation or invocation. +Ordinary scope calls omit the timestamp and retain Relay's current-time behavior. ## Control a Registration from a Background Task diff --git a/python/plugin/src/nemo_relay_plugin/_api.py b/python/plugin/src/nemo_relay_plugin/_api.py index aeb0f159c..2923bc52d 100644 --- a/python/plugin/src/nemo_relay_plugin/_api.py +++ b/python/plugin/src/nemo_relay_plugin/_api.py @@ -76,6 +76,7 @@ import tomllib from collections.abc import AsyncIterator, Awaitable, Callable, Iterable, Iterator, Mapping from dataclasses import asdict, dataclass, field +from datetime import datetime, timezone from enum import Enum from importlib import metadata from pathlib import Path @@ -1917,6 +1918,7 @@ async def push_scope( input: Json | None = None, scope_stack_id: str | None = None, parent_scope_id: str | None = None, + timestamp: datetime | None = None, ) -> str: """Start a scope on a Relay host-owned stack. @@ -1930,6 +1932,9 @@ async def push_scope( the current local binding is used. parent_scope_id: Optional parent scope. When omitted, the parent from the selected stack binding is used. + timestamp: Optional timezone-aware start time recorded on the + scope handle and start event. When omitted, Relay uses the + current time. Returns: An opaque scope handle to pass to :meth:`pop_scope`. @@ -1937,21 +1942,25 @@ async def push_scope( Raises: WorkerSdkError: The scope selection is invalid or the host rejects the request. - TypeError: A payload is not JSON-serializable. - ValueError: ``scope_type`` is not a supported :class:`ScopeType`. + TypeError: A payload is not JSON-serializable or ``timestamp`` is + not a :class:`datetime.datetime`. + ValueError: ``scope_type`` is unsupported or ``timestamp`` is + timezone-naive. """ - response = await self._host_stub.PushScope( - pb.PushScopeRequest( - activation_id=self._activation_id, - auth_token=self._auth_token, - scope=self._scope_context(scope_stack_id, parent_scope_id), - name=name, - scope_type=_proto_scope_type(scope_type), - data=_optional_json_envelope(data), - metadata=_optional_json_envelope(metadata), - input=_optional_json_envelope(input), - ) + request = pb.PushScopeRequest( + activation_id=self._activation_id, + auth_token=self._auth_token, + scope=self._scope_context(scope_stack_id, parent_scope_id), + name=name, + scope_type=_proto_scope_type(scope_type), + data=_optional_json_envelope(data), + metadata=_optional_json_envelope(metadata), + input=_optional_json_envelope(input), ) + timestamp_unix_micros = _datetime_to_unix_micros(timestamp) + if timestamp_unix_micros is not None: + request.timestamp_unix_micros = timestamp_unix_micros + response = await self._host_stub.PushScope(request) if response.HasField("error"): raise _worker_error_to_sdk(response.error) return response.scope_handle_id @@ -1962,6 +1971,7 @@ async def pop_scope( *, output: Json | None = None, metadata: Json | None = None, + timestamp: datetime | None = None, ) -> None: """End a host scope by its handle identifier. @@ -1970,20 +1980,26 @@ async def pop_scope( output: Optional JSON semantic output attached to the scope end event. metadata: Optional JSON metadata attached to the end event. + timestamp: Optional timezone-aware time recorded on the end event. + When omitted, Relay uses its default end time. Raises: WorkerSdkError: The host rejects the request. - TypeError: A payload is not JSON-serializable. + TypeError: A payload is not JSON-serializable or ``timestamp`` is + not a :class:`datetime.datetime`. + ValueError: ``timestamp`` is timezone-naive. """ - response = await self._host_stub.PopScope( - pb.PopScopeRequest( - activation_id=self._activation_id, - auth_token=self._auth_token, - scope_handle_id=scope_handle_id, - output=_optional_json_envelope(output), - metadata=_optional_json_envelope(metadata), - ) + request = pb.PopScopeRequest( + activation_id=self._activation_id, + auth_token=self._auth_token, + scope_handle_id=scope_handle_id, + output=_optional_json_envelope(output), + metadata=_optional_json_envelope(metadata), ) + timestamp_unix_micros = _datetime_to_unix_micros(timestamp) + if timestamp_unix_micros is not None: + request.timestamp_unix_micros = timestamp_unix_micros + response = await self._host_stub.PopScope(request) _ack_to_result(response) @contextlib.contextmanager @@ -2800,6 +2816,18 @@ def _optional_json_envelope(value: Json | None, schema: str = JSON_SCHEMA) -> An return _json_envelope(schema, value) +def _datetime_to_unix_micros(value: datetime | None) -> int | None: + if value is None: + return None + if not isinstance(value, datetime): + raise TypeError("timestamp must be a datetime.datetime object") + if value.tzinfo is None or value.utcoffset() is None: + raise ValueError("timestamp datetime must be timezone-aware") + epoch = datetime(1970, 1, 1, tzinfo=timezone.utc) + delta = value.astimezone(timezone.utc) - epoch + return (delta.days * 86_400 + delta.seconds) * 1_000_000 + delta.microseconds + + def _data_schema_json(value: DataSchema | Mapping[str, Json] | None) -> dict[str, str] | None: if value is None: return None diff --git a/python/tests/plugin/test_worker_sdk.py b/python/tests/plugin/test_worker_sdk.py index bc9269a12..794fe59c4 100644 --- a/python/tests/plugin/test_worker_sdk.py +++ b/python/tests/plugin/test_worker_sdk.py @@ -12,6 +12,7 @@ import socket import tempfile from collections.abc import AsyncIterator +from datetime import datetime, timezone from pathlib import Path from typing import Any, Literal, cast from unittest import mock @@ -2374,7 +2375,25 @@ async def test_runtime_host_calls_and_scope_context(host_stub: RecordingHostStub await runtime.emit_mark("mark", {"ok": True}) await runtime.emit_mark("override-parent", parent_scope_id="parent-2") scope_id = await runtime.push_scope("scope", scope_type=ScopeType.TOOL, input={"in": True}) + ordinary_push = _last_request(host_stub, pb.PushScopeRequest) + assert not ordinary_push.HasField("timestamp_unix_micros") await runtime.pop_scope(scope_id, output={"out": True}) + ordinary_pop = _last_request(host_stub, pb.PopScopeRequest) + assert not ordinary_pop.HasField("timestamp_unix_micros") + historical_scope_id = await runtime.push_scope( + "historical-scope", + timestamp=datetime(1969, 12, 31, 23, 59, 59, 999998, tzinfo=timezone.utc), + ) + historical_push = _last_request(host_stub, pb.PushScopeRequest) + assert historical_push.HasField("timestamp_unix_micros") + assert historical_push.timestamp_unix_micros == -2 + await runtime.pop_scope( + historical_scope_id, + timestamp=datetime(1970, 1, 1, tzinfo=timezone.utc), + ) + historical_pop = _last_request(host_stub, pb.PopScopeRequest) + assert historical_pop.HasField("timestamp_unix_micros") + assert historical_pop.timestamp_unix_micros == 0 tool_next = await ToolNext(runtime, "tool-next").call({"value": 1}) llm_next = await _llm_next(runtime, {"content": {"prompt": "hello"}}) stream_next = [chunk async for chunk in _llm_stream_next(runtime, {"content": {"prompt": "hello"}})] @@ -2694,6 +2713,30 @@ async def GetRuntimeDiagnostics(self, request: Any) -> Any: pass +@pytest.mark.parametrize( + ("timestamp", "error_type", "message"), + [ + (datetime(2026, 1, 1), ValueError, "timezone-aware"), + ("2026-01-01T00:00:00Z", TypeError, "datetime.datetime"), + ], +) +async def test_runtime_scope_timestamps_reject_invalid_values_before_host_call( + host_stub: RecordingHostStub, + timestamp: Any, + error_type: type[Exception], + message: str, +) -> None: + runtime = PluginRuntime(activation_id=ACTIVATION_ID, auth_token=AUTH_TOKEN, host_stub=host_stub) + request_count = len(host_stub.requests) + + with pytest.raises(error_type, match=message): + await runtime.push_scope("scope", timestamp=timestamp) + with pytest.raises(error_type, match=message): + await runtime.pop_scope("scope", timestamp=timestamp) + + assert len(host_stub.requests) == request_count + + async def test_lifecycle_acks(service: _WorkerService) -> None: cancel = await service.CancelInvocation( pb.CancelInvocationRequest(