use super::config::{
ClientInit, NOT_READY_PAUSE_MS, PULL_HOLD_MS, PULL_LIMIT, PULL_PAUSE_MS, RunState,
SHUTDOWN_CLOSE_BUDGET,
};
use super::link::{link_loop, transient};
use super::sessions::{accept_delivery, outcome_loop};
use crate::session::adapter_socket::AdapterSocket;
use crate::session::dispatch::{self, ClientLink, DispatchState};
use anyhow::{Result, anyhow};
use onlyne_layout::RoleWorkspace;
use onlyne_proto::{AckArgs, ClientOp, ControlOp, Delivery, PullArgs, PullReply};
use onlyne_store::ClientStore;
use std::sync::atomic::Ordering;
use std::time::Duration;
use tokio::time::sleep;
pub async fn run(init: ClientInit) -> Result<()> {
let workspace = RoleWorkspace::resolve(&init.workspace);
let store = ClientStore::open(workspace.client_db_path())?;
let state = RunState::new(&init, store)?;
let mut acceptor = tokio::spawn(acceptor(init.clone(), state.clone()));
let closing = tokio::spawn(close_on_signal(state.dispatch.clone()));
let mut outcomes = tokio::spawn(outcome_loop(state.clone()));
let outcome = tokio::select! {
link = link_loop(&init, &state) => link,
served = &mut acceptor => match served {
Ok(Ok(())) => Err(anyhow!("the adapter surface stopped serving")),
Ok(Err(error)) => Err(error),
Err(error) => Err(anyhow!("adapter acceptor task ended: {error}")),
},
pumped = &mut outcomes => match pumped {
Ok(Ok(())) => Err(anyhow!("the session outcome pump stopped")),
Ok(Err(error)) => Err(error),
Err(error) => Err(anyhow!("session outcome pump task ended: {error}")),
},
};
acceptor.abort();
closing.abort();
outcomes.abort();
outcome
}
#[cfg(unix)]
pub(super) async fn close_on_signal(dispatch: DispatchState) {
use tokio::signal::unix::{SignalKind, signal};
let mut terminate = match signal(SignalKind::terminate()) {
Ok(stream) => stream,
Err(error) => {
tracing::warn!(error = %error, "SIGTERM handler was not installed");
return;
}
};
let mut interrupt = match signal(SignalKind::interrupt()) {
Ok(stream) => stream,
Err(error) => {
tracing::warn!(error = %error, "SIGINT handler was not installed");
return;
}
};
tokio::select! {
_ = terminate.recv() => tracing::info!("SIGTERM: closing live sessions"),
_ = interrupt.recv() => tracing::info!("SIGINT: closing live sessions"),
}
dispatch::close_all(
&dispatch,
onlyne_session::CloseReason::Shutdown,
SHUTDOWN_CLOSE_BUDGET,
);
std::process::exit(0);
}
#[cfg(windows)]
pub(super) async fn close_on_signal(dispatch: DispatchState) {
let mut interrupt = match tokio::signal::windows::ctrl_c() {
Ok(stream) => stream,
Err(error) => {
tracing::warn!(error = %error, "Ctrl-C handler was not installed");
return;
}
};
interrupt.recv().await;
tracing::info!("Ctrl-C: closing live sessions");
dispatch::close_all(
&dispatch,
onlyne_session::CloseReason::Shutdown,
SHUTDOWN_CLOSE_BUDGET,
);
std::process::exit(0);
}
pub(super) async fn acceptor(init: ClientInit, state: RunState) -> Result<()> {
let workspace = RoleWorkspace::resolve(&init.workspace);
let cluster = state
.store
.config("cluster")
.ok()
.flatten()
.unwrap_or_default();
let socket = AdapterSocket {
workspace: workspace.root().to_path_buf(),
role: init.role.clone(),
cluster,
server: init.server.clone(),
dispatch: state.dispatch.clone(),
};
let (listener, _endpoint) = socket.bind().await?;
socket.accept_loop(listener).await
}
pub(super) async fn pull_ack_loop(
init: ClientInit,
link: ClientLink,
state: RunState,
) -> Result<()> {
loop {
if !state.accept_new.load(Ordering::SeqCst) {
sleep(Duration::from_millis(PULL_PAUSE_MS)).await;
continue;
}
let control_only = !state.dispatch.has_capacity();
let reply = match link
.request(ClientOp::Pull(PullArgs {
role: Some(init.role.clone()),
limit: PULL_LIMIT,
hold_ms: Some(PULL_HOLD_MS),
control_only: control_only.then_some(true),
}))
.await
{
Ok(reply) => reply,
Err(error) if transient(&error) => {
sleep(Duration::from_millis(NOT_READY_PAUSE_MS)).await;
continue;
}
Err(error) => return Err(anyhow!(error)),
};
if !reply.ok {
tracing::warn!(error = ?reply.error, "pull refused");
sleep(Duration::from_millis(PULL_PAUSE_MS)).await;
continue;
}
let Some(data) = reply.data else { continue };
let pulled: PullReply = serde_json::from_value(data)?;
for delivery in pulled.deliveries {
accept_delivery(&state, &delivery).await;
}
state.set_cursor(pulled.seq);
sleep(Duration::from_millis(PULL_PAUSE_MS)).await;
}
}
pub(super) async fn settle_control(state: &RunState, delivery: &Delivery) {
let ack = |accepted: bool, reason: Option<String>| AckArgs {
msg_id: delivery.msg_id.clone(),
op_id: None,
accepted,
reason,
};
let Some(op) = delivery.envelope.control.clone() else {
tracing::warn!(msg_id = %delivery.msg_id, "control delivery carried no command");
state.dispatch.push_settled(ack(
false,
Some("control delivery carried no control op".to_string()),
));
return;
};
match dispatch::on_control(&state.dispatch, &op).await {
Ok(held) => {
tracing::info!(
op = op.name(),
task = %op.task_id(),
held,
"control command applied"
);
state.dispatch.push_settled(ack(true, None));
if held && matches!(op, ControlOp::Recycle { .. } | ControlOp::Cancel { .. }) {
if let Err(error) = dispatch::sync_session(&state.dispatch, op.task_id()).await {
tracing::warn!(
error = %error,
task = %op.task_id(),
"a closed control session was not published"
);
}
}
}
Err(error) => {
tracing::warn!(error = %error, op = op.name(), "control command refused");
state
.dispatch
.push_settled(ack(false, Some(error.to_string())));
}
}
}
#[cfg(test)]
mod tests;