use std::{
hash::{Hash, Hasher},
sync::Arc,
};
use tokio::sync::Mutex;
use tokio_util::sync::CancellationToken;
#[cfg(test)]
use std::sync::{
OnceLock,
atomic::{AtomicBool, Ordering},
};
#[cfg(test)]
use tokio::sync::Notify;
use crate::{
domain::{
chat_completion::ChatCommandCompletion,
codec::decode_base64,
errors::{AgentError, AgentResult, ErrorCode},
pi_rpc::{ExtensionUiResponsePayload, ThinkingLevel},
policy::CommandPolicy,
protocol::{AgentMessage, BackendMessage, ChatExtensionResponse, SkillScope},
resize::TerminalSize,
skills::MAX_SEARCH_LIMIT,
},
infrastructure::{
directory_browser::DirectoryBrowser, pi_history::PiHistoryStore, pi_models::PiModels,
project_file_search::ProjectFileSearch, pty_session::PtySession,
},
operational::{
pi_rpc_manager::{CreatePiChat, ExtensionResponse, PiRpcManager, ResumePiChat},
session_manager::{SessionManager, SessionOutput},
skill_service::SkillRequest,
},
pairing::SourceIdentity,
};
const DEFAULT_COLS: u16 = 80;
const DEFAULT_ROWS: u16 = 24;
#[cfg(test)]
struct FileSearchTestPause {
enabled: AtomicBool,
paused: AtomicBool,
paused_notify: Notify,
resume_notify: Notify,
}
#[cfg(test)]
fn file_search_test_pause() -> &'static FileSearchTestPause {
static PAUSE: OnceLock<FileSearchTestPause> = OnceLock::new();
PAUSE.get_or_init(|| FileSearchTestPause {
enabled: AtomicBool::new(false),
paused: AtomicBool::new(false),
paused_notify: Notify::new(),
resume_notify: Notify::new(),
})
}
#[cfg(test)]
pub(crate) fn set_test_file_search_pause(enabled: bool) {
let pause = file_search_test_pause();
pause.enabled.store(enabled, Ordering::Release);
if !enabled {
pause.resume_notify.notify_waiters();
}
}
#[cfg(test)]
pub(crate) async fn wait_for_test_file_search_paused() {
let pause = file_search_test_pause();
while !pause.paused.load(Ordering::Acquire) {
pause.paused_notify.notified().await;
}
}
#[cfg(test)]
async fn pause_file_search_for_test() {
let pause = file_search_test_pause();
if !pause.enabled.load(Ordering::Acquire) {
return;
}
pause.paused.store(true, Ordering::Release);
pause.paused_notify.notify_waiters();
while pause.enabled.load(Ordering::Acquire) {
pause.resume_notify.notified().await;
}
pause.paused.store(false, Ordering::Release);
}
#[derive(Clone, Debug)]
pub(crate) enum AuthenticatedRequestOwner {
Ui {
source: SourceIdentity,
client_id: u64,
liveness: CancellationToken,
},
Legacy {
connection_id: String,
liveness: CancellationToken,
},
}
impl AuthenticatedRequestOwner {
pub(crate) fn ui(source: SourceIdentity, client_id: u64) -> Self {
Self::ui_with_liveness(source, client_id, CancellationToken::new())
}
pub(crate) fn ui_with_liveness(
source: SourceIdentity,
client_id: u64,
liveness: CancellationToken,
) -> Self {
Self::Ui {
source,
client_id,
liveness,
}
}
pub(crate) fn legacy(connection_id: String) -> Self {
Self::legacy_with_liveness(connection_id, CancellationToken::new())
}
pub(crate) fn legacy_with_liveness(connection_id: String, liveness: CancellationToken) -> Self {
Self::Legacy {
connection_id,
liveness,
}
}
pub(crate) fn invalidate(&self) {
match self {
Self::Ui { liveness, .. } | Self::Legacy { liveness, .. } => liveness.cancel(),
}
}
pub(crate) fn ensure_live(&self) -> AgentResult<()> {
if !self.is_live() {
return Err(AgentError::new(
ErrorCode::InvalidMessage,
"authenticated chat requester is no longer connected",
));
}
Ok(())
}
pub(crate) fn is_live(&self) -> bool {
match self {
Self::Ui { liveness, .. } | Self::Legacy { liveness, .. } => !liveness.is_cancelled(),
}
}
pub(crate) fn liveness_token(&self) -> CancellationToken {
match self {
Self::Ui { liveness, .. } | Self::Legacy { liveness, .. } => liveness.clone(),
}
}
pub(crate) fn can_replace_after_ui_reconnect(&self, successor: &Self) -> bool {
matches!(
(self, successor),
(
Self::Ui {
source: current_source,
client_id: current_client_id,
..
},
Self::Ui {
source: successor_source,
client_id: successor_client_id,
..
},
) if current_source == successor_source && successor_client_id > current_client_id
)
}
pub(crate) fn can_transfer_to_other_ui(&self, successor: &Self) -> bool {
matches!(
(self, successor),
(
Self::Ui {
source: current_source,
..
},
Self::Ui {
source: successor_source,
..
},
) if current_source != successor_source
)
}
}
impl PartialEq for AuthenticatedRequestOwner {
fn eq(&self, other: &Self) -> bool {
match (self, other) {
(
Self::Ui {
source: left_source,
client_id: left_client_id,
..
},
Self::Ui {
source: right_source,
client_id: right_client_id,
..
},
) => left_source == right_source && left_client_id == right_client_id,
(
Self::Legacy {
connection_id: left_connection_id,
..
},
Self::Legacy {
connection_id: right_connection_id,
..
},
) => left_connection_id == right_connection_id,
_ => false,
}
}
}
impl Eq for AuthenticatedRequestOwner {}
impl Hash for AuthenticatedRequestOwner {
fn hash<H: Hasher>(&self, state: &mut H) {
match self {
Self::Ui {
source, client_id, ..
} => {
0_u8.hash(state);
source.hash(state);
client_id.hash(state);
}
Self::Legacy { connection_id, .. } => {
1_u8.hash(state);
connection_id.hash(state);
}
}
}
}
pub(crate) struct RouteResponse {
pub message: AgentMessage,
}
pub(crate) enum RouteAction {
Direct {
response: Option<RouteResponse>,
output: Option<SessionOutput>,
},
Background(BackgroundRequest),
Chat(ChatRequest),
Skill(SkillRequest),
}
pub(crate) enum ChatRequest {
Create {
manager: Arc<PiRpcManager>,
history_operations: Arc<Mutex<()>>,
owner: AuthenticatedRequestOwner,
request: CreatePiChat,
},
Resume {
manager: Arc<PiRpcManager>,
history_operations: Arc<Mutex<()>>,
store: PiHistoryStore,
policy: CommandPolicy,
session_id: String,
request_id: String,
conversation_id: String,
owner: AuthenticatedRequestOwner,
},
Snapshot {
manager: Arc<PiRpcManager>,
session_id: String,
request_id: String,
owner: AuthenticatedRequestOwner,
},
SetModel {
manager: Arc<PiRpcManager>,
session_id: String,
request_id: String,
provider: String,
model_id: String,
owner: AuthenticatedRequestOwner,
},
SetThinking {
manager: Arc<PiRpcManager>,
session_id: String,
request_id: String,
level: ThinkingLevel,
owner: AuthenticatedRequestOwner,
},
TranscriptPage {
manager: Arc<PiRpcManager>,
session_id: String,
request_id: String,
cursor: String,
direction: crate::domain::chat::TranscriptDirection,
owner: AuthenticatedRequestOwner,
},
CommandsList {
manager: Arc<PiRpcManager>,
session_id: String,
request_id: String,
owner: AuthenticatedRequestOwner,
},
FilesSearch {
manager: Arc<PiRpcManager>,
session_id: String,
request_id: String,
query: String,
owner: AuthenticatedRequestOwner,
},
Prompt {
manager: Arc<PiRpcManager>,
session_id: String,
request_id: String,
text: String,
},
Abort {
manager: Arc<PiRpcManager>,
session_id: String,
request_id: String,
prompt_request_id: Option<String>,
},
Close {
manager: Arc<PiRpcManager>,
session_id: String,
request_id: String,
},
ExtensionRespond {
manager: Arc<PiRpcManager>,
session_id: String,
request_id: String,
extension_request_id: String,
response: ExtensionUiResponsePayload,
},
}
impl ChatRequest {
pub(crate) fn session_id(&self) -> &str {
match self {
Self::Create { request, .. } => &request.session_id,
Self::Resume { session_id, .. }
| Self::Snapshot { session_id, .. }
| Self::SetModel { session_id, .. }
| Self::SetThinking { session_id, .. }
| Self::TranscriptPage { session_id, .. }
| Self::CommandsList { session_id, .. }
| Self::FilesSearch { session_id, .. }
| Self::Prompt { session_id, .. }
| Self::Abort { session_id, .. }
| Self::Close { session_id, .. }
| Self::ExtensionRespond { session_id, .. } => session_id,
}
}
pub(crate) fn request_id(&self) -> &str {
match self {
Self::Create { request, .. } => &request.request_id,
Self::Resume { request_id, .. }
| Self::Snapshot { request_id, .. }
| Self::SetModel { request_id, .. }
| Self::SetThinking { request_id, .. }
| Self::TranscriptPage { request_id, .. }
| Self::CommandsList { request_id, .. }
| Self::FilesSearch { request_id, .. }
| Self::Prompt { request_id, .. }
| Self::Abort { request_id, .. }
| Self::Close { request_id, .. }
| Self::ExtensionRespond { request_id, .. } => request_id,
}
}
pub(crate) async fn execute(self) -> AgentResult<Option<AgentMessage>> {
match self {
Self::Create {
manager,
history_operations,
owner,
request,
} => {
let _history = history_operations.lock_owned().await;
owner.ensure_live()?;
manager.create_for_owner(request, owner).await?;
Ok(None)
}
Self::Resume {
manager,
history_operations,
store,
policy,
session_id,
request_id,
conversation_id,
owner,
} => {
let _history = history_operations.lock_owned().await;
owner.ensure_live()?;
let resolve_store = store.clone();
let resolved = tokio::task::spawn_blocking(move || {
resolve_store.resolve_session(&conversation_id)
})
.await
.map_err(|_| background_task_failed())??;
owner.ensure_live()?;
let cwd = policy.validate_chat_cwd(&resolved.cwd().to_string_lossy())?;
let snapshot = manager
.resume_for_owner_with_request_id(
ResumePiChat {
session_id: session_id.clone(),
cwd: cwd.into(),
session_path: resolved.canonical_path().to_path_buf(),
history_store: Some(store),
history_session: Some(resolved),
},
owner,
&request_id,
)
.await?;
Ok(Some(
manager
.finalize_snapshot_response(&session_id, &request_id, snapshot)
.await?,
))
}
Self::Snapshot {
manager,
session_id,
request_id,
owner,
} => {
let snapshot = manager
.snapshot_for_owner_with_request_id(&session_id, &owner, &request_id)
.await?;
Ok(Some(
manager
.finalize_snapshot_response(&session_id, &request_id, snapshot)
.await?,
))
}
Self::SetModel {
manager,
session_id,
request_id,
provider,
model_id,
owner,
} => {
manager
.set_model_for_owner(&session_id, &provider, &model_id, &owner)
.await?;
let snapshot = manager
.snapshot_for_owner_with_request_id(&session_id, &owner, &request_id)
.await?;
Ok(Some(
manager
.finalize_snapshot_response(&session_id, &request_id, snapshot)
.await?,
))
}
Self::SetThinking {
manager,
session_id,
request_id,
level,
owner,
} => {
manager
.set_thinking_level_for_owner(&session_id, level, &owner)
.await?;
let snapshot = manager
.snapshot_for_owner_with_request_id(&session_id, &owner, &request_id)
.await?;
Ok(Some(
manager
.finalize_snapshot_response(&session_id, &request_id, snapshot)
.await?,
))
}
Self::TranscriptPage {
manager,
session_id,
request_id,
cursor,
direction,
owner,
} => Ok(Some(
manager
.transcript_page_for_owner(&session_id, &request_id, &cursor, direction, &owner)
.await?,
)),
Self::CommandsList {
manager,
session_id,
request_id,
owner,
} => {
let commands = manager.commands_for_owner(&session_id, &owner).await?;
Ok(Some(AgentMessage::ChatCommandsListed {
session_id,
request_id,
commands: commands
.into_iter()
.map(ChatCommandCompletion::from)
.collect(),
}))
}
Self::FilesSearch {
manager,
session_id,
request_id,
query,
owner,
} => {
let workspace = manager.workspace_for_owner(&session_id, &owner).await?;
#[cfg(test)]
pause_file_search_for_test().await;
let search_query = query.clone();
let workspace_path = workspace.path().to_path_buf();
let result = tokio::task::spawn_blocking(move || {
ProjectFileSearch::default().search(&workspace_path, &search_query)
})
.await
.map_err(|_| {
AgentError::new(
ErrorCode::ChatFileSearchFailed,
"project file search task failed",
)
})??;
manager
.ensure_workspace_for_owner(&workspace, &owner)
.await?;
Ok(Some(AgentMessage::ChatFilesSearched {
session_id,
request_id,
query,
matches: result.matches,
truncated: result.truncated,
}))
}
Self::Prompt {
manager,
session_id,
request_id,
text,
} => {
manager.prompt(&session_id, &request_id, &text).await?;
Ok(None)
}
Self::Abort {
manager,
session_id,
request_id,
prompt_request_id,
} => {
manager
.abort(&session_id, &request_id, prompt_request_id.as_deref())
.await?;
Ok(Some(AgentMessage::ChatAbortAccepted {
event_sequence: None,
session_id,
request_id,
}))
}
Self::Close {
manager,
session_id,
request_id,
} => {
manager.close(&session_id).await?;
Ok(Some(AgentMessage::ChatSessionClosed {
event_sequence: None,
session_id,
request_id,
}))
}
Self::ExtensionRespond {
manager,
session_id,
request_id,
extension_request_id,
response,
} => {
manager
.extension_response(
&session_id,
ExtensionResponse {
id: extension_request_id.clone(),
response,
},
)
.await?;
Ok(Some(AgentMessage::ChatExtensionResponded {
event_sequence: None,
session_id,
request_id,
extension_request_id,
}))
}
}
}
}
pub(crate) enum BackgroundRequest {
ChatHistoryList {
store: PiHistoryStore,
request_id: String,
},
ChatHistoryDelete {
store: PiHistoryStore,
request_id: String,
conversation_id: String,
},
DirectoryList {
browser: DirectoryBrowser,
request_id: String,
path: String,
},
ModelsList {
models: PiModels,
request_id: String,
},
}
impl BackgroundRequest {
pub(crate) fn request_id(&self) -> Option<&str> {
match self {
Self::DirectoryList { request_id, .. } | Self::ModelsList { request_id, .. } => {
Some(request_id)
}
Self::ChatHistoryList { request_id, .. }
| Self::ChatHistoryDelete { request_id, .. } => Some(request_id),
}
}
pub(crate) async fn execute(self) -> AgentResult<AgentMessage> {
match self {
Self::ChatHistoryList { store, request_id } => {
tokio::task::spawn_blocking(move || store.listed_message(request_id))
.await
.map_err(|_| background_task_failed())?
}
Self::ChatHistoryDelete {
store,
request_id,
conversation_id,
} => tokio::task::spawn_blocking(move || {
store.deleted_message(request_id, &conversation_id)
})
.await
.map_err(|_| background_task_failed())?,
Self::DirectoryList {
browser,
request_id,
path,
} => {
let message = tokio::task::spawn_blocking(move || {
directory_list_message(browser, request_id, path)
})
.await
.map_err(|_| background_task_failed())?;
Ok(message)
}
Self::ModelsList { models, request_id } => Ok(AgentMessage::ModelsListed {
request_id,
models: models.list().await?,
}),
}
}
}
fn direct(response: Option<RouteResponse>, output: Option<SessionOutput>) -> RouteAction {
RouteAction::Direct { response, output }
}
impl RouteResponse {
fn direct(message: AgentMessage) -> Self {
Self { message }
}
}
pub(crate) struct ProcessContext<'a> {
pub manager: &'a Arc<Mutex<SessionManager>>,
pub pi_manager: Option<&'a Arc<PiRpcManager>>,
pub policy: &'a CommandPolicy,
pub history_operations: &'a Arc<Mutex<()>>,
pub history_store: Option<&'a PiHistoryStore>,
pub directory_browser: Option<&'a DirectoryBrowser>,
pub pi_models: Option<&'a PiModels>,
pub request_owner: AuthenticatedRequestOwner,
}
pub(crate) async fn process_message(
message: BackendMessage,
context: ProcessContext<'_>,
) -> AgentResult<RouteAction> {
let ProcessContext {
manager,
pi_manager,
policy,
history_operations,
history_store,
directory_browser,
pi_models,
request_owner,
} = context;
match message {
BackendMessage::SessionCreate {
session_id,
command,
cwd,
cols,
rows,
} => {
let size =
TerminalSize::new(cols.unwrap_or(DEFAULT_COLS), rows.unwrap_or(DEFAULT_ROWS))?;
let spawn = manager
.lock()
.await
.prepare_session(session_id.clone(), command, cwd)?;
let session =
PtySession::spawn(spawn.session_id.clone(), spawn.command, spawn.cwd, size).await?;
let pid = session.pid();
let output = manager
.lock()
.await
.insert_session(spawn.session_id.clone(), session)?;
Ok(direct(
Some(RouteResponse::direct(AgentMessage::SessionStarted {
session_id,
pid,
})),
Some(output),
))
}
BackendMessage::TerminalInput {
session_id,
data_base64,
..
} => {
let bytes = decode_base64(&data_base64)?;
let output = manager.lock().await.session_output(&session_id)?;
output.session.write_input(&bytes).await?;
Ok(direct(None, None))
}
BackendMessage::TerminalResize {
session_id,
cols,
rows,
} => {
let size = TerminalSize::new(cols, rows)?;
let output = manager.lock().await.session_output(&session_id)?;
output.session.resize(size).await?;
Ok(direct(None, None))
}
BackendMessage::SessionKill { session_id } => {
let output = manager.lock().await.session_output(&session_id)?;
output.session.kill().await?;
Ok(direct(None, None))
}
BackendMessage::ChatSessionCreate {
session_id,
request_id,
cwd,
provider,
model_id,
thinking_level,
prompt,
} => {
validate_chat_ids(&session_id, &request_id)?;
validate_non_empty(&prompt, "prompt")?;
validate_model_pair(provider.as_deref(), model_id.as_deref())?;
let cwd = policy.validate_chat_cwd(&cwd)?;
let manager = configured_pi_manager(pi_manager)?.clone();
Ok(RouteAction::Chat(ChatRequest::Create {
manager,
history_operations: history_operations.clone(),
owner: request_owner.clone(),
request: CreatePiChat {
session_id,
cwd: cwd.into(),
provider,
model_id,
thinking_level,
request_id,
prompt,
},
}))
}
BackendMessage::ChatSessionResume {
session_id,
request_id,
conversation_id,
} => {
validate_chat_ids(&session_id, &request_id)?;
validate_non_empty(&conversation_id, "conversationId")?;
Ok(RouteAction::Chat(ChatRequest::Resume {
manager: configured_pi_manager(pi_manager)?.clone(),
history_operations: history_operations.clone(),
store: configured_history_store(history_store)?.clone(),
policy: policy.clone(),
session_id,
request_id,
conversation_id,
owner: request_owner.clone(),
}))
}
BackendMessage::ChatSessionSnapshot {
session_id,
request_id,
} => {
validate_chat_ids(&session_id, &request_id)?;
Ok(RouteAction::Chat(ChatRequest::Snapshot {
manager: configured_pi_manager(pi_manager)?.clone(),
session_id,
request_id,
owner: request_owner.clone(),
}))
}
BackendMessage::ChatModelSet {
session_id,
request_id,
provider,
model_id,
} => {
validate_chat_ids(&session_id, &request_id)?;
validate_non_empty(&provider, "provider")?;
validate_non_empty(&model_id, "modelId")?;
Ok(RouteAction::Chat(ChatRequest::SetModel {
manager: configured_pi_manager(pi_manager)?.clone(),
session_id,
request_id,
provider,
model_id,
owner: request_owner.clone(),
}))
}
BackendMessage::ChatThinkingSet {
session_id,
request_id,
level,
} => {
validate_chat_ids(&session_id, &request_id)?;
Ok(RouteAction::Chat(ChatRequest::SetThinking {
manager: configured_pi_manager(pi_manager)?.clone(),
session_id,
request_id,
level,
owner: request_owner.clone(),
}))
}
BackendMessage::ChatTranscriptPage {
session_id,
request_id,
cursor,
direction,
} => {
validate_chat_ids(&session_id, &request_id)?;
validate_non_empty(&cursor, "cursor")?;
Ok(RouteAction::Chat(ChatRequest::TranscriptPage {
manager: configured_pi_manager(pi_manager)?.clone(),
session_id,
request_id,
cursor,
direction,
owner: request_owner,
}))
}
BackendMessage::ChatCommandsList {
session_id,
request_id,
} => {
validate_chat_ids(&session_id, &request_id)?;
Ok(RouteAction::Chat(ChatRequest::CommandsList {
manager: configured_pi_manager(pi_manager)?.clone(),
session_id,
request_id,
owner: request_owner,
}))
}
BackendMessage::ChatFilesSearch {
session_id,
request_id,
query,
} => {
validate_chat_ids(&session_id, &request_id)?;
Ok(RouteAction::Chat(ChatRequest::FilesSearch {
manager: configured_pi_manager(pi_manager)?.clone(),
session_id,
request_id,
query,
owner: request_owner,
}))
}
BackendMessage::ChatPrompt {
session_id,
request_id,
text,
} => {
validate_chat_ids(&session_id, &request_id)?;
validate_non_empty(&text, "text")?;
Ok(RouteAction::Chat(ChatRequest::Prompt {
manager: configured_pi_manager(pi_manager)?.clone(),
session_id,
request_id,
text,
}))
}
BackendMessage::ChatAbort {
session_id,
request_id,
prompt_request_id,
} => {
validate_chat_ids(&session_id, &request_id)?;
if let Some(prompt_request_id) = prompt_request_id.as_deref() {
validate_non_empty(prompt_request_id, "promptRequestId")?;
}
Ok(RouteAction::Chat(ChatRequest::Abort {
manager: configured_pi_manager(pi_manager)?.clone(),
session_id,
request_id,
prompt_request_id,
}))
}
BackendMessage::ChatSessionClose {
session_id,
request_id,
} => {
validate_chat_ids(&session_id, &request_id)?;
Ok(RouteAction::Chat(ChatRequest::Close {
manager: configured_pi_manager(pi_manager)?.clone(),
session_id,
request_id,
}))
}
BackendMessage::ChatExtensionRespond {
session_id,
request_id,
extension_request_id,
response,
} => {
validate_chat_ids(&session_id, &request_id)?;
validate_non_empty(&extension_request_id, "extensionRequestId")?;
Ok(RouteAction::Chat(ChatRequest::ExtensionRespond {
manager: configured_pi_manager(pi_manager)?.clone(),
session_id,
request_id,
extension_request_id,
response: match response {
ChatExtensionResponse::Value(value) => ExtensionUiResponsePayload::Value(value),
ChatExtensionResponse::Confirmed(confirmed) => {
ExtensionUiResponsePayload::Confirmed(confirmed)
}
ChatExtensionResponse::Cancelled => ExtensionUiResponsePayload::Cancelled,
},
}))
}
BackendMessage::ChatHistoryList { request_id } => {
validate_request_id(&request_id)?;
let store = configured_history_store(history_store)?.clone();
Ok(RouteAction::Background(
BackgroundRequest::ChatHistoryList { store, request_id },
))
}
BackendMessage::ChatHistoryDelete {
request_id,
conversation_id,
} => {
validate_request_id(&request_id)?;
validate_non_empty(&conversation_id, "conversationId")?;
let store = configured_history_store(history_store)?.clone();
Ok(RouteAction::Background(
BackgroundRequest::ChatHistoryDelete {
store,
request_id,
conversation_id,
},
))
}
BackendMessage::DirectoryList { request_id, path } => {
validate_request_id(&request_id)?;
let browser = directory_browser.ok_or_else(|| {
AgentError::new(
ErrorCode::InvalidMessage,
"directory browser is not configured on this agent",
)
})?;
Ok(RouteAction::Background(BackgroundRequest::DirectoryList {
browser: browser.clone(),
request_id,
path,
}))
}
BackendMessage::ModelsList { request_id } => {
validate_request_id(&request_id)?;
let models = pi_models.ok_or_else(|| {
AgentError::new(
ErrorCode::InvalidMessage,
"models listing is not configured on this agent",
)
})?;
Ok(RouteAction::Background(BackgroundRequest::ModelsList {
models: models.clone(),
request_id,
}))
}
BackendMessage::SkillsSearch {
request_id,
provider,
query,
limit,
offset,
} => {
validate_request_id(&request_id)?;
if provider
.as_deref()
.is_some_and(|value| value.trim().is_empty())
{
return Err(invalid_skill_request("provider must not be empty"));
}
if limit == 0 || limit > MAX_SEARCH_LIMIT {
return Err(invalid_skill_request(format!(
"limit must be between 1 and {MAX_SEARCH_LIMIT}"
)));
}
Ok(RouteAction::Skill(SkillRequest::Search {
request_id,
provider,
query: query.unwrap_or_default(),
limit,
offset,
}))
}
BackendMessage::SkillsInstalledList {
request_id,
workspace,
include_all,
workspaces,
} => {
validate_request_id(&request_id)?;
validate_optional_workspace(workspace.as_deref())?;
for workspace in &workspaces {
validate_non_empty(workspace, "workspaces[]")?;
}
Ok(RouteAction::Skill(SkillRequest::ListInstalled {
request_id,
workspace,
include_all,
workspaces,
}))
}
BackendMessage::SkillsInstall {
request_id,
provider,
skill_id,
scope,
workspace,
confirm_update,
} => {
validate_request_id(&request_id)?;
validate_non_empty(&provider, "provider")?;
validate_non_empty(&skill_id, "skillId")?;
validate_optional_workspace(workspace.as_deref())?;
match (&scope, &workspace) {
(SkillScope::Global, Some(_)) => {
return Err(invalid_skill_request(
"workspace is not allowed for global scope",
));
}
(SkillScope::Workspace, None) => {
return Err(invalid_skill_request(
"workspace is required for workspace scope",
));
}
_ => {}
}
Ok(RouteAction::Skill(SkillRequest::Install {
request_id,
provider,
skill_id,
scope,
workspace,
confirm_update,
}))
}
BackendMessage::SkillsUninstall {
request_id,
installation_id,
} => {
validate_request_id(&request_id)?;
validate_non_empty(&installation_id, "installationId")?;
Ok(RouteAction::Skill(SkillRequest::Uninstall {
request_id,
installation_id,
}))
}
}
}
fn validate_request_id(request_id: &str) -> AgentResult<()> {
validate_non_empty(request_id, "requestId")
}
fn validate_chat_ids(session_id: &str, request_id: &str) -> AgentResult<()> {
validate_non_empty(session_id, "sessionId")?;
validate_request_id(request_id)
}
fn validate_model_pair(provider: Option<&str>, model_id: Option<&str>) -> AgentResult<()> {
match (provider, model_id) {
(None, None) => Ok(()),
(Some(provider), Some(model_id)) => {
validate_non_empty(provider, "provider")?;
validate_non_empty(model_id, "modelId")
}
_ => Err(AgentError::new(
ErrorCode::InvalidMessage,
"provider and modelId must be supplied together",
)),
}
}
fn validate_optional_workspace(workspace: Option<&str>) -> AgentResult<()> {
if let Some(workspace) = workspace {
validate_non_empty(workspace, "workspace")?;
}
Ok(())
}
fn validate_non_empty(value: &str, field: &str) -> AgentResult<()> {
if value.trim().is_empty() {
return Err(invalid_skill_request(format!("{field} must not be empty")));
}
Ok(())
}
fn invalid_skill_request(message: impl Into<String>) -> AgentError {
AgentError::new(ErrorCode::InvalidMessage, message)
}
fn configured_history_store(
history_store: Option<&PiHistoryStore>,
) -> AgentResult<&PiHistoryStore> {
history_store.ok_or_else(|| {
AgentError::new(
ErrorCode::HistoryNotConfigured,
"Pi session history is not available",
)
})
}
fn configured_pi_manager(
pi_manager: Option<&Arc<PiRpcManager>>,
) -> AgentResult<&Arc<PiRpcManager>> {
pi_manager.ok_or_else(|| {
AgentError::new(
ErrorCode::InvalidMessage,
"Pi RPC chat is not configured on this agent",
)
})
}
fn directory_list_message(
browser: DirectoryBrowser,
request_id: String,
path: String,
) -> AgentMessage {
match browser.list(std::path::Path::new(&path)) {
Ok(entries) => AgentMessage::DirectoryListed {
request_id,
path,
entries,
},
Err(err) => AgentMessage::DirectoryError {
request_id,
code: err.code().as_str().to_string(),
message: err.message().to_string(),
},
}
}
fn background_task_failed() -> AgentError {
AgentError::new(
ErrorCode::BackendDisconnected,
"background response task stopped",
)
}