impl TurnHost {
async fn handle_web_command(&mut self, command: WireCommand, turn: &mut TurnState) {
match command {
WireCommand::Submit {
session_id,
text,
images,
interrupt,
} => {
if session_id.is_empty() || session_id == self.session.id {
self.submit_web_text(text, images, interrupt, turn).await;
} else {
self.submit_web_text_for_session(&session_id, text, images, interrupt)
.await;
}
}
WireCommand::TriggerRuleNow { id } => self.trigger_web_rule_now(id, turn).await,
WireCommand::Abort { session_id } => {
if session_id.is_empty() || session_id == self.session.id {
self.request_abort(turn);
} else if !self.sessions.contains(&session_id) {
self.error_line(format!(
"abort ignored: session {session_id} is not an active or registered session {}",
self.session.id
));
} else {
self.cancel_session(&session_id);
}
}
WireCommand::ResolveControlPlane {
session_id,
approve,
} => {
let session_id = if session_id.is_empty() {
self.session.id.clone()
} else {
session_id
};
let decision = if approve {
theway_core::ControlPlanePromptDecision::Allow
} else {
theway_core::ControlPlanePromptDecision::Deny {
reason: Some("denied by user".into()),
}
};
self.resolve_control_plane_prompt_for_session(&session_id, decision);
}
WireCommand::SetModel {
session_id,
spec,
response,
} => {
let active_session = session_id.is_empty() || session_id == self.session.id;
let ok = if active_session {
self.set_model_from_spec(&spec).await
} else {
self.set_model_for_session(&session_id, &spec).await
};
if ok && active_session {
self.start_next_queued_turn(turn);
}
let _ = response.send(ok);
}
WireCommand::SetThinking {
session_id,
level,
response,
} => {
let ok = if session_id.is_empty() || session_id == self.session.id {
self.set_thinking_level(&level).await
} else {
self.set_thinking_for_session(&session_id, &level).await
};
let _ = response.send(ok);
}
WireCommand::SetSkillDirs { dirs } => self.handle_set_skill_dirs(dirs, turn).await,
WireCommand::Configure { config } => self.handle_configure(config, turn).await,
WireCommand::InvokeExtensionCommand {
name,
arguments,
has_interactive_client,
response,
} => {
let result = self
.handle_extension_command(name, arguments, has_interactive_client)
.await;
let _ = response.send(result);
}
WireCommand::ReloadExtensions {
cancel_active,
response,
} => {
let result = self
.handle_extension_reload(cancel_active, turn)
.await;
let _ = response.send(result);
}
WireCommand::DecideExtensionTrust { request, response } => {
let result = self.handle_extension_trust(request).await;
let _ = response.send(result);
}
WireCommand::ActivateSession { request, response } => {
let result = self.handle_activate_session(request, turn).await;
let _ = response.send(result);
}
WireCommand::SetCredential { request, response } => {
let result = self.handle_set_credential(request);
let _ = response.send(result);
}
WireCommand::ClearCredential { request, response } => {
let result = self.handle_clear_credential(request);
let _ = response.send(result);
}
WireCommand::SessionDeleted { id } => {
self.handle_session_deleted(&id, turn).await;
}
}
}
async fn handle_session_deleted(&mut self, id: &str, turn: &mut TurnState) {
if id != self.session.id {
self.sessions.remove(id);
return;
}
let remaining = match self.session.repository.list().await {
Ok(records) => records,
Err(error) => {
self.error_line(format!(
"active session deleted: fallback lookup failed: {error}"
));
return;
}
};
let Some(fallback_id) = remaining.last().map(|record| record.id.clone()) else {
self.system_line(format!("deleted active session {}; no sessions remain", id));
return;
};
let runtime = match (self.session.factory)(fallback_id.clone()).await {
Ok(runtime) => runtime,
Err(error) => {
self.error_line(format!(
"active session deleted: cannot open fallback session {fallback_id}: {error:#}"
));
return;
}
};
if turn.fut.is_some() {
self.request_abort(turn);
if let Some(future) = turn.fut.take() {
let _ = future.await;
}
}
let cwd = runtime.cwd.clone();
let new_state = SessionRuntimeState::from_runtime(
runtime,
self.session.factory.clone(),
self.session.repository.clone(),
self.session.retry.clone(),
self.session.log_path.clone(),
FeedProjectionState::new(
self.projection.capabilities.clone(),
self.projection.thinking_summary.clone(),
),
);
let old = std::mem::replace(&mut self.session, new_state);
self.sessions.insert(old);
if self.runtime.cwd != cwd {
self.runtime.cwd = cwd;
}
self.session
.kernel
.harness()
.session_switched(&self.session.id)
.await;
self.automation
.reload
.set_trigger_executor(self.session.kernel.trigger_executor().clone());
self.clear_feed();
crate::feed_replay::replay_transcript(
&mut self.projection.feed,
&self.session.kernel.harness().agent().state().messages,
self.runtime.feed_history_limit,
);
self.system_line(format!(
"deleted session {}; switched to {}",
id, self.session.id
));
self.session.busy = false;
self.session.queue.clear();
self.session.cumulative_usage = WireContextUsage::default();
self.projection.control_plane_prompt = None;
turn.aborted = false;
turn.prefix = "";
self.refresh_goal_state().await;
self.publish_current_snapshot().await;
}
fn handle_set_credential(
&mut self,
request: WireSetCredentialRequest,
) -> Result<(), WireRpcError> {
if request.session_id.trim().is_empty() {
return Err(WireRpcError {
code: "invalid_argument".into(),
message: "session_id must not be empty".into(),
});
}
if request.provider.trim().is_empty() {
return Err(WireRpcError {
code: "invalid_argument".into(),
message: "provider must not be empty".into(),
});
}
self.automation
.services
.session_execution
.set_credential(&request.session_id, &request.provider, request.secret)
.map_err(|error| match error {
crate::session_execution::RegistryError::SessionNotRegistered(id) => WireRpcError {
code: "not_found".into(),
message: format!("session {id} is not registered; activate it first"),
},
other => WireRpcError {
code: "failed_precondition".into(),
message: other.to_string(),
},
})
}
fn handle_clear_credential(
&mut self,
request: WireClearCredentialRequest,
) -> Result<(), WireRpcError> {
let session_id = request.session_id.trim();
if session_id.is_empty() {
return Err(WireRpcError {
code: "invalid_argument".into(),
message: "session_id must not be empty".into(),
});
}
if self
.automation
.services
.session_execution
.get(session_id)
.is_none()
{
return Err(WireRpcError {
code: "not_found".into(),
message: format!("session {session_id} is not registered; activate it first"),
});
}
match request.provider.as_deref() {
Some(provider) if provider.trim().is_empty() => Err(WireRpcError {
code: "invalid_argument".into(),
message: "provider must not be empty".into(),
}),
Some(provider) => {
self.automation
.services
.session_execution
.clear_credential(session_id, provider);
Ok(())
}
None => {
self.automation
.services
.session_execution
.clear_credentials(session_id);
Ok(())
}
}
}
async fn handle_activate_session(
&mut self,
request: WireActivateSessionRequest,
turn: &mut TurnState,
) -> Result<WireActivateSessionResponse, WireRpcError> {
let current_harness = self.session.kernel.harness().clone();
let activator = self
.automation
.services
.session_activator
.get()
.ok_or_else(|| WireRpcError {
code: "failed_precondition".into(),
message: "session activator is not installed".into(),
})?;
let activation = activator
.activate(&request, ¤t_harness)
.await?;
let response = WireActivateSessionResponse {
session: Some(activation.summary.clone()),
created: activation.created,
};
self.apply_activation(activation, turn).await;
Ok(response)
}
async fn handle_configure(&mut self, mut config: WireDaemonConfig, turn: &mut TurnState) {
tracing::info!(
target: "mcp",
"configure received: mcp_servers={} clear={:?}",
config.mcp_servers.len(),
config.clear_fields
);
let unknown = config.unknown_clear_fields();
if !unknown.is_empty() {
self.error_line(format!(
"configure: unknown clear field(s): {}",
unknown.join(", ")
));
return;
}
let mut applied = WireDaemonConfig::default();
if !config.models.is_empty() {
crate::model_defaults::register_models(&config.models);
applied.models = config.models.clone();
self.runtime.model_catalog = model_catalog();
}
if let Some(raw_key) = config.api_key.as_deref() {
let provider = config
.provider
.as_deref()
.map(str::trim)
.filter(|provider| !provider.is_empty())
.map(str::to_string)
.or_else(|| {
self.session
.kernel
.harness()
.agent()
.state()
.model
.as_ref()
.map(|model| model.provider.0.clone())
});
match provider {
Some(provider) => {
let key = raw_key.trim();
let mut keys = self
.automation
.services
.configured_api_keys
.write()
.expect("configured api keys poisoned");
if key.is_empty() {
keys.remove(&provider);
drop(keys);
applied.clear_fields.push("api_key".into());
} else {
keys.insert(provider, key.to_string());
drop(keys);
applied.api_key = Some(raw_key.to_string());
}
}
None => self.error_line("configure: api_key requires a provider"),
}
}
if config.auto_fetch_models == Some(true) {
let provider = config
.provider
.as_deref()
.map(str::trim)
.filter(|provider| !provider.is_empty())
.map(str::to_string)
.or_else(|| {
self.session
.kernel
.harness()
.agent()
.state()
.model
.as_ref()
.map(|model| model.provider.0.clone())
});
let base_url = config.base_url.clone().or_else(|| {
self.runtime
.config
.read()
.expect("daemon config poisoned")
.base_url
.clone()
});
match (provider, base_url) {
(Some(provider), Some(base_url)) => {
let configured_key = self
.automation
.services
.configured_api_keys
.read()
.expect("configured api keys poisoned")
.get(&provider)
.cloned();
let api_key = theway_transport::auth::AuthStore::load()
.unwrap_or_default()
.resolve_for_provider_with(&provider, configured_key.as_deref());
match crate::model_fetch::fetch_models(&base_url, api_key.as_deref(), &provider)
.await
{
Ok(models) if !models.is_empty() => {
let first = models[0].id.clone();
crate::model_defaults::register_models(&models);
self.runtime.model_catalog = model_catalog();
if config.model.is_none() {
config.provider.get_or_insert(provider);
config.model = Some(first);
}
applied.auto_fetch_models = Some(true);
applied.models = models;
}
Ok(_) => self
.error_line("configure: auto-fetch models returned an empty catalog"),
Err(err) => {
self.error_line(format!("configure: auto-fetch models: {err}"));
}
}
}
(None, _) => {
self.error_line("configure: auto-fetch models requires a provider")
}
(_, None) => {
self.error_line("configure: auto-fetch models requires a base_url")
}
}
}
if (config.clears("provider") && config.provider.is_none())
|| (config.clears("model") && config.model.is_none())
{
self.error_line("configure: the active provider/model cannot be cleared");
} else if config.provider.is_some() != config.model.is_some() {
self.error_line("configure: provider and model must be supplied together");
} else if config.provider.is_some()
|| config.base_url.is_some()
|| config.clears("base_url")
{
let mut model = match (config.provider.as_deref(), config.model.as_deref()) {
(Some(provider), Some(id)) => theway_llm_provider::get_model(
&theway_llm_provider::Provider::from(provider),
id,
),
_ => self.session.kernel.harness().agent().state().model.clone(),
};
if config.clears("base_url")
&& let Some(current) = model.as_ref()
{
model = theway_llm_provider::get_model(¤t.provider, ¤t.id)
.or_else(|| Some(current.clone()));
}
if let Some(model) = model.as_mut()
&& let Some(base_url) = config.base_url.as_ref()
{
model.base_url = base_url.clone();
}
match model {
Some(model) if self.apply_model(model.clone()).await => {
applied.provider = Some(model.provider.0.clone());
applied.model = Some(model.id.clone());
if model.base_url.is_empty() {
applied.clear_fields.push("base_url".into());
} else {
applied.base_url = Some(model.base_url);
}
}
Some(_) => {}
None => self.error_line("configure: no active or matching model to update"),
}
}
if config.thinking.is_some() || config.clears("thinking") {
let enabled = config.thinking.unwrap_or(false);
let level = if enabled {
theway_core::ThinkingLevel::High
} else {
theway_core::ThinkingLevel::Off
};
match self.session.kernel.harness().set_thinking_level(level).await {
Ok(_) if config.thinking.is_none() => applied.clear_fields.push("thinking".into()),
Ok(_) => applied.thinking = Some(enabled),
Err(err) => self.error_line(format!("configure thinking: {err}")),
}
}
if config.thinking_level.is_some() || config.clears("thinking_level") {
let requested = config.thinking_level.as_deref();
let level = match requested {
Some(raw) => match raw.parse::<theway_core::ThinkingLevel>() {
Ok(level) => Some(level),
Err(err) => {
self.error_line(format!(
"configure thinking_level: invalid level {raw:?}: {err}"
));
None
}
},
None => Some(theway_core::ThinkingLevel::Off),
};
if let Some(level) = level {
match self.session.kernel.harness().set_thinking_level(level).await {
Ok(_) if requested.is_none() => {
applied.clear_fields.push("thinking_level".into())
}
Ok(_) => applied.thinking_level = config.thinking_level.clone(),
Err(err) => self.error_line(format!("configure thinking_level: {err}")),
}
}
}
if !config.skills.is_empty() || config.clears("skills") {
let provisioned: Vec<theway_core::Skill> = config
.skills
.iter()
.map(|skill| theway_core::Skill {
name: skill.name.clone(),
description: skill.description.clone(),
file_path: skill.file_path.clone(),
content: skill.content.clone(),
disable_model_invocation: skill.disable_model_invocation,
source: if skill.source == "project" {
theway_core::SkillSource::Project
} else {
theway_core::SkillSource::User
},
})
.collect();
let builtins: Vec<_> = self
.session
.kernel
.harness()
.skills()
.into_iter()
.filter(|skill| matches!(skill.source, theway_core::SkillSource::Builtin))
.collect();
let mut merged =
crate::builtin_skills::merge_with_user_project(builtins, &provisioned);
let overrides = crate::skill_overrides::load(&self.runtime.paths.base).await;
crate::skill_overrides::apply(&overrides, &mut merged);
self.session.kernel.harness().replace_skills(merged);
*self.runtime.provisioned_skills.write().unwrap() = provisioned.clone();
if provisioned.is_empty() {
applied.clear_fields.push("skills".into());
} else {
applied.skills = config.skills.clone();
}
}
if !config.templates.is_empty() || config.clears("templates") {
let provisioned: Vec<theway_core::PromptTemplate> = config
.templates
.iter()
.map(|template| theway_core::PromptTemplate {
name: template.name.clone(),
description: if template.description.trim().is_empty() {
None
} else {
Some(template.description.clone())
},
content: template.content.clone(),
file_path: template.file_path.clone(),
})
.collect();
self.session.kernel.harness().replace_templates(provisioned.clone());
*self.runtime.provisioned_templates.write().unwrap() = provisioned.clone();
if provisioned.is_empty() {
applied.clear_fields.push("templates".into());
} else {
applied.templates = config.templates.clone();
}
}
if !config.mcp_servers.is_empty() || config.clears("mcp_servers") {
let requested: Vec<crate::mcp_loader::ServerConfig> = config
.mcp_servers
.iter()
.map(crate::mcp_loader::server_config_from_wire)
.collect();
match crate::mcp_loader::validate_unique_names(&requested) {
Ok(()) => {
let old_tools = self
.runtime
.mcp_provision
.read()
.unwrap()
.tools
.clone();
let result = crate::mcp_loader::connect_servers(
&requested,
&self.runtime.cwd,
&self.runtime.paths.base.join("auth.json"),
)
.await;
let (new_tools, new_hooks, capabilities_update) = {
let mut slot = self.runtime.mcp_provision.write().unwrap();
slot.replace_connection_result(requested, result);
let capabilities_update = (
slot.server_names.len(),
slot.tool_names.len(),
slot.server_names.clone(),
slot.tool_names.clone(),
slot.errors.clone(),
);
(slot.tools.clone(), slot.hooks.clone(), capabilities_update)
};
if self.session.mcp_overlay.is_some() {
self.remerge_active_session_mcp();
} else {
self.session
.kernel
.harness()
.replace_mcp_tools(&old_tools, new_tools);
self.projection.capabilities.mcp_servers = capabilities_update.0;
self.projection.capabilities.mcp_tools = capabilities_update.1;
self.projection.capabilities.mcp_server_names = capabilities_update.2;
self.projection.capabilities.mcp_tool_names = capabilities_update.3;
self.projection.capabilities.mcp_server_errors = capabilities_update.4;
{
use crate::orchestration::session::NotificationHookSink;
use crate::trigger_engine::notification_hook::NotificationHook;
let executor = self.runtime.trigger_executor.clone();
let mut slot = self.runtime.mcp_provision.write().unwrap();
for hook in &new_hooks {
let label = hook.label().to_string();
if slot.registered_labels.insert(label) {
executor.register(hook.clone());
}
}
}
}
if config.mcp_servers.is_empty() {
applied.clear_fields.push("mcp_servers".into());
} else {
applied.mcp_servers = config.mcp_servers.clone();
}
let (connected, failed_count, tool_count, failed_names) = {
let slot = self.runtime.mcp_provision.read().unwrap();
(
slot.server_names.len(),
slot.errors.len(),
slot.tool_names.len(),
slot.errors
.iter()
.map(|(name, _)| name.as_str())
.collect::<Vec<_>>()
.join(", "),
)
};
if failed_count == 0 {
self.system_line(format!(
"MCP: connected {connected} server(s), {tool_count} tool(s)"
));
} else {
self.system_line(format!(
"MCP: connected {connected} server(s), {failed_count} failed: {failed_names}"
));
}
}
Err(message) => self.error_line(format!("configure mcp_servers: {message}")),
}
}
if !config.builtin_skills.is_empty() || config.clears("builtin_skills") {
let requested = if config.clears("builtin_skills") && config.builtin_skills.is_empty() {
Vec::new()
} else {
config.builtin_skills.clone()
};
let resolved = crate::builtin_skills::resolve_builtins(&[], &requested)
.expect("an empty CLI list cannot produce a hard builtin error");
for diagnostic in resolved.diagnostics {
self.error_line(diagnostic);
}
let enabled: Vec<String> = resolved
.skills
.iter()
.map(|skill| skill.name.clone())
.collect();
let non_builtin: Vec<_> = self
.session
.kernel
.harness()
.skills()
.into_iter()
.filter(|skill| !matches!(skill.source, theway_core::SkillSource::Builtin))
.collect();
self.session.kernel
.harness()
.replace_skills(crate::builtin_skills::merge_with_user_project(
resolved.skills,
&non_builtin,
));
if enabled.is_empty() {
applied.clear_fields.push("builtin_skills".into());
} else {
applied.builtin_skills = enabled;
}
}
if !config.skills_dirs.is_empty() || config.clears("skills_dirs") {
let dirs = if config.skills_dirs.is_empty() {
Vec::new()
} else {
config.skills_dirs.clone()
};
self.handle_set_skill_dirs(dirs, turn).await;
let actual = self.runtime.path_context.read().unwrap().skills_dirs.clone();
if actual.is_empty() {
applied.clear_fields.push("skills_dirs".into());
} else {
applied.skills_dirs = actual;
}
}
if let Some(secs) = config.trigger_poll_secs {
if secs == 0 {
self.error_line("configure: trigger_poll_secs must be greater than zero");
} else {
self.automation.services.dynamic_triggers.set_poll_interval_secs(secs);
applied.trigger_poll_secs = Some(secs);
}
} else if config.clears("trigger_poll_secs") {
self.automation.services.dynamic_triggers.set_poll_interval_secs(
theway_transport::triggers::DEFAULT_DYNAMIC_TRIGGER_POLL_INTERVAL_SECS,
);
applied.clear_fields.push("trigger_poll_secs".into());
}
if let Some(lines) = config.tui_max_feed_lines {
if lines == 0 {
self.error_line("configure: tui_max_feed_lines must be greater than zero");
} else {
self.runtime.feed_history_limit = Some(lines);
applied.tui_max_feed_lines = Some(lines);
}
} else if config.clears("tui_max_feed_lines") {
self.runtime.feed_history_limit = None;
applied.clear_fields.push("tui_max_feed_lines".into());
}
if let Some(addr) = config.tool_service_addr.as_ref() {
if addr.trim().is_empty() {
self.error_line("configure: tool_service_addr must not be empty; clear it instead");
} else {
applied.tool_service_addr = Some(addr.clone());
}
} else if config.clears("tool_service_addr") {
applied.clear_fields.push("tool_service_addr".into());
}
if config.storage_service_addr.is_some() || config.clears("storage_service_addr") {
self.error_line(
"configure: storage_service_addr is startup-only and cannot be changed at runtime",
);
}
if config.executor_kind.is_some() || config.clears("executor_kind") {
self.error_line(
"configure: executor_kind is startup-only and cannot be changed at runtime; set `[executor] kind` in config.toml and restart the daemon",
);
}
if config.tgrep.is_some() || config.clears("tgrep") {
self.error_line(
"configure: tgrep is startup-only and cannot be changed at runtime; set `[tools] tgrep` in config.toml and restart the daemon",
);
}
let touched = self.runtime.config.write().unwrap().merge_from(&applied);
if touched == 0 {
self.system_line("configure: no applicable settings changed");
} else {
self.system_line(format!("configure: applied {touched} setting(s)"));
}
self.start_next_queued_turn(turn);
}
async fn handle_set_skill_dirs(&mut self, dirs: Vec<String>, turn: &mut TurnState) {
let dirs: Vec<PathBuf> = dirs.into_iter().map(PathBuf::from).collect();
self.runtime.paths.set_extra_skill_dirs(dirs);
self.runtime.path_context.write().unwrap().skills_dirs = self
.runtime
.paths
.current_extra_skill_dirs()
.into_iter()
.map(|dir| dir.to_string_lossy().into_owned())
.collect();
if turn.fut.is_some() {
self.request_abort(turn);
}
match self.session.kernel.harness().reload_skills_from_disk().await {
Ok(out) => self.system_line(format!(
"set skill dirs: {} loaded, {} diagnostics",
out.skills.len(),
out.diagnostics.len()
)),
Err(e) => self.error_line(format!("set skill dirs: {e:#}")),
}
}
async fn trigger_web_rule_now(&mut self, id: String, turn: &mut TurnState) {
let id = id.trim();
if id.is_empty() {
self.error_line("trigger: missing rule id");
return;
}
let Some(rule) = self.automation.services.dynamic_triggers
.list()
.into_iter()
.find(|rule| rule.id == id)
else {
self.error_line(format!("trigger: no dynamic trigger rule with id `{id}`"));
return;
};
let display = format!(
"trigger now {}: {}",
feed::truncate_chars(&rule.id, 18),
wire_preview(&rule.action)
);
if turn.fut.is_some() {
self.queue_user_prompt(display, rule.action, Vec::new())
.await;
} else {
self.projection.feed.push_user(display);
self.start_user_prompt_turn(rule.action, Vec::new(), turn);
}
}
}