#![deny(rustdoc::broken_intra_doc_links)]
use std::collections::HashMap;
use std::sync::RwLock;
pub const PROTOCOL_GENERATION: &str = "rlmesh-wire-v1";
pub const CURRENT_WORKFLOW_EDITION: &str = env!("RLMESH_CURRENT_WORKFLOW_EDITION");
pub const WORKFLOW_EDITION_BASE: &str = env!("RLMESH_WORKFLOW_EDITION_BASE");
pub const BUILD_COHORT: &str = env!("RLMESH_BUILD_COHORT");
pub const BUILD_SOURCE: &str = env!("RLMESH_BUILD_SOURCE");
pub const SUPPORTED_WORKFLOW_EDITIONS: &[&str] = &SUPPORTED_WORKFLOW_EDITION_ARRAY;
const SUPPORTED_WORKFLOW_EDITION_LIST: &str = env!("RLMESH_SUPPORTED_WORKFLOW_EDITIONS");
const SUPPORTED_WORKFLOW_EDITION_ARRAY: [&str; edition_count(SUPPORTED_WORKFLOW_EDITION_LIST)] =
split_editions(SUPPORTED_WORKFLOW_EDITION_LIST);
const fn edition_count(list: &str) -> usize {
let bytes = list.as_bytes();
let mut count = 1;
let mut index = 0;
while index < bytes.len() {
if bytes[index] == b',' {
count += 1;
}
index += 1;
}
count
}
const fn split_editions<const N: usize>(list: &'static str) -> [&'static str; N] {
let bytes = list.as_bytes();
let mut editions = [""; N];
let mut start = 0;
let mut index = 0;
let mut slot = 0;
while index < bytes.len() {
if bytes[index] == b',' {
editions[slot] = edition_at(bytes, start, index);
slot += 1;
start = index + 1;
}
index += 1;
}
editions[slot] = edition_at(bytes, start, bytes.len());
editions
}
const fn edition_at(bytes: &'static [u8], start: usize, end: usize) -> &'static str {
let (_, tail) = bytes.split_at(start);
let (edition, _) = tail.split_at(end - start);
match std::str::from_utf8(edition) {
Ok(edition) => edition,
Err(_) => panic!("RLMESH_SUPPORTED_WORKFLOW_EDITIONS is not UTF-8"),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum Edition {
E2026_06,
}
impl Edition {
const ALL: &'static [Edition] = &[Edition::E2026_06];
pub const fn base(&self) -> &'static str {
match self {
Edition::E2026_06 => "2026.06",
}
}
pub fn current() -> Edition {
Edition::parse(CURRENT_WORKFLOW_EDITION)
.expect("CURRENT_WORKFLOW_EDITION names an arm (every_arm_is_listed_and_has_a_row)")
}
pub fn parse(edition: &str) -> Result<Edition, UnknownEdition> {
let name = edition.trim();
let (base, _, _) = edition_sort_key(name);
Edition::ALL
.iter()
.copied()
.find(|edition| edition.base() == base)
.ok_or_else(|| UnknownEdition(name.to_string()))
}
}
impl std::fmt::Display for Edition {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.base())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UnknownEdition(String);
impl UnknownEdition {
pub fn name(&self) -> &str {
&self.0
}
}
impl std::fmt::Display for UnknownEdition {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "unknown workflow edition {:?}", self.0)
}
}
impl std::error::Error for UnknownEdition {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct EditionDefaults {
pub default_max_episode_steps: i64,
pub trial_index_option_key: &'static str,
pub driver_owned_reset_modes: &'static [core::v1::AutoresetMode],
pub success_info_keys: &'static [&'static str],
pub conformance_warning_info_key: &'static str,
}
static DEFAULTS_2026_06: EditionDefaults = EditionDefaults {
default_max_episode_steps: 100_000,
trial_index_option_key: "trial_index",
driver_owned_reset_modes: &[
core::v1::AutoresetMode::Disabled,
core::v1::AutoresetMode::Unspecified,
],
success_info_keys: &["is_success", "success", "task_success"],
conformance_warning_info_key: "rlmesh.conformance.warning",
};
pub const fn defaults(edition: Edition) -> &'static EditionDefaults {
match edition {
Edition::E2026_06 => &DEFAULTS_2026_06,
}
}
pub fn is_retained_edition(edition: Edition) -> bool {
SUPPORTED_WORKFLOW_EDITIONS
.iter()
.any(|retained| Edition::parse(retained) == Ok(edition))
}
pub fn parse_retained_edition(edition: &str) -> Result<Edition, String> {
Edition::parse(edition)
.ok()
.filter(|edition| is_retained_edition(*edition))
.ok_or_else(|| {
format!(
"runtime cannot drive workflow edition {:?}; this build implements {:?}",
edition.trim(),
SUPPORTED_WORKFLOW_EDITIONS
)
})
}
pub fn parse_declared_edition(edition: &str) -> Result<Edition, String> {
let parsed = parse_retained_edition(edition)?;
let want = edition.trim();
if SUPPORTED_WORKFLOW_EDITIONS
.iter()
.any(|can| want_admits(want, can))
{
return Ok(parsed);
}
Err(format!(
"workflow edition {want:?} admits none of the editions this build offers \
({SUPPORTED_WORKFLOW_EDITIONS:?}), so declaring it would refuse every session; \
declare one of those, or the bare {:?} base, instead",
parsed.base()
))
}
pub mod capabilities {
pub const MODEL_CONCURRENT_PREDICT_V1: &str = "rlmesh.model.concurrent_predict.v1";
pub const MODEL_OBSERVATION_HISTORY_V1: &str = "rlmesh.model.observation_history.v1";
pub const ENV_SUBSET_STEP: &str = "subset_step";
}
pub fn is_protocol_generation_supported(generation: &str) -> bool {
generation.trim() == PROTOCOL_GENERATION
}
pub fn edition_sort_key(edition: &str) -> (&str, bool, &str) {
match edition.split_once('-') {
Some((base, suffix)) => (base, true, suffix),
None => (edition, false, ""),
}
}
pub fn want_admits(want: &str, edition: &str) -> bool {
match edition_sort_key(want) {
(base, false, _) => edition_sort_key(edition).0 <= base,
want => edition_sort_key(edition) <= want,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SessionOffer {
pub editions: Vec<String>,
pub preferred: Option<String>,
}
impl SessionOffer {
pub fn new(editions: &[&str]) -> Self {
Self {
editions: editions.iter().map(|e| e.trim().to_string()).collect(),
preferred: None,
}
}
pub fn this_build(declared: Option<&str>) -> Self {
Self {
editions: supported_workflow_editions(),
preferred: Some(declared_workflow_edition(declared).to_string()),
}
}
fn can(&self) -> impl Iterator<Item = &str> {
self.editions
.iter()
.map(|edition| edition.trim())
.filter(|edition| !edition.is_empty())
}
fn want(&self) -> Option<&str> {
self.preferred
.as_deref()
.map(str::trim)
.filter(|want| !want.is_empty())
}
}
pub fn declared_workflow_edition(declared: Option<&str>) -> &str {
declared
.map(str::trim)
.filter(|declared| !declared.is_empty())
.unwrap_or(CURRENT_WORKFLOW_EDITION)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TierOffer {
pub tier: String,
pub can: Vec<String>,
pub want: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct EditionRefusal {
pub tiers: Vec<TierOffer>,
}
impl std::fmt::Display for EditionRefusal {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
for (index, tier) in self.tiers.iter().enumerate() {
if index > 0 {
f.write_str("; ")?;
}
match &tier.want {
Some(want) => write!(f, "{} wants {want:?} and can {:?}", tier.tier, tier.can)?,
None => write!(
f,
"{} declares no edition and can {:?}",
tier.tier, tier.can
)?,
}
}
Ok(())
}
}
impl std::error::Error for EditionRefusal {}
fn select_workflow_edition(tiers: &[(&str, &SessionOffer)]) -> Result<String, EditionRefusal> {
let refuse = || EditionRefusal {
tiers: tiers
.iter()
.map(|(tier, offer)| TierOffer {
tier: (*tier).to_string(),
can: offer.editions.clone(),
want: offer.want().map(str::to_string),
})
.collect(),
};
let (_, first) = tiers.first().ok_or_else(&refuse)?;
first
.can()
.filter(|edition| {
tiers
.iter()
.all(|(_, offer)| offer.can().any(|other| other == *edition))
})
.filter(|edition| {
tiers
.iter()
.all(|(_, offer)| offer.want().is_none_or(|want| want_admits(want, edition)))
})
.max_by_key(|edition| edition_sort_key(edition))
.map(str::to_string)
.ok_or_else(&refuse)
}
pub fn negotiate_workflow_edition(
peer: &SessionOffer,
runtime: &SessionOffer,
) -> Result<String, EditionRefusal> {
select_workflow_edition(&[("env", peer), ("runtime", runtime)])
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RuntimeCap {
Capability,
Declaration,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SessionFloor {
pub selected_workflow_edition: String,
pub desired_workflow_edition: String,
pub runtime_cap: Option<RuntimeCap>,
}
impl SessionFloor {
pub fn runtime_limited(&self) -> bool {
self.runtime_cap.is_some()
}
}
pub fn negotiate_session_floor(
env: &SessionOffer,
model: &SessionOffer,
runtime: &SessionOffer,
) -> Result<SessionFloor, EditionRefusal> {
let selected_workflow_edition =
select_workflow_edition(&[("env", env), ("model", model), ("runtime", runtime)])?;
let desired_workflow_edition = select_workflow_edition(&[("env", env), ("model", model)])
.unwrap_or_else(|_| selected_workflow_edition.clone());
let runtime_cap = (desired_workflow_edition != selected_workflow_edition).then(|| {
let undeclared = SessionOffer {
editions: runtime.editions.clone(),
preferred: None,
};
let capable =
select_workflow_edition(&[("env", env), ("model", model), ("runtime", &undeclared)])
.unwrap_or_else(|_| selected_workflow_edition.clone());
if edition_sort_key(&selected_workflow_edition) < edition_sort_key(&capable) {
RuntimeCap::Declaration
} else {
RuntimeCap::Capability
}
});
Ok(SessionFloor {
selected_workflow_edition,
desired_workflow_edition,
runtime_cap,
})
}
pub fn evaluate_handshake(client_protocol_generation: &str) -> bool {
is_protocol_generation_supported(client_protocol_generation)
}
pub fn core_handshake_request(
component: &str,
capabilities: &[&str],
declared: Option<&str>,
) -> core::v1::HandshakeRequest {
core::v1::HandshakeRequest {
protocol_generation: PROTOCOL_GENERATION.to_string(),
peer_info: Some(peer_info(component)),
capabilities: capability_map(capabilities),
supported_workflow_editions: supported_workflow_editions(),
preferred_workflow_edition: declared_workflow_edition(declared).to_string(),
}
}
pub fn generation_mismatch_message(client_protocol_generation: &str) -> String {
format!(
"protocol generation {client_protocol_generation} not compatible with server \
{PROTOCOL_GENERATION}"
)
}
pub fn supported_workflow_editions() -> Vec<String> {
SUPPORTED_WORKFLOW_EDITIONS
.iter()
.map(|edition| (*edition).to_string())
.collect()
}
pub fn capability_map(names: &[&str]) -> HashMap<String, String> {
names
.iter()
.map(|name| ((*name).to_string(), "true".to_string()))
.collect()
}
pub fn has_capability(map: &HashMap<String, String>, name: &str) -> bool {
map.get(name).is_some_and(|value| value == "true")
}
pub mod core {
pub mod v1 {
tonic::include_proto!("rlmesh.core.v1");
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct PeerInfoOverride {
pub language: String,
pub language_version: String,
pub package_version: String,
pub os: String,
pub os_version: String,
pub arch: String,
pub framework_versions: HashMap<String, String>,
pub extra: HashMap<String, String>,
}
static PEER_INFO_OVERRIDE: RwLock<Option<PeerInfoOverride>> = RwLock::new(None);
pub fn set_peer_info_override(info: PeerInfoOverride) {
if let Ok(mut guard) = PEER_INFO_OVERRIDE.write() {
*guard = Some(info);
}
}
pub fn peer_info(component: &str) -> core::v1::PeerInfo {
let mut extra = HashMap::new();
extra.insert(
"rlmesh.workflow.base".to_string(),
WORKFLOW_EDITION_BASE.to_string(),
);
extra.insert(
"rlmesh.workflow.edition".to_string(),
CURRENT_WORKFLOW_EDITION.to_string(),
);
extra.insert("rlmesh.build.cohort".to_string(), BUILD_COHORT.to_string());
extra.insert("rlmesh.build.source".to_string(), BUILD_SOURCE.to_string());
let mut info = core::v1::PeerInfo {
component: component.to_string(),
package_version: env!("CARGO_PKG_VERSION").to_string(),
language: "rust".to_string(),
language_version: String::new(),
os: std::env::consts::OS.to_string(),
os_version: String::new(),
arch: std::env::consts::ARCH.to_string(),
framework_versions: HashMap::new(),
extra,
};
if let Ok(guard) = PEER_INFO_OVERRIDE.read()
&& let Some(over) = guard.as_ref()
{
if !over.language.is_empty() {
info.language = over.language.clone();
}
if !over.language_version.is_empty() {
info.language_version = over.language_version.clone();
}
if !over.package_version.is_empty() {
info.package_version = over.package_version.clone();
}
if !over.os.is_empty() {
info.os = over.os.clone();
}
if !over.os_version.is_empty() {
info.os_version = over.os_version.clone();
}
if !over.arch.is_empty() {
info.arch = over.arch.clone();
}
if !over.framework_versions.is_empty() {
info.framework_versions = over.framework_versions.clone();
}
if !over.extra.is_empty() {
info.extra.extend(over.extra.clone());
}
}
info
}
pub mod env {
pub mod v1 {
tonic::include_proto!("rlmesh.env.v1");
}
}
pub mod spaces {
pub mod v1 {
tonic::include_proto!("rlmesh.spaces.v1");
}
}
pub mod model {
pub mod v1 {
tonic::include_proto!("rlmesh.model.v1");
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct EndpointPhases {
pub decode_ns: u64,
pub user_ns: u64,
pub encode_ns: u64,
pub queue_ns: u64,
pub adapter_ns: u64,
pub held_episodes: Option<u32>,
pub held_state_bytes: Option<u64>,
pub in_flight: u32,
pub lane_skew_ns: Option<u64>,
}
impl EndpointPhases {
pub fn reported(ns: u64) -> Option<u64> {
(ns != 0).then_some(ns)
}
pub fn from_durations(
decode: std::time::Duration,
user: std::time::Duration,
encode: std::time::Duration,
) -> Self {
fn ns(duration: std::time::Duration) -> u64 {
duration.as_nanos().min(u128::from(u64::MAX)) as u64
}
Self {
decode_ns: ns(decode),
user_ns: ns(user),
encode_ns: ns(encode),
..Self::default()
}
}
pub fn nest(decode_ns: u64, call_ns: u64, encode_ns: u64, inner: Self) -> Self {
let inner_wire = inner.decode_ns.saturating_add(inner.encode_ns);
if inner_wire > call_ns {
return Self {
decode_ns,
user_ns: call_ns,
encode_ns,
lane_skew_ns: inner.lane_skew_ns,
..Self::default()
};
}
Self {
decode_ns: decode_ns.saturating_add(inner.decode_ns),
user_ns: call_ns - inner_wire,
encode_ns: encode_ns.saturating_add(inner.encode_ns),
lane_skew_ns: inner.lane_skew_ns,
..Self::default()
}
}
pub fn from_env_response(response: &env::v1::JoinResponse) -> Self {
Self {
decode_ns: response.decode_ns.unwrap_or(0),
user_ns: response.user_ns.unwrap_or(0),
encode_ns: response.encode_ns.unwrap_or(0),
queue_ns: response.queue_ns.unwrap_or(0),
lane_skew_ns: response.lane_skew_ns,
..Self::default()
}
}
pub fn from_model_response(response: &model::v1::JoinResponse) -> Self {
Self {
decode_ns: response.decode_ns.unwrap_or(0),
user_ns: response.user_ns.unwrap_or(0),
encode_ns: response.encode_ns.unwrap_or(0),
queue_ns: response.queue_ns.unwrap_or(0),
in_flight: response.in_flight.unwrap_or(0),
adapter_ns: response.adapter_ns.unwrap_or(0),
held_episodes: response.held_episodes,
held_state_bytes: response.held_state_bytes,
..Self::default()
}
}
}
pub fn elapsed_ns(started_at: std::time::Instant) -> u64 {
started_at.elapsed().as_nanos().min(u128::from(u64::MAX)) as u64
}
pub fn lane_skew_ns(lane_ns: &mut [u64]) -> u64 {
if lane_ns.len() < 2 {
return 0;
}
let (_, median, slower) = lane_ns.select_nth_unstable((lane_ns.len() - 1) / 2);
let median = *median;
slower
.iter()
.copied()
.max()
.unwrap_or(median)
.saturating_sub(median)
}
pub fn lane_final_info(
info: Option<&spaces::v1::MetaMap>,
lane: usize,
width: usize,
) -> Option<spaces::v1::MetaMap> {
use spaces::v1::meta_value::Kind;
let info = info?;
let Some(final_info) = info.entries.get("final_info") else {
if width == 1 {
return Some(info.clone());
}
let sliced = vector_info_lane(info, lane, width);
return (!sliced.entries.is_empty()).then_some(sliced);
};
let is_present = match info.entries.get("_final_info") {
Some(mask) => meta_bool_at(mask, lane).unwrap_or(false),
None => width == 1,
};
if !is_present {
return None;
}
match &final_info.kind {
Some(Kind::Map(map)) => Some(map.clone()),
Some(Kind::List(list)) => match &list.items.get(lane)?.kind {
Some(Kind::Map(map)) => Some(map.clone()),
_ => None,
},
_ => None,
}
}
fn vector_info_lane(info: &spaces::v1::MetaMap, lane: usize, width: usize) -> spaces::v1::MetaMap {
use spaces::v1::meta_value::Kind;
let entries = info
.entries
.iter()
.filter(|(key, _)| {
!key.strip_prefix('_')
.is_some_and(|masked| info.entries.contains_key(masked))
})
.filter(|(key, _)| {
info.entries
.get(&format!("_{key}"))
.is_none_or(|mask| meta_bool_at(mask, lane) == Some(true))
})
.filter_map(|(key, value)| {
let lane_value = match &value.kind {
Some(Kind::Map(map)) => spaces::v1::MetaValue {
kind: Some(Kind::Map(vector_info_lane(map, lane, width))),
},
Some(Kind::List(list)) if list.items.len() == width => list.items[lane].clone(),
_ => return None,
};
Some((key.clone(), lane_value))
})
.collect();
spaces::v1::MetaMap { entries }
}
fn meta_bool_at(value: &spaces::v1::MetaValue, lane: usize) -> Option<bool> {
use spaces::v1::meta_value::Kind;
match &value.kind {
Some(Kind::List(list)) => match list.items.get(lane)?.kind {
Some(Kind::Bool(flag)) => Some(flag),
_ => None,
},
_ => None,
}
}
#[cfg(test)]
#[path = "../build_manifest.rs"]
mod build_manifest;
#[cfg(test)]
mod tests {
use super::{
CURRENT_WORKFLOW_EDITION, PROTOCOL_GENERATION, SUPPORTED_WORKFLOW_EDITIONS, SessionOffer,
evaluate_handshake, is_protocol_generation_supported, negotiate_session_floor,
negotiate_workflow_edition, supported_workflow_editions,
};
fn tier(can: &[&str], want: Option<&str>) -> SessionOffer {
SessionOffer {
editions: offer(can),
preferred: want.map(str::to_string),
}
}
#[test]
fn lane_final_info_reads_each_gymnasium_info_layout() {
use super::lane_final_info;
use super::spaces::v1::meta_value::Kind;
use super::spaces::v1::{MetaList, MetaMap, MetaValue};
let value = |kind: Kind| MetaValue { kind: Some(kind) };
let list = |items: Vec<MetaValue>| value(Kind::List(MetaList { items }));
let map = |entries: Vec<(&str, MetaValue)>| MetaMap {
entries: entries
.into_iter()
.map(|(key, value)| (key.to_string(), value))
.collect(),
};
let success = |flag: bool| map(vec![("is_success", value(Kind::Bool(flag)))]);
let scalar = success(false);
assert_eq!(lane_final_info(Some(&scalar), 0, 1), Some(success(false)));
assert_eq!(lane_final_info(None, 0, 2), None);
let explicit = map(vec![
(
"final_info",
list(vec![
value(Kind::Map(success(true))),
value(Kind::Map(MetaMap::default())),
]),
),
(
"_final_info",
list(vec![value(Kind::Bool(true)), value(Kind::Bool(false))]),
),
]);
assert_eq!(lane_final_info(Some(&explicit), 0, 2), Some(success(true)));
assert_eq!(lane_final_info(Some(&explicit), 1, 2), None);
let unmasked = map(vec![(
"is_success",
list(vec![value(Kind::Bool(false)), value(Kind::Bool(true))]),
)]);
assert_eq!(lane_final_info(Some(&unmasked), 1, 2), Some(success(true)));
let masked_out = map(vec![
(
"is_success",
list(vec![value(Kind::Bool(true)), value(Kind::Bool(false))]),
),
(
"_is_success",
list(vec![value(Kind::Bool(true)), value(Kind::Bool(false))]),
),
]);
assert_eq!(lane_final_info(Some(&masked_out), 1, 2), None);
}
#[test]
fn phases_from_a_peer_that_stamps_nothing_read_back_as_zero() {
use super::{EndpointPhases, env, model};
let env_response = env::v1::JoinResponse {
endpoint_total_ns: Some(1_000),
..Default::default()
};
assert_eq!(
EndpointPhases::from_env_response(&env_response),
EndpointPhases::default()
);
let model_response = model::v1::JoinResponse {
endpoint_total_ns: Some(1_000),
..Default::default()
};
assert_eq!(
EndpointPhases::from_model_response(&model_response),
EndpointPhases::default()
);
let stamped = model::v1::JoinResponse {
endpoint_total_ns: Some(1_000),
decode_ns: Some(10),
user_ns: Some(20),
encode_ns: Some(30),
queue_ns: Some(40),
in_flight: Some(3),
adapter_ns: Some(15),
held_episodes: Some(5),
held_state_bytes: Some(6_000),
..Default::default()
};
assert_eq!(
EndpointPhases::from_model_response(&stamped),
EndpointPhases {
decode_ns: 10,
user_ns: 20,
encode_ns: 30,
queue_ns: 40,
in_flight: 3,
adapter_ns: 15,
held_episodes: Some(5),
held_state_bytes: Some(6_000),
lane_skew_ns: None,
}
);
let skewed = env::v1::JoinResponse {
endpoint_total_ns: Some(1_000),
lane_skew_ns: Some(880),
..Default::default()
};
assert_eq!(
EndpointPhases::from_env_response(&skewed).lane_skew_ns,
Some(880)
);
let even = env::v1::JoinResponse {
lane_skew_ns: Some(0),
..Default::default()
};
assert_eq!(
EndpointPhases::from_env_response(&even).lane_skew_ns,
Some(0)
);
assert_eq!(EndpointPhases::reported(0), None);
assert_eq!(EndpointPhases::reported(7), Some(7));
}
#[test]
fn nesting_folds_an_inner_split_into_the_outer_one() {
use super::EndpointPhases;
assert_eq!(
EndpointPhases::nest(1, 10, 2, EndpointPhases::default()),
EndpointPhases {
decode_ns: 1,
user_ns: 10,
encode_ns: 2,
..EndpointPhases::default()
}
);
let inner = EndpointPhases {
decode_ns: 3,
user_ns: 5,
encode_ns: 2,
..EndpointPhases::default()
};
assert_eq!(
EndpointPhases::nest(1, 10, 2, inner),
EndpointPhases {
decode_ns: 4,
user_ns: 5,
encode_ns: 4,
..EndpointPhases::default()
}
);
assert_eq!(
EndpointPhases::nest(1, 14, 2, inner),
EndpointPhases {
decode_ns: 4,
user_ns: 9,
encode_ns: 4,
..EndpointPhases::default()
}
);
assert_eq!(
EndpointPhases::nest(1, 4, 2, inner),
EndpointPhases {
decode_ns: 1,
user_ns: 4,
encode_ns: 2,
..EndpointPhases::default()
}
);
}
#[test]
fn lane_skew_is_the_slowest_lane_over_the_median_one() {
use super::lane_skew_ns;
assert_eq!(lane_skew_ns(&mut [10, 10, 10, 10, 10, 10, 10, 900]), 890);
assert_eq!(lane_skew_ns(&mut [900; 8]), 0);
assert_eq!(lane_skew_ns(&mut [10, 100]), 90);
assert_eq!(lane_skew_ns(&mut [1, 100, 100, 100]), 0);
assert_eq!(lane_skew_ns(&mut [42]), 0);
assert_eq!(lane_skew_ns(&mut []), 0);
}
#[test]
fn nesting_keeps_the_inner_envs_lane_skew() {
use super::EndpointPhases;
let measured = EndpointPhases {
user_ns: 5,
lane_skew_ns: Some(77),
..EndpointPhases::default()
};
assert_eq!(
EndpointPhases::nest(1, 10, 2, measured).lane_skew_ns,
Some(77)
);
let skew_only = EndpointPhases {
lane_skew_ns: Some(77),
..EndpointPhases::default()
};
assert_eq!(
EndpointPhases::nest(1, 10, 2, skew_only).lane_skew_ns,
Some(77)
);
}
fn offer(editions: &[&str]) -> Vec<String> {
editions.iter().map(|edition| edition.to_string()).collect()
}
#[test]
fn peer_info_default_then_override_merges_python_with_rust_fallback() {
use super::{PeerInfoOverride, peer_info, set_peer_info_override};
use std::collections::HashMap;
let rust_info = peer_info("rlmesh-env");
assert_eq!(rust_info.component, "rlmesh-env");
assert_eq!(rust_info.language, "rust");
assert!(rust_info.language_version.is_empty());
assert!(rust_info.framework_versions.is_empty());
let detected_os = rust_info.os.clone();
let detected_arch = rust_info.arch.clone();
let detected_pkg = rust_info.package_version.clone();
let mut frameworks = HashMap::new();
frameworks.insert("numpy".to_string(), "1.26.4".to_string());
set_peer_info_override(PeerInfoOverride {
language: "python".to_string(),
language_version: "3.11.4".to_string(),
package_version: String::new(),
os: String::new(),
os_version: "ubuntu-22.04".to_string(),
arch: "aarch64".to_string(),
framework_versions: frameworks,
extra: HashMap::from([("rlmesh.startup.listen_ms".to_string(), "4200".to_string())]),
});
let py_info = peer_info("rlmesh-env");
assert_eq!(
py_info
.extra
.get("rlmesh.startup.listen_ms")
.map(String::as_str),
Some("4200")
);
assert!(py_info.extra.contains_key("rlmesh.build.cohort"));
assert_eq!(py_info.component, "rlmesh-env");
assert_eq!(py_info.language, "python");
assert_eq!(py_info.language_version, "3.11.4");
assert_eq!(py_info.os_version, "ubuntu-22.04");
assert_eq!(py_info.arch, "aarch64");
assert_eq!(
py_info.framework_versions.get("numpy").map(String::as_str),
Some("1.26.4")
);
assert_eq!(py_info.os, detected_os);
assert_eq!(py_info.package_version, detected_pkg);
assert_eq!(
py_info
.extra
.get("rlmesh.workflow.edition")
.map(String::as_str),
Some(CURRENT_WORKFLOW_EDITION)
);
let _ = detected_arch;
}
#[test]
fn has_capability_reads_advertised_features() {
use super::{capabilities, capability_map, has_capability};
let map = capability_map(&[capabilities::MODEL_CONCURRENT_PREDICT_V1]);
assert!(has_capability(
&map,
capabilities::MODEL_CONCURRENT_PREDICT_V1
));
assert!(!has_capability(&map, "rlmesh.not.advertised.v1"));
assert_eq!(map[capabilities::MODEL_CONCURRENT_PREDICT_V1], "true");
for value in ["1", "yes", "TRUE", ""] {
let odd =
std::collections::HashMap::from([("rlmesh.odd.v1".to_string(), value.to_string())]);
assert!(!has_capability(&odd, "rlmesh.odd.v1"), "{value:?}");
}
}
#[test]
fn env_subset_step_keeps_its_shipped_wire_spelling() {
assert_eq!(super::capabilities::ENV_SUBSET_STEP, "subset_step");
}
#[test]
fn protocol_generation_is_plain_equality() {
assert!(is_protocol_generation_supported(PROTOCOL_GENERATION));
assert!(is_protocol_generation_supported(&format!(
" {PROTOCOL_GENERATION} "
)));
assert!(!is_protocol_generation_supported("rlmesh-wire-v2"));
assert!(!is_protocol_generation_supported(""));
assert!(!is_protocol_generation_supported("0.1.0"));
}
#[test]
fn split_editions_recovers_every_supported_edition() {
const LIST: &str = "2026.09-0.2.0-rc.1,2026.06";
const EDITIONS: [&str; super::edition_count(LIST)] = super::split_editions(LIST);
assert_eq!(EDITIONS, ["2026.09-0.2.0-rc.1", "2026.06"]);
assert_eq!(super::edition_count("2026.06"), 1);
}
#[test]
fn manifest_string_list_reads_both_array_spellings() {
use super::build_manifest::manifest_string_list;
let single = "[workflow]\nsupported_editions = [\"2026.09-0.2.0-rc.1\", \"2026.06\"]\n";
let multi = "[workflow]\nsupported_editions = [\n \"2026.09-0.2.0-rc.1\", # cohort\n \
\"2026.06\",\n]\n";
let want = ["2026.09-0.2.0-rc.1".to_string(), "2026.06".to_string()];
assert_eq!(manifest_string_list(single, "supported_editions"), want);
assert_eq!(manifest_string_list(multi, "supported_editions"), want);
assert!(manifest_string_list(single, "current_edition").is_empty());
}
#[test]
fn supported_workflow_editions_lead_with_current() {
assert!(!SUPPORTED_WORKFLOW_EDITIONS.is_empty());
assert_eq!(SUPPORTED_WORKFLOW_EDITIONS[0], CURRENT_WORKFLOW_EDITION);
let unique: std::collections::BTreeSet<&&str> =
SUPPORTED_WORKFLOW_EDITIONS.iter().collect();
assert_eq!(unique.len(), SUPPORTED_WORKFLOW_EDITIONS.len());
assert!(
SUPPORTED_WORKFLOW_EDITIONS
.iter()
.all(|edition| !edition.trim().is_empty())
);
assert_eq!(
supported_workflow_editions(),
SUPPORTED_WORKFLOW_EDITIONS
.iter()
.map(|edition| (*edition).to_string())
.collect::<Vec<_>>()
);
}
#[test]
fn negotiation_selects_mutual_edition() {
let runtime = SessionOffer::this_build(None);
assert_eq!(
negotiate_workflow_edition(&SessionOffer::new(&[CURRENT_WORKFLOW_EDITION]), &runtime),
Ok(CURRENT_WORKFLOW_EDITION.to_string())
);
assert_eq!(
negotiate_workflow_edition(
&SessionOffer::new(&["2025.01", CURRENT_WORKFLOW_EDITION, "2031.12"]),
&runtime
),
Ok(CURRENT_WORKFLOW_EDITION.to_string())
);
}
#[test]
fn negotiation_trims_offered_editions() {
let padded = SessionOffer {
editions: vec![format!(" {CURRENT_WORKFLOW_EDITION} ")],
preferred: None,
};
assert_eq!(
negotiate_workflow_edition(&padded, &SessionOffer::this_build(None)),
Ok(CURRENT_WORKFLOW_EDITION.to_string())
);
}
#[test]
fn negotiation_rejects_unknown_or_empty_offers() {
let runtime = SessionOffer::this_build(None);
for offered in [
&[][..],
&[""][..],
&["2026"][..],
&["next"][..],
&["2026.11", "2027.01"][..],
] {
let refusal = negotiate_workflow_edition(&SessionOffer::new(offered), &runtime)
.expect_err("no mutual edition");
let message = refusal.to_string();
assert!(message.contains("env"), "{message}");
assert!(message.contains("runtime"), "{message}");
assert!(message.contains(CURRENT_WORKFLOW_EDITION), "{message}");
}
}
#[test]
fn evaluate_handshake_gates_generation_only() {
assert!(evaluate_handshake(PROTOCOL_GENERATION));
assert!(!evaluate_handshake("rlmesh-wire-v2"));
}
#[test]
fn is_retained_edition_matches_the_window() {
use super::{Edition, is_retained_edition};
let current = Edition::parse(CURRENT_WORKFLOW_EDITION).expect("the current edition parses");
assert!(is_retained_edition(current));
assert_eq!(
Edition::parse(&format!(" {CURRENT_WORKFLOW_EDITION} ")),
Ok(current)
);
assert_eq!(Edition::parse(super::WORKFLOW_EDITION_BASE), Ok(current));
assert!(Edition::parse("2099.01").is_err());
assert!(Edition::parse("").is_err());
}
#[test]
fn parse_resolves_any_cohort_spelling_of_a_known_base() {
use super::{Edition, WORKFLOW_EDITION_BASE, parse_retained_edition};
let current = Edition::parse(CURRENT_WORKFLOW_EDITION).expect("the current edition parses");
assert_eq!(Edition::parse(WORKFLOW_EDITION_BASE), Ok(current));
assert_eq!(Edition::parse(CURRENT_WORKFLOW_EDITION), Ok(current));
for foreign in [
format!("{WORKFLOW_EDITION_BASE}-0.0.1-rc.1"),
format!("{WORKFLOW_EDITION_BASE}-dev.deadbeef"),
] {
assert_eq!(
Edition::parse(&foreign),
Ok(current),
"{foreign} is a cohort of a known base"
);
}
assert!(Edition::parse("2099.01").is_err());
assert!(Edition::parse("2099.01-0.1.0-rc.12").is_err());
assert_eq!(
parse_retained_edition(CURRENT_WORKFLOW_EDITION),
Ok(current)
);
let error = parse_retained_edition(" 2099.01 ").expect_err("no arm implements it");
assert!(
error.contains("\"2099.01\"") && error.contains(CURRENT_WORKFLOW_EDITION),
"expected the arrived name and the retained list, got: {error}"
);
}
#[test]
fn a_declaration_must_leave_this_build_something_to_run() {
use super::{Edition, WORKFLOW_EDITION_BASE, parse_declared_edition};
let current = Edition::parse(CURRENT_WORKFLOW_EDITION).expect("the current edition parses");
for offered in SUPPORTED_WORKFLOW_EDITIONS {
assert_eq!(
parse_declared_edition(offered),
Ok(current),
"{offered} is offered, so it is declarable"
);
}
assert_eq!(parse_declared_edition(WORKFLOW_EDITION_BASE), Ok(current));
assert_eq!(
parse_declared_edition(&format!(" {WORKFLOW_EDITION_BASE} ")),
Ok(current)
);
let stale_cohort = format!("{WORKFLOW_EDITION_BASE}-0.0.0");
match parse_declared_edition(&stale_cohort) {
Ok(edition) => {
assert_eq!(edition, current);
assert!(SUPPORTED_WORKFLOW_EDITIONS.contains(&WORKFLOW_EDITION_BASE));
}
Err(error) => {
assert!(
error.contains(&stale_cohort)
&& error.contains(CURRENT_WORKFLOW_EDITION)
&& error.contains(&format!("{WORKFLOW_EDITION_BASE:?} base")),
"expected both halves of the mismatch, got: {error}"
);
}
}
assert!(parse_declared_edition("2099.01").is_err());
}
#[test]
fn want_admits_by_base_when_bare_and_by_key_when_suffixed() {
use super::want_admits;
assert!(want_admits("2026.06", "2026.06"));
assert!(want_admits("2026.06", "2026.06-dev.aaa"));
assert!(want_admits("2026.06", "2026.06-0.1.0-rc.12"));
assert!(want_admits("2026.06", "2026.01"));
assert!(!want_admits("2026.06", "2026.08"));
assert!(!want_admits("2026.06", "2026.08-dev.aaa"));
assert!(want_admits("2026.06-dev.bbb", "2026.06-dev.bbb"));
assert!(want_admits("2026.06-dev.bbb", "2026.06-dev.aaa"));
assert!(want_admits("2026.06-dev.bbb", "2026.06"));
assert!(!want_admits("2026.06-dev.bbb", "2026.06-dev.ccc"));
assert!(!want_admits("2026.06-dev.bbb", "2026.08"));
}
#[test]
fn every_arm_is_listed_and_has_a_row() {
use super::{Edition, defaults, is_retained_edition};
for edition in Edition::ALL {
match edition {
Edition::E2026_06 => {
assert!(Edition::ALL.contains(&Edition::E2026_06));
}
}
let row = defaults(*edition);
assert!(!row.trial_index_option_key.is_empty());
assert!(!row.conformance_warning_info_key.is_empty());
assert_eq!(Edition::parse(edition.base()), Ok(*edition));
assert_eq!(edition.to_string(), edition.base());
}
assert_eq!(
Edition::parse(CURRENT_WORKFLOW_EDITION).map(is_retained_edition),
Ok(true)
);
assert_eq!(
Edition::parse(CURRENT_WORKFLOW_EDITION),
Ok(Edition::current())
);
for retained in SUPPORTED_WORKFLOW_EDITIONS {
let edition = Edition::parse(retained)
.unwrap_or_else(|err| panic!("retained edition {retained:?} has no arm: {err}"));
assert!(is_retained_edition(edition));
}
}
#[test]
fn session_floor_picks_highest_all_three_support() {
let env = SessionOffer::new(&["2026.01", " 2026.06 "]);
let model = SessionOffer::new(&["2026.06", "2026.01"]);
let runtime = SessionOffer::new(&["2026.06"]);
let floor = negotiate_session_floor(&env, &model, &runtime).expect("a floor");
assert_eq!(floor.selected_workflow_edition, "2026.06");
assert_eq!(floor.desired_workflow_edition, "2026.06");
assert!(!floor.runtime_limited());
}
#[test]
fn session_floor_flags_runtime_as_limiting_tier() {
let env = SessionOffer::new(&["2026.06", "2026.08"]);
let model = SessionOffer::new(&["2026.06", "2026.08"]);
let runtime = SessionOffer::new(&["2026.06"]);
let floor = negotiate_session_floor(&env, &model, &runtime).expect("a floor");
assert_eq!(floor.selected_workflow_edition, "2026.06");
assert_eq!(floor.desired_workflow_edition, "2026.08");
assert!(floor.runtime_limited());
}
#[test]
fn session_floor_is_none_when_no_common_edition() {
let env = SessionOffer::new(&["2026.08"]);
let model = SessionOffer::new(&["2026.08"]);
let runtime = SessionOffer::new(&["2026.06"]);
assert!(negotiate_session_floor(&env, &model, &runtime).is_err());
let empty = SessionOffer::new(&[""]);
let ok = SessionOffer::new(&["2026.06"]);
assert!(negotiate_session_floor(&empty, &ok, &ok).is_err());
}
#[test]
fn edition_ordering_prefers_exact_cohort_then_newer_date() {
use super::edition_sort_key;
assert!(edition_sort_key("2026.06-0.1.0-rc.1") > edition_sort_key("2026.06"));
assert!(edition_sort_key("2026.09-0.2.0-beta.1") > edition_sort_key("2026.06"));
assert!(edition_sort_key("2026.09") > edition_sort_key("2026.06-0.1.0-rc.1"));
assert!(edition_sort_key("2026.06-0.1.0-rc.2") > edition_sort_key("2026.06-0.1.0-rc.1"));
assert_eq!(
negotiate_workflow_edition(
&SessionOffer::new(&["2025.01", CURRENT_WORKFLOW_EDITION, "2099.12"]),
&SessionOffer::this_build(None)
),
Ok(CURRENT_WORKFLOW_EDITION.to_string())
);
}
#[test]
fn session_floor_prefers_exact_edition_cohort_over_sealed_fallback() {
let editions = &["2026.06", "2026.06-0.1.0-rc.1"];
let offer = SessionOffer::new(editions);
let floor = negotiate_session_floor(&offer, &offer, &offer).expect("a floor");
assert_eq!(floor.selected_workflow_edition, "2026.06-0.1.0-rc.1");
}
#[test]
fn want_can_selection_table() {
type Tier<'a> = (&'a str, &'a [&'a str], Option<&'a str>);
let cases: &[(&str, &[Tier<'_>], Option<&str>)] = &[
(
"two-party, undeclared: highest mutual",
&[
("env", &["2026.01", "2026.06"], None),
("runtime", &["2026.06"], None),
],
Some("2026.06"),
),
(
"three-party, undeclared: highest all three share",
&[
("env", &["2026.01", " 2026.06 "], None),
("model", &["2026.06", "2026.01"], None),
("runtime", &["2026.06"], None),
],
Some("2026.06"),
),
(
"three-party, undeclared: the runtime's CAN set caps the floor",
&[
("env", &["2026.06", "2026.08"], None),
("model", &["2026.06", "2026.08"], None),
("runtime", &["2026.06"], None),
],
Some("2026.06"),
),
(
"three-party, undeclared: empty intersection refuses",
&[
("env", &["2026.08"], None),
("model", &["2026.08"], None),
("runtime", &["2026.06"], None),
],
None,
),
(
"one peer declares below the mutual max: the declared one wins",
&[
("env", &["2026.06", "2026.08"], Some("2026.06")),
("model", &["2026.06", "2026.08"], None),
("runtime", &["2026.06", "2026.08"], None),
],
Some("2026.06"),
),
(
"a WANT above what another peer CAN never lifts the floor",
&[
("env", &["2026.06", "2026.08", "2026.10"], Some("2026.10")),
("model", &["2026.06", "2026.08"], None),
("runtime", &["2026.06", "2026.08", "2026.10"], None),
],
Some("2026.08"),
),
(
"a WANT the intersection does not contain still selects the \
highest mutual below it: a pin is a ceiling, not an exact demand",
&[
("env", &["2026.06", "2026.08"], Some("2026.07")),
("model", &["2026.06", "2026.08"], None),
("runtime", &["2026.06", "2026.08"], None),
],
Some("2026.06"),
),
(
"a WANT below everything mutual refuses rather than running higher",
&[
("env", &["2026.06"], Some("2025.01")),
("model", &["2026.06"], None),
("runtime", &["2026.06"], None),
],
None,
),
(
"the lowest WANT wins when several are declared",
&[
("env", &["2026.06", "2026.08", "2026.10"], Some("2026.10")),
("model", &["2026.06", "2026.08", "2026.10"], Some("2026.08")),
(
"runtime",
&["2026.06", "2026.08", "2026.10"],
Some("2026.10"),
),
],
Some("2026.08"),
),
(
"undeclared, matching dev cohorts: the exact cohort beats its base",
&[
("env", &["2026.06", "2026.06-dev.aaa"], None),
("model", &["2026.06", "2026.06-dev.aaa"], None),
("runtime", &["2026.06", "2026.06-dev.aaa"], None),
],
Some("2026.06-dev.aaa"),
),
(
"differing dev cohorts do not match: both fall back to the sealed base",
&[
("env", &["2026.06", "2026.06-dev.aaa"], None),
("model", &["2026.06", "2026.06-dev.bbb"], None),
("runtime", &["2026.06", "2026.06-dev.aaa"], None),
],
Some("2026.06"),
),
(
"a WANT of the bare base admits every cohort of it: the exact cohort still wins",
&[
("env", &["2026.06", "2026.06-dev.aaa"], Some("2026.06")),
("model", &["2026.06", "2026.06-dev.aaa"], None),
("runtime", &["2026.06", "2026.06-dev.aaa"], None),
],
Some("2026.06-dev.aaa"),
),
(
"bare base WANT on both sides, cohort-only CANs: selects the cohort",
&[
("env", &["2026.06-dev.aaa"], Some("2026.06")),
("runtime", &["2026.06-dev.aaa"], Some("2026.06")),
],
Some("2026.06-dev.aaa"),
),
(
"bare base WANT on one side, three cohort-only tiers: selects the cohort",
&[
("env", &["2026.06-dev.aaa"], Some("2026.06")),
("model", &["2026.06-dev.aaa"], None),
("runtime", &["2026.06-dev.aaa"], None),
],
Some("2026.06-dev.aaa"),
),
(
"bare base WANT admits its own cohorts but not a newer base's",
&[
(
"env",
&["2026.06-dev.aaa", "2026.08-dev.aaa"],
Some("2026.06"),
),
("runtime", &["2026.06-dev.aaa", "2026.08-dev.aaa"], None),
],
Some("2026.06-dev.aaa"),
),
(
"bare base WANT against a newer base only refuses",
&[
("env", &["2026.08"], Some("2026.06")),
("model", &["2026.08"], None),
("runtime", &["2026.08", "2026.08-dev.aaa"], None),
],
None,
),
(
"cohort WANT against a differing cohort of the same base refuses",
&[
("env", &["2026.06-dev.aaa"], Some("2026.06-dev.aaa")),
("runtime", &["2026.06-dev.bbb"], None),
],
None,
),
(
"cohort WANT keeps the exact order: a higher cohort in the intersection is excluded",
&[
("env", &["2026.06-dev.bbb"], Some("2026.06-dev.aaa")),
("runtime", &["2026.06-dev.bbb"], None),
],
None,
),
(
"a 0.1.0-shaped peer (one CAN, no WANT) pins a newer build to 2026.06",
&[
("env", &["2026.06"], None),
("runtime", &["2026.06", "2099.01"], None),
],
Some("2026.06"),
),
(
"the same peer against a runtime that declares the newer edition",
&[
("env", &["2026.06"], None),
("model", &["2026.06"], None),
("runtime", &["2026.06", "2099.01"], Some("2099.01")),
],
Some("2026.06"),
),
(
"a tier that CAN nothing refuses",
&[("env", &[], None), ("runtime", &["2026.06"], None)],
None,
),
];
for (name, tiers, expected) in cases {
let offers: Vec<(&str, SessionOffer)> = tiers
.iter()
.map(|(tier_name, can, want)| (*tier_name, tier(can, *want)))
.collect();
let borrowed: Vec<(&str, &SessionOffer)> = offers
.iter()
.map(|(tier_name, o)| (*tier_name, o))
.collect();
let selected = super::select_workflow_edition(&borrowed);
match expected {
Some(edition) => assert_eq!(selected.as_deref(), Ok(*edition), "{name}"),
None => {
let refusal = selected.expect_err(name).to_string();
for (tier_name, can, _) in tiers.iter() {
assert!(refusal.contains(tier_name), "{name}: {refusal}");
for edition in can.iter() {
assert!(refusal.contains(edition), "{name}: {refusal}");
}
}
}
}
}
}
#[test]
fn undeclared_everywhere_is_the_old_highest_mutual() {
fn highest_mutual(sets: &[&[String]]) -> Option<String> {
let (first, rest) = sets.split_first()?;
first
.iter()
.map(|value| value.trim())
.filter(|value| !value.is_empty())
.filter(|value| {
rest.iter()
.all(|set| set.iter().any(|other| other.trim() == *value))
})
.max_by_key(|value| super::edition_sort_key(value))
.map(|value| value.to_string())
}
let shapes: &[&[&[&str]]] = &[
&[&["2026.06"], &["2026.06"]],
&[&["2026.01", " 2026.06 "], &["2026.06", "2026.01"]],
&[&["2025.01", "2026.06", "2031.12"], &["2026.06"]],
&[&["2026.11", "2027.01"], &["2026.06"]],
&[&[""], &["2026.06"]],
&[&[], &["2026.06"]],
&[
&["2026.06", "2026.08"],
&["2026.06", "2026.08"],
&["2026.06"],
],
&[&["2026.08"], &["2026.08"], &["2026.06"]],
&[
&["2026.06", "2026.06-0.1.0-rc.1"],
&["2026.06", "2026.06-0.1.0-rc.1"],
&["2026.06", "2026.06-0.1.0-rc.1"],
],
&[
&["2026.06", "2026.06-dev.aaa"],
&["2026.06", "2026.06-dev.bbb"],
&["2026.06", "2026.06-dev.aaa"],
],
&[&["2026.06"], &["2026.06", "2099.01"]],
];
for shape in shapes {
let offers: Vec<SessionOffer> = shape.iter().map(|can| tier(can, None)).collect();
let owned: Vec<Vec<String>> = shape.iter().map(|can| offer(can)).collect();
let sets: Vec<&[String]> = owned.iter().map(Vec::as_slice).collect();
let borrowed: Vec<(&str, &SessionOffer)> = offers.iter().map(|o| ("tier", o)).collect();
assert_eq!(
super::select_workflow_edition(&borrowed).ok(),
highest_mutual(&sets),
"{shape:?}"
);
}
}
#[test]
fn session_floor_reports_why_the_runtime_capped_the_session() {
use super::RuntimeCap;
let peers = tier(&["2026.06", "2026.08"], None);
let floor =
negotiate_session_floor(&peers, &peers, &tier(&["2026.06"], None)).expect("a floor");
assert_eq!(floor.runtime_cap, Some(RuntimeCap::Capability));
let runtime = tier(&["2026.06", "2026.08"], Some("2026.06"));
let floor = negotiate_session_floor(&peers, &peers, &runtime).expect("a floor");
assert_eq!(floor.selected_workflow_edition, "2026.06");
assert_eq!(floor.desired_workflow_edition, "2026.08");
assert_eq!(floor.runtime_cap, Some(RuntimeCap::Declaration));
let wide = tier(&["2026.06", "2026.08", "2026.10"], None);
let runtime = tier(&["2026.06", "2026.08"], Some("2026.06"));
let floor = negotiate_session_floor(&wide, &wide, &runtime).expect("a floor");
assert_eq!(floor.selected_workflow_edition, "2026.06");
assert_eq!(floor.desired_workflow_edition, "2026.10");
assert_eq!(floor.runtime_cap, Some(RuntimeCap::Declaration));
let pinned = tier(&["2026.06", "2026.08"], Some("2026.06"));
let floor = negotiate_session_floor(&pinned, &peers, &peers).expect("a floor");
assert_eq!(floor.selected_workflow_edition, "2026.06");
assert_eq!(floor.desired_workflow_edition, "2026.06");
assert_eq!(floor.runtime_cap, None);
assert!(!floor.runtime_limited());
}
#[test]
fn this_build_declares_its_current_edition_by_default() {
use super::{core_handshake_request, declared_workflow_edition, edition_sort_key};
assert_eq!(
SUPPORTED_WORKFLOW_EDITIONS
.iter()
.max_by_key(|edition| edition_sort_key(edition))
.copied(),
Some(CURRENT_WORKFLOW_EDITION)
);
assert_eq!(declared_workflow_edition(None), CURRENT_WORKFLOW_EDITION);
assert_eq!(
declared_workflow_edition(Some(" ")),
CURRENT_WORKFLOW_EDITION
);
assert_eq!(declared_workflow_edition(Some(" 2026.06 ")), "2026.06");
let offer = SessionOffer::this_build(None);
assert_eq!(offer.editions, supported_workflow_editions());
assert_eq!(
offer.preferred.as_deref(),
Some(CURRENT_WORKFLOW_EDITION),
"an undeclared build wants its current edition"
);
let request = core_handshake_request("rlmesh-env", &[], None);
assert_eq!(request.preferred_workflow_edition, CURRENT_WORKFLOW_EDITION);
assert_eq!(
core_handshake_request("rlmesh-env", &[], Some("2026.06")).preferred_workflow_edition,
"2026.06"
);
}
}