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}