use super::*;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use async_trait::async_trait;
use cf_system_sdks::directory::{
DirectoryInvalidArgument, DirectoryPermissionDenied, ServiceEndpoint, ServiceInstanceInfo,
};
#[derive(Clone, Copy)]
enum Rejection {
InvalidArgument,
PermissionDenied,
Conflict,
}
struct StubDirectory {
fail_register: AtomicUsize,
register_calls: AtomicUsize,
register_ok: AtomicUsize,
heartbeats: AtomicUsize,
fail_resolve: AtomicUsize,
reject_kind: Option<Rejection>,
rejecting: AtomicBool,
}
impl StubDirectory {
fn new(fail_register: usize, fail_resolve: usize) -> Self {
Self {
fail_register: AtomicUsize::new(fail_register),
register_calls: AtomicUsize::new(0),
register_ok: AtomicUsize::new(0),
heartbeats: AtomicUsize::new(0),
fail_resolve: AtomicUsize::new(fail_resolve),
reject_kind: None,
rejecting: AtomicBool::new(false),
}
}
fn rejecting(kind: Rejection) -> Self {
Self {
reject_kind: Some(kind),
rejecting: AtomicBool::new(true),
..Self::new(0, 0)
}
}
fn stop_rejecting(&self) {
self.rejecting.store(false, Ordering::SeqCst);
}
}
#[async_trait]
impl DirectoryClient for StubDirectory {
async fn resolve_grpc_service(&self, _service: &str) -> anyhow::Result<ServiceEndpoint> {
Ok(ServiceEndpoint::new("http://grpc"))
}
async fn resolve_rest_service(&self, gear: &str) -> anyhow::Result<ServiceEndpoint> {
if self.fail_resolve.load(Ordering::SeqCst) > 0 {
self.fail_resolve.fetch_sub(1, Ordering::SeqCst);
anyhow::bail!("not yet available");
}
Ok(ServiceEndpoint::new(format!("http://{gear}:8080")))
}
async fn get_openapi_spec(&self, _gear: &str) -> anyhow::Result<String> {
Ok(String::new())
}
async fn list_instances(&self, _gear: &str) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
Ok(vec![])
}
async fn list_all_instances(&self) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
Ok(vec![])
}
async fn register_instance(&self, _info: RegisterInstanceInfo) -> anyhow::Result<()> {
self.register_calls.fetch_add(1, Ordering::SeqCst);
if self.rejecting.load(Ordering::SeqCst) {
match self.reject_kind {
Some(Rejection::InvalidArgument) => {
return Err(DirectoryInvalidArgument::new("bad label").into());
}
Some(Rejection::PermissionDenied) => {
return Err(DirectoryPermissionDenied::new("peer not authorized").into());
}
Some(Rejection::Conflict) => {
anyhow::bail!(
"directory register_instance failed: gRPC FailedPrecondition: \
a gRPC service name in this registration is already owned by another gear"
);
}
None => {}
}
}
if self.fail_register.load(Ordering::SeqCst) > 0 {
self.fail_register.fetch_sub(1, Ordering::SeqCst);
anyhow::bail!("directory unavailable");
}
self.register_ok.fetch_add(1, Ordering::SeqCst);
Ok(())
}
async fn deregister_instance(&self, _gear: &str, _instance: &str) -> anyhow::Result<()> {
Ok(())
}
async fn send_heartbeat(&self, _gear: &str, _instance: &str) -> anyhow::Result<()> {
self.heartbeats.fetch_add(1, Ordering::SeqCst);
Ok(())
}
}
fn info() -> RegisterInstanceInfo {
RegisterInstanceInfo::new("billing", "i-1")
.with_version("1.0.0")
.with_rest_endpoint(ServiceEndpoint::new("http://billing:8080"))
.with_openapi_spec("{}")
}
#[test]
fn backoff_doubles_and_caps() {
assert_eq!(
next_backoff(Duration::from_millis(100)),
Duration::from_millis(200)
);
assert_eq!(
next_backoff(Duration::from_millis(200)),
Duration::from_millis(400)
);
assert_eq!(next_backoff(Duration::from_secs(20)), MAX_BACKOFF);
assert_eq!(next_backoff(MAX_BACKOFF), MAX_BACKOFF);
}
#[tokio::test]
async fn registration_retries_until_success() {
let directory: Arc<dyn DirectoryClient> = Arc::new(StubDirectory::new(2, 0));
let cancel = CancellationToken::new();
let ok = register_once_with_backoff(&directory, &info(), &cancel).await;
assert!(ok);
}
#[tokio::test]
async fn registration_keeps_loop_alive_on_permanent_rejection() {
let stub = Arc::new(StubDirectory::rejecting(Rejection::InvalidArgument));
let directory: Arc<dyn DirectoryClient> = stub.clone();
let cancel = CancellationToken::new();
let ok = register_once_with_backoff(&directory, &info(), &cancel).await;
assert!(ok, "a permanent rejection must not stop the presence loop");
assert_eq!(
stub.register_calls.load(Ordering::SeqCst),
1,
"a permanent rejection must not be retried"
);
}
#[tokio::test]
async fn registration_does_not_retry_permission_denied() {
let stub = Arc::new(StubDirectory::rejecting(Rejection::PermissionDenied));
let directory: Arc<dyn DirectoryClient> = stub.clone();
let cancel = CancellationToken::new();
let ok = register_once_with_backoff(&directory, &info(), &cancel).await;
assert!(ok, "a permanent rejection must not stop the presence loop");
assert_eq!(
stub.register_calls.load(Ordering::SeqCst),
1,
"a permission-denied rejection must not be retried"
);
}
#[tokio::test]
async fn registration_returns_false_when_cancelled() {
let directory: Arc<dyn DirectoryClient> = Arc::new(StubDirectory::new(1000, 0));
let cancel = CancellationToken::new();
cancel.cancel();
let ok = register_once_with_backoff(&directory, &info(), &cancel).await;
assert!(!ok, "cancelled registration must return false");
}
#[tokio::test]
async fn presence_loop_registers_once_then_heartbeats_until_cancel() {
let stub = Arc::new(StubDirectory::new(0, 0));
let directory: Arc<dyn DirectoryClient> = stub.clone();
let cancel = CancellationToken::new();
let task = tokio::spawn(presence_loop(
Arc::clone(&directory),
info(),
Duration::from_secs(1),
cancel.clone(),
));
tokio::time::sleep(Duration::from_millis(1500)).await;
assert_eq!(
stub.register_calls.load(Ordering::SeqCst),
1,
"steady state registers exactly once (self-heal re-register is 30s, not fired here)"
);
assert!(
stub.heartbeats.load(Ordering::SeqCst) >= 1,
"the single presence task must send heartbeats"
);
cancel.cancel();
task.await.unwrap();
}
#[tokio::test]
async fn presence_loop_survives_permanent_rejection() {
let stub = Arc::new(StubDirectory::rejecting(Rejection::InvalidArgument));
let directory: Arc<dyn DirectoryClient> = stub.clone();
let cancel = CancellationToken::new();
let task = tokio::spawn(presence_loop(
Arc::clone(&directory),
info(),
Duration::from_secs(1),
cancel.clone(),
));
tokio::time::sleep(Duration::from_millis(1500)).await;
assert!(
!task.is_finished(),
"a permanent rejection must not stop the presence loop"
);
assert!(
stub.heartbeats.load(Ordering::SeqCst) >= 1,
"the loop must keep heartbeating despite a permanent rejection"
);
cancel.cancel();
task.await.unwrap();
}
#[tokio::test(start_paused = true)]
async fn presence_loop_recovers_when_conflict_clears() {
let stub = Arc::new(StubDirectory::rejecting(Rejection::Conflict));
let directory: Arc<dyn DirectoryClient> = stub.clone();
let cancel = CancellationToken::new();
let task = tokio::spawn(presence_loop(
Arc::clone(&directory),
info(),
Duration::from_secs(1),
cancel.clone(),
));
for _ in 0..10 {
if stub.register_calls.load(Ordering::SeqCst) >= 1 {
break;
}
tokio::task::yield_now().await;
}
assert_eq!(
stub.register_calls.load(Ordering::SeqCst),
1,
"the initial registration is rejected but the loop keeps retrying"
);
assert_eq!(
stub.register_ok.load(Ordering::SeqCst),
0,
"while rejecting, no registration reaches the Ok path"
);
stub.stop_rejecting();
tokio::time::advance(MAX_BACKOFF + Duration::from_secs(1)).await;
for _ in 0..10 {
if stub.register_ok.load(Ordering::SeqCst) >= 1 {
break;
}
tokio::task::yield_now().await;
}
assert_eq!(
stub.register_ok.load(Ordering::SeqCst),
1,
"once the conflict clears the retry must succeed, not merely re-attempt"
);
cancel.cancel();
task.await.unwrap();
}