Skip to content

Commit cf36934

Browse files
committed
perf(gateway): Phase 1 search_tools — drop readiness resolve, Arc cache
Warm path no longer full-space feature-resolves for enrichment; session index is Arc-shared and warmer embeds in one batch. Measured warm HogQL ~346ms → ~33ms. Signed-off-by: crimsonsunset <jsangio1@gmail.com>
1 parent 102bb92 commit cf36934

7 files changed

Lines changed: 169 additions & 131 deletions

File tree

crates/mcpmux-gateway/src/services/embedding_warmer.rs

Lines changed: 24 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -181,20 +181,37 @@ impl EmbeddingWarmer {
181181
return Ok(());
182182
}
183183

184-
let mut records = Vec::new();
185184
let embed_started = Instant::now();
185+
let mut ordered_hashes = Vec::with_capacity(missing_hashes.len());
186+
let mut ordered_haystacks = Vec::with_capacity(missing_hashes.len());
186187
for content_hash in missing_hashes {
187188
let Some(haystack) = haystacks_by_hash.get(&content_hash).cloned() else {
188189
continue;
189190
};
191+
ordered_hashes.push(content_hash);
192+
ordered_haystacks.push(haystack);
193+
}
190194

191-
let Some(vectors) = self.embeddings.embed_documents(&[haystack], None) else {
192-
continue;
193-
};
194-
let Some(vector) = vectors.into_iter().next() else {
195-
continue;
196-
};
195+
let Some(vectors) = self
196+
.embeddings
197+
.embed_documents(&ordered_haystacks, None)
198+
else {
199+
info!(
200+
space_id = %space_id,
201+
server_id,
202+
embedded = 0,
203+
skipped_present,
204+
missing = ordered_hashes.len(),
205+
embed_ms = embed_started.elapsed().as_millis() as u64,
206+
model_version = self.embeddings.model_version(),
207+
model_state = ?self.embeddings.state(),
208+
"[embed] warm batch done"
209+
);
210+
return Ok(());
211+
};
197212

