1use std::sync::Arc;
7
8use nexql_index::{
9 BuildDepth, BuildMode, BuildRequest, CatalogDb, Embedder, IndexQueryService, IndexScope,
10 IndexStore, PgCatalogDb, QueryPolicyFilter, RefResolution, SearchOptions, build_index,
11};
12use nexql_policy::{
13 PolicyFilter, SqlDecision, enforce_read_table_policy, select_table_refs, validate_readonly_sql,
14};
15use serde_json::{Value, json};
16
17use crate::cell_json::{redact_pii_in_payload, rows_to_json_array};
18use crate::error::ToolError;
19use crate::export::{ExportFormat, columns_from_rows, rows_to_csv, rows_to_sql_insert};
20use crate::plan::{analyze_deep_plan, build_explain_sql, extract_plan_metrics};
21use crate::registry::ToolName;
22use crate::schema::{ToolSpec, active_tools};
23use crate::session::ToolSession;
24use crate::sql::{self, REPORT_LIMIT_DEFAULT, SLOW_QUERIES_DEFAULT, parse_ref};
25use crate::write::{
26 apply_ddl, create_index_concurrently, edit_row, execute_sql, import_data, run_maintenance,
27 terminate_query,
28};
29
30const SEARCH_SCHEMA_LIMIT: usize = 10;
32
33const NO_INDEX_HINT: &str =
34 "No schema index configured — call the 'rebuild_index' tool to build an index.";
35
36#[derive(Debug, Clone)]
37pub struct ToolOutcome {
38 pub text: String,
39 pub structured: Option<Value>,
40 pub is_error: bool,
41}
42
43impl ToolOutcome {
44 pub fn ok_json(value: Value) -> Self {
50 let value = ensure_structured_object(value);
51 let text = serde_json::to_string_pretty(&value).unwrap_or_else(|_| value.to_string());
52 Self {
53 text,
54 structured: Some(value),
55 is_error: false,
56 }
57 }
58
59 pub fn err(msg: impl Into<String>) -> Self {
60 let message = msg.into();
61 Self {
62 text: message.clone(),
63 structured: Some(json!({ "error": message })),
64 is_error: true,
65 }
66 }
67}
68
69fn ensure_structured_object(value: Value) -> Value {
71 match value {
72 Value::Array(rows) => json!({ "rows": rows }),
73 other => other,
74 }
75}
76
77pub struct ToolRouter {
78 session: Arc<ToolSession>,
79 index_override: Option<Option<IndexStore>>,
81 use_semantic: bool,
83 embedder: Option<Arc<dyn Embedder>>,
84 specs: Vec<ToolSpec>,
85 managed_extension: bool,
86}
87
88impl ToolRouter {
89 pub fn new(session: Arc<ToolSession>) -> Self {
90 Self {
91 session,
92 index_override: None,
93 use_semantic: false,
94 embedder: None,
95 specs: active_tools(),
96 managed_extension: false,
97 }
98 }
99
100 pub fn with_index_store(session: Arc<ToolSession>, store: Option<IndexStore>) -> Self {
102 Self {
103 session,
104 index_override: Some(store),
105 use_semantic: false,
106 embedder: None,
107 specs: active_tools(),
108 managed_extension: false,
109 }
110 }
111
112 pub fn with_semantic(
114 mut self,
115 use_semantic: bool,
116 embedder: Option<Arc<dyn Embedder>>,
117 ) -> Self {
118 self.use_semantic = use_semantic;
119 self.embedder = embedder;
120 self
121 }
122
123 pub fn with_profile(mut self, profile: crate::registry::ToolProfile) -> Self {
125 self.specs = crate::schema::tools_for_profile(profile);
126 self
127 }
128
129 pub fn with_managed_extension(mut self, enabled: bool) -> Self {
131 self.managed_extension = enabled;
132 if enabled {
133 const BLOCKED: &[ToolName] = &[
134 ToolName::SetupConnection,
135 ToolName::SaveProfile,
136 ToolName::TestProfile,
137 ToolName::ExportProfile,
138 ToolName::ImportProfile,
139 ];
140 self.specs.retain(|s| !BLOCKED.contains(&s.name));
141 }
142 self
143 }
144
145 pub fn specs(&self) -> &[ToolSpec] {
146 &self.specs
147 }
148
149 fn index_store(&self) -> Option<&IndexStore> {
150 match &self.index_override {
151 Some(inner) => inner.as_ref(),
152 None => self.session.index_store.as_ref(),
153 }
154 }
155
156 fn query_filter(&self) -> QueryPolicyFilter {
157 policy_to_query_filter(&self.session.filter())
158 }
159
160 pub async fn call(&self, name: &str, args: Value) -> ToolOutcome {
161 let outcome = match self.call_inner(name, args).await {
162 Ok(out) => out,
163 Err(e) => ToolOutcome::err(e.to_string()),
164 };
165 self.tag_outcome_with_context(outcome).await
166 }
167
168 async fn tag_outcome_with_context(&self, mut outcome: ToolOutcome) -> ToolOutcome {
169 let (connection_id, database) = self.session.active_context().await;
170 let access_mode = match self.session.access_mode() {
171 nexql_policy::AccessMode::Read => "read",
172 nexql_policy::AccessMode::Write => "write",
173 nexql_policy::AccessMode::Admin => "admin",
174 };
175 let mut freshness: Option<serde_json::Value> = None;
176 if let Some(store) = self.session.index_store.as_ref() {
177 let base = store.base_dir(&connection_id, &database);
178 if let Ok(Some(manifest)) = store.read_manifest(&base) {
179 let stale = self.session.is_index_stale(&connection_id, &database);
180 let mut freshness_obj = json!({
181 "indexedAt": manifest.indexed_at,
182 "schemaFingerprint": manifest.schema_fingerprint,
183 "stale": stale,
184 });
185 if stale {
186 freshness_obj["reason"] = json!("schema_changed");
187 }
188 freshness = Some(freshness_obj);
189 } else {
190 freshness = Some(json!({ "stale": true, "reason": "no_index" }));
191 }
192 }
193 if let Some(ref mut structured) = outcome.structured
194 && let Some(obj) = structured.as_object_mut()
195 {
196 if !obj.contains_key("connectionId") {
197 obj.insert("connectionId".into(), json!(connection_id));
198 }
199 if !obj.contains_key("database") {
200 obj.insert("database".into(), json!(database));
201 }
202 if !obj.contains_key("accessMode") {
203 obj.insert("accessMode".into(), json!(access_mode));
204 }
205 if let Some(ref f) = freshness {
206 obj.insert("freshness".into(), f.clone());
207 }
208 }
209 let header = format!(
210 "[context connectionId={connection_id} database={database} accessMode={access_mode}]\n"
211 );
212 if !outcome.text.starts_with("[context ") {
213 outcome.text = format!("{header}{}", outcome.text);
214 }
215 outcome
216 }
217
218 async fn call_inner(&self, name: &str, args: Value) -> Result<ToolOutcome, ToolError> {
219 let tool = ToolName::parse(name).ok_or_else(|| ToolError::Unknown(name.to_string()))?;
220 match tool {
221 ToolName::ListConnections => Ok(self.list_connections()),
222 ToolName::ListDatabases => self.list_databases(&args).await,
223 ToolName::ListSchemas => self.list_schemas().await,
224 ToolName::ListObjects => self.list_objects(&args).await,
225 ToolName::GetCurrentContext => self.get_current_context().await,
226 ToolName::SwitchConnection => self.switch_connection(&args).await,
227 ToolName::RunSelect => self.run_select(&args).await,
228 ToolName::ExplainQuery => self.explain_query(&args).await,
229 ToolName::SearchSchema => self.search_schema(&args).await,
230 ToolName::DescribeObject => self.describe_object(&args).await,
231 ToolName::GetJoinPath => self.get_join_path(&args).await,
232 ToolName::SampleValues => self.sample_values(&args).await,
233 ToolName::GetDdl => self.get_ddl(&args).await,
234 ToolName::TableStats => self.table_stats(&args).await,
235 ToolName::IndexUsage => self.index_usage(&args).await,
236 ToolName::ListRunningQueries => self.list_running_queries().await,
237 ToolName::FindBlockingLocks => self.find_blocking_locks().await,
238 ToolName::SlowQueries => self.slow_queries(&args).await,
239 ToolName::DbHealthCheck => self.db_health_check().await,
240 ToolName::ExplainAnalyze => self.explain_analyze(&args).await,
241 ToolName::AnalyzeQueryPlan => self.analyze_query_plan(&args).await,
242 ToolName::GetIndexStatus => self.get_index_status().await,
243 ToolName::ListExtensions => self.list_extensions().await,
244 ToolName::ServerSettings => self.server_settings().await,
245 ToolName::SuggestIndexes => self.suggest_indexes(&args).await,
246 ToolName::FindUnusedIndexes => self.find_unused_indexes(&args).await,
247 ToolName::BloatReport => self.bloat_report(&args).await,
248 ToolName::FindMissingFks => self.find_missing_fks(&args).await,
249 ToolName::ExportQuery => self.export_query(&args).await,
250 ToolName::ListRoles => self.list_roles(&args).await,
251 ToolName::DbDashboard => self.db_dashboard().await,
252 ToolName::DeepPlanAnalysis => self.deep_plan_analysis(&args).await,
253 ToolName::SchemaDiff => self.schema_diff(&args).await,
254 ToolName::GenerateMigration => self.generate_migration(&args).await,
255 ToolName::ExecuteSql => self.execute_sql_tool(&args).await,
256 ToolName::EditRow => self.edit_row_tool(&args).await,
257 ToolName::ImportData => self.import_data_tool(&args).await,
258 ToolName::ApplyDdl => self.apply_ddl_tool(&args).await,
259 ToolName::CreateIndexConcurrently => self.create_index_concurrently_tool(&args).await,
260 ToolName::RunMaintenance => self.run_maintenance_tool(&args).await,
261 ToolName::TerminateQuery => self.terminate_query_tool(&args).await,
262 ToolName::ResolveTarget => self.resolve_target(&args).await,
263 ToolName::DiscoverTools => self.discover_tools(&args).await,
264 ToolName::AutoTuneQuery => self.auto_tune_query(&args).await,
265 ToolName::CheckDdlSafety => self.check_ddl_safety_tool(&args).await,
266 ToolName::RebuildIndex => self.rebuild_index_tool(&args).await,
267 ToolName::RefreshIndex => self.refresh_index_tool(&args).await,
268 ToolName::RunDoctor => self.run_doctor_tool().await,
269 ToolName::SetupConnection => self.setup_connection_tool(&args).await,
270 ToolName::SaveProfile => self.save_profile_tool(&args).await,
271 ToolName::TestProfile => self.test_profile_tool(&args).await,
272 ToolName::ExportProfile => self.export_profile_tool(&args).await,
273 ToolName::ImportProfile => self.import_profile_tool(&args).await,
274 }
275 }
276
277 fn require_write(&self) -> Result<(), ToolError> {
278 if !self.session.access_mode().allows_writes() {
279 return Err(ToolError::Execution(
280 "write tools require --access-mode write or admin (current session: read)".into(),
281 ));
282 }
283 Ok(())
284 }
285
286 fn require_admin(&self) -> Result<(), ToolError> {
287 if !self.session.access_mode().allows_admin() {
288 return Err(ToolError::Execution(
289 "admin tools require --access-mode admin".into(),
290 ));
291 }
292 Ok(())
293 }
294
295 async fn execute_sql_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
296 self.require_write()?;
297 let sql = args
298 .get("sql")
299 .and_then(|v| v.as_str())
300 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
301 let dry_run = args
302 .get("dry_run")
303 .and_then(|v| v.as_bool())
304 .unwrap_or(false);
305 execute_sql(&self.session, sql, dry_run).await
306 }
307
308 async fn edit_row_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
309 self.require_write()?;
310 edit_row(&self.session, args).await
311 }
312
313 async fn import_data_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
314 self.require_write()?;
315 import_data(&self.session, args).await
316 }
317
318 async fn apply_ddl_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
319 self.require_admin()?;
320 let sql = args
321 .get("sql")
322 .and_then(|v| v.as_str())
323 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
324 let dry_run = args
325 .get("dry_run")
326 .and_then(|v| v.as_bool())
327 .unwrap_or(false);
328 apply_ddl(&self.session, sql, dry_run).await
329 }
330
331 async fn create_index_concurrently_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
332 self.require_admin()?;
333 let sql = args
334 .get("sql")
335 .and_then(|v| v.as_str())
336 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
337 create_index_concurrently(&self.session, sql).await
338 }
339
340 async fn run_maintenance_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
341 self.require_admin()?;
342 run_maintenance(&self.session, args).await
343 }
344
345 async fn terminate_query_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
346 self.require_admin()?;
347 terminate_query(&self.session, args).await
348 }
349
350 async fn live_databases_by_connection(
353 &self,
354 ) -> std::collections::HashMap<String, std::collections::HashSet<String>> {
355 use std::collections::HashMap;
356 let mut map = HashMap::new();
357 for conn in self.session.connections() {
358 if let Ok(names) = self.list_database_names_for(&conn).await {
359 map.insert(conn.id.clone(), names.into_iter().collect());
360 }
361 }
362 map
363 }
364
365 async fn list_database_names_for(
366 &self,
367 conn: &crate::session::ConnectionInfo,
368 ) -> Result<Vec<String>, ToolError> {
369 let client = if self.session.active_context().await.0 == conn.id {
370 self.session.checkout().await?
371 } else {
372 let pool_opts = self.session.pool_opts();
373 let pool = nexql_conn::create_pool(&conn.params, &pool_opts).await?;
374 nexql_conn::checkout_guarded(&pool, &pool_opts).await?
375 };
376 let rows = client
377 .query(
378 "SELECT datname FROM pg_database WHERE datistemplate = false ORDER BY datname",
379 &[],
380 )
381 .await?;
382 Ok(rows.iter().map(|r| r.get(0)).collect())
383 }
384
385 async fn resolve_target(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
386 let hint = args
387 .get("hint")
388 .and_then(|v| v.as_str())
389 .map(str::trim)
390 .filter(|s| !s.is_empty());
391 let object_hint = args
392 .get("objectHint")
393 .and_then(|v| v.as_str())
394 .map(str::trim)
395 .filter(|s| !s.is_empty());
396 if hint.is_none() && object_hint.is_none() {
397 return Err(ToolError::InvalidArgs(
398 "At least one of \"hint\" or \"objectHint\" is required.".into(),
399 ));
400 }
401
402 let connections = self.session.connections();
403 if connections.is_empty() {
404 return Ok(ToolOutcome::err("No connections configured."));
405 }
406
407 #[derive(Clone)]
408 struct Candidate {
409 connection_id: String,
410 database: String,
411 }
412 fn key_of(c: &Candidate) -> String {
413 format!("{}\u{0}{}", c.connection_id, c.database)
414 }
415
416 let indexed: Vec<(String, String)> = self
417 .index_store()
418 .map(|store| store.list_indexed_databases().unwrap_or_default())
419 .unwrap_or_default();
420
421 let live_dbs = self.live_databases_by_connection().await;
422
423 let mut seen = std::collections::HashSet::new();
424 let mut candidates: Vec<Candidate> = Vec::new();
425 let mut add_candidate = |connection_id: &str, database: &str| {
426 if !connections.iter().any(|c| c.id == connection_id) {
427 return;
428 }
429 let key = format!("{connection_id}\u{0}{database}");
430 if !seen.insert(key) {
431 return;
432 }
433 candidates.push(Candidate {
434 connection_id: connection_id.to_string(),
435 database: database.to_string(),
436 });
437 };
438 for (cid, db) in &indexed {
439 if let Some(set) = live_dbs.get(cid) {
440 if set.contains(db) {
441 add_candidate(cid, db);
442 } else if let Some(store) = self.index_store() {
443 let _ = store.clear_index(cid, db);
444 }
445 }
446 }
447 for c in &connections {
448 if let Some(dbs) = live_dbs.get(&c.id) {
449 for db in dbs {
450 add_candidate(&c.id, db);
451 }
452 } else {
453 let db = c.database.clone().unwrap_or_else(|| "postgres".into());
454 add_candidate(&c.id, &db);
455 }
456 }
457
458 let mut scored: std::collections::HashMap<String, (Candidate, f64, Vec<String>)> =
459 std::collections::HashMap::new();
460
461 if let Some(hint) = hint {
462 for c in &candidates {
463 let Some(conn) = connections.iter().find(|x| x.id == c.connection_id) else {
464 continue;
465 };
466 let fields: [(&str, &str); 3] = [
467 ("connection name", conn.name.as_str()),
468 ("host", conn.host.as_deref().unwrap_or("")),
469 ("database", c.database.as_str()),
470 ];
471 let mut best = 0.0f64;
472 let mut best_field = "";
473 for (label, value) in fields {
474 let s = fuzzy_score(hint, value);
475 if s > best {
476 best = s;
477 best_field = label;
478 }
479 }
480 if best > 0.0 {
481 let entry = scored
482 .entry(key_of(c))
483 .or_insert_with(|| (c.clone(), 0.0, Vec::new()));
484 entry.1 += best;
485 entry
486 .2
487 .push(format!("{best_field} matched hint \"{hint}\" ({best:.0})"));
488 }
489 }
490 }
491
492 if let Some(object_hint) = object_hint
493 && let Some(store) = self.index_store()
494 {
495 let filter = self.query_filter();
496 for (cid, db) in &indexed {
497 let svc = IndexQueryService::new(store, cid.clone(), db.clone());
498 if let Ok(hits) = svc.search_schema(
499 object_hint,
500 3,
501 Some(&filter),
502 SearchOptions {
503 use_semantic: self.use_semantic,
504 embedder: self.embedder.as_deref(),
505 },
506 ) && let Some(top) = hits.first()
507 {
508 let c = Candidate {
509 connection_id: cid.clone(),
510 database: db.clone(),
511 };
512 let entry = scored
513 .entry(key_of(&c))
514 .or_insert_with(|| (c.clone(), 0.0, Vec::new()));
515 entry.1 += top.score * 10.0;
516 entry.2.push(format!(
517 "schema search for \"{object_hint}\" found {} (score {:.2})",
518 top.ref_, top.score
519 ));
520 }
521 }
522 }
523
524 let mut ranked: Vec<(Candidate, f64, Vec<String>)> = scored.into_values().collect();
525 ranked.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
526
527 if ranked.is_empty() {
528 let candidates_json: Vec<Value> = connections
529 .iter()
530 .map(|c| {
531 json!({
532 "connectionId": c.id,
533 "connectionName": c.name,
534 "database": c.database.clone().unwrap_or_else(|| "postgres".into()),
535 })
536 })
537 .collect();
538 return Ok(ToolOutcome::ok_json(json!({
539 "ambiguous": true,
540 "message": format!(
541 "No connection/database matched \"{}\". Choose from the configured connections.",
542 hint.or(object_hint).unwrap_or_default()
543 ),
544 "candidates": candidates_json
545 })));
546 }
547
548 let winner = &ranked[0];
549 let is_tied = ranked
550 .get(1)
551 .is_some_and(|runner_up| runner_up.1 >= winner.1 * 0.85);
552
553 if is_tied {
554 let threshold = winner.1 * 0.85;
555 let tied: Vec<&(Candidate, f64, Vec<String>)> =
556 ranked.iter().filter(|r| r.1 >= threshold).take(5).collect();
557 let candidates_json: Vec<Value> = tied
558 .iter()
559 .filter_map(|(c, score, evidence)| {
560 connections
561 .iter()
562 .find(|x| x.id == c.connection_id)
563 .map(|conn| {
564 json!({
565 "connectionId": c.connection_id,
566 "connectionName": conn.name,
567 "database": c.database,
568 "score": score,
569 "evidence": evidence,
570 })
571 })
572 })
573 .collect();
574 return Ok(ToolOutcome::ok_json(json!({
575 "ambiguous": true,
576 "message": format!("{} equally-plausible candidates matched.", tied.len()),
577 "candidates": candidates_json
578 })));
579 }
580
581 let (winner_candidate, winner_score, winner_evidence) = winner;
582
583 if let Some(object_hint) = object_hint
584 && let Some(store) = self.index_store()
585 {
586 let filter = self.query_filter();
587 let svc = IndexQueryService::new(
588 store,
589 &winner_candidate.connection_id,
590 &winner_candidate.database,
591 );
592 if let Ok(hits) = svc.search_schema(
593 object_hint,
594 5,
595 Some(&filter),
596 SearchOptions {
597 use_semantic: self.use_semantic,
598 embedder: self.embedder.as_deref(),
599 },
600 ) && hits.len() >= 2
601 {
602 let top_score = hits[0].score;
603 let tied: Vec<&nexql_index::RankedHit> = hits
604 .iter()
605 .filter(|h| scores_equal(h.score, top_score))
606 .collect();
607 if tied.len() > 1 {
608 let candidates_json: Vec<Value> = tied
609 .iter()
610 .map(|h| {
611 json!({
612 "ref": h.ref_,
613 "score": h.score,
614 "kind": h.kind,
615 "connectionId": winner_candidate.connection_id,
616 "database": winner_candidate.database,
617 })
618 })
619 .collect();
620 return Ok(ToolOutcome::ok_json(json!({
621 "ambiguous": true,
622 "message": format!(
623 "{} objects matched \"{object_hint}\" with equal scores — choose explicitly.",
624 tied.len()
625 ),
626 "candidates": candidates_json,
627 })));
628 }
629 }
630 }
631
632 self.session
633 .switch(
634 &winner_candidate.connection_id,
635 Some(winner_candidate.database.clone()),
636 )
637 .await?;
638 let conn = connections
639 .iter()
640 .find(|x| x.id == winner_candidate.connection_id)
641 .ok_or_else(|| ToolError::Execution("resolved connection vanished".into()))?;
642
643 Ok(ToolOutcome::ok_json(json!({
644 "resolved": true,
645 "connectionId": winner_candidate.connection_id,
646 "connectionName": conn.name,
647 "database": winner_candidate.database,
648 "confidence": winner_score,
649 "evidence": winner_evidence,
650 })))
651 }
652
653 async fn discover_tools(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
654 let query = args
655 .get("query")
656 .and_then(|v| v.as_str())
657 .map(str::to_lowercase);
658 let category = args
659 .get("category")
660 .and_then(|v| v.as_str())
661 .map(str::to_lowercase);
662
663 let all_specs = active_tools();
665 let filtered: Vec<Value> = all_specs
666 .into_iter()
667 .filter(|spec| {
668 if spec.name == ToolName::DiscoverTools {
669 return false;
670 }
671 if let Some(ref cat) = category {
672 match cat.as_str() {
673 "query" if !ToolName::QUERY_PROFILE.contains(&spec.name) => return false,
674 "dba" if !ToolName::DBA_PROFILE.contains(&spec.name) => return false,
675 "write" if !ToolName::PHASE9.contains(&spec.name) => return false,
676 _ => {}
677 }
678 }
679 if let Some(ref q) = query {
680 let name_match = spec.name.as_str().contains(q.as_str());
681 let desc_match = spec.description.to_lowercase().contains(q.as_str());
682 if !name_match && !desc_match {
683 return false;
684 }
685 }
686 true
687 })
688 .map(|spec| {
689 json!({
690 "name": spec.name.as_str(),
691 "description": spec.description,
692 "input_schema": spec.input_schema,
693 })
694 })
695 .collect();
696
697 Ok(ToolOutcome::ok_json(json!({
698 "query": args.get("query"),
699 "category": args.get("category"),
700 "count": filtered.len(),
701 "tools": filtered,
702 })))
703 }
704
705 fn build_tuning_summary(plan_structured: &Option<Value>, suggestions: &Value) -> String {
706 let mut parts = Vec::new();
707 if let Some(structured) = plan_structured
708 && let Some(metrics) = structured.get("metrics")
709 {
710 if let Some(exec_time) = metrics.get("executionTime").and_then(|v| v.as_f64()) {
711 parts.push(format!("Query executed in {:.2}ms.", exec_time));
712 }
713 if let Some(seq_scans) = metrics.get("sequentialScans").and_then(|v| v.as_u64())
714 && seq_scans > 0
715 {
716 parts.push(format!("Found {seq_scans} sequential scan(s)."));
717 }
718 }
719
720 let candidate_count = suggestions
721 .get("high_seq_scan_tables")
722 .and_then(|v| v.as_array())
723 .map(|a| a.len())
724 .unwrap_or(0)
725 + suggestions
726 .get("unindexed_fk_columns")
727 .and_then(|v| v.as_array())
728 .map(|a| a.len())
729 .unwrap_or(0);
730
731 if candidate_count > 0 {
732 parts.push(format!(
733 "{candidate_count} index recommendation(s) identified."
734 ));
735 } else {
736 parts.push("No explicit index candidate recommendations generated.".into());
737 }
738
739 if parts.is_empty() {
740 "Auto-tune evaluation complete. Inspect execution plan and index recommendations."
741 .into()
742 } else {
743 parts.join(" ")
744 }
745 }
746
747 async fn auto_tune_query(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
748 let sql = args
749 .get("sql")
750 .and_then(|v| v.as_str())
751 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
752
753 let deep_plan = self
754 .deep_plan_analysis(&json!({ "sql": sql, "analyze": true }))
755 .await?;
756
757 let suggestions_res = self.suggest_indexes(&json!({ "sql": sql })).await;
758 let (suggestions_data, suggestions_error) = match suggestions_res {
759 Ok(outcome) => (outcome.structured.unwrap_or(json!([])), None),
760 Err(e) => (json!([]), Some(e.to_string())),
761 };
762
763 let summary_text = Self::build_tuning_summary(&deep_plan.structured, &suggestions_data);
764
765 let mut payload = json!({
766 "target_query": sql,
767 "deep_plan_analysis": deep_plan.structured,
768 "index_suggestions": suggestions_data,
769 "tuning_summary": summary_text,
770 });
771
772 if let Some(err) = suggestions_error {
773 payload["suggestions_error"] = json!(err);
774 }
775
776 Ok(ToolOutcome::ok_json(payload))
777 }
778
779 async fn check_ddl_safety_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
780 let ddl = args
781 .get("ddl")
782 .and_then(|v| v.as_str())
783 .ok_or_else(|| ToolError::InvalidArgs("ddl is required".into()))?;
784
785 let report = crate::dba_guard::analyze_ddl_safety(ddl);
786 Ok(ToolOutcome::ok_json(report))
787 }
788
789 async fn rebuild_index_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
790 let store = self
791 .index_store()
792 .ok_or_else(|| ToolError::Execution("Index store unavailable".into()))?;
793 let (connection_id, database) = self.session.active_context().await;
794 let depth_str = args
795 .get("depth")
796 .and_then(|v| v.as_str())
797 .unwrap_or("structure");
798 let depth: BuildDepth = match depth_str.to_lowercase().as_str() {
799 "profiles" | "full" => BuildDepth::Profiles,
800 _ => BuildDepth::Structure,
801 };
802
803 let req = BuildRequest {
804 connection_id: connection_id.clone(),
805 database: database.clone(),
806 scope: IndexScope {
807 included_schemas: vec![],
808 excluded_objects: vec![],
809 pii_excluded_columns: vec![],
810 },
811 depth,
812 build_mode: BuildMode::Guided,
813 environment: "development".into(),
814 embeddings: self.use_semantic,
815 };
816
817 let client = self.session.checkout().await?;
818 let db = PgCatalogDb::new(&client);
819 let manifest = build_index(store, &db, &req, None, None, self.embedder.as_deref())
820 .await
821 .map_err(|e| ToolError::Execution(format!("Index build failed: {e}")))?;
822 self.session.clear_index_stale(&connection_id, &database);
823
824 Ok(ToolOutcome::ok_json(json!({
825 "status": "completed",
826 "connection_id": connection_id,
827 "database": database,
828 "schema_fingerprint": manifest.schema_fingerprint,
829 "counts": manifest.counts,
830 "build_ms": manifest.stats.build_ms,
831 })))
832 }
833
834 async fn refresh_index_tool(&self, _args: &Value) -> Result<ToolOutcome, ToolError> {
835 let store = self
836 .index_store()
837 .ok_or_else(|| ToolError::Execution("Index store unavailable".into()))?;
838 let (connection_id, database) = self.session.active_context().await;
839 let base = store.base_dir(&connection_id, &database);
840 let manifest = store.read_manifest(&base)?.ok_or_else(|| {
841 ToolError::Execution(
842 "No existing index manifest to refresh — call 'rebuild_index'.".into(),
843 )
844 })?;
845
846 let req = BuildRequest {
847 connection_id: connection_id.clone(),
848 database: database.clone(),
849 scope: manifest.scope,
850 depth: manifest.build_depth,
851 build_mode: manifest.build_mode,
852 environment: manifest.environment,
853 embeddings: self.use_semantic,
854 };
855
856 let client = self.session.checkout().await?;
857 let db = PgCatalogDb::new(&client);
858 let new_manifest = build_index(store, &db, &req, None, None, self.embedder.as_deref())
859 .await
860 .map_err(|e| ToolError::Execution(format!("Index refresh failed: {e}")))?;
861 self.session.clear_index_stale(&connection_id, &database);
862
863 Ok(ToolOutcome::ok_json(json!({
864 "status": "refreshed",
865 "connection_id": connection_id,
866 "database": database,
867 "schema_fingerprint": new_manifest.schema_fingerprint,
868 "counts": new_manifest.counts,
869 "build_ms": new_manifest.stats.build_ms,
870 })))
871 }
872
873 async fn run_doctor_tool(&self) -> Result<ToolOutcome, ToolError> {
874 let (connection_id, database) = self.session.active_context().await;
875 let client = self.session.checkout().await?;
876
877 let version: String = client
878 .query_one("SELECT version()", &[])
879 .await
880 .map_err(|e| ToolError::Execution(e.to_string()))?
881 .get(0);
882
883 let is_super: String = client
884 .query_one("SELECT current_setting('is_superuser')", &[])
885 .await
886 .map_err(|e| ToolError::Execution(e.to_string()))?
887 .get(0);
888 let is_superuser = is_super.eq_ignore_ascii_case("on");
889
890 let ro: String = client
891 .query_one("SHOW default_transaction_read_only", &[])
892 .await
893 .map_err(|e| ToolError::Execution(e.to_string()))?
894 .get(0);
895
896 let timeout: String = client
897 .query_one("SHOW statement_timeout", &[])
898 .await
899 .map_err(|e| ToolError::Execution(e.to_string()))?
900 .get(0);
901
902 let pgs_present: bool = match client
903 .query_one(
904 "SELECT EXISTS (SELECT 1 FROM pg_extension WHERE extname = 'pg_stat_statements')",
905 &[],
906 )
907 .await
908 {
909 Ok(row) => row.get(0),
910 Err(_) => false,
911 };
912
913 let index_status = if let Some(store) = self.index_store() {
914 let base = store.base_dir(&connection_id, &database);
915 match store.read_manifest(&base) {
916 Ok(Some(m)) => json!({
917 "present": true,
918 "indexed_at": m.indexed_at,
919 "fingerprint": m.schema_fingerprint,
920 "tables": m.counts.tables,
921 }),
922 _ => json!({ "present": false }),
923 }
924 } else {
925 json!({ "present": false, "reason": "no_index_store" })
926 };
927
928 let recent_errors = read_recent_log_errors();
929
930 Ok(ToolOutcome::ok_json(json!({
931 "status": "ok",
932 "connection_id": connection_id,
933 "database": database,
934 "version": version.split(',').next().unwrap_or(&version),
935 "access_mode": format!("{:?}", self.session.access_mode()),
936 "superuser": is_superuser,
937 "read_only": ro,
938 "statement_timeout": timeout,
939 "pg_stat_statements": pgs_present,
940 "index": index_status,
941 "recent_errors": recent_errors,
942 })))
943 }
944
945 fn register_profile_in_session(
946 &self,
947 name: &str,
948 profile: &nexql_conn::ProfileConfig,
949 ) -> Result<(), ToolError> {
950 self.session.register_profile(
951 name,
952 profile,
953 self.session.access_mode(),
954 self.session.caps(),
955 )
956 }
957
958 async fn setup_connection_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
959 let profile_name = args
960 .get("name")
961 .and_then(|v| v.as_str())
962 .unwrap_or("default");
963
964 let candidates = crate::detect::ConnectionDetector::detect_all(None);
965
966 let url = args.get("url").and_then(|v| v.as_str());
967 let host = args.get("host").and_then(|v| v.as_str());
968 let port = args.get("port").and_then(|v| v.as_u64()).map(|n| n as u16);
969 let dbname = args.get("dbname").and_then(|v| v.as_str());
970 let user = args.get("user").and_then(|v| v.as_str());
971 let password = args.get("password").and_then(|v| v.as_str());
972 let sslmode = args.get("sslmode").and_then(|v| v.as_str());
973
974 let best_cand = candidates
975 .iter()
976 .find(|c| c.is_complete)
977 .or_else(|| candidates.first());
978
979 let res_host = host.or_else(|| best_cand.and_then(|c| c.host.as_deref()));
980 let res_port = port.or_else(|| best_cand.and_then(|c| c.port));
981 let res_dbname = dbname.or_else(|| best_cand.and_then(|c| c.dbname.as_deref()));
982 let res_user = user.or_else(|| best_cand.and_then(|c| c.user.as_deref()));
983 let res_password = password.or_else(|| best_cand.and_then(|c| c.password.as_deref()));
984 let res_url = url.or_else(|| best_cand.and_then(|c| c.url.as_deref()));
985 let res_sslmode = sslmode.or_else(|| best_cand.and_then(|c| c.sslmode.as_deref()));
986
987 if res_url.is_none() && (res_host.is_none() || res_dbname.is_none() || res_user.is_none()) {
988 let missing: Vec<&str> = vec![
989 if res_host.is_none() {
990 Some("host")
991 } else {
992 None
993 },
994 if res_dbname.is_none() {
995 Some("dbname")
996 } else {
997 None
998 },
999 if res_user.is_none() {
1000 Some("user")
1001 } else {
1002 None
1003 },
1004 ]
1005 .into_iter()
1006 .flatten()
1007 .collect();
1008
1009 return Ok(ToolOutcome::ok_json(json!({
1010 "status": "needs_input",
1011 "message": "Insufficient connection details. Please supply missing fields.",
1012 "detectedCandidates": candidates.iter().map(|c| c.redacted_json()).collect::<Vec<_>>(),
1013 "missingFields": missing
1014 })));
1015 }
1016
1017 let params = nexql_conn::ConnectionParams {
1018 url: res_url.map(String::from),
1019 host: res_host.map(String::from),
1020 port: res_port,
1021 dbname: res_dbname.map(String::from),
1022 user: res_user.map(String::from),
1023 password: res_password.map(String::from),
1024 sslmode: res_sslmode.map(String::from),
1025 ..Default::default()
1026 };
1027
1028 match nexql_conn::test_connection(¶ms).await {
1029 Ok(report) => {
1030 let p_config = nexql_conn::ProfileConfig {
1031 url: params.url.clone(),
1032 host: params.host.clone(),
1033 port: params.port,
1034 dbname: params.dbname.clone(),
1035 user: params.user.clone(),
1036 password: params.password.clone(),
1037 sslmode: params.sslmode.clone(),
1038 ..Default::default()
1039 };
1040
1041 let path = nexql_conn::ConfigFile::default_path().ok_or_else(|| {
1042 ToolError::Execution("Could not resolve config directory".into())
1043 })?;
1044 let mut cfg = nexql_conn::ConfigFile::load_path(&path).unwrap_or_default();
1045 cfg.upsert_profile(profile_name, p_config.clone());
1046 let backup = cfg
1047 .save(&path)
1048 .map_err(|e| ToolError::Execution(e.to_string()))?;
1049 self.register_profile_in_session(profile_name, &p_config)?;
1050
1051 Ok(ToolOutcome::ok_json(json!({
1052 "status": "configured",
1053 "profileName": profile_name,
1054 "serverVersion": report.server_version,
1055 "isSuperuser": report.is_superuser,
1056 "latencyMs": report.latency.as_millis(),
1057 "configPath": path.to_string_lossy().to_string(),
1058 "backup": backup.map(|b| b.to_string_lossy().to_string()),
1059 "sessionReloaded": true,
1060 })))
1061 }
1062 Err(e) => Ok(ToolOutcome::ok_json(json!({
1063 "status": "failed",
1064 "error": e.to_string(),
1065 "detectedCandidates": candidates.iter().map(|c| c.redacted_json()).collect::<Vec<_>>()
1066 }))),
1067 }
1068 }
1069
1070 async fn save_profile_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1071 let name = args
1072 .get("name")
1073 .and_then(|v| v.as_str())
1074 .ok_or_else(|| ToolError::InvalidArgs("name parameter is required".into()))?;
1075
1076 let p_config = nexql_conn::ProfileConfig {
1077 url: args.get("url").and_then(|v| v.as_str()).map(String::from),
1078 host: args.get("host").and_then(|v| v.as_str()).map(String::from),
1079 port: args.get("port").and_then(|v| v.as_u64()).map(|n| n as u16),
1080 dbname: args
1081 .get("dbname")
1082 .and_then(|v| v.as_str())
1083 .map(String::from),
1084 user: args.get("user").and_then(|v| v.as_str()).map(String::from),
1085 password: args
1086 .get("password")
1087 .and_then(|v| v.as_str())
1088 .map(String::from),
1089 sslmode: args
1090 .get("sslmode")
1091 .and_then(|v| v.as_str())
1092 .map(String::from),
1093 access_mode: args
1094 .get("access_mode")
1095 .and_then(|v| v.as_str())
1096 .map(String::from),
1097 max_rows: args
1098 .get("max_rows")
1099 .and_then(|v| v.as_u64())
1100 .map(|n| n as u32),
1101 ..Default::default()
1102 };
1103
1104 let path = nexql_conn::ConfigFile::default_path()
1105 .ok_or_else(|| ToolError::Execution("Could not resolve config directory".into()))?;
1106
1107 let mut cfg = nexql_conn::ConfigFile::load_path(&path).unwrap_or_default();
1108 cfg.upsert_profile(name, p_config.clone());
1109 let backup = cfg
1110 .save(&path)
1111 .map_err(|e| ToolError::Execution(e.to_string()))?;
1112 self.register_profile_in_session(name, &p_config)?;
1113
1114 Ok(ToolOutcome::ok_json(json!({
1115 "status": "saved",
1116 "profile": name,
1117 "configPath": path.to_string_lossy().to_string(),
1118 "backup": backup.map(|b| b.to_string_lossy().to_string()),
1119 "sessionReloaded": true,
1120 })))
1121 }
1122
1123 async fn test_profile_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1124 let name = args.get("name").and_then(|v| v.as_str());
1125
1126 let params = if let Some(pname) = name {
1127 let conn = self
1128 .session
1129 .connections()
1130 .into_iter()
1131 .find(|c| c.id == pname)
1132 .ok_or_else(|| ToolError::InvalidArgs(format!("Profile '{pname}' not found")))?;
1133 conn.params.clone()
1134 } else {
1135 nexql_conn::ConnectionParams {
1136 url: args.get("url").and_then(|v| v.as_str()).map(String::from),
1137 host: args.get("host").and_then(|v| v.as_str()).map(String::from),
1138 port: args.get("port").and_then(|v| v.as_u64()).map(|n| n as u16),
1139 dbname: args
1140 .get("dbname")
1141 .and_then(|v| v.as_str())
1142 .map(String::from),
1143 user: args.get("user").and_then(|v| v.as_str()).map(String::from),
1144 password: args
1145 .get("password")
1146 .and_then(|v| v.as_str())
1147 .map(String::from),
1148 sslmode: args
1149 .get("sslmode")
1150 .and_then(|v| v.as_str())
1151 .map(String::from),
1152 ..Default::default()
1153 }
1154 };
1155
1156 match nexql_conn::test_connection(¶ms).await {
1157 Ok(report) => Ok(ToolOutcome::ok_json(json!({
1158 "success": true,
1159 "serverVersion": report.server_version,
1160 "isSuperuser": report.is_superuser,
1161 "latencyMs": report.latency.as_millis()
1162 }))),
1163 Err(e) => Ok(ToolOutcome::ok_json(json!({
1164 "success": false,
1165 "error": e.to_string()
1166 }))),
1167 }
1168 }
1169
1170 async fn export_profile_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1171 let format = args
1172 .get("format")
1173 .and_then(|v| v.as_str())
1174 .unwrap_or("full");
1175 let path = nexql_conn::ConfigFile::default_path()
1176 .ok_or_else(|| ToolError::Execution("Could not resolve config directory".into()))?;
1177 let cfg = nexql_conn::ConfigFile::load_path(&path).unwrap_or_default();
1178
1179 if format == "project" {
1180 let proj = cfg.export_shareable();
1181 let toml_str =
1182 toml::to_string_pretty(&proj).map_err(|e| ToolError::Execution(e.to_string()))?;
1183 Ok(ToolOutcome::ok_json(json!({
1184 "format": "project",
1185 "filename": ".nexql/config.toml",
1186 "description": "Project policy overlay (no credentials). Use format=full for shareable connection profiles.",
1187 "content": toml_str,
1188 })))
1189 } else {
1190 let sanitized = cfg.export_full_sanitized();
1191 let toml_str = sanitized
1192 .to_toml_string()
1193 .map_err(|e| ToolError::Execution(e.to_string()))?;
1194 Ok(ToolOutcome::ok_json(json!({
1195 "format": "full",
1196 "description": "Full user config with passwords and secrets stripped — suitable for team sharing.",
1197 "profileCount": sanitized.profiles.len(),
1198 "content": toml_str,
1199 })))
1200 }
1201 }
1202
1203 async fn import_profile_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1204 let content = if let Some(c) = args.get("content").and_then(|v| v.as_str()) {
1205 c.to_string()
1206 } else if let Some(p) = args.get("path").and_then(|v| v.as_str()) {
1207 std::fs::read_to_string(p)
1208 .map_err(|e| ToolError::Execution(format!("failed to read file {p}: {e}")))?
1209 } else {
1210 return Err(ToolError::Execution(
1211 "either 'content' or 'path' must be specified".into(),
1212 ));
1213 };
1214
1215 let path = nexql_conn::ConfigFile::default_path()
1216 .ok_or_else(|| ToolError::Execution("Could not resolve config directory".into()))?;
1217 let mut cfg = nexql_conn::ConfigFile::load_path(&path).unwrap_or_default();
1218
1219 let imported: nexql_conn::ConfigFile = toml::from_str(&content)
1220 .map_err(|e| ToolError::Execution(format!("failed to parse TOML content: {e}")))?;
1221
1222 let mut count = 0;
1223 let mut imported_names: Vec<String> = Vec::new();
1224 for (name, prof) in imported.profiles {
1225 cfg.upsert_profile(name.clone(), prof.clone());
1226 self.register_profile_in_session(&name, &prof)?;
1227 imported_names.push(name);
1228 count += 1;
1229 }
1230 if imported.default_profile.is_some() {
1231 cfg.default_profile = imported.default_profile;
1232 }
1233
1234 let backup = cfg
1235 .save(&path)
1236 .map_err(|e| ToolError::Execution(e.to_string()))?;
1237
1238 Ok(ToolOutcome::ok_json(json!({
1239 "status": "imported",
1240 "imported_profiles": count,
1241 "profiles": imported_names,
1242 "configPath": path.to_string_lossy().to_string(),
1243 "backup": backup.map(|b| b.to_string_lossy().to_string()),
1244 "sessionReloaded": true,
1245 })))
1246 }
1247
1248 async fn index_service(&self) -> Result<(&IndexStore, String, String), ToolError> {
1249 let store = self
1250 .index_store()
1251 .ok_or_else(|| ToolError::Execution(NO_INDEX_HINT.into()))?;
1252 let (connection_id, database) = self.session.active_context().await;
1253 let base = store.base_dir(&connection_id, &database);
1254 if store.read_manifest(&base)?.is_none() {
1255 return Err(ToolError::Execution(format!(
1256 "No schema index for database \"{database}\" — call the 'rebuild_index' tool to build an index."
1257 )));
1258 }
1259 Ok((store, connection_id, database))
1260 }
1261
1262 fn ref_resolution_error(ref_: &str, resolution: &RefResolution) -> Option<ToolError> {
1266 match resolution {
1267 RefResolution::Resolved(_) => None,
1268 RefResolution::Ambiguous(candidates) => Some(ToolError::InvalidArgs(format!(
1269 "ambiguous relation \"{ref_}\": {}",
1270 candidates.join(", ")
1271 ))),
1272 RefResolution::Unknown {
1273 suggestion: Some(s),
1274 } => Some(ToolError::InvalidArgs(format!(
1275 "unknown relation \"{ref_}\" — did you mean \"{s}\"?"
1276 ))),
1277 RefResolution::Unknown { suggestion: None } => Some(ToolError::InvalidArgs(format!(
1278 "unknown relation \"{ref_}\" — call search_schema to find valid refs."
1279 ))),
1280 }
1281 }
1282
1283 fn resolve_indexed_ref_strict(
1287 svc: &IndexQueryService<'_>,
1288 ref_: &str,
1289 ) -> Result<String, ToolError> {
1290 let resolution = svc.resolve_ref(ref_)?;
1291 if let Some(err) = Self::ref_resolution_error(ref_, &resolution) {
1292 return Err(err);
1293 }
1294 match resolution {
1295 RefResolution::Resolved(r) => Ok(r),
1296 _ => unreachable!("ref_resolution_error covers every non-Resolved case"),
1297 }
1298 }
1299
1300 fn resolve_indexed_ref_soft(
1305 svc: &IndexQueryService<'_>,
1306 ref_: &str,
1307 ) -> Result<String, ToolError> {
1308 match svc.resolve_ref(ref_)? {
1309 RefResolution::Resolved(r) => Ok(r),
1310 RefResolution::Ambiguous(candidates) => Err(ToolError::InvalidArgs(format!(
1311 "ambiguous relation \"{ref_}\": {}",
1312 candidates.join(", ")
1313 ))),
1314 RefResolution::Unknown { .. } => Ok(ref_.to_owned()),
1315 }
1316 }
1317
1318 async fn resolve_ref_best_effort(&self, ref_: &str) -> Result<String, ToolError> {
1324 let Some(store) = self.index_store() else {
1325 return Ok(ref_.to_owned());
1326 };
1327 let (connection_id, database) = self.session.active_context().await;
1328 let base = store.base_dir(&connection_id, &database);
1329 if store.read_manifest(&base)?.is_none() {
1330 return Ok(ref_.to_owned());
1331 }
1332 let svc = IndexQueryService::new(store, &connection_id, &database);
1333 Self::resolve_indexed_ref_soft(&svc, ref_)
1334 }
1335
1336 async fn search_schema(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1337 let query = args
1338 .get("query")
1339 .and_then(|v| v.as_str())
1340 .unwrap_or("")
1341 .trim();
1342 if query.is_empty() {
1343 return Ok(ToolOutcome::ok_json(json!([])));
1344 }
1345 let (store, connection_id, database) = self.index_service().await?;
1346 let svc = IndexQueryService::new(store, &connection_id, &database);
1347 let filter = self.query_filter();
1348 let hits = svc.search_schema(
1349 query,
1350 SEARCH_SCHEMA_LIMIT,
1351 Some(&filter),
1352 SearchOptions {
1353 use_semantic: self.use_semantic,
1354 embedder: self.embedder.as_deref(),
1355 },
1356 )?;
1357 let rows: Vec<Value> = hits
1358 .into_iter()
1359 .map(|h| {
1360 json!({
1361 "ref": h.ref_,
1362 "score": h.score,
1363 "kind": h.kind,
1364 })
1365 })
1366 .collect();
1367 Ok(ToolOutcome::ok_json(json!(rows)))
1368 }
1369
1370 async fn describe_object(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1371 let ref_ = args
1372 .get("ref")
1373 .and_then(|v| v.as_str())
1374 .ok_or_else(|| ToolError::InvalidArgs("ref is required".into()))?;
1375 let (store, connection_id, database) = self.index_service().await?;
1376 let svc = IndexQueryService::new(store, &connection_id, &database);
1377 let resolved = Self::resolve_indexed_ref_soft(&svc, ref_)?;
1378 let filter = self.query_filter();
1379 let entry = svc.describe_object(&resolved, Some(&filter))?;
1380 let value = serde_json::to_value(entry).map_err(|e| ToolError::Execution(e.to_string()))?;
1381 Ok(ToolOutcome::ok_json(value))
1382 }
1383
1384 async fn get_join_path(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1385 let a = args
1386 .get("a")
1387 .and_then(|v| v.as_str())
1388 .ok_or_else(|| ToolError::InvalidArgs("a is required".into()))?;
1389 let b = args
1390 .get("b")
1391 .and_then(|v| v.as_str())
1392 .ok_or_else(|| ToolError::InvalidArgs("b is required".into()))?;
1393 let (store, connection_id, database) = self.index_service().await?;
1394 let svc = IndexQueryService::new(store, &connection_id, &database);
1395 let resolved_a = Self::resolve_indexed_ref_strict(&svc, a)?;
1400 let resolved_b = Self::resolve_indexed_ref_strict(&svc, b)?;
1401 let path = svc.get_join_path(&resolved_a, &resolved_b)?;
1402 let path_value =
1403 serde_json::to_value(path).map_err(|e| ToolError::Execution(e.to_string()))?;
1404 Ok(ToolOutcome::ok_json(json!({
1405 "path": path_value,
1406 "resolved_a": resolved_a,
1407 "resolved_b": resolved_b,
1408 })))
1409 }
1410
1411 async fn sample_values(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1412 let ref_ = args
1413 .get("ref")
1414 .and_then(|v| v.as_str())
1415 .ok_or_else(|| ToolError::InvalidArgs("ref is required".into()))?;
1416 let col = args
1417 .get("col")
1418 .and_then(|v| v.as_str())
1419 .ok_or_else(|| ToolError::InvalidArgs("col is required".into()))?;
1420 let (store, connection_id, database) = self.index_service().await?;
1421 let svc = IndexQueryService::new(store, &connection_id, &database);
1422 let resolved = Self::resolve_indexed_ref_soft(&svc, ref_)?;
1423 let ref_ = resolved.as_str();
1424 let filter = self.query_filter();
1425 let result = svc.sample_values(ref_, col, Some(&filter), None)?;
1426
1427 let mut values = result.values;
1428 let mut message = result.message;
1429
1430 if values.is_empty()
1431 && let Ok(client) = self.session.checkout().await
1432 {
1433 let parts: Vec<&str> = ref_.split('.').collect();
1434 let (schema, table) = match parts.as_slice() {
1435 [s, t] => (*s, *t),
1436 _ => ("public", ref_),
1437 };
1438 let safe_schema = schema.replace('"', "\"\"");
1439 let safe_table = table.replace('"', "\"\"");
1440 let safe_col = col.replace('"', "\"\"");
1441 let query = format!(
1442 "SELECT DISTINCT \"{safe_col}\"::text FROM \"{safe_schema}\".\"{safe_table}\" WHERE \"{safe_col}\" IS NOT NULL LIMIT 20"
1443 );
1444 if let Ok(rows) = client.query(&query, &[]).await {
1445 let sampled: Vec<String> = rows
1446 .iter()
1447 .filter_map(|r| r.get::<_, Option<String>>(0))
1448 .collect();
1449 if !sampled.is_empty() {
1450 values = sampled;
1451 message = None;
1452 }
1453 }
1454 }
1455
1456 let mut payload = json!({ "values": values });
1457 if let Some(msg) = message {
1458 payload["message"] = json!(msg);
1459 }
1460 Ok(ToolOutcome::ok_json(payload))
1461 }
1462
1463 fn list_connections(&self) -> ToolOutcome {
1464 let rows: Vec<Value> = self
1465 .session
1466 .connections()
1467 .iter()
1468 .map(|c| {
1469 json!({
1470 "id": c.id,
1471 "name": c.name,
1472 "host": c.host,
1473 "port": c.port,
1474 "database": c.database,
1475 })
1476 })
1477 .collect();
1478 ToolOutcome::ok_json(json!(rows))
1479 }
1480
1481 async fn list_databases(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1482 let connection_id = args
1483 .get("connectionId")
1484 .and_then(|v| v.as_str())
1485 .ok_or_else(|| ToolError::InvalidArgs("connectionId is required".into()))?;
1486 let conn = self
1487 .session
1488 .connections()
1489 .into_iter()
1490 .find(|c| c.id == connection_id)
1491 .ok_or_else(|| {
1492 ToolError::Execution(format!(
1493 "Connection not found for ID: {connection_id} — call list_connections"
1494 ))
1495 })?;
1496 let client = {
1498 if self.session.active_context().await.0 == connection_id {
1500 self.session.checkout().await?
1501 } else {
1502 let pool_opts = self.session.pool_opts();
1503 let pool = nexql_conn::create_pool(&conn.params, &pool_opts).await?;
1504 nexql_conn::checkout_guarded(&pool, &pool_opts).await?
1505 }
1506 };
1507 let rows = client
1508 .query(
1509 "SELECT datname FROM pg_database WHERE datistemplate = false ORDER BY datname",
1510 &[],
1511 )
1512 .await?;
1513 let names: Vec<String> = rows.iter().map(|r| r.get(0)).collect();
1514 Ok(ToolOutcome::ok_json(json!(names)))
1515 }
1516
1517 async fn list_schemas(&self) -> Result<ToolOutcome, ToolError> {
1518 let client = self.session.checkout().await?;
1519 let rows = client
1520 .query(
1521 r#"
1522 SELECT nspname AS schema_name
1523 FROM pg_namespace
1524 WHERE nspname NOT IN ('pg_catalog', 'information_schema', 'pg_toast')
1525 AND nspname NOT LIKE 'pg_%'
1526 ORDER BY nspname
1527 "#,
1528 &[],
1529 )
1530 .await?;
1531 let out: Vec<Value> = rows
1532 .iter()
1533 .filter(|r| {
1534 let name: String = r.get(0);
1535 self.session.filter().allows_schema(&name)
1536 })
1537 .map(|r| json!({ "schema_name": r.get::<_, String>(0) }))
1538 .collect();
1539 Ok(ToolOutcome::ok_json(json!(out)))
1540 }
1541
1542 async fn list_objects(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1543 let schema = args
1544 .get("schema")
1545 .and_then(|v| v.as_str())
1546 .unwrap_or("public");
1547 if !schema
1548 .chars()
1549 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
1550 {
1551 return Err(ToolError::InvalidArgs(
1552 "Invalid or missing schema name format".into(),
1553 ));
1554 }
1555 if !self.session.filter().allows_schema(schema) {
1556 return Ok(ToolOutcome::ok_json(json!([])));
1557 }
1558 let kind = args.get("kind").and_then(|v| v.as_str());
1559 let mut queries = Vec::new();
1560 let push_rel = |queries: &mut Vec<String>, relkinds: &[&str], label: &str| {
1561 let kinds = relkinds
1562 .iter()
1563 .map(|k| format!("'{k}'"))
1564 .collect::<Vec<_>>()
1565 .join(",");
1566 queries.push(format!(
1567 r#"
1568 SELECT n.nspname AS schema, c.relname AS name, '{label}' AS kind,
1569 d.description AS comment
1570 FROM pg_class c
1571 JOIN pg_namespace n ON n.oid = c.relnamespace
1572 LEFT JOIN pg_description d ON d.objoid = c.oid AND d.objsubid = 0
1573 WHERE n.nspname = $1 AND c.relkind IN ({kinds})
1574 "#
1575 ));
1576 };
1577 if kind.is_none() || kind == Some("table") {
1578 push_rel(&mut queries, &["r", "f", "p"], "table");
1579 }
1580 if kind.is_none() || kind == Some("view") {
1581 push_rel(&mut queries, &["v"], "view");
1582 }
1583 if kind.is_none() || kind == Some("matview") {
1584 push_rel(&mut queries, &["m"], "matview");
1585 }
1586 if queries.is_empty() {
1587 return Ok(ToolOutcome::ok_json(json!([])));
1588 }
1589 let sql = queries.join("\nUNION ALL\n") + "\nORDER BY kind, name";
1590 let client = self.session.checkout().await?;
1591 let rows = client.query(&sql, &[&schema]).await?;
1592 let out: Vec<Value> = rows
1593 .iter()
1594 .filter(|r| {
1595 let s: String = r.get("schema");
1596 let name: String = r.get("name");
1597 self.session.filter().allows_table(&s, &name)
1598 })
1599 .map(|r| {
1600 json!({
1601 "schema": r.get::<_, String>("schema"),
1602 "name": r.get::<_, String>("name"),
1603 "kind": r.get::<_, String>("kind"),
1604 "comment": r.get::<_, Option<String>>("comment"),
1605 })
1606 })
1607 .collect();
1608 Ok(ToolOutcome::ok_json(json!(out)))
1609 }
1610
1611 async fn get_current_context(&self) -> Result<ToolOutcome, ToolError> {
1612 let (connection_id, database) = self.session.active_context().await;
1613 let conn = self
1614 .session
1615 .connections()
1616 .into_iter()
1617 .find(|c| c.id == connection_id);
1618 Ok(ToolOutcome::ok_json(json!({
1619 "connectionId": connection_id,
1620 "connectionName": conn.as_ref().map(|c| c.name.clone()).unwrap_or_else(|| "Unknown".into()),
1621 "database": database,
1622 "host": conn.as_ref().and_then(|c| c.host.clone()),
1623 "port": conn.as_ref().and_then(|c| c.port),
1624 "access_mode": match self.session.access_mode() {
1625 nexql_policy::AccessMode::Read => "read",
1626 nexql_policy::AccessMode::Write => "write",
1627 nexql_policy::AccessMode::Admin => "admin",
1628 },
1629 })))
1630 }
1631
1632 async fn switch_connection(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1633 let connection_id = args
1634 .get("connectionId")
1635 .and_then(|v| v.as_str())
1636 .ok_or_else(|| ToolError::InvalidArgs("connectionId is required".into()))?;
1637 let database = args
1638 .get("database")
1639 .and_then(|v| v.as_str())
1640 .map(str::to_owned);
1641 self.session.switch(connection_id, database).await?;
1642 self.get_current_context().await
1643 }
1644
1645 async fn run_select(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1646 let sql = args
1647 .get("sql")
1648 .and_then(|v| v.as_str())
1649 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
1650 match validate_readonly_sql(sql)? {
1651 SqlDecision::Allow => {}
1652 SqlDecision::Reject => {
1653 return Err(ToolError::Execution(
1654 "Security Error: Only read-only SELECT, WITH, or EXPLAIN statements are permitted."
1655 .into(),
1656 ));
1657 }
1658 }
1659 enforce_read_table_policy(&self.session.filter(), sql)?;
1660 let trimmed = sql.trim().to_ascii_lowercase();
1661 if trimmed.starts_with("explain") {
1662 return self.run_select_internal(sql, None).await;
1663 }
1664 let max_rows = self.session.caps().max_rows;
1665 self.run_select_internal(sql, Some(max_rows)).await
1666 }
1667
1668 async fn explain_query(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1669 let sql = args
1670 .get("sql")
1671 .and_then(|v| v.as_str())
1672 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
1673 match validate_readonly_sql(sql)? {
1674 SqlDecision::Allow => {}
1675 SqlDecision::Reject => {
1676 return Err(ToolError::Execution(
1677 "Security Error: Only SELECT, WITH, or EXPLAIN statements can be analyzed."
1678 .into(),
1679 ));
1680 }
1681 }
1682 enforce_read_table_policy(&self.session.filter(), sql)?;
1683 let clean = if sql.trim().to_ascii_lowercase().starts_with("explain") {
1684 sql.to_string()
1685 } else {
1686 format!("EXPLAIN {sql}")
1687 };
1688 if validate_readonly_sql(&clean)? == SqlDecision::Reject {
1690 return Err(ToolError::Execution(
1691 "Security Error: EXPLAIN target is not read-only.".into(),
1692 ));
1693 }
1694 self.run_select_internal(&clean, None).await
1695 }
1696
1697 async fn get_ddl(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1698 let ref_ = args
1699 .get("ref")
1700 .and_then(|v| v.as_str())
1701 .ok_or_else(|| ToolError::InvalidArgs("ref is required".into()))?;
1702 let resolved = self.resolve_ref_best_effort(ref_).await?;
1703 let (schema, name) = parse_ref(&resolved).map_err(ToolError::InvalidArgs)?;
1704 let kind = args.get("kind").and_then(|v| v.as_str()).unwrap_or("table");
1705 let reg = sql::regclass_literal(&schema, &name);
1706 let client = self.session.checkout().await?;
1707
1708 match kind {
1709 "view" | "matview" => {
1710 let sql = format!("SELECT pg_get_viewdef({reg}, true) AS definition");
1711 let rows = client.query(&sql, &[]).await?;
1712 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1713 }
1714 "function" => {
1715 let sql = format!(
1716 r#"SELECT p.proname AS name, pg_get_functiondef(p.oid) AS definition
1717 FROM pg_proc p
1718 JOIN pg_namespace n ON n.oid = p.pronamespace
1719 WHERE n.nspname = '{schema}' AND p.proname = '{name}'"#
1720 );
1721 let rows = client.query(&sql, &[]).await?;
1722 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1723 }
1724 "index" => {
1725 let sql = format!("SELECT pg_get_indexdef({reg}) AS definition");
1726 let rows = client.query(&sql, &[]).await?;
1727 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1728 }
1729 "table" => {
1730 let columns = client
1731 .query(&sql::column_details(&schema, &name), &[])
1732 .await?;
1733 let constraints = client
1734 .query(
1735 &format!(
1736 r#"SELECT conname AS name, pg_get_constraintdef(oid) AS definition
1737 FROM pg_constraint WHERE conrelid = {reg} ORDER BY conname"#
1738 ),
1739 &[],
1740 )
1741 .await?;
1742 let indexes = client
1743 .query(
1744 &format!(
1745 r#"SELECT indexname AS name, indexdef AS definition
1746 FROM pg_indexes
1747 WHERE schemaname = '{schema}' AND tablename = '{name}'
1748 ORDER BY indexname"#
1749 ),
1750 &[],
1751 )
1752 .await?;
1753 Ok(ToolOutcome::ok_json(json!({
1754 "table": format!("{schema}.{name}"),
1755 "columns": rows_to_json(&columns),
1756 "constraints": rows_to_json(&constraints),
1757 "indexes": rows_to_json(&indexes),
1758 })))
1759 }
1760 other => Err(ToolError::InvalidArgs(format!(
1761 "Unsupported DDL kind \"{other}\". Use table, view, matview, function, or index."
1762 ))),
1763 }
1764 }
1765
1766 async fn table_stats(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1767 let ref_ = args
1768 .get("ref")
1769 .and_then(|v| v.as_str())
1770 .ok_or_else(|| ToolError::InvalidArgs("ref is required".into()))?;
1771 let resolved = self.resolve_ref_best_effort(ref_).await?;
1772 let (schema, name) = parse_ref(&resolved).map_err(ToolError::InvalidArgs)?;
1773 let client = self.session.checkout().await?;
1774 let stats = client.query(&sql::table_stats(&schema, &name), &[]).await?;
1775 let activity = client
1776 .query(&sql::table_activity(&schema, &name), &[])
1777 .await?;
1778 let columns = client
1779 .query(&sql::column_stats(&schema, &name), &[])
1780 .await?;
1781 let size = rows_to_json(&stats)
1782 .as_array()
1783 .and_then(|a| a.first())
1784 .cloned()
1785 .unwrap_or(Value::Null);
1786 let activity = rows_to_json(&activity)
1787 .as_array()
1788 .and_then(|a| a.first())
1789 .cloned()
1790 .unwrap_or(Value::Null);
1791 Ok(ToolOutcome::ok_json(json!({
1792 "size": size,
1793 "activity": activity,
1794 "columns": rows_to_json(&columns),
1795 })))
1796 }
1797
1798 async fn index_usage(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1799 let ref_ = args
1800 .get("ref")
1801 .and_then(|v| v.as_str())
1802 .ok_or_else(|| ToolError::InvalidArgs("ref is required".into()))?;
1803 let (schema, name) = parse_ref(ref_).map_err(ToolError::InvalidArgs)?;
1804 let client = self.session.checkout().await?;
1805 let rows = client.query(&sql::index_usage(&schema, &name), &[]).await?;
1806 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1807 }
1808
1809 async fn list_running_queries(&self) -> Result<ToolOutcome, ToolError> {
1810 let client = self.session.checkout().await?;
1811 let rows = client.query(sql::running_queries(), &[]).await?;
1812 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1813 }
1814
1815 async fn find_blocking_locks(&self) -> Result<ToolOutcome, ToolError> {
1816 let client = self.session.checkout().await?;
1817 let rows = client.query(sql::blocking_locks(), &[]).await?;
1818 let values = rows_to_json(&rows);
1819 if values.as_array().map(|a| a.is_empty()).unwrap_or(true) {
1820 return Ok(ToolOutcome::ok_json(json!({
1821 "message": "No blocking locks found.",
1822 "locks": [],
1823 })));
1824 }
1825 Ok(ToolOutcome::ok_json(values))
1826 }
1827
1828 async fn slow_queries(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1829 let limit = args
1830 .get("limit")
1831 .and_then(|v| v.as_u64())
1832 .map(|n| n as u32)
1833 .unwrap_or(SLOW_QUERIES_DEFAULT);
1834 let client = self.session.checkout().await?;
1835 match client.query(&sql::slow_queries(limit), &[]).await {
1836 Ok(rows) => Ok(ToolOutcome::ok_json(rows_to_json(&rows))),
1837 Err(e) => {
1838 if let Some(message) = sql::map_stat_statements_error(&e) {
1839 Ok(ToolOutcome::ok_json(json!({
1840 "error": message,
1841 "hint": message,
1842 })))
1843 } else {
1844 Err(ToolError::Postgres(e))
1845 }
1846 }
1847 }
1848 }
1849
1850 async fn db_health_check(&self) -> Result<ToolOutcome, ToolError> {
1851 let client = self.session.checkout().await?;
1852 let sections: &[(&str, &str)] = &[
1853 ("overview", sql::database_stats()),
1854 ("cache", sql::cache_hit_ratio()),
1855 ("dead_tuples", sql::database_maintenance_stats()),
1856 ("connection_states", sql::connection_states()),
1857 ("blocking_locks", sql::blocking_locks()),
1858 ];
1859 let mut report = serde_json::Map::new();
1860 for (key, q) in sections {
1861 match client.query(*q, &[]).await {
1862 Ok(rows) => {
1863 report.insert((*key).into(), rows_to_json(&rows));
1864 }
1865 Err(e) => {
1866 report.insert((*key).into(), json!({ "error": e.to_string() }));
1867 }
1868 }
1869 }
1870 let lock_count = report
1871 .get("blocking_locks")
1872 .and_then(|v| v.as_array())
1873 .map(|a| a.len() as u64);
1874 report.insert("blocking_lock_count".into(), json!(lock_count));
1875 Ok(ToolOutcome::ok_json(Value::Object(report)))
1876 }
1877
1878 async fn explain_analyze(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1879 let sql = args
1880 .get("sql")
1881 .and_then(|v| v.as_str())
1882 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
1883 require_select_or_with(&self.session.filter(), sql)?;
1884 let explain = build_explain_sql(sql, true);
1885 self.run_explain_in_transaction(&explain).await
1886 }
1887
1888 async fn analyze_query_plan(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1889 let sql = args
1890 .get("sql")
1891 .and_then(|v| v.as_str())
1892 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
1893 require_select_or_with(&self.session.filter(), sql)?;
1894 let analyze = args
1895 .get("analyze")
1896 .and_then(|v| v.as_bool())
1897 .unwrap_or(false);
1898 let explain = build_explain_sql(sql, analyze);
1899 let outcome = self.run_explain_in_transaction(&explain).await?;
1900 let rows = outcome.structured.unwrap_or(Value::Null);
1901 let row_array = rows
1902 .get("rows")
1903 .and_then(|v| v.as_array())
1904 .or_else(|| rows.as_array());
1905 let plan = row_array
1906 .and_then(|a| a.first())
1907 .and_then(|r| r.get("QUERY PLAN"))
1908 .cloned()
1909 .unwrap_or(Value::Null);
1910 let metrics = extract_plan_metrics(&plan).or_else(|| extract_plan_metrics(&rows));
1911 let recommendations = metrics
1912 .as_ref()
1913 .and_then(|m| m.get("recommendations"))
1914 .cloned()
1915 .unwrap_or_else(|| json!([]));
1916 Ok(ToolOutcome::ok_json(json!({
1917 "metrics": metrics,
1918 "recommendations": recommendations,
1919 "plan": plan,
1920 })))
1921 }
1922
1923 async fn run_explain_in_transaction(
1925 &self,
1926 explain_sql: &str,
1927 ) -> Result<ToolOutcome, ToolError> {
1928 let client = self.session.checkout().await?;
1929 client
1930 .batch_execute("SET statement_timeout = '30s'")
1931 .await?;
1932 client.batch_execute("BEGIN").await?;
1933 let result = async {
1934 client.batch_execute("SET TRANSACTION READ ONLY").await?;
1935 let rows = client.query(explain_sql, &[]).await?;
1936 Ok::<_, ToolError>(rows_to_json(&rows))
1937 }
1938 .await;
1939 let _ = client.batch_execute("ROLLBACK").await;
1941 match result {
1942 Ok(values) => Ok(ToolOutcome::ok_json(values)),
1943 Err(e) => Err(e),
1944 }
1945 }
1946
1947 async fn get_index_status(&self) -> Result<ToolOutcome, ToolError> {
1948 let (store, connection_id, database) = self.index_service().await?;
1949 let base = store.base_dir(&connection_id, &database);
1950 let Some(manifest) = store.read_manifest(&base)? else {
1951 return Err(ToolError::Execution(format!(
1952 "No schema index for database \"{database}\" — run `nexql-mcp index build`."
1953 )));
1954 };
1955
1956 let mut live_fingerprint: Option<String> = None;
1957 let mut drift: Option<bool> = None;
1958 if let Ok(client) = self.session.checkout().await {
1959 let db = PgCatalogDb::new(&client);
1960 if let Ok(fp) = db.schema_fingerprint().await {
1961 drift = Some(fp != manifest.schema_fingerprint);
1962 live_fingerprint = Some(fp);
1963 }
1964 }
1965
1966 Ok(ToolOutcome::ok_json(json!({
1967 "connectionId": manifest.connection_id,
1968 "database": manifest.database,
1969 "indexedAt": manifest.indexed_at,
1970 "fingerprint": manifest.schema_fingerprint,
1971 "liveFingerprint": live_fingerprint,
1972 "drift": drift,
1973 "pgVersion": manifest.pg_version,
1974 "counts": {
1975 "tables": manifest.counts.tables,
1976 "views": manifest.counts.views,
1977 "functions": manifest.counts.functions,
1978 "enums": manifest.counts.enums,
1979 },
1980 "buildMs": manifest.stats.build_ms,
1981 "warnings": manifest.stats.warnings,
1982 })))
1983 }
1984
1985 async fn list_extensions(&self) -> Result<ToolOutcome, ToolError> {
1986 let client = self.session.checkout().await?;
1987 let rows = client.query(sql::list_extensions(), &[]).await?;
1988 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1989 }
1990
1991 async fn server_settings(&self) -> Result<ToolOutcome, ToolError> {
1992 let client = self.session.checkout().await?;
1993 let rows = client.query(sql::server_settings(), &[]).await?;
1994 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1995 }
1996
1997 async fn suggest_indexes(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1998 let limit = args
1999 .get("limit")
2000 .and_then(|v| v.as_u64())
2001 .map(|n| n as u32)
2002 .unwrap_or(REPORT_LIMIT_DEFAULT);
2003 let client = self.session.checkout().await?;
2004 let mut query_errors = serde_json::Map::new();
2005
2006 let high_seq_json = match client.query(&sql::high_seq_scan_tables(limit), &[]).await {
2007 Ok(rows) => rows_to_json(&rows),
2008 Err(e) => {
2009 query_errors.insert(
2010 "high_seq_scan_tables".into(),
2011 json!(nexql_conn::format_postgres_error(&e)),
2012 );
2013 Value::Null
2014 }
2015 };
2016
2017 let unindexed_json = match client.query(&sql::unindexed_fk_columns(limit), &[]).await {
2018 Ok(rows) => rows_to_json(&rows),
2019 Err(e) => {
2020 query_errors.insert(
2021 "unindexed_fk_columns".into(),
2022 json!(nexql_conn::format_postgres_error(&e)),
2023 );
2024 Value::Null
2025 }
2026 };
2027
2028 let mut pg_stat_available = false;
2029 let mut slow_queries = Value::Null;
2030 let mut pg_stat_note: Option<String> = None;
2031 match client.query(&sql::slow_queries(limit.min(10)), &[]).await {
2032 Ok(rows) => {
2033 pg_stat_available = true;
2034 slow_queries = rows_to_json(&rows);
2035 }
2036 Err(e) => {
2037 if let Some(message) = sql::map_stat_statements_error(&e) {
2038 pg_stat_note = Some(message);
2039 } else {
2040 query_errors.insert(
2041 "slow_queries".into(),
2042 json!(nexql_conn::format_postgres_error(&e)),
2043 );
2044 }
2045 }
2046 }
2047
2048 let mut plan_heuristics = Value::Null;
2049 if let Some(sql_text) = args.get("sql").and_then(|v| v.as_str()) {
2050 require_select_or_with(&self.session.filter(), sql_text)?;
2051 let explain = build_explain_sql(sql_text, false);
2052 match self.run_explain_in_transaction(&explain).await {
2053 Ok(outcome) => {
2054 let rows = outcome.structured.unwrap_or(Value::Null);
2055 let plan = rows
2056 .as_array()
2057 .and_then(|a| a.first())
2058 .and_then(|r| r.get("QUERY PLAN"))
2059 .cloned()
2060 .unwrap_or(Value::Null);
2061 let metrics =
2062 extract_plan_metrics(&plan).or_else(|| extract_plan_metrics(&rows));
2063 plan_heuristics = json!({
2064 "metrics": metrics,
2065 "hint": "Use analyze_query_plan with analyze=true for actual timings before creating indexes.",
2066 });
2067 }
2068 Err(e) => {
2069 query_errors.insert("plan_heuristics".into(), json!(e.to_string()));
2070 }
2071 }
2072 }
2073
2074 let has_candidates = high_seq_json
2075 .as_array()
2076 .map(|a| !a.is_empty())
2077 .unwrap_or(false)
2078 || unindexed_json
2079 .as_array()
2080 .map(|a| !a.is_empty())
2081 .unwrap_or(false)
2082 || plan_heuristics != Value::Null;
2083
2084 let mut payload = if !has_candidates && !pg_stat_available {
2085 json!({
2086 "suggestions": [],
2087 "message": "No index suggestions yet. Either table stats show healthy index use, or there is not enough scan history. Enable pg_stat_statements and/or pass a sql argument for EXPLAIN plan heuristics.",
2088 "hint": pg_stat_note,
2089 })
2090 } else if !has_candidates {
2091 json!({
2092 "high_seq_scan_tables": high_seq_json,
2093 "unindexed_fk_columns": unindexed_json,
2094 "slow_queries": slow_queries,
2095 "plan_heuristics": plan_heuristics,
2096 "message": "No strong index candidates from sequential-scan or unindexed-FK heuristics. Review slow_queries / pass sql for plan-level advice.",
2097 "hint": "CREATE INDEX CONCURRENTLY after validating with EXPLAIN (ANALYZE, BUFFERS).",
2098 })
2099 } else {
2100 json!({
2101 "high_seq_scan_tables": high_seq_json,
2102 "unindexed_fk_columns": unindexed_json,
2103 "slow_queries": slow_queries,
2104 "plan_heuristics": plan_heuristics,
2105 "pg_stat_statements": pg_stat_available,
2106 "hint": pg_stat_note.unwrap_or_else(|| {
2107 "Validate candidates with analyze_query_plan / EXPLAIN before CREATE INDEX CONCURRENTLY.".into()
2108 }),
2109 })
2110 };
2111
2112 if !query_errors.is_empty() {
2113 payload["query_errors"] = Value::Object(query_errors);
2114 }
2115
2116 Ok(ToolOutcome::ok_json(payload))
2117 }
2118
2119 async fn find_unused_indexes(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2120 let limit = args
2121 .get("limit")
2122 .and_then(|v| v.as_u64())
2123 .map(|n| n as u32)
2124 .unwrap_or(REPORT_LIMIT_DEFAULT);
2125 let client = self.session.checkout().await?;
2126 let rows = client.query(&sql::find_unused_indexes(limit), &[]).await?;
2127 let indexes = rows_to_json(&rows);
2128 if indexes.as_array().map(|a| a.is_empty()).unwrap_or(true) {
2129 return Ok(ToolOutcome::ok_json(json!({
2130 "indexes": [],
2131 "message": "No unused non-constraint indexes found (idx_scan = 0). Note: pg_stat_reset / server restart clears scan counts — treat never-scanned indexes cautiously on fresh stats.",
2132 })));
2133 }
2134 Ok(ToolOutcome::ok_json(json!({
2135 "indexes": indexes,
2136 "hint": "Prefer DROP INDEX CONCURRENTLY after confirming the workload (and that stats are mature).",
2137 })))
2138 }
2139
2140 async fn bloat_report(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2141 let limit = args
2142 .get("limit")
2143 .and_then(|v| v.as_u64())
2144 .map(|n| n as u32)
2145 .unwrap_or(REPORT_LIMIT_DEFAULT);
2146 let client = self.session.checkout().await?;
2147 let rows = client.query(&sql::bloat_report(limit), &[]).await?;
2148 let tables = rows_to_json(&rows);
2149 if tables.as_array().map(|a| a.is_empty()).unwrap_or(true) {
2150 return Ok(ToolOutcome::ok_json(json!({
2151 "tables": [],
2152 "method": "dead_tuple_ratio",
2153 "message": "No tables with significant dead-tuple pressure (>1000 dead tuples). This is a simplified estimate from pg_stat_user_tables, not physical page bloat.",
2154 })));
2155 }
2156 Ok(ToolOutcome::ok_json(json!({
2157 "tables": tables,
2158 "method": "dead_tuple_ratio",
2159 "note": "Approximate bloat via n_dead_tup / (n_live_tup + n_dead_tup). Not a physical page-bloat estimate (pgstattuple / check_postgres). Consider VACUUM / VACUUM FULL only after confirming impact.",
2160 "hint": "VACUUM ANALYZE on high bloat_pct tables; investigate autovacuum settings if last_autovacuum is stale.",
2161 })))
2162 }
2163
2164 async fn find_missing_fks(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2165 let limit = args
2166 .get("limit")
2167 .and_then(|v| v.as_u64())
2168 .map(|n| n as u32)
2169 .unwrap_or(REPORT_LIMIT_DEFAULT);
2170 let capped = limit.clamp(1, sql::REPORT_LIMIT_MAX) as usize;
2171
2172 if let Ok((store, connection_id, database)) = self.index_service().await {
2174 let base = store.base_dir(&connection_id, &database);
2175 if let Ok(Some(manifest)) = store.read_manifest(&base)
2176 && let Ok(Some(graph)) = store.read_join_graph(&base, &manifest)
2177 {
2178 let candidates: Vec<Value> = graph
2179 .edges
2180 .into_iter()
2181 .filter(|e| e.inferred == Some(true) && e.disabled != Some(true))
2182 .take(capped)
2183 .map(|e| {
2184 let cols: Vec<Value> = e
2185 .cols
2186 .iter()
2187 .map(|(a, b)| json!({ "from": a, "to": b }))
2188 .collect();
2189 json!({
2190 "from_table": e.from,
2191 "to_table": e.to,
2192 "via": e.via,
2193 "columns": cols,
2194 "detection": "join_graph_inferred",
2195 })
2196 })
2197 .collect();
2198 if !candidates.is_empty() {
2199 return Ok(ToolOutcome::ok_json(json!({
2200 "candidates": candidates,
2201 "source": "join_graph",
2202 "hint": "These edges were inferred by naming convention and have no declared FK. Review before ALTER TABLE … ADD FOREIGN KEY.",
2203 })));
2204 }
2205 }
2206 }
2207
2208 let client = self.session.checkout().await?;
2209 let rows = client
2210 .query(&sql::find_missing_fks_catalog(limit), &[])
2211 .await?;
2212 let candidates = rows_to_json(&rows);
2213 if candidates.as_array().map(|a| a.is_empty()).unwrap_or(true) {
2214 return Ok(ToolOutcome::ok_json(json!({
2215 "candidates": [],
2216 "source": "catalog",
2217 "message": "No missing FK candidates found via join-graph inferred edges or *_id naming against single-column PKs.",
2218 })));
2219 }
2220 Ok(ToolOutcome::ok_json(json!({
2221 "candidates": candidates,
2222 "source": "catalog",
2223 "hint": "Naming-inferred only — verify referential integrity and nullability before adding constraints. Run `nexql-mcp index build` for join-graph inferred edges.",
2224 })))
2225 }
2226
2227 async fn list_roles(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2228 let client = self.session.checkout().await?;
2229 let role = args
2230 .get("role")
2231 .and_then(|v| v.as_str())
2232 .map(str::trim)
2233 .filter(|s| !s.is_empty());
2234
2235 let Some(role_name) = role else {
2236 let rows = client.query(sql::list_roles(), &[]).await?;
2237 return Ok(ToolOutcome::ok_json(rows_to_json(&rows)));
2238 };
2239
2240 let details = client.query(sql::role_details(), &[&role_name]).await?;
2241 if details.is_empty() {
2242 return Err(ToolError::Execution(format!(
2243 "Role \"{role_name}\" not found"
2244 )));
2245 }
2246 let member_of = client.query(sql::role_member_of(), &[&role_name]).await?;
2247 let has_members = client.query(sql::role_has_members(), &[&role_name]).await?;
2248 let privileges = client
2249 .query(sql::role_table_privileges(), &[&role_name])
2250 .await?;
2251
2252 Ok(ToolOutcome::ok_json(json!({
2253 "role": rows_to_json(&details).as_array().and_then(|a| a.first().cloned()).unwrap_or(Value::Null),
2254 "member_of": rows_to_json(&member_of),
2255 "has_members": rows_to_json(&has_members),
2256 "table_privileges": rows_to_json(&privileges),
2257 })))
2258 }
2259
2260 async fn export_query(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2261 let sql = args
2262 .get("sql")
2263 .and_then(|v| v.as_str())
2264 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
2265 require_select_or_with(&self.session.filter(), sql)?;
2266
2267 let format = args
2268 .get("format")
2269 .and_then(|v| v.as_str())
2270 .map(|s| {
2271 ExportFormat::parse(s).ok_or_else(|| {
2272 ToolError::InvalidArgs(format!(
2273 "Unsupported format \"{s}\". Use csv, json, or sqlinsert."
2274 ))
2275 })
2276 })
2277 .transpose()?
2278 .unwrap_or(ExportFormat::Csv);
2279
2280 let table_target = match args.get("table").and_then(|v| v.as_str()) {
2281 Some(t) if !t.trim().is_empty() => Some(parse_ref(t).map_err(ToolError::InvalidArgs)?),
2282 _ => None,
2283 };
2284
2285 if format == ExportFormat::SqlInsert && table_target.is_none() {
2286 return Err(ToolError::InvalidArgs(
2287 "table (schema.name) is required when format=sqlinsert".into(),
2288 ));
2289 }
2290
2291 let max_rows = self.session.caps().max_rows;
2292 let outcome = self.run_select_internal(sql, Some(max_rows)).await?;
2293 if outcome.is_error {
2294 return Ok(outcome);
2295 }
2296
2297 let structured = outcome.structured.unwrap_or(Value::Null);
2298 let rows_val = structured
2299 .get("rows")
2300 .cloned()
2301 .or_else(|| structured.get("data").and_then(|d| d.get("rows").cloned()))
2302 .unwrap_or(Value::Array(vec![]));
2303 let rows = rows_val.as_array().cloned().unwrap_or_default();
2304 let columns = columns_from_rows(&rows);
2305 let truncated = structured
2306 .get("truncated")
2307 .and_then(|v| v.as_bool())
2308 .unwrap_or(false);
2309
2310 let payload = match format {
2311 ExportFormat::Json => json!({
2312 "format": format.as_str(),
2313 "rowCount": rows.len(),
2314 "truncated": truncated,
2315 "columns": columns,
2316 "rows": rows,
2317 }),
2318 ExportFormat::Csv => {
2319 let content = rows_to_csv(&rows, &columns);
2320 let caps = self.session.caps();
2321 let (char_trunc, content) = caps.truncate_chars(&content);
2322 json!({
2323 "format": format.as_str(),
2324 "rowCount": rows.len(),
2325 "truncated": truncated || char_trunc,
2326 "columns": columns,
2327 "content": content,
2328 })
2329 }
2330 ExportFormat::SqlInsert => {
2331 let (schema, table) = table_target.expect("checked above");
2332 let content = rows_to_sql_insert(&rows, &columns, &schema, &table);
2333 let caps = self.session.caps();
2334 let (char_trunc, content) = caps.truncate_chars(&content);
2335 json!({
2336 "format": format.as_str(),
2337 "rowCount": rows.len(),
2338 "truncated": truncated || char_trunc,
2339 "table": format!("{schema}.{table}"),
2340 "columns": columns,
2341 "content": content,
2342 })
2343 }
2344 };
2345
2346 Ok(ToolOutcome::ok_json(payload))
2347 }
2348
2349 async fn db_dashboard(&self) -> Result<ToolOutcome, ToolError> {
2350 let client = self.session.checkout().await?;
2351 let sections: &[(&str, &str)] = &[
2352 ("db_info", sql::dashboard_db_info()),
2353 ("connection_states", sql::connection_states()),
2354 ("top_tables", sql::dashboard_top_tables()),
2355 ("object_counts", sql::dashboard_object_counts()),
2356 ("active_queries", sql::dashboard_active_queries()),
2357 ("blocking_locks", sql::blocking_locks()),
2358 ("max_connections", sql::dashboard_max_connections()),
2359 ("extension_count", sql::dashboard_extension_count()),
2360 ("cache", sql::cache_hit_ratio()),
2361 ];
2362 let mut report = serde_json::Map::new();
2363 for (key, q) in sections {
2364 match client.query(*q, &[]).await {
2365 Ok(rows) => {
2366 report.insert((*key).into(), rows_to_json(&rows));
2367 }
2368 Err(e) => {
2369 report.insert((*key).into(), json!({ "error": e.to_string() }));
2370 }
2371 }
2372 }
2373
2374 for key in ["db_info", "object_counts", "extension_count", "cache"] {
2376 if let Some(Value::Array(arr)) = report.get(key).cloned()
2377 && arr.len() == 1
2378 {
2379 report.insert(key.into(), arr.into_iter().next().unwrap());
2380 }
2381 }
2382 if let Some(Value::Array(arr)) = report.get("max_connections").cloned()
2383 && let Some(row) = arr.first()
2384 {
2385 report.insert(
2386 "max_connections".into(),
2387 row.get("max_connections")
2388 .cloned()
2389 .unwrap_or_else(|| row.clone()),
2390 );
2391 }
2392
2393 Ok(ToolOutcome::ok_json(Value::Object(report)))
2394 }
2395
2396 async fn deep_plan_analysis(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2397 let sql = args
2398 .get("sql")
2399 .and_then(|v| v.as_str())
2400 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
2401 require_select_or_with(&self.session.filter(), sql)?;
2402 let analyze = args
2403 .get("analyze")
2404 .and_then(|v| v.as_bool())
2405 .unwrap_or(true);
2406 let explain = build_explain_sql(sql, analyze);
2407 let outcome = self.run_explain_in_transaction(&explain).await?;
2408 let rows = outcome.structured.unwrap_or(Value::Null);
2409 let row_array = rows
2410 .get("rows")
2411 .and_then(|v| v.as_array())
2412 .or_else(|| rows.as_array());
2413 let plan = row_array
2414 .and_then(|a| a.first())
2415 .and_then(|r| r.get("QUERY PLAN"))
2416 .cloned()
2417 .unwrap_or(Value::Null);
2418 let deep = analyze_deep_plan(&plan, sql)
2419 .or_else(|| analyze_deep_plan(&rows, sql))
2420 .ok_or_else(|| {
2421 ToolError::Execution("Could not parse EXPLAIN JSON plan for deep analysis".into())
2422 })?;
2423 let metrics = extract_plan_metrics(&plan).or_else(|| extract_plan_metrics(&rows));
2424 Ok(ToolOutcome::ok_json(json!({
2425 "deep": deep,
2426 "metrics": metrics,
2427 "plan": plan,
2428 "analyzed": analyze,
2429 })))
2430 }
2431
2432 async fn schema_diff(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2433 let source_schema = args
2434 .get("sourceSchema")
2435 .and_then(|v| v.as_str())
2436 .ok_or_else(|| ToolError::InvalidArgs("sourceSchema is required".into()))?;
2437 let target_schema = args
2438 .get("targetSchema")
2439 .and_then(|v| v.as_str())
2440 .ok_or_else(|| ToolError::InvalidArgs("targetSchema is required".into()))?;
2441 crate::schema_diff::require_safe_schema(source_schema)?;
2442 crate::schema_diff::require_safe_schema(target_schema)?;
2443
2444 let client = self.session.checkout().await?;
2445 let source = crate::schema_diff::load_schema_snapshot(&client, source_schema).await?;
2446 let target = crate::schema_diff::load_schema_snapshot(&client, target_schema).await?;
2447 let diffs = crate::schema_diff::compute_schema_diff(&source, &target);
2448 let changed = diffs
2449 .iter()
2450 .filter(|d| d.status != crate::schema_diff::DiffStatus::Unchanged)
2451 .count();
2452 Ok(ToolOutcome::ok_json(json!({
2453 "sourceSchema": source_schema,
2454 "targetSchema": target_schema,
2455 "tableCount": diffs.len(),
2456 "changedCount": changed,
2457 "diffs": crate::schema_diff::diffs_to_json(&diffs),
2458 })))
2459 }
2460
2461 async fn generate_migration(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2462 let source_schema = args
2463 .get("sourceSchema")
2464 .and_then(|v| v.as_str())
2465 .ok_or_else(|| ToolError::InvalidArgs("sourceSchema is required".into()))?;
2466 let target_schema = args
2467 .get("targetSchema")
2468 .and_then(|v| v.as_str())
2469 .ok_or_else(|| ToolError::InvalidArgs("targetSchema is required".into()))?;
2470 crate::schema_diff::require_safe_schema(source_schema)?;
2471 crate::schema_diff::require_safe_schema(target_schema)?;
2472
2473 let client = self.session.checkout().await?;
2474 let source = crate::schema_diff::load_schema_snapshot(&client, source_schema).await?;
2475 let target = crate::schema_diff::load_schema_snapshot(&client, target_schema).await?;
2476 let diffs = crate::schema_diff::compute_schema_diff(&source, &target);
2477 let statements =
2478 crate::schema_diff::build_migration_statements(source_schema, target_schema, &diffs);
2479 let sql = if statements.is_empty() {
2480 format!("-- No differences between {source_schema} and {target_schema}")
2481 } else {
2482 statements.join("\n\n")
2483 };
2484 Ok(ToolOutcome::ok_json(json!({
2485 "sourceSchema": source_schema,
2486 "targetSchema": target_schema,
2487 "statementCount": statements.len(),
2488 "sql": sql,
2489 "hint": "Read-only: review and run via execute_sql / apply_ddl only with --access-mode write|admin. Destructive drops are commented out.",
2490 })))
2491 }
2492
2493 async fn run_select_internal(
2494 &self,
2495 sql: &str,
2496 max_rows: Option<u32>,
2497 ) -> Result<ToolOutcome, ToolError> {
2498 let client = self.session.checkout().await?;
2499 let Some(max_rows) = max_rows else {
2500 let rows = client.query(sql, &[]).await?;
2501 let values = rows_to_json(&rows);
2502 let payload = self.apply_pii_redaction(sql, ensure_structured_object(values));
2503 let text = serde_json::to_string_pretty(&payload)
2504 .map_err(|e| ToolError::Execution(e.to_string()))?;
2505 let caps = self.session.caps();
2506 let (trunc, text) = caps.truncate_chars(&text);
2507 let structured = if trunc {
2508 json!({ "truncated_chars": true, "data": payload })
2509 } else {
2510 payload
2511 };
2512 return Ok(ToolOutcome {
2513 text: text.to_string(),
2514 structured: Some(structured),
2515 is_error: false,
2516 });
2517 };
2518
2519 let cleaned = sql.trim().trim_end_matches(';').trim();
2520 let wrapped = format!(
2521 "SELECT * FROM ({cleaned}) AS nexql_limited LIMIT {}",
2522 max_rows + 1
2523 );
2524 let rows = client.query(&wrapped, &[]).await.map_err(|e| {
2525 ToolError::Execution(format!(
2526 "Failed to execute row-limited query (refusing unbounded fallback): {}",
2527 nexql_conn::format_postgres_error(&e)
2528 ))
2529 })?;
2530 let truncated = rows.len() as u32 > max_rows;
2531 let keep = if truncated {
2532 &rows[..max_rows as usize]
2533 } else {
2534 &rows[..]
2535 };
2536 let values = rows_to_json(keep);
2537 let mut payload = self.apply_pii_redaction(sql, ensure_structured_object(values));
2539 if truncated && let Some(obj) = payload.as_object_mut() {
2540 obj.insert("truncated".into(), json!(true));
2541 obj.insert("maxRows".into(), json!(max_rows));
2542 }
2543 let text = serde_json::to_string_pretty(&payload)
2544 .map_err(|e| ToolError::Execution(e.to_string()))?;
2545 let caps = self.session.caps();
2546 let (char_trunc, text) = caps.truncate_chars(&text);
2547 let structured = if char_trunc {
2548 json!({ "truncated_chars": true, "data": payload })
2549 } else {
2550 payload
2551 };
2552 Ok(ToolOutcome {
2553 text: text.to_string(),
2554 structured: Some(structured),
2555 is_error: false,
2556 })
2557 }
2558
2559 fn apply_pii_redaction(&self, sql: &str, mut payload: Value) -> Value {
2560 let filter = self.session.filter();
2561 if filter.pii_columns.is_empty() {
2562 return payload;
2563 }
2564 let Ok(tables) = select_table_refs(sql) else {
2565 return payload;
2566 };
2567 let (redacted, cols) = redact_pii_in_payload(payload, &filter.pii_columns, &tables);
2568 payload = redacted;
2569 if !cols.is_empty()
2570 && let Some(obj) = payload.as_object_mut()
2571 {
2572 obj.insert("piiRedactedColumns".into(), json!(cols));
2573 }
2574 payload
2575 }
2576}
2577
2578fn normalize_for_match(s: &str) -> String {
2580 let mut out = String::new();
2581 let mut last_was_sep = true;
2582 for ch in s.to_lowercase().chars() {
2583 if ch.is_ascii_alphanumeric() {
2584 out.push(ch);
2585 last_was_sep = false;
2586 } else if !last_was_sep {
2587 out.push(' ');
2588 last_was_sep = true;
2589 }
2590 }
2591 out.trim().to_string()
2592}
2593
2594fn fuzzy_score(hint: &str, candidate: &str) -> f64 {
2596 let h = normalize_for_match(hint);
2597 let c = normalize_for_match(candidate);
2598 if h.is_empty() || c.is_empty() {
2599 return 0.0;
2600 }
2601 if h == c {
2602 return 100.0;
2603 }
2604 if c.contains(&h) || h.contains(&c) {
2605 return 75.0;
2606 }
2607 let h_tokens: std::collections::HashSet<&str> =
2608 h.split(' ').filter(|s| !s.is_empty()).collect();
2609 let c_tokens: std::collections::HashSet<&str> =
2610 c.split(' ').filter(|s| !s.is_empty()).collect();
2611 let overlap = h_tokens.intersection(&c_tokens).count();
2612 if overlap == 0 {
2613 return 0.0;
2614 }
2615 (overlap as f64 / h_tokens.len().max(c_tokens.len()) as f64) * 60.0
2616}
2617
2618fn policy_to_query_filter(filter: &PolicyFilter) -> QueryPolicyFilter {
2619 QueryPolicyFilter {
2620 allow_schemas: filter.allow_schemas.clone(),
2621 deny_schemas: filter.deny_schemas.clone(),
2622 deny_tables: filter.deny_tables.clone(),
2623 pii_columns: filter.pii_columns.clone(),
2624 }
2625}
2626
2627fn require_select_or_with(filter: &PolicyFilter, sql: &str) -> Result<(), ToolError> {
2628 match validate_readonly_sql(sql)? {
2629 SqlDecision::Allow => {}
2630 SqlDecision::Reject => {
2631 return Err(ToolError::Execution(
2632 "Security Error: Only SELECT or WITH statements can be analyzed.".into(),
2633 ));
2634 }
2635 }
2636 enforce_read_table_policy(filter, sql)?;
2637 let trimmed = sql.trim().to_ascii_lowercase();
2638 if !(trimmed.starts_with("select") || trimmed.starts_with("with")) {
2639 return Err(ToolError::Execution(
2640 "Security Error: Only SELECT or WITH statements can be analyzed.".into(),
2641 ));
2642 }
2643 Ok(())
2644}
2645
2646fn rows_to_json(rows: &[tokio_postgres::Row]) -> Value {
2647 rows_to_json_array(rows)
2648}
2649
2650fn scores_equal(a: f64, b: f64) -> bool {
2651 (a - b).abs() <= f64::EPSILON * a.abs().max(b.abs()).max(1.0)
2652}
2653
2654fn read_recent_log_errors() -> Vec<String> {
2655 let path = std::env::var("NEXQL_MCP_LOG")
2656 .map(std::path::PathBuf::from)
2657 .ok()
2658 .or_else(|| {
2659 std::env::var_os("HOME").map(|h| {
2660 std::path::PathBuf::from(h)
2661 .join(".config")
2662 .join("nexql-mcp")
2663 .join("logs")
2664 .join("nexql-mcp.log")
2665 })
2666 });
2667
2668 let Some(log_path) = path else {
2669 return Vec::new();
2670 };
2671
2672 let Ok(content) = std::fs::read_to_string(&log_path) else {
2673 return Vec::new();
2674 };
2675
2676 content
2677 .lines()
2678 .rev()
2679 .take(50)
2680 .filter(|line| {
2681 line.contains("ERROR")
2682 || line.contains("WARN")
2683 || line.contains("failed")
2684 || line.contains("Error")
2685 })
2686 .map(String::from)
2687 .collect()
2688}
2689
2690#[cfg(test)]
2691mod tests {
2692 use super::*;
2693 use crate::plan::build_explain_sql;
2694 use nexql_policy::PolicyFilter;
2695 use serde_json::json;
2696
2697 use crate::session::{ConnectionInfo, ConnectionPolicy, ToolSession};
2698 use nexql_policy::{AccessMode, PolicyCaps};
2699
2700 fn test_conn() -> ConnectionInfo {
2701 ConnectionInfo {
2702 id: "conn-1".into(),
2703 name: "conn-1".into(),
2704 host: Some("127.0.0.1".into()),
2705 port: Some(5432),
2706 database: Some("appdb".into()),
2707 params: Default::default(),
2708 policy: ConnectionPolicy {
2709 access_mode: AccessMode::Read,
2710 caps: PolicyCaps::default(),
2711 filter: PolicyFilter::default(),
2712 environment: None,
2713 },
2714 }
2715 }
2716
2717 #[test]
2718 fn scores_equal_treats_near_duplicates_as_tied() {
2719 let s = 3.295836866004329_f64;
2720 assert!(super::scores_equal(s, s));
2721 assert!(super::scores_equal(s, s + f64::EPSILON));
2722 }
2723
2724 #[test]
2725 fn policy_maps_one_to_one() {
2726 let f = PolicyFilter {
2727 allow_schemas: vec!["public".into()],
2728 deny_schemas: vec!["pgboss".into()],
2729 deny_tables: vec!["auth.*".into()],
2730 pii_columns: vec!["public.users.ssn".into()],
2731 };
2732 let q = policy_to_query_filter(&f);
2733 assert_eq!(q.allow_schemas, f.allow_schemas);
2734 assert_eq!(q.deny_schemas, f.deny_schemas);
2735 assert_eq!(q.deny_tables, f.deny_tables);
2736 assert_eq!(q.pii_columns, f.pii_columns);
2737 }
2738
2739 #[test]
2740 fn ok_json_wraps_arrays_for_cursor_structured_content() {
2741 let out = ToolOutcome::ok_json(json!([{ "id": 1 }, { "id": 2 }]));
2742 assert!(!out.is_error);
2743 let s = out.structured.as_ref().unwrap();
2744 assert!(s.is_object(), "structuredContent must be object, got {s}");
2745 assert_eq!(s["rows"].as_array().unwrap().len(), 2);
2746 assert!(out.text.contains("\"rows\""));
2747 }
2748
2749 #[test]
2750 fn ok_json_leaves_objects_unchanged() {
2751 let out = ToolOutcome::ok_json(json!({ "kind": "table", "name": "orders" }));
2752 let s = out.structured.as_ref().unwrap();
2753 assert_eq!(s["kind"], "table");
2754 assert!(s.get("rows").is_none());
2755 }
2756
2757 #[test]
2758 fn router_specs_include_phase4_and_phase9() {
2759 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2760 let router = ToolRouter::with_index_store(session, None);
2761 assert_eq!(router.specs().len(), ToolName::ACTIVE.len());
2762 let names: Vec<_> = router.specs().iter().map(|s| s.name.as_str()).collect();
2763 assert!(names.contains(&"search_schema"));
2764 assert!(names.contains(&"get_ddl"));
2765 assert!(names.contains(&"explain_analyze"));
2766 assert!(names.contains(&"get_index_status"));
2767 assert!(names.contains(&"list_extensions"));
2768 assert!(names.contains(&"server_settings"));
2769 assert!(names.contains(&"suggest_indexes"));
2770 assert!(names.contains(&"find_unused_indexes"));
2771 assert!(names.contains(&"bloat_report"));
2772 assert!(names.contains(&"find_missing_fks"));
2773 assert!(names.contains(&"export_query"));
2774 assert!(names.contains(&"list_roles"));
2775 assert!(names.contains(&"db_dashboard"));
2776 assert!(names.contains(&"deep_plan_analysis"));
2777 assert!(names.contains(&"execute_sql"));
2778 assert!(names.contains(&"edit_row"));
2779 assert!(names.contains(&"import_data"));
2780 assert!(names.contains(&"apply_ddl"));
2781 assert!(names.contains(&"create_index_concurrently"));
2782 assert!(names.contains(&"run_maintenance"));
2783 assert!(names.contains(&"terminate_query"));
2784 }
2785
2786 #[tokio::test]
2787 async fn write_tools_refuse_read_mode() {
2788 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2789 let router = ToolRouter::with_index_store(session, None);
2790 for tool in [
2791 "execute_sql",
2792 "edit_row",
2793 "import_data",
2794 "apply_ddl",
2795 "create_index_concurrently",
2796 "run_maintenance",
2797 "terminate_query",
2798 ] {
2799 let out = router
2800 .call(tool, json!({ "sql": "SELECT 1", "table": "public.t", "rows": [], "action": "insert", "values": {}, "pid": 1 }))
2801 .await;
2802 assert!(out.is_error, "{tool}: {}", out.text);
2803 assert!(
2804 out.text.contains("write") || out.text.contains("admin"),
2805 "{tool}: {}",
2806 out.text
2807 );
2808 }
2809 }
2810
2811 #[tokio::test]
2812 async fn table_stats_rejects_injection_ref() {
2813 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2814 let router = ToolRouter::with_index_store(session, None);
2815 let out = router
2816 .call("table_stats", json!({ "ref": "public.users; DROP" }))
2817 .await;
2818 assert!(out.is_error, "{}", out.text);
2819 assert!(
2820 out.text.contains("Invalid object reference") || out.text.contains("invalid arguments"),
2821 "expected ref validation error, got: {}",
2822 out.text
2823 );
2824 }
2825
2826 #[test]
2827 fn explain_transaction_path_builds_readonly_sequence() {
2828 let explain = build_explain_sql("SELECT 1", true);
2830 assert!(explain.starts_with("EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON)"));
2831 assert!(!explain.to_ascii_lowercase().contains("commit"));
2832 let steps = ["BEGIN", "SET TRANSACTION READ ONLY", &explain, "ROLLBACK"];
2833 assert_eq!(steps.len(), 4);
2834 assert_eq!(steps[0], "BEGIN");
2835 assert_eq!(steps[1], "SET TRANSACTION READ ONLY");
2836 assert_eq!(steps[3], "ROLLBACK");
2837 }
2838
2839 #[tokio::test]
2840 async fn missing_index_returns_actionable_error() {
2841 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2842 let router = ToolRouter::with_index_store(session, None);
2843 let out = router
2844 .call("search_schema", json!({ "query": "users" }))
2845 .await;
2846 assert!(out.is_error, "{}", out.text);
2847 assert!(
2848 out.text.contains("rebuild_index"),
2849 "expected actionable hint, got: {}",
2850 out.text
2851 );
2852 }
2853
2854 fn write_join_path_fixture(store: &IndexStore) {
2858 use nexql_index::{
2859 BuildDepth, BuildMode, ColumnEntry, DbObjectKind, IndexCounts, IndexDerived,
2860 IndexManifest, IndexScope, IndexStats, JOIN_GRAPH_FILE, JoinEdge, JoinGraph,
2861 ObjectEntry, ObjectShard, TOKENS_FILE,
2862 };
2863 use std::collections::HashMap;
2864
2865 let base = store.base_dir("conn-1", "appdb");
2866 let manifest = IndexManifest {
2867 format_version: 1,
2868 connection_id: "conn-1".into(),
2869 database: "appdb".into(),
2870 indexed_at: "2026-08-08T00:00:00.000Z".into(),
2871 build_mode: BuildMode::Auto,
2872 build_depth: BuildDepth::Structure,
2873 schema_fingerprint: "fp".into(),
2874 pg_version: "18.4".into(),
2875 environment: "development".into(),
2876 scope: IndexScope {
2877 included_schemas: vec!["public".into()],
2878 excluded_objects: vec![],
2879 pii_excluded_columns: vec![],
2880 },
2881 counts: IndexCounts {
2882 tables: 3,
2883 views: 0,
2884 functions: 0,
2885 enums: 0,
2886 },
2887 shards: vec![ObjectShard {
2888 file: "objects-public-0.json".into(),
2889 schema: "public".into(),
2890 objects: 3,
2891 bytes: 512,
2892 hash: "abc".into(),
2893 }],
2894 derived: IndexDerived {
2895 tokens: TOKENS_FILE.into(),
2896 join_graph: JOIN_GRAPH_FILE.into(),
2897 values: None,
2898 embeddings: None,
2899 embeddings_meta: None,
2900 },
2901 stats: IndexStats {
2902 build_ms: 1,
2903 queries_run: 1,
2904 warnings: vec![],
2905 },
2906 };
2907 store.write_manifest(&base, &manifest).unwrap();
2908
2909 fn entry(oid: u32) -> ObjectEntry {
2910 ObjectEntry {
2911 kind: DbObjectKind::Table,
2912 oid,
2913 object_hash: format!("hash{oid}"),
2914 comment: None,
2915 row_estimate: 10.0,
2916 size_bytes: 8192,
2917 columns: vec![ColumnEntry {
2918 name: "id".into(),
2919 type_name: "integer".into(),
2920 not_null: true,
2921 default_value: None,
2922 comment: None,
2923 ordinal: 1,
2924 is_pk: Some(true),
2925 profile: None,
2926 pii: None,
2927 }],
2928 primary_key: Some(vec!["id".into()]),
2929 foreign_keys: None,
2930 indexes: None,
2931 checks: None,
2932 excluded: None,
2933 definition: None,
2934 signature: None,
2935 language: None,
2936 volatility: None,
2937 body: None,
2938 values: None,
2939 base_type: None,
2940 constraint: None,
2941 }
2942 }
2943 let mut shard = HashMap::new();
2944 shard.insert("public.orders".into(), entry(1));
2945 shard.insert("public.customers".into(), entry(2));
2946 shard.insert("public.order_items".into(), entry(3));
2947 store
2948 .write_shard_entries(&base, "objects-public-0.json", &shard)
2949 .unwrap();
2950
2951 let graph = JoinGraph {
2952 edges: vec![
2953 JoinEdge {
2954 from: "public.order_items".into(),
2955 to: "public.orders".into(),
2956 via: "order_items_order_id_fkey".into(),
2957 cols: vec![("order_id".into(), "id".into())],
2958 inferred: None,
2959 disabled: None,
2960 },
2961 JoinEdge {
2962 from: "public.orders".into(),
2963 to: "public.customers".into(),
2964 via: "orders_customer_id_fkey".into(),
2965 cols: vec![("customer_id".into(), "id".into())],
2966 inferred: None,
2967 disabled: None,
2968 },
2969 ],
2970 };
2971 store.write_join_graph(&base, &graph).unwrap();
2972 }
2973
2974 fn join_path_router() -> ToolRouter {
2975 let tmp = tempfile::TempDir::new().unwrap();
2976 let store = IndexStore::new(tmp.path());
2977 write_join_path_fixture(&store);
2978 std::mem::forget(tmp); let session = ToolSession::for_tests(
2980 vec![test_conn()],
2981 PolicyFilter::default(),
2982 Some(IndexStore::new(store.root())),
2983 );
2984 ToolRouter::with_index_store(session, Some(store))
2985 }
2986
2987 #[tokio::test]
2991 async fn get_join_path_resolves_unqualified_names() {
2992 let router = join_path_router();
2993 let out = router
2994 .call(
2995 "get_join_path",
2996 json!({ "a": "order_items", "b": "orders" }),
2997 )
2998 .await;
2999 assert!(!out.is_error, "{}", out.text);
3000 let structured = out.structured.unwrap();
3001 assert_eq!(structured["resolved_a"], "public.order_items");
3002 assert_eq!(structured["resolved_b"], "public.orders");
3003 assert_eq!(structured["path"][0]["via"], "order_items_order_id_fkey");
3004 }
3005
3006 #[tokio::test]
3007 async fn get_join_path_qualified_still_works() {
3008 let router = join_path_router();
3009 let out = router
3010 .call(
3011 "get_join_path",
3012 json!({ "a": "public.order_items", "b": "public.customers" }),
3013 )
3014 .await;
3015 assert!(!out.is_error, "{}", out.text);
3016 let structured = out.structured.unwrap();
3017 assert_eq!(structured["path"].as_array().unwrap().len(), 2);
3018 }
3019
3020 #[tokio::test]
3021 async fn get_join_path_unknown_name_errors_with_suggestion() {
3022 let router = join_path_router();
3023 let out = router
3024 .call("get_join_path", json!({ "a": "ordrs", "b": "customers" }))
3025 .await;
3026 assert!(out.is_error, "{}", out.text);
3027 assert!(
3028 out.text.contains("did you mean") && out.text.contains("public.orders"),
3029 "expected a did-you-mean suggestion, got: {}",
3030 out.text
3031 );
3032 }
3033
3034 #[tokio::test]
3035 async fn empty_index_dir_returns_build_hint() {
3036 let tmp = tempfile::TempDir::new().unwrap();
3037 let store = IndexStore::new(tmp.path());
3038 let session = ToolSession::for_tests(
3039 vec![test_conn()],
3040 PolicyFilter::default(),
3041 Some(IndexStore::new(tmp.path())),
3042 );
3043 let router = ToolRouter::with_index_store(session, Some(store));
3044 let out = router
3045 .call("describe_object", json!({ "ref": "public.users" }))
3046 .await;
3047 assert!(out.is_error, "{}", out.text);
3048 assert!(
3049 out.text.contains("rebuild_index"),
3050 "expected build hint, got: {}",
3051 out.text
3052 );
3053 }
3054
3055 #[tokio::test]
3056 async fn outcome_tagged_with_connection_id_and_database() {
3057 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
3058 let router = ToolRouter::new(session);
3059 let out = router.call("list_connections", json!({})).await;
3060 let structured = out.structured.expect("structured outcome");
3061 assert_eq!(
3062 structured.get("connectionId").and_then(|v| v.as_str()),
3063 Some("conn-1")
3064 );
3065 assert_eq!(
3066 structured.get("database").and_then(|v| v.as_str()),
3067 Some("appdb")
3068 );
3069 }
3070
3071 #[tokio::test]
3072 async fn setup_connection_returns_needs_input_when_incomplete() {
3073 unsafe {
3074 std::env::remove_var("DATABASE_URL");
3075 std::env::remove_var("POSTGRES_URL");
3076 std::env::remove_var("PGHOST");
3077 }
3078 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
3079 let router = ToolRouter::new(session);
3080 let out = router.call("setup_connection", json!({})).await;
3081 let structured = out.structured.expect("structured outcome");
3082 assert!(structured.get("status").is_some());
3083 }
3084
3085 #[tokio::test]
3086 async fn save_profile_persists_config() {
3087 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
3088 let router = ToolRouter::new(session);
3089 let temp_dir = tempfile::tempdir().unwrap();
3090 let cfg_path = temp_dir.path().join("config.toml");
3091 unsafe {
3092 std::env::set_var("NEXQL_MCP_CONFIG", &cfg_path);
3093 }
3094
3095 let out = router
3096 .call(
3097 "save_profile",
3098 json!({
3099 "name": "staging",
3100 "host": "127.0.0.1",
3101 "port": 5432,
3102 "dbname": "stage_db",
3103 "user": "stage_user"
3104 }),
3105 )
3106 .await;
3107
3108 let structured = out.structured.expect("structured outcome");
3109 assert_eq!(
3110 structured.get("status").and_then(|v| v.as_str()),
3111 Some("saved")
3112 );
3113 assert_eq!(
3114 structured.get("profile").and_then(|v| v.as_str()),
3115 Some("staging")
3116 );
3117 }
3118
3119 #[tokio::test]
3120 async fn check_ddl_safety_tool_dispatches_ast_report() {
3121 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
3122 let router = ToolRouter::new(session);
3123 let out = router
3124 .call(
3125 "check_ddl_safety",
3126 json!({ "ddl": "CREATE INDEX idx_col ON users(col);" }),
3127 )
3128 .await;
3129 let structured = out.structured.expect("structured outcome");
3130 assert_eq!(
3131 structured.get("overall_risk").and_then(|v| v.as_str()),
3132 Some("CRITICAL")
3133 );
3134 }
3135}