Skip to main content

rlmesh_runtime/hooks/
traits.rs

1//! The [`RuntimeHooks`] observer trait and its no-op default.
2
3use async_trait::async_trait;
4use prost::bytes::Bytes;
5
6use super::{
7    ActionReceivedEvent, EnvConnectedEvent, EpisodeCompletedEvent, EpisodeStartedEvent, LogEvent,
8    ModelConnectedEvent, ObservationEmittedEvent, RelayAdvisoryEvent, SessionEndedEvent,
9    SessionFailedEvent, SessionStartedEvent, StepCompletedEvent, TelemetrySnapshotEvent,
10};
11
12/// Error returned by a [`RuntimeHooks`] callback.
13#[derive(Debug, thiserror::Error)]
14#[non_exhaustive]
15pub enum HookError {
16    #[error("{0}")]
17    Message(String),
18}
19
20/// A [`RuntimeHooks`] that ignores every event — the default when a session
21/// installs no observer.
22#[derive(Debug, Default)]
23pub struct NoopRuntimeHooks;
24
25#[async_trait]
26impl RuntimeHooks for NoopRuntimeHooks {}
27
28/// Callbacks the driver fans out around the `reset -> predict -> step` loop.
29///
30/// Lifecycle, progress, and log hooks are best-effort: the driver logs a failure
31/// and keeps the route moving. The two transform hooks
32/// ([`Self::transform_action`], [`Self::transform_observation`]) are fatal — a
33/// failed transform leaves the next wire payload undefined, so the route fails
34/// and shuts down. One shared instance serves every concurrent route, so each
35/// event carries its route/session identity inline.
36#[async_trait]
37pub trait RuntimeHooks: Send + Sync {
38    async fn env_connected(&self, _event: EnvConnectedEvent) -> Result<(), HookError> {
39        Ok(())
40    }
41
42    async fn model_connected(&self, _event: ModelConnectedEvent) -> Result<(), HookError> {
43        Ok(())
44    }
45
46    async fn session_started(&self, _event: SessionStartedEvent) -> Result<(), HookError> {
47        Ok(())
48    }
49
50    async fn episode_started(&self, _event: EpisodeStartedEvent) -> Result<(), HookError> {
51        Ok(())
52    }
53
54    /// The last event under its episode id: no observation or step of that
55    /// episode follows it. Under `NEXT_STEP` autoreset the ended episode is
56    /// observed and stepped once more (the autoreset step) before this fires.
57    async fn episode_completed(&self, _event: EpisodeCompletedEvent) -> Result<(), HookError> {
58        Ok(())
59    }
60
61    async fn action_received(&self, _event: ActionReceivedEvent) -> Result<(), HookError> {
62        Ok(())
63    }
64
65    /// Fatal hook: a failed transform leaves the next action payload undefined,
66    /// so the route fails rather than send something undefined to the env.
67    async fn transform_action(
68        &self,
69        event: ActionReceivedEvent,
70    ) -> Result<Option<Vec<Bytes>>, HookError> {
71        Ok(event.action)
72    }
73
74    async fn step_completed(&self, _event: StepCompletedEvent) -> Result<(), HookError> {
75        Ok(())
76    }
77
78    async fn observation_emitted(&self, _event: ObservationEmittedEvent) -> Result<(), HookError> {
79        Ok(())
80    }
81
82    /// Fatal hook; see [`Self::transform_action`].
83    async fn transform_observation(
84        &self,
85        event: ObservationEmittedEvent,
86    ) -> Result<Option<Vec<Bytes>>, HookError> {
87        Ok(event.observation)
88    }
89
90    async fn session_ended(&self, _event: SessionEndedEvent) -> Result<(), HookError> {
91        Ok(())
92    }
93
94    /// Live telemetry push, best-effort. A background ticker streams a `Window`
95    /// snapshot (the live tier, cleared each `RuntimeLimits::telemetry_window`)
96    /// while the session runs; one final cumulative `Session` snapshot is
97    /// delivered at session end (the durable tier, also returned on
98    /// `RuntimeReport.telemetry`). Branch on `event.snapshot.horizon` for window
99    /// vs session. Dispatched from a separate task, so this may run concurrently
100    /// with the other hooks — do not assume serialized delivery.
101    async fn on_telemetry(&self, _event: TelemetrySnapshotEvent) -> Result<(), HookError> {
102        Ok(())
103    }
104
105    async fn session_failed(&self, _event: SessionFailedEvent) -> Result<(), HookError> {
106        Ok(())
107    }
108
109    async fn log(&self, _event: LogEvent) -> Result<(), HookError> {
110        Ok(())
111    }
112
113    /// Best-effort; see [`RelayAdvisoryEvent`].
114    async fn relay_advisory(&self, _event: RelayAdvisoryEvent) -> Result<(), HookError> {
115        Ok(())
116    }
117}