Skip to main content

onlyne_client/runtime/runloop/
run.rs

1use 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
19/// Run the role runtime until a permanent handshake failure or a dead local
20/// surface stops it.
21///
22/// The acceptor task owns the workspace socket bind, and its failure ends this
23/// run the same way a permanent link failure does: a client that keeps a TLS
24/// link while its socket is unbound reads as connected from the server while
25/// every verb the workspace issues fails, so the exit status and the message on
26/// stderr are the operator's only notice.
27pub 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    // The registration names this process as the client for this workspace, so
54    // the run's end takes it with it on the path that returns. The signal path
55    // leaves through `process::exit` and takes its own copy.
56    deregister_client(workspace.root());
57    outcome
58}
59
60/// Remove the registration the client published for `root`.
61///
62/// A file left naming a dead pid is what an external runtime's plugin reads as
63/// a live client, and this is the one place every non-signal exit passes
64/// through, so the removal is not optional on the way out.
65fn 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/// Close live sessions when the operator stops the client.
76///
77/// `SIGTERM` ends the foreground client, and the default disposition would
78/// kill the process with every tab it opened still running: the resources
79/// would outlive the only thing that can address them. Each session closes with
80/// [`crate::backend::CloseReason::Shutdown`] first, so the backend record and
81/// the plugin-facing tab map end truthfully.
82#[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/// Close live sessions when the operator stops the client.
113///
114/// Windows has no SIGTERM; operators use `onlyne shutdown` for a graceful
115/// daemon stop. Ctrl-C is the console interrupt, and it runs the same
116/// close_all budget the unix SIGINT path uses.
117#[cfg(windows)]
118pub(super) async fn close_on_signal(dispatch: DispatchState, root: PathBuf) {
119    // `ctrl_c()` installs synchronously and hands back the watch stream; the
120    // await belongs on `recv`, which yields once per console interrupt.
121    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
139/// Bind the adapter socket and serve it for the life of the process.
140///
141/// The bind is the first act, and its error leaves this task: a socket that
142/// never opened is a workspace whose plugins cannot mount, and only the exit
143/// status says so. Once the listener exists the surface stays up — a failed
144/// `accept` is logged and retried by [`AdapterSocket::accept_loop`] — so this
145/// task ends the run exactly when the local surface could not start.
146pub(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
165/// The pull-ack task: drain what the server queued, then settle what the local
166/// side finished.
167pub(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        // A role at `max_sessions` stops asking for work it has nowhere to run
178        // (plan §5), and the same pause must not stop it hearing the command that
179        // frees a slot: a full role is exactly the one whose operator wants to
180        // `recycle` or `focus`. Control rows still travel on the ordinary pull
181        // when the role has capacity.
182        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
214/// Apply one delivered control command and settle its row.
215///
216/// The row settles whether or not this role still holds the task it names. A
217/// command whose session already ended has nothing left to act on, and leaving
218/// the row in flight would report an operator's `control` as undelivered.
219pub(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                // A published exit releases the task's in-flight rows. The
245                // command's own row is one of them until its ack is enqueued.
246                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}