use crate::backend::{
AcpOptions, ProcessRunner, Runner, SessionBackend, SessionPlacement, WorktreePolicy,
backend_for, detect_placement, process_env,
};
use crate::runtime::intent::IntentMachine;
use crate::session::dispatch::DispatchState;
use anyhow::Result;
use onlyne_net::backoff::Backoff;
use onlyne_proto::Welcome;
use onlyne_store::ClientStore;
use std::path::PathBuf;
use std::sync::{Arc, atomic::AtomicBool};
use std::time::Duration;
use tokio::sync::Mutex;
pub const DEFAULT_INTENT_ATTEMPTS: u32 = 3;
pub const RECONNECT_LADDER_SECONDS: [u64; 7] = [1, 2, 4, 8, 16, 32, 60];
pub const PULL_PAUSE_MS: u64 = 200;
pub const FLUSH_PAUSE_MS: u64 = 200;
pub const PULL_HOLD_MS: u64 = 1_000;
pub const PULL_LIMIT: u32 = 32;
pub const READINESS_POLL_MS: u64 = 250;
pub const SHUTDOWN_CLOSE_BUDGET: Duration = Duration::from_secs(8);
pub const EVENT_CURSOR_KEY: &str = "event_seq";
pub const NOT_READY_PAUSE_MS: u64 = 200;
pub const OUTCOME_POLL_MS: u64 = 100;
pub fn default_intent_backoff() -> Vec<u64> {
vec![1_000, 2_000, 4_000]
}
pub fn reconnect_backoff() -> Backoff {
Backoff::with_limits(
Duration::from_secs(RECONNECT_LADDER_SECONDS[0]),
Duration::from_secs(RECONNECT_LADDER_SECONDS[6]),
)
}
#[derive(Debug, Clone)]
pub struct ClientInit {
pub workspace: PathBuf,
pub role: String,
pub server: String,
pub key_path: PathBuf,
pub cert_pin: String,
pub orca_worktree: String,
pub stall_report_secs: u64,
pub reconnect_grace_secs: u64,
pub placement: Option<SessionPlacement>,
pub acp: onlyne_config::AcpSection,
pub session: onlyne_config::SessionPolicy,
}
impl ClientInit {
pub fn new(
workspace: impl Into<PathBuf>,
role: impl Into<String>,
server: impl Into<String>,
key_path: impl Into<PathBuf>,
cert_pin: impl Into<String>,
) -> Self {
Self {
workspace: workspace.into(),
role: role.into(),
server: server.into(),
key_path: key_path.into(),
cert_pin: cert_pin.into(),
orca_worktree: "host".to_string(),
stall_report_secs: onlyne_config::DEFAULT_STALL_REPORT_SECS,
reconnect_grace_secs: onlyne_config::DEFAULT_RECONNECT_GRACE_SECS,
placement: None,
acp: onlyne_config::AcpSection::default(),
session: onlyne_config::SessionPolicy::default(),
}
}
pub fn with_orca_worktree(mut self, worktree: impl Into<String>) -> Self {
self.orca_worktree = worktree.into();
self
}
pub fn with_stall_report_secs(mut self, secs: u64) -> Self {
self.stall_report_secs = secs;
self
}
pub fn with_reconnect_grace_secs(mut self, secs: u64) -> Self {
self.reconnect_grace_secs = secs;
self
}
pub fn with_placement(mut self, placement: Option<SessionPlacement>) -> Self {
self.placement = placement;
self
}
pub fn with_acp(mut self, acp: onlyne_config::AcpSection) -> Self {
self.acp = acp;
self
}
pub fn with_session(mut self, session: onlyne_config::SessionPolicy) -> Self {
self.session = session;
self
}
}
pub fn acp_options(acp: &onlyne_config::AcpSection) -> AcpOptions {
AcpOptions {
mode: acp.mode.clone(),
model: acp.model.clone(),
reasoning_effort: acp.reasoning_effort.clone(),
allow_permissions: acp.permission == "allow",
}
}
#[derive(Clone)]
pub struct BackendSelector {
pub placement: SessionPlacement,
pub runner: Arc<dyn Runner>,
pub worktree: WorktreePolicy,
pub acp: AcpOptions,
}
impl BackendSelector {
pub fn build(&self, drive: onlyne_config::Drive) -> Result<Arc<dyn SessionBackend>> {
backend_for(
drive,
self.placement,
Arc::clone(&self.runner),
self.worktree.clone(),
&self.acp,
)
.map(Arc::from)
}
}
#[derive(Clone)]
pub struct RunState {
pub accept_new: Arc<AtomicBool>,
pub store: ClientStore,
pub intents: Arc<parking_lot::Mutex<IntentMachine>>,
pub dispatch: DispatchState,
pub welcome: Arc<Mutex<Option<Welcome>>>,
pub stall_report_secs: u64,
pub reconnect_grace_secs: u64,
pub selector: BackendSelector,
}
impl RunState {
pub fn new(init: &ClientInit, store: ClientStore) -> Result<Self> {
let detected = detect_placement(&process_env(), init.placement)?;
let selector = BackendSelector {
placement: detected.placement,
runner: Arc::new(ProcessRunner),
worktree: WorktreePolicy::from_config(&init.orca_worktree),
acp: acp_options(&init.acp),
};
tracing::info!(
placement = %selector.placement,
source = ?detected.source,
explicit = ?detected.explicit,
"placement resolved"
);
let backend = selector.build(onlyne_config::Drive::Plugin)?;
let dispatch = DispatchState::new(
init.role.clone(),
init.workspace.clone(),
Vec::new(),
1,
Arc::clone(&backend),
store.clone(),
)
.with_session_policy(init.session.clone())
.with_placement(selector.placement);
dispatch.set_drive(onlyne_config::Drive::Plugin, None);
let intents = IntentMachine::new(
store.clone(),
DEFAULT_INTENT_ATTEMPTS,
default_intent_backoff(),
);
let accept_new = dispatch.accept_new();
Ok(Self {
accept_new,
store,
intents: Arc::new(parking_lot::Mutex::new(intents)),
dispatch,
welcome: Arc::new(Mutex::new(None)),
stall_report_secs: init.stall_report_secs,
reconnect_grace_secs: init.reconnect_grace_secs,
selector,
})
}
pub(super) async fn adopt(&self, welcome: &Welcome) {
let slice = crate::session::slice::RoleSlice::from_welcome(welcome);
self.install_runtime(slice.drive);
self.dispatch.reconfigure(slice);
self.dispatch.set_topology(&welcome.cluster);
{
let mut intents = self.intents.lock();
if let Some(attempts) = welcome.intent_attempts {
intents.attempts = attempts;
}
if let Some(ladder) = welcome
.intent_backoff_ms
.as_ref()
.filter(|ladder| !ladder.is_empty())
{
intents.backoff_ms = ladder.clone();
}
}
*self.welcome.lock().await = Some(welcome.clone());
}
pub(super) fn install_runtime(&self, drive: onlyne_config::Drive) {
if self.dispatch.drive() == Some(drive) {
return;
}
match self.selector.build(drive) {
Ok(backend) => {
if self.dispatch.set_backend(backend) {
self.dispatch.set_drive(drive, None);
tracing::info!(
drive = %drive,
placement = %self.selector.placement,
backend = %self.dispatch.session_backend(),
"session backend selected"
);
self.republish_registration();
} else {
tracing::warn!(
drive = %drive,
placement = %self.selector.placement,
live = self.dispatch.session_count(),
held = %self.dispatch.session_backend(),
"the role's drive changed while it holds live sessions: they keep the \
backend they were opened under, and the new drive lands once the role \
is quiet"
);
}
}
Err(error) => {
tracing::error!(
drive = %drive,
placement = %self.selector.placement,
error = %error,
"the drive this role's spec declares cannot run under this machine's \
placement; every delivery for this role will be refused with that sentence"
);
self.dispatch.set_drive(drive, Some(error.to_string()));
}
}
}
fn republish_registration(&self) {
if let Err(error) = crate::session::adapter_socket::republish_registration(
&self.dispatch.workspace(),
&self.dispatch.role(),
self.dispatch.session_backend(),
self.dispatch.placement_name(),
) {
tracing::warn!(error = %error, "the client registration was not republished");
}
}
pub(super) fn cursor(&self) -> u64 {
self.store
.config(EVENT_CURSOR_KEY)
.ok()
.flatten()
.and_then(|value| value.parse().ok())
.unwrap_or(0)
}
pub(super) fn set_cursor(&self, seq: u64) {
if let Err(error) = self.store.put_config(EVENT_CURSOR_KEY, &seq.to_string()) {
tracing::warn!(error = %error, "event cursor was not stored");
}
}
}