#![allow(clippy::needless_range_loop)]
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use cudarc::driver::{CudaEvent, CudaSlice};
use crate::Engine;
pub const ROOT: usize = 0;
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum TpTransport {
HostCanonical,
PeerPull,
}
impl TpTransport {
pub fn name(self) -> &'static str {
match self {
Self::HostCanonical => "host-canonical",
Self::PeerPull => "peer-pull",
}
}
}
pub const TRANSPORT_ENV: &str = "MEMRA_TP_TRANSPORT";
pub const TRANSPORT_ENV_GLM5: &str = "MEMRA_GLM5_TP_TRANSPORT";
pub fn parse_transport(flag: &str, value: Option<&str>) -> Result<TpTransport, String> {
match value {
None | Some("") | Some("0") | Some("host-canonical") => Ok(TpTransport::HostCanonical),
Some("1") | Some("peer-pull") => Ok(TpTransport::PeerPull),
Some(other) => Err(format!(
"{flag}={other:?} is not a known transport \
(host-canonical | 0 = the v1 host-staged arm, the default; peer-pull | 1 = the \
consumer-issued device peer copy arm)"
)),
}
}
pub fn resolve_transport(
general: Option<&str>,
alias: Option<&str>,
) -> Result<(TpTransport, &'static str), String> {
let (armed, value) = match (general, alias) {
(Some(g), Some(a)) if g != a => {
return Err(format!(
"{TRANSPORT_ENV}={g:?} and {TRANSPORT_ENV_GLM5}={a:?} disagree — the alias and \
the general flag name ONE transport (unset one); refused rather than silently \
picking a precedence winner"
));
}
(Some(g), _) => (TRANSPORT_ENV, Some(g)),
(None, Some(a)) => (TRANSPORT_ENV_GLM5, Some(a)),
(None, None) => (TRANSPORT_ENV, None),
};
Ok((parse_transport(armed, value)?, armed))
}
pub fn set_transport_names() -> Vec<&'static str> {
[TRANSPORT_ENV, TRANSPORT_ENV_GLM5]
.into_iter()
.filter(|k| std::env::var_os(k).is_some_and(|v| !v.is_empty()))
.collect()
}
fn rollback_advice() -> String {
let names = set_transport_names();
if names.is_empty() {
return format!("{TRANSPORT_ENV}=0");
}
names
.iter()
.map(|n| format!("{n}=0"))
.collect::<Vec<_>>()
.join(" ")
}
pub fn transport_env() -> Result<(TpTransport, &'static str), String> {
resolve_transport(
std::env::var(TRANSPORT_ENV).ok().as_deref(),
std::env::var(TRANSPORT_ENV_GLM5).ok().as_deref(),
)
}
pub static TP_HOST_LEGS: AtomicU64 = AtomicU64::new(0);
pub static TP_HOST_SYNCS: AtomicU64 = AtomicU64::new(0);
pub static TP_PEER_PULLS: AtomicU64 = AtomicU64::new(0);
pub static TP_PUB_EVENTS: AtomicU64 = AtomicU64::new(0);
pub static TP_LOCAL_COPIES: AtomicU64 = AtomicU64::new(0);
pub static TP_XFER_BYTES: AtomicU64 = AtomicU64::new(0);
macro_rules! snapshot_fns {
($($name:ident => $counter:ident),* $(,)?) => {
$(
pub fn $name() -> u64 { $counter.load(Ordering::Relaxed) }
)*
};
}
snapshot_fns! {
tp_host_legs => TP_HOST_LEGS,
tp_host_syncs => TP_HOST_SYNCS,
tp_peer_pulls => TP_PEER_PULLS,
tp_pub_events => TP_PUB_EVENTS,
tp_local_copies => TP_LOCAL_COPIES,
tp_xfer_bytes => TP_XFER_BYTES,
}
pub fn transport_census_line(tag: &str, transport: TpTransport) -> String {
format!(
"[{tag}] census transport={} host_legs={} host_syncs={} peer_pulls={} \
pub_events={} local_copies={} xfer_bytes={}",
transport.name(),
tp_host_legs(),
tp_host_syncs(),
tp_peer_pulls(),
tp_pub_events(),
tp_local_copies(),
tp_xfer_bytes(),
)
}
fn charge_host_leg(bytes: usize, sync: bool) {
TP_HOST_LEGS.fetch_add(1, Ordering::Relaxed);
if sync {
TP_HOST_SYNCS.fetch_add(1, Ordering::Relaxed);
}
TP_XFER_BYTES.fetch_add(bytes as u64, Ordering::Relaxed);
}
fn charge_peer_pull(bytes: usize) {
TP_PEER_PULLS.fetch_add(1, Ordering::Relaxed);
TP_XFER_BYTES.fetch_add(bytes as u64, Ordering::Relaxed);
}
fn charge_local_copy() {
TP_LOCAL_COPIES.fetch_add(1, Ordering::Relaxed);
}
pub struct PeerPullLink {
pub_ev: Vec<CudaEvent>,
rel_ev: Vec<CudaEvent>,
}
impl PeerPullLink {
pub fn new(engines: &[&Engine]) -> Result<Self, Box<dyn std::error::Error>> {
let mut pub_ev = Vec::with_capacity(engines.len());
let mut rel_ev = Vec::with_capacity(engines.len());
for e in engines {
pub_ev.push(e.ctx().new_event(None)?);
rel_ev.push(e.ctx().new_event(None)?);
}
Ok(Self { pub_ev, rel_ev })
}
}
pub struct Hop<'a> {
pub engines: Vec<&'a Engine>,
pub transport: TpTransport,
pub link: Option<&'a PeerPullLink>,
}
impl<'a> Hop<'a> {
pub fn ranks(&self) -> usize {
self.engines.len()
}
pub fn engine(&self, r: usize) -> &'a Engine {
self.engines[r]
}
fn link(&self) -> Result<&'a PeerPullLink, Box<dyn std::error::Error>> {
self.link.ok_or_else(|| {
"glm5-tp transport: the peer-pull arm is armed without a publication link \
(runtime construction bug — refused rather than issuing an unordered peer read)"
.into()
})
}
fn publish(&self, producer: usize, consumer: usize) -> Result<(), Box<dyn std::error::Error>> {
let link = self.link()?;
let ev = &link.pub_ev[producer];
{
let prod = self.engine(producer);
let _main = prod.gpu.enter_main()?;
ev.record(&prod.stream())?;
}
{
let cons = self.engine(consumer);
let _main = cons.gpu.enter_main()?;
cons.stream().wait(ev)?;
}
TP_PUB_EVENTS.fetch_add(2, Ordering::Relaxed);
Ok(())
}
fn release(&self, producer: usize, consumer: usize) -> Result<(), Box<dyn std::error::Error>> {
let link = self.link()?;
let ev = &link.rel_ev[consumer];
{
let cons = self.engine(consumer);
let _main = cons.gpu.enter_main()?;
ev.record(&cons.stream())?;
}
{
let prod = self.engine(producer);
let _main = prod.gpu.enter_main()?;
prod.stream().wait(ev)?;
}
TP_PUB_EVENTS.fetch_add(2, Ordering::Relaxed);
Ok(())
}
#[allow(clippy::too_many_arguments)] fn pull_f32(
&self,
producer: usize,
src: &CudaSlice<f32>,
src_off: usize,
consumer: usize,
dst: &mut CudaSlice<f32>,
dst_off: usize,
n: usize,
) -> Result<(), Box<dyn std::error::Error>> {
if src.len() < src_off + n || dst.len() < dst_off + n {
return Err(format!(
"glm5-tp peer pull geometry: src {}..{} of {} -> dst {}..{} of {}",
src_off,
src_off + n,
src.len(),
dst_off,
dst_off + n,
dst.len(),
)
.into());
}
{
let consumer_engine = self.engine(consumer);
let _main = consumer_engine.gpu.enter_main()?;
let mut view = dst.slice_mut(dst_off..dst_off + n);
consumer_engine
.stream()
.memcpy_dtod(&src.slice(src_off..src_off + n), &mut view)?;
}
self.release(producer, consumer)?;
charge_peer_pull(n * std::mem::size_of::<f32>());
Ok(())
}
}
pub fn fanout_f32(
hop: &Hop<'_>,
src: &CudaSlice<f32>,
n: usize,
) -> Result<Vec<CudaSlice<f32>>, Box<dyn std::error::Error>> {
let bytes = n * std::mem::size_of::<f32>();
let mut out = Vec::with_capacity(hop.ranks() - 1);
match hop.transport {
TpTransport::HostCanonical => {
let host = hop.engine(ROOT).dtoh_view(&src.slice(0..n))?;
charge_host_leg(bytes, true);
for to in 1..hop.ranks() {
out.push(hop.engine(to).htod(&host)?);
charge_host_leg(bytes, false);
}
}
TpTransport::PeerPull => {
for to in 1..hop.ranks() {
let mut buf = hop.engine(to).uninit(n)?;
hop.publish(ROOT, to)?;
hop.pull_f32(ROOT, src, 0, to, &mut buf, 0, n)?;
out.push(buf);
}
}
}
Ok(out)
}
pub fn fanout_f32_to(
hop: &Hop<'_>,
to: usize,
src: &CudaSlice<f32>,
n: usize,
) -> Result<CudaSlice<f32>, Box<dyn std::error::Error>> {
let bytes = n * std::mem::size_of::<f32>();
match hop.transport {
TpTransport::HostCanonical => {
let host = hop.engine(ROOT).dtoh_view(&src.slice(0..n))?;
charge_host_leg(bytes, true);
let out = hop.engine(to).htod(&host)?;
charge_host_leg(bytes, false);
Ok(out)
}
TpTransport::PeerPull => {
let mut buf = hop.engine(to).uninit(n)?;
hop.publish(ROOT, to)?;
hop.pull_f32(ROOT, src, 0, to, &mut buf, 0, n)?;
Ok(buf)
}
}
}
pub fn fanout_i32(
hop: &Hop<'_>,
src: &CudaSlice<i32>,
n: usize,
) -> Result<Vec<CudaSlice<i32>>, Box<dyn std::error::Error>> {
let bytes = n * std::mem::size_of::<i32>();
let mut out = Vec::with_capacity(hop.ranks() - 1);
match hop.transport {
TpTransport::HostCanonical => {
let host = hop.engine(ROOT).dtoh_i32(src)?;
charge_host_leg(bytes, true);
for to in 1..hop.ranks() {
out.push(hop.engine(to).htod_i32(&host)?);
charge_host_leg(bytes, false);
}
}
TpTransport::PeerPull => {
for to in 1..hop.ranks() {
let mut buf = hop.engine(to).uninit_i32(n)?;
hop.publish(ROOT, to)?;
{
let consumer = hop.engine(to);
let _main = consumer.gpu.enter_main()?;
let mut view = buf.slice_mut(0..n);
consumer.stream().memcpy_dtod(&src.slice(0..n), &mut view)?;
}
hop.release(ROOT, to)?;
charge_peer_pull(bytes);
out.push(buf);
}
}
}
Ok(out)
}
pub fn gather_parts(
hop: &Hop<'_>,
parts: &[&CudaSlice<f32>],
t: usize,
part: usize,
) -> Result<Vec<CudaSlice<f32>>, Box<dyn std::error::Error>> {
let ranks = hop.ranks();
if t == 0 || part == 0 {
return Err("glm5-tp gather: zero geometry".into());
}
if parts.len() != ranks {
return Err(format!("glm5-tp gather: {} parts for {ranks} ranks", parts.len()).into());
}
let full = ranks * part;
let span = t * part;
for (r, p) in parts.iter().enumerate() {
if p.len() < span {
return Err(format!(
"glm5-tp gather geometry: part[{r}] {} for t={t} part={part}",
p.len()
)
.into());
}
}
let bytes = span * std::mem::size_of::<f32>();
match hop.transport {
TpTransport::HostCanonical => {
let mut hosts: Vec<Option<Vec<f32>>> = (0..ranks).map(|_| None).collect();
for r in (0..ranks).rev() {
hosts[r] = Some(hop.engine(r).dtoh_view(&parts[r].slice(0..span))?);
charge_host_leg(bytes, true);
}
let mut full_host = vec![0f32; t * full];
for (r, h) in hosts.iter().enumerate() {
let h = h.as_ref().expect("drained above");
for tok in 0..t {
full_host[tok * full + r * part..tok * full + (r + 1) * part]
.copy_from_slice(&h[tok * part..(tok + 1) * part]);
}
}
let mut out = Vec::with_capacity(ranks);
for r in 0..ranks {
out.push(hop.engine(r).htod(&full_host)?);
charge_host_leg(t * full * std::mem::size_of::<f32>(), false);
}
Ok(out)
}
TpTransport::PeerPull => {
for s in 0..ranks {
for d in 0..ranks {
if s != d {
hop.publish(s, d)?;
}
}
}
let mut out = Vec::with_capacity(ranks);
for d in 0..ranks {
let mut full_d = hop.engine(d).uninit(t * full)?;
if t == 1 {
for s in 0..ranks {
if s == d {
hop.engine(d).copy_range_into(
&mut full_d,
s * part,
parts[s],
0,
part,
)?;
charge_local_copy();
} else {
hop.pull_f32(s, parts[s], 0, d, &mut full_d, s * part, part)?;
}
}
} else {
for s in 0..ranks {
if s == d {
hop.engine(d).place_rows_strided(
parts[s],
&mut full_d,
part,
t,
full,
s * part,
)?;
charge_local_copy();
} else {
let mut foreign = hop.engine(d).uninit(span)?;
hop.pull_f32(s, parts[s], 0, d, &mut foreign, 0, span)?;
hop.engine(d).place_rows_strided(
&foreign,
&mut full_d,
part,
t,
full,
s * part,
)?;
charge_local_copy();
}
}
}
out.push(full_d);
}
Ok(out)
}
}
}
pub fn concat_parts_on_root(
hop: &Hop<'_>,
parts: &[&CudaSlice<f32>],
t: usize,
part: usize,
) -> Result<CudaSlice<f32>, Box<dyn std::error::Error>> {
let ranks = hop.ranks();
if t == 0 || part == 0 {
return Err("glm5-tp concat: zero geometry".into());
}
if parts.len() != ranks {
return Err(format!("glm5-tp concat: {} parts for {ranks} ranks", parts.len()).into());
}
let full = ranks * part;
let span = t * part;
for (r, p) in parts.iter().enumerate() {
if p.len() < span {
return Err(format!(
"glm5-tp concat geometry: part[{r}] {} for t={t} part={part}",
p.len()
)
.into());
}
}
let bytes = span * std::mem::size_of::<f32>();
match hop.transport {
TpTransport::HostCanonical => {
let mut hosts: Vec<Option<Vec<f32>>> = (0..ranks).map(|_| None).collect();
for r in (0..ranks).rev() {
hosts[r] = Some(hop.engine(r).dtoh_view(&parts[r].slice(0..span))?);
charge_host_leg(bytes, true);
}
let mut out_host = vec![0f32; t * full];
for (r, h) in hosts.iter().enumerate() {
let h = h.as_ref().expect("drained above");
for tok in 0..t {
out_host[tok * full + r * part..tok * full + (r + 1) * part]
.copy_from_slice(&h[tok * part..(tok + 1) * part]);
}
}
let out = hop.engine(ROOT).htod(&out_host)?;
charge_host_leg(t * full * std::mem::size_of::<f32>(), false);
Ok(out)
}
TpTransport::PeerPull => {
let mut out = hop.engine(ROOT).uninit(t * full)?;
for s in 1..ranks {
hop.publish(s, ROOT)?;
}
if t == 1 {
for s in 0..ranks {
if s == ROOT {
hop.engine(ROOT)
.copy_range_into(&mut out, s * part, parts[s], 0, part)?;
charge_local_copy();
} else {
hop.pull_f32(s, parts[s], 0, ROOT, &mut out, s * part, part)?;
}
}
} else {
for s in 0..ranks {
if s == ROOT {
hop.engine(ROOT).place_rows_strided(
parts[s],
&mut out,
part,
t,
full,
s * part,
)?;
charge_local_copy();
} else {
let mut foreign = hop.engine(ROOT).uninit(span)?;
hop.pull_f32(s, parts[s], 0, ROOT, &mut foreign, 0, span)?;
hop.engine(ROOT).place_rows_strided(
&foreign,
&mut out,
part,
t,
full,
s * part,
)?;
charge_local_copy();
}
}
}
Ok(out)
}
}
}
pub fn host_stage_block(
hop: &Hop<'_>,
from: usize,
src: &CudaSlice<f32>,
n: usize,
) -> Result<Vec<f32>, Box<dyn std::error::Error>> {
let host = hop.engine(from).dtoh_view(&src.slice(0..n))?;
charge_host_leg(n * std::mem::size_of::<f32>(), true);
Ok(host)
}
pub fn host_row_to(
hop: &Hop<'_>,
to: usize,
row: &[f32],
) -> Result<CudaSlice<f32>, Box<dyn std::error::Error>> {
let out = hop.engine(to).htod(row)?;
charge_host_leg(std::mem::size_of_val(row), false);
Ok(out)
}
pub fn return_block_to_root(
hop: &Hop<'_>,
from: usize,
blk: &CudaSlice<f32>,
dst: &mut CudaSlice<f32>,
dst_off: usize,
n: usize,
) -> Result<(), Box<dyn std::error::Error>> {
let bytes = n * std::mem::size_of::<f32>();
match hop.transport {
TpTransport::HostCanonical => {
let host = hop.engine(from).dtoh_view(&blk.slice(0..n))?;
charge_host_leg(bytes, true);
hop.engine(ROOT).htod_f32_into_at(&host, dst, dst_off)?;
charge_host_leg(bytes, false);
Ok(())
}
TpTransport::PeerPull => {
hop.publish(from, ROOT)?;
hop.pull_f32(from, blk, 0, ROOT, dst, dst_off, n)
}
}
}
pub fn return_row_to_root(
hop: &Hop<'_>,
from: usize,
y: &CudaSlice<f32>,
n: usize,
) -> Result<CudaSlice<f32>, Box<dyn std::error::Error>> {
match hop.transport {
TpTransport::HostCanonical => {
let host = hop.engine(from).dtoh_view(&y.slice(0..n))?;
charge_host_leg(n * std::mem::size_of::<f32>(), true);
let out = hop.engine(ROOT).htod(&host)?;
charge_host_leg(n * std::mem::size_of::<f32>(), false);
Ok(out)
}
TpTransport::PeerPull => {
let mut out = hop.engine(ROOT).uninit(n)?;
hop.publish(from, ROOT)?;
hop.pull_f32(from, y, 0, ROOT, &mut out, 0, n)?;
Ok(out)
}
}
}
const PULL_PROBE_WORDS: &[usize] = &[4_096, 16_384, 262_144, 16_777_216];
static TRANSPORT_MARKED: AtomicBool = AtomicBool::new(false);
pub fn arm_transport(
transport: TpTransport,
armed_flag: &str,
tag: &str,
engines: &[&Engine],
same_device_gate: bool,
) -> Result<Option<PeerPullLink>, Box<dyn std::error::Error>> {
if transport == TpTransport::HostCanonical {
announce(
tag,
transport,
same_device_gate,
"host staging, no peer mapping",
);
return Ok(None);
}
if same_device_gate {
eprintln!(
"[{tag}] same-device gate: peer-access grant SKIPPED (one device, {} \
contexts); the ladder below proves bit-preservation only, never fabric engagement",
engines.len(),
);
} else {
for (i, a) in engines.iter().enumerate() {
for (j, b) in engines.iter().enumerate() {
if i != j {
crate::tp::grant_peer_access(a, b, &format!("{armed_flag}=peer-pull"))?;
}
}
}
}
let link = PeerPullLink::new(engines)?;
let hop = Hop {
engines: engines.to_vec(),
transport,
link: Some(&link),
};
let census_before = (
TP_PEER_PULLS.load(Ordering::Relaxed),
TP_PUB_EVENTS.load(Ordering::Relaxed),
TP_XFER_BYTES.load(Ordering::Relaxed),
TP_HOST_LEGS.load(Ordering::Relaxed),
TP_HOST_SYNCS.load(Ordering::Relaxed),
TP_LOCAL_COPIES.load(Ordering::Relaxed),
);
let ranks = engines.len();
for &words in PULL_PROBE_WORDS {
for producer in 0..ranks {
for consumer in 0..ranks {
if consumer == producer {
continue;
}
let expected: Vec<f32> = (0..words)
.map(|i| {
f32::from_bits(
(i as u32)
.wrapping_mul(0x9e37_79b9)
.wrapping_add(((producer as u32) + 1) << 16)
& 0x7f7f_ffff,
)
})
.collect();
let poison: Vec<f32> = expected.iter().map(|v| -*v - 1.0).collect();
let src = hop.engine(producer).htod(&expected)?;
let mut dst = hop.engine(consumer).htod(&poison)?;
hop.publish(producer, consumer)?;
hop.pull_f32(producer, &src, 0, consumer, &mut dst, 0, words)?;
let actual = hop.engine(consumer).dtoh(&dst)?;
let mismatches = actual
.iter()
.zip(&expected)
.filter(|(a, b)| a.to_bits() != b.to_bits())
.count();
if mismatches != 0 {
return Err(format!(
"{armed_flag}=peer-pull byte-integrity ladder FAILED: \
rank{producer}->rank{consumer} at {} bytes, {mismatches}/{words} \
words differ (refused before any layer was sharded; roll back \
with: {})",
words * std::mem::size_of::<f32>(),
rollback_advice(),
)
.into());
}
}
}
}
let restore = |counter: &AtomicU64, before: u64| {
let now = counter.load(Ordering::Relaxed);
counter.fetch_sub(now.saturating_sub(before), Ordering::Relaxed);
};
restore(&TP_PEER_PULLS, census_before.0);
restore(&TP_PUB_EVENTS, census_before.1);
restore(&TP_XFER_BYTES, census_before.2);
restore(&TP_HOST_LEGS, census_before.3);
restore(&TP_HOST_SYNCS, census_before.4);
restore(&TP_LOCAL_COPIES, census_before.5);
eprintln!(
"[{tag}] peer-pull byte-integrity ladder PASS: directions={} \
byte_ladder={:?} mismatches=0 same_device_gate={same_device_gate} \
census_excluded=arm-time-ladder-traffic",
ranks * (ranks - 1),
PULL_PROBE_WORDS
.iter()
.map(|w| w * std::mem::size_of::<f32>())
.collect::<Vec<_>>(),
);
announce(
tag,
transport,
same_device_gate,
"consumer-issued cuMemcpyDtoDAsync, event-published, atomics-free",
);
Ok(Some(link))
}
fn announce(tag: &str, transport: TpTransport, same_device_gate: bool, how: &str) {
if TRANSPORT_MARKED.swap(true, Ordering::Relaxed) {
return;
}
eprintln!(
"[{tag}] armed transport={} shape={how} same_device_gate={same_device_gate} \
performance_claim=false",
transport.name(),
);
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn transport_parse_is_literal_and_fail_closed() {
for off in [None, Some(""), Some("0"), Some("host-canonical")] {
assert_eq!(
parse_transport(TRANSPORT_ENV, off).unwrap(),
TpTransport::HostCanonical
);
}
for on in [Some("1"), Some("peer-pull")] {
assert_eq!(
parse_transport(TRANSPORT_ENV, on).unwrap(),
TpTransport::PeerPull
);
}
for bad in ["peer_pull", "peerpull", "p2p", "native", "2", "on", "true"] {
for flag in [TRANSPORT_ENV, TRANSPORT_ENV_GLM5] {
let err =
parse_transport(flag, Some(bad)).expect_err("unknown transport must refuse");
assert!(err.contains(flag), "{err}");
assert!(err.contains(bad), "{err}");
assert!(err.contains("host-canonical"), "{err}");
assert!(err.contains("peer-pull"), "{err}");
}
}
}
#[test]
fn alias_resolution_honors_both_names_and_refuses_a_disagreeing_pair() {
assert_eq!(
resolve_transport(None, None).unwrap(),
(TpTransport::HostCanonical, TRANSPORT_ENV)
);
assert_eq!(
resolve_transport(Some("peer-pull"), None).unwrap(),
(TpTransport::PeerPull, TRANSPORT_ENV)
);
assert_eq!(
resolve_transport(None, Some("peer-pull")).unwrap(),
(TpTransport::PeerPull, TRANSPORT_ENV_GLM5)
);
assert_eq!(
resolve_transport(Some("1"), Some("1")).unwrap(),
(TpTransport::PeerPull, TRANSPORT_ENV)
);
for (g, a) in [("peer-pull", "0"), ("0", "1"), ("1", "host-canonical")] {
let err = resolve_transport(Some(g), Some(a))
.expect_err("a disagreeing pair must refuse the load");
assert!(err.contains(TRANSPORT_ENV), "{err}");
assert!(err.contains(TRANSPORT_ENV_GLM5), "{err}");
assert!(err.contains("disagree"), "{err}");
}
assert!(resolve_transport(Some("0"), Some("host-canonical")).is_err());
}
#[test]
fn default_is_the_v1_arm() {
assert_eq!(
parse_transport(TRANSPORT_ENV, None).unwrap(),
TpTransport::HostCanonical
);
assert_eq!(TpTransport::HostCanonical.name(), "host-canonical");
assert_eq!(TpTransport::PeerPull.name(), "peer-pull");
}
#[test]
fn census_line_names_every_counter() {
let line = transport_census_line("glm5-tp-transport", TpTransport::PeerPull);
for field in [
"transport=peer-pull",
"host_legs=",
"host_syncs=",
"peer_pulls=",
"pub_events=",
"local_copies=",
"xfer_bytes=",
] {
assert!(line.contains(field), "{line} is missing {field}");
}
}
}