Skip to main content

rlmesh_runtime/hooks/
relay.rs

1//! The relay seam: what the driver does with a payload bound for a served peer.
2//!
3//! The runtime is the interpreter between two peers that never talk to each
4//! other, so every payload it relays is checked against the *target* leg's
5//! [`PeerCeiling`]. A [`RelayPolicy`] decides whether it goes out as-is, goes
6//! out converted with an advisory, or is refused. The OSS default,
7//! [`RefusingRelayPolicy`], never converts.
8
9use std::sync::Arc;
10
11use prost::bytes::Bytes;
12use rlmesh_proto::spaces::v1::{SpaceSpec, space_spec::Spec};
13use rlmesh_spaces::{Advisory, DType};
14
15use crate::spec::PeerCeiling;
16
17/// Which way a relayed payload travels.
18#[derive(Debug, Clone, Copy, PartialEq, Eq)]
19pub enum Leg {
20    /// The env contract or an observation, on its way to the model.
21    EnvToModel,
22    /// An action, on its way to the env.
23    ModelToEnv,
24}
25
26impl Leg {
27    /// The peer this leg delivers to.
28    pub fn target(self) -> &'static str {
29        match self {
30            Leg::EnvToModel => "model",
31            Leg::ModelToEnv => "env",
32        }
33    }
34}
35
36/// One relayed payload, for a [`RelayPolicy`] to check against the target
37/// leg's [`PeerCeiling`]. Capability gating is not checked here: it happens
38/// where a feature is emitted.
39#[derive(Debug, Clone, PartialEq)]
40pub struct PayloadFacts {
41    /// The space typing `leaves`, in leaf order. For the env contract, a
42    /// two-element Tuple of the observation space then the action space.
43    pub space: Arc<SpaceSpec>,
44    /// Summed leaf bytes, not the framed wire size; the encoded contract's
45    /// length for the env contract.
46    pub byte_len: usize,
47    /// The payload's leaves, for a policy that converts them; empty for the
48    /// env contract, which is checked before any leaf flows.
49    pub leaves: Vec<Bytes>,
50}
51
52impl PayloadFacts {
53    /// The facts of `leaves` typed by `space`.
54    pub fn new(space: Arc<SpaceSpec>, leaves: Vec<Bytes>) -> Self {
55        Self {
56            space,
57            byte_len: leaves.iter().map(Bytes::len).sum(),
58            leaves,
59        }
60    }
61
62    /// Why `ceiling` cannot carry this payload, or `None` when it can.
63    pub fn exceeds(&self, ceiling: &PeerCeiling) -> Option<String> {
64        if let Some(dtype) = dtype_outside(&self.space, ceiling) {
65            return Some(format!("dtype {} is outside its ceiling", dtype.name()));
66        }
67        if self.byte_len > ceiling.max_message_size {
68            return Some(format!(
69                "{} bytes exceeds its {}-byte message cap",
70                self.byte_len, ceiling.max_message_size
71            ));
72        }
73        None
74    }
75}
76
77/// The first leaf dtype of `space` that `ceiling` cannot decode.
78fn dtype_outside(space: &SpaceSpec, ceiling: &PeerCeiling) -> Option<DType> {
79    match &space.spec {
80        Some(Spec::Dict(dict)) => dict
81            .spaces
82            .iter()
83            .find_map(|space| dtype_outside(space, ceiling)),
84        Some(Spec::Tuple(tuple)) => tuple
85            .spaces
86            .iter()
87            .find_map(|space| dtype_outside(space, ceiling)),
88        // A value no DType names never decoded into this spec: the contract
89        // decode at the handshake refuses it first.
90        _ => DType::try_from(space.dtype)
91            .ok()
92            .filter(|dtype| !ceiling.dtypes.contains(dtype)),
93    }
94}
95
96/// A policy's verdict on one relayed payload.
97#[derive(Debug, Clone, PartialEq)]
98pub enum RelayDecision {
99    /// Relay the payload unchanged.
100    Forward,
101    /// Relay `leaves` in its place, recording `advisory` on the report.
102    /// `space` types the converted leaves when their layout changed; the
103    /// emitted hook event then carries it in place of the original leg's space.
104    ///
105    /// On the env contract (empty `leaves`) only `advisory` takes effect.
106    Convert {
107        leaves: Vec<Bytes>,
108        space: Option<Arc<SpaceSpec>>,
109        advisory: Advisory,
110    },
111    /// Fail the route with this reason.
112    Refuse(String),
113}
114
115/// Decides what reaches a served peer that may not decode a payload as sent.
116///
117/// Consulted for the env contract at session start and for every observation
118/// and action after the transform hooks, only on a leg whose ceiling the spec
119/// carries (an in-process leg has none). Install one with
120/// [`RuntimeDriver::with_relay_policy`](crate::RuntimeDriver::with_relay_policy);
121/// the default is [`RefusingRelayPolicy`].
122///
123/// The contract call is advisory-only: the driver records a
124/// [`RelayDecision::Convert`]'s advisory but relays nothing, so the policy
125/// owner must down-convert the contract it sends the model at resolve itself.
126pub trait RelayPolicy: Send + Sync {
127    fn reconcile(&self, leg: Leg, ceiling: &PeerCeiling, payload: &PayloadFacts) -> RelayDecision;
128}
129
130/// The OSS relay policy: forward what the target can decode, refuse the rest.
131#[derive(Debug, Default)]
132pub struct RefusingRelayPolicy;
133
134impl RelayPolicy for RefusingRelayPolicy {
135    fn reconcile(&self, _leg: Leg, ceiling: &PeerCeiling, payload: &PayloadFacts) -> RelayDecision {
136        payload
137            .exceeds(ceiling)
138            .map_or(RelayDecision::Forward, RelayDecision::Refuse)
139    }
140}