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}
17
18#[derive(Debug, Clone, PartialEq, Eq)]
19pub struct EnvConnectedEvent {
20    pub session_id: String,
21    pub route: RuntimeEnvContext,
22    pub env_id: String,
23}
24
25#[derive(Debug, Clone, PartialEq, Eq)]
26pub struct ModelConnectedEvent {
27    pub session_id: String,
28    pub route: RuntimeEnvContext,
29}
30
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct SessionStartedEvent {
33    pub session_id: String,
34    pub route: RuntimeEnvContext,
35    pub env_id: String,
36}
37
38#[derive(Debug, Clone, PartialEq, Eq)]
39pub struct SessionEndedEvent {
40    pub session_id: String,
41    pub route: RuntimeEnvContext,
42    pub reason: String,
43    pub total_steps: i64,
44    pub total_episodes: i64,
45}
46
47#[derive(Debug, Clone, PartialEq, Eq)]
48pub struct SessionFailedEvent {
49    pub session_id: String,
50    pub route: RuntimeEnvContext,
51    pub reason: String,
52}
53
54#[derive(Debug, Clone, PartialEq, Eq)]
55pub struct LogEvent {
56    pub session_id: String,
57    pub route: RuntimeEnvContext,
58    pub level: LogLevel,
59    pub message: String,
60    pub source: Option<String>,
61}
62
63#[derive(Debug, Clone, Copy, PartialEq, Eq)]
64pub enum LogLevel {
65    Debug,
66    Info,
67    Warn,
68    Error,
69}
70
71#[derive(Debug, Clone, PartialEq, Eq)]
72pub struct EpisodeStartedEvent {
73    pub session_id: String,
74    pub route: RuntimeEnvContext,
75    pub episode_id: String,
76    pub episode_record_id: String,
77    pub episode_index: i64,
78    pub env_index: i32,
79    pub started_from_auto_reset: bool,
80}
81
82#[derive(Debug, Clone, PartialEq)]
83pub struct EpisodeCompletedEvent {
84    pub session_id: String,
85    pub route: RuntimeEnvContext,
86    pub episode_id: String,
87    pub episode_record_id: String,
88    pub episode_index: i64,
89    pub env_index: i32,
90    pub step_count: i64,
91    pub cumulative_reward: f64,
92    pub terminated: bool,
93    pub truncated: bool,
94    pub duration_ms: i64,
95    pub final_info: Option<MetaMap>,
96}
97
98#[derive(Debug, Clone, PartialEq)]
99pub struct ActionReceivedEvent {
100    pub session_id: String,
101    pub route: RuntimeEnvContext,
102    pub episode_id: String,
103    pub episode_record_id: String,
104    pub episode_ids: Vec<String>,
105    pub episode_record_ids: Vec<String>,
106    pub step: i64,
107    pub env_index: i32,
108    /// Shared so the per-step, per-hook event fan-out clones an `Arc` pointer
109    /// rather than deep-copying the action space spec on every step.
110    pub action_space: Arc<SpaceSpec>,
111    /// Opaque per-leaf wire bytes; the relay is content-blind (§13).
112    pub action: Option<Vec<Bytes>>,
113}
114
115#[derive(Debug, Clone, PartialEq)]
116pub struct StepCompletedEvent {
117    pub session_id: String,
118    pub route: RuntimeEnvContext,
119    pub episode_id: String,
120    pub episode_record_id: String,
121    pub step: i64,
122    pub env_index: i32,
123    pub rewards: Vec<f64>,
124}
125
126#[derive(Debug, Clone, PartialEq)]
127pub struct ObservationEmittedEvent {
128    pub session_id: String,
129    pub route: RuntimeEnvContext,
130    pub episode_id: String,
131    pub episode_record_id: String,
132    pub episode_ids: Vec<String>,
133    pub episode_record_ids: Vec<String>,
134    pub step: i64,
135    pub env_index: i32,
136    pub is_reset: bool,
137    pub num_envs: u32,
138    /// Shared so the per-step, per-hook event fan-out clones an `Arc` pointer
139    /// rather than deep-copying the observation space spec on every step.
140    pub observation_space: Arc<SpaceSpec>,
141    /// Opaque per-leaf wire bytes; the relay is content-blind (§13).
142    pub observation: Option<Vec<Bytes>>,
143}
144
145/// A live telemetry snapshot tagged with the route/session it belongs to.
146///
147/// The managed runner shares one `RuntimeHooks` across concurrent routes, so —
148/// like every other event — the snapshot carries identity inline. Identity
149/// lives on the event, not on the `Snapshot` itself, which stays a pure metrics
150/// payload (it is also the durable `RuntimeReport.telemetry`). Use
151/// `snapshot.horizon` to tell a window tick from the cumulative session.
152#[derive(Debug, Clone, PartialEq)]
153pub struct TelemetrySnapshotEvent {
154    pub session_id: String,
155    pub route: RuntimeEnvContext,
156    pub snapshot: crate::telemetry::Snapshot,
157}