rlmesh-runtime 0.1.0-rc.14

Internal RLMesh crate (unstable Rust API): runtime driver for evaluation sessions.
Documentation
//! Event payloads the driver hands to each [`RuntimeHooks`](super::RuntimeHooks)
//! callback. Each carries its route/session identity inline, since one hooks
//! instance serves every concurrent route.

use std::sync::Arc;

use prost::bytes::Bytes;
use rlmesh_proto::spaces::v1::MetaMap;
use rlmesh_proto::spaces::v1::SpaceSpec;

#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct RuntimeEnvContext {
    pub env_id: String,
    pub env_component_id: String,
    pub model_component_id: String,
    /// The lane this session drives when the env is lane-driven (one session
    /// per lane over a shared endpoint); `None` for a whole-env session. Lets
    /// telemetry and lifecycle events be sliced per lane after the fact.
    pub lane: Option<u32>,
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct EnvConnectedEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub env_id: String,
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ModelConnectedEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SessionStartedEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub env_id: String,
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SessionEndedEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub reason: String,
    pub total_steps: i64,
    pub total_episodes: i64,
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SessionFailedEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub reason: String,
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LogEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub level: LogLevel,
    pub message: String,
    pub source: Option<String>,
}

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LogLevel {
    Debug,
    Info,
    Warn,
    Error,
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct EpisodeStartedEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub episode_id: String,
    pub episode_record_id: String,
    pub episode_index: i64,
    pub env_index: i32,
    pub started_from_auto_reset: bool,
    pub seed: Option<i64>,
    /// The trial ordinal this episode walks (`trial_index_base` + its
    /// episode-start position); `None` only under `NEXT_STEP` autoreset, where
    /// no ordinal is minted.
    pub trial_index: Option<u64>,
}

#[derive(Debug, Clone, PartialEq)]
pub struct EpisodeCompletedEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub episode_id: String,
    pub episode_record_id: String,
    pub episode_index: i64,
    pub env_index: i32,
    pub step_count: i64,
    pub cumulative_reward: f64,
    pub terminated: bool,
    pub truncated: bool,
    pub duration_ms: i64,
    pub final_info: Option<MetaMap>,
    /// The env-reported task outcome from the final step's info (`is_success`,
    /// `success`, or `task_success`, first present wins; numbers coerce by
    /// truthiness), the same value the episode summary carries, so a hook
    /// never re-derives it. `None` when the env emits no such key.
    pub success: Option<bool>,
    pub seed: Option<i64>,
    /// The trial ordinal this episode walked (`trial_index_base` + its
    /// episode-start position); `None` only under `NEXT_STEP` autoreset, where
    /// no ordinal is minted.
    pub trial_index: Option<u64>,
    /// Per-step means over the episode of the predict and env-step wall time
    /// (see `StepCompletedEvent::predict_ms`), in milliseconds; `None` for an
    /// episode that completed before its first step.
    pub predict_ms: Option<f64>,
    pub step_ms: Option<f64>,
}

#[derive(Debug, Clone, PartialEq)]
pub struct ActionReceivedEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub episode_id: String,
    pub episode_record_id: String,
    pub episode_ids: Vec<String>,
    pub episode_record_ids: Vec<String>,
    pub step: i64,
    pub env_index: i32,
    /// Shared so the per-step, per-hook event fan-out clones an `Arc` pointer
    /// rather than deep-copying the action space spec on every step. Types
    /// `action`: when the relay policy converted it into a new layout, this is
    /// the converted space and `raw_action` keeps the contract's.
    pub action_space: Arc<SpaceSpec>,
    /// Opaque per-leaf wire bytes; the relay is content-blind (§13).
    ///
    /// To persist as a single artifact, use `rlmesh-grpc`'s
    /// `wire::leaves_to_blob` — plain concatenation loses the leaf
    /// boundaries, which cannot be recovered for variable-length leaves.
    pub action: Option<Vec<Bytes>>,
    /// The model's action before `transform_action` ran; equals `action` when no
    /// transform changed it.
    pub raw_action: Option<Vec<Bytes>>,
}

#[derive(Debug, Clone, PartialEq)]
pub struct StepCompletedEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub episode_id: String,
    pub episode_record_id: String,
    pub step: i64,
    pub env_index: i32,
    pub rewards: Vec<f64>,
    pub infos: Option<MetaMap>,
    /// Per lane, aligned with `rewards`: whether this step ended the lane's
    /// episode in a terminal state / by truncation (env-reported or a runtime
    /// cap), known at the step itself, so a hook never waits for the
    /// `episode_completed` event to learn how a step ended.
    pub terminated: Vec<bool>,
    pub truncated: Vec<bool>,
    /// Per lane, aligned with `rewards`: true where this response is the
    /// lane's `NEXT_STEP` autoreset roll -- the new episode's reset observation,
    /// still attributed to the ended id, not a step of any episode. Lanes roll
    /// independently, so the other lanes' entries are real steps. `infos` is
    /// `None` when any lane rolled.
    pub autoreset_roll: Vec<bool>,
    /// Wall time, in milliseconds, of the predict(s) that landed since the
    /// group's previous step -- zero for a step served from chunk replay, and
    /// under prefetch the next chunk's predict, attributed to the step it
    /// landed on -- and of this step's env round trip. Every lane of the
    /// group experienced both.
    pub predict_ms: f64,
    pub step_ms: f64,
}

#[derive(Debug, Clone, PartialEq)]
pub struct ObservationEmittedEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub episode_id: String,
    pub episode_record_id: String,
    pub episode_ids: Vec<String>,
    pub episode_record_ids: Vec<String>,
    pub step: i64,
    pub env_index: i32,
    pub is_reset: bool,
    pub num_envs: u32,
    /// Shared so the per-step, per-hook event fan-out clones an `Arc` pointer
    /// rather than deep-copying the observation space spec on every step. Types
    /// `observation`: when the relay policy converted it into a new layout, this
    /// is the converted space and `raw_observation` keeps the contract's.
    pub observation_space: Arc<SpaceSpec>,
    /// Opaque per-leaf wire bytes; the relay is content-blind (§13).
    ///
    /// To persist as a single artifact, use `rlmesh-grpc`'s
    /// `wire::leaves_to_blob` — plain concatenation loses the leaf
    /// boundaries, which cannot be recovered for variable-length leaves.
    pub observation: Option<Vec<Bytes>>,
    /// The env's observation before `transform_observation` ran; equals
    /// `observation` when no transform changed it.
    pub raw_observation: Option<Vec<Bytes>>,
    /// The env's reset infos when `is_reset`; step infos ride on
    /// `StepCompletedEvent` instead.
    pub infos: Option<MetaMap>,
}

/// A relay policy converted a payload for a peer that could not decode it as
/// sent. Delivered once per distinct advisory, which the report also carries.
#[derive(Debug, Clone, PartialEq)]
pub struct RelayAdvisoryEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub leg: super::Leg,
    pub advisory: rlmesh_spaces::Advisory,
}

/// A live telemetry snapshot tagged with the route/session it belongs to.
///
/// The managed runner shares one `RuntimeHooks` across concurrent routes, so —
/// like every other event — the snapshot carries identity inline. Identity
/// lives on the event, not on the `Snapshot` itself, which stays a pure metrics
/// payload (it is also the durable `RuntimeReport.telemetry`). Use
/// `snapshot.horizon` to tell a window tick from the cumulative session.
#[derive(Debug, Clone, PartialEq)]
pub struct TelemetrySnapshotEvent {
    pub session_id: String,
    pub route: RuntimeEnvContext,
    pub snapshot: crate::telemetry::Snapshot,
}