Files
soft_usite/pkg/services/query_runner_service.go
T
Lizandro GuarnizoandCopilot 2779c0a15d fix: ping real de BD + campo Host en ConxDb + reactividad Alpine
Problema raíz:
1. El ping TCP usaba Servidor.IpServidor → fallaba si la BD solo escucha
   en localhost del servidor remoto (comportamiento normal en producción)
2. La reactividad Alpine v3.13 no detectaba keys nuevas en objeto vacío {}
   dentro de x-for loops, spinner nunca aparecía

Cambios:
- pkg/models/conx_db.go: nuevo campo Host (vacío = usa Servidor.IpServidor)
  + método HostEfectivo() como fuente única de verdad para host
- pkg/services/query_runner_service.go: todos los puntos (openDynamicDB,
  openDynamicDBWithName, redisConnect, mongoURI) usan c.HostEfectivo()
- rest/controllers/servidor_controller.go: PingConexion reescrito —
  primero prueba conexión real a la BD (no solo TCP), si falla prueba TCP,
  devuelve host+puerto testeado y mensajes de error descriptivos
- servidor_dashboard.html: hacerPing usa spread-replace para forzar
  reactividad + $nextTick, resultado muestra host:puerto testeado y
  mensaje de error completo en múltiples líneas

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
2026-05-21 11:56:12 -05:00

1190 lines
33 KiB
Go

