Skip to main content

onlyne_client/runtime/runloop/
config.rs

1use crate::backend::{
2    AcpOptions, ProcessRunner, Runner, SessionBackend, SessionPlacement, WorktreePolicy,
3    backend_for, detect_placement, process_env,
4};
5use crate::runtime::intent::IntentMachine;
6use crate::session::dispatch::DispatchState;
7use anyhow::Result;
8use onlyne_net::backoff::Backoff;
9use onlyne_proto::Welcome;
10use onlyne_store::ClientStore;
11use std::path::PathBuf;
12use std::sync::{Arc, atomic::AtomicBool};
13use std::time::Duration;
14use tokio::sync::Mutex;
15
16/// Attempt ceiling used until `welcome` carries the role's own value.
17pub const DEFAULT_INTENT_ATTEMPTS: u32 = 3;
18/// Reconnect ladder in seconds (§6: 1/2/4/8/16/32/60).
19pub const RECONNECT_LADDER_SECONDS: [u64; 7] = [1, 2, 4, 8, 16, 32, 60];
20/// Pause between pull attempts that returned nothing.
21pub const PULL_PAUSE_MS: u64 = 200;
22/// Pause between intent flush passes.
23pub const FLUSH_PAUSE_MS: u64 = 200;
24/// Long-poll window the client asks the server for.
25pub const PULL_HOLD_MS: u64 = 1_000;
26/// Deliveries drained per pull.
27pub const PULL_LIMIT: u32 = 32;
28/// Poll interval for the readiness watcher.
29pub const READINESS_POLL_MS: u64 = 250;
30/// Bound on the session sweep the SIGTERM/SIGINT handler runs, short enough
31/// that an operator's own grace period still sees the process leave.
32pub const SHUTDOWN_CLOSE_BUDGET: Duration = Duration::from_secs(8);
33/// Key holding the durable event cursor in `config_cache`.
34pub const EVENT_CURSOR_KEY: &str = "event_seq";
35/// Retry delay for a request that the transport answered `NotReady`.
36pub const NOT_READY_PAUSE_MS: u64 = 200;
37/// Poll cadence for terminal facts emitted by a self-driven session backend.
38pub const OUTCOME_POLL_MS: u64 = 100;
39
40/// Ladder used until `welcome` carries the role's own values.
41pub fn default_intent_backoff() -> Vec<u64> {
42    vec![1_000, 2_000, 4_000]
43}
44
45/// Reconnect ladder as durations, capped at the last rung.
46pub fn reconnect_backoff() -> Backoff {
47    Backoff::with_limits(
48        Duration::from_secs(RECONNECT_LADDER_SECONDS[0]),
49        Duration::from_secs(RECONNECT_LADDER_SECONDS[6]),
50    )
51}
52
53#[derive(Debug, Clone)]
54pub struct ClientInit {
55    pub workspace: PathBuf,
56    pub role: String,
57    pub server: String,
58    pub key_path: PathBuf,
59    pub cert_pin: String,
60    /// The workspace config's `[orca] worktree` value: `host`, `inherit`, or a
61    /// literal Orca worktree selector. Only an Orca session backend reads it.
62    pub orca_worktree: String,
63    /// Seconds a running session may sit without Applied progress before a stall
64    /// fault is reported. Zero disables the report.
65    pub stall_report_secs: u64,
66    /// Seconds a dropped plugin connection may stay away before this client
67    /// retires the task-free session it left behind. Zero disables the sweep.
68    pub reconnect_grace_secs: u64,
69    /// The placement this run was told to use: `onlyne-client run` passes the
70    /// workspace `config.toml`'s `placement` key, and an embedding passes
71    /// whatever it resolved — the scenario suite passes the in-process
72    /// `fake` runtime. `None` probes orca, then zellij, and falls back to
73    /// `headless`. `ONLYNE_BACKEND` in the process environment takes precedence
74    /// when it names a placement.
75    pub placement: Option<SessionPlacement>,
76    /// The workspace config's `[acp]` table. Only the ACP session backend reads
77    /// it: the mode, model and reasoning effort handed to the agent when a
78    /// session opens, and what to answer when the agent asks for permission.
79    pub acp: onlyne_config::AcpSection,
80    /// The workspace config's `[client.session]` table: which deliveries one
81    /// session of this role serves, and how long an idle one may wait before
82    /// this client releases its process (plan §10).
83    pub session: onlyne_config::SessionPolicy,
84}
85
86impl ClientInit {
87    pub fn new(
88        workspace: impl Into<PathBuf>,
89        role: impl Into<String>,
90        server: impl Into<String>,
91        key_path: impl Into<PathBuf>,
92        cert_pin: impl Into<String>,
93    ) -> Self {
94        Self {
95            workspace: workspace.into(),
96            role: role.into(),
97            server: server.into(),
98            key_path: key_path.into(),
99            cert_pin: cert_pin.into(),
100            orca_worktree: "host".to_string(),
101            stall_report_secs: onlyne_config::DEFAULT_STALL_REPORT_SECS,
102            reconnect_grace_secs: onlyne_config::DEFAULT_RECONNECT_GRACE_SECS,
103            placement: None,
104            acp: onlyne_config::AcpSection::default(),
105            session: onlyne_config::SessionPolicy::default(),
106        }
107    }
108
109    /// Adopt the `[orca] worktree` policy the workspace config carries.
110    pub fn with_orca_worktree(mut self, worktree: impl Into<String>) -> Self {
111        self.orca_worktree = worktree.into();
112        self
113    }
114    pub fn with_stall_report_secs(mut self, secs: u64) -> Self {
115        self.stall_report_secs = secs;
116        self
117    }
118    /// Adopt the `[client] reconnect_grace_secs` value the workspace config carries.
119    pub fn with_reconnect_grace_secs(mut self, secs: u64) -> Self {
120        self.reconnect_grace_secs = secs;
121        self
122    }
123    /// Adopt the `placement` this run was told to use, if it was told one.
124    pub fn with_placement(mut self, placement: Option<SessionPlacement>) -> Self {
125        self.placement = placement;
126        self
127    }
128    /// Adopt the `[acp]` table the workspace config carries.
129    pub fn with_acp(mut self, acp: onlyne_config::AcpSection) -> Self {
130        self.acp = acp;
131        self
132    }
133    /// Adopt the `[client.session]` table the workspace config carries.
134    pub fn with_session(mut self, session: onlyne_config::SessionPolicy) -> Self {
135        self.session = session;
136        self
137    }
138}
139
140/// The `[acp]` table in the shape a session backend can read without a config
141/// dependency. The only decision made here is the one the backend acts on: the
142/// permission word, already validated by the config loader, becomes whether this
143/// client grants an agent's request.
144pub fn acp_options(acp: &onlyne_config::AcpSection) -> AcpOptions {
145    AcpOptions {
146        mode: acp.mode.clone(),
147        model: acp.model.clone(),
148        reasoning_effort: acp.reasoning_effort.clone(),
149        allow_permissions: acp.permission == "allow",
150    }
151}
152
153/// What picking a session backend needs from this machine.
154///
155/// The placement is the machine's half of the pair and is known at startup; the
156/// drive is the runtime's half and arrives later, with `welcome`. Keeping the
157/// two apart here is the point of the slice: the backend is chosen when both
158/// halves are known, instead of by one fused value read out of one file.
159#[derive(Clone)]
160pub struct BackendSelector {
161    pub placement: SessionPlacement,
162    pub runner: Arc<dyn Runner>,
163    pub worktree: WorktreePolicy,
164    pub acp: AcpOptions,
165}
166
167impl BackendSelector {
168    /// The backend this role's drive and this machine's placement select.
169    ///
170    /// Refuses a pair the rule does not allow — `acp` anywhere but `headless` —
171    /// so a caller never gets a backend that would run the agent where its
172    /// channel cannot follow it.
173    pub fn build(&self, drive: onlyne_config::Drive) -> Result<Arc<dyn SessionBackend>> {
174        backend_for(
175            drive,
176            self.placement,
177            Arc::clone(&self.runner),
178            self.worktree.clone(),
179            &self.acp,
180        )
181        .map(Arc::from)
182    }
183}
184
185#[derive(Clone)]
186pub struct RunState {
187    pub accept_new: Arc<AtomicBool>,
188    pub store: ClientStore,
189    pub intents: Arc<parking_lot::Mutex<IntentMachine>>,
190    pub dispatch: DispatchState,
191    pub welcome: Arc<Mutex<Option<Welcome>>>,
192    pub stall_report_secs: u64,
193    /// Seconds a dropped plugin connection may stay away before this client
194    /// retires the task-free session it left behind. Zero disables the sweep.
195    pub reconnect_grace_secs: u64,
196    /// The machine's half of the runtime pair: the placement a role's drive is
197    /// resolved against, and the settings only a backend reads.
198    pub selector: BackendSelector,
199}
200
201impl RunState {
202    pub fn new(init: &ClientInit, store: ClientStore) -> Result<Self> {
203        let detected = detect_placement(&process_env(), init.placement)?;
204        let selector = BackendSelector {
205            placement: detected.placement,
206            runner: Arc::new(ProcessRunner),
207            worktree: WorktreePolicy::from_config(&init.orca_worktree),
208            acp: acp_options(&init.acp),
209        };
210        tracing::info!(
211            placement = %selector.placement,
212            source = ?detected.source,
213            explicit = ?detected.explicit,
214            "placement resolved"
215        );
216        // The backend a role's drive selects is installed when the drive
217        // arrives with `welcome`. Until then the default drive's backend is in
218        // place, because a dispatcher always holds one and no session can open
219        // before the first `hello`.
220        let backend = selector.build(onlyne_config::Drive::Plugin)?;
221        let dispatch = DispatchState::new(
222            init.role.clone(),
223            init.workspace.clone(),
224            Vec::new(),
225            1,
226            Arc::clone(&backend),
227            store.clone(),
228        )
229        .with_session_policy(init.session.clone())
230        .with_placement(selector.placement);
231        dispatch.set_drive(onlyne_config::Drive::Plugin, None);
232        let intents = IntentMachine::new(
233            store.clone(),
234            DEFAULT_INTENT_ATTEMPTS,
235            default_intent_backoff(),
236        );
237        let accept_new = dispatch.accept_new();
238        Ok(Self {
239            accept_new,
240            store,
241            intents: Arc::new(parking_lot::Mutex::new(intents)),
242            dispatch,
243            welcome: Arc::new(Mutex::new(None)),
244            stall_report_secs: init.stall_report_secs,
245            reconnect_grace_secs: init.reconnect_grace_secs,
246            selector,
247        })
248    }
249
250    /// Adopt the role slice the server sent with `welcome`.
251    pub(super) async fn adopt(&self, welcome: &Welcome) {
252        let slice = crate::session::slice::RoleSlice::from_welcome(welcome);
253        // The backend a drive selects is installed before the slice lands, so
254        // the command that arrives with the slice is never handed to the
255        // backend the previous drive left behind.
256        self.install_runtime(slice.drive);
257        self.dispatch.reconfigure(slice);
258        // The topology name is the address the host backends group sessions
259        // under, so it is recorded with the rest of what the server says about
260        // this role. `welcome.cluster` is the server's own `[server] name`.
261        self.dispatch.set_topology(&welcome.cluster);
262        {
263            let mut intents = self.intents.lock();
264            if let Some(attempts) = welcome.intent_attempts {
265                intents.attempts = attempts;
266            }
267            if let Some(ladder) = welcome
268                .intent_backoff_ms
269                .as_ref()
270                .filter(|ladder| !ladder.is_empty())
271            {
272                intents.backoff_ms = ladder.clone();
273            }
274        }
275        *self.welcome.lock().await = Some(welcome.clone());
276    }
277
278    /// Install the backend this role's drive selects on this machine.
279    ///
280    /// Runs on every `welcome` and on every spec reload, and does nothing when
281    /// the drive has not moved. Three answers:
282    ///
283    /// * the backend is installed, and the registration is republished with the
284    ///   name it reports;
285    /// * the drive moved while this role holds live sessions, so nothing moves:
286    ///   their panes, tabs, and children are the installed backend's to close
287    ///   and to probe, and the next attempt lands once the role is quiet;
288    /// * the drive and this machine's placement cannot be paired at all, which
289    ///   is recorded so every delivery meets the sentence instead of a backend
290    ///   the previous drive left behind.
291    pub(super) fn install_runtime(&self, drive: onlyne_config::Drive) {
292        if self.dispatch.drive() == Some(drive) {
293            return;
294        }
295        match self.selector.build(drive) {
296            Ok(backend) => {
297                if self.dispatch.set_backend(backend) {
298                    self.dispatch.set_drive(drive, None);
299                    tracing::info!(
300                        drive = %drive,
301                        placement = %self.selector.placement,
302                        backend = %self.dispatch.session_backend(),
303                        "session backend selected"
304                    );
305                    self.republish_registration();
306                } else {
307                    tracing::warn!(
308                        drive = %drive,
309                        placement = %self.selector.placement,
310                        live = self.dispatch.session_count(),
311                        held = %self.dispatch.session_backend(),
312                        "the role's drive changed while it holds live sessions: they keep the \
313                         backend they were opened under, and the new drive lands once the role \
314                         is quiet"
315                    );
316                }
317            }
318            Err(error) => {
319                tracing::error!(
320                    drive = %drive,
321                    placement = %self.selector.placement,
322                    error = %error,
323                    "the drive this role's spec declares cannot run under this machine's \
324                     placement; every delivery for this role will be refused with that sentence"
325                );
326                self.dispatch.set_drive(drive, Some(error.to_string()));
327            }
328        }
329    }
330
331    /// Republish this client's registration.
332    ///
333    /// The bind wrote one before the role's drive was known, so the `runtime`
334    /// field it carries is the default drive's backend. An external runtime's
335    /// plugin reads that file to find the clients it serves, and a stale name
336    /// there is the same class of fact this slice exists to stop trusting.
337    fn republish_registration(&self) {
338        if let Err(error) = crate::session::adapter_socket::republish_registration(
339            &self.dispatch.workspace(),
340            &self.dispatch.role(),
341            self.dispatch.session_backend(),
342            self.dispatch.placement_name(),
343        ) {
344            tracing::warn!(error = %error, "the client registration was not republished");
345        }
346    }
347
348    /// Durable event cursor for the next `subscribe`. Zero asks the server for
349    /// its current head.
350    pub(super) fn cursor(&self) -> u64 {
351        self.store
352            .config(EVENT_CURSOR_KEY)
353            .ok()
354            .flatten()
355            .and_then(|value| value.parse().ok())
356            .unwrap_or(0)
357    }
358
359    pub(super) fn set_cursor(&self, seq: u64) {
360        if let Err(error) = self.store.put_config(EVENT_CURSOR_KEY, &seq.to_string()) {
361            tracing::warn!(error = %error, "event cursor was not stored");
362        }
363    }
364}