Skip to main content

rlmesh_runtime/
driver.rs

1use std::future::Future;
2use std::sync::Arc;
3use std::time::{Duration, Instant};
4
5use async_trait::async_trait;
6use prost::Message;
7use rlmesh_proto::common::v1::MessageBytes;
8use rlmesh_proto::core::v1::OperationTelemetry;
9use rlmesh_proto::env::v1::{
10    EpisodeMetadata, ResetRequest, ResetResponse, StepRequest, StepResponse,
11};
12use rlmesh_proto::model::v1::{CloseRouteRequest, PredictRequest, PredictResponse};
13use rlmesh_proto::spaces::v1::SpaceValue;
14use tokio_util::sync::CancellationToken;
15
16use crate::episodes::EpisodeRecord;
17use crate::hooks::{
18    ActionReceivedEvent, EpisodeCompletedEvent, EpisodeStartedEvent, HookError, LogEvent, LogLevel,
19    NoopRuntimeHooks, ObservationEmittedEvent, RuntimeHooks, SessionEndedEvent, SessionFailedEvent,
20    SessionStartedEvent, StepCompletedEvent, TelemetrySummaryEvent, TelemetryWindowEvent,
21};
22use crate::route::requests::RequestPhase;
23use crate::spec::{RuntimeReport, RuntimeSessionSpec};
24use crate::state::{RouteSnapshot, RouteState, StartedEpisode};
25use crate::timing::{RuntimeTiming, StepTimingSample};
26
27#[derive(Debug, thiserror::Error)]
28pub enum RuntimeError {
29    #[error("invalid runtime session spec: {0}")]
30    InvalidSpec(String),
31
32    #[error(
33        "{operation} timed out on route {route_id} component {component_id} at runtime step {step} after {timeout:?}"
34    )]
35    OperationTimeout {
36        route_id: String,
37        component_id: String,
38        operation: &'static str,
39        step: i64,
40        timeout: Duration,
41    },
42
43    #[error("route {route_id} cancelled at runtime step {step}: {reason}")]
44    RouteCancelled {
45        route_id: String,
46        step: i64,
47        reason: String,
48    },
49
50    #[error(
51        "environment {operation} failed at runtime step {step}: {message}. If the source is 'transport error: connection closed', the environment server exited, crashed, or received SIGTERM before replying; inspect the environment container logs immediately before the runtime error timestamp"
52    )]
53    EnvRpc {
54        operation: &'static str,
55        step: i64,
56        message: String,
57    },
58
59    #[error("model endpoint {component_id} request failed: {message}")]
60    ModelRpc {
61        component_id: String,
62        message: String,
63    },
64
65    #[error(
66        "model endpoint {component_id} returned mismatched route identity for request {request_id}"
67    )]
68    ModelRouteMismatch {
69        component_id: String,
70        request_id: String,
71    },
72
73    #[error("protocol error: {0}")]
74    Protocol(String),
75
76    #[error("runtime hook failed: {0}")]
77    Hook(HookError),
78}
79
80impl RuntimeError {
81    pub fn operation_timeout(
82        route_id: impl Into<String>,
83        component_id: impl Into<String>,
84        operation: &'static str,
85        step: i64,
86        timeout: Duration,
87    ) -> Self {
88        Self::OperationTimeout {
89            route_id: route_id.into(),
90            component_id: component_id.into(),
91            operation,
92            step,
93            timeout,
94        }
95    }
96
97    pub fn route_cancelled(
98        route_id: impl Into<String>,
99        step: i64,
100        reason: impl Into<String>,
101    ) -> Self {
102        Self::RouteCancelled {
103            route_id: route_id.into(),
104            step,
105            reason: reason.into(),
106        }
107    }
108}
109
110pub struct RuntimeEnvReset {
111    pub response: ResetResponse,
112    pub telemetry: Option<OperationTelemetry>,
113}
114
115pub struct RuntimeEnvStep {
116    pub response: StepResponse,
117    pub telemetry: Option<OperationTelemetry>,
118}
119
120pub struct RuntimeModelPrediction {
121    pub response: PredictResponse,
122    pub telemetry: Option<OperationTelemetry>,
123}
124
125#[async_trait]
126pub trait RuntimeEnv: Send {
127    async fn reset(&mut self, request: ResetRequest) -> Result<RuntimeEnvReset, RuntimeError>;
128
129    async fn step(&mut self, request: StepRequest) -> Result<RuntimeEnvStep, RuntimeError>;
130
131    async fn close(&mut self, _timeout: Duration) -> Result<(), String> {
132        Ok(())
133    }
134}
135
136#[async_trait]
137pub trait RuntimeModel: Send {
138    async fn predict(
139        &mut self,
140        request: PredictRequest,
141    ) -> Result<RuntimeModelPrediction, RuntimeError>;
142
143    async fn close_route(
144        &mut self,
145        _request: CloseRouteRequest,
146        _timeout: Duration,
147    ) -> Result<(), String> {
148        Ok(())
149    }
150}
151
152pub struct RuntimeDriver<E, M> {
153    spec: RuntimeSessionSpec,
154    env: E,
155    model: M,
156    hooks: Arc<dyn RuntimeHooks>,
157}
158
159impl<E, M> RuntimeDriver<E, M>
160where
161    E: RuntimeEnv,
162    M: RuntimeModel,
163{
164    pub fn new(spec: RuntimeSessionSpec, env: E, model: M, hooks: Arc<dyn RuntimeHooks>) -> Self {
165        Self {
166            spec,
167            env,
168            model,
169            hooks,
170        }
171    }
172
173    pub fn without_hooks(spec: RuntimeSessionSpec, env: E, model: M) -> Self {
174        Self::new(spec, env, model, Arc::new(NoopRuntimeHooks))
175    }
176
177    pub async fn run(self) -> Result<RuntimeReport, RuntimeError> {
178        self.run_with_cancellation(CancellationToken::new()).await
179    }
180
181    pub async fn run_with_cancellation(
182        mut self,
183        cancellation: CancellationToken,
184    ) -> Result<RuntimeReport, RuntimeError> {
185        self.spec.validate().map_err(RuntimeError::InvalidSpec)?;
186        let mut state = RouteState::new(&self.spec);
187        let result = self.run_loop(&mut state, &cancellation).await;
188        if let Err(error) = &result {
189            self.shutdown_after_failure(&mut state, error).await;
190        }
191        result
192    }
193
194    async fn run_loop(
195        &mut self,
196        state: &mut RouteState,
197        cancellation: &CancellationToken,
198    ) -> Result<RuntimeReport, RuntimeError> {
199        let mut timings = RuntimeTiming::default();
200        self.invoke_session_started(state, &self.spec.env_id).await;
201
202        let reset_started = Instant::now();
203        let reset_timeout = self.spec.limits.env_reset_timeout;
204        let reset_timeout_ms = self.spec.limits.env_reset_timeout_ms();
205        let reset_ok = await_runtime_operation(
206            cancellation,
207            reset_timeout,
208            RuntimeError::operation_timeout(
209                state.route_id(),
210                state.env_component_id(),
211                "env.reset",
212                0,
213                reset_timeout,
214            ),
215            self.cancelled_error(state, 0),
216            self.env.reset(ResetRequest {
217                seeds: vec![],
218                options: None,
219                timeout_ms: reset_timeout_ms,
220            }),
221        )
222        .await?;
223        let reset_latency = reset_started.elapsed();
224        timings.reset.record(reset_latency);
225        timings
226            .window
227            .record_operation_telemetry(state.env_component_id(), reset_ok.telemetry.as_ref());
228        self.invoke_log(
229            state,
230            LogLevel::Info,
231            format!(
232                "env reset complete in {:.0}ms ({} episode(s) ready)",
233                reset_latency.as_secs_f64() * 1000.0,
234                reset_ok
235                    .response
236                    .episode_ids
237                    .iter()
238                    .filter(|value| !value.is_empty())
239                    .count()
240            ),
241        )
242        .await;
243
244        let reset_observation = value_bytes(reset_ok.response.observation.as_ref())?;
245        let started_episodes = state.start_episodes(reset_ok.response.episode_ids, false);
246        self.invoke_started_episodes(state, started_episodes).await;
247
248        let mut reset_msg =
249            state.predict_request(reset_observation.clone(), RequestPhase::ResetObservation);
250        let mut reset_event =
251            self.observation_event(state, state.snapshot(), true, reset_observation.clone());
252        let transformed_reset_observation = self
253            .invoke_transform_observation(reset_event.clone())
254            .await?;
255        reset_event.observation = transformed_reset_observation.clone();
256        reset_msg.observation = transformed_reset_observation.map(bytes_value);
257        self.invoke_observation_emitted(reset_event).await;
258
259        let mut pending_observation_msg = reset_msg;
260
261        loop {
262            if cancellation.is_cancelled() {
263                return Err(self.cancelled_error(state, state.snapshot().step));
264            }
265
266            let predict_snapshot = state.snapshot();
267            let model_wait_started = Instant::now();
268            let predict_timeout = self.spec.limits.model_predict_timeout;
269            let expected_context = pending_observation_msg.context.clone();
270            let action_msg = await_runtime_operation(
271                cancellation,
272                predict_timeout,
273                RuntimeError::operation_timeout(
274                    state.route_id(),
275                    state.model_component_id(),
276                    "model.predict",
277                    predict_snapshot.step,
278                    predict_timeout,
279                ),
280                self.cancelled_error(state, predict_snapshot.step),
281                self.model.predict(pending_observation_msg),
282            )
283            .await?;
284            if action_msg.response.context != expected_context {
285                let request_id = expected_context
286                    .as_ref()
287                    .map(|context| context.request_id.clone())
288                    .unwrap_or_default();
289                return Err(RuntimeError::ModelRouteMismatch {
290                    component_id: state.model_component_id().to_string(),
291                    request_id,
292                });
293            }
294            let model_action = value_bytes(action_msg.response.action.as_ref())?;
295            let model_wait_latency = model_wait_started.elapsed();
296            timings.model_wait.record(model_wait_latency);
297            timings.window.record_operation_telemetry(
298                state.model_component_id(),
299                action_msg.telemetry.as_ref(),
300            );
301
302            let action_step = predict_snapshot.step + 1;
303            let mut action_event = ActionReceivedEvent {
304                session_id: state.session_id().to_string(),
305                route: state.route_context(),
306                episode_id: predict_snapshot.episode_id.clone(),
307                episode_record_id: predict_snapshot.episode_record_id.clone(),
308                episode_ids: predict_snapshot.episode_ids.clone(),
309                episode_record_ids: predict_snapshot.episode_record_ids.clone(),
310                step: action_step,
311                env_index: predict_snapshot.env_index,
312                action_space: self.spec.action_space().clone(),
313                action: model_action,
314            };
315            action_event.action = self.invoke_transform_action(action_event.clone()).await?;
316            let request_bytes = action_event
317                .action
318                .as_ref()
319                .map(|action| action.data.len())
320                .unwrap_or_default();
321            self.invoke_action_received(action_event.clone()).await;
322
323            let env_step_started = Instant::now();
324            let step_timeout = self.spec.limits.env_step_timeout;
325            let step_timeout_ms = self.spec.limits.env_step_timeout_ms();
326            let step_ok = await_runtime_operation(
327                cancellation,
328                step_timeout,
329                RuntimeError::operation_timeout(
330                    state.route_id(),
331                    state.env_component_id(),
332                    "env.step",
333                    action_step,
334                    step_timeout,
335                ),
336                self.cancelled_error(state, action_step),
337                self.env.step(StepRequest {
338                    action: action_event.action.map(bytes_value),
339                    timeout_ms: step_timeout_ms,
340                }),
341            )
342            .await?;
343            let env_step_latency = env_step_started.elapsed();
344            let step_observation = value_bytes(step_ok.response.observation.as_ref())?;
345            timings.env_step.record(env_step_latency);
346            timings
347                .window
348                .record_operation_telemetry(state.env_component_id(), step_ok.telemetry.as_ref());
349            let response_bytes = step_observation
350                .as_ref()
351                .map(|obs| obs.data.len())
352                .unwrap_or_default()
353                + step_ok
354                    .response
355                    .infos
356                    .as_ref()
357                    .map(Message::encoded_len)
358                    .unwrap_or_default();
359            timings.window.record_step(StepTimingSample {
360                model_wait: model_wait_latency,
361                env_step: env_step_latency,
362                request_bytes,
363                response_bytes,
364                env_component_id: state.env_component_id(),
365                model_component_id: state.model_component_id(),
366            });
367
368            state.record_step();
369            let step_snapshot = state.snapshot();
370            self.invoke_step_completed(StepCompletedEvent {
371                session_id: state.session_id().to_string(),
372                route: state.route_context(),
373                episode_id: step_snapshot.episode_id.clone(),
374                episode_record_id: step_snapshot.episode_record_id.clone(),
375                step: step_snapshot.step,
376                env_index: step_snapshot.env_index,
377                rewards: step_ok.response.rewards.clone(),
378            })
379            .await;
380
381            if let Some(event) = timings.maybe_emit_window(
382                state.session_id(),
383                state.route_context(),
384                self.spec.limits.telemetry_window,
385            ) {
386                self.invoke_telemetry_window(event).await;
387            }
388
389            let observation_snapshot = if step_ok.response.episode_ids.is_empty() {
390                state.snapshot()
391            } else {
392                let started_episodes = state.observe_episode_ids(step_ok.response.episode_ids);
393                self.invoke_started_episodes(state, started_episodes).await;
394                state.snapshot()
395            };
396            self.invoke_observation_emitted(self.observation_event(
397                state,
398                observation_snapshot,
399                false,
400                step_observation.clone(),
401            ))
402            .await;
403
404            self.emit_completed_episodes(state, &step_ok.response.completed_episodes)
405                .await;
406
407            if self
408                .spec
409                .max_episodes
410                .is_some_and(|limit| state.total_episodes() >= limit as i64)
411            {
412                if let Some(event) = timings.flush_window(state.session_id(), state.route_context())
413                {
414                    self.invoke_telemetry_window(event).await;
415                }
416                if let Some(event) =
417                    timings.telemetry_summary(state.session_id(), state.route_context())
418                {
419                    self.invoke_telemetry_summary(event).await;
420                }
421                let close_request = state.close_route_request("completed requested episodes");
422                self.shutdown_terminal_route(state, "completed requested episodes", close_request)
423                    .await;
424                self.invoke_session_ended(
425                    state,
426                    "completed requested episodes",
427                    state.total_steps(),
428                    state.total_episodes(),
429                )
430                .await;
431                timings.log_summary(state.total_steps(), state.total_episodes());
432                return Ok(RuntimeReport {
433                    session_id: state.session_id().to_string(),
434                    route_id: self.spec.route_id.clone(),
435                    total_steps: state.total_steps(),
436                    total_episodes: state.total_episodes(),
437                });
438            }
439
440            let need_reset = !step_ok.response.completed_episodes.is_empty();
441            let (next_obs, phase, is_reset_msg) = if need_reset {
442                let reset_started = Instant::now();
443                let step = state.snapshot().step;
444                let reset_timeout = self.spec.limits.env_reset_timeout;
445                let reset_timeout_ms = self.spec.limits.env_reset_timeout_ms();
446                let reset_ok = await_runtime_operation(
447                    cancellation,
448                    reset_timeout,
449                    RuntimeError::operation_timeout(
450                        state.route_id(),
451                        state.env_component_id(),
452                        "env.reset",
453                        step,
454                        reset_timeout,
455                    ),
456                    self.cancelled_error(state, step),
457                    self.env.reset(ResetRequest {
458                        seeds: vec![],
459                        options: None,
460                        timeout_ms: reset_timeout_ms,
461                    }),
462                )
463                .await?;
464                timings.reset.record(reset_started.elapsed());
465                timings.window.record_operation_telemetry(
466                    state.env_component_id(),
467                    reset_ok.telemetry.as_ref(),
468                );
469                let next_obs = value_bytes(reset_ok.response.observation.as_ref())?;
470                let started_episodes = state.start_episodes(reset_ok.response.episode_ids, true);
471                self.invoke_started_episodes(state, started_episodes).await;
472                (next_obs, RequestPhase::ResetObservation, true)
473            } else {
474                (
475                    step_observation.clone(),
476                    RequestPhase::StepObservation,
477                    false,
478                )
479            };
480
481            let mut obs_msg = state.predict_request(next_obs.clone(), phase);
482            let mut outgoing_observation_event =
483                self.observation_event(state, state.snapshot(), is_reset_msg, next_obs);
484            let transformed_observation = self
485                .invoke_transform_observation(outgoing_observation_event.clone())
486                .await?;
487            outgoing_observation_event.observation = transformed_observation.clone();
488            obs_msg.observation = transformed_observation.map(bytes_value);
489            if is_reset_msg {
490                self.invoke_observation_emitted(outgoing_observation_event)
491                    .await;
492            }
493
494            pending_observation_msg = obs_msg;
495        }
496    }
497
498    async fn shutdown_after_failure(&mut self, state: &mut RouteState, error: &RuntimeError) {
499        let reason = error.to_string();
500        let request = state.close_route_request(reason.clone());
501        self.shutdown_terminal_route(state, &reason, request).await;
502
503        if let Err(err) = self
504            .hooks
505            .session_failed(SessionFailedEvent {
506                session_id: state.session_id().to_string(),
507                route: state.route_context(),
508                reason,
509            })
510            .await
511        {
512            tracing::warn!("runtime hook session_failed failed: {err}");
513        }
514    }
515
516    async fn shutdown_terminal_route(
517        &mut self,
518        state: &RouteState,
519        reason: &str,
520        request: CloseRouteRequest,
521    ) {
522        let timeout = self.spec.limits.service_close_timeout;
523        let model_close = self.model.close_route(request, timeout);
524        if self.spec.close_env_on_end {
525            let env_close = self.env.close(timeout);
526            let (env_result, model_result) = tokio::join!(env_close, model_close);
527            if let Err(err) = env_result {
528                tracing::warn!(error = %err, "environment close failed during route shutdown");
529            }
530            if let Err(err) = model_result {
531                tracing::warn!(
532                    error = %err,
533                    reason,
534                    "model route close failed during route shutdown; relying on owner shutdown"
535                );
536            }
537            return;
538        }
539
540        tracing::debug!(
541            route_id = %state.route_id(),
542            reason,
543            "skipping environment close for route; endpoint remains owned by the run"
544        );
545        if let Err(err) = model_close.await {
546            tracing::warn!(
547                error = %err,
548                reason,
549                "model route close failed during route shutdown; relying on owner shutdown"
550            );
551        }
552    }
553
554    fn cancelled_error(&self, state: &RouteState, step: i64) -> RuntimeError {
555        RuntimeError::route_cancelled(
556            state.route_id(),
557            step,
558            "cancelled after sibling route failure",
559        )
560    }
561
562    async fn invoke_session_started(&self, state: &RouteState, env_id: &str) {
563        if let Err(err) = self
564            .hooks
565            .session_started(SessionStartedEvent {
566                session_id: state.session_id().to_string(),
567                route: state.route_context(),
568                env_id: env_id.to_string(),
569            })
570            .await
571        {
572            tracing::warn!("runtime hook session_started failed: {err}");
573        }
574    }
575
576    async fn invoke_started_episodes(&self, state: &RouteState, episodes: Vec<StartedEpisode>) {
577        for episode in episodes {
578            self.invoke_episode_started(state, &episode.episode_id, &episode.record)
579                .await;
580        }
581    }
582
583    async fn invoke_episode_started(
584        &self,
585        state: &RouteState,
586        episode_id: &str,
587        record: &EpisodeRecord,
588    ) {
589        if let Err(err) = self
590            .hooks
591            .episode_started(EpisodeStartedEvent {
592                session_id: state.session_id().to_string(),
593                route: state.route_context(),
594                episode_id: episode_id.to_string(),
595                episode_record_id: record.record_id.clone(),
596                episode_index: record.index,
597                env_index: record.env_index,
598                started_from_auto_reset: record.started_from_auto_reset,
599            })
600            .await
601        {
602            tracing::warn!("runtime hook episode_started failed: {err}");
603        }
604    }
605
606    async fn emit_completed_episodes(&self, state: &mut RouteState, episodes: &[EpisodeMetadata]) {
607        for completed in episodes {
608            let record = state.complete_episode(&completed.episode_id);
609            self.invoke_episode_completed(EpisodeCompletedEvent {
610                session_id: state.session_id().to_string(),
611                route: state.route_context(),
612                episode_id: completed.episode_id.clone(),
613                episode_record_id: record
614                    .as_ref()
615                    .map(|record| record.record_id.clone())
616                    .unwrap_or_default(),
617                episode_index: record.as_ref().map_or(0, |record| record.index),
618                env_index: completed.env_index,
619                step_count: completed.step_count,
620                cumulative_reward: completed.cumulative_reward,
621                terminated: completed.terminated,
622                truncated: completed.truncated,
623                duration_ms: completed.duration_ms,
624                final_info: completed.final_info.clone(),
625            })
626            .await;
627        }
628    }
629
630    async fn invoke_episode_completed(&self, event: EpisodeCompletedEvent) {
631        if let Err(err) = self.hooks.episode_completed(event).await {
632            tracing::warn!("runtime hook episode_completed failed: {err}");
633        }
634    }
635
636    async fn invoke_action_received(&self, event: ActionReceivedEvent) {
637        if let Err(err) = self.hooks.action_received(event).await {
638            tracing::warn!("runtime hook action_received failed: {err}");
639        }
640    }
641
642    async fn invoke_transform_action(
643        &self,
644        event: ActionReceivedEvent,
645    ) -> Result<Option<MessageBytes>, RuntimeError> {
646        match self.hooks.transform_action(event).await {
647            Ok(action) => Ok(action),
648            Err(err) => {
649                tracing::warn!("runtime hook transform_action failed: {err}");
650                Err(RuntimeError::Hook(err))
651            }
652        }
653    }
654
655    async fn invoke_step_completed(&self, event: StepCompletedEvent) {
656        if let Err(err) = self.hooks.step_completed(event).await {
657            tracing::warn!("runtime hook step_completed failed: {err}");
658        }
659    }
660
661    async fn invoke_observation_emitted(&self, event: ObservationEmittedEvent) {
662        if let Err(err) = self.hooks.observation_emitted(event).await {
663            tracing::warn!("runtime hook observation_emitted failed: {err}");
664        }
665    }
666
667    async fn invoke_transform_observation(
668        &self,
669        event: ObservationEmittedEvent,
670    ) -> Result<Option<MessageBytes>, RuntimeError> {
671        match self.hooks.transform_observation(event).await {
672            Ok(observation) => Ok(observation),
673            Err(err) => {
674                tracing::warn!("runtime hook transform_observation failed: {err}");
675                Err(RuntimeError::Hook(err))
676            }
677        }
678    }
679
680    async fn invoke_telemetry_window(&self, event: TelemetryWindowEvent) {
681        if let Err(err) = self.hooks.telemetry_window(event).await {
682            tracing::warn!("runtime hook telemetry_window failed: {err}");
683        }
684    }
685
686    async fn invoke_telemetry_summary(&self, event: TelemetrySummaryEvent) {
687        if let Err(err) = self.hooks.telemetry_summary(event).await {
688            tracing::warn!("runtime hook telemetry_summary failed: {err}");
689        }
690    }
691
692    async fn invoke_log(&self, state: &RouteState, level: LogLevel, message: impl Into<String>) {
693        if let Err(err) = self
694            .hooks
695            .log(LogEvent {
696                session_id: state.session_id().to_string(),
697                route: state.route_context(),
698                level,
699                message: message.into(),
700                source: Some("runtime".to_string()),
701            })
702            .await
703        {
704            tracing::warn!("runtime hook log failed: {err}");
705        }
706    }
707
708    async fn invoke_session_ended(
709        &self,
710        state: &RouteState,
711        reason: &str,
712        total_steps: i64,
713        total_episodes: i64,
714    ) {
715        if let Err(err) = self
716            .hooks
717            .session_ended(SessionEndedEvent {
718                session_id: state.session_id().to_string(),
719                route: state.route_context(),
720                reason: reason.to_string(),
721                total_steps,
722                total_episodes,
723            })
724            .await
725        {
726            tracing::warn!("runtime hook session_ended failed: {err}");
727        }
728    }
729
730    fn observation_event(
731        &self,
732        state: &RouteState,
733        snapshot: RouteSnapshot,
734        is_reset: bool,
735        observation: Option<MessageBytes>,
736    ) -> ObservationEmittedEvent {
737        ObservationEmittedEvent {
738            session_id: state.session_id().to_string(),
739            route: state.route_context(),
740            episode_id: snapshot.episode_id,
741            episode_record_id: snapshot.episode_record_id,
742            episode_ids: snapshot.episode_ids,
743            episode_record_ids: snapshot.episode_record_ids,
744            step: snapshot.step,
745            env_index: snapshot.env_index,
746            is_reset,
747            num_envs: self.spec.num_envs as u32,
748            observation_space: self.spec.observation_space().clone(),
749            observation,
750        }
751    }
752}
753
754async fn await_runtime_operation<T, F>(
755    cancellation: &CancellationToken,
756    timeout: Duration,
757    timeout_error: RuntimeError,
758    cancelled_error: RuntimeError,
759    operation: F,
760) -> Result<T, RuntimeError>
761where
762    F: Future<Output = Result<T, RuntimeError>>,
763{
764    tokio::select! {
765        _ = cancellation.cancelled() => Err(cancelled_error),
766        result = tokio::time::timeout(timeout, operation) => match result {
767            Ok(result) => result,
768            Err(_) => Err(timeout_error),
769        },
770    }
771}
772
773fn bytes_value(value: MessageBytes) -> SpaceValue {
774    SpaceValue { bytes: Some(value) }
775}
776
777fn value_bytes(payload: Option<&SpaceValue>) -> Result<Option<MessageBytes>, RuntimeError> {
778    let Some(payload) = payload else {
779        return Ok(None);
780    };
781    Ok(payload.bytes.clone())
782}