impl TurnHost {
pub(crate) fn new(config: DaemonConfig) -> Self {
let registry = Arc::new(config.registry);
registry.set_file_commands(crate::file_commands::scan_file_commands(
&config.cwd,
&config.paths.home,
));
let completer = SlashCompleter::from_commands(slash_commands(®istry));
let reload_runtime = config.services.reload.install(ReloadRuntime::new(
registry.clone(),
config.cwd.clone(),
config.trigger_executor.clone(),
Arc::new(AtomicU64::new(0)),
));
let path_context = Arc::new(std::sync::RwLock::new(WirePathContext {
home: config.paths.home.to_string_lossy().into_owned(),
base: config.paths.base.to_string_lossy().into_owned(),
work_dir: config.paths.work_dir.to_string_lossy().into_owned(),
skills_dirs: config
.paths
.current_extra_skill_dirs()
.into_iter()
.map(|dir| dir.to_string_lossy().into_owned())
.collect(),
}));
let startup_state = config.harness.agent().state();
let startup_model = startup_state.model.clone();
let startup_thinking = startup_state
.thinking_level
.map(|level| level != theway_core::ThinkingLevel::Off);
drop(startup_state);
let daemon_config = Arc::new(std::sync::RwLock::new(WireDaemonConfig {
provider: startup_model.as_ref().map(|model| model.provider.0.clone()),
model: startup_model.as_ref().map(|model| model.id.clone()),
base_url: startup_model
.as_ref()
.map(|model| model.base_url.clone())
.filter(|url| !url.is_empty()),
thinking: startup_thinking,
builtin_skills: config.startup.builtin_skills.clone(),
skills_dirs: config
.paths
.current_extra_skill_dirs()
.into_iter()
.map(|dir| dir.to_string_lossy().into_owned())
.collect(),
trigger_poll_secs: Some(config.startup.trigger_poll_secs),
tui_max_feed_lines: config.startup.tui_max_feed_lines,
tool_service_addr: None,
storage_service_addr: config.startup.storage_service_addr.clone(),
clear_fields: Vec::new(),
}));
let tool_ops: Arc<dyn ToolOps> = Arc::new(ForwardingToolOps::new(daemon_config.clone()));
let mut kernel = ReplKernel::new(config.harness, config.trigger_executor, config.retry);
kernel.set_extension_host(config.extension_host);
Self {
session: SessionRuntimeState {
kernel,
id: config.session_id,
cwd: config.cwd.clone(),
log_path: config.log_path,
tool_count: config.tool_count,
factory: config.session_factory,
repository: config.session_repo,
shared_state: config.current_session_state,
busy: false,
queue: VecDeque::new(),
},
automation: AutomationRuntime {
services: config.services,
reload: reload_runtime,
dag: config.dag_engine,
subagents: config.subagent_registry,
},
runtime: RuntimeConfiguration {
registry,
completer,
cwd: config.cwd,
paths: config.paths,
path_context,
config: daemon_config,
tool_ops,
model_catalog: model_catalog(),
feed_history_limit: config.startup.tui_max_feed_lines,
latest: None,
snapshot_tx: None,
},
projection: FeedProjectionState {
feed: Feed::new(),
plain_lines_cache: theway_transport::feed::PlainLinesCache::new(100),
block_versions: Vec::new(),
dirty_blocks: BTreeSet::new(),
latest_trigger_poll: None,
latest_goal: None,
thinking_summary: config.thinking_summary,
thinking_burst: super::thinking_summary::ThinkingBurst::default(),
control_plane_prompt: None,
capabilities: config.capabilities,
},
inputs: RuntimeEventInputs {
feed_rx: Some(config.feed_rx),
feed_tx: config.feed_tx,
main_run_rx: Some(config.main_run_rx),
control_plane_prompt_rx: config.control_plane_prompt_rx,
},
}
}
fn system_line(&mut self, text: impl AsRef<str>) {
self.projection
.feed
.push_plain_untimed(text.as_ref(), Level::System);
}
fn error_line(&mut self, text: impl AsRef<str>) {
self.projection
.feed
.push_plain_untimed(text.as_ref(), Level::Error);
}
pub(crate) fn transport_endpoints(&mut self) -> TransportEndpoints {
let (command_tx, command_rx) = mpsc::unbounded_channel::<WireCommand>();
let (snapshot_tx, _) = broadcast::channel::<WireStatusUpdate>(128);
let latest = Arc::new(Mutex::new(self.wire_snapshot()));
let (event_tx, _) = broadcast::channel::<WireAgentEvent>(256);
let (dag_event_tx, _) = broadcast::channel::<WireDagEvent>(256);
let (core_dag_event_tx, _) = broadcast::channel::<DagEvent>(256);
self.automation.dag
.set_event_sender(Some(core_dag_event_tx.clone()));
let agent_fwd = {
let mut agent_rx = self.automation.subagents.subscribe();
let agent_tx = event_tx.clone();
let mut dag_rx = core_dag_event_tx.subscribe();
let dag_tx = dag_event_tx.clone();
tokio::spawn(async move {
let agent_loop = async move {
loop {
match agent_rx.recv().await {
Ok(event) => {
let _ = agent_tx.send(agent_event(event));
}
Err(broadcast::error::RecvError::Lagged(n)) => {
tracing::warn!("SubagentJobEvent broadcast lagged by {n}, skipping");
continue;
}
Err(broadcast::error::RecvError::Closed) => break,
}
}
tracing::debug!(
"SubagentJobEvent registry channel closed; forwarder task exiting"
);
};
let dag_loop = async move {
loop {
match dag_rx.recv().await {
Ok(event) => {
let _ = dag_tx.send(dag_event(event));
}
Err(broadcast::error::RecvError::Lagged(n)) => {
tracing::warn!("DagEvent broadcast lagged by {n}, skipping");
continue;
}
Err(broadcast::error::RecvError::Closed) => break,
}
}
tracing::debug!("DagEvent channel closed; forwarder task exiting");
};
tokio::join!(agent_loop, dag_loop);
})
.abort_handle()
};
TransportEndpoints {
command_tx,
command_rx,
snapshot_tx,
latest,
events: event_tx,
dag_events: dag_event_tx,
completer: self.runtime.completer.clone(),
job_ops: Arc::new(CoreJobOps::new(
self.automation.subagents.clone(),
self.automation.dag.clone(),
)),
graph_ops: Arc::new(CoreGraphOps::new(self.automation.dag.clone())),
session_ops: Arc::new(crate::session_ops::AppSessionOps::new(
self.session.repository.clone(),
self.automation.dag.clone(),
self.session.shared_state.clone(),
self.automation.services.session_execution.clone(),
)),
path_context: self.runtime.path_context.clone(),
daemon_config: self.runtime.config.clone(),
tool_ops: self.runtime.tool_ops.clone(),
storage_ops: std::sync::Arc::new(theway_transport::UnavailableStorageOps),
session_id: self.session.id.clone(),
agent_fwd,
}
}
pub(crate) async fn run_transport_loop(
mut self,
mode: TransportMode,
endpoints: TransportEndpoints,
mut server_task: tokio::task::JoinHandle<Result<()>>,
) -> Result<()> {
let label = mode.label();
let mut command_rx = endpoints.command_rx;
self.runtime.latest = Some(endpoints.latest.clone());
self.runtime.snapshot_tx = Some(endpoints.snapshot_tx.clone());
let latest = endpoints.latest;
let snapshot_tx = endpoints.snapshot_tx;
let mut feed_rx = self.inputs.feed_rx.take().expect("feed_rx taken once");
let mut main_run_rx = self.inputs.main_run_rx.take().expect("main_run_rx taken once");
let mut control_plane_prompt_rx = self.inputs.control_plane_prompt_rx.take();
let mut turn = TurnState::default();
self.refresh_goal_state().await;
self.publish_snapshot(&latest, &snapshot_tx, true).await;
let mut dirty = false;
let mut metadata_dirty = false;
let mut publish_tick = tokio::time::interval(Duration::from_millis(50));
publish_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
publish_tick.reset();
loop {
tokio::select! {
biased;
result = poll_turn(&mut turn.fut), if turn.fut.is_some() => {
self.finish_turn(&mut turn, result).await;
dirty = true;
metadata_dirty = true;
}
Some(command) = command_rx.recv() => {
self.handle_web_command(command, &mut turn).await;
dirty = true;
metadata_dirty = true;
}
Some(update) = feed_rx.recv() => {
metadata_dirty |= self.apply_feed_update(update);
while let Ok(update) = feed_rx.try_recv() {
metadata_dirty |= self.apply_feed_update(update);
}
dirty = true;
}
Some(trace_id) = main_run_rx.recv(), if turn.fut.is_none() => {
self.start_triggered_turn(trace_id, &mut turn);
dirty = true;
metadata_dirty = true;
}
Some(prompt) = async {
match control_plane_prompt_rx.as_mut() {
Some(rx) => rx.recv().await,
None => None,
}
}, if self.projection.control_plane_prompt.is_none() && control_plane_prompt_rx.is_some() => {
self.show_control_plane_prompt(prompt);
dirty = true;
metadata_dirty = true;
}
_ = publish_tick.tick(), if dirty => {
dirty = false;
self.publish_snapshot(&latest, &snapshot_tx, metadata_dirty).await;
metadata_dirty = false;
}
_ = tokio::signal::ctrl_c() => {
if turn.fut.is_some() {
self.request_abort(&mut turn);
self.publish_snapshot(&latest, &snapshot_tx, true).await;
}
break;
}
_ = sigterm_received() => {
self.system_line(format!("[{label}] received SIGTERM, shutting down"));
if turn.fut.is_some() {
self.request_abort(&mut turn);
self.publish_snapshot(&latest, &snapshot_tx, true).await;
}
break;
}
server_result = &mut server_task => {
match server_result {
Ok(Ok(())) => {}
Ok(Err(e)) => self.error_line(format!("{label} server: {e}")),
Err(e) => self.error_line(format!("{label} server task: {e}")),
}
break;
}
}
}
self.automation.services.session_execution.clear_all_credentials();
Ok(())
}
}
async fn sigterm_received() {
#[cfg(unix)]
if let Ok(mut terminate) =
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
{
terminate.recv().await;
return;
}
std::future::pending::<()>().await;
}