onlyne_client/runtime/runloop/
config.rs1use 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
16pub const DEFAULT_INTENT_ATTEMPTS: u32 = 3;
18pub const RECONNECT_LADDER_SECONDS: [u64; 7] = [1, 2, 4, 8, 16, 32, 60];
20pub const PULL_PAUSE_MS: u64 = 200;
22pub const FLUSH_PAUSE_MS: u64 = 200;
24pub const PULL_HOLD_MS: u64 = 1_000;
26pub const PULL_LIMIT: u32 = 32;
28pub const READINESS_POLL_MS: u64 = 250;
30pub const SHUTDOWN_CLOSE_BUDGET: Duration = Duration::from_secs(8);
33pub const EVENT_CURSOR_KEY: &str = "event_seq";
35pub const NOT_READY_PAUSE_MS: u64 = 200;
37pub const OUTCOME_POLL_MS: u64 = 100;
39
40pub fn default_intent_backoff() -> Vec<u64> {
42 vec![1_000, 2_000, 4_000]
43}
44
45pub 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 pub orca_worktree: String,
63 pub stall_report_secs: u64,
66 pub reconnect_grace_secs: u64,
69 pub placement: Option<SessionPlacement>,
76 pub acp: onlyne_config::AcpSection,
80 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 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 pub fn with_reconnect_grace_secs(mut self, secs: u64) -> Self {
120 self.reconnect_grace_secs = secs;
121 self
122 }
123 pub fn with_placement(mut self, placement: Option<SessionPlacement>) -> Self {
125 self.placement = placement;
126 self
127 }
128 pub fn with_acp(mut self, acp: onlyne_config::AcpSection) -> Self {
130 self.acp = acp;
131 self
132 }
133 pub fn with_session(mut self, session: onlyne_config::SessionPolicy) -> Self {
135 self.session = session;
136 self
137 }
138}
139
140pub 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#[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 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 pub reconnect_grace_secs: u64,
196 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 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 pub(super) async fn adopt(&self, welcome: &Welcome) {
252 let slice = crate::session::slice::RoleSlice::from_welcome(welcome);
253 self.install_runtime(slice.drive);
257 self.dispatch.reconfigure(slice);
258 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 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 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 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}