use std::sync::Arc;
use std::time::Duration;
use tokio_util::sync::CancellationToken;
use tracing::{error, info, warn};
use crate::bridge::config::BridgeApp;
mod backoff;
mod handle;
mod send;
mod session;
use send::{GetUpdatesOutcome, HubClient};
use session::SessionDispatcher;
use crate::bridge::ApprovalBroker;
#[cfg(test)]
use backoff::{backoff_for, backoff_for_test, MAX_BACKOFF_SECS};
#[cfg(test)]
use send::{
classify_sendoutcome, parse_sendoutcome, run_partial_forward_loop, sanitize_errmsg,
send_final_with_retry, ReplySender, SendOutcome,
};
#[cfg(test)]
use session::session_dispatch_key;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BridgeStop {
TokenRejected,
FatalCliError(String),
Shutdown,
}
pub async fn run_bridge_with_shutdown(
hub_url: String,
token: String,
app: BridgeApp,
shutdown: CancellationToken,
) -> BridgeStop {
let client = match HubClient::new(hub_url, token) {
Ok(c) => c,
Err(e) => return BridgeStop::FatalCliError(e.to_string()),
};
let app = Arc::new(app);
let (stop_tx, mut stop_rx) = tokio::sync::watch::channel(None::<BridgeStop>);
let dispatcher = Arc::new(SessionDispatcher::new(
client.clone(),
Arc::clone(&app),
stop_tx,
shutdown.clone(),
ApprovalBroker::new(),
));
let mut buf = String::new();
let mut backoff_secs: u64 = 3;
const MAX_BACKOFF_SECS: u64 = 60;
{
let dispatcher_weak = Arc::downgrade(&dispatcher);
let shutdown_clone = shutdown.clone();
tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(60));
loop {
tokio::select! {
biased;
_ = shutdown_clone.cancelled() => return,
_ = interval.tick() => {
if let Some(d) = dispatcher_weak.upgrade() {
d.evict_closed_senders();
} else {
return;
}
}
}
}
});
}
info!(
routing = %app.routing_label(),
profiles = ?app.profile_names(),
"ilink-hub-bridge connected; waiting for getupdates"
);
loop {
if stop_rx.has_changed().unwrap_or(false) {
if let Some(reason) = stop_rx.borrow_and_update().clone() {
return reason;
}
}
let getupdates_fut = client.getupdates(&mut buf);
let resp = tokio::select! {
biased;
_ = shutdown.cancelled() => {
tokio::time::sleep(Duration::from_secs(2)).await;
return BridgeStop::Shutdown;
}
r = getupdates_fut => match r {
Ok(GetUpdatesOutcome::Ok(r)) => {
backoff_secs = 3;
r
}
Ok(GetUpdatesOutcome::TokenRejected) => return BridgeStop::TokenRejected,
Err(e) => {
error!(error = %e, backoff_secs, "getupdates failed; retrying with backoff");
let sleep = tokio::time::sleep(Duration::from_secs(backoff_secs));
backoff_secs = (backoff_secs * 2).min(MAX_BACKOFF_SECS);
tokio::select! {
biased;
_ = shutdown.cancelled() => {
tokio::time::sleep(Duration::from_secs(2)).await;
return BridgeStop::Shutdown;
}
_ = sleep => {}
}
continue;
}
},
};
if resp.ret != Some(0) {
warn!(
ret = ?resp.ret,
errcode = ?resp.errcode,
errmsg = ?resp.errmsg,
"getupdates returned non-zero ret"
);
}
for msg in resp.msgs.unwrap_or_default() {
dispatcher.dispatch(msg).await;
}
}
}
pub async fn run_bridge(hub_url: String, token: String, app: BridgeApp) -> BridgeStop {
run_bridge_with_shutdown(hub_url, token, app, CancellationToken::new()).await
}
#[cfg(test)]
mod tests;