use std::collections::hash_map::DefaultHasher;
use std::collections::{HashMap, HashSet};
use std::hash::{Hash, Hasher};
use std::sync::Arc;
use std::time::Duration;
use bytes::Bytes;
use cf_system_sdks::directory::{DirectoryClient, ServiceInstanceInfo};
use tokio_util::sync::CancellationToken;
use toolkit_gateway::{
Endpoint, GatewayProvider, GearInstance, GearName, InstanceSpec, OpenApiSpec,
};
const MIN_SYNC_INTERVAL: Duration = Duration::from_secs(1);
const SYNC_TIMEOUT: Duration = Duration::from_secs(30);
pub fn spawn_directory_sync(
provider: Arc<dyn GatewayProvider>,
directory: Arc<dyn DirectoryClient>,
interval: Duration,
cancel: CancellationToken,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(directory_sync_loop(provider, directory, interval, cancel))
}
async fn directory_sync_loop(
provider: Arc<dyn GatewayProvider>,
directory: Arc<dyn DirectoryClient>,
interval: Duration,
cancel: CancellationToken,
) {
let mut ticker = tokio::time::interval(interval.max(MIN_SYNC_INTERVAL));
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let mut state = SyncState::default();
loop {
tokio::select! {
() = cancel.cancelled() => break,
_ = ticker.tick() => {
if reconcile_pass(provider.as_ref(), directory.as_ref(), &mut state, &cancel)
.await
.is_break()
{
break;
}
}
}
}
tracing::info!("gateway directory-sync stopping");
}
async fn reconcile_pass(
provider: &dyn GatewayProvider,
directory: &dyn DirectoryClient,
state: &mut SyncState,
cancel: &CancellationToken,
) -> std::ops::ControlFlow<()> {
tokio::select! {
() = cancel.cancelled() => std::ops::ControlFlow::Break(()),
result = tokio::time::timeout(SYNC_TIMEOUT, reconcile(provider, directory, state)) => {
if result.is_err() {
tracing::warn!(
timeout_secs = SYNC_TIMEOUT.as_secs(),
"gateway directory-sync reconcile timed out; abandoning pass",
);
}
std::ops::ControlFlow::Continue(())
}
}
}
async fn reconcile(
provider: &dyn GatewayProvider,
directory: &dyn DirectoryClient,
state: &mut SyncState,
) {
if let Err(err) = sync_once_stateful(provider, directory, state).await {
tracing::warn!(error = %err, "gateway directory-sync poll failed");
}
}
async fn sync_once_stateful(
provider: &dyn GatewayProvider,
directory: &dyn DirectoryClient,
state: &mut SyncState,
) -> anyhow::Result<()> {
let instances = directory.list_all_instances().await?;
reconcile_snapshot(provider, directory, instances, state).await;
Ok(())
}
#[cfg(test)]
async fn sync_once(
provider: &dyn GatewayProvider,
directory: &dyn DirectoryClient,
) -> anyhow::Result<()> {
sync_once_stateful(provider, directory, &mut SyncState::default()).await
}
#[derive(Default)]
struct SyncState {
last_applied: Option<u64>,
}
fn snapshot_fingerprint(instances: &[ServiceInstanceInfo]) -> u64 {
let mut entries: Vec<(&str, &str, &str, &str)> = instances
.iter()
.filter_map(|i| {
let rest = i.rest_endpoint.as_ref()?;
let hash = i.openapi_spec_hash.as_deref()?;
Some((
i.gear.as_str(),
i.instance_id.as_str(),
rest.uri.as_str(),
hash,
))
})
.collect();
entries.sort_unstable();
let mut hasher = DefaultHasher::new();
entries.hash(&mut hasher);
hasher.finish()
}
async fn reconcile_snapshot(
provider: &dyn GatewayProvider,
directory: &dyn DirectoryClient,
instances: Vec<ServiceInstanceInfo>,
state: &mut SyncState,
) {
let fingerprint = snapshot_fingerprint(&instances);
if state.last_applied == Some(fingerprint) {
tracing::debug!(
fingerprint,
"directory snapshot unchanged since last apply; skipping spec fetch and route rebuild",
);
return;
}
let Resolved {
desired,
to_register,
complete,
} = resolve_snapshot(directory, instances).await;
let (removals, empty_snapshot_valve) =
compute_removals(provider.registered_instances(), &desired);
provider.apply_snapshot(to_register, &removals).await;
state.last_applied = if complete && !empty_snapshot_valve {
Some(fingerprint)
} else {
None
};
}
struct Resolved {
desired: HashSet<GearInstance>,
to_register: Vec<InstanceSpec<'static>>,
complete: bool,
}
async fn resolve_snapshot(
directory: &dyn DirectoryClient,
instances: Vec<ServiceInstanceInfo>,
) -> Resolved {
let mut desired: HashSet<GearInstance> = HashSet::new();
let mut specs: HashMap<GearName, Bytes> = HashMap::new();
let mut to_register: Vec<InstanceSpec<'static>> = Vec::with_capacity(instances.len());
let mut complete = true;
for inst in instances {
let gear = GearName::from(inst.gear.as_str());
desired.insert(GearInstance::new(gear.clone(), &inst.instance_id));
let proxyable = inst.rest_endpoint.is_some() && inst.openapi_spec_hash.is_some();
match resolve_instance_spec(directory, &inst, &gear, &mut specs).await {
Some(spec) => to_register.push(spec),
None if proxyable => complete = false,
None => {}
}
}
Resolved {
desired,
to_register,
complete,
}
}
fn compute_removals(
registered: Vec<GearInstance>,
desired: &HashSet<GearInstance>,
) -> (Vec<GearInstance>, bool) {
if desired.is_empty() && !registered.is_empty() {
tracing::warn!(
registered = registered.len(),
"directory snapshot is empty while instances are registered; \
skipping deregistration this poll to avoid mass route withdrawal",
);
return (Vec::new(), true);
}
let removals = registered
.into_iter()
.filter(|instance| !desired.contains(instance))
.collect();
(removals, false)
}
async fn resolve_instance_spec(
directory: &dyn DirectoryClient,
inst: &ServiceInstanceInfo,
gear: &GearName,
specs: &mut HashMap<GearName, Bytes>,
) -> Option<InstanceSpec<'static>> {
let rest = inst.rest_endpoint.as_ref()?;
let endpoint = parse_endpoint(gear.as_str(), &rest.uri)?;
let spec = fetch_or_get_spec(directory, gear, specs).await?;
Some(InstanceSpec {
gear: gear.clone(),
instance_id: inst.instance_id.clone(),
spec: OpenApiSpec::SerializedJson(spec),
endpoint,
})
}
async fn fetch_or_get_spec(
directory: &dyn DirectoryClient,
gear: &GearName,
specs: &mut HashMap<GearName, Bytes>,
) -> Option<Bytes> {
if let Some(spec) = specs.get(gear) {
return Some(spec.clone());
}
let spec = Bytes::from(fetch_spec(directory, gear.as_str()).await?);
specs.insert(gear.clone(), spec.clone());
Some(spec)
}
fn parse_endpoint(gear: &str, uri: &str) -> Option<Endpoint> {
match Endpoint::parse(uri) {
Ok(endpoint) => Some(endpoint),
Err(err) => {
tracing::warn!(
gear = %gear,
uri = %uri,
error = %err,
"skipping gear with unparseable REST endpoint",
);
None
}
}
}
async fn fetch_spec(directory: &dyn DirectoryClient, gear: &str) -> Option<String> {
match directory.get_openapi_spec(gear).await {
Ok(spec) => Some(spec),
Err(err) => {
tracing::warn!(gear = %gear, error = %err, "skipping gear: could not fetch OpenAPI spec");
None
}
}
}
#[cfg(test)]
mod tests {
use super::{
SyncState, snapshot_fingerprint, spawn_directory_sync, sync_once, sync_once_stateful,
};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use toolkit_gateway::{
Endpoint, GatewayError, GatewayProvider, GearInstance, GearName, InstanceSpec, OpenApiSpec,
ProxyRegistry, ToolKitGatewayProvider,
};
use cf_system_sdks::directory::{
DirectoryClient, RegisterInstanceInfo, ServiceEndpoint, ServiceInstanceInfo,
};
use tokio_util::sync::CancellationToken;
use toolkit::directory::LocalDirectoryClient;
use toolkit::runtime::GearManager;
struct CountingProvider {
inner: ToolKitGatewayProvider,
applies: Arc<AtomicUsize>,
}
impl CountingProvider {
fn new(registry: Arc<ProxyRegistry>) -> (Self, Arc<AtomicUsize>) {
let applies = Arc::new(AtomicUsize::new(0));
let provider = Self {
inner: ToolKitGatewayProvider::new(registry),
applies: Arc::clone(&applies),
};
(provider, applies)
}
}
#[async_trait::async_trait]
impl GatewayProvider for CountingProvider {
async fn register_routes(
&self,
gear: &GearName,
instance_id: &str,
spec: OpenApiSpec<'_>,
endpoint: &Endpoint,
) -> Result<(), GatewayError> {
self.inner
.register_routes(gear, instance_id, spec, endpoint)
.await
}
async fn deregister_routes(
&self,
gear: &GearName,
instance_id: &str,
) -> Result<(), GatewayError> {
self.inner.deregister_routes(gear, instance_id).await
}
fn registered_instances(&self) -> Vec<GearInstance> {
self.inner.registered_instances()
}
async fn apply_snapshot<'a>(
&self,
instances: Vec<InstanceSpec<'a>>,
removals: &[GearInstance],
) {
self.applies.fetch_add(1, Ordering::SeqCst);
self.inner.apply_snapshot(instances, removals).await;
}
}
fn public_spec(gear: &str, path: &str) -> String {
serde_json::json!({
"openapi": "3.1.0",
"info": { "title": gear, "version": "1.0.0" },
"paths": { path: { "get": { "x-toolkit-visibility": "public", "responses": {} } } },
})
.to_string()
}
async fn register(dir: &dyn DirectoryClient, gear: &str, uri: &str, spec: String) -> String {
let instance_id = uuid::Uuid::new_v4().to_string();
dir.register_instance(RegisterInstanceInfo {
gear: gear.to_owned(),
instance_id: instance_id.clone(),
grpc_services: vec![],
version: None,
rest_endpoint: Some(ServiceEndpoint::new(uri)),
openapi_spec: Some(spec),
})
.await
.unwrap();
instance_id
}
#[tokio::test]
async fn sync_registers_public_routes_then_prunes_removed_gears() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
let calc_instance = register(
dir.as_ref(),
"calc",
"http://calc:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
register(
dir.as_ref(),
"billing",
"http://billing:9090",
public_spec("billing", "/billing/v1/pay"),
)
.await;
let registry = Arc::new(ProxyRegistry::new());
let provider = ToolKitGatewayProvider::new(Arc::clone(®istry));
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(registry.match_path("/calc/v1/add").is_some());
assert!(registry.match_path("/billing/v1/pay").is_some());
assert_eq!(registry.instance_count(), 2);
dir.deregister_instance("calc", &calc_instance)
.await
.unwrap();
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(registry.match_path("/calc/v1/add").is_none());
assert!(registry.match_path("/billing/v1/pay").is_some());
assert!(!registry.contains_gear(&GearName::from("calc")));
}
#[tokio::test]
async fn sync_picks_up_spec_change_on_resync() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
let id = register(
dir.as_ref(),
"calc",
"http://calc:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
let registry = Arc::new(ProxyRegistry::new());
let provider = ToolKitGatewayProvider::new(Arc::clone(®istry));
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(registry.match_path("/calc/v1/add").is_some());
assert!(registry.match_path("/calc/v1/sub").is_none());
dir.register_instance(RegisterInstanceInfo {
gear: "calc".to_owned(),
instance_id: id.clone(),
grpc_services: vec![],
version: None,
rest_endpoint: Some(ServiceEndpoint::new("http://calc:8080")),
openapi_spec: Some(public_spec("calc", "/calc/v1/sub")),
})
.await
.unwrap();
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(registry.match_path("/calc/v1/sub").is_some());
assert!(registry.match_path("/calc/v1/add").is_none());
assert_eq!(registry.instance_count(), 1);
}
#[tokio::test]
async fn sync_retains_routes_when_new_instance_fails_but_gear_still_present() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
let good = register(
dir.as_ref(),
"calc",
"http://calc:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
let registry = Arc::new(ProxyRegistry::new());
let provider = ToolKitGatewayProvider::new(Arc::clone(®istry));
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(registry.match_path("/calc/v1/add").is_some());
assert!(registry.contains_instance(&GearName::from("calc"), &good));
register(
dir.as_ref(),
"calc",
"/schemeless-uri",
public_spec("calc", "/calc/v1/add"),
)
.await;
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(registry.match_path("/calc/v1/add").is_some());
assert!(registry.contains_gear(&GearName::from("calc")));
assert!(registry.contains_instance(&GearName::from("calc"), &good));
}
#[tokio::test]
async fn sync_skips_instances_without_rest_or_spec() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
dir.register_instance(RegisterInstanceInfo {
gear: "worker".to_owned(),
instance_id: uuid::Uuid::new_v4().to_string(),
grpc_services: vec![(
"worker.Svc".to_owned(),
ServiceEndpoint::new("http://worker:7000"),
)],
version: None,
rest_endpoint: None,
openapi_spec: None,
})
.await
.unwrap();
let registry = Arc::new(ProxyRegistry::new());
let provider = ToolKitGatewayProvider::new(Arc::clone(®istry));
sync_once(&provider, dir.as_ref()).await.unwrap();
assert_eq!(registry.instance_count(), 0);
}
#[tokio::test]
async fn sync_skips_gear_with_rest_but_no_published_spec() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
dir.register_instance(RegisterInstanceInfo {
gear: "specless".to_owned(),
instance_id: uuid::Uuid::new_v4().to_string(),
grpc_services: vec![],
version: None,
rest_endpoint: Some(ServiceEndpoint::new("http://specless:8080")),
openapi_spec: None,
})
.await
.unwrap();
let registry = Arc::new(ProxyRegistry::new());
let provider = ToolKitGatewayProvider::new(Arc::clone(®istry));
sync_once(&provider, dir.as_ref()).await.unwrap();
assert_eq!(registry.instance_count(), 0);
}
#[tokio::test]
async fn sync_fetches_spec_lazily_from_slim_snapshot() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
register(
dir.as_ref(),
"calc",
"http://calc:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
let snapshot = dir.list_all_instances().await.unwrap();
assert!(snapshot.iter().all(|i| i.openapi_spec.is_none()));
let registry = Arc::new(ProxyRegistry::new());
let provider = ToolKitGatewayProvider::new(Arc::clone(®istry));
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(registry.match_path("/calc/v1/add").is_some());
}
#[tokio::test]
async fn sync_deregisters_only_one_instance_of_same_gear() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
let calc_a = register(
dir.as_ref(),
"calc",
"http://calc-a:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
let calc_b = register(
dir.as_ref(),
"calc",
"http://calc-b:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
register(
dir.as_ref(),
"billing",
"http://billing:9090",
public_spec("billing", "/billing/v1/pay"),
)
.await;
let registry = Arc::new(ProxyRegistry::new());
let provider = ToolKitGatewayProvider::new(Arc::clone(®istry));
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(registry.contains_instance(&GearName::from("calc"), &calc_a));
assert!(registry.contains_instance(&GearName::from("calc"), &calc_b));
assert_eq!(registry.instance_count(), 3);
let matched = registry.match_path("/calc/v1/add").expect("calc match");
assert_eq!(matched.gear.as_str(), "calc");
assert!(matches!(
matched.endpoint.authority().as_str(),
"calc-a:8080" | "calc-b:8080"
));
assert!(registry.match_path("/billing/v1/pay").is_some());
dir.deregister_instance("calc", &calc_a).await.unwrap();
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(!registry.contains_instance(&GearName::from("calc"), &calc_a));
assert!(registry.contains_instance(&GearName::from("calc"), &calc_b));
assert!(registry.contains_gear(&GearName::from("calc")));
assert_eq!(registry.instance_count(), 2);
let remaining = registry
.match_path("/calc/v1/add")
.expect("calc still proxied");
assert_eq!(remaining.endpoint.authority().as_str(), "calc-b:8080");
assert!(registry.match_path("/billing/v1/pay").is_some());
}
#[tokio::test]
async fn sync_skips_mass_deregistration_on_empty_snapshot() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
let calc = register(
dir.as_ref(),
"calc",
"http://calc:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
let registry = Arc::new(ProxyRegistry::new());
let provider = ToolKitGatewayProvider::new(Arc::clone(®istry));
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(registry.match_path("/calc/v1/add").is_some());
assert_eq!(registry.instance_count(), 1);
dir.deregister_instance("calc", &calc).await.unwrap();
assert!(dir.list_all_instances().await.unwrap().is_empty());
sync_once(&provider, dir.as_ref()).await.unwrap();
assert_eq!(registry.instance_count(), 1);
assert!(registry.match_path("/calc/v1/add").is_some());
}
#[tokio::test]
async fn sync_prunes_normally_when_snapshot_nonempty() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
let calc = register(
dir.as_ref(),
"calc",
"http://calc:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
register(
dir.as_ref(),
"billing",
"http://billing:9090",
public_spec("billing", "/billing/v1/pay"),
)
.await;
let registry = Arc::new(ProxyRegistry::new());
let provider = ToolKitGatewayProvider::new(Arc::clone(®istry));
sync_once(&provider, dir.as_ref()).await.unwrap();
assert_eq!(registry.instance_count(), 2);
dir.deregister_instance("calc", &calc).await.unwrap();
sync_once(&provider, dir.as_ref()).await.unwrap();
assert!(!registry.contains_gear(&GearName::from("calc")));
assert!(registry.match_path("/billing/v1/pay").is_some());
assert_eq!(registry.instance_count(), 1);
}
#[tokio::test]
async fn spawn_loop_syncs_then_stops_on_cancel() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
register(
dir.as_ref(),
"calc",
"http://calc:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
let registry = Arc::new(ProxyRegistry::new());
let provider: Arc<dyn GatewayProvider> =
Arc::new(ToolKitGatewayProvider::new(Arc::clone(®istry)));
let cancel = CancellationToken::new();
let handle = spawn_directory_sync(
provider,
Arc::clone(&dir),
Duration::from_secs(1),
cancel.clone(),
);
let mut synced = false;
for _ in 0..100 {
if registry.match_path("/calc/v1/add").is_some() {
synced = true;
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
assert!(synced, "background loop should reconcile on the first tick");
cancel.cancel();
handle.await.expect("sync loop joins cleanly after cancel");
}
#[tokio::test]
async fn sync_skips_rebuild_when_snapshot_unchanged() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
register(
dir.as_ref(),
"calc",
"http://calc:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
let registry = Arc::new(ProxyRegistry::new());
let (provider, applies) = CountingProvider::new(Arc::clone(®istry));
let mut state = SyncState::default();
sync_once_stateful(&provider, dir.as_ref(), &mut state)
.await
.unwrap();
assert_eq!(applies.load(Ordering::SeqCst), 1, "first poll reconciles");
assert!(registry.match_path("/calc/v1/add").is_some());
sync_once_stateful(&provider, dir.as_ref(), &mut state)
.await
.unwrap();
assert_eq!(
applies.load(Ordering::SeqCst),
1,
"unchanged second poll is skipped (no rebuild)"
);
assert!(registry.match_path("/calc/v1/add").is_some());
}
#[tokio::test]
async fn sync_reapplies_after_spec_change() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
let id = register(
dir.as_ref(),
"calc",
"http://calc:8080",
public_spec("calc", "/calc/v1/add"),
)
.await;
let registry = Arc::new(ProxyRegistry::new());
let (provider, applies) = CountingProvider::new(Arc::clone(®istry));
let mut state = SyncState::default();
sync_once_stateful(&provider, dir.as_ref(), &mut state)
.await
.unwrap();
assert_eq!(applies.load(Ordering::SeqCst), 1);
dir.register_instance(RegisterInstanceInfo {
gear: "calc".to_owned(),
instance_id: id,
grpc_services: vec![],
version: None,
rest_endpoint: Some(ServiceEndpoint::new("http://calc:8080")),
openapi_spec: Some(public_spec("calc", "/calc/v1/sub")),
})
.await
.unwrap();
sync_once_stateful(&provider, dir.as_ref(), &mut state)
.await
.unwrap();
assert_eq!(
applies.load(Ordering::SeqCst),
2,
"changed spec hash forces a reconcile"
);
assert!(registry.match_path("/calc/v1/sub").is_some());
assert!(registry.match_path("/calc/v1/add").is_none());
}
#[tokio::test]
async fn sync_retries_when_snapshot_unchanged_but_last_pass_incomplete() {
let mgr = Arc::new(GearManager::new());
let dir: Arc<dyn DirectoryClient> = Arc::new(LocalDirectoryClient::new(Arc::clone(&mgr)));
register(
dir.as_ref(),
"calc",
"/schemeless-uri",
public_spec("calc", "/calc/v1/add"),
)
.await;
let registry = Arc::new(ProxyRegistry::new());
let (provider, applies) = CountingProvider::new(Arc::clone(®istry));
let mut state = SyncState::default();
sync_once_stateful(&provider, dir.as_ref(), &mut state)
.await
.unwrap();
assert_eq!(applies.load(Ordering::SeqCst), 1);
assert!(
state.last_applied.is_none(),
"incomplete pass is not cached"
);
sync_once_stateful(&provider, dir.as_ref(), &mut state)
.await
.unwrap();
assert_eq!(
applies.load(Ordering::SeqCst),
2,
"unresolved proxyable instance forces a retry next poll"
);
}
fn instance_with_hash(gear: &str, uri: &str, hash: Option<&str>) -> ServiceInstanceInfo {
ServiceInstanceInfo {
gear: gear.to_owned(),
instance_id: "inst-1".to_owned(),
endpoint: ServiceEndpoint::new(uri),
version: None,
rest_endpoint: Some(ServiceEndpoint::new(uri)),
openapi_spec: None,
openapi_spec_hash: hash.map(ToOwned::to_owned),
}
}
#[test]
fn snapshot_fingerprint_is_stable_and_change_sensitive() {
let a = vec![instance_with_hash("calc", "http://calc:8080", Some("h1"))];
let a_again = vec![instance_with_hash("calc", "http://calc:8080", Some("h1"))];
let changed = vec![instance_with_hash("calc", "http://calc:8080", Some("h2"))];
assert_eq!(
snapshot_fingerprint(&a),
snapshot_fingerprint(&a_again),
"identical snapshots fingerprint equally"
);
assert_ne!(
snapshot_fingerprint(&a),
snapshot_fingerprint(&changed),
"a changed spec hash changes the fingerprint"
);
}
}