use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use crate::activity::ActivityRegistry;
use crate::config::WorkerConfig;
use crate::error::WorkerError;
use crate::runtime::liminal::LiminalActivityWorker;
use crate::runtime::liminal_redial::{
CandidateCursor, RedialBackoff, RedialError, ServeResult, run_redial_loop,
};
pub fn serve_with_redial<Ready>(
candidates: Vec<String>,
config: &WorkerConfig,
registry: &Arc<ActivityRegistry>,
initial_backoff: Duration,
max_backoff: Duration,
stop: &AtomicBool,
mut on_first_ready: Ready,
) -> Result<(), WorkerError>
where
Ready: FnMut() + Send,
{
let mut cursor = CandidateCursor::new(candidates).map_err(redial_setup_error)?;
let mut backoff = RedialBackoff::new(initial_backoff, max_backoff);
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(WorkerError::registration)?;
let mut announced_ready = false;
let connect = |address: &str| -> Result<LiminalActivityWorker, WorkerError> {
LiminalActivityWorker::connect(address, config, Arc::clone(registry))
};
let serve = |worker: LiminalActivityWorker| -> ServeResult {
if !announced_ready {
announced_ready = true;
on_first_ready();
}
runtime.block_on(worker.serve_until_drop(|| stop.load(Ordering::Relaxed)))
};
run_redial_loop(
&mut cursor,
&mut backoff,
connect,
serve,
std::thread::sleep,
|| stop.load(Ordering::Relaxed),
WorkerError::is_retryable,
)
}
fn redial_setup_error(error: RedialError) -> WorkerError {
WorkerError::registration(error)
}