use std::collections::BTreeMap;
use std::fmt;
use anyhow::{bail, Result};
use super::ingress::{IngressPlan, IngressRule};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DiscoveredRecord {
pub ident: String,
pub mesh_ip: String,
pub ports: Vec<u16>,
pub named_ports: BTreeMap<String, u16>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum UnknownReason {
Undeclared,
Unreachable(String),
EndpointAbsent(String),
BadAnswer(String),
}
impl fmt::Display for UnknownReason {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Undeclared => write!(
f,
"is not declared in .yah/infra/machines/, so there is no address to ask"
),
Self::Unreachable(detail) => write!(f, "did not answer ({detail})"),
Self::EndpointAbsent(detail) => write!(
f,
"answered 404 — its yubaba predates GET /service-records, so it CANNOT answer, \
which is not the same as having no records ({detail})"
),
Self::BadAnswer(detail) => write!(f, "answered unusably ({detail})"),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RecordVisibility {
Answered(Vec<DiscoveredRecord>),
Unknown(UnknownReason),
}
impl RecordVisibility {
pub fn as_str(&self) -> &'static str {
match self {
Self::Answered(_) => "answered",
Self::Unknown(_) => "unknown",
}
}
pub fn records(&self) -> &[DiscoveredRecord] {
match self {
Self::Answered(records) => records,
Self::Unknown(_) => &[],
}
}
pub fn unknown_reason(&self) -> Option<&UnknownReason> {
match self {
Self::Answered(_) => None,
Self::Unknown(reason) => Some(reason),
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ServiceRecordFanout {
pub nodes: BTreeMap<String, RecordVisibility>,
}
impl ServiceRecordFanout {
pub const CHANNEL: &'static str = "service-records";
pub fn push_answer(&mut self, machine: impl Into<String>, records: Vec<DiscoveredRecord>) {
self.nodes
.insert(machine.into(), RecordVisibility::Answered(records));
}
pub fn push_unknown(&mut self, machine: impl Into<String>, reason: UnknownReason) {
self.nodes
.insert(machine.into(), RecordVisibility::Unknown(reason));
}
pub fn asked(&self) -> usize {
self.nodes.len()
}
pub fn unknown(&self) -> impl Iterator<Item = (&str, &UnknownReason)> {
self.nodes.iter().filter_map(|(machine, visibility)| {
visibility
.unknown_reason()
.map(|reason| (machine.as_str(), reason))
})
}
pub fn is_partial(&self) -> bool {
self.unknown().next().is_some()
}
pub fn records(&self) -> impl Iterator<Item = (&str, &DiscoveredRecord)> {
self.nodes.iter().flat_map(|(machine, visibility)| {
visibility
.records()
.iter()
.map(move |r| (machine.as_str(), r))
})
}
pub fn answered_by(&self, machine: &str) -> bool {
matches!(
self.nodes.get(machine),
Some(RecordVisibility::Answered(_))
)
}
pub fn unknown_reason(&self, machine: &str) -> Option<&UnknownReason> {
self.nodes.get(machine)?.unknown_reason()
}
pub fn upstreams_for(&self, rule: &IngressRule) -> Vec<String> {
let mut out: Vec<String> = Vec::new();
let Some(port) = rule.port else {
return out;
};
for (_, visibility) in self
.nodes
.iter()
.filter(|(machine, _)| self.in_scope(rule, machine))
{
for record in visibility.records() {
if record.ports.contains(&port) && !out.contains(&record.mesh_ip) {
out.push(record.mesh_ip.clone());
}
}
}
out
}
pub fn port_for(&self, rule: &IngressRule, ident: &str) -> Option<u16> {
let mut candidates: Vec<u16> = Vec::new();
let mut serving: Vec<u16> = Vec::new();
for (_, visibility) in self
.nodes
.iter()
.filter(|(machine, _)| self.in_scope(rule, machine))
{
for record in visibility.records().iter().filter(|r| r.ident == ident) {
for port in &record.ports {
if !candidates.contains(port) {
candidates.push(*port);
}
}
if let Some(port) = record.named_ports.get(kamaji::DEFAULT_PORT_NAME) {
if !serving.contains(port) {
serving.push(*port);
}
}
}
}
match (serving.as_slice(), candidates.as_slice()) {
([one], _) => Some(*one),
(_, [one]) => Some(*one),
_ => None,
}
}
pub fn address_for_ident(&self, ident: &str) -> Option<String> {
let mut found: Vec<String> = Vec::new();
for (_, visibility) in self.nodes.iter() {
for record in visibility.records().iter().filter(|r| r.ident == ident) {
let port = record
.named_ports
.get(kamaji::DEFAULT_PORT_NAME)
.copied()
.or(match record.ports.as_slice() {
[one] => Some(*one),
_ => None,
});
let Some(port) = port else { continue };
let addr = format!("{}:{port}", record.mesh_ip);
if !found.contains(&addr) {
found.push(addr);
}
}
}
match found.as_slice() {
[one] => Some(one.clone()),
_ => None,
}
}
fn in_scope(&self, rule: &IngressRule, machine: &str) -> bool {
rule.machines.is_empty() || rule.machines.iter().any(|m| m == machine)
}
pub fn unknown_note(&self) -> Option<String> {
let list = self
.unknown()
.map(|(machine, reason)| format!("{machine} {reason}"))
.collect::<Vec<_>>();
if list.is_empty() {
return None;
}
Some(format!(
"{} of {} node(s) on the {} channel are `unknown`, so this read is a LOWER BOUND on \
the fleet, not a description of it: {}",
list.len(),
self.asked(),
Self::CHANNEL,
list.join("; ")
))
}
fn uncertainty_for(&self, rule: &IngressRule) -> Option<String> {
if !rule.machines.is_empty() {
for scope in &rule.machines {
if let Some(reason) = self.unknown_reason(scope) {
return Some(format!("the node this slot is placed on, {scope}, {reason}"));
}
}
if rule.machines.iter().all(|m| self.answered_by(m)) {
return None;
}
}
self.unknown_note()
}
}
impl IngressPlan {
pub fn resolve_upstreams_from(&mut self, fanout: &ServiceRecordFanout) -> Result<()> {
for rule in &mut self.rules {
if !rule.upstream_hosts.is_empty() {
continue;
}
rule.upstream_hosts = fanout.upstreams_for(rule);
if rule.upstream_hosts.is_empty() {
if let Some(uncertainty) = fanout.uncertainty_for(rule) {
bail!(
"slot [providers.{}] fronted at {} resolved no upstream for port {}, but \
the discovery read was PARTIAL and cannot tell that apart from a record \
it never got to see: {uncertainty}. Absence here is UNKNOWN, not empty — \
re-run once those node(s) answer, or pin `upstream_host` on the slot if \
a node is expected to stay dark.",
rule.slot,
rule.hostname,
rule.port_label()
);
}
}
rule.upstreams()?;
}
Ok(())
}
pub fn resolve_ports_from(&mut self, fanout: &ServiceRecordFanout, ident: &str) {
for rule in &mut self.rules {
if rule.port.is_some() {
continue;
}
rule.port = fanout.port_for(rule, ident);
}
}
pub fn workload_machines(&self) -> Vec<&str> {
let mut out: Vec<&str> = Vec::new();
for machine in self.rules.iter().flat_map(|r| r.machines.iter()) {
if !out.contains(&machine.as_str()) {
out.push(machine.as_str());
}
}
out
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::IngressProvider;
fn record(ip: &str, ports: &[u16]) -> DiscoveredRecord {
named_record("yah-marketing", ip, ports)
}
fn named_record(ident: &str, ip: &str, ports: &[u16]) -> DiscoveredRecord {
DiscoveredRecord {
ident: ident.into(),
mesh_ip: ip.into(),
ports: ports.to_vec(),
named_ports: kamaji::name_anonymous_ports(ports),
}
}
fn named_port_record(ident: &str, ip: &str, ports: &[(&str, u16)]) -> DiscoveredRecord {
DiscoveredRecord {
ident: ident.into(),
mesh_ip: ip.into(),
ports: ports.iter().map(|(_, p)| *p).collect(),
named_ports: ports
.iter()
.map(|(n, p)| ((*n).to_string(), *p))
.collect(),
}
}
fn rule(hostname: &str, port: u16, machines: &[&str]) -> IngressRule {
IngressRule {
hostname: hostname.into(),
port: Some(port),
slot: "compute".into(),
provider_id: None,
machines: machines.iter().map(|m| m.to_string()).collect(),
upstream_hosts: Vec::new(),
}
}
fn plan(rules: Vec<IngressRule>) -> IngressPlan {
IngressPlan {
provider: IngressProvider::Passway,
rules,
front_doors: Vec::new(),
tunnel_id: None,
edge_provider_id: None,
image: None,
auth: None,
via: None,
behind_tunnel: false,
tunnel_door: None,
}
}
#[test]
fn a_node_reporting_none_is_not_the_same_value_as_a_node_that_could_not_be_seen() {
let mut reported_none = ServiceRecordFanout::default();
reported_none.push_answer("us-east-001", Vec::new());
let mut could_not_see = ServiceRecordFanout::default();
could_not_see.push_unknown(
"us-east-001",
UnknownReason::Unreachable("connection refused".into()),
);
assert_eq!(reported_none.records().count(), 0);
assert_eq!(could_not_see.records().count(), 0);
assert_ne!(reported_none, could_not_see);
assert_eq!(reported_none.nodes["us-east-001"].as_str(), "answered");
assert_eq!(could_not_see.nodes["us-east-001"].as_str(), "unknown");
assert!(!reported_none.is_partial());
assert!(could_not_see.is_partial());
assert!(reported_none.unknown_note().is_none());
assert!(could_not_see.unknown_note().is_some());
}
#[test]
fn an_unseen_node_stays_in_the_map_rather_than_dropping_out_of_it() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![record("100.64.0.3", &[8080])]);
fanout.push_unknown(
"us-south-001",
UnknownReason::Unreachable("no route to host".into()),
);
assert_eq!(fanout.asked(), 2);
assert_eq!(
fanout.nodes.keys().collect::<Vec<_>>(),
vec!["us-east-001", "us-south-001"]
);
assert_eq!(fanout.nodes["us-south-001"].as_str(), "unknown");
assert!(fanout.nodes["us-south-001"].records().is_empty());
}
#[test]
fn the_three_live_cases_are_three_distinct_values() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![record("100.64.0.3", &[8080])]);
fanout.push_unknown(
"us-south-001",
UnknownReason::Unreachable("connect timed out".into()),
);
fanout.push_unknown(
"us-west-015",
UnknownReason::EndpointAbsent("GET /service-records returned 404".into()),
);
assert_eq!(fanout.asked(), 3);
assert_eq!(fanout.records().count(), 1);
assert!(matches!(
fanout.unknown_reason("us-south-001"),
Some(UnknownReason::Unreachable(_))
));
assert!(matches!(
fanout.unknown_reason("us-west-015"),
Some(UnknownReason::EndpointAbsent(_))
));
assert_eq!(fanout.unknown_reason("us-east-001"), None);
assert_ne!(
fanout.nodes["us-south-001"],
fanout.nodes["us-west-015"],
"an unreachable node and a node that cannot answer are different facts"
);
}
#[test]
fn a_404_is_cannot_answer_not_no_records() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_unknown(
"us-west-015",
UnknownReason::EndpointAbsent("GET /service-records?ready=true returned 404".into()),
);
assert!(fanout.is_partial());
let note = fanout.unknown_note().unwrap();
assert!(note.contains("us-west-015"), "got: {note}");
assert!(note.contains("CANNOT answer"), "got: {note}");
assert!(note.contains("predates"), "got: {note}");
}
#[test]
fn the_note_names_the_channel_every_unseen_node_and_its_own_reason() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![record("100.64.0.3", &[8080])]);
fanout.push_unknown("us-south-001", UnknownReason::Unreachable("timed out".into()));
fanout.push_unknown("us-west-015", UnknownReason::EndpointAbsent("404".into()));
fanout.push_unknown("typo-node", UnknownReason::Undeclared);
assert_eq!(fanout.asked(), 4);
let note = fanout.unknown_note().unwrap();
assert!(note.contains("3 of 4 node(s)"), "got: {note}");
assert!(note.contains("service-records channel"), "got: {note}");
for name in ["us-south-001", "us-west-015", "typo-node"] {
assert!(note.contains(name), "got: {note}");
}
assert!(note.contains("timed out"), "got: {note}");
assert!(note.contains(".yah/infra/machines/"), "got: {note}");
assert!(!note.contains("us-east-001"), "got: {note}");
}
#[test]
fn a_record_registered_on_one_node_resolves_through_a_read_that_asked_several() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![record("100.64.0.3", &[8080])]);
fanout.push_answer("us-west-001", vec![record("100.64.0.9", &[9090])]);
let mut p = plan(vec![
rule("a.yah.dev", 8080, &[]),
rule("b.yah.dev", 9090, &[]),
]);
p.resolve_upstreams_from(&fanout).unwrap();
assert_eq!(
p.passway_upstreams().unwrap(),
vec!["a.yah.dev=100.64.0.3:8080", "b.yah.dev=100.64.0.9:9090"]
);
}
#[test]
fn a_placed_rule_is_answered_only_from_its_own_node() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![record("100.64.0.3", &[8080])]);
fanout.push_answer("us-west-001", vec![record("100.64.0.9", &[8080])]);
let mut p = plan(vec![rule("b.yah.dev", 8080, &["us-west-001"])]);
p.resolve_upstreams_from(&fanout).unwrap();
assert_eq!(
p.passway_upstreams().unwrap(),
vec!["b.yah.dev=100.64.0.9:8080"]
);
}
#[test]
fn workload_machines_is_the_deduplicated_set_the_read_must_ask() {
let p = plan(vec![
rule("a.yah.dev", 8080, &["us-east-001"]),
rule("b.yah.dev", 9090, &["us-west-001"]),
rule("c.yah.dev", 7070, &["us-east-001"]),
rule("d.yah.dev", 6060, &[]),
]);
assert_eq!(p.workload_machines(), vec!["us-east-001", "us-west-001"]);
}
#[test]
fn one_rule_at_scale_two_contributes_both_its_nodes() {
let p = plan(vec![rule("a.yah.dev", 8080, &["us-east-001", "us-west-001"])]);
assert_eq!(p.workload_machines(), vec!["us-east-001", "us-west-001"]);
}
#[test]
fn a_rule_at_scale_two_resolves_every_replica() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![record("100.64.0.3", &[8080])]);
fanout.push_answer("us-west-001", vec![record("100.64.0.9", &[8080])]);
let mut p = plan(vec![rule("a.yah.dev", 8080, &["us-east-001", "us-west-001"])]);
p.resolve_upstreams_from(&fanout).unwrap();
assert_eq!(
p.rules[0].upstream_hosts,
vec!["100.64.0.3".to_string(), "100.64.0.9".to_string()]
);
assert_eq!(
p.passway_upstreams().unwrap(),
vec!["a.yah.dev=100.64.0.3:8080", "a.yah.dev=100.64.0.9:8080"]
);
}
#[test]
fn a_placement_set_still_excludes_a_node_it_does_not_name() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![record("100.64.0.3", &[8080])]);
fanout.push_answer("us-west-001", vec![record("100.64.0.9", &[8080])]);
fanout.push_answer("us-south-001", vec![record("100.64.0.7", &[8080])]);
let mut p = plan(vec![rule("a.yah.dev", 8080, &["us-east-001", "us-west-001"])]);
p.resolve_upstreams_from(&fanout).unwrap();
assert_eq!(
p.rules[0].upstream_hosts,
vec!["100.64.0.3".to_string(), "100.64.0.9".to_string()]
);
}
#[test]
fn one_unseen_member_of_a_placement_set_makes_the_answer_partial() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![record("100.64.0.3", &[8080])]);
fanout.push_unknown("us-west-001", UnknownReason::Unreachable("timed out".into()));
let p = plan(vec![rule("a.yah.dev", 8080, &["us-east-001", "us-west-001"])]);
assert_eq!(
fanout.upstreams_for(&p.rules[0]),
vec!["100.64.0.3".to_string()],
"the seen half still resolves"
);
assert!(
fanout.is_partial(),
"and the read reports itself as a lower bound"
);
assert!(fanout
.unknown_note()
.expect("an unseen node produces a note")
.contains("us-west-001"));
}
#[test]
fn an_unresolved_rule_on_an_unseen_node_says_unknown_not_not_up() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![record("100.64.0.3", &[8080])]);
fanout.push_unknown(
"us-west-002",
UnknownReason::Unreachable("no route to host".into()),
);
let mut p = plan(vec![rule("b.yah.dev", 9090, &["us-west-002"])]);
let err = p.resolve_upstreams_from(&fanout).unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("PARTIAL"), "got: {msg}");
assert!(msg.contains("UNKNOWN, not empty"), "got: {msg}");
assert!(msg.contains("us-west-002"), "got: {msg}");
assert!(msg.contains("no route to host"), "got: {msg}");
assert!(!msg.contains("no resolved upstream address"), "got: {msg}");
}
#[test]
fn an_unscoped_rule_is_unknown_while_any_node_is_unseen() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", Vec::new());
fanout.push_unknown("us-south-001", UnknownReason::Unreachable("timed out".into()));
let mut p = plan(vec![rule("a.yah.dev", 8080, &[])]);
let err = p.resolve_upstreams_from(&fanout).unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("PARTIAL"), "got: {msg}");
assert!(msg.contains("us-south-001"), "got: {msg}");
}
#[test]
fn a_rule_on_a_node_that_answered_still_gets_the_definite_error() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", Vec::new());
fanout.push_unknown("us-west-002", UnknownReason::Unreachable("offline".into()));
let mut p = plan(vec![rule("a.yah.dev", 8080, &["us-east-001"])]);
let err = p.resolve_upstreams_from(&fanout).unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("no resolved upstream address"), "got: {msg}");
assert!(!msg.contains("PARTIAL"), "got: {msg}");
}
#[test]
fn an_undeclared_placement_is_unknown_not_an_empty_answer() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_unknown("typo-node", UnknownReason::Undeclared);
let mut p = plan(vec![rule("a.yah.dev", 8080, &["typo-node"])]);
let err = p.resolve_upstreams_from(&fanout).unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("PARTIAL"), "got: {msg}");
assert!(msg.contains(".yah/infra/machines/"), "got: {msg}");
}
#[test]
fn a_pinned_upstream_survives_a_read_that_saw_nothing() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_unknown("us-west-002", UnknownReason::Unreachable("offline".into()));
let mut pinned = rule("a.yah.dev", 8080, &["us-west-002"]);
pinned.upstream_hosts = vec!["127.0.0.1".into()];
let mut p = plan(vec![pinned]);
p.resolve_upstreams_from(&fanout).unwrap();
assert_eq!(
p.passway_upstreams().unwrap(),
vec!["a.yah.dev=127.0.0.1:8080"]
);
}
#[test]
fn an_empty_read_asks_nobody_and_claims_nothing() {
let fanout = ServiceRecordFanout::default();
assert_eq!(fanout.asked(), 0);
assert!(!fanout.is_partial());
assert!(fanout.unknown_note().is_none());
assert_eq!(fanout.records().count(), 0);
assert_eq!(plan(Vec::new()).workload_machines(), Vec::<&str>::new());
}
fn crowded_node() -> ServiceRecordFanout {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer(
"us-east-001",
vec![
named_record("yah-marketing", "100.64.0.3", &[43117]),
named_record("yah-marketing-revalidate", "100.64.0.3", &[8081]),
named_record("yah-marketing-feed", "100.64.0.3", &[8081]),
],
);
fanout
}
#[test]
fn a_portless_rule_takes_the_port_its_own_record_reports() {
let fanout = crowded_node();
let mut portless = rule("yah.dev", 0, &["us-east-001"]);
portless.port = None;
let mut p = plan(vec![portless]);
p.resolve_ports_from(&fanout, "yah-marketing");
assert_eq!(p.rules[0].port, Some(43117));
p.resolve_upstreams_from(&fanout).unwrap();
assert_eq!(
p.passway_upstreams().unwrap(),
vec!["yah.dev=100.64.0.3:43117"]
);
}
#[test]
fn a_pinned_port_is_never_overwritten_by_discovery() {
let fanout = crowded_node();
let mut p = plan(vec![rule("yah.dev", 8080, &["us-east-001"])]);
p.resolve_ports_from(&fanout, "yah-marketing");
assert_eq!(p.rules[0].port, Some(8080), "an operator pin always wins");
}
#[test]
fn an_ident_nobody_reported_leaves_the_port_unresolved() {
let fanout = crowded_node();
let mut portless = rule("yah.dev", 0, &["us-east-001"]);
portless.port = None;
let mut p = plan(vec![portless]);
p.resolve_ports_from(&fanout, "yah-marketing-preview");
assert_eq!(p.rules[0].port, None);
let msg = format!("{:#}", p.resolve_upstreams_from(&fanout).unwrap_err());
assert!(msg.contains("no resolved port"), "got: {msg}");
assert!(msg.contains("no resolved upstream address"), "got: {msg}");
}
#[test]
fn two_ports_on_one_ident_is_ambiguous_rather_than_a_guess() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer(
"us-east-001",
vec![named_record("yah-marketing", "100.64.0.3", &[8080, 9090])],
);
let mut portless = rule("yah.dev", 0, &["us-east-001"]);
portless.port = None;
assert_eq!(fanout.port_for(&portless, "yah-marketing"), None);
}
#[test]
fn a_named_serving_port_disambiguates_a_multi_listener_workload() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer(
"us-east-001",
vec![named_port_record(
"yah-marketing",
"100.64.0.3",
&[("http", 8080), ("metrics", 9090)],
)],
);
let mut portless = rule("yah.dev", 0, &["us-east-001"]);
portless.port = None;
assert_eq!(fanout.port_for(&portless, "yah-marketing"), Some(8080));
}
#[test]
fn two_nodes_disagreeing_about_http_is_still_ambiguous() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer(
"us-east-001",
vec![named_port_record(
"yah-marketing",
"100.64.0.3",
&[("http", 8080), ("metrics", 9090)],
)],
);
fanout.push_answer(
"us-west-001",
vec![named_port_record(
"yah-marketing",
"100.64.0.9",
&[("http", 8081), ("metrics", 9090)],
)],
);
let mut portless = rule("yah.dev", 0, &["us-east-001", "us-west-001"]);
portless.port = None;
assert_eq!(fanout.port_for(&portless, "yah-marketing"), None);
}
#[test]
fn a_record_with_no_names_still_resolves_its_single_port() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer(
"us-east-001",
vec![DiscoveredRecord {
ident: "yah-marketing".into(),
mesh_ip: "100.64.0.3".into(),
ports: vec![43117],
named_ports: BTreeMap::new(),
}],
);
let mut portless = rule("yah.dev", 0, &["us-east-001"]);
portless.port = None;
assert_eq!(fanout.port_for(&portless, "yah-marketing"), Some(43117));
}
#[test]
fn a_record_on_a_node_outside_the_placement_does_not_answer() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer(
"us-west-001",
vec![named_record("yah-marketing", "100.64.0.9", &[43117])],
);
let mut portless = rule("yah.dev", 0, &["us-east-001"]);
portless.port = None;
assert_eq!(fanout.port_for(&portless, "yah-marketing"), None);
}
#[test]
fn an_unresolved_port_resolves_no_address_either() {
let fanout = crowded_node();
let mut portless = rule("yah.dev", 0, &["us-east-001"]);
portless.port = None;
assert!(fanout.upstreams_for(&portless).is_empty());
}
#[test]
fn an_ident_resolves_to_the_address_that_workload_registered() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer(
"us-east-001",
vec![
named_record("noisetable", "100.64.0.3", &[8080]),
named_record("noisetable-account", "100.64.0.3", &[43117]),
named_record("something-else", "100.64.0.3", &[9999]),
],
);
assert_eq!(
fanout.address_for_ident("noisetable-account").as_deref(),
Some("100.64.0.3:43117")
);
assert_eq!(fanout.address_for_ident("noisetable").as_deref(), Some("100.64.0.3:8080"));
assert_eq!(fanout.address_for_ident("not-deployed"), None);
}
#[test]
fn a_unit_on_two_nodes_is_ambiguous_rather_than_load_balanced() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![named_record("acct", "100.64.0.3", &[43117])]);
fanout.push_answer("us-south-001", vec![named_record("acct", "100.64.0.4", &[43117])]);
assert_eq!(fanout.address_for_ident("acct"), None);
}
#[test]
fn one_address_reported_twice_is_still_one_address() {
let mut fanout = ServiceRecordFanout::default();
fanout.push_answer("us-east-001", vec![named_record("acct", "100.64.0.3", &[43117])]);
fanout.push_answer("us-east-001-again", vec![named_record("acct", "100.64.0.3", &[43117])]);
assert_eq!(fanout.address_for_ident("acct").as_deref(), Some("100.64.0.3:43117"));
}
#[test]
fn a_multi_listener_unit_resolves_on_its_serving_port_name() {
let mut named = ServiceRecordFanout::default();
named.push_answer(
"us-east-001",
vec![DiscoveredRecord {
ident: "acct".into(),
mesh_ip: "100.64.0.3".into(),
ports: vec![43117, 9100],
named_ports: [("http".to_string(), 43117), ("metrics".to_string(), 9100)]
.into_iter()
.collect(),
}],
);
assert_eq!(named.address_for_ident("acct").as_deref(), Some("100.64.0.3:43117"));
let mut anonymous = ServiceRecordFanout::default();
anonymous.push_answer(
"us-east-001",
vec![named_record("acct", "100.64.0.3", &[43117, 9100])],
);
assert_eq!(anonymous.address_for_ident("acct"), None);
}
}