diff --git a/apps/desktop/src/lib/common/__tests__/jsonTree.spec.ts b/apps/desktop/src/lib/common/__tests__/jsonTree.spec.ts
new file mode 100644
index 000000000..18d898cb3
--- /dev/null
+++ b/apps/desktop/src/lib/common/__tests__/jsonTree.spec.ts
@@ -0,0 +1,46 @@
+import { describe, expect, it } from "vitest";
+import { appendJsonPointer, createJsonTreeRoot, getJsonTreeChildren, getVisibleJsonTreeNodes, isJsonTreeContainer, isJsonTreeInitiallyExpanded } from "../jsonTree";
+import { parseJsonPreservingLargeNumbers } from "../safeJsonFormat";
+
+describe("jsonTree", () => {
+ it("uses RFC 6901 paths for object keys and array positions", () => {
+ const root = createJsonTreeRoot({ "a/b~c": ["value"] });
+ const objectChild = getJsonTreeChildren(root)[0];
+ const arrayChild = getJsonTreeChildren(objectChild)[0];
+
+ expect(root.path).toBe("");
+ expect(objectChild.path).toBe("/a~1b~0c");
+ expect(arrayChild.path).toBe("/a~1b~0c/0");
+ expect(appendJsonPointer("/parent", "a/b~c")).toBe("/parent/a~1b~0c");
+ });
+
+ it("keeps lossless JSON numbers as scalar nodes", () => {
+ const value = parseJsonPreservingLargeNumbers('{"id":518400931654815740}') as Record;
+ const child = getJsonTreeChildren(createJsonTreeRoot(value))[0];
+
+ expect(isJsonTreeContainer(child.value)).toBe(false);
+ });
+
+ it("treats initial depth as expanded container levels from the root", () => {
+ expect(isJsonTreeInitiallyExpanded(0, 2)).toBe(true);
+ expect(isJsonTreeInitiallyExpanded(1, 2)).toBe(true);
+ expect(isJsonTreeInitiallyExpanded(2, 2)).toBe(false);
+ expect(isJsonTreeInitiallyExpanded(99, Number.POSITIVE_INFINITY)).toBe(true);
+ });
+
+ it("flattens expanded branches iteratively for virtual rendering", () => {
+ const root = createJsonTreeRoot({ first: { nested: true }, second: ["kept"] });
+ const nodes = getVisibleJsonTreeNodes(root, (node) => node.path !== "/first");
+
+ expect(nodes.map((node) => node.path)).toEqual(["", "/first", "/second", "/second/0"]);
+ });
+
+ it("handles deeply nested expanded JSON without recursive traversal", () => {
+ let value: unknown = true;
+ for (let depth = 0; depth < 2_000; depth += 1) value = { child: value };
+
+ const nodes = getVisibleJsonTreeNodes(createJsonTreeRoot(value), () => true);
+
+ expect(nodes).toHaveLength(2_001);
+ });
+});
diff --git a/apps/desktop/src/lib/common/__tests__/safeJsonFormat.spec.ts b/apps/desktop/src/lib/common/__tests__/safeJsonFormat.spec.ts
index c78e9cdc7..bbe77f5f4 100644
--- a/apps/desktop/src/lib/common/__tests__/safeJsonFormat.spec.ts
+++ b/apps/desktop/src/lib/common/__tests__/safeJsonFormat.spec.ts
@@ -73,4 +73,16 @@ describe("safeJsonFormat", () => {
expect(isLosslessJsonNumber(parsed.companyId) ? parsed.companyId.raw : null).toBe("518400931654815740");
expect(parsed.safe).toBe(42);
});
+
+ it("preserves fractional and exponent literals that Number would round or overflow", () => {
+ const input = '{"fraction":0.123456789012345678901234,"scientific":1.234567890123456789e20,"overflow":1e999}';
+ const parsed = parseJsonPreservingLargeNumbers(input) as Record;
+
+ expect(isLosslessJsonNumber(parsed.fraction) ? parsed.fraction.raw : null).toBe("0.123456789012345678901234");
+ expect(isLosslessJsonNumber(parsed.scientific) ? parsed.scientific.raw : null).toBe("1.234567890123456789e20");
+ expect(isLosslessJsonNumber(parsed.overflow) ? parsed.overflow.raw : null).toBe("1e999");
+ expect(safeJsonFormat(input, 2)).toContain("0.123456789012345678901234");
+ expect(safeJsonFormat(input, 2)).toContain("1.234567890123456789e20");
+ expect(safeJsonFormat(input, 2)).toContain("1e999");
+ });
});
diff --git a/apps/desktop/src/lib/common/jsonTree.ts b/apps/desktop/src/lib/common/jsonTree.ts
new file mode 100644
index 000000000..62792d3f9
--- /dev/null
+++ b/apps/desktop/src/lib/common/jsonTree.ts
@@ -0,0 +1,111 @@
+import { isLosslessJsonNumber } from "./safeJsonFormat";
+
+export type JsonTreeContainerKind = "array" | "object";
+export type JsonTreeParentKind = JsonTreeContainerKind | "root";
+
+export interface JsonTreeNode {
+ key: string;
+ label: string;
+ value: unknown;
+ /** RFC 6901 JSON Pointer. The root value is represented by an empty string. */
+ path: string;
+ /** Zero-based structural depth; the root value is at depth zero. */
+ depth: number;
+ parentKind: JsonTreeParentKind;
+}
+
+export type JsonTreeContainer = Record | unknown[];
+
+export function createJsonTreeRoot(value: unknown): JsonTreeNode {
+ return {
+ key: "$",
+ label: "$",
+ value,
+ path: "",
+ depth: 0,
+ parentKind: "root",
+ };
+}
+
+export function isJsonTreeContainer(value: unknown): value is JsonTreeContainer {
+ // LosslessJsonNumber is an object wrapper, but represents a scalar JSON number.
+ return value !== null && typeof value === "object" && !isLosslessJsonNumber(value);
+}
+
+export function jsonTreeContainerKind(value: JsonTreeContainer): JsonTreeContainerKind {
+ return Array.isArray(value) ? "array" : "object";
+}
+
+export function jsonTreeContainerSummary(value: JsonTreeContainer, includeObjectLength = true): string {
+ if (Array.isArray(value)) return `Array(${value.length})`;
+ return includeObjectLength ? `Object(${Object.keys(value).length})` : "Object";
+}
+
+/** Escape a reference token according to RFC 6901. */
+export function escapeJsonPointerSegment(segment: string): string {
+ return segment.replaceAll("~", "~0").replaceAll("/", "~1");
+}
+
+export function appendJsonPointer(path: string, segment: string | number): string {
+ return `${path}/${escapeJsonPointerSegment(String(segment))}`;
+}
+
+/**
+ * Create child nodes only for an already-expanded parent. Callers should not
+ * invoke this while a container is collapsed so large JSON payloads stay lazy.
+ */
+export function getJsonTreeChildren(node: JsonTreeNode): JsonTreeNode[] {
+ if (!isJsonTreeContainer(node.value)) return [];
+
+ if (Array.isArray(node.value)) {
+ return node.value.map((value, index) => ({
+ key: String(index),
+ label: String(index),
+ value,
+ path: appendJsonPointer(node.path, index),
+ depth: node.depth + 1,
+ parentKind: "array",
+ }));
+ }
+
+ return Object.entries(node.value).map(([key, value]) => ({
+ key,
+ label: key,
+ value,
+ path: appendJsonPointer(node.path, key),
+ depth: node.depth + 1,
+ parentKind: "object",
+ }));
+}
+
+/**
+ * Flatten the currently expanded tree with an iterative traversal. This lets
+ * virtual renderers keep every node logically expanded without recursive DOM
+ * creation or deep-call-stack failures.
+ */
+export function getVisibleJsonTreeNodes(root: JsonTreeNode, isExpanded: (node: JsonTreeNode) => boolean): JsonTreeNode[] {
+ const nodes: JsonTreeNode[] = [];
+ const pending = [root];
+
+ while (pending.length > 0) {
+ const node = pending.pop();
+ if (!node) continue;
+ nodes.push(node);
+
+ if (!isJsonTreeContainer(node.value) || !isExpanded(node)) continue;
+ const children = getJsonTreeChildren(node);
+ for (let index = children.length - 1; index >= 0; index -= 1) pending.push(children[index]);
+ }
+
+ return nodes;
+}
+
+/**
+ * `initialExpandedDepth` counts container levels from the root. For example,
+ * a value of 2 expands the root and its direct container children.
+ */
+export function isJsonTreeInitiallyExpanded(depth: number, initialExpandedDepth: number): boolean {
+ if (initialExpandedDepth === Number.POSITIVE_INFINITY) return true;
+ if (!Number.isFinite(initialExpandedDepth)) return false;
+ return depth < Math.max(0, Math.floor(initialExpandedDepth));
+}
diff --git a/apps/desktop/src/lib/common/safeJsonFormat.ts b/apps/desktop/src/lib/common/safeJsonFormat.ts
index 977c2165e..7721d79b8 100644
--- a/apps/desktop/src/lib/common/safeJsonFormat.ts
+++ b/apps/desktop/src/lib/common/safeJsonFormat.ts
@@ -18,7 +18,7 @@ export function isLosslessJsonNumber(value: unknown): value is LosslessJsonNumbe
/**
* Parses JSON while retaining numeric literals that JavaScript cannot safely
- * represent as numbers. Callers can render these values without adding quotes.
+ * represent exactly. Callers can render these values without adding quotes.
*/
export function parseJsonPreservingLargeNumbers(text: string): unknown {
const protectedJson = protectLargeJsonNumbers(text);
@@ -28,7 +28,8 @@ export function parseJsonPreservingLargeNumbers(text: string): unknown {
/**
* Parse and re-stringify JSON while preserving numeric literals whose integer
- * parts exceed Number.MAX_SAFE_INTEGER (2^53 - 1).
+ * parts exceed Number.MAX_SAFE_INTEGER (2^53 - 1), plus decimal and exponent
+ * forms that JavaScript may round or turn into Infinity.
*/
export function safeJsonFormat(text: string, indent?: number): string {
const protectedJson = protectLargeJsonNumbers(text);
@@ -62,7 +63,7 @@ function protectLargeJsonNumbers(text: string): ProtectedJsonNumbers {
const numberMatch = text.slice(index).match(/^-?(?:0|[1-9]\d*)(?:\.\d+)?(?:[eE][+-]?\d+)?/);
if (numberMatch) {
const raw = numberMatch[0];
- if (hasUnsafeIntegerPart(raw)) {
+ if (shouldPreserveJsonNumber(raw)) {
// A quoted placeholder lets the native parser validate the rest of the JSON.
const placeholder = `${placeholderPrefix}${numbers.size}__`;
numbers.set(placeholder, raw);
@@ -96,7 +97,12 @@ function findJsonStringEnd(text: string, start: number): number {
return text.length;
}
-function hasUnsafeIntegerPart(raw: string): boolean {
+function shouldPreserveJsonNumber(raw: string): boolean {
+ // Keep fractional/exponent forms verbatim. Even when a particular value is
+ // representable today, parsing it through Number can change its precision or
+ // spelling before the JSON viewer renders it.
+ if (raw.includes(".") || raw.includes("e") || raw.includes("E") || raw === "-0") return true;
+
const unsigned = raw.startsWith("-") ? raw.slice(1) : raw;
const integerPart = unsigned.split(/[.eE]/, 1)[0];
const normalized = integerPart.replace(/^0+(?=\d)/, "");
diff --git a/apps/desktop/src/lib/redis/redisJsonHighlighter.ts b/apps/desktop/src/lib/common/shikiJsonHighlighter.ts
similarity index 66%
rename from apps/desktop/src/lib/redis/redisJsonHighlighter.ts
rename to apps/desktop/src/lib/common/shikiJsonHighlighter.ts
index 3c615cb32..11420b2e2 100644
--- a/apps/desktop/src/lib/redis/redisJsonHighlighter.ts
+++ b/apps/desktop/src/lib/common/shikiJsonHighlighter.ts
@@ -1,8 +1,8 @@
import type { AppThemeAppearance } from "@/lib/app/appTheme";
-export type RedisJsonHighlighter = (content: string, appearance?: AppThemeAppearance) => string;
+export type JsonHighlighter = (content: string, appearance?: AppThemeAppearance) => string;
-interface RedisShikiJsonHighlighterOptions {
+interface ShikiJsonHighlighterOptions {
appearance: () => AppThemeAppearance;
}
@@ -15,8 +15,8 @@ type ShikiHighlighter = Awaited | undefined;
-export async function createRedisShikiJsonHighlighter(options: RedisShikiJsonHighlighterOptions): Promise {
- const highlighter = await getRedisShikiHighlighter();
+export async function createShikiJsonHighlighter(options: ShikiJsonHighlighterOptions): Promise {
+ const highlighter = await getShikiJsonHighlighter();
return (content, appearance = options.appearance()) =>
highlighter.codeToHtml(content, {
lang: "json",
@@ -25,12 +25,12 @@ export async function createRedisShikiJsonHighlighter(options: RedisShikiJsonHig
});
}
-function getRedisShikiHighlighter(): Promise {
- highlighterPromise ??= loadRedisShikiHighlighter();
+function getShikiJsonHighlighter(): Promise {
+ highlighterPromise ??= loadShikiJsonHighlighter();
return highlighterPromise;
}
-async function loadRedisShikiHighlighter(): Promise {
+async function loadShikiJsonHighlighter(): Promise {
const [{ createHighlighterCore }, { createJavaScriptRegexEngine }, githubDark, githubLight, json] = await Promise.all([import("shiki/core"), import("shiki/engine/javascript"), import("shiki/themes/github-dark.mjs"), import("shiki/themes/github-light.mjs"), import("shiki/langs/json.mjs")]);
return createHighlighterCore({
diff --git a/apps/desktop/src/lib/elasticsearch/elasticsearchJsonResponse.ts b/apps/desktop/src/lib/elasticsearch/elasticsearchJsonResponse.ts
new file mode 100644
index 000000000..2fd731ac8
--- /dev/null
+++ b/apps/desktop/src/lib/elasticsearch/elasticsearchJsonResponse.ts
@@ -0,0 +1,26 @@
+import type { DatabaseType, QueryResult } from "@/types/database";
+
+export interface ElasticsearchJsonResponse {
+ status: number;
+ body: string;
+}
+
+const ELASTICSEARCH_REST_STATEMENT = /^(?:GET|POST|PUT|DELETE)\s+\S+/i;
+
+/**
+ * Detect the result shape emitted for a JSON response to an explicit
+ * Elasticsearch REST request. SQL and text (such as CAT) results keep using
+ * the normal data-grid path.
+ */
+export function elasticsearchJsonResponseForResult(databaseType: DatabaseType | undefined, sourceStatement: string | undefined, result: QueryResult | undefined): ElasticsearchJsonResponse | undefined {
+ if (databaseType !== "elasticsearch" || !result || typeof sourceStatement !== "string") return undefined;
+ if (!ELASTICSEARCH_REST_STATEMENT.test(sourceStatement.trim())) return undefined;
+ if (result.columns.length !== 2 || result.columns[0] !== "status" || result.columns[1] !== "response" || result.rows.length !== 1) return undefined;
+
+ const row = result.rows[0];
+ if (!row || row.length !== 2) return undefined;
+
+ const [status, body] = row;
+ if (typeof status !== "number" || !Number.isInteger(status) || status < 100 || status > 599 || typeof body !== "string") return undefined;
+ return { status, body };
+}
diff --git a/crates/dbx-core/Cargo.toml b/crates/dbx-core/Cargo.toml
index 93aeb2fe1..5f6a8d3a2 100644
--- a/crates/dbx-core/Cargo.toml
+++ b/crates/dbx-core/Cargo.toml
@@ -39,7 +39,7 @@ sqlite-sqlcipher = ["rusqlite/bundled-sqlcipher-vendored-openssl"]
[dependencies]
serde = { version = "1.0", features = ["derive"] }
-serde_json = { version = "1.0", features = ["preserve_order"] }
+serde_json = { version = "1.0", features = ["arbitrary_precision", "preserve_order"] }
regex = "1"
rayon = "1"
percent-encoding = "2"
diff --git a/crates/dbx-core/src/db/elasticsearch_driver.rs b/crates/dbx-core/src/db/elasticsearch_driver.rs
index 1724edd35..0baad8947 100644
--- a/crates/dbx-core/src/db/elasticsearch_driver.rs
+++ b/crates/dbx-core/src/db/elasticsearch_driver.rs
@@ -790,7 +790,15 @@ pub async fn execute_rest_query(client: &EsClient, input: &str) -> Result client.delete(&path).send().await,
+ "DELETE" => {
+ let req = client.delete(&path);
+ if let Some(b) = body {
+ let json: serde_json::Value = serde_json::from_str(b).map_err(|e| format!("Invalid JSON body: {e}"))?;
+ req.json(&json).send().await
+ } else {
+ req.send().await
+ }
+ }
_ => return Err(format!("Unsupported HTTP method: {method}. Use GET, POST, PUT, or DELETE.")),
}
.map_err(|e| format!("Elasticsearch request failed: {e}"))?;
@@ -865,18 +873,7 @@ fn parse_elasticsearch_response(
has_more: false,
})
} else {
- let pretty = serde_json::to_string_pretty(&body).unwrap_or_else(|_| body.to_string());
- Ok(crate::types::QueryResult {
- columns: vec!["status".to_string(), "response".to_string()],
- column_types: Vec::new(),
- column_sortables: vec![],
- rows: vec![vec![serde_json::Value::Number(status.into()), serde_json::Value::String(pretty)]],
- affected_rows: 0,
- execution_time_ms: start.elapsed().as_millis(),
- truncated: false,
- session_id: None,
- has_more: false,
- })
+ Ok(json_response_result(status, &body, start))
}
} else if let Some(hits) = body.pointer("/hits/hits").and_then(|v| v.as_array()) {
// Treat any `_search`-shaped body as the hits result, even when empty —
@@ -941,18 +938,30 @@ fn parse_elasticsearch_response(
has_more: false,
})
} else {
- let pretty = serde_json::to_string_pretty(&body).unwrap_or_else(|_| body.to_string());
- Ok(crate::types::QueryResult {
- columns: vec!["status".to_string(), "response".to_string()],
- column_types: Vec::new(),
- column_sortables: vec![],
- rows: vec![vec![serde_json::Value::Number(status.into()), serde_json::Value::String(pretty)]],
- affected_rows: 0,
- execution_time_ms: start.elapsed().as_millis(),
- truncated: false,
- session_id: None,
- has_more: false,
- })
+ Ok(json_response_result(status, &body, start))
+ }
+}
+
+fn json_response_result(status: u16, body: &serde_json::Value, start: std::time::Instant) -> crate::types::QueryResult {
+ let body_text = serde_json::to_string_pretty(body).unwrap_or_else(|_| body.to_string());
+ raw_json_response_result(status, body_text, start)
+}
+
+fn raw_json_response_result(
+ status: u16,
+ body_text: impl Into,
+ start: std::time::Instant,
+) -> crate::types::QueryResult {
+ crate::types::QueryResult {
+ columns: vec!["status".to_string(), "response".to_string()],
+ column_types: Vec::new(),
+ column_sortables: vec![],
+ rows: vec![vec![serde_json::Value::Number(status.into()), serde_json::Value::String(body_text.into())]],
+ affected_rows: 0,
+ execution_time_ms: start.elapsed().as_millis(),
+ truncated: false,
+ session_id: None,
+ has_more: false,
}
}
@@ -962,11 +971,13 @@ fn parse_elasticsearch_rest_response(
start: std::time::Instant,
) -> Result {
if body_text.trim().is_empty() {
- return parse_elasticsearch_response(status, serde_json::Value::Null, start);
+ return Ok(json_response_result(status, &serde_json::Value::Null, start));
}
- if let Ok(body) = serde_json::from_str::(body_text) {
- return parse_elasticsearch_response(status, body, start);
+ if serde_json::from_str::(body_text).is_ok() {
+ // Validate the payload as JSON, but retain the HTTP body verbatim so
+ // numeric literals are not changed by a parse/serialize round trip.
+ return Ok(raw_json_response_result(status, body_text, start));
}
// CAT APIs default to text/plain for human-readable output. Keep those
@@ -1501,6 +1512,35 @@ mod tests {
use serde_json::json;
use std::time::Duration;
+ async fn read_http_request(socket: &mut tokio::net::TcpStream) -> String {
+ use tokio::io::AsyncReadExt;
+
+ let mut bytes = Vec::new();
+ let mut buffer = [0_u8; 1024];
+ loop {
+ let read = socket.read(&mut buffer).await.unwrap();
+ assert!(read > 0, "HTTP request ended before its body was received");
+ bytes.extend_from_slice(&buffer[..read]);
+
+ let Some(headers_end) = bytes.windows(4).position(|window| window == b"\r\n\r\n") else {
+ continue;
+ };
+ let content_length = std::str::from_utf8(&bytes[..headers_end])
+ .unwrap()
+ .lines()
+ .find_map(|line| {
+ let (name, value) = line.split_once(':')?;
+ name.eq_ignore_ascii_case("content-length").then(|| value.trim().parse::().unwrap())
+ })
+ .unwrap_or(0);
+ let request_end = headers_end + 4 + content_length;
+ if bytes.len() >= request_end {
+ bytes.truncate(request_end);
+ return String::from_utf8(bytes).unwrap();
+ }
+ }
+ }
+
#[test]
fn url_params_can_disable_elasticsearch_tls_verification() {
assert!(elasticsearch_accept_invalid_certs(false, Some("sslmode=disable")));
@@ -1718,10 +1758,29 @@ mod tests {
)
.unwrap();
+ assert_ne!(result.columns, vec!["status", "response"]);
+ let name_idx = result.columns.iter().position(|column| column == "name").unwrap();
+ assert_eq!(result.rows[0][name_idx], json!("Alice"));
let routing_idx = result.columns.iter().position(|column| column == "_routing").unwrap();
assert_eq!(result.rows[0][routing_idx], json!("tenant-1"));
}
+ #[test]
+ fn keeps_sql_api_response_tabular() {
+ let result = super::parse_elasticsearch_response(
+ 200,
+ json!({
+ "columns": [{ "name": "name" }],
+ "rows": [["Alice"]]
+ }),
+ std::time::Instant::now(),
+ )
+ .unwrap();
+
+ assert_eq!(result.columns, vec!["name"]);
+ assert_eq!(result.rows, vec![vec![json!("Alice")]]);
+ }
+
#[test]
fn parses_aggregation_response_before_empty_hits() {
let result = super::parse_elasticsearch_response(
@@ -1769,6 +1828,27 @@ mod tests {
assert_eq!(result.affected_rows, 2);
}
+ #[test]
+ fn keeps_mapping_rest_response_numeric_literals_lossless() {
+ let body = r#"{
+ "products": {
+ "mappings": {
+ "_meta": {
+ "largest_id": 123456789012345678901234567890,
+ "ratio": 0.123456789012345678901234567890,
+ "estimate": 1e400
+ },
+ "properties": { "name": { "type": "keyword" } }
+ }
+ }
+}"#;
+ let result = super::parse_elasticsearch_rest_response(200, body, std::time::Instant::now()).unwrap();
+
+ assert_eq!(result.columns, vec!["status", "response"]);
+ assert_eq!(result.rows[0][0], json!(200));
+ assert_eq!(result.rows[0][1].as_str(), Some(body));
+ }
+
#[tokio::test]
async fn execute_rest_query_keeps_plain_text_response_body() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
@@ -1800,6 +1880,181 @@ mod tests {
assert_eq!(result.rows[1][0], json!("green open app-log-2026-07 42 10mb"));
}
+ #[tokio::test]
+ async fn execute_rest_query_preserves_numeric_literals_from_http_body() {
+ use tokio::io::AsyncWriteExt;
+
+ let response_body = r#"{"largest_id":123456789012345678901234567890,"ratio":0.123456789012345678901234567890,"estimate":1e400}"#;
+ let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let addr = listener.local_addr().unwrap();
+ let server = tokio::spawn(async move {
+ let (mut socket, _) = listener.accept().await.unwrap();
+ let request = read_http_request(&mut socket).await;
+ assert!(request.starts_with("GET /products/_mapping "));
+ let response = format!(
+ "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
+ response_body.len(),
+ response_body
+ );
+ socket.write_all(response.as_bytes()).await.unwrap();
+ });
+
+ let client = EsClient::new(&format!("http://{addr}"), None, None, false, Duration::from_secs(1));
+ let result = super::execute_rest_query(&client, "GET /products/_mapping").await.unwrap();
+ server.await.unwrap();
+
+ assert_eq!(result.columns, vec!["status", "response"]);
+ assert_eq!(result.rows[0][0], json!(200));
+ assert_eq!(result.rows[0][1].as_str(), Some(response_body));
+ }
+
+ #[tokio::test]
+ async fn execute_rest_delete_sends_json_body() {
+ use tokio::io::AsyncWriteExt;
+
+ let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let addr = listener.local_addr().unwrap();
+ let server = tokio::spawn(async move {
+ let (mut socket, _) = listener.accept().await.unwrap();
+ let request = read_http_request(&mut socket).await;
+ let (headers, body) = request.split_once("\r\n\r\n").unwrap();
+ assert!(headers.starts_with("DELETE /_search/scroll "));
+ assert_eq!(serde_json::from_str::(body).unwrap(), json!({ "scroll_id": ["scroll-1"] }));
+
+ let response_body = r#"{"succeeded":true,"num_freed":1}"#;
+ let response = format!(
+ "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
+ response_body.len(),
+ response_body
+ );
+ socket.write_all(response.as_bytes()).await.unwrap();
+ });
+
+ let client = EsClient::new(&format!("http://{addr}"), None, None, false, Duration::from_secs(1));
+ let result =
+ super::execute_rest_query(&client, "DELETE /_search/scroll\n{\"scroll_id\":[\"scroll-1\"]}").await.unwrap();
+ server.await.unwrap();
+
+ assert_eq!(result.columns, vec!["status", "response"]);
+ assert_eq!(result.rows[0][0], json!(200));
+ assert_eq!(
+ serde_json::from_str::(result.rows[0][1].as_str().unwrap()).unwrap(),
+ json!({ "succeeded": true, "num_freed": 1 })
+ );
+ }
+
+ #[tokio::test]
+ async fn execute_rest_query_keeps_json_error_response() {
+ use tokio::io::AsyncWriteExt;
+
+ let response_body = r#"{"error":{"type":"index_not_found_exception","reason":"no such index"},"status":404}"#;
+ let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let addr = listener.local_addr().unwrap();
+ let server = tokio::spawn(async move {
+ let (mut socket, _) = listener.accept().await.unwrap();
+ let request = read_http_request(&mut socket).await;
+ assert!(request.starts_with("GET /missing/_mapping "));
+ let response = format!(
+ "HTTP/1.1 404 Not Found\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
+ response_body.len(),
+ response_body
+ );
+ socket.write_all(response.as_bytes()).await.unwrap();
+ });
+
+ let client = EsClient::new(&format!("http://{addr}"), None, None, false, Duration::from_secs(1));
+ let result = super::execute_rest_query(&client, "GET /missing/_mapping").await.unwrap();
+ server.await.unwrap();
+
+ assert_eq!(result.columns, vec!["status", "response"]);
+ assert_eq!(result.rows[0][0], json!(404));
+ assert_eq!(
+ serde_json::from_str::(result.rows[0][1].as_str().unwrap()).unwrap(),
+ json!({ "error": { "type": "index_not_found_exception", "reason": "no such index" }, "status": 404 })
+ );
+ }
+
+ #[tokio::test]
+ async fn execute_select_query_keeps_search_response_tabular() {
+ use tokio::io::AsyncWriteExt;
+
+ let response_body = r#"{"hits":{"hits":[{"_id":"product-1","_source":{"name":"Notebook"}}]}}"#;
+ let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let addr = listener.local_addr().unwrap();
+ let server = tokio::spawn(async move {
+ let (mut socket, _) = listener.accept().await.unwrap();
+ let request = read_http_request(&mut socket).await;
+ assert!(request.starts_with("POST /products/_search "));
+ let response = format!(
+ "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
+ response_body.len(),
+ response_body
+ );
+ socket.write_all(response.as_bytes()).await.unwrap();
+ });
+
+ let client = EsClient::new(&format!("http://{addr}"), None, None, false, Duration::from_secs(1));
+ let result = super::execute_rest_query(&client, "SELECT * FROM products LIMIT 1").await.unwrap();
+ server.await.unwrap();
+
+ assert_ne!(result.columns, vec!["status", "response"]);
+ let name_idx = result.columns.iter().position(|column| column == "name").unwrap();
+ assert_eq!(result.rows[0][name_idx], json!("Notebook"));
+ }
+
+ #[tokio::test]
+ async fn execute_rest_search_preserves_full_json_response() {
+ use tokio::io::{AsyncReadExt, AsyncWriteExt};
+
+ let body = json!({
+ "took": 3,
+ "hits": {
+ "total": { "value": 1, "relation": "eq" },
+ "max_score": 1.0,
+ "hits": [{
+ "_index": "products",
+ "_id": "product-1",
+ "_score": 1.0,
+ "_source": { "name": "Notebook", "price": 1299 },
+ "highlight": { "name": ["Notebook"] }
+ }]
+ },
+ "aggregations": {
+ "by_category": {
+ "buckets": [{ "key": "electronics", "doc_count": 1 }]
+ }
+ }
+ });
+ let response_body = body.to_string();
+ let server_response_body = response_body.clone();
+ let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let addr = listener.local_addr().unwrap();
+ let server = tokio::spawn(async move {
+ let (mut socket, _) = listener.accept().await.unwrap();
+ let mut request = [0_u8; 4096];
+ let read = socket.read(&mut request).await.unwrap();
+ let request = String::from_utf8_lossy(&request[..read]);
+ assert!(request.starts_with("POST /products/_search "));
+ let response = format!(
+ "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
+ server_response_body.len(),
+ server_response_body
+ );
+ socket.write_all(response.as_bytes()).await.unwrap();
+ });
+
+ let client = EsClient::new(&format!("http://{addr}"), None, None, false, Duration::from_secs(1));
+ let result =
+ super::execute_rest_query(&client, "POST /products/_search\n{\"query\":{\"match_all\":{}}}").await.unwrap();
+ server.await.unwrap();
+
+ assert_eq!(result.columns, vec!["status", "response"]);
+ assert_eq!(result.rows[0][0], json!(200));
+ let response = result.rows[0][1].as_str().unwrap();
+ assert_eq!(response, response_body);
+ assert_eq!(serde_json::from_str::(response).unwrap(), body);
+ }
+
#[test]
fn document_body_removes_elasticsearch_id_metadata() {
let doc = super::elasticsearch_document_body_from_json(r#"{"_id":"abc","_routing":"tenant-1","name":"Alice"}"#)
diff --git a/packages/app-tests/elasticsearchJsonResponse.test.ts b/packages/app-tests/elasticsearchJsonResponse.test.ts
new file mode 100644
index 000000000..e3eedb8f1
--- /dev/null
+++ b/packages/app-tests/elasticsearchJsonResponse.test.ts
@@ -0,0 +1,70 @@
+import { strict as assert } from "node:assert";
+import { test } from "vitest";
+import { elasticsearchJsonResponseForResult } from "../../apps/desktop/src/lib/elasticsearch/elasticsearchJsonResponse.ts";
+import type { QueryResult } from "../../apps/desktop/src/types/database.ts";
+
+function jsonResponse(overrides: Partial = {}): QueryResult {
+ return {
+ columns: ["status", "response"],
+ rows: [[200, '{\n "acknowledged": true\n}']],
+ affected_rows: 0,
+ execution_time_ms: 1,
+ ...overrides,
+ };
+}
+
+test("classifies Elasticsearch GET mapping and POST JSON responses", () => {
+ const mapping = jsonResponse();
+ assert.deepEqual(elasticsearchJsonResponseForResult("elasticsearch", "GET /products/_mapping", mapping), {
+ status: 200,
+ body: '{\n "acknowledged": true\n}',
+ });
+
+ const search = jsonResponse({ rows: [[201, '{"hits":{"hits":[]}}']] });
+ assert.deepEqual(elasticsearchJsonResponseForResult("elasticsearch", ' post /products/_search\n{"query":{"match_all":{}}}', search), {
+ status: 201,
+ body: '{"hits":{"hits":[]}}',
+ });
+});
+
+test("rejects SQL and non-JSON Elasticsearch result shapes", () => {
+ const response = jsonResponse();
+
+ assert.equal(elasticsearchJsonResponseForResult("elasticsearch", "SELECT * FROM products", response), undefined);
+ assert.equal(elasticsearchJsonResponseForResult("postgres", "GET /products/_mapping", response), undefined);
+ assert.equal(
+ elasticsearchJsonResponseForResult("elasticsearch", "GET /_cat/indices", {
+ columns: ["response"],
+ rows: [["green open products"]],
+ affected_rows: 1,
+ execution_time_ms: 1,
+ }),
+ undefined,
+ );
+});
+
+test("rejects invalid Elasticsearch JSON response status and row shapes", () => {
+ const invalidResults: QueryResult[] = [
+ jsonResponse({ columns: ["response", "status"] }),
+ jsonResponse({ rows: [[200, "{}", "unexpected"]] }),
+ jsonResponse({ rows: [[99, "{}"]] }),
+ jsonResponse({ rows: [[600, "{}"]] }),
+ jsonResponse({ rows: [["200", "{}"]] }),
+ jsonResponse({ rows: [[200, null]] }),
+ ];
+
+ for (const result of invalidResults) {
+ assert.equal(elasticsearchJsonResponseForResult("elasticsearch", "GET /products/_mapping", result), undefined);
+ }
+});
+
+test("uses the supplied result source statement to classify the response", () => {
+ const result = jsonResponse({ sourceStatement: "GET /products/_mapping" });
+
+ assert.deepEqual(elasticsearchJsonResponseForResult("elasticsearch", result.sourceStatement, result), {
+ status: 200,
+ body: '{\n "acknowledged": true\n}',
+ });
+ assert.equal(elasticsearchJsonResponseForResult("elasticsearch", "SELECT * FROM products", result), undefined);
+ assert.equal(elasticsearchJsonResponseForResult("elasticsearch", undefined, result), undefined);
+});
diff --git a/packages/app-tests/queryResultToolbar.test.ts b/packages/app-tests/queryResultToolbar.test.ts
index b237bb889..234306432 100644
--- a/packages/app-tests/queryResultToolbar.test.ts
+++ b/packages/app-tests/queryResultToolbar.test.ts
@@ -87,3 +87,12 @@ test("DataGrid marks toolbar refresh separately from current-result reloads", ()
assert.match(dataGrid, /emit\("reload", props\.sql,[^;]+"refresh"\);/);
assert.match(dataGrid, /function onToolbarRollback\(\)[\s\S]*?emit\("reload", props\.sql,[^;]+\);/);
});
+
+test("Elasticsearch JSON refresh preserves multi-result query groups", () => {
+ const contentArea = source(contentAreaPath);
+
+ assert.match(
+ contentArea,
+ /if \(activeElasticsearchJsonResponse\.value\) \{[\s\S]*?emit\("reload", activeResultSql\.value, undefined, undefined, undefined, undefined, undefined, "refresh"\);/,
+ );
+});