use std::sync::Arc;
use std::time::Duration;
use tokio_util::sync::CancellationToken;
use cf_system_sdks::directory::{
DirectoryClient, DirectoryInvalidArgument, DirectoryPermissionDenied, RegisterInstanceInfo,
};
const INITIAL_BACKOFF: Duration = Duration::from_millis(100);
const MAX_BACKOFF: Duration = Duration::from_secs(30);
const RE_REGISTER_INTERVAL: Duration = Duration::from_secs(30);
fn next_backoff(current: Duration) -> Duration {
current.saturating_mul(2).min(MAX_BACKOFF)
}
async fn sleep_or_cancel(dur: Duration, cancel: &CancellationToken) -> bool {
tokio::select! {
() = cancel.cancelled() => false,
() = tokio::time::sleep(dur) => true,
}
}
enum RejectionKind {
Permanent(Option<&'static str>),
Transient,
}
fn classify_rejection(err: &anyhow::Error) -> RejectionKind {
if err.is::<DirectoryInvalidArgument>() {
return RejectionKind::Permanent(Some("Check its configured labels and endpoint URIs"));
}
if err.is::<DirectoryPermissionDenied>() {
return RejectionKind::Permanent(Some(
"Check that this peer is authorized to register this gear (its ServiceAccount \
namespace / trust domain allowlist)",
));
}
RejectionKind::Transient
}
async fn register_once_with_backoff(
directory: &Arc<dyn DirectoryClient>,
info: &RegisterInstanceInfo,
cancel: &CancellationToken,
) -> bool {
let mut backoff = INITIAL_BACKOFF;
loop {
if cancel.is_cancelled() {
return false;
}
match directory.register_instance(info.clone()).await {
Ok(()) => {
tracing::info!(gear = %info.gear, instance = %info.instance_id, "registered with DirectoryService");
return true;
}
Err(e) => {
if let RejectionKind::Permanent(guidance) = classify_rejection(&e) {
tracing::error!(
gear = %info.gear,
instance = %info.instance_id,
error = %e,
guidance = guidance.unwrap_or("(no remediation hint available)"),
"registration permanently rejected by the directory; this instance \
may not be discoverable"
);
return true;
}
tracing::warn!(
gear = %info.gear,
error = %e,
backoff_ms = backoff.as_millis(),
"registration attempt failed; retrying"
);
if !sleep_or_cancel(backoff, cancel).await {
return false;
}
backoff = next_backoff(backoff);
}
}
}
}
pub(super) async fn presence_loop(
directory: Arc<dyn DirectoryClient>,
info: RegisterInstanceInfo,
heartbeat_interval: Duration,
cancel: CancellationToken,
) {
if !register_once_with_backoff(&directory, &info, &cancel).await {
return;
}
let heartbeat_interval = heartbeat_interval.max(Duration::from_secs(1));
let mut heartbeat = tokio::time::interval(heartbeat_interval);
heartbeat.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
heartbeat.tick().await;
let mut reregister = tokio::time::interval(RE_REGISTER_INTERVAL);
reregister.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
reregister.tick().await;
loop {
tokio::select! {
() = cancel.cancelled() => return,
_ = heartbeat.tick() => {
if let Err(e) = directory.send_heartbeat(&info.gear, &info.instance_id).await {
tracing::warn!(
gear = %info.gear,
instance = %info.instance_id,
error = %e,
"heartbeat failed; re-registering to self-heal"
);
if !register_once_with_backoff(&directory, &info, &cancel).await {
return;
}
} else {
tracing::trace!(gear = %info.gear, "heartbeat sent");
}
}
_ = reregister.tick() => {
if !register_once_with_backoff(&directory, &info, &cancel).await {
return;
}
}
}
}
}
#[cfg(test)]
#[cfg_attr(coverage_nightly, coverage(off))]
#[path = "oop_registration_tests.rs"]
mod tests;