Skip to main content

rlmesh_runtime/
spec.rs

1//! The runtime session spec, its limits, and the report a finished run returns.
2
3use std::sync::LazyLock;
4use std::time::Duration;
5
6use rlmesh_proto::core::v1::{AutoresetMode, EnvContract};
7use rlmesh_proto::spaces::v1::SpaceSpec;
8use serde::{Deserialize, Serialize};
9
10/// Empty fallback returned by the internal `*_validated` accessors only on the
11/// unreachable path where the space is absent despite validation (see their
12/// `debug_assert!`s). Lets those accessors stay panic-free and lint-clean.
13static EMPTY_SPACE_SPEC: LazyLock<SpaceSpec> = LazyLock::new(SpaceSpec::default);
14
15/// Everything one route needs to run: its identity, the negotiated env
16/// contract, and the per-op limits. [`validate`](Self::validate) gates a spec
17/// before the driver runs it.
18#[derive(Debug, Clone, PartialEq)]
19pub struct RuntimeSessionSpec {
20    /// Correlation label only; OSS does not key on it (the managed layer owns
21    /// session lifecycle). Kept as plumbing for telemetry/logs — its removal is
22    /// the deferred closed split, out of scope here.
23    pub session_id: String,
24    /// The connected env container, UUIDv7 (minted by the runtime on attach).
25    /// The single routing key: replaces the old `route_id` + positional lane.
26    /// (Repurposed from the former descriptive-name field; the human env name
27    /// now lives only in the language SDK's own contract type.)
28    pub env_id: String,
29    pub env_component_id: String,
30    pub model_component_id: String,
31    /// Workflow edition negotiated at the env handshake. The runtime refuses an
32    /// edition it was not built to drive (see [`RuntimeSessionSpec::validate`]).
33    pub workflow_edition: String,
34    pub env_contract: EnvContract,
35    pub num_envs: usize,
36    pub base_seed: Option<i64>,
37    /// Explicit per-episode reset seeds, consumed in episode-start order (a
38    /// vector reset claims one per lane). Overrides `base_seed` derivation when
39    /// non-empty; episodes beyond the list reset unseeded. Requires
40    /// driver-owned resets (autoreset `DISABLED`): under `NEXT_STEP` the env
41    /// seeds its own rolls, so the list would silently not apply.
42    pub episode_seeds: Vec<i64>,
43    pub max_episodes: Option<u64>,
44    /// Truncate any episode after this many steps (runtime-enforced; the lane
45    /// is reset and the episode reported `truncated`). Requires driver-owned
46    /// resets (autoreset `DISABLED`).
47    pub max_episode_steps: Option<i64>,
48    /// Truncate any episode after this wall-clock duration (seconds), same
49    /// semantics and autoreset requirement as `max_episode_steps`.
50    pub max_episode_seconds: Option<f64>,
51    pub close_env_on_end: bool,
52    pub limits: RuntimeLimits,
53}
54
55impl RuntimeSessionSpec {
56    pub fn validate(&self) -> Result<(), String> {
57        if self.session_id.trim().is_empty() {
58            return Err("runtime session_id must not be empty".to_string());
59        }
60        if self.env_id.trim().is_empty() {
61            return Err("runtime env_id must not be empty".to_string());
62        }
63        if self.env_component_id.trim().is_empty() {
64            return Err("runtime env_component_id must not be empty".to_string());
65        }
66        if self.model_component_id.trim().is_empty() {
67            return Err("runtime model_component_id must not be empty".to_string());
68        }
69        if self.num_envs == 0 {
70            return Err("runtime num_envs must be greater than zero".to_string());
71        }
72        if self.observation_space().is_none() {
73            return Err("runtime env_contract is missing observation_space".to_string());
74        }
75        if self.action_space().is_none() {
76            return Err("runtime env_contract is missing action_space".to_string());
77        }
78        if self.max_episodes == Some(0) {
79            return Err("runtime max_episodes must be greater than zero when set".to_string());
80        }
81        if self.max_episode_steps.is_some_and(|cap| cap <= 0) {
82            return Err("runtime max_episode_steps must be greater than zero when set".to_string());
83        }
84        if self.max_episode_seconds.is_some_and(|cap| cap <= 0.0) {
85            return Err(
86                "runtime max_episode_seconds must be greater than zero when set".to_string(),
87            );
88        }
89        let driver_owns_resets = matches!(
90            rlmesh_proto::core::v1::AutoresetMode::try_from(self.env_contract.autoreset_mode),
91            Ok(rlmesh_proto::core::v1::AutoresetMode::Disabled)
92                | Ok(rlmesh_proto::core::v1::AutoresetMode::Unspecified)
93        );
94        if !self.episode_seeds.is_empty()
95            && self.num_envs > 1
96            && !self.episode_seeds.len().is_multiple_of(self.num_envs)
97        {
98            return Err(format!(
99                "episode_seeds ({} seeds) must be a multiple of num_envs ({}) for a \
100                 vectorized env: driver-owned vector resets claim one seed per lane \
101                 per batch, so a partial batch would silently drop the tail seeds",
102                self.episode_seeds.len(),
103                self.num_envs
104            ));
105        }
106        if !driver_owns_resets {
107            if !self.episode_seeds.is_empty() {
108                return Err(
109                    "episode_seeds requires an env with autoreset disabled: under NEXT_STEP \
110                     autoreset the env seeds its own episode rolls, so explicit per-episode \
111                     seeds cannot apply"
112                        .to_string(),
113                );
114            }
115            if self.max_episode_steps.is_some() || self.max_episode_seconds.is_some() {
116                return Err(
117                    "max_episode_steps / max_episode_seconds require an env with autoreset \
118                     disabled: under NEXT_STEP autoreset the env owns lane resets, so the \
119                     runtime cannot truncate an episode"
120                        .to_string(),
121                );
122            }
123        }
124        // The runtime drives one of the editions it was built for; a session
125        // negotiated under an edition outside the support window is refused rather
126        // than run under semantics it never agreed to. Membership, not equality
127        // with CURRENT: a supported older edition (a graceful downgrade) is valid.
128        if !rlmesh_proto::is_supported_edition(&self.workflow_edition) {
129            return Err(format!(
130                "runtime cannot drive workflow edition {:?}; this build implements {:?}",
131                self.workflow_edition,
132                rlmesh_proto::SUPPORTED_WORKFLOW_EDITIONS
133            ));
134        }
135        // An autoreset mode this build does not understand (e.g. a newer peer's
136        // mode) must fail loudly at session setup, never silently fold to
137        // DISABLED and change lifecycle semantics.
138        if AutoresetMode::try_from(self.env_contract.autoreset_mode).is_err() {
139            return Err(format!(
140                "unknown autoreset mode {} on the wire; this build supports \
141                 UNSPECIFIED, NEXT_STEP, SAME_STEP, DISABLED only",
142                self.env_contract.autoreset_mode
143            ));
144        }
145        // SAME_STEP is reserved on the wire but not yet driven by the runtime:
146        // the driver currently aliases NEXT_STEP|SAME_STEP to a purely
147        // observational path, while the env server never rolls SAME_STEP episode
148        // ids -> done lanes would stall. Reject it here so it cannot reach the
149        // runtime under a false assumption of support.
150        if self.env_contract.autoreset_mode == AutoresetMode::SameStep as i32 {
151            return Err(
152                "SAME_STEP autoreset is reserved but not yet supported by the runtime; \
153                 construct the env with NEXT_STEP or DISABLED autoreset"
154                    .to_string(),
155            );
156        }
157        // Vectorized sessions require NEXT_STEP autoreset: the env resets each
158        // done lane itself, so the driver never needs per-lane reset. DISABLED
159        // (and the UNSPECIFIED default) would require resetting just the done
160        // lanes, which stock gymnasium vector envs cannot do. There is no
161        // partial-reset API, and a full reset clobbers the still-running lanes.
162        // Reject the combination up front instead of failing mid-run the first
163        // time lanes terminate at different steps. A future in-house vector
164        // engine with per-lane reset will lift this gate. (SAME_STEP is already
165        // rejected above, so the only mode that passes here is NEXT_STEP.)
166        if self.num_envs > 1 && self.env_contract.autoreset_mode != AutoresetMode::NextStep as i32 {
167            return Err(
168                "vectorized runtime sessions (num_envs > 1) require NEXT_STEP autoreset; \
169                 DISABLED autoreset needs per-lane reset, which is unavailable for stock \
170                 gymnasium vector envs. Use NEXT_STEP autoreset, or run with num_envs == 1."
171                    .to_string(),
172            );
173        }
174        Ok(())
175    }
176
177    pub fn env_context(&self) -> crate::hooks::RuntimeEnvContext {
178        crate::hooks::RuntimeEnvContext {
179            env_id: self.env_id.clone(),
180            env_component_id: self.env_component_id.clone(),
181            model_component_id: self.model_component_id.clone(),
182        }
183    }
184
185    /// Returns the observation space, or `None` if the spec has not been
186    /// populated/validated (`env_contract.observation_space` is unset).
187    ///
188    /// All `RuntimeSessionSpec` fields are public, so an unvalidated spec is
189    /// trivial to construct; this accessor never panics. The driver validates
190    /// the spec before running and uses the infallible internal accessor.
191    pub fn observation_space(&self) -> Option<&SpaceSpec> {
192        self.env_contract
193            .spec
194            .as_ref()
195            .and_then(|spec| spec.observation_space.as_ref())
196    }
197
198    /// Returns the action space, or `None` if the spec has not been
199    /// populated/validated (`env_contract.action_space` is unset).
200    ///
201    /// See [`RuntimeSessionSpec::observation_space`] for why this is fallible.
202    pub fn action_space(&self) -> Option<&SpaceSpec> {
203        self.env_contract
204            .spec
205            .as_ref()
206            .and_then(|spec| spec.action_space.as_ref())
207    }
208
209    /// Observation space for internal use after [`validate`](Self::validate)
210    /// has confirmed it is present.
211    pub(crate) fn observation_space_validated(&self) -> &SpaceSpec {
212        debug_assert!(
213            self.observation_space().is_some(),
214            "observation_space accessed before validate()"
215        );
216        // LazyLock<SpaceSpec> derefs to &SpaceSpec on the unreachable None path.
217        self.observation_space()
218            .unwrap_or_else(|| &EMPTY_SPACE_SPEC)
219    }
220
221    /// Action space for internal use after [`validate`](Self::validate) has
222    /// confirmed it is present.
223    pub(crate) fn action_space_validated(&self) -> &SpaceSpec {
224        debug_assert!(
225            self.action_space().is_some(),
226            "action_space accessed before validate()"
227        );
228        // LazyLock<SpaceSpec> derefs to &SpaceSpec on the unreachable None path.
229        self.action_space().unwrap_or_else(|| &EMPTY_SPACE_SPEC)
230    }
231}
232
233/// One completed episode's summary, recorded in completion order across all
234/// lanes. The pull counterpart of the `episode_completed` hook event, so a
235/// caller without hooks (e.g. the in-process `run_local` loop) still gets
236/// per-episode results on the report.
237#[derive(Debug, Clone, PartialEq)]
238pub struct EpisodeSummary {
239    /// 0-based completion index within the session.
240    pub episode_index: i64,
241    /// The vector lane the episode ran on (0 for a single env).
242    pub env_index: i32,
243    /// The explicit seed this episode was reset with (`episode_seeds` /
244    /// `base_seed` derivation), `None` for an unseeded or autoreset-rolled one.
245    pub seed: Option<i64>,
246    pub step_count: i64,
247    pub cumulative_reward: f64,
248    pub terminated: bool,
249    pub truncated: bool,
250    pub duration_ms: i64,
251    /// Env-reported task outcome from the final step's info (Gymnasium's
252    /// `is_success` / `success` key); `None` when the env emits no such signal.
253    pub success: Option<bool>,
254}
255
256/// What a finished or aborted session returns: totals plus the durable
257/// telemetry aggregate.
258#[derive(Debug, Clone, PartialEq)]
259pub struct RuntimeReport {
260    pub session_id: String,
261    pub env_id: String,
262    pub total_steps: i64,
263    pub total_episodes: i64,
264    /// Every completed episode, in completion order.
265    pub episodes: Vec<EpisodeSummary>,
266    /// Session-total telemetry aggregate (per-op latency/percentiles/bytes) —
267    /// the durable pull counterpart to the live `RuntimeHooks::on_telemetry` push.
268    pub telemetry: crate::telemetry::Snapshot,
269}
270
271/// Per-op timeouts and the telemetry window for one session. Serialized with
272/// explicit millisecond field names (`*Ms`); legacy unsuffixed fields are rejected.
273#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
274#[serde(rename_all = "camelCase", deny_unknown_fields)]
275pub struct RuntimeLimits {
276    #[serde(
277        default = "default_connect_timeout",
278        rename = "envConnectTimeoutMs",
279        serialize_with = "duration_millis::serialize",
280        deserialize_with = "duration_millis::deserialize"
281    )]
282    pub env_connect_timeout: Duration,
283    #[serde(
284        default = "default_model_connect_timeout",
285        rename = "modelConnectTimeoutMs",
286        serialize_with = "duration_millis::serialize",
287        deserialize_with = "duration_millis::deserialize"
288    )]
289    pub model_connect_timeout: Duration,
290    #[serde(
291        default = "default_configure_route_timeout",
292        rename = "configureRouteTimeoutMs",
293        serialize_with = "duration_millis::serialize",
294        deserialize_with = "duration_millis::deserialize"
295    )]
296    pub configure_route_timeout: Duration,
297    #[serde(
298        default = "default_env_reset_timeout",
299        rename = "envResetTimeoutMs",
300        serialize_with = "duration_millis::serialize",
301        deserialize_with = "duration_millis::deserialize"
302    )]
303    pub env_reset_timeout: Duration,
304    #[serde(
305        default = "default_model_predict_timeout",
306        rename = "modelPredictTimeoutMs",
307        serialize_with = "duration_millis::serialize",
308        deserialize_with = "duration_millis::deserialize"
309    )]
310    pub model_predict_timeout: Duration,
311    #[serde(
312        default = "default_env_step_timeout",
313        rename = "envStepTimeoutMs",
314        serialize_with = "duration_millis::serialize",
315        deserialize_with = "duration_millis::deserialize"
316    )]
317    pub env_step_timeout: Duration,
318    #[serde(
319        default = "default_service_close_timeout",
320        rename = "serviceCloseTimeoutMs",
321        serialize_with = "duration_millis::serialize",
322        deserialize_with = "duration_millis::deserialize"
323    )]
324    pub service_close_timeout: Duration,
325    /// How often the background ticker pushes Window + Session telemetry
326    /// snapshots to `on_telemetry`, on a wall clock. `0` disables live streaming
327    /// (the final session snapshot is still delivered at session end); any
328    /// non-zero value below 1ms is floored to 1ms.
329    #[serde(
330        default = "default_telemetry_window",
331        rename = "telemetryWindowMs",
332        serialize_with = "duration_millis::serialize",
333        deserialize_with = "duration_millis::deserialize"
334    )]
335    pub telemetry_window: Duration,
336}
337
338impl Default for RuntimeLimits {
339    fn default() -> Self {
340        Self {
341            env_connect_timeout: default_connect_timeout(),
342            model_connect_timeout: default_model_connect_timeout(),
343            configure_route_timeout: default_configure_route_timeout(),
344            env_reset_timeout: default_env_reset_timeout(),
345            model_predict_timeout: default_model_predict_timeout(),
346            env_step_timeout: default_env_step_timeout(),
347            service_close_timeout: default_service_close_timeout(),
348            telemetry_window: default_telemetry_window(),
349        }
350    }
351}
352
353impl RuntimeLimits {
354    pub fn env_step_timeout_ms(&self) -> i64 {
355        duration_ms_i64(self.env_step_timeout)
356    }
357
358    pub fn env_reset_timeout_ms(&self) -> i64 {
359        duration_ms_i64(self.env_reset_timeout)
360    }
361}
362
363fn default_connect_timeout() -> Duration {
364    Duration::from_secs(60)
365}
366
367fn default_model_connect_timeout() -> Duration {
368    Duration::from_secs(600)
369}
370
371fn default_configure_route_timeout() -> Duration {
372    Duration::from_secs(600)
373}
374
375fn default_env_reset_timeout() -> Duration {
376    Duration::from_secs(300)
377}
378
379fn default_model_predict_timeout() -> Duration {
380    Duration::from_secs(300)
381}
382
383fn default_env_step_timeout() -> Duration {
384    Duration::from_secs(300)
385}
386
387fn default_service_close_timeout() -> Duration {
388    Duration::from_secs(5)
389}
390
391fn default_telemetry_window() -> Duration {
392    Duration::from_secs(1)
393}
394
395fn duration_ms_i64(duration: Duration) -> i64 {
396    duration.as_millis().try_into().unwrap_or(i64::MAX)
397}
398
399mod duration_millis {
400    use std::time::Duration;
401
402    use serde::{Deserialize, Deserializer, Serializer};
403
404    pub(super) fn serialize<S>(duration: &Duration, serializer: S) -> Result<S::Ok, S::Error>
405    where
406        S: Serializer,
407    {
408        let millis = duration.as_millis().try_into().unwrap_or(u64::MAX);
409        serializer.serialize_u64(millis)
410    }
411
412    pub(super) fn deserialize<'de, D>(deserializer: D) -> Result<Duration, D::Error>
413    where
414        D: Deserializer<'de>,
415    {
416        Ok(Duration::from_millis(u64::deserialize(deserializer)?))
417    }
418}
419
420#[cfg(test)]
421mod tests {
422    use std::time::Duration;
423
424    use rlmesh_proto::core::v1::{AutoresetMode, EnvContract, EnvSpec};
425    use rlmesh_proto::spaces::v1::SpaceSpec;
426    use serde_json::json;
427
428    use super::{RuntimeLimits, RuntimeSessionSpec};
429
430    fn valid_spec() -> RuntimeSessionSpec {
431        RuntimeSessionSpec {
432            session_id: "session".to_string(),
433            env_id: "env-id".to_string(),
434            env_component_id: "env".to_string(),
435            model_component_id: "model".to_string(),
436            workflow_edition: rlmesh_proto::CURRENT_WORKFLOW_EDITION.to_string(),
437            env_contract: EnvContract {
438                spec: Some(EnvSpec {
439                    observation_space: Some(SpaceSpec::default()),
440                    action_space: Some(SpaceSpec::default()),
441                    ..Default::default()
442                }),
443                num_envs: 1,
444                ..Default::default()
445            },
446            num_envs: 1,
447            episode_seeds: Vec::new(),
448            base_seed: None,
449            max_episodes: Some(1),
450            max_episode_steps: None,
451            max_episode_seconds: None,
452            close_env_on_end: true,
453            limits: RuntimeLimits::default(),
454        }
455    }
456
457    #[test]
458    fn validate_rejects_vector_episode_seeds_that_do_not_cover_whole_batches() {
459        let mut spec = valid_spec();
460        spec.num_envs = 2;
461        spec.env_contract.autoreset_mode = AutoresetMode::NextStep as i32;
462        spec.episode_seeds = vec![7, 8, 9];
463        let error = spec.validate().unwrap_err();
464        assert!(
465            error.contains("multiple of num_envs"),
466            "expected the whole-batch seed rule, got: {error}"
467        );
468
469        spec.episode_seeds = vec![7, 8, 9, 10];
470        let error = spec.validate().unwrap_err();
471        assert!(
472            error.contains("autoreset disabled"),
473            "a covering list still needs driver-owned resets, got: {error}"
474        );
475    }
476
477    #[test]
478    fn validate_rejects_an_edition_the_runtime_cannot_drive() {
479        let mut spec = valid_spec();
480        spec.workflow_edition = "2099.01".to_string();
481        let error = spec.validate().unwrap_err();
482        assert!(
483            error.contains("2099.01") && error.contains("cannot drive"),
484            "expected an edition-refusal error, got: {error}"
485        );
486
487        // The edition the build implements is accepted.
488        spec.workflow_edition = rlmesh_proto::CURRENT_WORKFLOW_EDITION.to_string();
489        assert!(spec.validate().is_ok());
490    }
491
492    #[test]
493    fn space_accessors_return_none_on_unvalidated_spec() {
494        let mut spec = valid_spec();
495        // An unvalidated spec is trivially constructible since all fields are
496        // public; the accessors must not panic.
497        spec.env_contract = EnvContract::default();
498
499        assert!(spec.observation_space().is_none());
500        assert!(spec.action_space().is_none());
501    }
502
503    #[test]
504    fn space_accessors_return_some_on_populated_spec() {
505        let spec = valid_spec();
506        assert!(spec.observation_space().is_some());
507        assert!(spec.action_space().is_some());
508    }
509
510    #[test]
511    fn validate_accepts_vectorized_next_step_runtime_sessions() {
512        // num_envs > 1 is supported with NEXT_STEP autoreset: the env resets each
513        // done lane itself, so the driver never needs per-lane reset.
514        let mut spec = valid_spec();
515        spec.num_envs = 4;
516        spec.env_contract.num_envs = 4;
517        spec.env_contract.autoreset_mode = AutoresetMode::NextStep as i32;
518
519        assert!(spec.validate().is_ok());
520    }
521
522    #[test]
523    fn validate_rejects_disabled_vectorized_sessions() {
524        // DISABLED (and the UNSPECIFIED default) with num_envs > 1 needs per-lane
525        // reset, which stock gymnasium vector envs cannot do. Reject up front
526        // rather than failing mid-run on the first staggered termination.
527        for mode in [AutoresetMode::Disabled, AutoresetMode::Unspecified] {
528            let mut spec = valid_spec();
529            spec.num_envs = 4;
530            spec.env_contract.num_envs = 4;
531            spec.env_contract.autoreset_mode = mode as i32;
532
533            let error = spec.validate().unwrap_err();
534            assert!(
535                error.contains("NEXT_STEP"),
536                "expected a NEXT_STEP-guidance rejection for {mode:?}, got: {error}"
537            );
538        }
539    }
540
541    #[test]
542    fn validate_rejects_same_step_autoreset() {
543        // SAME_STEP is reserved but unsupported; validation must reject it so it
544        // cannot reach the runtime and stall lanes.
545        let mut spec = valid_spec();
546        spec.env_contract.autoreset_mode = AutoresetMode::SameStep as i32;
547
548        let error = spec.validate().unwrap_err();
549        assert!(
550            error.contains("SAME_STEP"),
551            "expected SAME_STEP rejection, got: {error}"
552        );
553    }
554
555    #[test]
556    fn runtime_limits_json_uses_explicit_millisecond_fields() {
557        let value = serde_json::to_value(RuntimeLimits::default()).unwrap();
558
559        assert_eq!(value["envConnectTimeoutMs"], json!(60_000));
560        assert_eq!(value["modelConnectTimeoutMs"], json!(600_000));
561        assert_eq!(value["configureRouteTimeoutMs"], json!(600_000));
562        assert_eq!(value["envResetTimeoutMs"], json!(300_000));
563        assert_eq!(value["modelPredictTimeoutMs"], json!(300_000));
564        assert_eq!(value["envStepTimeoutMs"], json!(300_000));
565        assert_eq!(value["serviceCloseTimeoutMs"], json!(5_000));
566        assert_eq!(value["telemetryWindowMs"], json!(1_000));
567        assert!(value.get("envConnectTimeout").is_none());
568
569        let parsed: RuntimeLimits = serde_json::from_value(json!({
570            "envConnectTimeoutMs": 1,
571            "modelConnectTimeoutMs": 2,
572            "configureRouteTimeoutMs": 3,
573            "envResetTimeoutMs": 4,
574            "modelPredictTimeoutMs": 5,
575            "envStepTimeoutMs": 6,
576            "serviceCloseTimeoutMs": 7,
577            "telemetryWindowMs": 8
578        }))
579        .unwrap();
580
581        assert_eq!(parsed.env_connect_timeout, Duration::from_millis(1));
582        assert_eq!(parsed.model_connect_timeout, Duration::from_millis(2));
583        assert_eq!(parsed.configure_route_timeout, Duration::from_millis(3));
584        assert_eq!(parsed.env_reset_timeout, Duration::from_millis(4));
585        assert_eq!(parsed.model_predict_timeout, Duration::from_millis(5));
586        assert_eq!(parsed.env_step_timeout, Duration::from_millis(6));
587        assert_eq!(parsed.service_close_timeout, Duration::from_millis(7));
588        assert_eq!(parsed.telemetry_window, Duration::from_millis(8));
589    }
590
591    #[test]
592    fn runtime_limits_reject_legacy_unsuffixed_fields() {
593        let error = serde_json::from_value::<RuntimeLimits>(json!({
594            "envConnectTimeout": 1
595        }))
596        .unwrap_err();
597
598        assert!(error.to_string().contains("envConnectTimeout"));
599    }
600}