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}