1use std::sync::Arc;
7
8use nexql_index::{
9 BuildDepth, BuildMode, BuildRequest, CatalogDb, Embedder, IndexQueryService, IndexScope,
10 IndexStore, PgCatalogDb, QueryPolicyFilter, SearchOptions, build_index,
11};
12use nexql_policy::{PolicyFilter, SqlDecision, enforce_read_table_policy, select_table_refs, validate_readonly_sql};
13use serde_json::{Value, json};
14
15use crate::cell_json::{redact_pii_in_payload, rows_to_json_array};
16use crate::error::ToolError;
17use crate::export::{ExportFormat, columns_from_rows, rows_to_csv, rows_to_sql_insert};
18use crate::plan::{analyze_deep_plan, build_explain_sql, extract_plan_metrics};
19use crate::registry::ToolName;
20use crate::schema::{ToolSpec, active_tools};
21use crate::session::ToolSession;
22use crate::sql::{self, REPORT_LIMIT_DEFAULT, SLOW_QUERIES_DEFAULT, parse_ref};
23use crate::write::{
24 apply_ddl, create_index_concurrently, edit_row, execute_sql, import_data, run_maintenance,
25 terminate_query,
26};
27
28const SEARCH_SCHEMA_LIMIT: usize = 10;
30
31const NO_INDEX_HINT: &str =
32 "No schema index configured — call the 'rebuild_index' tool to build an index.";
33
34#[derive(Debug, Clone)]
35pub struct ToolOutcome {
36 pub text: String,
37 pub structured: Option<Value>,
38 pub is_error: bool,
39}
40
41impl ToolOutcome {
42 pub fn ok_json(value: Value) -> Self {
48 let value = ensure_structured_object(value);
49 let text = serde_json::to_string_pretty(&value).unwrap_or_else(|_| value.to_string());
50 Self {
51 text,
52 structured: Some(value),
53 is_error: false,
54 }
55 }
56
57 pub fn err(msg: impl Into<String>) -> Self {
58 let message = msg.into();
59 Self {
60 text: message.clone(),
61 structured: Some(json!({ "error": message })),
62 is_error: true,
63 }
64 }
65}
66
67fn ensure_structured_object(value: Value) -> Value {
69 match value {
70 Value::Array(rows) => json!({ "rows": rows }),
71 other => other,
72 }
73}
74
75pub struct ToolRouter {
76 session: Arc<ToolSession>,
77 index_override: Option<Option<IndexStore>>,
79 use_semantic: bool,
81 embedder: Option<Arc<dyn Embedder>>,
82 specs: Vec<ToolSpec>,
83 managed_extension: bool,
84}
85
86impl ToolRouter {
87 pub fn new(session: Arc<ToolSession>) -> Self {
88 Self {
89 session,
90 index_override: None,
91 use_semantic: false,
92 embedder: None,
93 specs: active_tools(),
94 managed_extension: false,
95 }
96 }
97
98 pub fn with_index_store(session: Arc<ToolSession>, store: Option<IndexStore>) -> Self {
100 Self {
101 session,
102 index_override: Some(store),
103 use_semantic: false,
104 embedder: None,
105 specs: active_tools(),
106 managed_extension: false,
107 }
108 }
109
110 pub fn with_semantic(
112 mut self,
113 use_semantic: bool,
114 embedder: Option<Arc<dyn Embedder>>,
115 ) -> Self {
116 self.use_semantic = use_semantic;
117 self.embedder = embedder;
118 self
119 }
120
121 pub fn with_profile(mut self, profile: crate::registry::ToolProfile) -> Self {
123 self.specs = crate::schema::tools_for_profile(profile);
124 self
125 }
126
127 pub fn with_managed_extension(mut self, enabled: bool) -> Self {
129 self.managed_extension = enabled;
130 if enabled {
131 const BLOCKED: &[ToolName] = &[
132 ToolName::SetupConnection,
133 ToolName::SaveProfile,
134 ToolName::TestProfile,
135 ToolName::ExportProfile,
136 ToolName::ImportProfile,
137 ];
138 self.specs.retain(|s| !BLOCKED.contains(&s.name));
139 }
140 self
141 }
142
143 pub fn specs(&self) -> &[ToolSpec] {
144 &self.specs
145 }
146
147 fn index_store(&self) -> Option<&IndexStore> {
148 match &self.index_override {
149 Some(inner) => inner.as_ref(),
150 None => self.session.index_store.as_ref(),
151 }
152 }
153
154 fn query_filter(&self) -> QueryPolicyFilter {
155 policy_to_query_filter(&self.session.filter())
156 }
157
158 pub async fn call(&self, name: &str, args: Value) -> ToolOutcome {
159 let outcome = match self.call_inner(name, args).await {
160 Ok(out) => out,
161 Err(e) => ToolOutcome::err(e.to_string()),
162 };
163 self.tag_outcome_with_context(outcome).await
164 }
165
166 async fn tag_outcome_with_context(&self, mut outcome: ToolOutcome) -> ToolOutcome {
167 let (connection_id, database) = self.session.active_context().await;
168 let access_mode = match self.session.access_mode() {
169 nexql_policy::AccessMode::Read => "read",
170 nexql_policy::AccessMode::Write => "write",
171 nexql_policy::AccessMode::Admin => "admin",
172 };
173 let mut freshness: Option<serde_json::Value> = None;
174 if let Some(store) = self.session.index_store.as_ref() {
175 let base = store.base_dir(&connection_id, &database);
176 if let Ok(Some(manifest)) = store.read_manifest(&base) {
177 let stale = self.session.is_index_stale(&connection_id, &database);
178 let mut freshness_obj = json!({
179 "indexedAt": manifest.indexed_at,
180 "schemaFingerprint": manifest.schema_fingerprint,
181 "stale": stale,
182 });
183 if stale {
184 freshness_obj["reason"] = json!("schema_changed");
185 }
186 freshness = Some(freshness_obj);
187 } else {
188 freshness = Some(json!({ "stale": true, "reason": "no_index" }));
189 }
190 }
191 if let Some(ref mut structured) = outcome.structured
192 && let Some(obj) = structured.as_object_mut()
193 {
194 if !obj.contains_key("connectionId") {
195 obj.insert("connectionId".into(), json!(connection_id));
196 }
197 if !obj.contains_key("database") {
198 obj.insert("database".into(), json!(database));
199 }
200 if !obj.contains_key("accessMode") {
201 obj.insert("accessMode".into(), json!(access_mode));
202 }
203 if let Some(ref f) = freshness {
204 obj.insert("freshness".into(), f.clone());
205 }
206 }
207 let header = format!(
208 "[context connectionId={connection_id} database={database} accessMode={access_mode}]\n"
209 );
210 if !outcome.text.starts_with("[context ") {
211 outcome.text = format!("{header}{}", outcome.text);
212 }
213 outcome
214 }
215
216 async fn call_inner(&self, name: &str, args: Value) -> Result<ToolOutcome, ToolError> {
217 let tool = ToolName::parse(name).ok_or_else(|| ToolError::Unknown(name.to_string()))?;
218 match tool {
219 ToolName::ListConnections => Ok(self.list_connections()),
220 ToolName::ListDatabases => self.list_databases(&args).await,
221 ToolName::ListSchemas => self.list_schemas().await,
222 ToolName::ListObjects => self.list_objects(&args).await,
223 ToolName::GetCurrentContext => self.get_current_context().await,
224 ToolName::SwitchConnection => self.switch_connection(&args).await,
225 ToolName::RunSelect => self.run_select(&args).await,
226 ToolName::ExplainQuery => self.explain_query(&args).await,
227 ToolName::SearchSchema => self.search_schema(&args).await,
228 ToolName::DescribeObject => self.describe_object(&args).await,
229 ToolName::GetJoinPath => self.get_join_path(&args).await,
230 ToolName::SampleValues => self.sample_values(&args).await,
231 ToolName::GetDdl => self.get_ddl(&args).await,
232 ToolName::TableStats => self.table_stats(&args).await,
233 ToolName::IndexUsage => self.index_usage(&args).await,
234 ToolName::ListRunningQueries => self.list_running_queries().await,
235 ToolName::FindBlockingLocks => self.find_blocking_locks().await,
236 ToolName::SlowQueries => self.slow_queries(&args).await,
237 ToolName::DbHealthCheck => self.db_health_check().await,
238 ToolName::ExplainAnalyze => self.explain_analyze(&args).await,
239 ToolName::AnalyzeQueryPlan => self.analyze_query_plan(&args).await,
240 ToolName::GetIndexStatus => self.get_index_status().await,
241 ToolName::ListExtensions => self.list_extensions().await,
242 ToolName::ServerSettings => self.server_settings().await,
243 ToolName::SuggestIndexes => self.suggest_indexes(&args).await,
244 ToolName::FindUnusedIndexes => self.find_unused_indexes(&args).await,
245 ToolName::BloatReport => self.bloat_report(&args).await,
246 ToolName::FindMissingFks => self.find_missing_fks(&args).await,
247 ToolName::ExportQuery => self.export_query(&args).await,
248 ToolName::ListRoles => self.list_roles(&args).await,
249 ToolName::DbDashboard => self.db_dashboard().await,
250 ToolName::DeepPlanAnalysis => self.deep_plan_analysis(&args).await,
251 ToolName::SchemaDiff => self.schema_diff(&args).await,
252 ToolName::GenerateMigration => self.generate_migration(&args).await,
253 ToolName::ExecuteSql => self.execute_sql_tool(&args).await,
254 ToolName::EditRow => self.edit_row_tool(&args).await,
255 ToolName::ImportData => self.import_data_tool(&args).await,
256 ToolName::ApplyDdl => self.apply_ddl_tool(&args).await,
257 ToolName::CreateIndexConcurrently => self.create_index_concurrently_tool(&args).await,
258 ToolName::RunMaintenance => self.run_maintenance_tool(&args).await,
259 ToolName::TerminateQuery => self.terminate_query_tool(&args).await,
260 ToolName::ResolveTarget => self.resolve_target(&args).await,
261 ToolName::DiscoverTools => self.discover_tools(&args).await,
262 ToolName::AutoTuneQuery => self.auto_tune_query(&args).await,
263 ToolName::CheckDdlSafety => self.check_ddl_safety_tool(&args).await,
264 ToolName::RebuildIndex => self.rebuild_index_tool(&args).await,
265 ToolName::RefreshIndex => self.refresh_index_tool(&args).await,
266 ToolName::RunDoctor => self.run_doctor_tool().await,
267 ToolName::SetupConnection => self.setup_connection_tool(&args).await,
268 ToolName::SaveProfile => self.save_profile_tool(&args).await,
269 ToolName::TestProfile => self.test_profile_tool(&args).await,
270 ToolName::ExportProfile => self.export_profile_tool(&args).await,
271 ToolName::ImportProfile => self.import_profile_tool(&args).await,
272 }
273 }
274
275 fn require_write(&self) -> Result<(), ToolError> {
276 if !self.session.access_mode().allows_writes() {
277 return Err(ToolError::Execution(
278 "write tools require --access-mode write or admin (current session: read)".into(),
279 ));
280 }
281 Ok(())
282 }
283
284 fn require_admin(&self) -> Result<(), ToolError> {
285 if !self.session.access_mode().allows_admin() {
286 return Err(ToolError::Execution(
287 "admin tools require --access-mode admin".into(),
288 ));
289 }
290 Ok(())
291 }
292
293 async fn execute_sql_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
294 self.require_write()?;
295 let sql = args
296 .get("sql")
297 .and_then(|v| v.as_str())
298 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
299 let dry_run = args
300 .get("dry_run")
301 .and_then(|v| v.as_bool())
302 .unwrap_or(false);
303 execute_sql(&self.session, sql, dry_run).await
304 }
305
306 async fn edit_row_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
307 self.require_write()?;
308 edit_row(&self.session, args).await
309 }
310
311 async fn import_data_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
312 self.require_write()?;
313 import_data(&self.session, args).await
314 }
315
316 async fn apply_ddl_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
317 self.require_admin()?;
318 let sql = args
319 .get("sql")
320 .and_then(|v| v.as_str())
321 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
322 let dry_run = args
323 .get("dry_run")
324 .and_then(|v| v.as_bool())
325 .unwrap_or(false);
326 apply_ddl(&self.session, sql, dry_run).await
327 }
328
329 async fn create_index_concurrently_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
330 self.require_admin()?;
331 let sql = args
332 .get("sql")
333 .and_then(|v| v.as_str())
334 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
335 create_index_concurrently(&self.session, sql).await
336 }
337
338 async fn run_maintenance_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
339 self.require_admin()?;
340 run_maintenance(&self.session, args).await
341 }
342
343 async fn terminate_query_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
344 self.require_admin()?;
345 terminate_query(&self.session, args).await
346 }
347
348 async fn live_databases_by_connection(
351 &self,
352 ) -> std::collections::HashMap<String, std::collections::HashSet<String>> {
353 use std::collections::HashMap;
354 let mut map = HashMap::new();
355 for conn in self.session.connections() {
356 if let Ok(names) = self.list_database_names_for(&conn).await {
357 map.insert(conn.id.clone(), names.into_iter().collect());
358 }
359 }
360 map
361 }
362
363 async fn list_database_names_for(
364 &self,
365 conn: &crate::session::ConnectionInfo,
366 ) -> Result<Vec<String>, ToolError> {
367 let client = if self.session.active_context().await.0 == conn.id {
368 self.session.checkout().await?
369 } else {
370 let pool_opts = self.session.pool_opts();
371 let pool = nexql_conn::create_pool(&conn.params, &pool_opts).await?;
372 nexql_conn::checkout_guarded(&pool, &pool_opts).await?
373 };
374 let rows = client
375 .query(
376 "SELECT datname FROM pg_database WHERE datistemplate = false ORDER BY datname",
377 &[],
378 )
379 .await?;
380 Ok(rows.iter().map(|r| r.get(0)).collect())
381 }
382
383 async fn resolve_target(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
384 let hint = args
385 .get("hint")
386 .and_then(|v| v.as_str())
387 .map(str::trim)
388 .filter(|s| !s.is_empty());
389 let object_hint = args
390 .get("objectHint")
391 .and_then(|v| v.as_str())
392 .map(str::trim)
393 .filter(|s| !s.is_empty());
394 if hint.is_none() && object_hint.is_none() {
395 return Err(ToolError::InvalidArgs(
396 "At least one of \"hint\" or \"objectHint\" is required.".into(),
397 ));
398 }
399
400 let connections = self.session.connections();
401 if connections.is_empty() {
402 return Ok(ToolOutcome::err("No connections configured."));
403 }
404
405 #[derive(Clone)]
406 struct Candidate {
407 connection_id: String,
408 database: String,
409 }
410 fn key_of(c: &Candidate) -> String {
411 format!("{}\u{0}{}", c.connection_id, c.database)
412 }
413
414 let indexed: Vec<(String, String)> = self
415 .index_store()
416 .map(|store| store.list_indexed_databases().unwrap_or_default())
417 .unwrap_or_default();
418
419 let live_dbs = self.live_databases_by_connection().await;
420
421 let mut seen = std::collections::HashSet::new();
422 let mut candidates: Vec<Candidate> = Vec::new();
423 let mut add_candidate = |connection_id: &str, database: &str| {
424 if !connections.iter().any(|c| c.id == connection_id) {
425 return;
426 }
427 let key = format!("{connection_id}\u{0}{database}");
428 if !seen.insert(key) {
429 return;
430 }
431 candidates.push(Candidate {
432 connection_id: connection_id.to_string(),
433 database: database.to_string(),
434 });
435 };
436 for (cid, db) in &indexed {
437 if let Some(set) = live_dbs.get(cid) {
438 if set.contains(db) {
439 add_candidate(cid, db);
440 } else if let Some(store) = self.index_store() {
441 let _ = store.clear_index(cid, db);
442 }
443 }
444 }
445 for c in &connections {
446 if let Some(dbs) = live_dbs.get(&c.id) {
447 for db in dbs {
448 add_candidate(&c.id, db);
449 }
450 } else {
451 let db = c.database.clone().unwrap_or_else(|| "postgres".into());
452 add_candidate(&c.id, &db);
453 }
454 }
455
456 let mut scored: std::collections::HashMap<String, (Candidate, f64, Vec<String>)> =
457 std::collections::HashMap::new();
458
459 if let Some(hint) = hint {
460 for c in &candidates {
461 let Some(conn) = connections.iter().find(|x| x.id == c.connection_id) else {
462 continue;
463 };
464 let fields: [(&str, &str); 3] = [
465 ("connection name", conn.name.as_str()),
466 ("host", conn.host.as_deref().unwrap_or("")),
467 ("database", c.database.as_str()),
468 ];
469 let mut best = 0.0f64;
470 let mut best_field = "";
471 for (label, value) in fields {
472 let s = fuzzy_score(hint, value);
473 if s > best {
474 best = s;
475 best_field = label;
476 }
477 }
478 if best > 0.0 {
479 let entry = scored
480 .entry(key_of(c))
481 .or_insert_with(|| (c.clone(), 0.0, Vec::new()));
482 entry.1 += best;
483 entry
484 .2
485 .push(format!("{best_field} matched hint \"{hint}\" ({best:.0})"));
486 }
487 }
488 }
489
490 if let Some(object_hint) = object_hint
491 && let Some(store) = self.index_store()
492 {
493 let filter = self.query_filter();
494 for (cid, db) in &indexed {
495 let svc = IndexQueryService::new(store, cid.clone(), db.clone());
496 if let Ok(hits) = svc.search_schema(
497 object_hint,
498 3,
499 Some(&filter),
500 SearchOptions {
501 use_semantic: self.use_semantic,
502 embedder: self.embedder.as_deref(),
503 },
504 ) && let Some(top) = hits.first()
505 {
506 let c = Candidate {
507 connection_id: cid.clone(),
508 database: db.clone(),
509 };
510 let entry = scored
511 .entry(key_of(&c))
512 .or_insert_with(|| (c.clone(), 0.0, Vec::new()));
513 entry.1 += top.score * 10.0;
514 entry.2.push(format!(
515 "schema search for \"{object_hint}\" found {} (score {:.2})",
516 top.ref_, top.score
517 ));
518 }
519 }
520 }
521
522 let mut ranked: Vec<(Candidate, f64, Vec<String>)> = scored.into_values().collect();
523 ranked.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
524
525 if ranked.is_empty() {
526 let candidates_json: Vec<Value> = connections
527 .iter()
528 .map(|c| {
529 json!({
530 "connectionId": c.id,
531 "connectionName": c.name,
532 "database": c.database.clone().unwrap_or_else(|| "postgres".into()),
533 })
534 })
535 .collect();
536 return Ok(ToolOutcome::ok_json(json!({
537 "ambiguous": true,
538 "message": format!(
539 "No connection/database matched \"{}\". Choose from the configured connections.",
540 hint.or(object_hint).unwrap_or_default()
541 ),
542 "candidates": candidates_json
543 })));
544 }
545
546 let winner = &ranked[0];
547 let is_tied = ranked
548 .get(1)
549 .is_some_and(|runner_up| runner_up.1 >= winner.1 * 0.85);
550
551 if is_tied {
552 let threshold = winner.1 * 0.85;
553 let tied: Vec<&(Candidate, f64, Vec<String>)> =
554 ranked.iter().filter(|r| r.1 >= threshold).take(5).collect();
555 let candidates_json: Vec<Value> = tied
556 .iter()
557 .filter_map(|(c, score, evidence)| {
558 connections
559 .iter()
560 .find(|x| x.id == c.connection_id)
561 .map(|conn| {
562 json!({
563 "connectionId": c.connection_id,
564 "connectionName": conn.name,
565 "database": c.database,
566 "score": score,
567 "evidence": evidence,
568 })
569 })
570 })
571 .collect();
572 return Ok(ToolOutcome::ok_json(json!({
573 "ambiguous": true,
574 "message": format!("{} equally-plausible candidates matched.", tied.len()),
575 "candidates": candidates_json
576 })));
577 }
578
579 let (winner_candidate, winner_score, winner_evidence) = winner;
580
581 if let Some(object_hint) = object_hint
582 && let Some(store) = self.index_store()
583 {
584 let filter = self.query_filter();
585 let svc = IndexQueryService::new(
586 store,
587 &winner_candidate.connection_id,
588 &winner_candidate.database,
589 );
590 if let Ok(hits) = svc.search_schema(
591 object_hint,
592 5,
593 Some(&filter),
594 SearchOptions {
595 use_semantic: self.use_semantic,
596 embedder: self.embedder.as_deref(),
597 },
598 ) && hits.len() >= 2
599 {
600 let top_score = hits[0].score;
601 let tied: Vec<&nexql_index::RankedHit> = hits
602 .iter()
603 .filter(|h| scores_equal(h.score, top_score))
604 .collect();
605 if tied.len() > 1 {
606 let candidates_json: Vec<Value> = tied
607 .iter()
608 .map(|h| {
609 json!({
610 "ref": h.ref_,
611 "score": h.score,
612 "kind": h.kind,
613 "connectionId": winner_candidate.connection_id,
614 "database": winner_candidate.database,
615 })
616 })
617 .collect();
618 return Ok(ToolOutcome::ok_json(json!({
619 "ambiguous": true,
620 "message": format!(
621 "{} objects matched \"{object_hint}\" with equal scores — choose explicitly.",
622 tied.len()
623 ),
624 "candidates": candidates_json,
625 })));
626 }
627 }
628 }
629
630 self.session
631 .switch(
632 &winner_candidate.connection_id,
633 Some(winner_candidate.database.clone()),
634 )
635 .await?;
636 let conn = connections
637 .iter()
638 .find(|x| x.id == winner_candidate.connection_id)
639 .ok_or_else(|| ToolError::Execution("resolved connection vanished".into()))?;
640
641 Ok(ToolOutcome::ok_json(json!({
642 "resolved": true,
643 "connectionId": winner_candidate.connection_id,
644 "connectionName": conn.name,
645 "database": winner_candidate.database,
646 "confidence": winner_score,
647 "evidence": winner_evidence,
648 })))
649 }
650
651 async fn discover_tools(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
652 let query = args
653 .get("query")
654 .and_then(|v| v.as_str())
655 .map(str::to_lowercase);
656 let category = args
657 .get("category")
658 .and_then(|v| v.as_str())
659 .map(str::to_lowercase);
660
661 let all_specs = active_tools();
663 let filtered: Vec<Value> = all_specs
664 .into_iter()
665 .filter(|spec| {
666 if spec.name == ToolName::DiscoverTools {
667 return false;
668 }
669 if let Some(ref cat) = category {
670 match cat.as_str() {
671 "query" if !ToolName::QUERY_PROFILE.contains(&spec.name) => return false,
672 "dba" if !ToolName::DBA_PROFILE.contains(&spec.name) => return false,
673 "write" if !ToolName::PHASE9.contains(&spec.name) => return false,
674 _ => {}
675 }
676 }
677 if let Some(ref q) = query {
678 let name_match = spec.name.as_str().contains(q.as_str());
679 let desc_match = spec.description.to_lowercase().contains(q.as_str());
680 if !name_match && !desc_match {
681 return false;
682 }
683 }
684 true
685 })
686 .map(|spec| {
687 json!({
688 "name": spec.name.as_str(),
689 "description": spec.description,
690 "input_schema": spec.input_schema,
691 })
692 })
693 .collect();
694
695 Ok(ToolOutcome::ok_json(json!({
696 "query": args.get("query"),
697 "category": args.get("category"),
698 "count": filtered.len(),
699 "tools": filtered,
700 })))
701 }
702
703 fn build_tuning_summary(plan_structured: &Option<Value>, suggestions: &Value) -> String {
704 let mut parts = Vec::new();
705 if let Some(structured) = plan_structured
706 && let Some(metrics) = structured.get("metrics")
707 {
708 if let Some(exec_time) = metrics.get("executionTime").and_then(|v| v.as_f64()) {
709 parts.push(format!("Query executed in {:.2}ms.", exec_time));
710 }
711 if let Some(seq_scans) = metrics.get("sequentialScans").and_then(|v| v.as_u64())
712 && seq_scans > 0
713 {
714 parts.push(format!("Found {seq_scans} sequential scan(s)."));
715 }
716 }
717
718 let candidate_count = suggestions
719 .get("high_seq_scan_tables")
720 .and_then(|v| v.as_array())
721 .map(|a| a.len())
722 .unwrap_or(0)
723 + suggestions
724 .get("unindexed_fk_columns")
725 .and_then(|v| v.as_array())
726 .map(|a| a.len())
727 .unwrap_or(0);
728
729 if candidate_count > 0 {
730 parts.push(format!(
731 "{candidate_count} index recommendation(s) identified."
732 ));
733 } else {
734 parts.push("No explicit index candidate recommendations generated.".into());
735 }
736
737 if parts.is_empty() {
738 "Auto-tune evaluation complete. Inspect execution plan and index recommendations."
739 .into()
740 } else {
741 parts.join(" ")
742 }
743 }
744
745 async fn auto_tune_query(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
746 let sql = args
747 .get("sql")
748 .and_then(|v| v.as_str())
749 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
750
751 let deep_plan = self
752 .deep_plan_analysis(&json!({ "sql": sql, "analyze": true }))
753 .await?;
754
755 let suggestions_res = self.suggest_indexes(&json!({ "sql": sql })).await;
756 let (suggestions_data, suggestions_error) = match suggestions_res {
757 Ok(outcome) => (outcome.structured.unwrap_or(json!([])), None),
758 Err(e) => (json!([]), Some(e.to_string())),
759 };
760
761 let summary_text = Self::build_tuning_summary(&deep_plan.structured, &suggestions_data);
762
763 let mut payload = json!({
764 "target_query": sql,
765 "deep_plan_analysis": deep_plan.structured,
766 "index_suggestions": suggestions_data,
767 "tuning_summary": summary_text,
768 });
769
770 if let Some(err) = suggestions_error {
771 payload["suggestions_error"] = json!(err);
772 }
773
774 Ok(ToolOutcome::ok_json(payload))
775 }
776
777 async fn check_ddl_safety_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
778 let ddl = args
779 .get("ddl")
780 .and_then(|v| v.as_str())
781 .ok_or_else(|| ToolError::InvalidArgs("ddl is required".into()))?;
782
783 let report = crate::dba_guard::analyze_ddl_safety(ddl);
784 Ok(ToolOutcome::ok_json(report))
785 }
786
787 async fn rebuild_index_tool(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
788 let store = self
789 .index_store()
790 .ok_or_else(|| ToolError::Execution("Index store unavailable".into()))?;
791 let (connection_id, database) = self.session.active_context().await;
792 let depth_str = args
793 .get("depth")
794 .and_then(|v| v.as_str())
795 .unwrap_or("structure");
796 let depth: BuildDepth = match depth_str.to_lowercase().as_str() {
797 "profiles" | "full" => BuildDepth::Profiles,
798 _ => BuildDepth::Structure,
799 };
800
801 let req = BuildRequest {
802 connection_id: connection_id.clone(),
803 database: database.clone(),
804 scope: IndexScope {
805 included_schemas: vec![],
806 excluded_objects: vec![],
807 pii_excluded_columns: vec![],
808 },
809 depth,
810 build_mode: BuildMode::Guided,
811 environment: "development".into(),
812 embeddings: self.use_semantic,
813 };
814
815 let client = self.session.checkout().await?;
816 let db = PgCatalogDb::new(&client);
817 let manifest = build_index(store, &db, &req, None, None, self.embedder.as_deref())
818 .await
819 .map_err(|e| ToolError::Execution(format!("Index build failed: {e}")))?;
820 self.session
821 .clear_index_stale(&connection_id, &database);
822
823 Ok(ToolOutcome::ok_json(json!({
824 "status": "completed",
825 "connection_id": connection_id,
826 "database": database,
827 "schema_fingerprint": manifest.schema_fingerprint,
828 "counts": manifest.counts,
829 "build_ms": manifest.stats.build_ms,
830 })))
831 }
832
833 async fn refresh_index_tool(&self, _args: &Value) -> Result<ToolOutcome, ToolError> {
834 let store = self
835 .index_store()
836 .ok_or_else(|| ToolError::Execution("Index store unavailable".into()))?;
837 let (connection_id, database) = self.session.active_context().await;
838 let base = store.base_dir(&connection_id, &database);
839 let manifest = store.read_manifest(&base)?.ok_or_else(|| {
840 ToolError::Execution(
841 "No existing index manifest to refresh — call 'rebuild_index'.".into(),
842 )
843 })?;
844
845 let req = BuildRequest {
846 connection_id: connection_id.clone(),
847 database: database.clone(),
848 scope: manifest.scope,
849 depth: manifest.build_depth,
850 build_mode: manifest.build_mode,
851 environment: manifest.environment,
852 embeddings: self.use_semantic,
853 };
854
855 let client = self.session.checkout().await?;
856 let db = PgCatalogDb::new(&client);
857 let new_manifest = build_index(store, &db, &req, None, None, self.embedder.as_deref())
858 .await
859 .map_err(|e| ToolError::Execution(format!("Index refresh failed: {e}")))?;
860 self.session
861 .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 async fn search_schema(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1263 let query = args
1264 .get("query")
1265 .and_then(|v| v.as_str())
1266 .unwrap_or("")
1267 .trim();
1268 if query.is_empty() {
1269 return Ok(ToolOutcome::ok_json(json!([])));
1270 }
1271 let (store, connection_id, database) = self.index_service().await?;
1272 let svc = IndexQueryService::new(store, &connection_id, &database);
1273 let filter = self.query_filter();
1274 let hits = svc.search_schema(
1275 query,
1276 SEARCH_SCHEMA_LIMIT,
1277 Some(&filter),
1278 SearchOptions {
1279 use_semantic: self.use_semantic,
1280 embedder: self.embedder.as_deref(),
1281 },
1282 )?;
1283 let rows: Vec<Value> = hits
1284 .into_iter()
1285 .map(|h| {
1286 json!({
1287 "ref": h.ref_,
1288 "score": h.score,
1289 "kind": h.kind,
1290 })
1291 })
1292 .collect();
1293 Ok(ToolOutcome::ok_json(json!(rows)))
1294 }
1295
1296 async fn describe_object(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1297 let ref_ = args
1298 .get("ref")
1299 .and_then(|v| v.as_str())
1300 .ok_or_else(|| ToolError::InvalidArgs("ref is required".into()))?;
1301 let (store, connection_id, database) = self.index_service().await?;
1302 let svc = IndexQueryService::new(store, &connection_id, &database);
1303 let filter = self.query_filter();
1304 let entry = svc.describe_object(ref_, Some(&filter))?;
1305 let value = serde_json::to_value(entry).map_err(|e| ToolError::Execution(e.to_string()))?;
1306 Ok(ToolOutcome::ok_json(value))
1307 }
1308
1309 async fn get_join_path(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1310 let a = args
1311 .get("a")
1312 .and_then(|v| v.as_str())
1313 .ok_or_else(|| ToolError::InvalidArgs("a is required".into()))?;
1314 let b = args
1315 .get("b")
1316 .and_then(|v| v.as_str())
1317 .ok_or_else(|| ToolError::InvalidArgs("b is required".into()))?;
1318 let (store, connection_id, database) = self.index_service().await?;
1319 let svc = IndexQueryService::new(store, &connection_id, &database);
1320 let path = svc.get_join_path(a, b)?;
1321 let value = serde_json::to_value(path).map_err(|e| ToolError::Execution(e.to_string()))?;
1322 Ok(ToolOutcome::ok_json(value))
1323 }
1324
1325 async fn sample_values(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1326 let ref_ = args
1327 .get("ref")
1328 .and_then(|v| v.as_str())
1329 .ok_or_else(|| ToolError::InvalidArgs("ref is required".into()))?;
1330 let col = args
1331 .get("col")
1332 .and_then(|v| v.as_str())
1333 .ok_or_else(|| ToolError::InvalidArgs("col is required".into()))?;
1334 let (store, connection_id, database) = self.index_service().await?;
1335 let svc = IndexQueryService::new(store, &connection_id, &database);
1336 let filter = self.query_filter();
1337 let result = svc.sample_values(ref_, col, Some(&filter), None)?;
1338
1339 let mut values = result.values;
1340 let mut message = result.message;
1341
1342 if values.is_empty()
1343 && let Ok(client) = self.session.checkout().await
1344 {
1345 let parts: Vec<&str> = ref_.split('.').collect();
1346 let (schema, table) = match parts.as_slice() {
1347 [s, t] => (*s, *t),
1348 _ => ("public", ref_),
1349 };
1350 let safe_schema = schema.replace('"', "\"\"");
1351 let safe_table = table.replace('"', "\"\"");
1352 let safe_col = col.replace('"', "\"\"");
1353 let query = format!(
1354 "SELECT DISTINCT \"{safe_col}\"::text FROM \"{safe_schema}\".\"{safe_table}\" WHERE \"{safe_col}\" IS NOT NULL LIMIT 20"
1355 );
1356 if let Ok(rows) = client.query(&query, &[]).await {
1357 let sampled: Vec<String> = rows
1358 .iter()
1359 .filter_map(|r| r.get::<_, Option<String>>(0))
1360 .collect();
1361 if !sampled.is_empty() {
1362 values = sampled;
1363 message = None;
1364 }
1365 }
1366 }
1367
1368 let mut payload = json!({ "values": values });
1369 if let Some(msg) = message {
1370 payload["message"] = json!(msg);
1371 }
1372 Ok(ToolOutcome::ok_json(payload))
1373 }
1374
1375 fn list_connections(&self) -> ToolOutcome {
1376 let rows: Vec<Value> = self
1377 .session
1378 .connections()
1379 .iter()
1380 .map(|c| {
1381 json!({
1382 "id": c.id,
1383 "name": c.name,
1384 "host": c.host,
1385 "port": c.port,
1386 "database": c.database,
1387 })
1388 })
1389 .collect();
1390 ToolOutcome::ok_json(json!(rows))
1391 }
1392
1393 async fn list_databases(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1394 let connection_id = args
1395 .get("connectionId")
1396 .and_then(|v| v.as_str())
1397 .ok_or_else(|| ToolError::InvalidArgs("connectionId is required".into()))?;
1398 let conn = self
1399 .session
1400 .connections()
1401 .into_iter()
1402 .find(|c| c.id == connection_id)
1403 .ok_or_else(|| {
1404 ToolError::Execution(format!(
1405 "Connection not found for ID: {connection_id} — call list_connections"
1406 ))
1407 })?;
1408 let client = {
1410 if self.session.active_context().await.0 == connection_id {
1412 self.session.checkout().await?
1413 } else {
1414 let pool_opts = self.session.pool_opts();
1415 let pool = nexql_conn::create_pool(&conn.params, &pool_opts).await?;
1416 nexql_conn::checkout_guarded(&pool, &pool_opts).await?
1417 }
1418 };
1419 let rows = client
1420 .query(
1421 "SELECT datname FROM pg_database WHERE datistemplate = false ORDER BY datname",
1422 &[],
1423 )
1424 .await?;
1425 let names: Vec<String> = rows.iter().map(|r| r.get(0)).collect();
1426 Ok(ToolOutcome::ok_json(json!(names)))
1427 }
1428
1429 async fn list_schemas(&self) -> Result<ToolOutcome, ToolError> {
1430 let client = self.session.checkout().await?;
1431 let rows = client
1432 .query(
1433 r#"
1434 SELECT nspname AS schema_name
1435 FROM pg_namespace
1436 WHERE nspname NOT IN ('pg_catalog', 'information_schema', 'pg_toast')
1437 AND nspname NOT LIKE 'pg_%'
1438 ORDER BY nspname
1439 "#,
1440 &[],
1441 )
1442 .await?;
1443 let out: Vec<Value> = rows
1444 .iter()
1445 .filter(|r| {
1446 let name: String = r.get(0);
1447 self.session.filter().allows_schema(&name)
1448 })
1449 .map(|r| json!({ "schema_name": r.get::<_, String>(0) }))
1450 .collect();
1451 Ok(ToolOutcome::ok_json(json!(out)))
1452 }
1453
1454 async fn list_objects(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1455 let schema = args
1456 .get("schema")
1457 .and_then(|v| v.as_str())
1458 .unwrap_or("public");
1459 if !schema
1460 .chars()
1461 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
1462 {
1463 return Err(ToolError::InvalidArgs(
1464 "Invalid or missing schema name format".into(),
1465 ));
1466 }
1467 if !self.session.filter().allows_schema(schema) {
1468 return Ok(ToolOutcome::ok_json(json!([])));
1469 }
1470 let kind = args.get("kind").and_then(|v| v.as_str());
1471 let mut queries = Vec::new();
1472 let push_rel = |queries: &mut Vec<String>, relkinds: &[&str], label: &str| {
1473 let kinds = relkinds
1474 .iter()
1475 .map(|k| format!("'{k}'"))
1476 .collect::<Vec<_>>()
1477 .join(",");
1478 queries.push(format!(
1479 r#"
1480 SELECT n.nspname AS schema, c.relname AS name, '{label}' AS kind,
1481 d.description AS comment
1482 FROM pg_class c
1483 JOIN pg_namespace n ON n.oid = c.relnamespace
1484 LEFT JOIN pg_description d ON d.objoid = c.oid AND d.objsubid = 0
1485 WHERE n.nspname = $1 AND c.relkind IN ({kinds})
1486 "#
1487 ));
1488 };
1489 if kind.is_none() || kind == Some("table") {
1490 push_rel(&mut queries, &["r", "f", "p"], "table");
1491 }
1492 if kind.is_none() || kind == Some("view") {
1493 push_rel(&mut queries, &["v"], "view");
1494 }
1495 if kind.is_none() || kind == Some("matview") {
1496 push_rel(&mut queries, &["m"], "matview");
1497 }
1498 if queries.is_empty() {
1499 return Ok(ToolOutcome::ok_json(json!([])));
1500 }
1501 let sql = queries.join("\nUNION ALL\n") + "\nORDER BY kind, name";
1502 let client = self.session.checkout().await?;
1503 let rows = client.query(&sql, &[&schema]).await?;
1504 let out: Vec<Value> = rows
1505 .iter()
1506 .filter(|r| {
1507 let s: String = r.get("schema");
1508 let name: String = r.get("name");
1509 self.session.filter().allows_table(&s, &name)
1510 })
1511 .map(|r| {
1512 json!({
1513 "schema": r.get::<_, String>("schema"),
1514 "name": r.get::<_, String>("name"),
1515 "kind": r.get::<_, String>("kind"),
1516 "comment": r.get::<_, Option<String>>("comment"),
1517 })
1518 })
1519 .collect();
1520 Ok(ToolOutcome::ok_json(json!(out)))
1521 }
1522
1523 async fn get_current_context(&self) -> Result<ToolOutcome, ToolError> {
1524 let (connection_id, database) = self.session.active_context().await;
1525 let conn = self
1526 .session
1527 .connections()
1528 .into_iter()
1529 .find(|c| c.id == connection_id);
1530 Ok(ToolOutcome::ok_json(json!({
1531 "connectionId": connection_id,
1532 "connectionName": conn.as_ref().map(|c| c.name.clone()).unwrap_or_else(|| "Unknown".into()),
1533 "database": database,
1534 "host": conn.as_ref().and_then(|c| c.host.clone()),
1535 "port": conn.as_ref().and_then(|c| c.port),
1536 "access_mode": match self.session.access_mode() {
1537 nexql_policy::AccessMode::Read => "read",
1538 nexql_policy::AccessMode::Write => "write",
1539 nexql_policy::AccessMode::Admin => "admin",
1540 },
1541 })))
1542 }
1543
1544 async fn switch_connection(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1545 let connection_id = args
1546 .get("connectionId")
1547 .and_then(|v| v.as_str())
1548 .ok_or_else(|| ToolError::InvalidArgs("connectionId is required".into()))?;
1549 let database = args
1550 .get("database")
1551 .and_then(|v| v.as_str())
1552 .map(str::to_owned);
1553 self.session.switch(connection_id, database).await?;
1554 self.get_current_context().await
1555 }
1556
1557 async fn run_select(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1558 let sql = args
1559 .get("sql")
1560 .and_then(|v| v.as_str())
1561 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
1562 match validate_readonly_sql(sql)? {
1563 SqlDecision::Allow => {}
1564 SqlDecision::Reject => {
1565 return Err(ToolError::Execution(
1566 "Security Error: Only read-only SELECT, WITH, or EXPLAIN statements are permitted."
1567 .into(),
1568 ));
1569 }
1570 }
1571 enforce_read_table_policy(&self.session.filter(), sql)?;
1572 let trimmed = sql.trim().to_ascii_lowercase();
1573 if trimmed.starts_with("explain") {
1574 return self.run_select_internal(sql, None).await;
1575 }
1576 let max_rows = self.session.caps().max_rows;
1577 self.run_select_internal(sql, Some(max_rows)).await
1578 }
1579
1580 async fn explain_query(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1581 let sql = args
1582 .get("sql")
1583 .and_then(|v| v.as_str())
1584 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
1585 match validate_readonly_sql(sql)? {
1586 SqlDecision::Allow => {}
1587 SqlDecision::Reject => {
1588 return Err(ToolError::Execution(
1589 "Security Error: Only SELECT, WITH, or EXPLAIN statements can be analyzed."
1590 .into(),
1591 ));
1592 }
1593 }
1594 enforce_read_table_policy(&self.session.filter(), sql)?;
1595 let clean = if sql.trim().to_ascii_lowercase().starts_with("explain") {
1596 sql.to_string()
1597 } else {
1598 format!("EXPLAIN {sql}")
1599 };
1600 if validate_readonly_sql(&clean)? == SqlDecision::Reject {
1602 return Err(ToolError::Execution(
1603 "Security Error: EXPLAIN target is not read-only.".into(),
1604 ));
1605 }
1606 self.run_select_internal(&clean, None).await
1607 }
1608
1609 async fn get_ddl(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1610 let ref_ = args
1611 .get("ref")
1612 .and_then(|v| v.as_str())
1613 .ok_or_else(|| ToolError::InvalidArgs("ref is required".into()))?;
1614 let (schema, name) = parse_ref(ref_).map_err(ToolError::InvalidArgs)?;
1615 let kind = args.get("kind").and_then(|v| v.as_str()).unwrap_or("table");
1616 let reg = sql::regclass_literal(&schema, &name);
1617 let client = self.session.checkout().await?;
1618
1619 match kind {
1620 "view" | "matview" => {
1621 let sql = format!("SELECT pg_get_viewdef({reg}, true) AS definition");
1622 let rows = client.query(&sql, &[]).await?;
1623 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1624 }
1625 "function" => {
1626 let sql = format!(
1627 r#"SELECT p.proname AS name, pg_get_functiondef(p.oid) AS definition
1628 FROM pg_proc p
1629 JOIN pg_namespace n ON n.oid = p.pronamespace
1630 WHERE n.nspname = '{schema}' AND p.proname = '{name}'"#
1631 );
1632 let rows = client.query(&sql, &[]).await?;
1633 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1634 }
1635 "index" => {
1636 let sql = format!("SELECT pg_get_indexdef({reg}) AS definition");
1637 let rows = client.query(&sql, &[]).await?;
1638 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1639 }
1640 "table" => {
1641 let columns = client
1642 .query(&sql::column_details(&schema, &name), &[])
1643 .await?;
1644 let constraints = client
1645 .query(
1646 &format!(
1647 r#"SELECT conname AS name, pg_get_constraintdef(oid) AS definition
1648 FROM pg_constraint WHERE conrelid = {reg} ORDER BY conname"#
1649 ),
1650 &[],
1651 )
1652 .await?;
1653 let indexes = client
1654 .query(
1655 &format!(
1656 r#"SELECT indexname AS name, indexdef AS definition
1657 FROM pg_indexes
1658 WHERE schemaname = '{schema}' AND tablename = '{name}'
1659 ORDER BY indexname"#
1660 ),
1661 &[],
1662 )
1663 .await?;
1664 Ok(ToolOutcome::ok_json(json!({
1665 "table": format!("{schema}.{name}"),
1666 "columns": rows_to_json(&columns),
1667 "constraints": rows_to_json(&constraints),
1668 "indexes": rows_to_json(&indexes),
1669 })))
1670 }
1671 other => Err(ToolError::InvalidArgs(format!(
1672 "Unsupported DDL kind \"{other}\". Use table, view, matview, function, or index."
1673 ))),
1674 }
1675 }
1676
1677 async fn table_stats(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1678 let ref_ = args
1679 .get("ref")
1680 .and_then(|v| v.as_str())
1681 .ok_or_else(|| ToolError::InvalidArgs("ref is required".into()))?;
1682 let (schema, name) = parse_ref(ref_).map_err(ToolError::InvalidArgs)?;
1683 let client = self.session.checkout().await?;
1684 let stats = client.query(&sql::table_stats(&schema, &name), &[]).await?;
1685 let activity = client
1686 .query(&sql::table_activity(&schema, &name), &[])
1687 .await?;
1688 let columns = client
1689 .query(&sql::column_stats(&schema, &name), &[])
1690 .await?;
1691 let size = rows_to_json(&stats)
1692 .as_array()
1693 .and_then(|a| a.first())
1694 .cloned()
1695 .unwrap_or(Value::Null);
1696 let activity = rows_to_json(&activity)
1697 .as_array()
1698 .and_then(|a| a.first())
1699 .cloned()
1700 .unwrap_or(Value::Null);
1701 Ok(ToolOutcome::ok_json(json!({
1702 "size": size,
1703 "activity": activity,
1704 "columns": rows_to_json(&columns),
1705 })))
1706 }
1707
1708 async fn index_usage(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1709 let ref_ = args
1710 .get("ref")
1711 .and_then(|v| v.as_str())
1712 .ok_or_else(|| ToolError::InvalidArgs("ref is required".into()))?;
1713 let (schema, name) = parse_ref(ref_).map_err(ToolError::InvalidArgs)?;
1714 let client = self.session.checkout().await?;
1715 let rows = client.query(&sql::index_usage(&schema, &name), &[]).await?;
1716 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1717 }
1718
1719 async fn list_running_queries(&self) -> Result<ToolOutcome, ToolError> {
1720 let client = self.session.checkout().await?;
1721 let rows = client.query(sql::running_queries(), &[]).await?;
1722 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1723 }
1724
1725 async fn find_blocking_locks(&self) -> Result<ToolOutcome, ToolError> {
1726 let client = self.session.checkout().await?;
1727 let rows = client.query(sql::blocking_locks(), &[]).await?;
1728 let values = rows_to_json(&rows);
1729 if values.as_array().map(|a| a.is_empty()).unwrap_or(true) {
1730 return Ok(ToolOutcome::ok_json(json!({
1731 "message": "No blocking locks found.",
1732 "locks": [],
1733 })));
1734 }
1735 Ok(ToolOutcome::ok_json(values))
1736 }
1737
1738 async fn slow_queries(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1739 let limit = args
1740 .get("limit")
1741 .and_then(|v| v.as_u64())
1742 .map(|n| n as u32)
1743 .unwrap_or(SLOW_QUERIES_DEFAULT);
1744 let client = self.session.checkout().await?;
1745 match client.query(&sql::slow_queries(limit), &[]).await {
1746 Ok(rows) => Ok(ToolOutcome::ok_json(rows_to_json(&rows))),
1747 Err(e) => {
1748 if let Some(message) = sql::map_stat_statements_error(&e) {
1749 Ok(ToolOutcome::ok_json(json!({
1750 "error": message,
1751 "hint": message,
1752 })))
1753 } else {
1754 Err(ToolError::Postgres(e))
1755 }
1756 }
1757 }
1758 }
1759
1760 async fn db_health_check(&self) -> Result<ToolOutcome, ToolError> {
1761 let client = self.session.checkout().await?;
1762 let sections: &[(&str, &str)] = &[
1763 ("overview", sql::database_stats()),
1764 ("cache", sql::cache_hit_ratio()),
1765 ("dead_tuples", sql::database_maintenance_stats()),
1766 ("connection_states", sql::connection_states()),
1767 ("blocking_locks", sql::blocking_locks()),
1768 ];
1769 let mut report = serde_json::Map::new();
1770 for (key, q) in sections {
1771 match client.query(*q, &[]).await {
1772 Ok(rows) => {
1773 report.insert((*key).into(), rows_to_json(&rows));
1774 }
1775 Err(e) => {
1776 report.insert((*key).into(), json!({ "error": e.to_string() }));
1777 }
1778 }
1779 }
1780 let lock_count = report
1781 .get("blocking_locks")
1782 .and_then(|v| v.as_array())
1783 .map(|a| a.len() as u64);
1784 report.insert("blocking_lock_count".into(), json!(lock_count));
1785 Ok(ToolOutcome::ok_json(Value::Object(report)))
1786 }
1787
1788 async fn explain_analyze(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1789 let sql = args
1790 .get("sql")
1791 .and_then(|v| v.as_str())
1792 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
1793 require_select_or_with(&self.session.filter(), sql)?;
1794 let explain = build_explain_sql(sql, true);
1795 self.run_explain_in_transaction(&explain).await
1796 }
1797
1798 async fn analyze_query_plan(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1799 let sql = args
1800 .get("sql")
1801 .and_then(|v| v.as_str())
1802 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
1803 require_select_or_with(&self.session.filter(), sql)?;
1804 let analyze = args
1805 .get("analyze")
1806 .and_then(|v| v.as_bool())
1807 .unwrap_or(false);
1808 let explain = build_explain_sql(sql, analyze);
1809 let outcome = self.run_explain_in_transaction(&explain).await?;
1810 let rows = outcome.structured.unwrap_or(Value::Null);
1811 let row_array = rows
1812 .get("rows")
1813 .and_then(|v| v.as_array())
1814 .or_else(|| rows.as_array());
1815 let plan = row_array
1816 .and_then(|a| a.first())
1817 .and_then(|r| r.get("QUERY PLAN"))
1818 .cloned()
1819 .unwrap_or(Value::Null);
1820 let metrics = extract_plan_metrics(&plan).or_else(|| extract_plan_metrics(&rows));
1821 let recommendations = metrics
1822 .as_ref()
1823 .and_then(|m| m.get("recommendations"))
1824 .cloned()
1825 .unwrap_or_else(|| json!([]));
1826 Ok(ToolOutcome::ok_json(json!({
1827 "metrics": metrics,
1828 "recommendations": recommendations,
1829 "plan": plan,
1830 })))
1831 }
1832
1833 async fn run_explain_in_transaction(
1835 &self,
1836 explain_sql: &str,
1837 ) -> Result<ToolOutcome, ToolError> {
1838 let client = self.session.checkout().await?;
1839 client
1840 .batch_execute("SET statement_timeout = '30s'")
1841 .await?;
1842 client.batch_execute("BEGIN").await?;
1843 let result = async {
1844 client.batch_execute("SET TRANSACTION READ ONLY").await?;
1845 let rows = client.query(explain_sql, &[]).await?;
1846 Ok::<_, ToolError>(rows_to_json(&rows))
1847 }
1848 .await;
1849 let _ = client.batch_execute("ROLLBACK").await;
1851 match result {
1852 Ok(values) => Ok(ToolOutcome::ok_json(values)),
1853 Err(e) => Err(e),
1854 }
1855 }
1856
1857 async fn get_index_status(&self) -> Result<ToolOutcome, ToolError> {
1858 let (store, connection_id, database) = self.index_service().await?;
1859 let base = store.base_dir(&connection_id, &database);
1860 let Some(manifest) = store.read_manifest(&base)? else {
1861 return Err(ToolError::Execution(format!(
1862 "No schema index for database \"{database}\" — run `nexql-mcp index build`."
1863 )));
1864 };
1865
1866 let mut live_fingerprint: Option<String> = None;
1867 let mut drift: Option<bool> = None;
1868 if let Ok(client) = self.session.checkout().await {
1869 let db = PgCatalogDb::new(&client);
1870 if let Ok(fp) = db.schema_fingerprint().await {
1871 drift = Some(fp != manifest.schema_fingerprint);
1872 live_fingerprint = Some(fp);
1873 }
1874 }
1875
1876 Ok(ToolOutcome::ok_json(json!({
1877 "connectionId": manifest.connection_id,
1878 "database": manifest.database,
1879 "indexedAt": manifest.indexed_at,
1880 "fingerprint": manifest.schema_fingerprint,
1881 "liveFingerprint": live_fingerprint,
1882 "drift": drift,
1883 "pgVersion": manifest.pg_version,
1884 "counts": {
1885 "tables": manifest.counts.tables,
1886 "views": manifest.counts.views,
1887 "functions": manifest.counts.functions,
1888 "enums": manifest.counts.enums,
1889 },
1890 "buildMs": manifest.stats.build_ms,
1891 "warnings": manifest.stats.warnings,
1892 })))
1893 }
1894
1895 async fn list_extensions(&self) -> Result<ToolOutcome, ToolError> {
1896 let client = self.session.checkout().await?;
1897 let rows = client.query(sql::list_extensions(), &[]).await?;
1898 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1899 }
1900
1901 async fn server_settings(&self) -> Result<ToolOutcome, ToolError> {
1902 let client = self.session.checkout().await?;
1903 let rows = client.query(sql::server_settings(), &[]).await?;
1904 Ok(ToolOutcome::ok_json(rows_to_json(&rows)))
1905 }
1906
1907 async fn suggest_indexes(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
1908 let limit = args
1909 .get("limit")
1910 .and_then(|v| v.as_u64())
1911 .map(|n| n as u32)
1912 .unwrap_or(REPORT_LIMIT_DEFAULT);
1913 let client = self.session.checkout().await?;
1914 let mut query_errors = serde_json::Map::new();
1915
1916 let high_seq_json = match client
1917 .query(&sql::high_seq_scan_tables(limit), &[])
1918 .await
1919 {
1920 Ok(rows) => rows_to_json(&rows),
1921 Err(e) => {
1922 query_errors.insert(
1923 "high_seq_scan_tables".into(),
1924 json!(nexql_conn::format_postgres_error(&e)),
1925 );
1926 Value::Null
1927 }
1928 };
1929
1930 let unindexed_json = match client
1931 .query(&sql::unindexed_fk_columns(limit), &[])
1932 .await
1933 {
1934 Ok(rows) => rows_to_json(&rows),
1935 Err(e) => {
1936 query_errors.insert(
1937 "unindexed_fk_columns".into(),
1938 json!(nexql_conn::format_postgres_error(&e)),
1939 );
1940 Value::Null
1941 }
1942 };
1943
1944 let mut pg_stat_available = false;
1945 let mut slow_queries = Value::Null;
1946 let mut pg_stat_note: Option<String> = None;
1947 match client.query(&sql::slow_queries(limit.min(10)), &[]).await {
1948 Ok(rows) => {
1949 pg_stat_available = true;
1950 slow_queries = rows_to_json(&rows);
1951 }
1952 Err(e) => {
1953 if let Some(message) = sql::map_stat_statements_error(&e) {
1954 pg_stat_note = Some(message);
1955 } else {
1956 query_errors.insert(
1957 "slow_queries".into(),
1958 json!(nexql_conn::format_postgres_error(&e)),
1959 );
1960 }
1961 }
1962 }
1963
1964 let mut plan_heuristics = Value::Null;
1965 if let Some(sql_text) = args.get("sql").and_then(|v| v.as_str()) {
1966 require_select_or_with(&self.session.filter(), sql_text)?;
1967 let explain = build_explain_sql(sql_text, false);
1968 match self.run_explain_in_transaction(&explain).await {
1969 Ok(outcome) => {
1970 let rows = outcome.structured.unwrap_or(Value::Null);
1971 let plan = rows
1972 .as_array()
1973 .and_then(|a| a.first())
1974 .and_then(|r| r.get("QUERY PLAN"))
1975 .cloned()
1976 .unwrap_or(Value::Null);
1977 let metrics =
1978 extract_plan_metrics(&plan).or_else(|| extract_plan_metrics(&rows));
1979 plan_heuristics = json!({
1980 "metrics": metrics,
1981 "hint": "Use analyze_query_plan with analyze=true for actual timings before creating indexes.",
1982 });
1983 }
1984 Err(e) => {
1985 query_errors.insert("plan_heuristics".into(), json!(e.to_string()));
1986 }
1987 }
1988 }
1989
1990 let has_candidates = high_seq_json
1991 .as_array()
1992 .map(|a| !a.is_empty())
1993 .unwrap_or(false)
1994 || unindexed_json
1995 .as_array()
1996 .map(|a| !a.is_empty())
1997 .unwrap_or(false)
1998 || plan_heuristics != Value::Null;
1999
2000 let mut payload = if !has_candidates && !pg_stat_available {
2001 json!({
2002 "suggestions": [],
2003 "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.",
2004 "hint": pg_stat_note,
2005 })
2006 } else if !has_candidates {
2007 json!({
2008 "high_seq_scan_tables": high_seq_json,
2009 "unindexed_fk_columns": unindexed_json,
2010 "slow_queries": slow_queries,
2011 "plan_heuristics": plan_heuristics,
2012 "message": "No strong index candidates from sequential-scan or unindexed-FK heuristics. Review slow_queries / pass sql for plan-level advice.",
2013 "hint": "CREATE INDEX CONCURRENTLY after validating with EXPLAIN (ANALYZE, BUFFERS).",
2014 })
2015 } else {
2016 json!({
2017 "high_seq_scan_tables": high_seq_json,
2018 "unindexed_fk_columns": unindexed_json,
2019 "slow_queries": slow_queries,
2020 "plan_heuristics": plan_heuristics,
2021 "pg_stat_statements": pg_stat_available,
2022 "hint": pg_stat_note.unwrap_or_else(|| {
2023 "Validate candidates with analyze_query_plan / EXPLAIN before CREATE INDEX CONCURRENTLY.".into()
2024 }),
2025 })
2026 };
2027
2028 if !query_errors.is_empty() {
2029 payload["query_errors"] = Value::Object(query_errors);
2030 }
2031
2032 Ok(ToolOutcome::ok_json(payload))
2033 }
2034
2035 async fn find_unused_indexes(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2036 let limit = args
2037 .get("limit")
2038 .and_then(|v| v.as_u64())
2039 .map(|n| n as u32)
2040 .unwrap_or(REPORT_LIMIT_DEFAULT);
2041 let client = self.session.checkout().await?;
2042 let rows = client.query(&sql::find_unused_indexes(limit), &[]).await?;
2043 let indexes = rows_to_json(&rows);
2044 if indexes.as_array().map(|a| a.is_empty()).unwrap_or(true) {
2045 return Ok(ToolOutcome::ok_json(json!({
2046 "indexes": [],
2047 "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.",
2048 })));
2049 }
2050 Ok(ToolOutcome::ok_json(json!({
2051 "indexes": indexes,
2052 "hint": "Prefer DROP INDEX CONCURRENTLY after confirming the workload (and that stats are mature).",
2053 })))
2054 }
2055
2056 async fn bloat_report(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2057 let limit = args
2058 .get("limit")
2059 .and_then(|v| v.as_u64())
2060 .map(|n| n as u32)
2061 .unwrap_or(REPORT_LIMIT_DEFAULT);
2062 let client = self.session.checkout().await?;
2063 let rows = client.query(&sql::bloat_report(limit), &[]).await?;
2064 let tables = rows_to_json(&rows);
2065 if tables.as_array().map(|a| a.is_empty()).unwrap_or(true) {
2066 return Ok(ToolOutcome::ok_json(json!({
2067 "tables": [],
2068 "method": "dead_tuple_ratio",
2069 "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.",
2070 })));
2071 }
2072 Ok(ToolOutcome::ok_json(json!({
2073 "tables": tables,
2074 "method": "dead_tuple_ratio",
2075 "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.",
2076 "hint": "VACUUM ANALYZE on high bloat_pct tables; investigate autovacuum settings if last_autovacuum is stale.",
2077 })))
2078 }
2079
2080 async fn find_missing_fks(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2081 let limit = args
2082 .get("limit")
2083 .and_then(|v| v.as_u64())
2084 .map(|n| n as u32)
2085 .unwrap_or(REPORT_LIMIT_DEFAULT);
2086 let capped = limit.clamp(1, sql::REPORT_LIMIT_MAX) as usize;
2087
2088 if let Ok((store, connection_id, database)) = self.index_service().await {
2090 let base = store.base_dir(&connection_id, &database);
2091 if let Ok(Some(manifest)) = store.read_manifest(&base)
2092 && let Ok(Some(graph)) = store.read_join_graph(&base, &manifest)
2093 {
2094 let candidates: Vec<Value> = graph
2095 .edges
2096 .into_iter()
2097 .filter(|e| e.inferred == Some(true) && e.disabled != Some(true))
2098 .take(capped)
2099 .map(|e| {
2100 let cols: Vec<Value> = e
2101 .cols
2102 .iter()
2103 .map(|(a, b)| json!({ "from": a, "to": b }))
2104 .collect();
2105 json!({
2106 "from_table": e.from,
2107 "to_table": e.to,
2108 "via": e.via,
2109 "columns": cols,
2110 "detection": "join_graph_inferred",
2111 })
2112 })
2113 .collect();
2114 if !candidates.is_empty() {
2115 return Ok(ToolOutcome::ok_json(json!({
2116 "candidates": candidates,
2117 "source": "join_graph",
2118 "hint": "These edges were inferred by naming convention and have no declared FK. Review before ALTER TABLE … ADD FOREIGN KEY.",
2119 })));
2120 }
2121 }
2122 }
2123
2124 let client = self.session.checkout().await?;
2125 let rows = client
2126 .query(&sql::find_missing_fks_catalog(limit), &[])
2127 .await?;
2128 let candidates = rows_to_json(&rows);
2129 if candidates.as_array().map(|a| a.is_empty()).unwrap_or(true) {
2130 return Ok(ToolOutcome::ok_json(json!({
2131 "candidates": [],
2132 "source": "catalog",
2133 "message": "No missing FK candidates found via join-graph inferred edges or *_id naming against single-column PKs.",
2134 })));
2135 }
2136 Ok(ToolOutcome::ok_json(json!({
2137 "candidates": candidates,
2138 "source": "catalog",
2139 "hint": "Naming-inferred only — verify referential integrity and nullability before adding constraints. Run `nexql-mcp index build` for join-graph inferred edges.",
2140 })))
2141 }
2142
2143 async fn list_roles(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2144 let client = self.session.checkout().await?;
2145 let role = args
2146 .get("role")
2147 .and_then(|v| v.as_str())
2148 .map(str::trim)
2149 .filter(|s| !s.is_empty());
2150
2151 let Some(role_name) = role else {
2152 let rows = client.query(sql::list_roles(), &[]).await?;
2153 return Ok(ToolOutcome::ok_json(rows_to_json(&rows)));
2154 };
2155
2156 let details = client.query(sql::role_details(), &[&role_name]).await?;
2157 if details.is_empty() {
2158 return Err(ToolError::Execution(format!(
2159 "Role \"{role_name}\" not found"
2160 )));
2161 }
2162 let member_of = client.query(sql::role_member_of(), &[&role_name]).await?;
2163 let has_members = client.query(sql::role_has_members(), &[&role_name]).await?;
2164 let privileges = client
2165 .query(sql::role_table_privileges(), &[&role_name])
2166 .await?;
2167
2168 Ok(ToolOutcome::ok_json(json!({
2169 "role": rows_to_json(&details).as_array().and_then(|a| a.first().cloned()).unwrap_or(Value::Null),
2170 "member_of": rows_to_json(&member_of),
2171 "has_members": rows_to_json(&has_members),
2172 "table_privileges": rows_to_json(&privileges),
2173 })))
2174 }
2175
2176 async fn export_query(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2177 let sql = args
2178 .get("sql")
2179 .and_then(|v| v.as_str())
2180 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
2181 require_select_or_with(&self.session.filter(), sql)?;
2182
2183 let format = args
2184 .get("format")
2185 .and_then(|v| v.as_str())
2186 .map(|s| {
2187 ExportFormat::parse(s).ok_or_else(|| {
2188 ToolError::InvalidArgs(format!(
2189 "Unsupported format \"{s}\". Use csv, json, or sqlinsert."
2190 ))
2191 })
2192 })
2193 .transpose()?
2194 .unwrap_or(ExportFormat::Csv);
2195
2196 let table_target = match args.get("table").and_then(|v| v.as_str()) {
2197 Some(t) if !t.trim().is_empty() => Some(parse_ref(t).map_err(ToolError::InvalidArgs)?),
2198 _ => None,
2199 };
2200
2201 if format == ExportFormat::SqlInsert && table_target.is_none() {
2202 return Err(ToolError::InvalidArgs(
2203 "table (schema.name) is required when format=sqlinsert".into(),
2204 ));
2205 }
2206
2207 let max_rows = self.session.caps().max_rows;
2208 let outcome = self.run_select_internal(sql, Some(max_rows)).await?;
2209 if outcome.is_error {
2210 return Ok(outcome);
2211 }
2212
2213 let structured = outcome.structured.unwrap_or(Value::Null);
2214 let rows_val = structured
2215 .get("rows")
2216 .cloned()
2217 .or_else(|| structured.get("data").and_then(|d| d.get("rows").cloned()))
2218 .unwrap_or(Value::Array(vec![]));
2219 let rows = rows_val.as_array().cloned().unwrap_or_default();
2220 let columns = columns_from_rows(&rows);
2221 let truncated = structured
2222 .get("truncated")
2223 .and_then(|v| v.as_bool())
2224 .unwrap_or(false);
2225
2226 let payload = match format {
2227 ExportFormat::Json => json!({
2228 "format": format.as_str(),
2229 "rowCount": rows.len(),
2230 "truncated": truncated,
2231 "columns": columns,
2232 "rows": rows,
2233 }),
2234 ExportFormat::Csv => {
2235 let content = rows_to_csv(&rows, &columns);
2236 let caps = self.session.caps();
2237 let (char_trunc, content) = caps.truncate_chars(&content);
2238 json!({
2239 "format": format.as_str(),
2240 "rowCount": rows.len(),
2241 "truncated": truncated || char_trunc,
2242 "columns": columns,
2243 "content": content,
2244 })
2245 }
2246 ExportFormat::SqlInsert => {
2247 let (schema, table) = table_target.expect("checked above");
2248 let content = rows_to_sql_insert(&rows, &columns, &schema, &table);
2249 let caps = self.session.caps();
2250 let (char_trunc, content) = caps.truncate_chars(&content);
2251 json!({
2252 "format": format.as_str(),
2253 "rowCount": rows.len(),
2254 "truncated": truncated || char_trunc,
2255 "table": format!("{schema}.{table}"),
2256 "columns": columns,
2257 "content": content,
2258 })
2259 }
2260 };
2261
2262 Ok(ToolOutcome::ok_json(payload))
2263 }
2264
2265 async fn db_dashboard(&self) -> Result<ToolOutcome, ToolError> {
2266 let client = self.session.checkout().await?;
2267 let sections: &[(&str, &str)] = &[
2268 ("db_info", sql::dashboard_db_info()),
2269 ("connection_states", sql::connection_states()),
2270 ("top_tables", sql::dashboard_top_tables()),
2271 ("object_counts", sql::dashboard_object_counts()),
2272 ("active_queries", sql::dashboard_active_queries()),
2273 ("blocking_locks", sql::blocking_locks()),
2274 ("max_connections", sql::dashboard_max_connections()),
2275 ("extension_count", sql::dashboard_extension_count()),
2276 ("cache", sql::cache_hit_ratio()),
2277 ];
2278 let mut report = serde_json::Map::new();
2279 for (key, q) in sections {
2280 match client.query(*q, &[]).await {
2281 Ok(rows) => {
2282 report.insert((*key).into(), rows_to_json(&rows));
2283 }
2284 Err(e) => {
2285 report.insert((*key).into(), json!({ "error": e.to_string() }));
2286 }
2287 }
2288 }
2289
2290 for key in ["db_info", "object_counts", "extension_count", "cache"] {
2292 if let Some(Value::Array(arr)) = report.get(key).cloned()
2293 && arr.len() == 1
2294 {
2295 report.insert(key.into(), arr.into_iter().next().unwrap());
2296 }
2297 }
2298 if let Some(Value::Array(arr)) = report.get("max_connections").cloned()
2299 && let Some(row) = arr.first()
2300 {
2301 report.insert(
2302 "max_connections".into(),
2303 row.get("max_connections")
2304 .cloned()
2305 .unwrap_or_else(|| row.clone()),
2306 );
2307 }
2308
2309 Ok(ToolOutcome::ok_json(Value::Object(report)))
2310 }
2311
2312 async fn deep_plan_analysis(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2313 let sql = args
2314 .get("sql")
2315 .and_then(|v| v.as_str())
2316 .ok_or_else(|| ToolError::InvalidArgs("sql is required".into()))?;
2317 require_select_or_with(&self.session.filter(), sql)?;
2318 let analyze = args
2319 .get("analyze")
2320 .and_then(|v| v.as_bool())
2321 .unwrap_or(true);
2322 let explain = build_explain_sql(sql, analyze);
2323 let outcome = self.run_explain_in_transaction(&explain).await?;
2324 let rows = outcome.structured.unwrap_or(Value::Null);
2325 let row_array = rows
2326 .get("rows")
2327 .and_then(|v| v.as_array())
2328 .or_else(|| rows.as_array());
2329 let plan = row_array
2330 .and_then(|a| a.first())
2331 .and_then(|r| r.get("QUERY PLAN"))
2332 .cloned()
2333 .unwrap_or(Value::Null);
2334 let deep = analyze_deep_plan(&plan, sql)
2335 .or_else(|| analyze_deep_plan(&rows, sql))
2336 .ok_or_else(|| {
2337 ToolError::Execution("Could not parse EXPLAIN JSON plan for deep analysis".into())
2338 })?;
2339 let metrics = extract_plan_metrics(&plan).or_else(|| extract_plan_metrics(&rows));
2340 Ok(ToolOutcome::ok_json(json!({
2341 "deep": deep,
2342 "metrics": metrics,
2343 "plan": plan,
2344 "analyzed": analyze,
2345 })))
2346 }
2347
2348 async fn schema_diff(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2349 let source_schema = args
2350 .get("sourceSchema")
2351 .and_then(|v| v.as_str())
2352 .ok_or_else(|| ToolError::InvalidArgs("sourceSchema is required".into()))?;
2353 let target_schema = args
2354 .get("targetSchema")
2355 .and_then(|v| v.as_str())
2356 .ok_or_else(|| ToolError::InvalidArgs("targetSchema is required".into()))?;
2357 crate::schema_diff::require_safe_schema(source_schema)?;
2358 crate::schema_diff::require_safe_schema(target_schema)?;
2359
2360 let client = self.session.checkout().await?;
2361 let source = crate::schema_diff::load_schema_snapshot(&client, source_schema).await?;
2362 let target = crate::schema_diff::load_schema_snapshot(&client, target_schema).await?;
2363 let diffs = crate::schema_diff::compute_schema_diff(&source, &target);
2364 let changed = diffs
2365 .iter()
2366 .filter(|d| d.status != crate::schema_diff::DiffStatus::Unchanged)
2367 .count();
2368 Ok(ToolOutcome::ok_json(json!({
2369 "sourceSchema": source_schema,
2370 "targetSchema": target_schema,
2371 "tableCount": diffs.len(),
2372 "changedCount": changed,
2373 "diffs": crate::schema_diff::diffs_to_json(&diffs),
2374 })))
2375 }
2376
2377 async fn generate_migration(&self, args: &Value) -> Result<ToolOutcome, ToolError> {
2378 let source_schema = args
2379 .get("sourceSchema")
2380 .and_then(|v| v.as_str())
2381 .ok_or_else(|| ToolError::InvalidArgs("sourceSchema is required".into()))?;
2382 let target_schema = args
2383 .get("targetSchema")
2384 .and_then(|v| v.as_str())
2385 .ok_or_else(|| ToolError::InvalidArgs("targetSchema is required".into()))?;
2386 crate::schema_diff::require_safe_schema(source_schema)?;
2387 crate::schema_diff::require_safe_schema(target_schema)?;
2388
2389 let client = self.session.checkout().await?;
2390 let source = crate::schema_diff::load_schema_snapshot(&client, source_schema).await?;
2391 let target = crate::schema_diff::load_schema_snapshot(&client, target_schema).await?;
2392 let diffs = crate::schema_diff::compute_schema_diff(&source, &target);
2393 let statements =
2394 crate::schema_diff::build_migration_statements(source_schema, target_schema, &diffs);
2395 let sql = if statements.is_empty() {
2396 format!("-- No differences between {source_schema} and {target_schema}")
2397 } else {
2398 statements.join("\n\n")
2399 };
2400 Ok(ToolOutcome::ok_json(json!({
2401 "sourceSchema": source_schema,
2402 "targetSchema": target_schema,
2403 "statementCount": statements.len(),
2404 "sql": sql,
2405 "hint": "Read-only: review and run via execute_sql / apply_ddl only with --access-mode write|admin. Destructive drops are commented out.",
2406 })))
2407 }
2408
2409 async fn run_select_internal(
2410 &self,
2411 sql: &str,
2412 max_rows: Option<u32>,
2413 ) -> Result<ToolOutcome, ToolError> {
2414 let client = self.session.checkout().await?;
2415 let Some(max_rows) = max_rows else {
2416 let rows = client.query(sql, &[]).await?;
2417 let values = rows_to_json(&rows);
2418 let payload = self.apply_pii_redaction(sql, ensure_structured_object(values));
2419 let text = serde_json::to_string_pretty(&payload)
2420 .map_err(|e| ToolError::Execution(e.to_string()))?;
2421 let caps = self.session.caps();
2422 let (trunc, text) = caps.truncate_chars(&text);
2423 let structured = if trunc {
2424 json!({ "truncated_chars": true, "data": payload })
2425 } else {
2426 payload
2427 };
2428 return Ok(ToolOutcome {
2429 text: text.to_string(),
2430 structured: Some(structured),
2431 is_error: false,
2432 });
2433 };
2434
2435 let cleaned = sql.trim().trim_end_matches(';').trim();
2436 let wrapped = format!(
2437 "SELECT * FROM ({cleaned}) AS nexql_limited LIMIT {}",
2438 max_rows + 1
2439 );
2440 let rows = client.query(&wrapped, &[]).await.map_err(|e| {
2441 ToolError::Execution(format!(
2442 "Failed to execute row-limited query (refusing unbounded fallback): {}",
2443 nexql_conn::format_postgres_error(&e)
2444 ))
2445 })?;
2446 let truncated = rows.len() as u32 > max_rows;
2447 let keep = if truncated {
2448 &rows[..max_rows as usize]
2449 } else {
2450 &rows[..]
2451 };
2452 let values = rows_to_json(keep);
2453 let mut payload = self.apply_pii_redaction(sql, ensure_structured_object(values));
2455 if truncated
2456 && let Some(obj) = payload.as_object_mut()
2457 {
2458 obj.insert("truncated".into(), json!(true));
2459 obj.insert("maxRows".into(), json!(max_rows));
2460 }
2461 let text = serde_json::to_string_pretty(&payload)
2462 .map_err(|e| ToolError::Execution(e.to_string()))?;
2463 let caps = self.session.caps();
2464 let (char_trunc, text) = caps.truncate_chars(&text);
2465 let structured = if char_trunc {
2466 json!({ "truncated_chars": true, "data": payload })
2467 } else {
2468 payload
2469 };
2470 Ok(ToolOutcome {
2471 text: text.to_string(),
2472 structured: Some(structured),
2473 is_error: false,
2474 })
2475 }
2476
2477 fn apply_pii_redaction(&self, sql: &str, mut payload: Value) -> Value {
2478 let filter = self.session.filter();
2479 if filter.pii_columns.is_empty() {
2480 return payload;
2481 }
2482 let Ok(tables) = select_table_refs(sql) else {
2483 return payload;
2484 };
2485 let (redacted, cols) = redact_pii_in_payload(payload, &filter.pii_columns, &tables);
2486 payload = redacted;
2487 if !cols.is_empty()
2488 && let Some(obj) = payload.as_object_mut()
2489 {
2490 obj.insert("piiRedactedColumns".into(), json!(cols));
2491 }
2492 payload
2493 }
2494}
2495
2496fn normalize_for_match(s: &str) -> String {
2498 let mut out = String::new();
2499 let mut last_was_sep = true;
2500 for ch in s.to_lowercase().chars() {
2501 if ch.is_ascii_alphanumeric() {
2502 out.push(ch);
2503 last_was_sep = false;
2504 } else if !last_was_sep {
2505 out.push(' ');
2506 last_was_sep = true;
2507 }
2508 }
2509 out.trim().to_string()
2510}
2511
2512fn fuzzy_score(hint: &str, candidate: &str) -> f64 {
2514 let h = normalize_for_match(hint);
2515 let c = normalize_for_match(candidate);
2516 if h.is_empty() || c.is_empty() {
2517 return 0.0;
2518 }
2519 if h == c {
2520 return 100.0;
2521 }
2522 if c.contains(&h) || h.contains(&c) {
2523 return 75.0;
2524 }
2525 let h_tokens: std::collections::HashSet<&str> =
2526 h.split(' ').filter(|s| !s.is_empty()).collect();
2527 let c_tokens: std::collections::HashSet<&str> =
2528 c.split(' ').filter(|s| !s.is_empty()).collect();
2529 let overlap = h_tokens.intersection(&c_tokens).count();
2530 if overlap == 0 {
2531 return 0.0;
2532 }
2533 (overlap as f64 / h_tokens.len().max(c_tokens.len()) as f64) * 60.0
2534}
2535
2536fn policy_to_query_filter(filter: &PolicyFilter) -> QueryPolicyFilter {
2537 QueryPolicyFilter {
2538 allow_schemas: filter.allow_schemas.clone(),
2539 deny_schemas: filter.deny_schemas.clone(),
2540 deny_tables: filter.deny_tables.clone(),
2541 pii_columns: filter.pii_columns.clone(),
2542 }
2543}
2544
2545fn require_select_or_with(filter: &PolicyFilter, sql: &str) -> Result<(), ToolError> {
2546 match validate_readonly_sql(sql)? {
2547 SqlDecision::Allow => {}
2548 SqlDecision::Reject => {
2549 return Err(ToolError::Execution(
2550 "Security Error: Only SELECT or WITH statements can be analyzed.".into(),
2551 ));
2552 }
2553 }
2554 enforce_read_table_policy(filter, sql)?;
2555 let trimmed = sql.trim().to_ascii_lowercase();
2556 if !(trimmed.starts_with("select") || trimmed.starts_with("with")) {
2557 return Err(ToolError::Execution(
2558 "Security Error: Only SELECT or WITH statements can be analyzed.".into(),
2559 ));
2560 }
2561 Ok(())
2562}
2563
2564fn rows_to_json(rows: &[tokio_postgres::Row]) -> Value {
2565 rows_to_json_array(rows)
2566}
2567
2568fn scores_equal(a: f64, b: f64) -> bool {
2569 (a - b).abs() <= f64::EPSILON * a.abs().max(b.abs()).max(1.0)
2570}
2571
2572fn read_recent_log_errors() -> Vec<String> {
2573 let path = std::env::var("NEXQL_MCP_LOG")
2574 .map(std::path::PathBuf::from)
2575 .ok()
2576 .or_else(|| {
2577 std::env::var_os("HOME").map(|h| {
2578 std::path::PathBuf::from(h)
2579 .join(".config")
2580 .join("nexql-mcp")
2581 .join("logs")
2582 .join("nexql-mcp.log")
2583 })
2584 });
2585
2586 let Some(log_path) = path else {
2587 return Vec::new();
2588 };
2589
2590 let Ok(content) = std::fs::read_to_string(&log_path) else {
2591 return Vec::new();
2592 };
2593
2594 content
2595 .lines()
2596 .rev()
2597 .take(50)
2598 .filter(|line| {
2599 line.contains("ERROR")
2600 || line.contains("WARN")
2601 || line.contains("failed")
2602 || line.contains("Error")
2603 })
2604 .map(String::from)
2605 .collect()
2606}
2607
2608#[cfg(test)]
2609mod tests {
2610 use super::*;
2611 use crate::plan::build_explain_sql;
2612 use nexql_policy::PolicyFilter;
2613 use serde_json::json;
2614
2615 use crate::session::{ConnectionInfo, ConnectionPolicy, ToolSession};
2616 use nexql_policy::{AccessMode, PolicyCaps};
2617
2618 fn test_conn() -> ConnectionInfo {
2619 ConnectionInfo {
2620 id: "conn-1".into(),
2621 name: "conn-1".into(),
2622 host: Some("127.0.0.1".into()),
2623 port: Some(5432),
2624 database: Some("appdb".into()),
2625 params: Default::default(),
2626 policy: ConnectionPolicy {
2627 access_mode: AccessMode::Read,
2628 caps: PolicyCaps::default(),
2629 filter: PolicyFilter::default(),
2630 environment: None,
2631 },
2632 }
2633 }
2634
2635 #[test]
2636 fn scores_equal_treats_near_duplicates_as_tied() {
2637 let s = 3.295836866004329_f64;
2638 assert!(super::scores_equal(s, s));
2639 assert!(super::scores_equal(s, s + f64::EPSILON));
2640 }
2641
2642 #[test]
2643 fn policy_maps_one_to_one() {
2644 let f = PolicyFilter {
2645 allow_schemas: vec!["public".into()],
2646 deny_schemas: vec!["pgboss".into()],
2647 deny_tables: vec!["auth.*".into()],
2648 pii_columns: vec!["public.users.ssn".into()],
2649 };
2650 let q = policy_to_query_filter(&f);
2651 assert_eq!(q.allow_schemas, f.allow_schemas);
2652 assert_eq!(q.deny_schemas, f.deny_schemas);
2653 assert_eq!(q.deny_tables, f.deny_tables);
2654 assert_eq!(q.pii_columns, f.pii_columns);
2655 }
2656
2657 #[test]
2658 fn ok_json_wraps_arrays_for_cursor_structured_content() {
2659 let out = ToolOutcome::ok_json(json!([{ "id": 1 }, { "id": 2 }]));
2660 assert!(!out.is_error);
2661 let s = out.structured.as_ref().unwrap();
2662 assert!(s.is_object(), "structuredContent must be object, got {s}");
2663 assert_eq!(s["rows"].as_array().unwrap().len(), 2);
2664 assert!(out.text.contains("\"rows\""));
2665 }
2666
2667 #[test]
2668 fn ok_json_leaves_objects_unchanged() {
2669 let out = ToolOutcome::ok_json(json!({ "kind": "table", "name": "orders" }));
2670 let s = out.structured.as_ref().unwrap();
2671 assert_eq!(s["kind"], "table");
2672 assert!(s.get("rows").is_none());
2673 }
2674
2675 #[test]
2676 fn router_specs_include_phase4_and_phase9() {
2677 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2678 let router = ToolRouter::with_index_store(session, None);
2679 assert_eq!(router.specs().len(), ToolName::ACTIVE.len());
2680 let names: Vec<_> = router.specs().iter().map(|s| s.name.as_str()).collect();
2681 assert!(names.contains(&"search_schema"));
2682 assert!(names.contains(&"get_ddl"));
2683 assert!(names.contains(&"explain_analyze"));
2684 assert!(names.contains(&"get_index_status"));
2685 assert!(names.contains(&"list_extensions"));
2686 assert!(names.contains(&"server_settings"));
2687 assert!(names.contains(&"suggest_indexes"));
2688 assert!(names.contains(&"find_unused_indexes"));
2689 assert!(names.contains(&"bloat_report"));
2690 assert!(names.contains(&"find_missing_fks"));
2691 assert!(names.contains(&"export_query"));
2692 assert!(names.contains(&"list_roles"));
2693 assert!(names.contains(&"db_dashboard"));
2694 assert!(names.contains(&"deep_plan_analysis"));
2695 assert!(names.contains(&"execute_sql"));
2696 assert!(names.contains(&"edit_row"));
2697 assert!(names.contains(&"import_data"));
2698 assert!(names.contains(&"apply_ddl"));
2699 assert!(names.contains(&"create_index_concurrently"));
2700 assert!(names.contains(&"run_maintenance"));
2701 assert!(names.contains(&"terminate_query"));
2702 }
2703
2704 #[tokio::test]
2705 async fn write_tools_refuse_read_mode() {
2706 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2707 let router = ToolRouter::with_index_store(session, None);
2708 for tool in [
2709 "execute_sql",
2710 "edit_row",
2711 "import_data",
2712 "apply_ddl",
2713 "create_index_concurrently",
2714 "run_maintenance",
2715 "terminate_query",
2716 ] {
2717 let out = router
2718 .call(tool, json!({ "sql": "SELECT 1", "table": "public.t", "rows": [], "action": "insert", "values": {}, "pid": 1 }))
2719 .await;
2720 assert!(out.is_error, "{tool}: {}", out.text);
2721 assert!(
2722 out.text.contains("write") || out.text.contains("admin"),
2723 "{tool}: {}",
2724 out.text
2725 );
2726 }
2727 }
2728
2729 #[tokio::test]
2730 async fn table_stats_rejects_injection_ref() {
2731 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2732 let router = ToolRouter::with_index_store(session, None);
2733 let out = router
2734 .call("table_stats", json!({ "ref": "public.users; DROP" }))
2735 .await;
2736 assert!(out.is_error, "{}", out.text);
2737 assert!(
2738 out.text.contains("Invalid object reference") || out.text.contains("invalid arguments"),
2739 "expected ref validation error, got: {}",
2740 out.text
2741 );
2742 }
2743
2744 #[test]
2745 fn explain_transaction_path_builds_readonly_sequence() {
2746 let explain = build_explain_sql("SELECT 1", true);
2748 assert!(explain.starts_with("EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON)"));
2749 assert!(!explain.to_ascii_lowercase().contains("commit"));
2750 let steps = ["BEGIN", "SET TRANSACTION READ ONLY", &explain, "ROLLBACK"];
2751 assert_eq!(steps.len(), 4);
2752 assert_eq!(steps[0], "BEGIN");
2753 assert_eq!(steps[1], "SET TRANSACTION READ ONLY");
2754 assert_eq!(steps[3], "ROLLBACK");
2755 }
2756
2757 #[tokio::test]
2758 async fn missing_index_returns_actionable_error() {
2759 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2760 let router = ToolRouter::with_index_store(session, None);
2761 let out = router
2762 .call("search_schema", json!({ "query": "users" }))
2763 .await;
2764 assert!(out.is_error, "{}", out.text);
2765 assert!(
2766 out.text.contains("rebuild_index"),
2767 "expected actionable hint, got: {}",
2768 out.text
2769 );
2770 }
2771
2772 #[tokio::test]
2773 async fn empty_index_dir_returns_build_hint() {
2774 let tmp = tempfile::TempDir::new().unwrap();
2775 let store = IndexStore::new(tmp.path());
2776 let session = ToolSession::for_tests(
2777 vec![test_conn()],
2778 PolicyFilter::default(),
2779 Some(IndexStore::new(tmp.path())),
2780 );
2781 let router = ToolRouter::with_index_store(session, Some(store));
2782 let out = router
2783 .call("describe_object", json!({ "ref": "public.users" }))
2784 .await;
2785 assert!(out.is_error, "{}", out.text);
2786 assert!(
2787 out.text.contains("rebuild_index"),
2788 "expected build hint, got: {}",
2789 out.text
2790 );
2791 }
2792
2793 #[tokio::test]
2794 async fn outcome_tagged_with_connection_id_and_database() {
2795 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2796 let router = ToolRouter::new(session);
2797 let out = router.call("list_connections", json!({})).await;
2798 let structured = out.structured.expect("structured outcome");
2799 assert_eq!(
2800 structured.get("connectionId").and_then(|v| v.as_str()),
2801 Some("conn-1")
2802 );
2803 assert_eq!(
2804 structured.get("database").and_then(|v| v.as_str()),
2805 Some("appdb")
2806 );
2807 }
2808
2809 #[tokio::test]
2810 async fn setup_connection_returns_needs_input_when_incomplete() {
2811 unsafe {
2812 std::env::remove_var("DATABASE_URL");
2813 std::env::remove_var("POSTGRES_URL");
2814 std::env::remove_var("PGHOST");
2815 }
2816 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2817 let router = ToolRouter::new(session);
2818 let out = router.call("setup_connection", json!({})).await;
2819 let structured = out.structured.expect("structured outcome");
2820 assert!(structured.get("status").is_some());
2821 }
2822
2823 #[tokio::test]
2824 async fn save_profile_persists_config() {
2825 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2826 let router = ToolRouter::new(session);
2827 let temp_dir = tempfile::tempdir().unwrap();
2828 let cfg_path = temp_dir.path().join("config.toml");
2829 unsafe {
2830 std::env::set_var("NEXQL_MCP_CONFIG", &cfg_path);
2831 }
2832
2833 let out = router
2834 .call(
2835 "save_profile",
2836 json!({
2837 "name": "staging",
2838 "host": "127.0.0.1",
2839 "port": 5432,
2840 "dbname": "stage_db",
2841 "user": "stage_user"
2842 }),
2843 )
2844 .await;
2845
2846 let structured = out.structured.expect("structured outcome");
2847 assert_eq!(
2848 structured.get("status").and_then(|v| v.as_str()),
2849 Some("saved")
2850 );
2851 assert_eq!(
2852 structured.get("profile").and_then(|v| v.as_str()),
2853 Some("staging")
2854 );
2855 }
2856
2857 #[tokio::test]
2858 async fn check_ddl_safety_tool_dispatches_ast_report() {
2859 let session = ToolSession::for_tests(vec![test_conn()], PolicyFilter::default(), None);
2860 let router = ToolRouter::new(session);
2861 let out = router
2862 .call(
2863 "check_ddl_safety",
2864 json!({ "ddl": "CREATE INDEX idx_col ON users(col);" }),
2865 )
2866 .await;
2867 let structured = out.structured.expect("structured outcome");
2868 assert_eq!(
2869 structured.get("overall_risk").and_then(|v| v.as_str()),
2870 Some("CRITICAL")
2871 );
2872 }
2873}