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}