213+
let mut records = Vec::with_capacity(vectors.len());
214+
for (content_hash, vector) in ordered_hashes.into_iter().zip(vectors) {
198215
records.push(EmbeddingRecord {
199216
content_hash: content_hash.clone(),
200217
model_version: self.embeddings.model_version().to_string(),

crates/mcpmux-gateway/src/services/meta_tools/meta_tool_common.rs

Lines changed: 38 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -181,11 +181,32 @@ pub(crate) fn format_invoke_not_ready_action_with_name(
181181
}
182182
}
183183

184-
/// Display names and pre-configured default param keys per installed server.
185-
pub(crate) async fn build_installed_server_meta_maps(
184+
/// Search-hit enrichment built from one installed-server list + pool statuses.
185+
///
186+
/// `binding_server_ids` should be the unique `server_id`s from the caller's
187+
/// **active** tool index (not a full `resolve_feature_sets` catalog scan).
188+
pub(crate) struct SearchServerEnrichment {
189+
pub readiness_map: HashMap<String, &'static str>,
190+
pub display_names: HashMap<String, String>,
191+
pub prefilled_params_by_server: HashMap<String, Vec<String>>,
192+
}
193+
194+
/// Whether the caller omitted or blanked the search query.
195+
pub(crate) fn is_query_empty(query: Option<&str>) -> bool {
196+
query.map(str::trim).is_none_or(str::is_empty)
197+
}
198+
199+
/// Build readiness + display/prefill maps without re-resolving FeatureSets.
200+
///
201+
/// One `list_for_space` for installed servers and one `get_all_statuses` pass.
202+
/// Extra `server_ids` (e.g. inactive widen) are labeled `bindable` unless they
203+
/// also appear in `binding_server_ids`.
204+
pub(crate) async fn build_search_server_enrichment(
186205
call: &MetaToolCall<'_>,
187206
space_id: &Uuid,
188-
) -> Result<(HashMap<String, String>, HashMap<String, Vec<String>>), MetaToolError> {
207+
binding_server_ids: &HashSet<String>,
208+
extra_server_ids: &HashSet<String>,
209+
) -> Result<SearchServerEnrichment, MetaToolError> {
189210
let installed = call
190211
.ctx
191212
.installed_server_repo
@@ -194,63 +215,31 @@ pub(crate) async fn build_installed_server_meta_maps(
194215
.map_err(|e| MetaToolError::Internal(e.to_string()))?;
195216

196217
let mut display_names = HashMap::new();
197-
let mut prefilled_params = HashMap::new();
218+
let mut prefilled_params_by_server = HashMap::new();
219+
let mut installed_by_id: HashMap<String, mcpmux_core::InstalledServer> = HashMap::new();
198220
for server in installed {
199221
display_names.insert(server.server_id.clone(), server.display_name().to_string());
200222
if !server.default_params.is_empty() {
201223
let mut keys: Vec<String> = server.default_params.keys().cloned().collect();
202224
keys.sort();
203-
prefilled_params.insert(server.server_id.clone(), keys);
225+
prefilled_params_by_server.insert(server.server_id.clone(), keys);
204226
}
227+
installed_by_id.insert(server.server_id.clone(), server);
205228
}
206229

207-
Ok((display_names, prefilled_params))
208-
}
209-
210-
/// Whether the caller omitted or blanked the search query.
211-
pub(crate) fn is_query_empty(query: Option<&str>) -> bool {
212-
query.map(str::trim).is_none_or(str::is_empty)
213-
}
214-
215-
/// Point-in-time `readiness` label per server for search hit enrichment.
216-
pub(crate) async fn build_server_readiness_map(
217-
call: &MetaToolCall<'_>,
218-
space_id: &Uuid,
219-
resolved: &ResolvedFeatureSet,
220-
) -> Result<HashMap<String, &'static str>, MetaToolError> {
221-
let binding_features = call
222-
.ctx
223-
.feature_service
224-
.resolve_feature_sets(&space_id.to_string(), &resolved.feature_set_ids)
225-
.await?;
226-
let binding_servers: HashSet<String> = binding_features
227-
.iter()
228-
.map(|f| f.server_id.clone())
229-
.collect();
230-
231-
let installed = call
232-
.ctx
233-
.installed_server_repo
234-
.list_for_space(&space_id.to_string())
235-
.await
236-
.map_err(|e| MetaToolError::Internal(e.to_string()))?;
237-
let installed_by_id: HashMap<String, mcpmux_core::InstalledServer> = installed
238-
.into_iter()
239-
.map(|s| (s.server_id.clone(), s))
240-
.collect();
241-
242230
let pool_statuses = call.ctx.server_manager.get_all_statuses(*space_id).await;
243231

244-
let server_ids: HashSet<String> = binding_servers
232+
let server_ids: HashSet<String> = binding_server_ids
245233
.iter()
234+
.chain(extra_server_ids.iter())
246235
.chain(installed_by_id.keys())
247236
.cloned()
248237
.collect();
249238

250-
let map = server_ids
239+
let readiness_map = server_ids
251240
.into_iter()
252241
.map(|server_id| {
253-
let in_binding = binding_servers.contains(&server_id);
242+
let in_binding = binding_server_ids.contains(&server_id);
254243
let connection_status = pool_statuses
255244
.get(&server_id)
256245
.map(|(status, _, _, _)| *status)
@@ -264,7 +253,12 @@ pub(crate) async fn build_server_readiness_map(
264253
(server_id, readiness)
265254
})
266255
.collect();
267-
Ok(map)
256+
257+
Ok(SearchServerEnrichment {
258+
readiness_map,
259+
display_names,
260+
prefilled_params_by_server,
261+
})
268262
}
269263

270264
/// Common path for every write tool: build payload, ask broker, run the

crates/mcpmux-gateway/src/services/meta_tools/registry.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -82,7 +82,7 @@ pub struct MetaToolContext {
8282
/// Per-server log tail reader (`current.log`); same source as the desktop UI.
8383
pub log_manager: Arc<ServerLogManager>,
8484
/// Per-session active tool index for `mcpmux_search_tools` (fingerprint-keyed).
85-
pub search_cache: Arc<DashMap<String, (u64, ToolIndex)>>,
85+
pub search_cache: Arc<DashMap<String, (u64, Arc<ToolIndex>)>>,
8686
/// Global embedding vectors keyed by content hash.
8787
pub embedding_store: Arc<DashMap<String, Vec<f32>>>,
8888
/// Persistent embedding repository backing `embedding_store` hydration.

crates/mcpmux-gateway/src/services/meta_tools/search_tools.rs

Lines changed: 61 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -1,18 +1,19 @@
11
//! `mcpmux_search_tools` — hybrid search and browse over the active tool index.
22
3+
use std::collections::HashSet;
4+
use std::sync::Arc;
5+
use std::time::Instant;
6+
37
use async_trait::async_trait;
48
use rmcp::model::CallToolResult;
59
use serde_json::{json, Value};
6-
use std::collections::HashSet;
7-
use std::time::Instant;
810
use tracing::{debug, info};
911
use uuid::Uuid;
1012

1113
use crate::services::embedding::model_state_label;
1214

1315
use super::meta_tool_common::{
14-
build_installed_server_meta_maps, build_server_readiness_map, caller_resolution,
15-
is_query_empty, text_result,
16+
build_search_server_enrichment, caller_resolution, is_query_empty, text_result,
1617
};
1718
use super::registry::{feature_set_ids_fingerprint, MetaTool, MetaToolCall, MetaToolError};
1819
use super::search_tools_index::{
@@ -161,25 +162,28 @@ impl MetaTool for SearchToolsTool {
161162
debug!(query_id = %query_id, query, "[search] query text");
162163
}
163164

164-
let readiness_started = Instant::now();
165-
let readiness_map = build_server_readiness_map(&call, &space_id, &resolved).await?;
166-
let readiness_ms = readiness_started.elapsed().as_millis() as u64;
167-
168-
let installed_meta_started = Instant::now();
169-
let (server_display_names, prefilled_params_by_server) =
170-
build_installed_server_meta_maps(&call, &space_id).await?;
171-
let installed_meta_ms = installed_meta_started.elapsed().as_millis() as u64;
172-
173165
let mut index_cache_hit = false;
174166
let active_index_started = Instant::now();
175-
let active_index = if let Some(session_id) = call.session_id {
176-
if let Some(entry) = call.ctx.search_cache.get(session_id) {
177-
let (cached_fp, cached_index) = entry.value();
178-
if *cached_fp == fingerprint {
179-
index_cache_hit = true;
180-
cached_index.clone()
167+
let active_index: Arc<crate::services::ToolIndex> =
168+
if let Some(session_id) = call.session_id {
169+
if let Some(entry) = call.ctx.search_cache.get(session_id) {
170+
let (cached_fp, cached_index) = entry.value();
171+
if *cached_fp == fingerprint {
172+
index_cache_hit = true;
173+
Arc::clone(cached_index)
174+
} else {
175+
drop(entry);
176+
build_and_cache_active_index(
177+
&call,
178+
&space_id,
179+
&resolved,
180+
fingerprint,
181+
session_id,
182+
query_id.as_str(),
183+
)
184+
.await?
185+
}
181186
} else {
182-
drop(entry);
183187
build_and_cache_active_index(
184188
&call,
185189
&space_id,
@@ -191,19 +195,8 @@ impl MetaTool for SearchToolsTool {
191195
.await?
192196
}
193197
} else {
194-
build_and_cache_active_index(
195-
&call,
196-
&space_id,
197-
&resolved,
198-
fingerprint,
199-
session_id,
200-
query_id.as_str(),
201-
)
202-
.await?
203-
}
204-
} else {
205-
build_active_index(&call, &space_id, &resolved, query_id.as_str()).await?
206-
};
198+
build_active_index(&call, &space_id, &resolved, query_id.as_str()).await?
199+
};
207200
let active_index_ms = active_index_started.elapsed().as_millis() as u64;
208201

209202
debug!(
@@ -214,13 +207,18 @@ impl MetaTool for SearchToolsTool {
214207
"[search] active index ready"
215208
);
216209

217-
let clone_started = Instant::now();
218-
let mut index = active_index.clone();
219-
let index_clone_ms = clone_started.elapsed().as_millis() as u64;
210+
let binding_server_ids: HashSet<String> = active_index
211+
.iter()
212+
.map(|entry| entry.server_id.clone())
213+
.collect();
220214

221215
let mut inactive_tool_count = 0usize;
222216
let mut inactive_widen_ms = 0_u64;
217+
let mut index_clone_ms = 0_u64;
218+
let mut extra_server_ids: HashSet<String> = HashSet::new();
223219

220+
// Copy-on-write only when inactive widen mutates the haystack.
221+
let mut widened_index: Option<crate::services::ToolIndex> = None;
224222
if include_inactive {
225223
debug!(
226224
query_id = %query_id,
@@ -242,12 +240,16 @@ impl MetaTool for SearchToolsTool {
242240
crate::services::tool_discovery::ToolDiscoveryService::build_inactive_index(
243241
&inactive,
244242
);
243+
let clone_started = Instant::now();
244+
let mut index = (*active_index).clone();
245+
index_clone_ms = clone_started.elapsed().as_millis() as u64;
245246
let active_keys: HashSet<(String, String)> = index
246247
.iter()
247248
.map(|e| (e.server_id.clone(), e.feature_name.clone()))
248249
.collect();
249250
let before_merge = index.len();
250251
for entry in inactive_index {
252+
extra_server_ids.insert(entry.server_id.clone());
251253
let key = (entry.server_id.clone(), entry.feature_name.clone());
252254
if !active_keys.contains(&key) {
253255
index.push(entry);
@@ -263,8 +265,29 @@ impl MetaTool for SearchToolsTool {
263265
inactive_widen_ms,
264266
"[search] inactive widen complete"
265267
);
268+
widened_index = Some(index);
266269
}
267270

271+
let index: &[crate::services::ToolIndexEntry] = widened_index
272+
.as_deref()
273+
.unwrap_or(active_index.as_slice());
274+
let merged_index_len = index.len();
275+
276+
let enrichment_started = Instant::now();
277+
let enrichment = build_search_server_enrichment(
278+
&call,
279+
&space_id,
280+
&binding_server_ids,
281+
&extra_server_ids,
282+
)
283+
.await?;
284+
let readiness_ms = enrichment_started.elapsed().as_millis() as u64;
285+
// installed list is folded into enrichment — keep the timing field for log continuity.
286+
let installed_meta_ms = 0_u64;
287+
let readiness_map = enrichment.readiness_map;
288+
let server_display_names = enrichment.display_names;
289+
let prefilled_params_by_server = enrichment.prefilled_params_by_server;
290+
268291
let (hydrate_ms, hydrated_missing_count) = if effective_query.is_some() {
269292
let hydrate =
270293
hydrate_active_embeddings(&call, query_id.as_str(), active_index.as_slice())
@@ -283,7 +306,7 @@ impl MetaTool for SearchToolsTool {
283306

284307
let rank_started = Instant::now();
285308
let result = crate::services::tool_discovery::ToolDiscoveryService::search(
286-
&index,
309+
index,
287310
effective_query,
288311
server_id_filter,
289312
detail_level,
@@ -472,7 +495,7 @@ impl MetaTool for SearchToolsTool {
472495
post_ms,
473496
accounted_ms,
474497
unaccounted_ms = total_ms.saturating_sub(accounted_ms),
475-
merged_index = index.len(),
498+
merged_index = merged_index_len,
476499
include_inactive,
477500
scope_all,
478501
server_id_set,

0 commit comments

Comments
 (0)