#![allow(dead_code)]
use std::sync::Arc;
use xds_cache::{Cache, ShardedCache, Snapshot};
use xds_core::{BoxResource, NodeHash, Resource, ResourceRegistry, TypeUrl};
use crate::delta::{ClientResourceState, DeltaHandler};
use crate::services::{AdsConfig, AdsService};
use crate::sotw::SotwHandler;
use crate::stream::StreamContext;
#[derive(Debug, Clone)]
struct TestCluster {
name: String,
version: String,
}
impl TestCluster {
fn new(name: &str) -> Self {
Self {
name: name.to_string(),
version: "1".to_string(),
}
}
}
impl Resource for TestCluster {
fn type_url(&self) -> &str {
TypeUrl::CLUSTER
}
fn name(&self) -> &str {
&self.name
}
fn encode(&self) -> Result<prost_types::Any, Box<dyn std::error::Error + Send + Sync>> {
Ok(prost_types::Any {
type_url: self.type_url().to_string(),
value: self.name.as_bytes().to_vec(),
})
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
#[derive(Debug, Clone)]
struct TestListener {
name: String,
}
impl TestListener {
fn new(name: &str) -> Self {
Self {
name: name.to_string(),
}
}
}
impl Resource for TestListener {
fn type_url(&self) -> &str {
TypeUrl::LISTENER
}
fn name(&self) -> &str {
&self.name
}
fn encode(&self) -> Result<prost_types::Any, Box<dyn std::error::Error + Send + Sync>> {
Ok(prost_types::Any {
type_url: self.type_url().to_string(),
value: self.name.as_bytes().to_vec(),
})
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
fn setup_cache_with_clusters(node_id: &str, clusters: Vec<&str>) -> (Arc<ShardedCache>, NodeHash) {
let cache = Arc::new(ShardedCache::new());
let node_hash = NodeHash::from_id(node_id);
let resources: Vec<BoxResource> = clusters
.into_iter()
.map(|name| Arc::new(TestCluster::new(name)) as BoxResource)
.collect();
let snapshot = Snapshot::builder()
.version("v1")
.resources(TypeUrl::CLUSTER.into(), resources)
.build();
cache.set_snapshot(node_hash, snapshot);
(cache, node_hash)
}
mod sotw_protocol {
use super::*;
#[test]
fn initial_request_empty_version() {
let (cache, node_hash) =
setup_cache_with_clusters("node-1", vec!["cluster-a", "cluster-b"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache, registry);
let ctx = StreamContext::new();
let result = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "", &[], node_hash)
.unwrap();
assert!(result.is_some(), "Initial request should receive response");
let response = result.unwrap();
assert_eq!(response.version_info, "v1");
assert_eq!(response.resources.len(), 2);
}
#[test]
fn request_with_current_version_no_response() {
let (cache, node_hash) = setup_cache_with_clusters("node-1", vec!["cluster-a"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache, registry);
let ctx = StreamContext::new();
let result = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "v1", &[], node_hash)
.unwrap();
assert!(result.is_none(), "Same version should not trigger response");
}
#[test]
fn request_with_old_version_gets_response() {
let (cache, node_hash) = setup_cache_with_clusters("node-1", vec!["cluster-a"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache, registry);
let ctx = StreamContext::new();
let result = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "v0", &[], node_hash)
.unwrap();
assert!(result.is_some(), "Old version should trigger response");
assert_eq!(result.unwrap().version_info, "v1");
}
#[test]
fn wildcard_subscription_returns_all() {
let (cache, node_hash) =
setup_cache_with_clusters("node-1", vec!["cluster-a", "cluster-b", "cluster-c"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache, registry);
let ctx = StreamContext::new();
let result = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "", &[], node_hash)
.unwrap();
let response = result.unwrap();
assert_eq!(response.resources.len(), 3);
}
#[test]
fn explicit_subscription_returns_requested() {
let (cache, node_hash) =
setup_cache_with_clusters("node-1", vec!["cluster-a", "cluster-b", "cluster-c"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache, registry);
let ctx = StreamContext::new();
let result = handler
.process_request(
&ctx,
TypeUrl::CLUSTER.into(),
"",
&["cluster-a".to_string(), "cluster-c".to_string()],
node_hash,
)
.unwrap();
let response = result.unwrap();
assert_eq!(response.resources.len(), 2);
}
#[test]
fn response_nonce_is_unique() {
let (cache, node_hash) = setup_cache_with_clusters("node-1", vec!["cluster-a"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache, registry);
let ctx = StreamContext::new();
let r1 = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "", &[], node_hash)
.unwrap()
.unwrap();
let snapshot = Snapshot::builder()
.version("v2")
.resources(
TypeUrl::CLUSTER.into(),
vec![Arc::new(TestCluster::new("cluster-a")) as BoxResource],
)
.build();
handler.cache().set_snapshot(node_hash, snapshot);
let r2 = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "v1", &[], node_hash)
.unwrap()
.unwrap();
assert_ne!(r1.nonce, r2.nonce, "Each response should have unique nonce");
}
#[test]
fn unknown_node_no_response() {
let cache = Arc::new(ShardedCache::new());
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache, registry);
let ctx = StreamContext::new();
let unknown_node = NodeHash::from_id("unknown-node");
let result = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "", &[], unknown_node)
.unwrap();
assert!(result.is_none(), "Unknown node should not receive response");
}
#[test]
fn response_type_url_matches_request() {
let (cache, node_hash) = setup_cache_with_clusters("node-1", vec!["cluster-a"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache, registry);
let ctx = StreamContext::new();
let result = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "", &[], node_hash)
.unwrap()
.unwrap();
assert_eq!(result.type_url.as_str(), TypeUrl::CLUSTER);
}
}
mod delta_protocol {
use super::*;
#[test]
fn initial_wildcard_returns_all_resources() {
let (cache, node_hash) =
setup_cache_with_clusters("delta-node", vec!["cluster-a", "cluster-b"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = DeltaHandler::new(cache, registry);
let ctx = StreamContext::new();
let mut state = ClientResourceState::new();
let response = handler
.process_request(
&ctx,
TypeUrl::CLUSTER.into(),
&mut state,
vec![],
vec![],
node_hash,
)
.expect("delta process_request should not error")
.expect("first wildcard request should produce a response");
assert_eq!(response.resources.len(), 2);
assert!(response.removed_resources.is_empty());
assert!(!response.nonce.is_empty());
}
#[test]
fn same_version_returns_no_response() {
let (cache, node_hash) = setup_cache_with_clusters("delta-node", vec!["cluster-a"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = DeltaHandler::new(cache, registry);
let ctx = StreamContext::new();
let mut state = ClientResourceState::new();
let _first = handler
.process_request(
&ctx,
TypeUrl::CLUSTER.into(),
&mut state,
vec![],
vec![],
node_hash,
)
.unwrap()
.unwrap();
let second = handler
.process_request(
&ctx,
TypeUrl::CLUSTER.into(),
&mut state,
vec![],
vec![],
node_hash,
)
.unwrap();
assert!(second.is_none());
}
#[test]
fn dropped_resource_is_reported_as_removed() {
let (cache, node_hash) =
setup_cache_with_clusters("delta-node", vec!["cluster-a", "cluster-b"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = DeltaHandler::new(Arc::clone(&cache), registry);
let ctx = StreamContext::new();
let mut state = ClientResourceState::new();
let _first = handler
.process_request(
&ctx,
TypeUrl::CLUSTER.into(),
&mut state,
vec![],
vec![],
node_hash,
)
.unwrap()
.unwrap();
let resources: Vec<BoxResource> = vec![Arc::new(TestCluster::new("cluster-a"))];
let snapshot = Snapshot::builder()
.version("v2")
.resources(TypeUrl::CLUSTER.into(), resources)
.build();
cache.set_snapshot(node_hash, snapshot);
let response = handler
.process_request(
&ctx,
TypeUrl::CLUSTER.into(),
&mut state,
vec![],
vec![],
node_hash,
)
.unwrap()
.expect("snapshot change should produce a delta response");
assert!(response.removed_resources.contains(&"cluster-b".to_string()));
}
}
mod ads_protocol {
use super::*;
#[test]
fn ads_handles_multiple_types() {
let cache = Arc::new(ShardedCache::new());
let registry = Arc::new(ResourceRegistry::new());
let service = AdsService::new(cache.clone(), registry);
let node_hash = NodeHash::from_id("node-1");
let cluster: BoxResource = Arc::new(TestCluster::new("cluster-1"));
let listener: BoxResource = Arc::new(TestListener::new("listener-1"));
let snapshot = Snapshot::builder()
.version("v1")
.resources(TypeUrl::CLUSTER.into(), vec![cluster])
.resources(TypeUrl::LISTENER.into(), vec![listener])
.build();
cache.set_snapshot(node_hash, snapshot);
let ctx = StreamContext::new();
let cluster_response = service
.process_sotw_request(&ctx, TypeUrl::CLUSTER, "", &[], node_hash, "", None)
.unwrap();
assert!(cluster_response.is_some());
let listener_response = service
.process_sotw_request(&ctx, TypeUrl::LISTENER, "", &[], node_hash, "", None)
.unwrap();
assert!(listener_response.is_some());
}
#[test]
fn ads_config_customization() {
let cache = Arc::new(ShardedCache::new());
let registry = Arc::new(ResourceRegistry::new());
let config = AdsConfig {
max_concurrent_streams: 50,
response_buffer_size: 8,
enable_delta: false,
};
let service = AdsService::with_config(cache, registry, config);
assert!(!service.config().enable_delta);
assert_eq!(service.config().max_concurrent_streams, 50);
}
}
mod resource_types {
use super::*;
#[test]
fn standard_type_urls() {
assert_eq!(
TypeUrl::CLUSTER,
"type.googleapis.com/envoy.config.cluster.v3.Cluster"
);
assert_eq!(
TypeUrl::LISTENER,
"type.googleapis.com/envoy.config.listener.v3.Listener"
);
assert_eq!(
TypeUrl::ROUTE,
"type.googleapis.com/envoy.config.route.v3.RouteConfiguration"
);
assert_eq!(
TypeUrl::ENDPOINT,
"type.googleapis.com/envoy.config.endpoint.v3.ClusterLoadAssignment"
);
assert_eq!(
TypeUrl::SECRET,
"type.googleapis.com/envoy.extensions.transport_sockets.tls.v3.Secret"
);
}
#[test]
fn type_url_parsing() {
let cluster_url = TypeUrl::new(TypeUrl::CLUSTER);
assert_eq!(cluster_url.as_str(), TypeUrl::CLUSTER);
let custom_url = TypeUrl::new("custom.type/Resource");
assert_eq!(custom_url.as_str(), "custom.type/Resource");
}
}
mod node_identification {
use super::*;
#[test]
fn different_nodes_different_hashes() {
let hash1 = NodeHash::from_id("node-1");
let hash2 = NodeHash::from_id("node-2");
assert_ne!(hash1, hash2);
}
#[test]
fn same_node_same_hash() {
let hash1 = NodeHash::from_id("node-1");
let hash2 = NodeHash::from_id("node-1");
assert_eq!(hash1, hash2);
}
#[test]
fn node_isolation() {
let cache = Arc::new(ShardedCache::new());
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache.clone(), registry);
let node1 = NodeHash::from_id("node-1");
let node2 = NodeHash::from_id("node-2");
let snapshot1 = Snapshot::builder()
.version("v1-node1")
.resources(
TypeUrl::CLUSTER.into(),
vec![Arc::new(TestCluster::new("cluster-for-node1")) as BoxResource],
)
.build();
let snapshot2 = Snapshot::builder()
.version("v1-node2")
.resources(
TypeUrl::CLUSTER.into(),
vec![Arc::new(TestCluster::new("cluster-for-node2")) as BoxResource],
)
.build();
cache.set_snapshot(node1, snapshot1);
cache.set_snapshot(node2, snapshot2);
let ctx = StreamContext::new();
let r1 = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "", &[], node1)
.unwrap()
.unwrap();
let r2 = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "", &[], node2)
.unwrap()
.unwrap();
assert_eq!(r1.version_info, "v1-node1");
assert_eq!(r2.version_info, "v1-node2");
}
}
mod versioning {
use super::*;
#[test]
fn version_string_comparison() {
let (cache, node_hash) = setup_cache_with_clusters("node-1", vec!["cluster-a"]);
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache.clone(), registry);
let ctx = StreamContext::new();
let r1 = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "", &[], node_hash)
.unwrap()
.unwrap();
let r2 = handler
.process_request(
&ctx,
TypeUrl::CLUSTER.into(),
&r1.version_info,
&[],
node_hash,
)
.unwrap();
assert!(
r2.is_none(),
"Exact version match should not trigger response"
);
}
#[test]
fn semantic_versions() {
let cache = Arc::new(ShardedCache::new());
let registry = Arc::new(ResourceRegistry::new());
let handler = SotwHandler::new(cache.clone(), registry);
let ctx = StreamContext::new();
let node_hash = NodeHash::from_id("node-1");
let snapshot = Snapshot::builder()
.version("1.2.3")
.resources(
TypeUrl::CLUSTER.into(),
vec![Arc::new(TestCluster::new("cluster-a")) as BoxResource],
)
.build();
cache.set_snapshot(node_hash, snapshot);
let response = handler
.process_request(&ctx, TypeUrl::CLUSTER.into(), "1.2.2", &[], node_hash)
.unwrap()
.unwrap();
assert_eq!(response.version_info, "1.2.3");
}
}