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, SessionEndedEvent, SessionFailedEvent,
9    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    async fn episode_completed(&self, _event: EpisodeCompletedEvent) -> Result<(), HookError> {
55        Ok(())
56    }
57
58    async fn action_received(&self, _event: ActionReceivedEvent) -> Result<(), HookError> {
59        Ok(())
60    }
61
62    /// Fatal hook: a failed transform leaves the next action payload undefined,
63    /// so the route fails rather than send something undefined to the env.
64    async fn transform_action(
65        &self,
66        event: ActionReceivedEvent,
67    ) -> Result<Option<Vec<Bytes>>, HookError> {
68        Ok(event.action)
69    }
70
71    async fn step_completed(&self, _event: StepCompletedEvent) -> Result<(), HookError> {
72        Ok(())
73    }
74
75    async fn observation_emitted(&self, _event: ObservationEmittedEvent) -> Result<(), HookError> {
76        Ok(())
77    }
78
79    /// Fatal hook; see [`Self::transform_action`].
80    async fn transform_observation(
81        &self,
82        event: ObservationEmittedEvent,
83    ) -> Result<Option<Vec<Bytes>>, HookError> {
84        Ok(event.observation)
85    }
86
87    async fn session_ended(&self, _event: SessionEndedEvent) -> Result<(), HookError> {
88        Ok(())
89    }
90
91    /// Live telemetry push, best-effort. A background ticker streams a `Window`
92    /// snapshot (the live tier, cleared each `RuntimeLimits::telemetry_window`)
93    /// while the session runs; one final cumulative `Session` snapshot is
94    /// delivered at session end (the durable tier, also returned on
95    /// `RuntimeReport.telemetry`). Branch on `event.snapshot.horizon` for window
96    /// vs session. Dispatched from a separate task, so this may run concurrently
97    /// with the other hooks — do not assume serialized delivery.
98    async fn on_telemetry(&self, _event: TelemetrySnapshotEvent) -> Result<(), HookError> {
99        Ok(())
100    }
101
102    async fn session_failed(&self, _event: SessionFailedEvent) -> Result<(), HookError> {
103        Ok(())
104    }
105
106    async fn log(&self, _event: LogEvent) -> Result<(), HookError> {
107        Ok(())
108    }
109}