Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
60 changes: 51 additions & 9 deletions src/gateway_integrations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -117,36 +117,78 @@ pub fn already_registered(name: &str) -> bool {
/// would be invalid TOML), so it's safe to call directly, not only behind the
/// caller's `already_registered` check. Returns a status line.
pub fn register(intg: &GatewayIntegration) -> String {
register_block(intg.name, intg.toml_block)
}

/// `register` for a block built at runtime (e.g. a path derived from the
/// current home dir). Same idempotency and append semantics.
pub fn register_block(name: &str, toml_block: &str) -> String {
let path = gateway_toml_path();
if already_registered(intg.name) {
return format!(
"skip {} already registered ({})",
intg.name,
path.display()
);
if already_registered(name) {
return format!("skip {name} already registered ({})", path.display());
}
let existing = fs::read_to_string(&path).unwrap_or_default();

let mut out = existing.trim_end().to_string();
if !out.is_empty() {
out.push_str("\n\n");
}
out.push_str(intg.toml_block.trim());
out.push_str(toml_block.trim());
out.push('\n');

if let Some(parent) = path.parent() {
let _ = fs::create_dir_all(parent);
}
match fs::write(&path, out) {
Ok(_) => format!(
"ok {} MCP registered behind the gateway ({})",
intg.name,
"ok {name} MCP registered behind the gateway ({})",
path.display()
),
Err(e) => format!("fail writing {}: {e}", path.display()),
}
}

/// The local rivalsearch MCP checkout, if present. Path is derived from the
/// current home dir at runtime -- never hardcoded into the binary.
fn rivalsearch_dir() -> PathBuf {
home()
.join(".config")
.join("opencode")
.join("mcp-servers")
.join("RivalSearchMCP")
}

fn rivalsearch_available() -> bool {
rivalsearch_dir().join("server.py").exists()
&& std::process::Command::new("uv")
.arg("--version")
.output()
.is_ok_and(|o| o.status.success())
}

fn rivalsearch_block() -> String {
let dir = rivalsearch_dir().to_string_lossy().replace('\\', "/");
format!(
"[servers.rivalsearch]\nkind = \"mcp_stdio\"\ncommand = \"uv\"\nargs = [\"--directory\", \"{dir}\", \"run\", \"fastmcp\", \"run\", \"server.py\"]"
)
}

/// Auto-register LOCAL stdio servers whose binaries the user already has
/// (leanctx, rivalsearch) -- no consent prompt: they spawn locally-installed
/// software and carry no secrets or network endpoints. Remote integrations
/// (github, ...) stay behind `init`'s consent flow. Idempotent; returns
/// status lines for the entries it actually added.
pub fn auto_register_local() -> Vec<String> {
let mut added = Vec::new();
if leanctx_installed() && !already_registered("leanctx") {
added.push(register(&LEANCTX));
}
if rivalsearch_available() && !already_registered("rivalsearch") {
added.push(register_block("rivalsearch", &rivalsearch_block()));
}
added
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down
14 changes: 14 additions & 0 deletions src/mcp_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ mod handoff;
pub(crate) mod item;
mod memory_tool;
mod review;
pub(crate) mod search;
pub(crate) mod types;

use crate::optimize;
Expand Down Expand Up @@ -874,6 +875,9 @@ impl AgentflareMcp {
}

fn load_gateway_config() -> gateway_registry::GatewayConfig {
for msg in crate::gateway_integrations::auto_register_local() {
eprintln!("agentflare: gateway {msg}");
}
let path = crate::paths::home()
.join(".agentflare")
.join("gateway.toml");
Expand Down Expand Up @@ -1371,6 +1375,16 @@ impl AgentflareMcp {
fn asset(&self, Parameters(req): Parameters<AssetRequest>) -> Result<String, ErrorData> {
self.asset_impl(req)
}

#[tool(
description = "Global unified search across four sources. type='store' (default) — FTS search across indexed store documents (artifacts, notes), grouped by doc_type. type='memory' — FTS search across brain.db observations (decisions, findings, patterns). type='code' — code search delegated to the gateway's leanctx ctx_search (regex, compressed output). type='web' — internet search via rivalsearch web_search tool. Returns { query, source, total, groups|results }."
)]
async fn search(
&self,
Parameters(req): Parameters<SearchRequest>,
) -> Result<String, ErrorData> {
self.search_impl(req).await
}
}

impl AgentflareMcp {
Expand Down
208 changes: 208 additions & 0 deletions src/mcp_server/search.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,208 @@
use super::*;

impl AgentflareMcp {
pub async fn search_impl(&self, req: SearchRequest) -> Result<String, ErrorData> {
let search_type = req.r#type.as_deref().unwrap_or("store");
match search_type {
"code" => self.search_code(&req).await,
"memory" => self.search_memory(&req),
"web" => self.search_web(&req).await,
"store" => self.search_store(&req),
other => Err(ErrorData::invalid_params(
format!("unknown type '{other}' — use store|memory|code|web"),
None,
)),
}
}

fn search_store(&self, req: &SearchRequest) -> Result<String, ErrorData> {
let q = req.query.trim();
if q.is_empty() {
return Err(ErrorData::invalid_params("query must not be empty", None));
}
let limit = req.limit.unwrap_or(20);

let ws_id = match self.with_backend_db(Self::resolve_workspace_id) {
Ok(Ok(id)) => id,
Ok(Err(e)) => return Err(ErrorData::internal_error(e.to_string(), None)),
Err(e) => return Err(e),
};

self.with_store(|store| -> Result<String, ErrorData> {
let matches = store
.doc_search(&ws_id, q, limit)
.map_err(|e| ErrorData::internal_error(e.to_string(), None))?;

let mut grouped: std::collections::BTreeMap<String, Vec<serde_json::Value>> =
std::collections::BTreeMap::new();

for m in matches {
let Some(doc) = store
.doc_get(&m.id)
.map_err(|e| ErrorData::internal_error(e.to_string(), None))?
else {
continue; // stale FTS row / doc deleted between search and get
};

let entry = serde_json::json!({
"id": doc.id,
"path": doc.path,
"title": doc.title,
"doc_type": doc.doc_type,
"snippet": m.snippet,
"score": m.score,
"source": doc.source,
"mime": doc.mime,
"size": doc.size,
"created_at": doc.created_at,
"updated_at": doc.updated_at,
});
grouped
.entry(if doc.doc_type.is_empty() {
"unknown".into()
} else {
doc.doc_type.clone()
})
.or_default()
.push(entry);
}

let result = serde_json::json!({
"query": q,
"source": "store",
"total": grouped.values().map(|v| v.len()).sum::<usize>(),
"groups": grouped,
});
Ok(serde_json::to_string_pretty(&result).unwrap_or_default())
})?
}

fn search_memory(&self, req: &SearchRequest) -> Result<String, ErrorData> {
let q = req.query.trim();
if q.is_empty() {
return Err(ErrorData::invalid_params("query must not be empty", None));
}
let limit = req.limit.unwrap_or(20);

let brain = match crate::memory::store::open() {
Ok(conn) => conn,
Err(e) => {
return Err(ErrorData::internal_error(
format!("failed to open brain.db: {e}"),
None,
));
}
};

let observations = match crate::memory::search::search(&brain, q, None, None, limit) {
Ok(obs) => obs,
Err(e) => {
return Err(ErrorData::internal_error(
format!("memory search failed: {e}"),
None,
));
}
};

let mut grouped: std::collections::BTreeMap<String, Vec<serde_json::Value>> =
std::collections::BTreeMap::new();

for obs in observations {
let entry = serde_json::json!({
"id": obs.id,
"type": obs.r#type,
"title": obs.title,
"content": obs.content,
"project": obs.project,
"session_id": obs.session_id,
"created_at": obs.created_at,
"updated_at": obs.updated_at,
"pinned": obs.pinned,
"topic_key": obs.topic_key,
});
let key = if obs.r#type.is_empty() {
"unknown".into()
} else {
obs.r#type.clone()
};
grouped.entry(key).or_default().push(entry);
}

Ok(serde_json::json!({
"query": q,
"source": "memory",
"total": grouped.values().map(|v| v.len()).sum::<usize>(),
"groups": grouped,
})
.to_string())
}

/// Delegates to the gateway's `leanctx` server (`ctx_search`, regex
/// action) -- same pattern as the web arm; no subprocess, no output
/// parsing. Unregistered/unavailable server degrades to an error payload.
async fn search_code(&self, req: &SearchRequest) -> Result<String, ErrorData> {
let q = req.query.trim();
if q.is_empty() {
return Err(ErrorData::invalid_params("query must not be empty", None));
}
let limit = req.limit.unwrap_or(50);
let root = Self::repo_root();

let guard = self.ensure_gateway_registry().await?;
let reg = guard.as_ref().expect("ensured above");

let args = serde_json::json!({
"pattern": q,
"path": root.to_string_lossy(),
"max_results": limit,
});

match reg.execute("leanctx", "ctx_search", args).await {
Ok(val) => Ok(serde_json::json!({
"source": "code",
"query": q,
"results": val,
})
.to_string()),
Err(e) => Ok(serde_json::json!({
"source": "code",
"query": q,
"error": format!("leanctx ctx_search failed: {e}"),
"results": [],
})
.to_string()),
}
}

async fn search_web(&self, req: &SearchRequest) -> Result<String, ErrorData> {
let q = req.query.trim();
if q.is_empty() {
return Err(ErrorData::invalid_params("query must not be empty", None));
}
let limit = req.limit.unwrap_or(10);

let guard = self.ensure_gateway_registry().await?;
let reg = guard.as_ref().expect("ensured above");

let args = serde_json::json!({
"query": q,
"max_results": limit,
});

match reg.execute("rivalsearch", "web_search", args).await {
Ok(val) => Ok(serde_json::json!({
"source": "web",
"query": q,
"results": val,
})
.to_string()),
Err(e) => Ok(serde_json::json!({
"source": "web",
"query": q,
"error": format!("rivalsearch web_search failed: {e}"),
"results": [],
})
.to_string()),
}
}
}
1 change: 1 addition & 0 deletions src/mcp_server/tests/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -351,3 +351,4 @@ mod action_tests;
mod artifact_tests;
mod asset_tests;
mod item_tests;
mod search_tests;
Loading
Loading