Skip to main content

nexql_tools/
exec.rs

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