Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 19 additions & 28 deletions apps/web/src/app/api/sql-client/connect/route.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
import { requireBackendSession } from "@/lib/require-backend-session";
import { NextResponse } from "next/server";
import { Pool as PgPool } from "pg";
import mysql from "mysql2/promise";
import { getSqlPool, releaseSqlPool, SqlDbType } from "@/lib/sql-client-pool";

export async function POST(request: Request) {
const authError = await requireBackendSession(request);
Expand All @@ -14,37 +13,29 @@ export async function POST(request: Request) {
return NextResponse.json({ error: "type, host, and database are required" }, { status: 400 });
}

if (type === "postgresql") {
const pool = new PgPool({
if (type === "postgresql" || type === "mysql" || type === "mariadb") {
const handle = await getSqlPool({
type: type as SqlDbType,
host,
port: port || 5432,
port: Number(port) || 0,
database,
user: username,
username,
password,
ssl: ssl ? { rejectUnauthorized: false } : false,
connectionTimeoutMillis: 10000,
ssl: Boolean(ssl),
});
const client = await pool.connect();
const res = await client.query("SELECT version()");
client.release();
await pool.end();
return NextResponse.json({ success: true, version: res.rows[0]?.version });
}

if (type === "mysql" || type === "mariadb") {
const conn = await mysql.createConnection({
host,
port: port || 3306,
database,
user: username,
password,
ssl: ssl ? { rejectUnauthorized: false } : undefined,
connectTimeout: 10000,
});
const [rows] = await conn.execute("SELECT VERSION() as version");
await conn.end();
const version = (rows as { version: string }[])[0]?.version;
return NextResponse.json({ success: true, version });
try {
if (handle.pg) {
const res = await handle.pg.query("SELECT version()");
return NextResponse.json({ success: true, version: res.rows[0]?.version });
}

const [rows] = await handle.mysql!.query("SELECT VERSION() as version");
const version = (rows as { version: string }[])[0]?.version;
return NextResponse.json({ success: true, version });
} finally {
releaseSqlPool(handle.key);
}
}

return NextResponse.json({ error: `Unsupported database type: ${type}` }, { status: 400 });
Expand Down
174 changes: 77 additions & 97 deletions apps/web/src/app/api/sql-client/query/route.ts
Original file line number Diff line number Diff line change
@@ -1,109 +1,89 @@
import { requireBackendSession } from "@/lib/require-backend-session";
import { NextResponse } from "next/server";
import { Pool as PgPool } from "pg";
import mysql from "mysql2/promise";
import { getSqlPool, releaseSqlPool, SqlDbType } from "@/lib/sql-client-pool";
import { splitSqlStatements } from "@/lib/sql-split";

const MAX_ROWS = 5000;
const TIMEOUT_MS = 30000;

export async function POST(request: Request) {
const authError = await requireBackendSession(request);
if (authError) return authError;

try {
const { type, host, port, database, username, password, ssl, query, limit } =
await request.json();

if (!type || !host || !database || !query) {
return NextResponse.json(
{ error: "type, host, database, and query are required" },
{ status: 400 }
);
}

const rowLimit = Math.min(Number(limit) || 500, MAX_ROWS);
const start = Date.now();

if (type === "postgresql") {
const pool = new PgPool({
host,
port: port || 5432,
database,
user: username,
password,
ssl: ssl ? { rejectUnauthorized: false } : false,
connectionTimeoutMillis: TIMEOUT_MS,
statement_timeout: TIMEOUT_MS,
});

const client = await pool.connect();
try {
// Split statements by semicolon and execute each, returning last result
const statements = query
.split(";")
.map((s: string) => s.trim())
.filter(Boolean);

let result = null;
for (const stmt of statements) {
result = await client.query(stmt);
}

const elapsed = Date.now() - start;
const rows = result?.rows?.slice(0, rowLimit) ?? [];
const columns = result?.fields?.map((f: { name: string }) => f.name) ?? [];
const rowCount = result?.rowCount ?? rows.length;

return NextResponse.json({ rows, columns, rowCount, executionTime: elapsed });
} finally {
client.release();
await pool.end();
}
}

if (type === "mysql" || type === "mariadb") {
const conn = await mysql.createConnection({
host,
port: port || 3306,
database,
user: username,
password,
ssl: ssl ? { rejectUnauthorized: false } : undefined,
connectTimeout: TIMEOUT_MS,
multipleStatements: true,
});

try {
const [rawRows, rawFields] = await conn.execute(query);
const elapsed = Date.now() - start;
interface StatementResult {
rows: Record<string, unknown>[];
columns: string[];
rowCount: number;
}

// multipleStatements may return arrays of result sets
const isMulti = Array.isArray(rawRows) && Array.isArray(rawRows[0]);
const rows = isMulti
? (((rawRows as unknown) as unknown[][]).at(-1) as Record<string, unknown>[]) ?? []
: (rawRows as Record<string, unknown>[]);
export async function POST(request: Request) {
const authError = await requireBackendSession(request);
if (authError) return authError;

try {
const { type, host, port, database, username, password, ssl, query, limit } =
await request.json();

if (!type || !host || !database || !query) {
return NextResponse.json(
{ error: "type, host, database, and query are required" },
{ status: 400 }
);
}

const fields = isMulti
? (((rawFields as unknown) as unknown[][]).at(-1) as { name: string }[]) ?? []
: (rawFields as { name: string }[]);
const rowLimit = Math.min(Number(limit) || 500, MAX_ROWS);
const start = Date.now();
const statements = splitSqlStatements(query);
if (statements.length === 0) {
return NextResponse.json({ error: "No SQL statement provided" }, { status: 400 });
}

const slicedRows = Array.isArray(rows) ? rows.slice(0, rowLimit) : [];
const columns = Array.isArray(fields) ? fields.map((f) => f.name) : [];
const handle = await getSqlPool({
type: type as SqlDbType,
host,
port: Number(port) || 0,
database,
username,
password,
ssl: Boolean(ssl),
});

return NextResponse.json({
rows: slicedRows,
columns,
rowCount: Array.isArray(rows) ? rows.length : 0,
executionTime: elapsed,
});
} finally {
await conn.end();
}
try {
const results: StatementResult[] = [];

for (const stmt of statements) {
// ponytail: cooperative abort between statements only. True mid-statement
// cancel needs pg_cancel_backend / conn.destroy(); add if long queries need killing.
if (request.signal.aborted) throw new Error("Query aborted");

if (handle.pg) {
const r = await handle.pg.query(stmt);
const rows = (r.rows ?? []) as Record<string, unknown>[];
results.push({
rows: rows.slice(0, rowLimit),
columns: r.fields?.map((f: { name: string }) => f.name) ?? [],
rowCount: r.rowCount ?? rows.length,
});
} else if (handle.mysql) {
const [rawRows, rawFields] = await handle.mysql.query(stmt);
const rows = Array.isArray(rawRows) ? (rawRows as Record<string, unknown>[]) : [];
const fields = Array.isArray(rawFields) ? (rawFields as { name: string }[]) : [];
results.push({
rows: rows.slice(0, rowLimit),
columns: fields.map((f) => f.name),
rowCount: rows.length,
});
}

return NextResponse.json({ error: `Unsupported database type: ${type}` }, { status: 400 });
} catch (error: unknown) {
const message = error instanceof Error ? error.message : String(error);
return NextResponse.json({ error: message }, { status: 500 });
}

const last = results[results.length - 1] ?? { rows: [], columns: [], rowCount: 0 };
return NextResponse.json({
rows: last.rows,
columns: last.columns,
rowCount: last.rowCount,
executionTime: Date.now() - start,
results,
});
} finally {
releaseSqlPool(handle.key);
}
} catch (error: unknown) {
const message = error instanceof Error ? error.message : String(error);
return NextResponse.json({ error: message }, { status: 500 });
}
}
Loading