onlyne_client/runtime/runloop/
run.rs1use super::config::{
2 ClientInit, NOT_READY_PAUSE_MS, PULL_HOLD_MS, PULL_LIMIT, PULL_PAUSE_MS, RunState,
3 SHUTDOWN_CLOSE_BUDGET,
4};
5use super::link::{link_loop, transient};
6use super::sessions::{accept_delivery, outcome_loop};
7use crate::session::adapter_socket::AdapterSocket;
8use crate::session::dispatch::{self, ClientLink, DispatchState};
9use anyhow::{Result, anyhow};
10use onlyne_config::layout::RoleWorkspace;
11use onlyne_proto::{AckArgs, ClientOp, ControlOp, Delivery, PullArgs, PullReply};
12use onlyne_store::ClientStore;
13use onlyne_wire::socket::{registration_path, remove_registration};
14use std::path::{Path, PathBuf};
15use std::sync::atomic::Ordering;
16use std::time::Duration;
17use tokio::time::sleep;
18
19pub async fn run(init: ClientInit) -> Result<()> {
28 let workspace = RoleWorkspace::resolve(&init.workspace);
29 let store = ClientStore::open(workspace.client_db_path())?;
30 let state = RunState::new(&init, store)?;
31 let mut acceptor = tokio::spawn(acceptor(init.clone(), state.clone()));
32 let closing = tokio::spawn(close_on_signal(
33 state.dispatch.clone(),
34 workspace.root().to_path_buf(),
35 ));
36 let mut outcomes = tokio::spawn(outcome_loop(state.clone()));
37 let outcome = tokio::select! {
38 link = link_loop(&init, &state) => link,
39 served = &mut acceptor => match served {
40 Ok(Ok(())) => Err(anyhow!("the adapter surface stopped serving")),
41 Ok(Err(error)) => Err(error),
42 Err(error) => Err(anyhow!("adapter acceptor task ended: {error}")),
43 },
44 pumped = &mut outcomes => match pumped {
45 Ok(Ok(())) => Err(anyhow!("the session outcome pump stopped")),
46 Ok(Err(error)) => Err(error),
47 Err(error) => Err(anyhow!("session outcome pump task ended: {error}")),
48 },
49 };
50 acceptor.abort();
51 closing.abort();
52 outcomes.abort();
53 deregister_client(workspace.root());
57 outcome
58}
59
60fn deregister_client(root: &Path) {
66 if let Err(error) = remove_registration(root) {
67 tracing::warn!(
68 error = %error,
69 registration = %registration_path(root).display(),
70 "could not remove the client registration"
71 );
72 }
73}
74
75#[cfg(unix)]
83pub(super) async fn close_on_signal(dispatch: DispatchState, root: PathBuf) {
84 use tokio::signal::unix::{SignalKind, signal};
85 let mut terminate = match signal(SignalKind::terminate()) {
86 Ok(stream) => stream,
87 Err(error) => {
88 tracing::warn!(error = %error, "SIGTERM handler was not installed");
89 return;
90 }
91 };
92 let mut interrupt = match signal(SignalKind::interrupt()) {
93 Ok(stream) => stream,
94 Err(error) => {
95 tracing::warn!(error = %error, "SIGINT handler was not installed");
96 return;
97 }
98 };
99 tokio::select! {
100 _ = terminate.recv() => tracing::info!("SIGTERM: closing live sessions"),
101 _ = interrupt.recv() => tracing::info!("SIGINT: closing live sessions"),
102 }
103 dispatch::close_all(
104 &dispatch,
105 crate::backend::CloseReason::Shutdown,
106 SHUTDOWN_CLOSE_BUDGET,
107 );
108 deregister_client(&root);
109 std::process::exit(0);
110}
111
112#[cfg(windows)]
118pub(super) async fn close_on_signal(dispatch: DispatchState, root: PathBuf) {
119 let mut interrupt = match tokio::signal::windows::ctrl_c() {
122 Ok(stream) => stream,
123 Err(error) => {
124 tracing::warn!(error = %error, "Ctrl-C handler was not installed");
125 return;
126 }
127 };
128 interrupt.recv().await;
129 tracing::info!("Ctrl-C: closing live sessions");
130 dispatch::close_all(
131 &dispatch,
132 crate::backend::CloseReason::Shutdown,
133 SHUTDOWN_CLOSE_BUDGET,
134 );
135 deregister_client(&root);
136 std::process::exit(0);
137}
138
139pub(super) async fn acceptor(init: ClientInit, state: RunState) -> Result<()> {
147 let workspace = RoleWorkspace::resolve(&init.workspace);
148 let cluster = state
149 .store
150 .config("cluster")
151 .ok()
152 .flatten()
153 .unwrap_or_default();
154 let socket = AdapterSocket {
155 workspace: workspace.root().to_path_buf(),
156 role: init.role.clone(),
157 cluster,
158 server: init.server.clone(),
159 dispatch: state.dispatch.clone(),
160 };
161 let (listener, _endpoint) = socket.bind().await?;
162 socket.accept_loop(listener).await
163}
164
165pub(super) async fn pull_ack_loop(
168 init: ClientInit,
169 link: ClientLink,
170 state: RunState,
171) -> Result<()> {
172 loop {
173 if !state.accept_new.load(Ordering::SeqCst) {
174 sleep(Duration::from_millis(PULL_PAUSE_MS)).await;
175 continue;
176 }
177 let control_only = !state.dispatch.has_capacity();
183 let reply = match link
184 .request(ClientOp::Pull(PullArgs {
185 role: Some(init.role.clone()),
186 limit: PULL_LIMIT,
187 hold_ms: Some(PULL_HOLD_MS),
188 control_only: control_only.then_some(true),
189 }))
190 .await
191 {
192 Ok(reply) => reply,
193 Err(error) if transient(&error) => {
194 sleep(Duration::from_millis(NOT_READY_PAUSE_MS)).await;
195 continue;
196 }
197 Err(error) => return Err(anyhow!(error)),
198 };
199 if !reply.ok {
200 tracing::warn!(error = ?reply.error, "pull refused");
201 sleep(Duration::from_millis(PULL_PAUSE_MS)).await;
202 continue;
203 }
204 let Some(data) = reply.data else { continue };
205 let pulled: PullReply = serde_json::from_value(data)?;
206 for delivery in pulled.deliveries {
207 accept_delivery(&state, &delivery).await;
208 }
209 state.set_cursor(pulled.seq);
210 sleep(Duration::from_millis(PULL_PAUSE_MS)).await;
211 }
212}
213
214pub(super) async fn settle_control(state: &RunState, delivery: &Delivery) {
220 let ack = |accepted: bool, reason: Option<String>| AckArgs {
221 msg_id: delivery.msg_id.clone(),
222 op_id: None,
223 accepted,
224 reason,
225 };
226 let Some(op) = delivery.envelope.control.clone() else {
227 tracing::warn!(msg_id = %delivery.msg_id, "control delivery carried no command");
228 state.dispatch.push_settled(ack(
229 false,
230 Some("control delivery carried no control op".to_string()),
231 ));
232 return;
233 };
234 match dispatch::on_control(&state.dispatch, &op).await {
235 Ok(held) => {
236 tracing::info!(
237 op = op.name(),
238 task = %op.task_id(),
239 held,
240 "control command applied"
241 );
242 state.dispatch.push_settled(ack(true, None));
243 if held && matches!(op, ControlOp::Recycle { .. } | ControlOp::Cancel { .. }) {
244 if let Err(error) = dispatch::sync_session(&state.dispatch, op.task_id()).await {
247 tracing::warn!(
248 error = %error,
249 task = %op.task_id(),
250 "a closed control session was not published"
251 );
252 }
253 }
254 }
255 Err(error) => {
256 tracing::warn!(error = %error, op = op.name(), "control command refused");
257 state
258 .dispatch
259 .push_settled(ack(false, Some(error.to_string())));
260 }
261 }
262}