Skip to main content

rlmesh_runtime/hooks/
events.rs

1//! Event payloads the driver hands to each [`RuntimeHooks`](super::RuntimeHooks)
2//! callback. Each carries its route/session identity inline, since one hooks
3//! instance serves every concurrent route.
4
5use std::sync::Arc;
6
7use prost::bytes::Bytes;
8use rlmesh_proto::spaces::v1::MetaMap;
9use rlmesh_proto::spaces::v1::SpaceSpec;
10
11#[derive(Debug, Clone, Default, PartialEq, Eq)]
12pub struct RuntimeEnvContext {
13    pub env_id: String,
14    pub env_component_id: String,
15    pub model_component_id: String,
16    /// The lane this session drives when the env is lane-driven (one session
17    /// per lane over a shared endpoint); `None` for a whole-env session. Lets
18    /// telemetry and lifecycle events be sliced per lane after the fact.
19    pub lane: Option<u32>,
20}
21
22#[derive(Debug, Clone, PartialEq, Eq)]
23pub struct EnvConnectedEvent {
24    pub session_id: String,
25    pub route: RuntimeEnvContext,
26    pub env_id: String,
27}
28
29#[derive(Debug, Clone, PartialEq, Eq)]
30pub struct ModelConnectedEvent {
31    pub session_id: String,
32    pub route: RuntimeEnvContext,
33}
34
35#[derive(Debug, Clone, PartialEq, Eq)]
36pub struct SessionStartedEvent {
37    pub session_id: String,
38    pub route: RuntimeEnvContext,
39    pub env_id: String,
40}
41
42#[derive(Debug, Clone, PartialEq, Eq)]
43pub struct SessionEndedEvent {
44    pub session_id: String,
45    pub route: RuntimeEnvContext,
46    pub reason: String,
47    pub total_steps: i64,
48    pub total_episodes: i64,
49}
50
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct SessionFailedEvent {
53    pub session_id: String,
54    pub route: RuntimeEnvContext,
55    pub reason: String,
56}
57
58#[derive(Debug, Clone, PartialEq, Eq)]
59pub struct LogEvent {
60    pub session_id: String,
61    pub route: RuntimeEnvContext,
62    pub level: LogLevel,
63    pub message: String,
64    pub source: Option<String>,
65}
66
67#[derive(Debug, Clone, Copy, PartialEq, Eq)]
68pub enum LogLevel {
69    Debug,
70    Info,
71    Warn,
72    Error,
73}
74
75#[derive(Debug, Clone, PartialEq, Eq)]
76pub struct EpisodeStartedEvent {
77    pub session_id: String,
78    pub route: RuntimeEnvContext,
79    pub episode_id: String,
80    pub episode_record_id: String,
81    pub episode_index: i64,
82    pub env_index: i32,
83    pub started_from_auto_reset: bool,
84    pub seed: Option<i64>,
85    /// The trial ordinal this episode walks (`trial_index_base` + its
86    /// episode-start position); `None` only under `NEXT_STEP` autoreset, where
87    /// no ordinal is minted.
88    pub trial_index: Option<u64>,
89}
90
91#[derive(Debug, Clone, PartialEq)]
92pub struct EpisodeCompletedEvent {
93    pub session_id: String,
94    pub route: RuntimeEnvContext,
95    pub episode_id: String,
96    pub episode_record_id: String,
97    pub episode_index: i64,
98    pub env_index: i32,
99    pub step_count: i64,
100    pub cumulative_reward: f64,
101    pub terminated: bool,
102    pub truncated: bool,
103    pub duration_ms: i64,
104    pub final_info: Option<MetaMap>,
105    /// The env-reported task outcome from the final step's info (`is_success`,
106    /// `success`, or `task_success`, first present wins; numbers coerce by
107    /// truthiness), the same value the episode summary carries, so a hook
108    /// never re-derives it. `None` when the env emits no such key.
109    pub success: Option<bool>,
110    pub seed: Option<i64>,
111    /// The trial ordinal this episode walked (`trial_index_base` + its
112    /// episode-start position); `None` only under `NEXT_STEP` autoreset, where
113    /// no ordinal is minted.
114    pub trial_index: Option<u64>,
115    /// Per-step means over the episode of the predict and env-step wall time
116    /// (see `StepCompletedEvent::predict_ms`), in milliseconds; `None` for an
117    /// episode that completed before its first step.
118    pub predict_ms: Option<f64>,
119    pub step_ms: Option<f64>,
120}
121
122#[derive(Debug, Clone, PartialEq)]
123pub struct ActionReceivedEvent {
124    pub session_id: String,
125    pub route: RuntimeEnvContext,
126    pub episode_id: String,
127    pub episode_record_id: String,
128    pub episode_ids: Vec<String>,
129    pub episode_record_ids: Vec<String>,
130    pub step: i64,
131    pub env_index: i32,
132    /// Shared so the per-step, per-hook event fan-out clones an `Arc` pointer
133    /// rather than deep-copying the action space spec on every step. Types
134    /// `action`: when the relay policy converted it into a new layout, this is
135    /// the converted space and `raw_action` keeps the contract's.
136    pub action_space: Arc<SpaceSpec>,
137    /// Opaque per-leaf wire bytes; the relay is content-blind (§13).
138    ///
139    /// To persist as a single artifact, use `rlmesh-grpc`'s
140    /// `wire::leaves_to_blob` — plain concatenation loses the leaf
141    /// boundaries, which cannot be recovered for variable-length leaves.
142    pub action: Option<Vec<Bytes>>,
143    /// The model's action before `transform_action` ran; equals `action` when no
144    /// transform changed it.
145    pub raw_action: Option<Vec<Bytes>>,
146}
147
148#[derive(Debug, Clone, PartialEq)]
149pub struct StepCompletedEvent {
150    pub session_id: String,
151    pub route: RuntimeEnvContext,
152    pub episode_id: String,
153    pub episode_record_id: String,
154    pub step: i64,
155    pub env_index: i32,
156    pub rewards: Vec<f64>,
157    pub infos: Option<MetaMap>,
158    /// Per lane, aligned with `rewards`: whether this step ended the lane's
159    /// episode in a terminal state / by truncation (env-reported or a runtime
160    /// cap), known at the step itself, so a hook never waits for the
161    /// `episode_completed` event to learn how a step ended.
162    pub terminated: Vec<bool>,
163    pub truncated: Vec<bool>,
164    /// Per lane, aligned with `rewards`: true where this response is the
165    /// lane's `NEXT_STEP` autoreset roll -- the new episode's reset observation,
166    /// still attributed to the ended id, not a step of any episode. Lanes roll
167    /// independently, so the other lanes' entries are real steps. `infos` is
168    /// `None` when any lane rolled.
169    pub autoreset_roll: Vec<bool>,
170    /// Wall time, in milliseconds, of the predict(s) that landed since the
171    /// group's previous step -- zero for a step served from chunk replay, and
172    /// under prefetch the next chunk's predict, attributed to the step it
173    /// landed on -- and of this step's env round trip. Every lane of the
174    /// group experienced both.
175    pub predict_ms: f64,
176    pub step_ms: f64,
177}
178
179#[derive(Debug, Clone, PartialEq)]
180pub struct ObservationEmittedEvent {
181    pub session_id: String,
182    pub route: RuntimeEnvContext,
183    pub episode_id: String,
184    pub episode_record_id: String,
185    pub episode_ids: Vec<String>,
186    pub episode_record_ids: Vec<String>,
187    pub step: i64,
188    pub env_index: i32,
189    pub is_reset: bool,
190    pub num_envs: u32,
191    /// Shared so the per-step, per-hook event fan-out clones an `Arc` pointer
192    /// rather than deep-copying the observation space spec on every step. Types
193    /// `observation`: when the relay policy converted it into a new layout, this
194    /// is the converted space and `raw_observation` keeps the contract's.
195    pub observation_space: Arc<SpaceSpec>,
196    /// Opaque per-leaf wire bytes; the relay is content-blind (§13).
197    ///
198    /// To persist as a single artifact, use `rlmesh-grpc`'s
199    /// `wire::leaves_to_blob` — plain concatenation loses the leaf
200    /// boundaries, which cannot be recovered for variable-length leaves.
201    pub observation: Option<Vec<Bytes>>,
202    /// The env's observation before `transform_observation` ran; equals
203    /// `observation` when no transform changed it.
204    pub raw_observation: Option<Vec<Bytes>>,
205    /// The env's reset infos when `is_reset`; step infos ride on
206    /// `StepCompletedEvent` instead.
207    pub infos: Option<MetaMap>,
208}
209
210/// A relay policy converted a payload for a peer that could not decode it as
211/// sent. Delivered once per distinct advisory, which the report also carries.
212#[derive(Debug, Clone, PartialEq)]
213pub struct RelayAdvisoryEvent {
214    pub session_id: String,
215    pub route: RuntimeEnvContext,
216    pub leg: super::Leg,
217    pub advisory: rlmesh_spaces::Advisory,
218}
219
220/// A live telemetry snapshot tagged with the route/session it belongs to.
221///
222/// The managed runner shares one `RuntimeHooks` across concurrent routes, so —
223/// like every other event — the snapshot carries identity inline. Identity
224/// lives on the event, not on the `Snapshot` itself, which stays a pure metrics
225/// payload (it is also the durable `RuntimeReport.telemetry`). Use
226/// `snapshot.horizon` to tell a window tick from the cumulative session.
227#[derive(Debug, Clone, PartialEq)]
228pub struct TelemetrySnapshotEvent {
229    pub session_id: String,
230    pub route: RuntimeEnvContext,
231    pub snapshot: crate::telemetry::Snapshot,
232}