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::collections::{HashMap, HashSet};
4use std::sync::LazyLock;
5use std::time::Duration;
6
7use rlmesh_proto::core::v1::{AutoresetMode, EnvContract};
8use rlmesh_proto::spaces::v1::meta_value::Kind as MetaKind;
9use rlmesh_proto::spaces::v1::{MetaList, MetaMap, MetaValue, SpaceSpec};
10use rlmesh_proto::{Edition, EditionDefaults};
11use rlmesh_spaces::{Advisory, DType};
12use serde::{Deserialize, Serialize};
13
14/// Empty fallback returned by the internal `*_validated` accessors only on the
15/// unreachable path where the space is absent despite validation (see their
16/// `debug_assert!`s). Lets those accessors stay panic-free and lint-clean.
17static EMPTY_SPACE_SPEC: LazyLock<SpaceSpec> = LazyLock::new(SpaceSpec::default);
18
19/// Key under which an env declares, in its contract metadata, which reserved
20/// reset-option keys it wants delivered (a list of strings). The runtime sends a
21/// reserved `ResetRequest.options` key only to an env that named it here, so an
22/// env that forwards `options` blindly into a third-party `reset` never receives
23/// one it cannot interpret. Mirrored in Python as `rlmesh.ENV_RESET_OPTIONS_KEY`.
24pub const ENV_RESET_OPTIONS_KEY: &str = "rlmesh.env.v1.reset_options";
25
26/// The reserved reset option carrying the ordinal of the episode a reset
27/// starts. Declared by an env under [`ENV_RESET_OPTIONS_KEY`]; read in Python
28/// as `rlmesh.trial_index(options)`.
29///
30/// The env-side spelling of the key. The driver never reads it: a session takes
31/// its trial-index key from [`EditionDefaults::trial_index_option_key`], so the
32/// edition governs what goes on the wire.
33pub const TRIAL_INDEX_OPTION: &str = "trial_index";
34
35/// Whether `contract` named `key` in its metadata under
36/// [`ENV_RESET_OPTIONS_KEY`], as either a list of strings or a bare string.
37///
38/// The declaration gate for every reserved `ResetRequest.options` key: an env
39/// that forwards `options` blindly into a third-party `reset` must never
40/// receive a reserved key it cannot interpret.
41pub fn declares_reset_option(contract: &EnvContract, key: &str) -> bool {
42    let Some(declared) = contract
43        .spec
44        .as_ref()
45        .and_then(|spec| spec.metadata.as_ref())
46        .and_then(|metadata| metadata.entries.get(ENV_RESET_OPTIONS_KEY))
47        .and_then(|declared| declared.kind.as_ref())
48    else {
49        return false;
50    };
51    match declared {
52        MetaKind::List(list) => list.items.iter().any(
53            |item| matches!(item.kind.as_ref(), Some(MetaKind::Text(declared)) if declared == key),
54        ),
55        MetaKind::Text(declared) => declared == key,
56        _ => false,
57    }
58}
59
60/// The `ResetRequest.options` map delivering `trials` under `option_key` — the
61/// trial-index key the session's edition governs
62/// ([`EditionDefaults::trial_index_option_key`], mirrored for env-side callers by
63/// [`TRIAL_INDEX_OPTION`]) — or `None` when there are no trials or the env never
64/// declared the key.
65///
66/// A single lane sends the bare integer; a multi-lane reset sends the list, in
67/// the same lane order as `seeds` and `episode_ids`.
68pub fn reset_options_for(
69    contract: &EnvContract,
70    option_key: &str,
71    trials: &[u64],
72) -> Option<MetaMap> {
73    if trials.is_empty() || !declares_reset_option(contract, option_key) {
74        return None;
75    }
76    let value = if trials.len() == 1 {
77        MetaValue {
78            kind: Some(MetaKind::Integer(trials[0] as i64)),
79        }
80    } else {
81        MetaValue {
82            kind: Some(MetaKind::List(MetaList {
83                items: trials
84                    .iter()
85                    .map(|trial| MetaValue {
86                        kind: Some(MetaKind::Integer(*trial as i64)),
87                    })
88                    .collect(),
89            })),
90        }
91    };
92    Some(MetaMap {
93        entries: [(option_key.to_string(), value)].into_iter().collect(),
94    })
95}
96
97/// What one served peer can decode, learned at its handshake. The driver checks
98/// every payload it relays against the target leg's ceiling through the
99/// session's [`RelayPolicy`](crate::RelayPolicy), not against the session
100/// edition alone.
101#[derive(Debug, Clone, PartialEq, Eq)]
102pub struct PeerCeiling {
103    /// The highest edition this peer and this build both implement, which may
104    /// exceed the session edition
105    /// ([`workflow_edition`](RuntimeSessionSpec::workflow_edition)).
106    pub edition: Edition,
107    pub dtypes: HashSet<DType>,
108    /// What the peer advertised; query with [`rlmesh_proto::has_capability`].
109    pub capabilities: HashMap<String, String>,
110    /// The largest message the peer accepts, in bytes.
111    pub max_message_size: usize,
112}
113
114impl PeerCeiling {
115    /// The highest edition this build retains among a peer's advertised
116    /// `supported` names, or `None` when they share none.
117    pub fn highest_shared_edition(supported: &[String]) -> Option<Edition> {
118        supported
119            .iter()
120            .filter(|name| rlmesh_proto::parse_retained_edition(name).is_ok())
121            .max_by(|a, b| {
122                rlmesh_proto::edition_sort_key(a).cmp(&rlmesh_proto::edition_sort_key(b))
123            })
124            .and_then(|name| rlmesh_proto::parse_retained_edition(name).ok())
125    }
126
127    /// A wire-v1 peer's ceiling. Wire-v1 advertises no dtype set, so a peer
128    /// decodes every dtype the wire defines.
129    pub fn wire_v1(
130        edition: Edition,
131        capabilities: HashMap<String, String>,
132        max_message_size: usize,
133    ) -> Self {
134        Self {
135            edition,
136            dtypes: DType::ALL.into_iter().collect(),
137            capabilities,
138            max_message_size,
139        }
140    }
141}
142
143/// Everything one route needs to run: its identity, the negotiated env
144/// contract, and the per-op limits. [`validate`](Self::validate) gates a spec
145/// before the driver runs it.
146#[derive(Debug, Clone, PartialEq)]
147pub struct RuntimeSessionSpec {
148    /// Correlation label only; OSS does not key on it (the managed layer owns
149    /// session lifecycle). Kept as plumbing for telemetry/logs — its removal is
150    /// the deferred closed split, out of scope here.
151    pub session_id: String,
152    /// The connected env container, UUIDv7 (minted by the runtime on attach).
153    /// The single routing key: replaces the old `route_id` + positional lane.
154    /// (Repurposed from the former descriptive-name field; the human env name
155    /// now lives only in the language SDK's own contract type.)
156    pub env_id: String,
157    pub env_component_id: String,
158    pub model_component_id: String,
159    /// Workflow edition negotiated at the env handshake, already resolved to the
160    /// typed arm whose semantics this session runs under. Wire names are parsed
161    /// once, where they arrive, by
162    /// [`rlmesh_proto::parse_retained_edition`] — that is where the runtime
163    /// refuses an edition it was not built to drive, so a name outside the
164    /// retained list never reaches this field.
165    pub workflow_edition: Edition,
166    pub env_contract: EnvContract,
167    pub num_envs: usize,
168    pub base_seed: Option<i64>,
169    /// Explicit per-episode reset seeds, consumed in episode-start order (a
170    /// vector reset claims one per lane). Overrides `base_seed` derivation when
171    /// non-empty; episodes beyond the list reset unseeded. Requires
172    /// driver-owned resets (autoreset `DISABLED`): under `NEXT_STEP` the env
173    /// seeds its own rolls, so the list would silently not apply.
174    pub episode_seeds: Vec<i64>,
175    pub max_episodes: Option<u64>,
176    /// First trial ordinal this route's episodes walk; `None` is the default
177    /// base, 0 (so is an explicit `Some(0)`). Under driver-owned resets
178    /// (autoreset `DISABLED`) the driver always mints one ordinal per episode
179    /// start (`base`, `base + 1`, ...), reports it on the episode events and
180    /// summaries, and delivers it as `ResetRequest.options["trial_index"]` to an
181    /// env that declared the key (see [`ENV_RESET_OPTIONS_KEY`]) -- an env that
182    /// did not never sees it. A sharded run gives each shard its own base so the
183    /// shards together walk a benchmark's trials once each, instead of every
184    /// shard re-deriving an index from a hashed seed. Under `NEXT_STEP` autoreset
185    /// the env restarts its own lanes, so no ordinal is minted at all: the
186    /// default base is inert there, and a non-zero base is rejected by
187    /// [`validate`](Self::validate) since it could never be walked.
188    pub trial_index_base: Option<u64>,
189    /// Truncate any episode after this many steps (runtime-enforced; the lane
190    /// is reset and the episode reported `truncated`). Requires driver-owned
191    /// resets (autoreset `DISABLED`).
192    pub max_episode_steps: Option<i64>,
193    /// Truncate any episode after this wall-clock duration (seconds), same
194    /// semantics and autoreset requirement as `max_episode_steps`.
195    pub max_episode_seconds: Option<f64>,
196    pub close_env_on_end: bool,
197    pub limits: RuntimeLimits,
198    /// The env advertised the `subset_step` handshake capability: every lane
199    /// can be reset and stepped on its own. The driver then runs one episode
200    /// loop per lane (a lane is its own group) instead of stepping the vector
201    /// in lockstep, and episode seeds/indices come from a route-global slot
202    /// counter so the scored set is fixed by the budget alone.
203    pub subset_step: bool,
204    /// The served env's ceiling; `None` when the env runs in-process.
205    pub env_ceiling: Option<PeerCeiling>,
206    /// The served model's ceiling; `None` when the model runs in-process.
207    pub model_ceiling: Option<PeerCeiling>,
208}
209
210impl RuntimeSessionSpec {
211    /// The first trial ordinal this route walks: `trial_index_base`, or 0 when
212    /// the session left it unset.
213    pub fn trial_index_base(&self) -> u64 {
214        self.trial_index_base.unwrap_or(0)
215    }
216
217    pub fn validate(&self) -> Result<(), String> {
218        if self.session_id.trim().is_empty() {
219            return Err("runtime session_id must not be empty".to_string());
220        }
221        if self.env_id.trim().is_empty() {
222            return Err("runtime env_id must not be empty".to_string());
223        }
224        if self.env_component_id.trim().is_empty() {
225            return Err("runtime env_component_id must not be empty".to_string());
226        }
227        if self.model_component_id.trim().is_empty() {
228            return Err("runtime model_component_id must not be empty".to_string());
229        }
230        if self.num_envs == 0 {
231            return Err("runtime num_envs must be greater than zero".to_string());
232        }
233        if self.observation_space().is_none() {
234            return Err("runtime env_contract is missing observation_space".to_string());
235        }
236        if self.action_space().is_none() {
237            return Err("runtime env_contract is missing action_space".to_string());
238        }
239        if self.max_episodes == Some(0) {
240            return Err("runtime max_episodes must be greater than zero when set".to_string());
241        }
242        if self.max_episode_steps.is_some_and(|cap| cap <= 0) {
243            return Err("runtime max_episode_steps must be greater than zero when set".to_string());
244        }
245        if self.max_episode_seconds.is_some_and(|cap| cap <= 0.0) {
246            return Err(
247                "runtime max_episode_seconds must be greater than zero when set".to_string(),
248            );
249        }
250        let driver_owns_resets = AutoresetMode::try_from(self.env_contract.autoreset_mode)
251            .is_ok_and(|mode| {
252                self.edition_defaults()
253                    .driver_owned_reset_modes
254                    .contains(&mode)
255            });
256        if !driver_owns_resets {
257            if !self.episode_seeds.is_empty() {
258                return Err(
259                    "episode_seeds requires an env with autoreset disabled: under NEXT_STEP \
260                     autoreset the env seeds its own episode rolls, so explicit per-episode \
261                     seeds cannot apply"
262                        .to_string(),
263                );
264            }
265            if self.trial_index_base() != 0 {
266                return Err(
267                    "trial_index_base requires an env with autoreset disabled: under \
268                     NEXT_STEP autoreset the env restarts its own lanes, so the runtime \
269                     has no reset on which to deliver a trial ordinal"
270                        .to_string(),
271                );
272            }
273            if self.max_episode_steps.is_some() || self.max_episode_seconds.is_some() {
274                return Err(
275                    "max_episode_steps / max_episode_seconds require an env with autoreset \
276                     disabled: under NEXT_STEP autoreset the env owns lane resets, so the \
277                     runtime cannot truncate an episode"
278                        .to_string(),
279                );
280            }
281        }
282        // The runtime drives one of the editions it was built for. Wire names are
283        // already refused at the parse boundary, so this re-checks the typed value
284        // for the one route that skips it: every field here is public, so a caller
285        // can name an arm this build implements but no longer retains. Same
286        // authority, same wording — membership, not equality with CURRENT, since a
287        // retained older edition (a graceful downgrade) is a valid session floor.
288        rlmesh_proto::parse_retained_edition(self.workflow_edition.base())?;
289        // An autoreset mode this build does not understand (e.g. a newer peer's
290        // mode) must fail loudly at session setup, never silently fold to
291        // DISABLED and change lifecycle semantics.
292        if AutoresetMode::try_from(self.env_contract.autoreset_mode).is_err() {
293            return Err(format!(
294                "unknown autoreset mode {} on the wire; this build supports \
295                 UNSPECIFIED, NEXT_STEP, SAME_STEP, DISABLED only",
296                self.env_contract.autoreset_mode
297            ));
298        }
299        // SAME_STEP is reserved on the wire but not yet driven by the runtime:
300        // the driver currently aliases NEXT_STEP|SAME_STEP to a purely
301        // observational path, while the env server never rolls SAME_STEP episode
302        // ids -> done lanes would stall. Reject it here so it cannot reach the
303        // runtime under a false assumption of support.
304        if self.env_contract.autoreset_mode == AutoresetMode::SameStep as i32 {
305            return Err(
306                "SAME_STEP autoreset is reserved but not yet supported by the runtime; \
307                 construct the env with NEXT_STEP or DISABLED autoreset"
308                    .to_string(),
309            );
310        }
311        // A lockstep vectorized session (one group of N lanes) requires
312        // NEXT_STEP autoreset: the env resets each done lane itself. Under
313        // DISABLED the driver would have to reset just the done lanes, which a
314        // stock gymnasium vector env cannot do (a full reset clobbers the
315        // still-running lanes). A lane endpoint (`subset_step`) is driven one
316        // group per lane instead, where DISABLED is the norm. Reject the
317        // combination up front instead of failing mid-run the first time lanes
318        // terminate at different steps. (SAME_STEP is already rejected above.)
319        if self.num_envs > 1
320            && !self.subset_step
321            && self.env_contract.autoreset_mode != AutoresetMode::NextStep as i32
322        {
323            return Err(
324                "vectorized runtime sessions (num_envs > 1) require NEXT_STEP autoreset unless \
325                 the env steps lanes individually (the `subset_step` capability); DISABLED \
326                 autoreset needs per-lane reset, which a stock gymnasium vector env cannot do. \
327                 Use NEXT_STEP autoreset, serve lanes, or run with num_envs == 1."
328                    .to_string(),
329            );
330        }
331        Ok(())
332    }
333
334    /// The edition-governed defaults this session runs under. Every value the
335    /// driver would otherwise hardcode comes from here, so a future edition
336    /// changes a table row instead of a code path.
337    pub fn edition_defaults(&self) -> &'static EditionDefaults {
338        rlmesh_proto::defaults(self.workflow_edition)
339    }
340
341    pub fn env_context(&self) -> crate::hooks::RuntimeEnvContext {
342        crate::hooks::RuntimeEnvContext {
343            env_id: self.env_id.clone(),
344            env_component_id: self.env_component_id.clone(),
345            model_component_id: self.model_component_id.clone(),
346            lane: None,
347        }
348    }
349
350    /// Returns the observation space, or `None` if the spec has not been
351    /// populated/validated (`env_contract.observation_space` is unset).
352    ///
353    /// All `RuntimeSessionSpec` fields are public, so an unvalidated spec is
354    /// trivial to construct; this accessor never panics. The driver validates
355    /// the spec before running and uses the infallible internal accessor.
356    pub fn observation_space(&self) -> Option<&SpaceSpec> {
357        self.env_contract
358            .spec
359            .as_ref()
360            .and_then(|spec| spec.observation_space.as_ref())
361    }
362
363    /// Returns the action space, or `None` if the spec has not been
364    /// populated/validated (`env_contract.action_space` is unset).
365    ///
366    /// See [`RuntimeSessionSpec::observation_space`] for why this is fallible.
367    pub fn action_space(&self) -> Option<&SpaceSpec> {
368        self.env_contract
369            .spec
370            .as_ref()
371            .and_then(|spec| spec.action_space.as_ref())
372    }
373
374    /// Observation space for internal use after [`validate`](Self::validate)
375    /// has confirmed it is present.
376    pub(crate) fn observation_space_validated(&self) -> &SpaceSpec {
377        debug_assert!(
378            self.observation_space().is_some(),
379            "observation_space accessed before validate()"
380        );
381        // LazyLock<SpaceSpec> derefs to &SpaceSpec on the unreachable None path.
382        self.observation_space()
383            .unwrap_or_else(|| &EMPTY_SPACE_SPEC)
384    }
385
386    /// Action space for internal use after [`validate`](Self::validate) has
387    /// confirmed it is present.
388    pub(crate) fn action_space_validated(&self) -> &SpaceSpec {
389        debug_assert!(
390            self.action_space().is_some(),
391            "action_space accessed before validate()"
392        );
393        // LazyLock<SpaceSpec> derefs to &SpaceSpec on the unreachable None path.
394        self.action_space().unwrap_or_else(|| &EMPTY_SPACE_SPEC)
395    }
396}
397
398/// One completed episode's summary, recorded in completion order across all
399/// lanes. The pull counterpart of the `episode_completed` hook event, so a
400/// caller without hooks (e.g. the in-process `run_local` loop) still gets
401/// per-episode results on the report.
402#[derive(Debug, Clone, PartialEq)]
403pub struct EpisodeSummary {
404    /// 1-based slot ordinal within the session, in completion order.
405    pub episode_index: i64,
406    /// The vector lane the episode ran on (0 for a single env).
407    pub env_index: i32,
408    /// The explicit seed this episode was reset with (`episode_seeds` /
409    /// `base_seed` derivation), `None` for an unseeded or autoreset-rolled one.
410    pub seed: Option<i64>,
411    /// The trial ordinal this episode walked (`trial_index_base` + its
412    /// episode-start position). Minted for every driver-owned reset whether or
413    /// not the env declared the reset option, so a coverage audit can read the
414    /// sweep off the report either way; `None` only under `NEXT_STEP` autoreset,
415    /// where the env restarts its own lanes and no ordinal is minted.
416    pub trial_index: Option<u64>,
417    pub step_count: i64,
418    pub cumulative_reward: f64,
419    pub terminated: bool,
420    pub truncated: bool,
421    pub duration_ms: i64,
422    /// Env-reported task outcome from the final step's info (Gymnasium's
423    /// `is_success` / `success`, or `task_success`); `None` when the env emits
424    /// no such signal.
425    pub success: Option<bool>,
426    /// Per-step means over the episode of the predict and env-step wall time,
427    /// in milliseconds (the same accounting as the `step_completed` hook
428    /// event); `None` for an episode that completed before its first step.
429    pub predict_ms: Option<f64>,
430    pub step_ms: Option<f64>,
431}
432
433/// What a finished or aborted session returns: totals plus the durable
434/// telemetry aggregate.
435#[derive(Debug, Clone, PartialEq)]
436pub struct RuntimeReport {
437    pub session_id: String,
438    pub env_id: String,
439    pub total_steps: i64,
440    pub total_episodes: i64,
441    /// Every completed episode, in completion order.
442    pub episodes: Vec<EpisodeSummary>,
443    /// Session-total telemetry aggregate (per-op latency/percentiles/bytes) —
444    /// the durable pull counterpart to the live `RuntimeHooks::on_telemetry` push.
445    pub telemetry: crate::telemetry::Snapshot,
446    /// Each distinct advisory the relay policy raised, in first-raised order.
447    pub advisories: Vec<Advisory>,
448}
449
450/// Per-op timeouts and the telemetry window for one session. Serialized with
451/// explicit millisecond field names (`*Ms`); legacy unsuffixed fields are rejected.
452#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
453#[serde(rename_all = "camelCase", deny_unknown_fields)]
454pub struct RuntimeLimits {
455    #[serde(
456        default = "default_connect_timeout",
457        rename = "envConnectTimeoutMs",
458        serialize_with = "duration_millis::serialize",
459        deserialize_with = "duration_millis::deserialize"
460    )]
461    pub env_connect_timeout: Duration,
462    #[serde(
463        default = "default_model_connect_timeout",
464        rename = "modelConnectTimeoutMs",
465        serialize_with = "duration_millis::serialize",
466        deserialize_with = "duration_millis::deserialize"
467    )]
468    pub model_connect_timeout: Duration,
469    #[serde(
470        default = "default_configure_route_timeout",
471        rename = "configureRouteTimeoutMs",
472        serialize_with = "duration_millis::serialize",
473        deserialize_with = "duration_millis::deserialize"
474    )]
475    pub configure_route_timeout: Duration,
476    #[serde(
477        default = "default_env_reset_timeout",
478        rename = "envResetTimeoutMs",
479        serialize_with = "duration_millis::serialize",
480        deserialize_with = "duration_millis::deserialize"
481    )]
482    pub env_reset_timeout: Duration,
483    #[serde(
484        default = "default_model_predict_timeout",
485        rename = "modelPredictTimeoutMs",
486        serialize_with = "duration_millis::serialize",
487        deserialize_with = "duration_millis::deserialize"
488    )]
489    pub model_predict_timeout: Duration,
490    #[serde(
491        default = "default_env_step_timeout",
492        rename = "envStepTimeoutMs",
493        serialize_with = "duration_millis::serialize",
494        deserialize_with = "duration_millis::deserialize"
495    )]
496    pub env_step_timeout: Duration,
497    #[serde(
498        default = "default_service_close_timeout",
499        rename = "serviceCloseTimeoutMs",
500        serialize_with = "duration_millis::serialize",
501        deserialize_with = "duration_millis::deserialize"
502    )]
503    pub service_close_timeout: Duration,
504    /// How often the background ticker pushes Window + Session telemetry
505    /// snapshots to `on_telemetry`, on a wall clock. `0` disables live streaming
506    /// (the final session snapshot is still delivered at session end); any
507    /// non-zero value below 1ms is floored to 1ms.
508    #[serde(
509        default = "default_telemetry_window",
510        rename = "telemetryWindowMs",
511        serialize_with = "duration_millis::serialize",
512        deserialize_with = "duration_millis::deserialize"
513    )]
514    pub telemetry_window: Duration,
515}
516
517impl Default for RuntimeLimits {
518    fn default() -> Self {
519        Self {
520            env_connect_timeout: default_connect_timeout(),
521            model_connect_timeout: default_model_connect_timeout(),
522            configure_route_timeout: default_configure_route_timeout(),
523            env_reset_timeout: default_env_reset_timeout(),
524            model_predict_timeout: default_model_predict_timeout(),
525            env_step_timeout: default_env_step_timeout(),
526            service_close_timeout: default_service_close_timeout(),
527            telemetry_window: default_telemetry_window(),
528        }
529    }
530}
531
532impl RuntimeLimits {
533    /// Clamped-non-negative i64 milliseconds; the proto field is uint64, so
534    /// callers `.max(0) as u64` without losing information.
535    pub fn env_step_timeout_ms(&self) -> i64 {
536        duration_ms_i64(self.env_step_timeout)
537    }
538
539    /// Clamped-non-negative i64 milliseconds; the proto field is uint64, so
540    /// callers `.max(0) as u64` without losing information.
541    pub fn env_reset_timeout_ms(&self) -> i64 {
542        duration_ms_i64(self.env_reset_timeout)
543    }
544}
545
546fn default_connect_timeout() -> Duration {
547    Duration::from_secs(60)
548}
549
550fn default_model_connect_timeout() -> Duration {
551    Duration::from_secs(600)
552}
553
554fn default_configure_route_timeout() -> Duration {
555    Duration::from_secs(600)
556}
557
558fn default_env_reset_timeout() -> Duration {
559    Duration::from_secs(300)
560}
561
562fn default_model_predict_timeout() -> Duration {
563    Duration::from_secs(300)
564}
565
566fn default_env_step_timeout() -> Duration {
567    Duration::from_secs(300)
568}
569
570fn default_service_close_timeout() -> Duration {
571    Duration::from_secs(5)
572}
573
574fn default_telemetry_window() -> Duration {
575    Duration::from_secs(1)
576}
577
578fn duration_ms_i64(duration: Duration) -> i64 {
579    duration.as_millis().try_into().unwrap_or(i64::MAX)
580}
581
582mod duration_millis {
583    use std::time::Duration;
584
585    use serde::{Deserialize, Deserializer, Serializer};
586
587    pub(super) fn serialize<S>(duration: &Duration, serializer: S) -> Result<S::Ok, S::Error>
588    where
589        S: Serializer,
590    {
591        let millis = duration.as_millis().try_into().unwrap_or(u64::MAX);
592        serializer.serialize_u64(millis)
593    }
594
595    pub(super) fn deserialize<'de, D>(deserializer: D) -> Result<Duration, D::Error>
596    where
597        D: Deserializer<'de>,
598    {
599        Ok(Duration::from_millis(u64::deserialize(deserializer)?))
600    }
601}
602
603#[cfg(test)]
604mod tests {
605    use std::time::Duration;
606
607    use rlmesh_proto::core::v1::{AutoresetMode, EnvContract, EnvSpec};
608    use rlmesh_proto::spaces::v1::SpaceSpec;
609    use serde_json::json;
610
611    use super::{
612        ENV_RESET_OPTIONS_KEY, MetaKind, MetaList, MetaMap, MetaValue, RuntimeLimits,
613        RuntimeSessionSpec, TRIAL_INDEX_OPTION, declares_reset_option, reset_options_for,
614    };
615
616    fn valid_spec() -> RuntimeSessionSpec {
617        RuntimeSessionSpec {
618            session_id: "session".to_string(),
619            env_id: "env-id".to_string(),
620            env_component_id: "env".to_string(),
621            model_component_id: "model".to_string(),
622            workflow_edition: rlmesh_proto::parse_retained_edition(
623                rlmesh_proto::CURRENT_WORKFLOW_EDITION,
624            )
625            .expect("this build drives its own edition"),
626            env_contract: EnvContract {
627                spec: Some(EnvSpec {
628                    observation_space: Some(SpaceSpec::default()),
629                    action_space: Some(SpaceSpec::default()),
630                    ..Default::default()
631                }),
632                num_envs: 1,
633                ..Default::default()
634            },
635            num_envs: 1,
636            episode_seeds: Vec::new(),
637            base_seed: None,
638            max_episodes: Some(1),
639            trial_index_base: None,
640            max_episode_steps: None,
641            max_episode_seconds: None,
642            close_env_on_end: true,
643            subset_step: false,
644            limits: RuntimeLimits::default(),
645            env_ceiling: None,
646            model_ceiling: None,
647        }
648    }
649
650    #[test]
651    fn validate_rejects_a_trial_base_the_env_owns_the_resets_for() {
652        let mut spec = valid_spec();
653        spec.trial_index_base = Some(5);
654        // Driver-owned resets (the default DISABLED/UNSPECIFIED) carry the ordinal.
655        assert!(spec.validate().is_ok());
656
657        spec.env_contract.autoreset_mode = AutoresetMode::NextStep as i32;
658        let error = spec.validate().unwrap_err();
659        assert!(
660            error.contains("trial_index_base requires an env with autoreset disabled"),
661            "expected the driver-owned-reset rule, got: {error}"
662        );
663    }
664
665    #[test]
666    fn validate_accepts_the_default_trial_base_under_next_step() {
667        // The ordinal is on by default, so the default base must not fail a run
668        // the env owns the resets for -- whether it arrives unset or as an
669        // explicit 0 (the Python surface always passes an integer).
670        let mut spec = valid_spec();
671        spec.env_contract.autoreset_mode = AutoresetMode::NextStep as i32;
672        for base in [None, Some(0)] {
673            spec.trial_index_base = base;
674            assert_eq!(spec.trial_index_base(), 0);
675            assert!(spec.validate().is_ok(), "base {base:?} must validate");
676        }
677    }
678
679    #[test]
680    fn reset_options_key_is_the_published_string() {
681        // Pinned against the Python mirror (`rlmesh.ENV_RESET_OPTIONS_KEY`, itself
682        // re-exported from this constant) by
683        // `python/rlmesh/tests/unit/test_reset_options.py`.
684        assert_eq!(super::ENV_RESET_OPTIONS_KEY, "rlmesh.env.v1.reset_options");
685    }
686
687    #[test]
688    fn a_ceiling_edition_is_the_highest_one_both_sides_implement() {
689        let names = |names: &[&str]| {
690            names
691                .iter()
692                .map(|name| name.to_string())
693                .collect::<Vec<_>>()
694        };
695        assert_eq!(
696            super::PeerCeiling::highest_shared_edition(&names(&[
697                "2099.01",
698                rlmesh_proto::CURRENT_WORKFLOW_EDITION,
699            ])),
700            Some(rlmesh_proto::Edition::current())
701        );
702        assert_eq!(
703            super::PeerCeiling::highest_shared_edition(&names(&["2099.01"])),
704            None
705        );
706    }
707
708    /// The spec carries a typed edition, so a made-up name is refused one step
709    /// earlier — at the shared parse boundary the env client and every other
710    /// string-carrying caller go through — naming the arrived string and the
711    /// retained list. That refusal is what an operator actually sees.
712    #[test]
713    fn an_edition_the_runtime_cannot_drive_is_refused_by_name() {
714        let error = rlmesh_proto::parse_retained_edition("2099.01")
715            .expect_err("an unimplemented edition is refused");
716        assert!(
717            error.contains("2099.01")
718                && error.contains("cannot drive")
719                && error.contains(rlmesh_proto::CURRENT_WORKFLOW_EDITION),
720            "expected an edition-refusal error naming the retained list, got: {error}"
721        );
722
723        // The edition the build implements is accepted, by base name and by this
724        // build's cohort spelling of it, and validates.
725        let mut spec = valid_spec();
726        for name in [
727            rlmesh_proto::CURRENT_WORKFLOW_EDITION,
728            rlmesh_proto::WORKFLOW_EDITION_BASE,
729        ] {
730            spec.workflow_edition = rlmesh_proto::parse_retained_edition(name)
731                .unwrap_or_else(|error| panic!("{name} must parse: {error}"));
732            assert!(spec.validate().is_ok());
733        }
734    }
735
736    #[test]
737    fn space_accessors_return_none_on_unvalidated_spec() {
738        let mut spec = valid_spec();
739        // An unvalidated spec is trivially constructible since all fields are
740        // public; the accessors must not panic.
741        spec.env_contract = EnvContract::default();
742
743        assert!(spec.observation_space().is_none());
744        assert!(spec.action_space().is_none());
745    }
746
747    #[test]
748    fn space_accessors_return_some_on_populated_spec() {
749        let spec = valid_spec();
750        assert!(spec.observation_space().is_some());
751        assert!(spec.action_space().is_some());
752    }
753
754    #[test]
755    fn validate_accepts_vectorized_next_step_runtime_sessions() {
756        // num_envs > 1 is supported with NEXT_STEP autoreset: the env resets each
757        // done lane itself, so the driver never needs per-lane reset.
758        let mut spec = valid_spec();
759        spec.num_envs = 4;
760        spec.env_contract.num_envs = 4;
761        spec.env_contract.autoreset_mode = AutoresetMode::NextStep as i32;
762
763        assert!(spec.validate().is_ok());
764    }
765
766    #[test]
767    fn validate_rejects_disabled_vectorized_sessions() {
768        // DISABLED (and the UNSPECIFIED default) with num_envs > 1 needs per-lane
769        // reset, which stock gymnasium vector envs cannot do. Reject up front
770        // rather than failing mid-run on the first staggered termination.
771        for mode in [AutoresetMode::Disabled, AutoresetMode::Unspecified] {
772            let mut spec = valid_spec();
773            spec.num_envs = 4;
774            spec.env_contract.num_envs = 4;
775            spec.env_contract.autoreset_mode = mode as i32;
776
777            let error = spec.validate().unwrap_err();
778            assert!(
779                error.contains("NEXT_STEP"),
780                "expected a NEXT_STEP-guidance rejection for {mode:?}, got: {error}"
781            );
782        }
783    }
784
785    #[test]
786    fn validate_rejects_same_step_autoreset() {
787        // SAME_STEP is reserved but unsupported; validation must reject it so it
788        // cannot reach the runtime and stall lanes.
789        let mut spec = valid_spec();
790        spec.env_contract.autoreset_mode = AutoresetMode::SameStep as i32;
791
792        let error = spec.validate().unwrap_err();
793        assert!(
794            error.contains("SAME_STEP"),
795            "expected SAME_STEP rejection, got: {error}"
796        );
797    }
798
799    #[test]
800    fn runtime_limits_json_uses_explicit_millisecond_fields() {
801        let value = serde_json::to_value(RuntimeLimits::default()).unwrap();
802
803        assert_eq!(value["envConnectTimeoutMs"], json!(60_000));
804        assert_eq!(value["modelConnectTimeoutMs"], json!(600_000));
805        assert_eq!(value["configureRouteTimeoutMs"], json!(600_000));
806        assert_eq!(value["envResetTimeoutMs"], json!(300_000));
807        assert_eq!(value["modelPredictTimeoutMs"], json!(300_000));
808        assert_eq!(value["envStepTimeoutMs"], json!(300_000));
809        assert_eq!(value["serviceCloseTimeoutMs"], json!(5_000));
810        assert_eq!(value["telemetryWindowMs"], json!(1_000));
811        assert!(value.get("envConnectTimeout").is_none());
812
813        let parsed: RuntimeLimits = serde_json::from_value(json!({
814            "envConnectTimeoutMs": 1,
815            "modelConnectTimeoutMs": 2,
816            "configureRouteTimeoutMs": 3,
817            "envResetTimeoutMs": 4,
818            "modelPredictTimeoutMs": 5,
819            "envStepTimeoutMs": 6,
820            "serviceCloseTimeoutMs": 7,
821            "telemetryWindowMs": 8
822        }))
823        .unwrap();
824
825        assert_eq!(parsed.env_connect_timeout, Duration::from_millis(1));
826        assert_eq!(parsed.model_connect_timeout, Duration::from_millis(2));
827        assert_eq!(parsed.configure_route_timeout, Duration::from_millis(3));
828        assert_eq!(parsed.env_reset_timeout, Duration::from_millis(4));
829        assert_eq!(parsed.model_predict_timeout, Duration::from_millis(5));
830        assert_eq!(parsed.env_step_timeout, Duration::from_millis(6));
831        assert_eq!(parsed.service_close_timeout, Duration::from_millis(7));
832        assert_eq!(parsed.telemetry_window, Duration::from_millis(8));
833    }
834
835    #[test]
836    fn runtime_limits_reject_legacy_unsuffixed_fields() {
837        let error = serde_json::from_value::<RuntimeLimits>(json!({
838            "envConnectTimeout": 1
839        }))
840        .unwrap_err();
841
842        assert!(error.to_string().contains("envConnectTimeout"));
843    }
844
845    /// An env contract whose metadata declares `reset_options = declared`.
846    fn contract_declaring(declared: MetaValue) -> EnvContract {
847        EnvContract {
848            spec: Some(EnvSpec {
849                metadata: Some(MetaMap {
850                    entries: [(ENV_RESET_OPTIONS_KEY.to_string(), declared)].into(),
851                }),
852                ..Default::default()
853            }),
854            ..Default::default()
855        }
856    }
857
858    fn text(value: &str) -> MetaValue {
859        MetaValue {
860            kind: Some(MetaKind::Text(value.to_string())),
861        }
862    }
863
864    fn list(items: Vec<MetaValue>) -> MetaValue {
865        MetaValue {
866            kind: Some(MetaKind::List(MetaList { items })),
867        }
868    }
869
870    fn trial_option(options: &MetaMap) -> Option<&MetaKind> {
871        options.entries.get(TRIAL_INDEX_OPTION)?.kind.as_ref()
872    }
873
874    #[test]
875    fn declares_reset_option_reads_a_list_or_a_bare_string() {
876        assert!(declares_reset_option(
877            &contract_declaring(list(vec![text("other"), text(TRIAL_INDEX_OPTION)])),
878            TRIAL_INDEX_OPTION,
879        ));
880        assert!(declares_reset_option(
881            &contract_declaring(text(TRIAL_INDEX_OPTION)),
882            TRIAL_INDEX_OPTION,
883        ));
884        assert!(!declares_reset_option(
885            &contract_declaring(list(vec![text("other")])),
886            TRIAL_INDEX_OPTION,
887        ));
888        assert!(!declares_reset_option(
889            &contract_declaring(MetaValue {
890                kind: Some(MetaKind::Integer(1)),
891            }),
892            TRIAL_INDEX_OPTION,
893        ));
894        assert!(!declares_reset_option(
895            &EnvContract::default(),
896            TRIAL_INDEX_OPTION,
897        ));
898    }
899
900    #[test]
901    fn reset_options_for_sends_an_integer_per_lane_and_a_list_for_many() {
902        let contract = contract_declaring(list(vec![text(TRIAL_INDEX_OPTION)]));
903
904        let single =
905            reset_options_for(&contract, TRIAL_INDEX_OPTION, &[7]).expect("single-lane options");
906        assert_eq!(trial_option(&single), Some(&MetaKind::Integer(7)));
907
908        let many = reset_options_for(&contract, TRIAL_INDEX_OPTION, &[7, 8, 9])
909            .expect("multi-lane options");
910        assert_eq!(
911            trial_option(&many),
912            Some(&MetaKind::List(MetaList {
913                items: [7, 8, 9]
914                    .into_iter()
915                    .map(|trial| MetaValue {
916                        kind: Some(MetaKind::Integer(trial)),
917                    })
918                    .collect(),
919            })),
920        );
921    }
922
923    #[test]
924    fn reset_options_for_withholds_without_trials_or_a_declaration() {
925        let contract = contract_declaring(list(vec![text(TRIAL_INDEX_OPTION)]));
926
927        assert_eq!(reset_options_for(&contract, TRIAL_INDEX_OPTION, &[]), None);
928        assert_eq!(
929            reset_options_for(
930                &contract_declaring(list(vec![text("other")])),
931                TRIAL_INDEX_OPTION,
932                &[7]
933            ),
934            None,
935        );
936        assert_eq!(
937            reset_options_for(&EnvContract::default(), TRIAL_INDEX_OPTION, &[7]),
938            None
939        );
940    }
941}