package services
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"regexp"
"strconv"
"strings"
"time"
_ "github.com/go-sql-driver/mysql"
_ "github.com/lib/pq"
_ "github.com/mattn/go-sqlite3"
_ "github.com/microsoft/go-mssqldb"
goredis "github.com/redis/go-redis/v9"
"go.mongodb.org/mongo-driver/bson"
"go.mongodb.org/mongo-driver/mongo"
"go.mongodb.org/mongo-driver/mongo/options"
"github.com/sujit-baniya/fiber-boilerplate/pkg/models"
)
// connTimeout es el tiempo máximo para establecer una conexión a la DB externa.
const connTimeout = 8 * time.Second
// QueryResult contiene el resultado de una consulta SQL.
type QueryResult struct {
Columns []string `json:"columns"`
Rows []map[string]any `json:"rows"`
RowCount int `json:"row_count"`
AffectedRows int64 `json:"affected_rows"`
DurationMs int64 `json:"duration_ms"`
IsSelect bool `json:"is_select"`
Error string `json:"error,omitempty"`
}
// openDynamicDB abre una conexión a la base de datos indicada por ConxDb.
func openDynamicDB(c models.ConxDb) (*sql.DB, error) {
driver := strings.ToLower(c.TipoDb.Nombre)
host := c.HostEfectivo()
port := c.Puerto
user := c.Usuario
pass := c.Password
var dsn string
var driverName string
switch {
case strings.Contains(driver, "postgres"):
driverName = "postgres"
dsn = fmt.Sprintf("host=%s port=%s user=%s password=%s sslmode=disable connect_timeout=8", host, port, user, pass)
case strings.Contains(driver, "mysql") || strings.Contains(driver, "mariadb"):
driverName = "mysql"
dsn = fmt.Sprintf("%s:%s@tcp(%s:%s)/?timeout=8s&readTimeout=20s&writeTimeout=20s", user, pass, host, port)
case strings.Contains(driver, "sqlite"):
driverName = "sqlite3"
dsn = host // para sqlite el host es la ruta del archivo
case strings.Contains(driver, "sqlserver") || strings.Contains(driver, "mssql"):
driverName = "sqlserver"
dsn = fmt.Sprintf("sqlserver://%s:%s@%s:%s?dial+timeout=8", user, pass, host, port)
default:
return nil, fmt.Errorf("driver no soportado: %s", driver)
}
db, err := sql.Open(driverName, dsn)
if err != nil {
return nil, err
}
db.SetConnMaxLifetime(30 * time.Second)
db.SetMaxOpenConns(2)
// Verificar conectividad inmediatamente para fallar rápido en lugar de bloquear al hacer la primera query
ctx, cancel := context.WithTimeout(context.Background(), connTimeout)
defer cancel()
if err := db.PingContext(ctx); err != nil {
db.Close()
return nil, fmt.Errorf("no se pudo conectar al servidor de base de datos: %w", err)
}
return db, nil
}
// ExecuteSQL ejecuta SQL arbitrario contra la conexión y devuelve QueryResult.
// También guarda en query_history.
func ExecuteSQL(conx models.ConxDb, database, sqlText string) QueryResult {
if isMongoDriver(strings.ToLower(conx.TipoDb.Nombre)) {
return mongoExecuteSQL(conx, database, sqlText)
}
if isRedisDriver(strings.ToLower(conx.TipoDb.Nombre)) {
return redisExecuteCommand(conx, database, sqlText)
}
start := time.Now()
db, err := openDynamicDB(conx)
if err != nil {
saveHistory(conx.ID, sqlText, "error", err.Error(), 0, time.Since(start).Milliseconds())
return QueryResult{Error: err.Error()}
}
defer db.Close()
// Si se especifica una base de datos para seleccionar
if database != "" {
driver := strings.ToLower(conx.TipoDb.Nombre)
if strings.Contains(driver, "postgres") {
// En postgres se cambia con SET search_path o reconectando con dbname en DSN
db2, err2 := openDynamicDBWithName(conx, database)
if err2 == nil {
db.Close()
db = db2
}
} else {
if _, err2 := db.Exec("USE " + quoteIdentifier(database, conx.TipoDb.Nombre)); err2 != nil {
saveHistory(conx.ID, sqlText, "error", err2.Error(), 0, time.Since(start).Milliseconds())
return QueryResult{Error: err2.Error()}
}
}
}
trimmed := strings.TrimSpace(sqlText)
isSelect := isSelectStatement(trimmed)
var result QueryResult
result.IsSelect = isSelect
if isSelect {
rows, err := db.Query(trimmed)
if err != nil {
elapsed := time.Since(start).Milliseconds()
saveHistory(conx.ID, sqlText, "error", err.Error(), 0, elapsed)
return QueryResult{Error: err.Error(), IsSelect: true}
}
defer rows.Close()
cols, _ := rows.Columns()
result.Columns = cols
for rows.Next() {
vals := make([]any, len(cols))
ptrs := make([]any, len(cols))
for i := range vals {
ptrs[i] = &vals[i]
}
rows.Scan(ptrs...)
row := make(map[string]any, len(cols))
for i, col := range cols {
v := vals[i]
if b, ok := v.([]byte); ok {
row[col] = string(b)
} else {
row[col] = v
}
}
result.Rows = append(result.Rows, row)
}
result.RowCount = len(result.Rows)
} else {
driver := strings.ToLower(conx.TipoDb.Nombre)
// PostgreSQL: cuando se envían múltiples sentencias en un solo Exec, el servidor
// las envuelve en una transacción implícita. CREATE DATABASE/DROP DATABASE no
// pueden correr dentro de una transacción, así que dividimos y ejecutamos c/u por separado.
if strings.Contains(driver, "postgres") {
stmts := splitPostgresStatements(trimmed)
if len(stmts) > 1 {
var totalAffected int64
for _, stmt := range stmts {
res, execErr := db.Exec(stmt)
if execErr != nil {
elapsed := time.Since(start).Milliseconds()
preview := stmt
if len(preview) > 80 {
preview = preview[:80] + "..."
}
saveHistory(conx.ID, sqlText, "error", execErr.Error(), 0, elapsed)
return QueryResult{Error: fmt.Sprintf("[%s]: %s", preview, execErr.Error())}
}
if affected, err2 := res.RowsAffected(); err2 == nil {
totalAffected += affected
}
}
elapsed := time.Since(start).Milliseconds()
result.AffectedRows = totalAffected
result.DurationMs = elapsed
saveHistory(conx.ID, sqlText, "ok", "", totalAffected, elapsed)
return result
}
}
res, err := db.Exec(trimmed)
elapsed := time.Since(start).Milliseconds()
if err != nil {
saveHistory(conx.ID, sqlText, "error", err.Error(), 0, elapsed)
return QueryResult{Error: err.Error()}
}
affected, _ := res.RowsAffected()
result.AffectedRows = affected
result.DurationMs = elapsed
saveHistory(conx.ID, sqlText, "ok", "", affected, elapsed)
return result
}
result.DurationMs = time.Since(start).Milliseconds()
saveHistory(conx.ID, sqlText, "ok", "", int64(result.RowCount), result.DurationMs)
return result
}
// ListDatabases devuelve la lista de bases de datos del servidor.
func ListDatabases(conx models.ConxDb) ([]string, error) {
driver := strings.ToLower(conx.TipoDb.Nombre)
if isMongoDriver(driver) {
return mongoListDatabases(conx)
}
if isRedisDriver(driver) {
return redisListDatabases(conx)
}
db, err := openDynamicDB(conx)
if err != nil {
return nil, err
}
defer db.Close()
var query string
switch {
case strings.Contains(driver, "postgres"):
query = "SELECT datname FROM pg_database WHERE datistemplate = false ORDER BY datname"
case strings.Contains(driver, "mysql") || strings.Contains(driver, "mariadb"):
query = "SHOW DATABASES"
case strings.Contains(driver, "sqlserver") || strings.Contains(driver, "mssql"):
query = "SELECT name FROM sys.databases ORDER BY name"
default:
return []string{"main"}, nil
}
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
rows, err := db.QueryContext(ctx, query)
if err != nil {
return nil, err
}
defer rows.Close()
var dbs []string
for rows.Next() {
var name string
rows.Scan(&name)
dbs = append(dbs, name)
}
return dbs, nil
}
// ListTables devuelve las tablas de una base de datos.
func ListTables(conx models.ConxDb, database string) ([]string, error) {
driver := strings.ToLower(conx.TipoDb.Nombre)
if isMongoDriver(driver) {
return mongoListCollections(conx, database)
}
if isRedisDriver(driver) {
return redisListKeys(conx, database)
}
var db *sql.DB
var err error
if strings.Contains(driver, "postgres") {
db, err = openDynamicDBWithName(conx, database)
} else {
db, err = openDynamicDB(conx)
}
if err != nil {
return nil, err
}
defer db.Close()
var query string
switch {
case strings.Contains(driver, "postgres"):
query = "SELECT tablename FROM pg_tables WHERE schemaname='public' ORDER BY tablename"
case strings.Contains(driver, "mysql") || strings.Contains(driver, "mariadb"):
if _, err := db.Exec("USE " + quoteIdentifier(database, conx.TipoDb.Nombre)); err != nil {
return nil, err
}
query = "SHOW TABLES"
case strings.Contains(driver, "sqlserver") || strings.Contains(driver, "mssql"):
query = fmt.Sprintf("USE [%s]; SELECT TABLE_NAME FROM INFORMATION_SCHEMA.TABLES WHERE TABLE_TYPE='BASE TABLE' ORDER BY TABLE_NAME", database)
default:
query = "SELECT name FROM sqlite_master WHERE type='table' ORDER BY name"
}
rows, err := db.Query(query)
if err != nil {
return nil, err
}
defer rows.Close()
var tables []string
for rows.Next() {
var name string
rows.Scan(&name)
tables = append(tables, name)
}
return tables, nil
}
// TestConnection verifica si la conexión es válida.
func TestDBConnection(conx models.ConxDb) error {
if isMongoDriver(strings.ToLower(conx.TipoDb.Nombre)) {
return mongoTestConnection(conx)
}
if isRedisDriver(strings.ToLower(conx.TipoDb.Nombre)) {
return redisTestConnection(conx)
}
db, err := openDynamicDB(conx)
if err != nil {
return err
}
defer db.Close()
return db.Ping()
}
// ── helpers ──────────────────────────────────────────────────────────────────
func openDynamicDBWithName(c models.ConxDb, dbName string) (*sql.DB, error) {
host := c.HostEfectivo()
port := c.Puerto
user := c.Usuario
pass := c.Password
driver := strings.ToLower(c.TipoDb.Nombre)
var dsn string
switch {
case strings.Contains(driver, "postgres"):
dsn = fmt.Sprintf("host=%s port=%s user=%s password=%s dbname=%s sslmode=disable connect_timeout=8", host, port, user, pass, dbName)
db, err := sql.Open("postgres", dsn)
if err != nil {
return nil, err
}
db.SetConnMaxLifetime(30 * time.Second)
db.SetMaxOpenConns(2)
ctx, cancel := context.WithTimeout(context.Background(), connTimeout)
defer cancel()
if err := db.PingContext(ctx); err != nil {
db.Close()
return nil, fmt.Errorf("no se pudo conectar a la base de datos '%s': %w", dbName, err)
}
return db, nil
default:
return openDynamicDB(c)
}
}
// splitPostgresStatements divide SQL multi-sentencia en sentencias individuales
// respetando dollar-quoting ($$ ... $$), strings con comillas simples y comentarios.
// Esto permite ejecutar cada sentencia por separado evitando que PostgreSQL
// envuelva múltiples statements en una transacción implícita (que bloquea CREATE DATABASE).
func splitPostgresStatements(sqlText string) []string {
var stmts []string
var cur strings.Builder
i, n := 0, len(sqlText)
for i < n {
ch := sqlText[i]
// Dollar-quoting: $tag$ ... $tag$ (tag puede ser vacío: $$)
if ch == '$' {
j := i + 1
for j < n && sqlText[j] != '$' && sqlText[j] != '\n' {
j++
}
if j < n && sqlText[j] == '$' {
tag := sqlText[i : j+1] // e.g. "$$" or "$func$"
cur.WriteString(tag)
i = j + 1
closeIdx := strings.Index(sqlText[i:], tag)
if closeIdx >= 0 {
cur.WriteString(sqlText[i : i+closeIdx+len(tag)])
i = i + closeIdx + len(tag)
} else {
cur.WriteString(sqlText[i:])
i = n
}
continue
}
}
// String con comillas simples
if ch == '\'' {
cur.WriteByte(ch)
i++
for i < n {
c := sqlText[i]
cur.WriteByte(c)
i++
if c == '\'' {
if i < n && sqlText[i] == '\'' {
cur.WriteByte(sqlText[i])
i++
} else {
break
}
}
}
continue
}
// Comentario de línea (--)
if ch == '-' && i+1 < n && sqlText[i+1] == '-' {
cur.WriteByte(ch)
i++
for i < n && sqlText[i] != '\n' {
cur.WriteByte(sqlText[i])
i++
}
continue
}
// Comentario de bloque /* ... */
if ch == '/' && i+1 < n && sqlText[i+1] == '*' {
cur.WriteString("/*")
i += 2
for i < n {
if sqlText[i] == '*' && i+1 < n && sqlText[i+1] == '/' {
cur.WriteString("*/")
i += 2
break
}
cur.WriteByte(sqlText[i])
i++
}
continue
}
// Punto y coma: fin de sentencia
if ch == ';' {
if stmt := strings.TrimSpace(cur.String()); stmt != "" {
stmts = append(stmts, stmt)
}
cur.Reset()
i++
continue
}
cur.WriteByte(ch)
i++
}
if stmt := strings.TrimSpace(cur.String()); stmt != "" {
stmts = append(stmts, stmt)
}
return stmts
}
func isSelectStatement(sql string) bool {
upper := strings.ToUpper(strings.TrimSpace(sql))
keywords := []string{"SELECT ", "SHOW ", "DESCRIBE ", "EXPLAIN ", "WITH ", "PRAGMA "}
for _, kw := range keywords {
if strings.HasPrefix(upper, kw) {
return true
}
}
return false
}
func quoteIdentifier(name, driver string) string {
d := strings.ToLower(driver)
if strings.Contains(d, "postgres") {
return `"` + strings.ReplaceAll(name, `"`, `""`) + `"`
}
return "`" + strings.ReplaceAll(name, "`", "``") + "`"
}
func saveHistory(conxID uint, sqlText, status, errMsg string, rows, durationMs int64) {
models.SaveQueryHistory(models.QueryHistory{
ConxDbID: conxID,
SQL: sqlText,
Status: status,
ErrorMsg: errMsg,
RowsAffect: rows,
DurationMs: durationMs,
ExecutedAt: time.Now(),
})
}
// ── Redis ─────────────────────────────────────────────────────────────────────
func isRedisDriver(driver string) bool {
return strings.Contains(driver, "redis") || strings.Contains(driver, "valkey")
}
func redisConnect(c models.ConxDb, dbIndex int) (*goredis.Client, error) {
host := c.HostEfectivo()
port := c.Puerto
pass := c.Password
if port == "" {
port = "6379"
}
client := goredis.NewClient(&goredis.Options{
Addr: fmt.Sprintf("%s:%s", host, port),
Password: pass,
DB: dbIndex,
DialTimeout: 8 * time.Second,
ReadTimeout: 15 * time.Second,
WriteTimeout: 15 * time.Second,
})
ctx, cancel := context.WithTimeout(context.Background(), connTimeout)
defer cancel()
if err := client.Ping(ctx).Err(); err != nil {
client.Close()
return nil, fmt.Errorf("no se pudo conectar a Redis: %w", err)
}
return client, nil
}
func redisTestConnection(c models.ConxDb) error {
client, err := redisConnect(c, 0)
if err != nil {
return err
}
client.Close()
return nil
}
// redisListDatabases devuelve db0..dbN según CONFIG GET databases.
func redisListDatabases(c models.ConxDb) ([]string, error) {
client, err := redisConnect(c, 0)
if err != nil {
return nil, err
}
defer client.Close()
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
total := 16 // default
vals, err2 := client.ConfigGet(ctx, "databases").Result()
if err2 == nil && len(vals) >= 2 {
if n, parseErr := strconv.Atoi(fmt.Sprint(vals["databases"])); parseErr == nil && n > 0 {
total = n
}
}
dbs := make([]string, total)
for i := range dbs {
dbs[i] = fmt.Sprintf("db%d", i)
}
return dbs, nil
}
// redisListKeys devuelve hasta 200 keys del db seleccionado (SCAN con limit).
func redisListKeys(c models.ConxDb, database string) ([]string, error) {
dbIndex := 0
if strings.HasPrefix(database, "db") {
if n, err := strconv.Atoi(database[2:]); err == nil {
dbIndex = n
}
}
client, err := redisConnect(c, dbIndex)
if err != nil {
return nil, err
}
defer client.Close()
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
var keys []string
iter := client.Scan(ctx, 0, "*", 200).Iterator()
for iter.Next(ctx) {
keys = append(keys, iter.Val())
if len(keys) >= 200 {
break
}
}
if err := iter.Err(); err != nil {
return nil, err
}
return keys, nil
}
// redisExecuteCommand parsea y ejecuta un comando Redis.
// Los comandos se escriben como en redis-cli: GET key / SET key value / etc.
// Soporta múltiples líneas: cada línea no vacía es un comando independiente.
func redisExecuteCommand(conx models.ConxDb, database, cmdText string) QueryResult {
start := time.Now()
dbIndex := 0
if strings.HasPrefix(database, "db") {
if n, err := strconv.Atoi(database[2:]); err == nil {
dbIndex = n
}
}
client, err := redisConnect(conx, dbIndex)
if err != nil {
saveHistory(conx.ID, cmdText, "error", err.Error(), 0, time.Since(start).Milliseconds())
return QueryResult{Error: err.Error()}
}
defer client.Close()
// Dividir en líneas, ignorar vacías y comentarios (#)
var lines []string
for _, line := range strings.Split(cmdText, "\n") {
trimmed := strings.TrimSpace(line)
if trimmed == "" || strings.HasPrefix(trimmed, "#") {
continue
}
lines = append(lines, trimmed)
}
if len(lines) == 0 {
return QueryResult{Error: "comando vacío"}
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
// Si hay múltiples comandos, ejecutarlos en pipeline y devolver tabla de resultados
if len(lines) > 1 {
var rows []map[string]any
for _, line := range lines {
args := redisParseArgs(line)
if len(args) == 0 {
continue
}
ifaces := make([]any, len(args))
for i, a := range args {
ifaces[i] = a
}
val, execErr := client.Do(ctx, ifaces...).Result()
row := map[string]any{
"command": line,
"result": redisValToString(val),
"error": "",
}
if execErr != nil && execErr != goredis.Nil {
row["error"] = execErr.Error()
}
rows = append(rows, row)
}
elapsed := time.Since(start).Milliseconds()
result := QueryResult{
Columns: []string{"command", "result", "error"},
Rows: rows,
RowCount: len(rows),
IsSelect: true,
DurationMs: elapsed,
}
saveHistory(conx.ID, cmdText, "ok", "", int64(len(rows)), elapsed)
return result
}
// Comando único
args := redisParseArgs(lines[0])
if len(args) == 0 {
return QueryResult{Error: "comando vacío"}
}
ifaces := make([]any, len(args))
for i, a := range args {
ifaces[i] = a
}
val, execErr := client.Do(ctx, ifaces...).Result()
elapsed := time.Since(start).Milliseconds()
if execErr != nil && execErr != goredis.Nil {
saveHistory(conx.ID, cmdText, "error", execErr.Error(), 0, elapsed)
return QueryResult{Error: execErr.Error(), DurationMs: elapsed}
}
result := redisResultToQueryResult(val, lines[0])
result.DurationMs = elapsed
saveHistory(conx.ID, cmdText, "ok", "", int64(result.RowCount), elapsed)
return result
}
// redisResultToQueryResult convierte la respuesta de Redis en QueryResult presentable.
func redisResultToQueryResult(val any, cmd string) QueryResult {
upper := strings.ToUpper(strings.Fields(cmd)[0])
switch v := val.(type) {
case nil:
return QueryResult{
IsSelect: true,
Columns: []string{"result"},
Rows: []map[string]any{{"result": "(nil)"}},
RowCount: 1,
}
case string:
return QueryResult{
IsSelect: true,
Columns: []string{"result"},
Rows: []map[string]any{{"result": v}},
RowCount: 1,
}
case int64:
label := "result"
if upper == "DEL" || upper == "EXISTS" || upper == "SREM" || upper == "LREM" {
label = "affected"
} else if upper == "TTL" || upper == "PTTL" {
label = "ttl_seconds"
} else if upper == "DBSIZE" || upper == "LLEN" || upper == "SCARD" || upper == "ZCARD" || upper == "HLEN" {
label = "count"
}
return QueryResult{
IsSelect: false,
AffectedRows: v,
Columns: []string{label},
Rows: []map[string]any{{label: v}},
RowCount: 1,
}
case []any:
// Lista o conjunto de valores
if upper == "HGETALL" && len(v)%2 == 0 {
// Alternar field/value → formato tabla
var rows []map[string]any
for i := 0; i+1 < len(v); i += 2 {
rows = append(rows, map[string]any{
"field": redisValToString(v[i]),
"value": redisValToString(v[i+1]),
})
}
return QueryResult{IsSelect: true, Columns: []string{"field", "value"}, Rows: rows, RowCount: len(rows)}
}
// KEYS, SMEMBERS, LRANGE, etc.
var rows []map[string]any
for _, item := range v {
rows = append(rows, map[string]any{"value": redisValToString(item)})
}
return QueryResult{IsSelect: true, Columns: []string{"value"}, Rows: rows, RowCount: len(rows)}
case map[any]any:
var rows []map[string]any
for k, mv := range v {
rows = append(rows, map[string]any{
"field": redisValToString(k),
"value": redisValToString(mv),
})
}
return QueryResult{IsSelect: true, Columns: []string{"field", "value"}, Rows: rows, RowCount: len(rows)}
default:
return QueryResult{
IsSelect: true,
Columns: []string{"result"},
Rows: []map[string]any{{"result": fmt.Sprintf("%v", val)}},
RowCount: 1,
}
}
}
func redisValToString(v any) string {
if v == nil {
return "(nil)"
}
return fmt.Sprintf("%v", v)
}
// redisParseArgs divide un comando Redis en tokens respetando comillas.
// Ej: SET mykey "hello world" → ["SET", "mykey", "hello world"]
func redisParseArgs(cmd string) []string {
var args []string
var cur strings.Builder
inQ := false
qChar := byte(0)
for i := 0; i < len(cmd); i++ {
ch := cmd[i]
if inQ {
if ch == qChar {
inQ = false
} else if ch == '\\' && i+1 < len(cmd) {
i++
cur.WriteByte(cmd[i])
} else {
cur.WriteByte(ch)
}
} else {
if ch == '"' || ch == '\'' {
inQ = true
qChar = ch
} else if ch == ' ' || ch == '\t' {
if cur.Len() > 0 {
args = append(args, cur.String())
cur.Reset()
}
} else {
cur.WriteByte(ch)
}
}
}
if cur.Len() > 0 {
args = append(args, cur.String())
}
return args
}
// ── MongoDB ───────────────────────────────────────────────────────────────────
func isMongoDriver(driver string) bool {
return strings.Contains(driver, "mongo")
}
func mongoURI(c models.ConxDb) string {
host := c.HostEfectivo()
port := c.Puerto
user := c.Usuario
pass := c.Password
if user != "" && pass != "" {
return fmt.Sprintf("mongodb://%s:%s@%s:%s/?directConnection=true&serverSelectionTimeoutMS=8000", user, pass, host, port)
}
return fmt.Sprintf("mongodb://%s:%s/?directConnection=true&serverSelectionTimeoutMS=8000", host, port)
}
func mongoConnect(c models.ConxDb) (*mongo.Client, error) {
ctx, cancel := context.WithTimeout(context.Background(), connTimeout)
defer cancel()
client, err := mongo.Connect(ctx, options.Client().ApplyURI(mongoURI(c)))
if err != nil {
return nil, err
}
if err := client.Ping(ctx, nil); err != nil {
client.Disconnect(context.Background()) //nolint
return nil, fmt.Errorf("no se pudo conectar a MongoDB: %w", err)
}
return client, nil
}
func mongoTestConnection(c models.ConxDb) error {
client, err := mongoConnect(c)
if err != nil {
return err
}
client.Disconnect(context.Background()) //nolint
return nil
}
func mongoListDatabases(c models.ConxDb) ([]string, error) {
client, err := mongoConnect(c)
if err != nil {
return nil, err
}
defer client.Disconnect(context.Background()) //nolint
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
return client.ListDatabaseNames(ctx, bson.M{})
}
func mongoListCollections(c models.ConxDb, database string) ([]string, error) {
client, err := mongoConnect(c)
if err != nil {
return nil, err
}
defer client.Disconnect(context.Background()) //nolint
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
return client.Database(database).ListCollectionNames(ctx, bson.M{})
}
var mongoQueryRe = regexp.MustCompile(`(?s)^db\.(\w[\w\d_]*)\((.*)\)\s*$|^db\.(\w[\w\d_]*)\.(\w+)\((.*)\)\s*$`)
func mongoStripComments(text string) string {
var lines []string
for _, line := range strings.Split(text, "\n") {
t := strings.TrimSpace(line)
if t == "" || strings.HasPrefix(t, "--") || strings.HasPrefix(t, "//") || strings.HasPrefix(t, "#") {
continue
}
lines = append(lines, t)
}
return strings.TrimSpace(strings.Join(lines, "\n"))
}
func mongoExecuteSQL(conx models.ConxDb, database, queryText string) QueryResult {
start := time.Now()
client, err := mongoConnect(conx)
if err != nil {
saveHistory(conx.ID, queryText, "error", err.Error(), 0, time.Since(start).Milliseconds())
return QueryResult{Error: err.Error()}
}
defer client.Disconnect(context.Background()) //nolint
db := client.Database(database)
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
result, err := mongoRunQuery(ctx, db, strings.TrimSpace(queryText))
elapsed := time.Since(start).Milliseconds()
if err != nil {
saveHistory(conx.ID, queryText, "error", err.Error(), 0, elapsed)
return QueryResult{Error: err.Error()}
}
result.DurationMs = elapsed
saveHistory(conx.ID, queryText, "ok", "", int64(result.RowCount)+result.AffectedRows, elapsed)
return *result
}
func mongoRunQuery(ctx context.Context, db *mongo.Database, query string) (*QueryResult, error) {
// Eliminar comentarios y líneas vacías
query = mongoStripComments(query)
m := mongoQueryRe.FindStringSubmatch(query)
if m == nil {
return nil, fmt.Errorf("sintaxis no reconocida. Use: db.coleccion.metodo({...})")
}
var collName, method, argsRaw string
// El regex tiene dos alternativas:
// Alt 1 (m[1],m[2]): db.metodo(args) → para db.runCommand(...)
// Alt 2 (m[3],m[4],m[5]): db.coleccion.metodo(args)
if m[1] != "" {
// db.runCommand(args)
collName = m[1]
method = strings.ToLower(m[1])
argsRaw = strings.TrimSpace(m[2])
} else {
collName = m[3]
method = strings.ToLower(m[4])
argsRaw = strings.TrimSpace(m[5])
}
if strings.ToLower(collName) == "runcommand" {
return mongoRunCommand(ctx, db, argsRaw)
}
coll := db.Collection(collName)
args := splitTopLevelArgs(argsRaw)
switch method {
case "find":
return mongoFind(ctx, coll, args, false)
case "findone":
return mongoFind(ctx, coll, args, true)
case "insertone":
return mongoInsertOne(ctx, coll, args)
case "insertmany":
return mongoInsertMany(ctx, coll, args)
case "updateone":
return mongoUpdate(ctx, coll, args, false)
case "updatemany":
return mongoUpdate(ctx, coll, args, true)
case "deleteone":
return mongoDelete(ctx, coll, args, false)
case "deletemany":
return mongoDelete(ctx, coll, args, true)
case "countdocuments":
return mongoCount(ctx, coll, args)
case "aggregate":
return mongoAggregate(ctx, coll, argsRaw)
case "drop":
if err := coll.Drop(ctx); err != nil {
return nil, err
}
return &QueryResult{IsSelect: false}, nil
default:
return nil, fmt.Errorf("método MongoDB no soportado: %s", m[2])
}
}
func mongoFind(ctx context.Context, coll *mongo.Collection, args []string, one bool) (*QueryResult, error) {
var filter bson.M
if len(args) > 0 && args[0] != "" {
if err := bson.UnmarshalExtJSON([]byte(args[0]), true, &filter); err != nil {
return nil, fmt.Errorf("filtro inválido: %w", err)
}
} else {
filter = bson.M{}
}
var docs []bson.M
if one {
var doc bson.M
if err := coll.FindOne(ctx, filter).Decode(&doc); err != nil {
if err == mongo.ErrNoDocuments {
return &QueryResult{IsSelect: true, Columns: []string{}, Rows: []map[string]any{}}, nil
}
return nil, err
}
docs = []bson.M{doc}
} else {
cursor, err := coll.Find(ctx, filter)
if err != nil {
return nil, err
}
defer cursor.Close(ctx)
if err := cursor.All(ctx, &docs); err != nil {
return nil, err
}
}
return bsonDocsToResult(docs), nil
}
func mongoInsertOne(ctx context.Context, coll *mongo.Collection, args []string) (*QueryResult, error) {
if len(args) == 0 || args[0] == "" {
return nil, fmt.Errorf("insertOne requiere un documento")
}
var doc bson.M
if err := bson.UnmarshalExtJSON([]byte(args[0]), true, &doc); err != nil {
return nil, fmt.Errorf("documento inválido: %w", err)
}
res, err := coll.InsertOne(ctx, doc)
if err != nil {
return nil, err
}
return &QueryResult{
IsSelect: true, AffectedRows: 1,
Columns: []string{"insertedId"},
Rows: []map[string]any{{"insertedId": fmt.Sprintf("%v", res.InsertedID)}},
RowCount: 1,
}, nil
}
func mongoInsertMany(ctx context.Context, coll *mongo.Collection, args []string) (*QueryResult, error) {
if len(args) == 0 || args[0] == "" {
return nil, fmt.Errorf("insertMany requiere un array de documentos")
}
var arr bson.A
if err := bson.UnmarshalExtJSON([]byte(args[0]), true, &arr); err != nil {
return nil, fmt.Errorf("documentos inválidos: %w", err)
}
docs := make([]interface{}, len(arr))
copy(docs, arr)
res, err := coll.InsertMany(ctx, docs)
if err != nil {
return nil, err
}
return &QueryResult{IsSelect: false, AffectedRows: int64(len(res.InsertedIDs))}, nil
}
func mongoUpdate(ctx context.Context, coll *mongo.Collection, args []string, many bool) (*QueryResult, error) {
if len(args) < 2 {
return nil, fmt.Errorf("update requiere filtro y documento de actualización")
}
var filter, update bson.M
if err := bson.UnmarshalExtJSON([]byte(args[0]), true, &filter); err != nil {
return nil, fmt.Errorf("filtro inválido: %w", err)
}
if err := bson.UnmarshalExtJSON([]byte(args[1]), true, &update); err != nil {
return nil, fmt.Errorf("update inválido: %w", err)
}
var affected int64
if many {
res, err := coll.UpdateMany(ctx, filter, update)
if err != nil {
return nil, err
}
affected = res.ModifiedCount
} else {
res, err := coll.UpdateOne(ctx, filter, update)
if err != nil {
return nil, err
}
affected = res.ModifiedCount
}
return &QueryResult{IsSelect: false, AffectedRows: affected}, nil
}
func mongoDelete(ctx context.Context, coll *mongo.Collection, args []string, many bool) (*QueryResult, error) {
var filter bson.M
if len(args) > 0 && args[0] != "" {
if err := bson.UnmarshalExtJSON([]byte(args[0]), true, &filter); err != nil {
return nil, fmt.Errorf("filtro inválido: %w", err)
}
} else {
filter = bson.M{}
}
var deleted int64
if many {
res, err := coll.DeleteMany(ctx, filter)
if err != nil {
return nil, err
}
deleted = res.DeletedCount
} else {
res, err := coll.DeleteOne(ctx, filter)
if err != nil {
return nil, err
}
deleted = res.DeletedCount
}
return &QueryResult{IsSelect: false, AffectedRows: deleted}, nil
}
func mongoCount(ctx context.Context, coll *mongo.Collection, args []string) (*QueryResult, error) {
var filter bson.M
if len(args) > 0 && args[0] != "" {
if err := bson.UnmarshalExtJSON([]byte(args[0]), true, &filter); err != nil {
return nil, fmt.Errorf("filtro inválido: %w", err)
}
} else {
filter = bson.M{}
}
count, err := coll.CountDocuments(ctx, filter)
if err != nil {
return nil, err
}
return &QueryResult{
IsSelect: true,
Columns: []string{"count"},
Rows: []map[string]any{{"count": count}},
RowCount: 1,
}, nil
}
func mongoAggregate(ctx context.Context, coll *mongo.Collection, pipelineStr string) (*QueryResult, error) {
var pipeline bson.A
if err := bson.UnmarshalExtJSON([]byte(pipelineStr), true, &pipeline); err != nil {
return nil, fmt.Errorf("pipeline inválido: %w", err)
}
cursor, err := coll.Aggregate(ctx, pipeline)
if err != nil {
return nil, err
}
defer cursor.Close(ctx)
var docs []bson.M
if err := cursor.All(ctx, &docs); err != nil {
return nil, err
}
return bsonDocsToResult(docs), nil
}
func mongoRunCommand(ctx context.Context, db *mongo.Database, cmdStr string) (*QueryResult, error) {
var cmd bson.D
if err := bson.UnmarshalExtJSON([]byte(cmdStr), true, &cmd); err != nil {
return nil, fmt.Errorf("comando inválido: %w", err)
}
var result bson.M
if err := db.RunCommand(ctx, cmd).Decode(&result); err != nil {
return nil, err
}
return bsonDocsToResult([]bson.M{result}), nil
}
func bsonDocsToResult(docs []bson.M) *QueryResult {
if len(docs) == 0 {
return &QueryResult{IsSelect: true, Columns: []string{}, Rows: []map[string]any{}}
}
keyIdx := make(map[string]int)
for _, doc := range docs {
for k := range doc {
if _, seen := keyIdx[k]; !seen {
keyIdx[k] = len(keyIdx)
}
}
}
cols := make([]string, len(keyIdx))
for k, i := range keyIdx {
cols[i] = k
}
rows := make([]map[string]any, len(docs))
for i, doc := range docs {
row := make(map[string]any, len(doc))
for k, v := range doc {
switch v.(type) {
case bson.M, bson.A, bson.D:
b, _ := json.Marshal(v)
row[k] = string(b)
default:
row[k] = fmt.Sprintf("%v", v)
}
}
rows[i] = row
}
return &QueryResult{IsSelect: true, Columns: cols, Rows: rows, RowCount: len(rows)}
}
func splitTopLevelArgs(s string) []string {
var args []string
depth := 0
inStr := false
start := 0
s = strings.TrimSpace(s)
for i := 0; i < len(s); i++ {
ch := s[i]
switch ch {
case '"':
if i == 0 || s[i-1] != '\\' {
inStr = !inStr
}
case '{', '[', '(':
if !inStr {
depth++
}
case '}', ']', ')':
if !inStr {
depth--
}
case ',':
if !inStr && depth == 0 {
args = append(args, strings.TrimSpace(s[start:i]))
start = i + 1
}
}
}
if last := strings.TrimSpace(s[start:]); last != "" {
args = append(args, last)
}
return args
}