feat(nebula): add NebulaGraph 3.x browser support

This commit is contained in:
Hand Sonic
2026-09-29 21:03:12 +08:00
committed by GitHub
parent e59a8c1bd4
commit ef05a6a19b
69 changed files with 1737 additions and 41 deletions
+2 -1
View File
@@ -57,6 +57,7 @@ const nativeDriverDirectories = {
kingbase: "kingbase-go",
iotdb: "iotdb",
neo4j: "neo4j-go",
nebula: "nebula-go",
vastbase: "vastbase-go",
rabbitmq: "rabbitmq",
rocketmq: "rocketmq",
@@ -68,7 +69,7 @@ const nativeDriverDirectories = {
const crateNativeDriverDirectories = {
"sqlite-worker": "crates/dbx-sqlite-worker",
};
const nativeDriverModules = new Set(["cassandra", "duckdb", "hive", "argo", "oracle", "xugu", "kingbase", "iotdb", "neo4j", "vastbase", "rabbitmq", "rocketmq", "zookeeper", "tdengine", "etcd", "etcd2", "sqlite-worker"]);
const nativeDriverModules = new Set(["cassandra", "duckdb", "hive", "argo", "oracle", "xugu", "kingbase", "iotdb", "neo4j", "nebula", "vastbase", "rabbitmq", "rocketmq", "zookeeper", "tdengine", "etcd", "etcd2", "sqlite-worker"]);
const nativeDriverSharedPaths = {
hive: [
"agents/go-common/go-gssapi",
+1
View File
@@ -21,6 +21,7 @@ export const goAgents = [
{ driver: "hive-go", binary: "hive", race: false },
{ driver: "vastbase-go", binary: "vastbase", race: false },
{ driver: "neo4j-go", binary: "neo4j", race: false },
{ driver: "nebula-go", binary: "nebula", race: false },
{ driver: "iotdb", binary: "iotdb", race: false },
];
+3 -3
View File
@@ -93,7 +93,7 @@ for (const file of ["Cargo.toml", "Cargo.lock", ".cargo/config.toml", "rust-tool
const result = plan([file]);
assert.deepEqual(groups(result), ["workspace"]);
assert.equal(result.rust_full, true);
assert.equal(result.agent_go.include.length, 10);
assert.equal(result.agent_go.include.length, goAgents.length);
assert.equal(result.agent_rust.include.length, 2);
assert.equal(result.agent_integration.include.length, 16);
assert.equal(result.agent_java, true);
@@ -150,12 +150,12 @@ test("shared Agent inputs and unknown native modules never silently lose coverag
for (const file of ["agents/common/src/main/java/Protocol.java", "agents/scripts/validate_agents.py", "agents/build.gradle",
"agents/drivers/new-driver/main.go", "crates/dbx-driver-agent/assets/agent-protocol-v2.json", ".github/workflows/agents-release.yml"]) {
const result = plan([file]);
assert.equal(result.agent_go.include.length, 10, file);
assert.equal(result.agent_go.include.length, goAgents.length, file);
assert.equal(result.agent_integration.include.length, 16, file);
assert.equal(result.agent_java, true, file);
}
const fallback = plan(["future-agent-filter-input"], { agentsChanged: true });
assert.equal(fallback.agent_go.include.length, 10);
assert.equal(fallback.agent_go.include.length, goAgents.length);
assert.equal(fallback.agent_integration.include.length, 16);
assert.equal(fallback.agent_java, true);
});
+1
View File
@@ -56,6 +56,7 @@ const DRIVER_DATABASE_ALIASES = {
kafka: "mq",
"kingbase-go": "kingbase",
"neo4j-go": "neo4j",
"nebula-go": "nebula",
"oracle-10g": "oracle",
"oracle-go": "oracle",
"oracle-legacy": "oracle",
@@ -15,7 +15,7 @@ import { basename, join } from "node:path";
import { tmpdir } from "node:os";
const REGISTRY_ASSET = "agent-registry.json";
const NATIVE_MODULES = new Set(["duckdb", "oracle", "xugu", "kingbase", "iotdb", "neo4j", "vastbase", "rabbitmq", "rocketmq", "zookeeper", "tdengine", "etcd", "etcd2"]);
const NATIVE_MODULES = new Set(["duckdb", "oracle", "xugu", "kingbase", "iotdb", "neo4j", "nebula", "vastbase", "rabbitmq", "rocketmq", "zookeeper", "tdengine", "etcd", "etcd2"]);
const PLATFORMS = [
"macos-aarch64",
"macos-x64",
+44 -2
View File
@@ -707,6 +707,44 @@ jobs:
name: neo4j-native
path: "release-native/dbx-agent-neo4j-*"
build-nebula-native:
needs: [bump-versions]
if: ${{ contains(fromJSON(needs.bump-versions.outputs.native_modules), 'nebula') }}
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version: "1.24.x"
- name: Test NebulaGraph native agent
working-directory: agents/drivers/nebula-go
run: go test ./...
- name: Cross-compile NebulaGraph native agent
shell: bash
run: |
mkdir -p release-native
cd agents/drivers/nebula-go
declare -A TARGETS=(
["macos-aarch64"]="darwin/arm64"
["macos-x64"]="darwin/amd64"
["linux-aarch64"]="linux/arm64"
["linux-x64"]="linux/amd64"
["windows-aarch64"]="windows/arm64"
["windows-x64"]="windows/amd64"
)
for platform in "${!TARGETS[@]}"; do
IFS=/ read -r goos goarch <<< "${TARGETS[$platform]}"
output="../../../release-native/dbx-agent-nebula-${platform}"
if [[ "$goos" == "windows" ]]; then
output="${output}.exe"
fi
CGO_ENABLED=0 GOOS="$goos" GOARCH="$goarch" go build -trimpath -ldflags="-s -w" -o "$output" .
done
- uses: actions/upload-artifact@v4
with:
name: nebula-native
path: "release-native/dbx-agent-nebula-*"
build-iotdb-native:
needs: [bump-versions]
if: ${{ contains(fromJSON(needs.bump-versions.outputs.native_modules), 'iotdb') }}
@@ -1257,7 +1295,7 @@ jobs:
retention-days: 1
release:
needs: [bump-versions, commit-versions, build-agents, build-oracle-native, build-xugu-native, build-rabbitmq-native, build-rocketmq-native, build-etcd-native, build-etcd2-native, build-zookeeper-native, build-cassandra-native, build-hive-native, build-argo-native, build-kingbase-native, build-vastbase-native, build-neo4j-native, build-iotdb-native, build-duckdb-native, build-sqlite-worker-native, build-tdengine-native, build-jre, reuse-previous-assets]
needs: [bump-versions, commit-versions, build-agents, build-oracle-native, build-xugu-native, build-rabbitmq-native, build-rocketmq-native, build-etcd-native, build-etcd2-native, build-zookeeper-native, build-cassandra-native, build-hive-native, build-argo-native, build-kingbase-native, build-vastbase-native, build-neo4j-native, build-nebula-native, build-iotdb-native, build-duckdb-native, build-sqlite-worker-native, build-tdengine-native, build-jre, reuse-previous-assets]
if: ${{ always() && !contains(needs.*.result, 'failure') && !contains(needs.*.result, 'cancelled') }}
runs-on: ubuntu-latest
steps:
@@ -1391,6 +1429,7 @@ jobs:
hive) echo "Apache Hive" ;;
argo) echo "星环Argo" ;;
neo4j) echo "Neo4j" ;;
nebula) echo "NebulaGraph" ;;
iotdb) echo "Apache IoTDB" ;;
tdengine) echo "TDengine" ;;
etcd) echo "etcd" ;;
@@ -1458,7 +1497,7 @@ jobs:
[ -n "$DRIVERS" ] && DRIVERS="${DRIVERS},"$'\n'
DRIVERS="${DRIVERS}$(generate_jar_entry "$name" "$label" "$f" "$jre_key" "$version" "$external_driver" "$native_json")"
done
for name in oracle xugu kingbase vastbase neo4j iotdb duckdb sqlite-worker rabbitmq rocketmq zookeeper cassandra hive argo tdengine etcd etcd2; do
for name in oracle xugu kingbase vastbase neo4j nebula iotdb duckdb sqlite-worker rabbitmq rocketmq zookeeper cassandra hive argo tdengine etcd etcd2; do
version=$(get_module_version "$name")
[ -f "release/dbx-agent-${name}-${version}.jar" ] && continue
native_json=$(generate_native_platforms "$name" "$version")
@@ -1548,6 +1587,7 @@ jobs:
hive) echo "Apache Hive" ;;
argo) echo "星环Argo" ;;
neo4j) echo "Neo4j" ;;
nebula) echo "NebulaGraph" ;;
iotdb) echo "Apache IoTDB" ;;
tdengine) echo "TDengine" ;;
etcd) echo "etcd" ;;
@@ -1600,6 +1640,8 @@ jobs:
LOG_PATH="agents/drivers/argo-go/"
elif [ "$name" = "neo4j" ]; then
LOG_PATH="agents/drivers/neo4j-go/"
elif [ "$name" = "nebula" ]; then
LOG_PATH="agents/drivers/nebula-go/"
elif [ "$name" = "iotdb" ]; then
LOG_PATH="agents/drivers/iotdb/"
elif [ "$name" = "sqlite-worker" ]; then
+1
View File
@@ -35,6 +35,7 @@ Each agent runs as a standalone process and communicates with DBX via stdin/stdo
| db2 | IBM DB2 | DB2 JDBC |
| informix | IBM Informix | Informix JDBC |
| neo4j | Neo4j | Official Neo4j Go Driver native agent |
| nebula | NebulaGraph 3.x | Official NebulaGraph Go Client native agent |
| cassandra | Apache Cassandra 2.1+ | Apache cassandra-gocql-driver native agent |
| bigquery | Google BigQuery | BigQuery JDBC |
| spanner | Google Cloud Spanner | Google Cloud Spanner JDBC |
+1
View File
@@ -35,6 +35,7 @@ DBX 的 Agent 驱动 —— 通过 JDBC 和原生数据库驱动支持各种数
| db2 | IBM DB2 | DB2 JDBC |
| informix | IBM Informix | Informix JDBC |
| neo4j | Neo4j | 官方 Neo4j Go Driver 原生 Agent |
| nebula | NebulaGraph 3.x | 官方 NebulaGraph Go Client 原生 Agent |
| cassandra | Apache Cassandra 2.1+ | Apache cassandra-gocql-driver 原生 Agent |
| bigquery | Google BigQuery | BigQuery JDBC |
| spanner | Google Cloud Spanner | Google Cloud Spanner JDBC |
+11
View File
@@ -0,0 +1,11 @@
module github.com/t8y2/dbx/agents/drivers/nebula-go
go 1.24
require github.com/vesoft-inc/nebula-go/v3 v3.8.0
require (
github.com/vesoft-inc/fbthrift v0.0.0-20230214024353-fa2f34755b28 // indirect
golang.org/x/net v0.17.0 // indirect
golang.org/x/text v0.13.0 // indirect
)
+55
View File
@@ -0,0 +1,55 @@
github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.7.0 h1:nwc3DEeHmmLAfoZucVR881uASk0Mfjw8xYJ99tb5CcY=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/vesoft-inc/fbthrift v0.0.0-20230214024353-fa2f34755b28 h1:gpoPCGeOEuk/TnoY9nLVK1FoBM5ie7zY3BPVG8q43ME=
github.com/vesoft-inc/fbthrift v0.0.0-20230214024353-fa2f34755b28/go.mod h1:xu7e9za8StcJhBZmCDwK1Hyv4/Y0xFsjS+uqp10ECJg=
github.com/vesoft-inc/nebula-go/v3 v3.8.0 h1:ecB87KMnMUcuKbgFESKIscdxA7Y1TcX7XEVqZQ1UqlA=
github.com/vesoft-inc/nebula-go/v3 v3.8.0/go.mod h1:YTNAQzimjXLXUaEDOzty/eCCye+9zkZRuUzXz9LQUpU=
github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY=
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc=
golang.org/x/crypto v0.14.0/go.mod h1:MVFd36DqK4CsrnJYDkBA3VC4m2GkXAM0PvzMCn4JQf4=
golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4=
golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg=
golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c=
golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs=
golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg=
golang.org/x/net v0.17.0 h1:pVaXccu2ozPjCXewfr1S7xza/zcXTity9cCdXQYSjIM=
golang.org/x/net v0.17.0/go.mod h1:NxSsAGuq816PNPmqtQdLE42eU2Fs7NoRIZrHJAlaCOE=
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.13.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8=
golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k=
golang.org/x/term v0.8.0/go.mod h1:xPskH00ivmX89bAKVGSKKtLOWNx2+17Eiy94tnKShWo=
golang.org/x/term v0.13.0/go.mod h1:LTmsnFJwVN6bCy1rVCoS+qHT1HhALEFxKncY3WNNh4U=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ=
golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8=
golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8=
golang.org/x/text v0.13.0 h1:ablQoSUd0tRdKxZewP80B+BaqeKJuVhuRxj/dkrun3k=
golang.org/x/text v0.13.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE=
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc=
golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU=
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c h1:dUUwHk2QECo/6vqA44rthZ8ie2QXMNeKRTHCNY2nXvo=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+347
View File
@@ -0,0 +1,347 @@
package main
import (
"bufio"
"encoding/json"
"errors"
"fmt"
"os"
"strings"
"sync"
"sync/atomic"
"time"
nebula "github.com/vesoft-inc/nebula-go/v3"
)
const (
legacySessionID = "__legacy__"
maxSessions = 256
)
type request struct {
ID json.RawMessage `json:"id"`
Method string `json:"method"`
Params map[string]json.RawMessage `json:"params"`
}
type response struct {
JSONRPC string `json:"jsonrpc"`
ID json.RawMessage `json:"id,omitempty"`
Result any `json:"result,omitempty"`
Error *rpcError `json:"error,omitempty"`
}
type connectParams struct {
DriverProfile string `json:"driver_profile"`
Host string `json:"host"`
Port int `json:"port"`
Database string `json:"database"`
Username string `json:"username"`
Password string `json:"password"`
SSL bool `json:"ssl"`
CACertPath string `json:"ca_cert_path"`
ClientCertPath string `json:"client_cert_path"`
ClientKeyPath string `json:"client_key_path"`
ConnectTimeoutSecs int `json:"connect_timeout_secs"`
}
type agentSession struct {
mu sync.Mutex
pool *nebula.ConnectionPool
conn *nebula.Session
params connectParams
cursors map[string]*queryCursor
nextCursorID uint64
closed bool
canceled atomic.Bool
closePool sync.Once
}
type runtimeServer struct {
mu sync.RWMutex
sessions map[string]*agentSession
}
func main() {
runtime := &runtimeServer{sessions: make(map[string]*agentSession)}
encoder := json.NewEncoder(os.Stdout)
var output sync.Mutex
var running sync.WaitGroup
fmt.Fprintln(os.Stdout, `{"ready":true}`)
scanner := bufio.NewScanner(os.Stdin)
scanner.Buffer(make([]byte, 0, 64*1024), 512*1024*1024)
for scanner.Scan() {
line := strings.TrimSpace(scanner.Text())
if line == "" {
continue
}
var envelope request
if json.Unmarshal([]byte(line), &envelope) == nil && envelope.Method == "shutdown" {
running.Wait()
result := runtime.handleLine(line)
output.Lock()
_ = encoder.Encode(result)
output.Unlock()
return
}
running.Add(1)
go func(line string) {
defer running.Done()
result := runtime.handleLine(line)
output.Lock()
defer output.Unlock()
if err := encoder.Encode(result); err != nil {
fmt.Fprintln(os.Stderr, err)
}
}(line)
}
running.Wait()
if err := scanner.Err(); err != nil {
fmt.Fprintln(os.Stderr, err)
}
_ = runtime.closeAll()
}
func (r *runtimeServer) handleLine(line string) response {
var req request
if err := json.Unmarshal([]byte(line), &req); err != nil {
return errorResponse(nil, "", "", err)
}
if len(req.ID) == 0 {
req.ID = json.RawMessage("1")
}
result, err := r.dispatch(req.Method, req.Params)
if err != nil {
return errorResponse(req.ID, req.Method, stringParam(req.Params, "agentSessionId"), err)
}
return response{JSONRPC: "2.0", ID: req.ID, Result: result}
}
func (r *runtimeServer) dispatch(method string, params map[string]json.RawMessage) (any, error) {
switch method {
case "handshake":
return map[string]any{
"protocolVersion": 2, "agentProtocolVersion": 2,
"capabilities": []string{"connect", "test_connection", "metadata", "query", "paged_query", "ddl", "multi_session", "structured_error_v1"},
}, nil
case "open_session", "connect":
var options connectParams
if err := decodeParams(params, &options); err != nil {
return nil, err
}
id := stringParam(params, "agentSessionId")
if method == "connect" {
id = legacySessionID
_ = r.closeSession(id)
}
if id == "" {
return nil, errors.New("agentSessionId is required")
}
return map[string]bool{"ok": true}, r.openSession(id, options)
case "test_connection":
var options connectParams
if err := decodeParams(params, &options); err != nil {
return nil, err
}
session, err := newAgentSession(options)
if err != nil {
return nil, err
}
defer session.close()
return map[string]bool{"ok": true}, nil
case "close_session", "disconnect":
id := stringParam(params, "agentSessionId")
if method == "disconnect" {
id = legacySessionID
}
return map[string]bool{"ok": true}, r.closeSession(id)
case "shutdown":
return map[string]bool{"ok": true}, r.closeAll()
case "cancel_session":
session, err := r.session(stringParam(params, "agentSessionId"))
if err != nil {
return nil, err
}
session.cancel()
return map[string]bool{"ok": true}, nil
}
id := stringParam(params, "agentSessionId")
if id == "" {
id = legacySessionID
}
session, err := r.session(id)
if err != nil {
return nil, err
}
session.mu.Lock()
defer session.mu.Unlock()
if session.closed || session.canceled.Load() {
return nil, errors.New("agent session is closed")
}
if method == "validate_session" || method == "validate_connection" {
_, err := session.execute("SHOW SPACES", "")
return map[string]bool{"ok": err == nil}, err
}
return session.dispatch(method, params)
}
func (r *runtimeServer) openSession(id string, options connectParams) error {
r.mu.Lock()
if _, exists := r.sessions[id]; exists {
r.mu.Unlock()
return fmt.Errorf("agent session already exists: %s", id)
}
if len(r.sessions) >= maxSessions {
r.mu.Unlock()
return errors.New("agent session limit reached")
}
r.mu.Unlock()
session, err := newAgentSession(options)
if err != nil {
return err
}
r.mu.Lock()
defer r.mu.Unlock()
if _, exists := r.sessions[id]; exists {
session.close()
return fmt.Errorf("agent session already exists: %s", id)
}
if len(r.sessions) >= maxSessions {
session.close()
return errors.New("agent session limit reached")
}
r.sessions[id] = session
return nil
}
func newAgentSession(options connectParams) (*agentSession, error) {
if options.DriverProfile != "" && !strings.EqualFold(options.DriverProfile, "nebula") && !strings.EqualFold(options.DriverProfile, "nebula-v3") {
return nil, fmt.Errorf("unsupported NebulaGraph driver profile %q; only NebulaGraph 3.x is supported", options.DriverProfile)
}
if strings.TrimSpace(options.Host) == "" {
return nil, errors.New("NebulaGraph graphd host is required")
}
if options.Port <= 0 || options.Port > 65535 {
return nil, errors.New("NebulaGraph graphd port is invalid")
}
config := nebula.GetDefaultConf()
config.MaxConnPoolSize = 1
config.MinConnPoolSize = 1
if options.ConnectTimeoutSecs > 0 {
config.TimeOut = time.Duration(options.ConnectTimeoutSecs) * time.Second
} else {
config.TimeOut = 15 * time.Second
}
hosts := []nebula.HostAddress{{Host: options.Host, Port: options.Port}}
var pool *nebula.ConnectionPool
var err error
if options.SSL {
if options.CACertPath == "" || options.ClientCertPath == "" || options.ClientKeyPath == "" {
return nil, errors.New("NebulaGraph TLS requires CA, client certificate and client key paths")
}
tlsConfig, tlsErr := nebula.GetDefaultSSLConfig(options.CACertPath, options.ClientCertPath, options.ClientKeyPath)
if tlsErr != nil {
return nil, tlsErr
}
pool, err = nebula.NewSslConnectionPool(hosts, config, tlsConfig, nebula.DefaultLogger{})
} else {
pool, err = nebula.NewConnectionPool(hosts, config, nebula.DefaultLogger{})
}
if err != nil {
return nil, err
}
conn, err := pool.GetSession(options.Username, options.Password)
if err != nil {
pool.Close()
return nil, err
}
session := &agentSession{pool: pool, conn: conn, params: options, cursors: make(map[string]*queryCursor)}
if _, err := session.execute("SHOW SPACES", ""); err != nil {
session.close()
return nil, err
}
return session, nil
}
func (r *runtimeServer) session(id string) (*agentSession, error) {
r.mu.RLock()
session := r.sessions[id]
r.mu.RUnlock()
if session == nil {
return nil, fmt.Errorf("agent session not found: %s", id)
}
return session, nil
}
func (r *runtimeServer) closeSession(id string) error {
r.mu.Lock()
session := r.sessions[id]
delete(r.sessions, id)
r.mu.Unlock()
if session != nil {
session.mu.Lock()
session.close()
session.mu.Unlock()
}
return nil
}
func (r *runtimeServer) closeAll() error {
r.mu.Lock()
sessions := r.sessions
r.sessions = make(map[string]*agentSession)
r.mu.Unlock()
for _, session := range sessions {
session.mu.Lock()
session.close()
session.mu.Unlock()
}
return nil
}
func (s *agentSession) cancel() {
s.canceled.Store(true)
s.closePool.Do(func() { s.pool.Close() })
}
func (s *agentSession) close() {
if s.closed {
return
}
s.closed = true
s.cursors = nil
s.closePool.Do(func() {
if !s.canceled.Load() {
s.conn.Release()
}
s.pool.Close()
})
}
func decodeParams(params map[string]json.RawMessage, target any) error {
data, err := json.Marshal(params)
if err != nil {
return err
}
return json.Unmarshal(data, target)
}
func stringParam(params map[string]json.RawMessage, key string) string {
var value string
_ = json.Unmarshal(params[key], &value)
return value
}
func intParam(params map[string]json.RawMessage, key string) int {
var value int
_ = json.Unmarshal(params[key], &value)
return value
}
func stringSliceParam(params map[string]json.RawMessage, key string) []string {
var values []string
_ = json.Unmarshal(params[key], &values)
return values
}
+135
View File
@@ -0,0 +1,135 @@
package main
import (
"encoding/json"
"reflect"
"strings"
"testing"
nebula "github.com/vesoft-inc/nebula-go/v3"
wire "github.com/vesoft-inc/nebula-go/v3/nebula"
"github.com/vesoft-inc/nebula-go/v3/nebula/graph"
)
func resultSet(t *testing.T, names []string, rows [][]*wire.Value) *nebula.ResultSet {
t.Helper()
columns := make([][]byte, len(names))
for index, name := range names {
columns[index] = []byte(name)
}
data := &wire.DataSet{ColumnNames: columns}
for _, row := range rows {
data.Rows = append(data.Rows, &wire.Row{Values: row})
}
result, err := nebula.GenResultSet(&graph.ExecutionResponse{ErrorCode: wire.ErrorCode_SUCCEEDED, Data: data})
if err != nil {
t.Fatal(err)
}
return result
}
func TestHandshakeMatchesAgentProtocol(t *testing.T) {
runtime := &runtimeServer{sessions: make(map[string]*agentSession)}
response := runtime.handleLine(`{"jsonrpc":"2.0","id":7,"method":"handshake","params":{}}`)
if response.Error != nil {
t.Fatalf("unexpected handshake error: %#v", response.Error)
}
result := response.Result.(map[string]any)
if result["protocolVersion"] != 2 {
t.Fatalf("unexpected protocol: %#v", result)
}
capabilities := result["capabilities"].([]string)
if !strings.Contains(strings.Join(capabilities, ","), "multi_session") {
t.Fatalf("missing multi-session capability: %#v", capabilities)
}
}
func TestDriverProfileSupportsV3AndLegacyConnections(t *testing.T) {
for _, profile := range []string{"", "nebula", "nebula-v3"} {
_, err := newAgentSession(connectParams{DriverProfile: profile})
if err == nil || !strings.Contains(err.Error(), "graphd host is required") {
t.Fatalf("profile %q unexpectedly rejected: %v", profile, err)
}
}
_, err := newAgentSession(connectParams{DriverProfile: "nebula-v5"})
if err == nil || !strings.Contains(err.Error(), "only NebulaGraph 3.x is supported") {
t.Fatalf("unsupported profile was not rejected: %v", err)
}
}
func TestMissingSessionReturnsValidErrorContract(t *testing.T) {
runtime := &runtimeServer{sessions: make(map[string]*agentSession)}
response := runtime.handleLine(`{"jsonrpc":"2.0","id":7,"method":"execute_query","params":{"agentSessionId":"missing","sql":"SHOW SPACES"}}`)
if response.Error == nil || response.Error.Data == nil {
t.Fatalf("missing structured error: %#v", response)
}
data := response.Error.Data
if data.ContractVersion != 1 || data.Stage != "execute" || data.OperationOutcome != "unknown" || data.AgentSessionID != "missing" || data.Category != "protocol" {
t.Fatalf("invalid error contract: %#v", data)
}
if _, err := json.Marshal(response); err != nil {
t.Fatal(err)
}
}
func TestQueryErrorClassifiesAsSQLWithoutSQLState(t *testing.T) {
response := errorResponse(json.RawMessage("1"), "execute_query", "session-1", &nebulaQueryError{code: "-1004", message: "syntax error"})
if response.Error.Data.Category != "sql" || response.Error.Data.SessionDisposition != "keep" {
t.Fatalf("unexpected SQL error: %#v", response.Error)
}
if !strings.Contains(response.Error.Message, "-1004") {
t.Fatalf("NebulaGraph error code missing: %q", response.Error.Message)
}
}
func TestQuoteNebulaIdentifier(t *testing.T) {
quoted, err := quoteNebulaIdentifier("edge`name\\part")
if err != nil || quoted != "`edge\\`name\\\\part`" {
t.Fatalf("quoted = %q, err = %v", quoted, err)
}
if _, err := quoteNebulaIdentifier("a\nb"); err == nil {
t.Fatal("control characters must not be accepted")
}
}
func TestReadRowsPreservesScalarAndGraphValues(t *testing.T) {
count := int64(42)
nullValue := wire.NullType___NULL__
value := resultSet(t, []string{"name", "count", "missing"}, [][]*wire.Value{{
{SVal: []byte("Ada")}, {IVal: &count}, {NVal: &nullValue},
}})
rows, types, err := readRows(value, 10)
if err != nil {
t.Fatal(err)
}
if !reflect.DeepEqual(rows, [][]any{{"Ada", "42", nil}}) || !reflect.DeepEqual(types, []string{"string", "int", "Unknown"}) {
t.Fatalf("rows=%#v types=%#v", rows, types)
}
metadata, err := metadataRows(value)
if err != nil || metadata[0]["name"] != "Ada" || metadata[0]["count"] != "42" {
t.Fatalf("metadata=%#v err=%v", metadata, err)
}
}
func TestFetchQueryPageAdvancesWithoutReexecuting(t *testing.T) {
session := &agentSession{cursors: map[string]*queryCursor{
"page-1": {result: queryResult{Columns: []string{"n"}, ColumnTypes: []string{"int"}, Rows: [][]any{{"1"}, {"2"}, {"3"}}}, offset: 1},
}}
page, err := session.fetchQueryPage("page-1", 1)
if err != nil || !page.HasMore || !reflect.DeepEqual(page.Rows, [][]any{{"2"}}) {
t.Fatalf("page=%#v err=%v", page, err)
}
last, err := session.fetchQueryPage("page-1", 1)
if err != nil || last.HasMore || !reflect.DeepEqual(last.Rows, [][]any{{"3"}}) || len(session.cursors) != 0 {
t.Fatalf("last=%#v err=%v cursors=%d", last, err, len(session.cursors))
}
}
func TestMetadataKindsRemainSeparate(t *testing.T) {
if !acceptsObjectType([]string{"TABLE"}, "TABLE") || acceptsObjectType([]string{"TABLE"}, "VIEW") {
t.Fatal("tag and edge types were mixed")
}
if !acceptsObjectType([]string{"EDGE"}, "VIEW") || !acceptsObjectType(nil, "TABLE") {
t.Fatal("edge alias or unrestricted metadata lookup failed")
}
}
+282
View File
@@ -0,0 +1,282 @@
package main
import (
"encoding/json"
"errors"
"fmt"
"sort"
"strings"
"unicode"
nebula "github.com/vesoft-inc/nebula-go/v3"
)
type databaseInfo struct {
Name string `json:"name"`
}
type tableInfo struct {
Name string `json:"name"`
TableType string `json:"table_type"`
Comment *string `json:"comment"`
}
type objectInfo struct {
Name string `json:"name"`
ObjectType string `json:"object_type"`
Schema string `json:"schema"`
Comment *string `json:"comment"`
}
type columnInfo struct {
Name string `json:"name"`
DataType string `json:"data_type"`
IsNullable bool `json:"is_nullable"`
ColumnDefault *string `json:"column_default"`
IsPrimaryKey bool `json:"is_primary_key"`
Extra *string `json:"extra"`
Comment *string `json:"comment"`
NumericPrecision *int `json:"numeric_precision"`
NumericScale *int `json:"numeric_scale"`
CharacterMaximumLength *int `json:"character_maximum_length"`
}
type objectSource struct {
Name string `json:"name"`
ObjectType string `json:"object_type"`
Schema *string `json:"schema"`
Source string `json:"source"`
}
func (s *agentSession) listDatabases() ([]databaseInfo, error) {
result, err := s.execute("SHOW SPACES", "")
if err != nil {
return nil, err
}
rows, err := metadataRows(result)
if err != nil {
return nil, err
}
spaces := make([]databaseInfo, 0, len(rows))
for _, row := range rows {
name := firstMetadataValue(row, "Name")
if name != "" {
spaces = append(spaces, databaseInfo{Name: name})
}
}
sort.Slice(spaces, func(i, j int) bool { return spaces[i].Name < spaces[j].Name })
return spaces, nil
}
func (s *agentSession) listTables(params map[string]json.RawMessage) ([]tableInfo, error) {
space := selectedSpace(stringParam(params, "database"), s.params.Database)
if space == "" {
return []tableInfo{}, nil
}
filter := strings.ToLower(strings.TrimSpace(stringParam(params, "filter")))
requestedTypes := stringSliceParam(params, "object_types")
if len(requestedTypes) == 0 {
requestedTypes = stringSliceParam(params, "objectTypes")
}
tables := make([]tableInfo, 0)
for _, kind := range []struct{ typeName, statement string }{{"TABLE", "SHOW TAGS"}, {"VIEW", "SHOW EDGES"}} {
if !acceptsObjectType(requestedTypes, kind.typeName) {
continue
}
result, err := s.execute(kind.statement, space)
if err != nil {
return nil, err
}
rows, err := metadataRows(result)
if err != nil {
return nil, err
}
for _, row := range rows {
name := firstMetadataValue(row, "Name")
if name != "" && strings.Contains(strings.ToLower(name), filter) {
tables = append(tables, tableInfo{Name: name, TableType: kind.typeName})
}
}
}
sort.Slice(tables, func(i, j int) bool {
if tables[i].TableType != tables[j].TableType {
return tables[i].TableType < tables[j].TableType
}
return tables[i].Name < tables[j].Name
})
offset, limit := intParam(params, "offset"), intParam(params, "limit")
if offset < 0 {
offset = 0
}
if offset >= len(tables) {
return []tableInfo{}, nil
}
tables = tables[offset:]
if limit > 0 && limit < len(tables) {
tables = tables[:limit]
}
return tables, nil
}
func (s *agentSession) listObjects(params map[string]json.RawMessage) ([]objectInfo, error) {
tables, err := s.listTables(params)
if err != nil {
return nil, err
}
objects := make([]objectInfo, 0, len(tables))
for _, table := range tables {
objects = append(objects, objectInfo{Name: table.Name, ObjectType: table.TableType, Schema: ""})
}
return objects, nil
}
func (s *agentSession) getColumns(database, name string) ([]columnInfo, error) {
space := selectedSpace(database, s.params.Database)
if space == "" {
return nil, errors.New("select a NebulaGraph space before reading object columns")
}
quoted, err := quoteNebulaIdentifier(name)
if err != nil {
return nil, err
}
result, err := s.execute("DESCRIBE TAG "+quoted, space)
if err != nil {
result, err = s.execute("DESCRIBE EDGE "+quoted, space)
if err != nil {
return nil, err
}
}
rows, err := metadataRows(result)
if err != nil {
return nil, err
}
columns := make([]columnInfo, 0, len(rows))
for _, row := range rows {
field := firstMetadataValue(row, "Field")
if field == "" {
continue
}
columns = append(columns, columnInfo{
Name: field, DataType: firstMetadataValue(row, "Type"),
IsNullable: !strings.EqualFold(firstMetadataValue(row, "Null"), "NO"),
ColumnDefault: optionalMetadataValue(row, "Default"), Comment: optionalMetadataValue(row, "Comment"),
})
}
return columns, nil
}
func (s *agentSession) getTableDDL(database, name, objectType string) (string, error) {
space := selectedSpace(database, s.params.Database)
if space == "" {
return "", errors.New("select a NebulaGraph space before reading object DDL")
}
quoted, err := quoteNebulaIdentifier(name)
if err != nil {
return "", err
}
kinds := []string{"TAG", "EDGE"}
if strings.EqualFold(objectType, "VIEW") || strings.EqualFold(objectType, "EDGE") {
kinds = []string{"EDGE"}
}
for _, kind := range kinds {
result, queryErr := s.execute("SHOW CREATE "+kind+" "+quoted, space)
if queryErr != nil {
if kind == kinds[len(kinds)-1] {
return "", queryErr
}
continue
}
rows, readErr := metadataRows(result)
if readErr != nil {
return "", readErr
}
if len(rows) > 0 {
for key, value := range rows[0] {
if strings.HasPrefix(strings.ToLower(key), "create ") {
return value, nil
}
}
}
}
return "", nil
}
func (s *agentSession) getObjectSource(params map[string]json.RawMessage) (objectSource, error) {
name := stringParam(params, "name")
objectType := strings.ToUpper(stringParam(params, "object_type"))
if objectType == "" {
objectType = strings.ToUpper(stringParam(params, "objectType"))
}
source, err := s.getTableDDL(stringParam(params, "database"), name, objectType)
if err != nil {
return objectSource{}, err
}
return objectSource{Name: name, ObjectType: objectType, Source: source}, nil
}
func selectedSpace(requested, configured string) string {
if strings.TrimSpace(requested) != "" {
return strings.TrimSpace(requested)
}
return strings.TrimSpace(configured)
}
func quoteNebulaIdentifier(name string) (string, error) {
if name == "" || strings.IndexFunc(name, unicode.IsControl) >= 0 {
return "", fmt.Errorf("invalid NebulaGraph identifier: %q", name)
}
escaped := strings.ReplaceAll(strings.ReplaceAll(name, `\`, `\\`), "`", "\\`")
return "`" + escaped + "`", nil
}
func acceptsObjectType(requested []string, kind string) bool {
if len(requested) == 0 {
return true
}
for _, candidate := range requested {
if strings.EqualFold(candidate, kind) || (kind == "TABLE" && strings.EqualFold(candidate, "TAG")) || (kind == "VIEW" && strings.EqualFold(candidate, "EDGE")) {
return true
}
}
return false
}
func metadataRows(result *nebula.ResultSet) ([]map[string]string, error) {
columns := result.GetColNames()
rows := make([]map[string]string, 0, len(result.GetRows()))
for index := range result.GetRows() {
record, err := result.GetRowValuesByIndex(index)
if err != nil {
return nil, err
}
row := make(map[string]string, len(columns))
for column, name := range columns {
value, err := record.GetValueByIndex(column)
if err != nil {
return nil, err
}
if normalized := normalizeValue(value); normalized != nil {
row[name] = fmt.Sprint(normalized)
}
}
rows = append(rows, row)
}
return rows, nil
}
func firstMetadataValue(row map[string]string, name string) string {
for key, value := range row {
if strings.EqualFold(key, name) {
return value
}
}
return ""
}
func optionalMetadataValue(row map[string]string, name string) *string {
value := firstMetadataValue(row, name)
if value == "" || strings.EqualFold(value, "NULL") {
return nil
}
return &value
}
@@ -0,0 +1,79 @@
package main
import (
"encoding/json"
"errors"
"io"
"net"
"strings"
)
type rpcError struct {
Code int `json:"code"`
Message string `json:"message"`
Data *rpcErrorData `json:"data"`
}
type rpcErrorData struct {
ContractVersion int `json:"contractVersion"`
Category string `json:"category"`
Retryable bool `json:"retryable"`
SessionDisposition string `json:"sessionDisposition"`
Stage string `json:"stage"`
OperationOutcome string `json:"operationOutcome"`
AgentSessionID string `json:"agentSessionId,omitempty"`
}
func errorResponse(id json.RawMessage, method, sessionID string, err error) response {
stage := errorStage(method)
data := &rpcErrorData{
ContractVersion: 1, Category: "protocol", SessionDisposition: "keep",
Stage: stage, OperationOutcome: errorOutcome(stage), AgentSessionID: strings.TrimSpace(sessionID),
}
var queryError *nebulaQueryError
if errors.As(err, &queryError) {
if stage == "execute" || stage == "fetch" {
data.Category = "sql"
} else {
data.Category = "connection"
}
} else if errors.Is(err, io.EOF) || isNetworkError(err) {
data.Category = "connection"
data.Retryable = stage == "connect" || stage == "validate"
if stage == "execute" || stage == "fetch" {
data.SessionDisposition = "quarantine"
}
}
return response{JSONRPC: "2.0", ID: id, Error: &rpcError{Code: -1, Message: err.Error(), Data: data}}
}
func errorStage(method string) string {
switch method {
case "open_session", "connect", "test_connection":
return "connect"
case "validate_session", "validate_connection":
return "validate"
case "fetch_query_page", "fetch_table_read_page":
return "fetch"
case "cancel_session":
return "cancel"
case "close_session", "disconnect", "close_query_session", "close_table_read_session", "shutdown":
return "close"
case "handshake", "":
return "request"
default:
return "execute"
}
}
func errorOutcome(stage string) string {
if stage == "request" || stage == "connect" || stage == "validate" {
return "not_started"
}
return "unknown"
}
func isNetworkError(err error) bool {
var networkError net.Error
return errors.As(err, &networkError)
}
+262
View File
@@ -0,0 +1,262 @@
package main
import (
"encoding/json"
"errors"
"fmt"
"strings"
"time"
nebula "github.com/vesoft-inc/nebula-go/v3"
)
const (
defaultMaxRows = 10000
maxQueryRows = 100000
defaultPage = 1000
)
type queryOptions struct {
SQL string `json:"sql"`
Database string `json:"database"`
MaxRows int `json:"maxRows"`
TimeoutSecs int `json:"timeoutSecs"`
}
type queryResult struct {
Columns []string `json:"columns"`
ColumnTypes []string `json:"column_types"`
Rows [][]any `json:"rows"`
AffectedRows int64 `json:"affected_rows"`
ExecutionTimeMS int64 `json:"execution_time_ms"`
Truncated bool `json:"truncated"`
}
type queryPageResult struct {
queryResult
SessionID *string `json:"session_id"`
HasMore bool `json:"has_more"`
}
type queryCursor struct {
result queryResult
offset int
}
func (s *agentSession) dispatch(method string, params map[string]json.RawMessage) (any, error) {
switch method {
case "connection_info":
return map[string]any{
"database": s.params.Database, "schema": "", "username": s.params.Username,
"identifierQuote": "`", "compatibilityMode": "ngql",
"databaseInfo": map[string]string{
"productName": "NebulaGraph", "driverName": "NebulaGraph Go Client", "driverVersion": "3.8.0",
"unquotedIdentifierCase": "mixed", "quotedIdentifierCase": "mixed",
},
}, nil
case "list_databases":
return s.listDatabases()
case "list_schemas":
return []string{}, nil
case "list_tables":
return s.listTables(params)
case "list_objects":
return s.listObjects(params)
case "list_data_types":
return []string{"bool", "int8", "int16", "int32", "int64", "float", "double", "string", "fixed_string", "date", "time", "datetime", "timestamp", "duration", "geography"}, nil
case "get_columns":
return s.getColumns(stringParam(params, "database"), stringParam(params, "table"))
case "get_table_ddl":
return s.getTableDDL(stringParam(params, "database"), stringParam(params, "table"), "")
case "get_object_source":
return s.getObjectSource(params)
case "get_table_comment":
return nil, nil
case "list_indexes", "list_foreign_keys", "list_triggers", "list_constraints", "list_partitions", "list_subpartitions":
return []any{}, nil
case "get_explain_info":
return map[string]any{"plan": "", "has_actual_stats": false}, nil
case "execute_query":
var options queryOptions
if err := decodeParams(params, &options); err != nil {
return nil, err
}
return s.executeQuery(options)
case "execute_query_page", "start_table_read":
var options queryOptions
if err := decodeParams(params, &options); err != nil {
return nil, err
}
return s.executeQueryPage(options, intParam(params, "pageSize"))
case "fetch_query_page", "fetch_table_read_page":
return s.fetchQueryPage(stringParam(params, "sessionId"), intParam(params, "pageSize"))
case "close_query_session", "close_table_read_session":
delete(s.cursors, stringParam(params, "sessionId"))
return map[string]bool{"ok": true}, nil
case "execute_batch":
started := time.Now()
for _, statement := range stringSliceParam(params, "statements") {
_, err := s.executeQuery(queryOptions{SQL: statement, Database: stringParam(params, "database")})
if err != nil {
return nil, err
}
}
return queryResult{Columns: []string{}, ColumnTypes: []string{}, Rows: [][]any{}, ExecutionTimeMS: time.Since(started).Milliseconds()}, nil
default:
return nil, fmt.Errorf("unsupported NebulaGraph agent method: %s", method)
}
}
func (s *agentSession) execute(statement, space string) (*nebula.ResultSet, error) {
if s.canceled.Load() {
return nil, errors.New("agent session was canceled")
}
if space = strings.TrimSpace(space); space != "" {
quoted, err := quoteNebulaIdentifier(space)
if err != nil {
return nil, err
}
selected, err := s.conn.Execute("USE " + quoted)
if err != nil {
return nil, err
}
if err := resultError(selected); err != nil {
return nil, err
}
}
result, err := s.conn.Execute(statement)
if err != nil {
return nil, err
}
if err := resultError(result); err != nil {
return nil, err
}
return result, nil
}
func (s *agentSession) executeQuery(options queryOptions) (queryResult, error) {
started := time.Now()
statement := strings.TrimSpace(options.SQL)
if statement == "" {
return queryResult{Columns: []string{}, ColumnTypes: []string{}, Rows: [][]any{}}, nil
}
space := options.Database
if space == "" {
space = s.params.Database
}
result, err := s.execute(statement, space)
if err != nil {
return queryResult{}, err
}
rows, types, err := readRows(result, effectiveMaxRows(options.MaxRows))
if err != nil {
return queryResult{}, err
}
return queryResult{
Columns: result.GetColNames(), ColumnTypes: types, Rows: rows,
ExecutionTimeMS: time.Since(started).Milliseconds(), Truncated: len(result.GetRows()) > len(rows),
}, nil
}
func (s *agentSession) executeQueryPage(options queryOptions, pageSize int) (queryPageResult, error) {
result, err := s.executeQuery(options)
if err != nil {
return queryPageResult{}, err
}
if pageSize <= 0 {
pageSize = defaultPage
}
if pageSize >= len(result.Rows) {
return queryPageResult{queryResult: result}, nil
}
s.nextCursorID++
id := fmt.Sprintf("nebula-query-%d", s.nextCursorID)
s.cursors[id] = &queryCursor{result: result, offset: pageSize}
result.Rows = result.Rows[:pageSize]
return queryPageResult{queryResult: result, SessionID: &id, HasMore: true}, nil
}
func (s *agentSession) fetchQueryPage(id string, pageSize int) (queryPageResult, error) {
cursor := s.cursors[id]
if cursor == nil {
return queryPageResult{}, fmt.Errorf("query session not found: %s", id)
}
if pageSize <= 0 {
pageSize = defaultPage
}
end := min(cursor.offset+pageSize, len(cursor.result.Rows))
result := cursor.result
result.Rows = result.Rows[cursor.offset:end]
cursor.offset = end
if end == len(cursor.result.Rows) {
delete(s.cursors, id)
return queryPageResult{queryResult: result}, nil
}
return queryPageResult{queryResult: result, SessionID: &id, HasMore: true}, nil
}
func effectiveMaxRows(requested int) int {
if requested <= 0 {
return defaultMaxRows
}
return min(requested, maxQueryRows)
}
func readRows(result *nebula.ResultSet, limit int) ([][]any, []string, error) {
columns := result.GetColNames()
rows := make([][]any, 0, min(len(result.GetRows()), limit))
types := make([]string, len(columns))
for i := range types {
types[i] = "Unknown"
}
for index := 0; index < len(result.GetRows()) && index < limit; index++ {
record, err := result.GetRowValuesByIndex(index)
if err != nil {
return nil, nil, err
}
row := make([]any, len(columns))
for column := range columns {
value, err := record.GetValueByIndex(column)
if err != nil {
return nil, nil, err
}
if types[column] == "Unknown" && value.GetType() != "null" {
types[column] = value.GetType()
}
row[column] = normalizeValue(value)
}
rows = append(rows, row)
}
return rows, types, nil
}
func normalizeValue(value *nebula.ValueWrapper) any {
if value == nil || value.GetType() == "null" {
return nil
}
if value.GetType() == "string" {
if text, err := value.AsString(); err == nil {
return text
}
}
return value.String()
}
type nebulaQueryError struct {
code string
message string
}
func (err *nebulaQueryError) Error() string {
return fmt.Sprintf("NebulaGraph error %s: %s", err.code, err.message)
}
func resultError(result *nebula.ResultSet) error {
if result == nil {
return errors.New("NebulaGraph returned no result")
}
if !result.IsSucceed() {
return &nebulaQueryError{code: fmt.Sprint(result.GetErrorCode()), message: result.GetErrorMsg()}
}
return nil
}
+1
View File
@@ -21,6 +21,7 @@ kingbase-go shared-fallback native-go Uses sys_catalog/information_schema for co
vastbase-go shared-fallback native-go Uses pg_catalog/information_schema for Vastbase metadata, then applies stable filtering and paging in the native agent.
kylin intentional-fallback java-jdbc-metadata Uses JDBC metadata from Kylin driver; no portable server-side paging API, common constraints filter locally.
neo4j-go shared-fallback native-go Uses Neo4j schema procedures with stable local filtering and paging because label discovery does not expose portable server-side metadata pagination.
nebula-go shared-fallback native-go Uses SHOW TAGS and SHOW EDGES with stable local filtering and paging because nGQL metadata listing has no server-side pagination.
oceanbase-oracle native-pushdown java-sql Uses Oracle-compatible metadata SQL with type, filter, stable order, and ROWNUM paging.
oracle-go native-pushdown native-go Uses Oracle-compatible metadata SQL with type, filter, stable order, and ROWNUM paging, with legacy profile guard.
snowflake native-pushdown java-sql Uses INFORMATION_SCHEMA with type, filter, stable order, and literal LIMIT/OFFSET.
1 driver strategy scope reason
21 vastbase-go shared-fallback native-go Uses pg_catalog/information_schema for Vastbase metadata, then applies stable filtering and paging in the native agent.
22 kylin intentional-fallback java-jdbc-metadata Uses JDBC metadata from Kylin driver; no portable server-side paging API, common constraints filter locally.
23 neo4j-go shared-fallback native-go Uses Neo4j schema procedures with stable local filtering and paging because label discovery does not expose portable server-side metadata pagination.
24 nebula-go shared-fallback native-go Uses SHOW TAGS and SHOW EDGES with stable local filtering and paging because nGQL metadata listing has no server-side pagination.
25 oceanbase-oracle native-pushdown java-sql Uses Oracle-compatible metadata SQL with type, filter, stable order, and ROWNUM paging.
26 oracle-go native-pushdown native-go Uses Oracle-compatible metadata SQL with type, filter, stable order, and ROWNUM paging, with legacy profile guard.
27 snowflake native-pushdown java-sql Uses INFORMATION_SCHEMA with type, filter, stable order, and literal LIMIT/OFFSET.
@@ -28,6 +28,8 @@ class DriverReleasePackagesTest(unittest.TestCase):
rocketmq_source.write_bytes(b"MZtest-rocketmq-agent")
cassandra_source = release_dir / "dbx-agent-cassandra-linux-x64"
cassandra_source.write_bytes(b"\x7fELFtest-cassandra-agent")
nebula_source = release_dir / "dbx-agent-nebula-linux-aarch64"
nebula_source.write_bytes(b"\x7fELFtest-nebula-agent")
tdengine_source = release_dir / "dbx-agent-tdengine-windows-aarch64.exe"
tdengine_source.write_bytes(b"MZtest-tdengine-agent")
etcd_source = release_dir / "dbx-agent-etcd-linux-x64"
@@ -43,6 +45,7 @@ class DriverReleasePackagesTest(unittest.TestCase):
"kingbase": "0.1.34",
"iotdb": "0.1.30",
"neo4j": "0.1.40",
"nebula": "0.1.0",
"vastbase": "0.1.37",
"duckdb": "0.1.0",
"rabbitmq": "0.1.0",
@@ -66,6 +69,7 @@ class DriverReleasePackagesTest(unittest.TestCase):
versioned_rabbitmq = release_dir / "dbx-agent-rabbitmq-0.1.0-linux-x64"
versioned_rocketmq = release_dir / "dbx-agent-rocketmq-0.1.0-windows-x64.exe"
versioned_cassandra = release_dir / "dbx-agent-cassandra-0.1.37-linux-x64"
versioned_nebula = release_dir / "dbx-agent-nebula-0.1.0-linux-aarch64"
versioned_tdengine = release_dir / "dbx-agent-tdengine-0.1.0-windows-aarch64.exe"
versioned_etcd = release_dir / "dbx-agent-etcd-0.1.40-linux-x64"
versioned_etcd2 = release_dir / "dbx-agent-etcd2-0.1.0-macos-aarch64"
@@ -75,6 +79,7 @@ class DriverReleasePackagesTest(unittest.TestCase):
versioned_java,
versioned_cassandra,
versioned_native,
versioned_nebula,
versioned_vastbase,
versioned_duckdb,
versioned_rabbitmq,
@@ -262,6 +267,7 @@ class DriverReleasePackagesTest(unittest.TestCase):
versioned_etcd2,
versioned_java,
versioned_native,
versioned_nebula,
versioned_rabbitmq,
versioned_rocketmq,
versioned_tdengine,
@@ -285,6 +291,20 @@ class DriverReleasePackagesTest(unittest.TestCase):
self.assertFalse(source.exists())
self.assertEqual(versioned.read_bytes(), b"\xcf\xfa\xed\xfetest-neo4j-agent")
def test_versions_nebula_native_artifacts(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
release_dir = Path(temp_dir)
source = release_dir / "dbx-agent-nebula-linux-aarch64"
source.write_bytes(b"\x7fELFtest-nebula-agent")
versions = {driver: "0.1.0" for driver in NATIVE_DRIVERS}
renamed = version_agent_artifacts(release_dir, versions)
versioned = release_dir / "dbx-agent-nebula-0.1.0-linux-aarch64"
self.assertEqual(renamed, [versioned])
self.assertFalse(source.exists())
self.assertEqual(versioned.read_bytes(), b"\x7fELFtest-nebula-agent")
def test_versions_iotdb_native_artifacts(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
release_dir = Path(temp_dir)
+1
View File
@@ -21,6 +21,7 @@ NATIVE_ONLY_AGENT_MODULES = {
"kingbase": "drivers/kingbase-go",
"iotdb": "drivers/iotdb",
"neo4j": "drivers/neo4j-go",
"nebula": "drivers/nebula-go",
"vastbase": "drivers/vastbase-go",
"tdengine": "drivers/tdengine",
"xugu": "drivers/xugu",
+1 -1
View File
@@ -4,7 +4,7 @@ import json
from pathlib import Path
NATIVE_DRIVERS = ("cassandra", "hive", "argo", "oracle", "xugu", "kingbase", "iotdb", "neo4j", "vastbase", "duckdb", "rabbitmq", "rocketmq", "zookeeper", "tdengine", "etcd", "etcd2", "sqlite-worker")
NATIVE_DRIVERS = ("cassandra", "hive", "argo", "oracle", "xugu", "kingbase", "iotdb", "neo4j", "nebula", "vastbase", "duckdb", "rabbitmq", "rocketmq", "zookeeper", "tdengine", "etcd", "etcd2", "sqlite-worker")
PLATFORMS = (
"macos-aarch64",
"macos-x64",
+1
View File
@@ -24,6 +24,7 @@
"kingbase": "0.1.59",
"kylin": "0.1.69",
"neo4j": "0.1.49",
"nebula": "0.1.0",
"oceanbase-oracle": "0.1.66",
"oracle": "0.1.65",
"snowflake": "0.1.69",
Binary file not shown.

After

Width:  |  Height:  |  Size: 9.1 KiB

@@ -151,7 +151,7 @@ import { resolveVisibleDatabaseSaveAction } from "@/components/sidebar/visibleDa
import { canSaveVisibleDatabaseSelection, connectionUsesVisibleSchemaFilter, filterDatabaseNamesForVisiblePicker, filterSchemaNamesForVisiblePicker, buildDraftVisibleSchemasConnectionId, normalizeVisibleSchemaSelection } from "@/lib/database/visibleDatabases";
import { isSchemaAware, isSingleDatabase, supportsDataDictionary } from "@/lib/database/databaseFeatureSupport";
import { normalizeConnectionScope, normalizeConnectionTimeouts } from "@/lib/connection/connectionSubmitNormalization";
import { databaseConnectionFormKind } from "@/lib/database/databaseDriverManifest";
import { databaseConnectionFormKind, databaseManifestEntry } from "@/lib/database/databaseDriverManifest";
import VisibleSchemasDialog from "@/components/sidebar/VisibleSchemasDialog.vue";
import CloudflareD1ConnectionFields from "@/components/connection/CloudflareD1ConnectionFields.vue";
import SpannerConnectionFields from "@/components/connection/SpannerConnectionFields.vue";
@@ -1184,6 +1184,8 @@ const driverProfiles: Record<string, ConnectionProfileDefinition> = {
...CONNECTION_PROFILES,
...jdbcProductDriverProfiles(),
};
const nebulaDriverProfiles = databaseManifestEntry("nebula")?.driverProfiles ?? [];
const nebulaDefaultDriverProfile = nebulaDriverProfiles[0]?.profile ?? "nebula";
function profileForConfig(config: ConnectionConfig) {
if (config.db_type === "plugin" && config.plugin_id && config.plugin_connection_provider) {
@@ -2804,7 +2806,7 @@ function applyProfile(val: string, preserveConnectionFields = false) {
const previousDatabaseType = form.value.db_type;
selectedType.value = val;
form.value.db_type = profile.type;
form.value.driver_profile = val;
form.value.driver_profile = val === "nebula" ? nebulaDefaultDriverProfile : val;
form.value.driver_label = isCustomCompatibleProfile() ? customDriverName.value.trim() || profile.label : profile.label;
const preserveMeilisearchConfig = preserveConnectionFields && previousDatabaseType === "meilisearch" && profile.type === "meilisearch";
if (profile.type !== "sqlserver" && !preserveMeilisearchConfig) {
@@ -3434,6 +3436,12 @@ function switchEtcdApiVersion(profile: "etcd" | "etcd-v2") {
resetTestState();
}
function switchNebulaDriverProfile(profile: unknown) {
if (typeof profile !== "string" || !nebulaDriverProfiles.some((entry) => entry.profile === profile)) return;
form.value.driver_profile = profile;
resetTestState();
}
function switchH2DriverProfile(profile: "h2" | "h2-v1" | "h2-v2" | "h2-v3" | "h2-custom") {
form.value.driver_profile = profile;
if (profile === "h2-custom") {
@@ -3650,11 +3658,12 @@ const tlsCapableDatabaseTypes = new Set<DatabaseType>([
"influxdb",
"victoriametrics",
"cassandra",
"nebula",
"zookeeper",
]);
const supportsTlsToggle = computed(() => tlsCapableDatabaseTypes.has(form.value.db_type) || supportsMysqlTlsTab(form.value.db_type, selectedType.value));
const supportsCaCertificatePath = computed(() => form.value.db_type === "clickhouse" || form.value.db_type === "victoriametrics");
const supportsGenericUrlParams = computed(() => form.value.db_type !== "manticoresearch" && form.value.db_type !== "hbase");
const supportsGenericUrlParams = computed(() => form.value.db_type !== "manticoresearch" && form.value.db_type !== "hbase" && form.value.db_type !== "nebula");
const showGenericUrlParamsHint = computed(() => form.value.db_type === "mysql" || form.value.db_type === "doris" || form.value.db_type === "starrocks");
const bareMysqlProfiles = new Set(["doris", "selectdb", "oceanbase"]);
const supportsMysqlTlsOptions = computed(() => mysqlTlsOptionsSupported(form.value.db_type, selectedType.value));
@@ -4628,6 +4637,9 @@ function connectionConfigForSubmit(id: string, generatedName = "", validatePlugi
} else {
config = { ...formValueForSubmit(), id } as LegacyConnectionConfig;
}
if (config.db_type === "nebula" && (!config.driver_profile || config.driver_profile === "nebula")) {
config.driver_profile = nebulaDefaultDriverProfile;
}
config.database_info = undefined;
config.database = normalizeStoredConnectionDatabase(config.db_type, config.database);
config.note = config.note?.trim() || undefined;
@@ -7031,6 +7043,20 @@ function openExternalUrl(url: string) {
</button>
</div>
<div v-if="form.db_type === 'nebula'" class="grid grid-cols-4 items-center gap-4">
<Label :class="connectionLabelClass">{{ t("connection.version") }}</Label>
<div class="col-span-3">
<Select :model-value="form.driver_profile === 'nebula' ? nebulaDefaultDriverProfile : form.driver_profile" @update:model-value="switchNebulaDriverProfile">
<SelectTrigger class="w-full">
<SelectValue />
</SelectTrigger>
<SelectContent>
<SelectItem v-for="profile in nebulaDriverProfiles" :key="profile.profile" :value="profile.profile">{{ profile.label }}</SelectItem>
</SelectContent>
</Select>
</div>
</div>
<!-- OceanBase mode toggle -->
<div v-if="selectedType === 'oceanbase'" class="grid grid-cols-4 items-center gap-4">
<Label :class="connectionLabelSmallClass">{{ t("connection.mode") }}</Label>
@@ -9322,7 +9348,7 @@ function openExternalUrl(url: string) {
</label>
</div>
<template v-if="form.db_type === 'etcd' || form.db_type === 'consul' || form.db_type === 'zookeeper' || form.db_type === 'elasticsearch' || form.db_type === 'easysearch'">
<template v-if="form.db_type === 'etcd' || form.db_type === 'consul' || form.db_type === 'zookeeper' || form.db_type === 'elasticsearch' || form.db_type === 'easysearch' || form.db_type === 'nebula'">
<div class="grid grid-cols-4 items-start gap-4">
<Label :class="connectionLabelSmallPaddedClass">
<span class="inline-flex items-center justify-end gap-1">
@@ -9375,7 +9401,7 @@ function openExternalUrl(url: string) {
<TooltipContent>{{ t("connection.etcdClientKeyBrowse") }}</TooltipContent>
</Tooltip>
</div>
<p class="text-[11px] leading-4 text-muted-foreground">
<p v-if="form.db_type !== 'nebula'" class="text-[11px] leading-4 text-muted-foreground">
{{ t("connection.etcdClientCertHint") }}
</p>
</div>
@@ -801,7 +801,7 @@ function sqlVariableSyntaxToggle(key: keyof SqlVariableSyntaxToggles): boolean {
function setSqlVariableSyntaxToggle(key: keyof SqlVariableSyntaxToggles, value: boolean) {
const dbType = editSqlVariableSyntaxDatabaseType.value;
if (dbType === "neo4j" && key === "named") return;
if ((dbType === "neo4j" || dbType === "nebula") && key === "named") return;
const merged: SqlVariableSyntaxToggles = {
...DEFAULT_SQL_VARIABLE_SYNTAX_TOGGLES,
...editSqlVariableSyntaxOverrides.value[dbType],
@@ -7002,7 +7002,7 @@ onUnmounted(() => {
<Switch
:id="`sql-var-syntax-${key}`"
:model-value="sqlVariableSyntaxToggle(key)"
:disabled="!editSqlVariableSubstitutionEnabled || (editSqlVariableSyntaxDatabaseType === 'neo4j' && key === 'named')"
:disabled="!editSqlVariableSubstitutionEnabled || (['neo4j', 'nebula'].includes(editSqlVariableSyntaxDatabaseType) && key === 'named')"
class="mt-0.5 shrink-0"
@update:model-value="(value) => setSqlVariableSyntaxToggle(key, value as boolean)"
/>
@@ -750,7 +750,7 @@ const resolvedDatabaseType = computed(() => props.databaseType ?? effectiveDatab
// editor and clipboard path must decode it whenever MongoDB values are on screen.
const usesMongoDocumentGridValues = computed(() => props.mongoCollectionGrid === true || resolvedDatabaseType.value === "mongodb");
const isResultsContext = computed(() => props.context === "results");
const canShowWhereSearch = computed(() => !!props.onExecuteSql && !isResultsContext.value && resolvedDatabaseType.value !== "victoriametrics");
const canShowWhereSearch = computed(() => !!props.onExecuteSql && !isResultsContext.value && resolvedDatabaseType.value !== "victoriametrics" && resolvedDatabaseType.value !== "nebula");
const canUseWhereSearch = computed(() => !!props.tableMeta && canShowWhereSearch.value);
const canUseServerColumnFilter = computed(() => canUseWhereSearch.value && !!props.connectionId && !!props.tableMeta);
const tableStructureCapabilities = computed(() => getTableStructureCapabilities(resolvedDatabaseType.value, resolvedConnectionConfig.value?.db_type));
@@ -87,6 +87,7 @@ const assetIcons: Record<string, string> = {
starrocks: "starrocks",
redshift: "redshift",
neo4j: "neo4j",
nebula: "nebula.png",
informix: "informix",
databricks: "databricks",
saphana: "saphana",
@@ -310,14 +310,12 @@ const effectiveDatabaseType = computed(() => effectiveDatabaseTypeForConnection(
const isGaussdbM = computed(() => effectiveDatabaseType.value === "gaussdb" && props.connection.driver_profile?.toLowerCase() === "gaussdb-m");
const isVictoriaMetrics = computed(() => effectiveDatabaseType.value === "victoriametrics");
const isMongodb = computed(() => props.connection.db_type === "mongodb");
// Victoria Metrics reports series instead of rows and has no byte size to show;
// every other engine (MongoDB collections included, via `collStats`) fills both
// the row and size columns.
const supportsObjectSizeStats = computed(() => !isVictoriaMetrics.value);
// Neither VictoriaMetrics series nor NebulaGraph tag/edge metadata has table byte-size stats.
const supportsObjectSizeStats = computed(() => !isVictoriaMetrics.value && effectiveDatabaseType.value !== "nebula");
// The batch table toolbar (export/copy/truncate/empty/drop selected) is SQL-only:
// MongoDB collections are not dropped or truncated through it.
const supportsBatchTableActions = computed(() => !isVictoriaMetrics.value && !isMongodb.value);
const showTableStatistics = computed(() => objectFilter.value === "all" || objectFilter.value === "tables");
const supportsBatchTableActions = computed(() => !isVictoriaMetrics.value && !isMongodb.value && effectiveDatabaseType.value !== "nebula");
const showTableStatistics = computed(() => effectiveDatabaseType.value !== "nebula" && (objectFilter.value === "all" || objectFilter.value === "tables"));
const showObjectRowStats = computed(() => showTableStatistics.value);
const showObjectSizeStats = computed(() => supportsObjectSizeStats.value && showTableStatistics.value);
const objectRowsLabel = computed(() => t(isVictoriaMetrics.value ? "objects.series" : "objects.rows"));
@@ -617,7 +615,7 @@ watch(
},
);
const showCheckboxColumn = computed(() => settingsStore.editorSettings.objectBrowserShowCheckbox || selectedTableCount.value > 0);
const showCheckboxColumn = computed(() => supportsBatchTableActions.value && (settingsStore.editorSettings.objectBrowserShowCheckbox || selectedTableCount.value > 0));
function toggleCheckboxColumn() {
const next = !settingsStore.editorSettings.objectBrowserShowCheckbox;
@@ -2576,7 +2574,7 @@ async function copySelectedTablesToClipboard() {
}
function canPasteTableClipboard(): boolean {
return !isVictoriaMetrics.value && !isMongodb.value && tableClipboardMatchesTarget(normalizedObjectBrowserTableClipboardEntries(), pasteTableTargetContext());
return supportsBatchTableActions.value && tableClipboardMatchesTarget(normalizedObjectBrowserTableClipboardEntries(), pasteTableTargetContext());
}
function normalizedObjectBrowserTableClipboardEntries() {
@@ -2668,7 +2666,7 @@ function openPasteTableDialog() {
}
function onObjectBrowserKeydown(event: KeyboardEvent) {
if (event.defaultPrevented) return;
if (event.defaultPrevented || !supportsBatchTableActions.value) return;
if (eventTargetAllowsAppClipboardShortcut(event, "c")) {
if (selectedTableCount.value === 0) return;
event.preventDefault();
@@ -3462,6 +3460,14 @@ function selectedBatchTableCountLabel(key: "batchDrop" | "batchTruncate" | "batc
}
function getTableMenuItems(item: ObjectBrowserRow): ContextMenuItem[] {
if (effectiveDatabaseType.value === "nebula") {
return [
{ label: t("contextMenu.viewData"), action: () => openViewData(item), icon: Table2 },
{ label: t("contextMenu.viewDdl"), action: () => openTableInfo(item, "ddl"), icon: FileCode },
{ label: "", separator: true },
{ label: t("contextMenu.copyName"), action: () => copyName(item), icon: Copy },
];
}
if (isVictoriaMetrics.value) {
return [
{ label: t("contextMenu.viewData"), action: () => openViewData(item), icon: Table2 },
@@ -3530,6 +3536,14 @@ function getTableMenuItems(item: ObjectBrowserRow): ContextMenuItem[] {
}
function getViewMenuItems(item: ObjectBrowserRow): ContextMenuItem[] {
if (effectiveDatabaseType.value === "nebula") {
return [
{ label: t("contextMenu.viewData"), action: () => openViewData(item), icon: Table2 },
{ label: t("contextMenu.viewDdl"), action: () => openTableInfo(item, "ddl"), icon: ScrollText },
{ label: "", separator: true },
{ label: t("contextMenu.copyName"), action: () => copyName(item), icon: Copy },
];
}
return [
{ label: t("contextMenu.viewData"), action: () => openViewData(item), icon: Table2 },
{ label: t("contextMenu.editView"), action: () => openSource(item), icon: PencilLine },
@@ -3709,7 +3723,7 @@ function getObjectBrowserMenuItems(item: ObjectBrowserRow): ContextMenuItem[] {
<LayoutGrid class="h-3.5 w-3.5" />
</button>
</div>
<Button v-if="showInlineCheckboxToggle" variant="ghost" size="icon" class="h-7 w-7" :class="{ 'text-primary': settingsStore.editorSettings.objectBrowserShowCheckbox }" :title="t('objects.toggleCheckbox')" @click="toggleCheckboxColumn">
<Button v-if="showInlineCheckboxToggle && supportsBatchTableActions" variant="ghost" size="icon" class="h-7 w-7" :class="{ 'text-primary': settingsStore.editorSettings.objectBrowserShowCheckbox }" :title="t('objects.toggleCheckbox')" @click="toggleCheckboxColumn">
<CheckSquare v-if="settingsStore.editorSettings.objectBrowserShowCheckbox" class="h-3.5 w-3.5" />
<Square v-else class="h-3.5 w-3.5" />
</Button>
@@ -3747,14 +3761,14 @@ function getObjectBrowserMenuItems(item: ObjectBrowserRow): ContextMenuItem[] {
<LayoutGrid class="h-3.5 w-3.5" />
{{ t("objects.viewGrid") }}
</DropdownMenuItem>
<DropdownMenuCheckboxItem :model-value="settingsStore.editorSettings.objectBrowserShowCheckbox" @select.prevent @update:model-value="toggleCheckboxColumn()">{{ t("objects.toggleCheckbox") }}</DropdownMenuCheckboxItem>
<DropdownMenuCheckboxItem v-if="supportsBatchTableActions" :model-value="settingsStore.editorSettings.objectBrowserShowCheckbox" @select.prevent @update:model-value="toggleCheckboxColumn()">{{ t("objects.toggleCheckbox") }}</DropdownMenuCheckboxItem>
<template v-if="showObjectFilter && toolbarTier >= 2">
<DropdownMenuSeparator />
<DropdownMenuCheckboxItem v-for="filter in objectFilters" :key="filter" :model-value="objectFilter === filter" @select.prevent @update:model-value="selectObjectFilter(filter)">{{ filterLabel(filter) }}</DropdownMenuCheckboxItem>
</template>
</ToolbarOverflowMenu>
</div>
<div v-if="selectedTableCount > 0" class="flex h-9 shrink-0 items-center gap-2 overflow-x-auto border-b bg-muted/30 px-3 text-xs">
<div v-if="selectedTableCount > 0 && supportsBatchTableActions" class="flex h-9 shrink-0 items-center gap-2 overflow-x-auto border-b bg-muted/30 px-3 text-xs">
<div class="min-w-0 flex-1 truncate text-muted-foreground">
{{ t("objects.selectedTables", { count: selectedTableCount }) }}
</div>
@@ -93,6 +93,7 @@ const mountedApps: Array<{ app: App; host: HTMLElement }> = [];
beforeEach(() => {
vi.clearAllMocks();
connection.db_type = "mysql";
mocks.store.treeClipboard = null;
mocks.listObjects.mockResolvedValue([
{ name: "orders", object_type: "TABLE" },
@@ -154,4 +155,18 @@ describe("ObjectBrowser table copy", () => {
await vi.waitFor(() => expect(mocks.toast).toHaveBeenCalledWith("grid.copyFailed", 5000));
expect(mocks.toast).not.toHaveBeenCalledWith("contextMenu.pasteTableClipboardUpdated", 2000);
});
it("does not expose relational table copy or batch actions for NebulaGraph", async () => {
connection.db_type = "nebula";
const host = await mountBrowser();
rowFor(host, "orders").dispatchEvent(new MouseEvent("click", { bubbles: true, ctrlKey: true }));
await nextTick();
expect(host.textContent).not.toContain("objects.exportSelected");
expect(host.textContent).not.toContain("objects.dropSelected");
expect(host.querySelector('[title="objects.toggleCheckbox"]')).toBeNull();
host.querySelector<HTMLElement>("[data-object-browser-root]")!.dispatchEvent(new KeyboardEvent("keydown", { key: "c", ctrlKey: true, bubbles: true }));
expect(mocks.copyToClipboard).not.toHaveBeenCalled();
expect(mocks.store.treeClipboard).toBeNull();
});
});
@@ -2787,6 +2787,7 @@ function requestDropTableChildObject() {
}
function canDropTreeNode(node: TreeNode): boolean {
if (databaseTypeForNode(node) === "nebula") return false;
if (isSqlServerLinkedNode(node)) return false;
if (node.type === "table") return !!node.connectionId && !!node.database;
if (node.type === "view" || node.type === "materialized_view" || node.type === "procedure" || node.type === "function" || node.type === "event") {
@@ -3129,6 +3130,7 @@ function requestDropSelectedNodes(): boolean {
}
function requestDropSelectedNode(): boolean {
if (currentDatabaseType() === "nebula") return false;
if (activeNode.value.type === "table") {
dropTable();
return true;
@@ -4753,7 +4755,7 @@ const canOpenSqlFileExecution = computed(() => {
const canExportAllDatabases = computed(() => {
if (activeNode.value.type !== "connection" || !activeNode.value.connectionId) return false;
const dbType = connectionStore.getConfig(activeNode.value.connectionId)?.db_type;
return !["redis", "mongodb", "dynamodb", "elasticsearch", "easysearch", "meilisearch", "solr", "qdrant", "milvus", "weaviate", "chromadb", "etcd", "zookeeper", "consul", "mq", "nacos", "plugin", "salesforce"].includes(dbType || "");
return !["redis", "mongodb", "dynamodb", "elasticsearch", "easysearch", "meilisearch", "solr", "qdrant", "milvus", "weaviate", "chromadb", "etcd", "zookeeper", "consul", "mq", "nacos", "plugin", "salesforce", "nebula"].includes(dbType || "");
});
const canOpenScheduledBackups = computed(() => {
@@ -5860,6 +5862,19 @@ function buildDatabaseSidebarMenu(context: SidebarMenuFactoryContext): boolean {
});
return true;
}
if (currentDatabaseType() === "nebula" && node.type === "database") {
if (canCloseDatabaseConnection.value) items.push({ label: t("contextMenu.closeDatabaseConnection"), action: closeDatabaseConnection, icon: Unplug });
items.push(copyNameMenuItem());
items.push({ label: "", separator: true });
if (canOpenObjectBrowser.value) items.push({ label: t("contextMenu.openObjectBrowser"), action: openObjectBrowser, icon: TableProperties });
items.push({ label: t("contextMenu.newQuery"), action: newQuery, icon: TerminalSquare });
const sqlHistoryMenu = savedSqlHistorySubmenu();
if (sqlHistoryMenu) items.push(sqlHistoryMenu);
items.push({ label: isNodeDefaultDatabase.value ? t("contextMenu.clearDefaultDatabase") : t("contextMenu.setDefaultDatabase"), action: isNodeDefaultDatabase.value ? clearNodeDefaultDatabase : setNodeAsDefaultDatabase, icon: Database });
items.push({ label: "", separator: true });
items.push({ label: t("contextMenu.refreshChildren"), action: refresh, icon: RefreshCw, shortcut: shortcutRefresh });
return true;
}
if (canCloseDatabaseConnection.value) {
items.unshift({ label: "", separator: true });
items.unshift({ label: t("contextMenu.closeDatabaseConnection"), action: closeDatabaseConnection, icon: Unplug });
@@ -6266,6 +6281,20 @@ function buildObjectSidebarMenu(context: SidebarMenuFactoryContext): boolean {
appendPluginTableMenuItems(items, node);
return true;
}
if (currentDatabaseType() === "nebula") {
items.push(copyNameMenuItem());
items.push({ label: t("contextMenu.newQuery"), action: newQuery, icon: TerminalSquare });
items.push({ label: "", separator: true });
items.push({ label: t("contextMenu.viewData"), action: openDataImmediately, icon: TableProperties });
items.push({ label: t("contextMenu.openInNewDataTab"), action: openDataInNewTabImmediately, icon: CopyPlus, shortcut: shortcutOpenDataInNewTab.value });
items.push({ label: t("contextMenu.viewDdl"), action: openDdl, icon: FileCode });
const sqlHistoryMenu = savedSqlHistorySubmenu();
if (sqlHistoryMenu) items.push(sqlHistoryMenu);
items.push({ label: "", separator: true });
items.push({ label: t("contextMenu.refreshChildren"), action: refresh, icon: RefreshCw, shortcut: shortcutRefresh });
appendPluginTableMenuItems(items, node);
return true;
}
const destructiveActions: ContextMenuItem[] = [];
items.push(copyNameMenuItem());
items.push({ label: t("contextMenu.newQuery"), action: newQuery, icon: TerminalSquare });
@@ -6618,7 +6647,7 @@ function buildObjectGroupSidebarMenu(context: SidebarMenuFactoryContext): boolea
const mysqlObjectTemplate = node.connectionId ? mysqlObjectTemplateForGroup(connectionStore.getConfig(node.connectionId), node) : null;
const hasMongoCreateIndexAction = node.type === "group-indexes" && canCreateMongoIndex.value;
const hasMongoDropAllIndexesAction = node.type === "group-indexes" && canDropAllMongoIndexes.value;
const hasGroupAction = (node.type === "group-tables" && canCreateTable.value) || (node.type === "group-views" && !!node.connectionId && !!node.database) || !!mysqlObjectTemplate || hasMongoCreateIndexAction || hasMongoDropAllIndexesAction;
const hasGroupAction = (node.type === "group-tables" && canCreateTable.value) || (node.type === "group-views" && !!node.connectionId && !!node.database && currentDatabaseType() !== "nebula") || !!mysqlObjectTemplate || hasMongoCreateIndexAction || hasMongoDropAllIndexesAction;
const canLoadAllObjectGroup = node.type === "group-tables" || node.type === "group-dolt-system-tables" || node.type === "group-views" || node.type === "group-materialized-views";
if (node.type === "group-tables" && canCreateTable.value) {
items.push({ label: t("contextMenu.createTable"), action: createTable, icon: Plus });
@@ -6629,7 +6658,7 @@ function buildObjectGroupSidebarMenu(context: SidebarMenuFactoryContext): boolea
items.push({ label: t("contextMenu.pasteTable"), action: openPasteTableDialog, icon: Clipboard });
}
}
if (node.type === "group-views" && node.connectionId && node.database) {
if (node.type === "group-views" && node.connectionId && node.database && currentDatabaseType() !== "nebula") {
items.push({ label: t("contextMenu.createView"), action: createView, icon: Plus });
}
if (node.type === "group-events" && node.connectionId && node.database) {
@@ -0,0 +1,80 @@
// @vitest-environment happy-dom
import { createApp, defineComponent, h, nextTick, ref, type App } from "vue";
import { createPinia, setActivePinia } from "pinia";
import { afterEach, describe, expect, it, vi } from "vitest";
import i18n from "@/i18n";
import type { ContextMenuItem } from "@/components/ui/CustomContextMenu.vue";
import type { TreeNode } from "@/types/database";
vi.mock("@/lib/backend/api", async (importOriginal) => {
const actual = await importOriginal<typeof import("@/lib/backend/api")>();
return { ...actual, listPlugins: vi.fn().mockResolvedValue([]) };
});
import SidebarTreeRuntimeHost from "@/components/sidebar/SidebarTreeRuntimeHost.vue";
import { useConnectionStore } from "@/stores/connectionStore";
const connection = { id: "nebula-1", name: "NebulaGraph", db_type: "nebula", driver_profile: "nebula", host: "localhost", port: 9669, username: "root", password: "" };
const node = (type: TreeNode["type"], label: string): TreeNode => ({ id: `nebula-1:space:${type}:${label}`, type, label, connectionId: connection.id, database: "space", tableName: label });
const mountedApps: App<Element>[] = [];
const labels = (items: ContextMenuItem[]): string[] => items.flatMap((item) => [...(item.label ? [item.label] : []), ...(item.children ? labels(item.children) : [])]);
const tr = (key: string) => i18n.global.t(key);
async function mountHost() {
const pinia = createPinia();
setActivePinia(pinia);
useConnectionStore().connections = [connection];
const host = ref<InstanceType<typeof SidebarTreeRuntimeHost> | null>(null);
const root = defineComponent({ setup: () => () => h(SidebarTreeRuntimeHost, { ref: host, node: node("connection", connection.name), depth: 0 }) });
const container = document.createElement("div");
document.body.appendChild(container);
const app = createApp(root);
app.use(pinia);
app.use(i18n);
app.mount(container);
mountedApps.push(app);
await nextTick();
await new Promise((resolve) => setTimeout(resolve, 0));
await nextTick();
return host.value as { buildContextMenu(node: TreeNode): ContextMenuItem[] };
}
describe("NebulaGraph sidebar context menus", () => {
afterEach(() => {
for (const app of mountedApps.splice(0)) app.unmount();
document.body.innerHTML = "";
});
it("keeps space browsing and queries without relational database operations", async () => {
const host = await mountHost();
const items = labels(host.buildContextMenu(node("database", "space")));
expect(items).toEqual(expect.arrayContaining([tr("contextMenu.copyName"), tr("contextMenu.newQuery"), tr("contextMenu.openObjectBrowser"), tr("contextMenu.refreshChildren")]));
for (const key of ["transfer.dataTransfer", "diff.title", "dataCompare.title", "contextMenu.exportDatabase", "dataDictionary.title", "sqlFile.title"]) {
expect(items).not.toContain(tr(key));
}
});
it.each(["table", "view"] as const)("keeps %s reads and DDL without SQL table mutations", async (type) => {
const host = await mountHost();
const items = labels(host.buildContextMenu(node(type, type === "table" ? "person" : "knows")));
expect(items).toEqual(expect.arrayContaining([tr("contextMenu.copyName"), tr("contextMenu.newQuery"), tr("contextMenu.viewData"), tr("contextMenu.viewDdl"), tr("contextMenu.refreshChildren")]));
for (const key of ["contextMenu.generateSql", "contextMenu.exportDatabase", "contextMenu.exportData", "contextMenu.editView", "contextMenu.dropView", "contextMenu.dropTable", "contextMenu.emptyTable", "contextMenu.duplicateStructure", "dataCompare.title"]) {
expect(items).not.toContain(tr(key));
}
});
it("does not offer SQL view creation on the edge group", async () => {
const host = await mountHost();
const items = labels(host.buildContextMenu(node("group-views", "Edges")));
expect(items).toContain(tr("contextMenu.refreshChildren"));
expect(items).not.toContain(tr("contextMenu.createView"));
});
it("does not offer SQL files or all-database export on the connection", async () => {
const host = await mountHost();
const items = labels(host.buildContextMenu(node("connection", connection.name)));
expect(items).toContain(tr("contextMenu.newQuery"));
expect(items).not.toContain(tr("sqlFile.title"));
expect(items).not.toContain(tr("contextMenu.exportAllDatabases"));
});
});
@@ -425,6 +425,40 @@ describe("useSidebarDataOpenRuntime", () => {
});
});
it("waits for NebulaGraph properties before building the first table query", async () => {
mocks.databaseType = "nebula";
let releaseMetadata: () => void = () => {};
const gate = new Promise<void>((resolve) => {
releaseMetadata = resolve;
});
mocks.loadTableMetadata.mockImplementation(async () => {
await gate;
return {
metadata: {
schema: "public",
tableName: "users",
tableType: "TABLE",
database: "app",
columns: [{ name: "name", data_type: "string", is_nullable: true, column_default: null, is_primary_key: false, extra: null }],
indexes: [],
primaryKeys: [],
cachedAt: Date.now(),
},
cacheStatus: "miss",
ageMs: 0,
};
});
const opening = useSidebarDataOpenRuntime().openData(tableNode);
await vi.waitFor(() => expect(mocks.loadTableMetadata).toHaveBeenCalled());
expect(mocks.buildTableSelectSql).not.toHaveBeenCalled();
releaseMetadata();
await opening;
expect(mocks.buildTableSelectSql).toHaveBeenCalledWith(expect.objectContaining({ databaseType: "nebula", columns: ["name"] }));
expect(mocks.callOrder).toEqual(["query"]);
});
it("keeps Dameng metadata deferred until after the table query", async () => {
mocks.databaseType = "dameng";
@@ -0,0 +1,21 @@
import { describe, expect, it } from "vitest";
import { connectionUrlPlaceholder } from "@/lib/connection/connectionPresentation";
import { parseConnectionUrl } from "@/lib/connection/connectionUrl";
describe("NebulaGraph connection URLs", () => {
it("parses graphd connection details and an optional space", () => {
expect(parseConnectionUrl("nebula://root:secret@graphd.example.com:9669/demo")).toMatchObject({
dbType: "nebula",
driverProfile: "nebula",
host: "graphd.example.com",
port: 9669,
username: "root",
password: "secret",
database: "demo",
});
});
it("shows a NebulaGraph-specific placeholder", () => {
expect(connectionUrlPlaceholder("nebula")).toBe("nebula://root:password@graphd:9669/space");
});
});
@@ -37,6 +37,10 @@ describe("databaseDriverManifest", () => {
});
});
it("offers only the supported NebulaGraph 3.x profile", () => {
expect(databaseManifestEntry("nebula")?.driverProfiles).toEqual([{ profile: "nebula-v3", agentKey: "nebula", label: "NebulaGraph 3.x", storeVisible: false }]);
});
it("routes only specialized connection forms explicitly", () => {
expect(databaseConnectionFormKind("mysql")).toBe("standard");
expect(databaseConnectionFormKind("jdbc")).toBe("jdbc");
@@ -46,6 +46,15 @@ describe("extractSqlParameters", () => {
expect(extractSqlParameters("MATCH (p:Person {name:${name}}) RETURN p", options)).toEqual(["name"]);
});
it("keeps NebulaGraph tags and edge types intact in nGQL", () => {
const ngql = 'MATCH (p:Person)-[:WORK_IN]->(c:Company{name:"星云科技"}) RETURN p.Person.name LIMIT 10';
const options = { databaseType: "nebula" as const };
expect(extractSqlParameterDescriptors(ngql, options)).toEqual([]);
expect(substituteSqlParameters(ngql, {}, options)).toBe(ngql);
expect(extractSqlParameters("MATCH (p:Person {name:${name}}) RETURN p LIMIT 10", options)).toEqual(["name"]);
expect(extractSqlParameters("select :customer_id", { databaseType: "mysql" })).toEqual(["customer_id"]);
});
it("extracts unique template parameters in order", () => {
const sql = "select * from t where pt_dt between ${start_date} and ${end_date} or pt_dt = ${start_date}";
expect(extractSqlParameters(sql)).toEqual(["start_date", "end_date"]);
@@ -19,6 +19,13 @@ describe("resolveSqlVariableSyntaxToggles", () => {
expect(resolveSqlVariableSyntaxToggles({ neo4j: { shell: false } }, "neo4j")).toEqual({ ...toggles, shell: false });
});
it("disables SQL named placeholders for NebulaGraph", () => {
const toggles = resolveSqlVariableSyntaxToggles(undefined, "nebula");
expect(toggles).toEqual({ ...DEFAULT_SQL_VARIABLE_SYNTAX_TOGGLES, named: false });
expect(enabledSqlParameterSyntaxes(toggles)).not.toContain("named");
expect(resolveSqlVariableSyntaxToggles(undefined, "mysql").named).toBe(true);
});
it("enables every syntax when the database type is unknown", () => {
expect(resolveSqlVariableSyntaxToggles({ mysql: { shell: false } }, undefined)).toEqual(DEFAULT_SQL_VARIABLE_SYNTAX_TOGGLES);
});
@@ -322,6 +322,7 @@ describe("quoteTableIdentifier", () => {
expect(requiresEagerTableMetadataForDataOpen("salesforce")).toBe(true);
expect(requiresEagerTableMetadataForDataOpen("mysql")).toBe(true);
expect(requiresEagerTableMetadataForDataOpen("postgres")).toBe(true);
expect(requiresEagerTableMetadataForDataOpen("nebula")).toBe(true);
// Drivers whose preview works from `SELECT *` keep loading metadata lazily.
expect(requiresEagerTableMetadataForDataOpen("sqlite")).toBe(false);
expect(requiresEagerTableMetadataForDataOpen(undefined)).toBe(false);
@@ -38,6 +38,15 @@ describe("Transwarp driver installation", () => {
});
});
describe("NebulaGraph driver installation", () => {
it("routes the v3 profile and legacy connections to the same Agent package", () => {
for (const profile of [undefined, "nebula", "nebula-v3"]) {
expect(agentDriverInstallKey("nebula", profile)).toBe("nebula");
expect(showAgentDriverInstallHint("nebula", [{ db_type: "nebula", installed: true }], profile)).toBe(false);
}
});
});
describe("shouldApplyDriverStoreFocus", () => {
it("applies when the driver first appears after a list load", () => {
expect(shouldApplyDriverStoreFocus(null, "driver:mysql", false)).toBe(true);
@@ -125,6 +125,7 @@ describe("AGENT_DRIVER_CATEGORY_MAP integrity", () => {
"mongodb",
// graph_ai
"neo4j",
"nebula",
// timeseries
"influxdb",
"iotdb",
@@ -1,5 +1,6 @@
import type { DatabaseType } from "@/types/database";
import { supportsDriverManagement } from "@/lib/database/databaseCapabilities";
import { databaseManifestEntry } from "@/lib/database/databaseDriverManifest";
export interface AgentDriverInstallState {
db_type: string;
@@ -38,6 +39,10 @@ export function agentDriverInstallKey(dbType: DatabaseType | undefined, driverPr
if (dbType === "oracle") return "oracle";
if (dbType === "h2") return "h2";
if (dbType === "transwarp") return "transwarp";
if (dbType === "nebula") {
const entry = databaseManifestEntry(dbType);
return entry?.driverProfiles?.find((profile) => profile.profile === driverProfile)?.agentKey ?? entry?.agentKey;
}
if (dbType === "mongodb") return "mongodb";
if (dbType === "dameng") return "dameng";
if (dbType === "gbase") return driverProfile === "gbase8s" ? "gbase8s" : "gbase8a";
@@ -40,6 +40,7 @@ const DRIVER_STARTUP_FLOOR_TYPES = new Set<DatabaseType>([
"db2",
"informix",
"neo4j",
"nebula",
"cassandra",
"bigquery",
"spanner",
@@ -212,6 +212,9 @@ export function connectionUrlPlaceholder(dbType: DatabaseType, driverProfile?: s
case "mongodb":
return "mongodb://user:password@host:port/database";
case "nebula":
return "nebula://root:password@graphd:9669/space";
case "dynamodb":
return "https://dynamodb.us-east-1.amazonaws.com";
@@ -54,6 +54,7 @@ const SCHEME_PROFILES: Record<string, ConnectionProfile> = {
"mongodb+srv": { type: "mongodb", profile: "mongodb", label: "MongoDB", defaultPort: 27017 },
dynamodb: { type: "dynamodb", profile: "dynamodb", label: "Amazon DynamoDB", defaultPort: 443 },
clickhouse: { type: "clickhouse", profile: "clickhouse", label: "ClickHouse", defaultPort: 8123 },
nebula: { type: "nebula", profile: "nebula", label: "NebulaGraph", defaultPort: 9669 },
sqlserver: { type: "sqlserver", profile: "sqlserver", label: "SQL Server", defaultPort: 1433 },
mssql: { type: "sqlserver", profile: "sqlserver", label: "SQL Server", defaultPort: 1433 },
oracle: { type: "oracle", profile: "oracle", label: "Oracle", defaultPort: 1521 },
@@ -46,6 +46,7 @@ export const AGENT_DRIVER_CATEGORY_MAP: Readonly<Record<string, DriverCategoryKe
ignite3: "analytics",
mongodb: "document",
neo4j: "graph_ai",
nebula: "graph_ai",
"oceanbase-oracle": "domestic",
oracle: "sql",
oscar: "domestic",
@@ -12,6 +12,7 @@ describe("supportsDataDictionary", () => {
it("is hidden for engines that cannot list table metadata", () => {
expect(supportsDataDictionary("redis")).toBe(false);
expect(supportsDataDictionary("nebula")).toBe(false);
expect(supportsDataDictionary("plugin")).toBe(false);
});
});
@@ -147,7 +147,7 @@ export function supportsQueryExecution(dbType?: DatabaseType): boolean {
* that hierarchy, so they must not be offered by sidebar "Add to AI" actions.
*/
export function supportsAiAssistantContext(dbType?: DatabaseType): boolean {
return supportsQueryExecution(dbType) && !usesConnectionOnlyQueryTarget(dbType);
return supportsQueryExecution(dbType) && !usesConnectionOnlyQueryTarget(dbType) && dbType !== "nebula";
}
export function supportsConnectionScopedQueryExecution(dbType?: DatabaseType): boolean {
@@ -182,7 +182,7 @@ export function supportsSqlFileExecution(dbType?: DatabaseType): boolean {
return supportsDatabaseFeature(dbType, "sqlFileExecution");
}
const NON_SQL_IN_LIST_PASTE_TYPES = new Set<DatabaseType>(["neo4j"]);
const NON_SQL_IN_LIST_PASTE_TYPES = new Set<DatabaseType>(["neo4j", "nebula"]);
export function supportsSqlInListPaste(dbType?: DatabaseType): boolean {
if (!dbType) return true;
@@ -200,7 +200,7 @@ export function supportsSchemaDiagram(dbType?: DatabaseType): boolean {
/** Relational engines that can list tables and columns. Independent of diagram support. */
export function supportsDataDictionary(dbType?: DatabaseType): boolean {
return supportsDatabaseFeature(dbType, "metadataBrowse");
return dbType !== "nebula" && supportsDatabaseFeature(dbType, "metadataBrowse");
}
export function supportsDatabaseSearch(dbType?: DatabaseType): boolean {
@@ -256,7 +256,19 @@ export function supportsObjectBrowserTreeNode(dbType: DatabaseType | undefined,
export function supportsTableTruncate(dbType?: DatabaseType): boolean {
return (
!!dbType && dbType !== "impala" && dbType !== "sqlite" && dbType !== "rqlite" && dbType !== "turso" && dbType !== "cloudflare-d1" && dbType !== "duckdb" && dbType !== "influxdb" && dbType !== "influxdb3" && dbType !== "victoriametrics" && dbType !== "manticoresearch" && dbType !== "salesforce"
!!dbType &&
dbType !== "impala" &&
dbType !== "sqlite" &&
dbType !== "rqlite" &&
dbType !== "turso" &&
dbType !== "cloudflare-d1" &&
dbType !== "duckdb" &&
dbType !== "influxdb" &&
dbType !== "influxdb3" &&
dbType !== "victoriametrics" &&
dbType !== "manticoresearch" &&
dbType !== "salesforce" &&
dbType !== "nebula"
);
}
@@ -77,6 +77,7 @@ export const DATABASE_NAMESPACE_CREATION_MATRIX = {
db2: { database: "schema" },
informix: { connection: "database" },
neo4j: { deferred: "database creation depends on edition/admin privileges" },
nebula: { deferred: "space creation requires partition, replica and VID type options" },
cassandra: { deferred: "keyspace creation requires replication options" },
bigquery: { deferred: "dataset creation needs project/location options" },
spanner: { deferred: "database creation requires the Cloud Spanner Admin API" },
@@ -97,6 +97,7 @@ const DATABASE_TYPE_OBJECTS = new Map<DatabaseType, SidebarObjectKind[]>([
["tdengine", TABLE_VIEW_OBJECTS],
["iotdb", TABLE_VIEW_OBJECTS],
["neo4j", TABLE_VIEW_OBJECTS],
["nebula", TABLE_VIEW_OBJECTS],
// others
["influxdb", ["TABLE"]],
["influxdb3", ["TABLE"]],
@@ -78,6 +78,7 @@ export const DATABASE_PROPERTY_EDITING_MATRIX = {
db2: { deferred: "schema properties need product-specific handling" },
informix: { deferred: "schema properties need product-specific handling" },
neo4j: { deferred: "database properties depend on edition/admin privileges" },
nebula: { deferred: "space properties require NebulaGraph-specific administration" },
cassandra: { deferred: "keyspace properties require replication option handling" },
bigquery: { deferred: "dataset properties need project/location-specific handling" },
spanner: { deferred: "database/schema properties are Cloud Spanner Admin API operations" },
+1 -1
View File
@@ -98,7 +98,7 @@ export interface ResolveNewQueryInitialSqlInput extends ResolveNewQueryTableInpu
// Database types whose "table" view does not use standard SQL `SELECT * FROM <table>`
// (e.g. Neo4j uses Cypher). The new-query prefill is skipped for these.
const NEW_QUERY_PREFILL_DISABLED_TYPES: ReadonlySet<DatabaseType | undefined> = new Set<DatabaseType | undefined>(["neo4j"]);
const NEW_QUERY_PREFILL_DISABLED_TYPES: ReadonlySet<DatabaseType | undefined> = new Set<DatabaseType | undefined>(["neo4j", "nebula"]);
export function isNewQueryPrefillSupported(databaseType: DatabaseType | undefined): boolean {
return !NEW_QUERY_PREFILL_DISABLED_TYPES.has(databaseType);
+1 -1
View File
@@ -378,7 +378,7 @@ function findSqlParameterOccurrences(sql: string, options?: SqlParameterOptions)
const occurrences: ParameterOccurrence[] = [];
const databaseType = options?.databaseType;
const nativeSqlServerParameters = collectNativeSqlServerParameters(sql, databaseType);
const supportsNamedParameters = databaseType !== "saphana" && databaseType !== "neo4j";
const supportsNamedParameters = databaseType !== "saphana" && databaseType !== "neo4j" && databaseType !== "nebula";
const enabledSyntaxes = options?.enabledSyntaxes ? new Set(options.enabledSyntaxes) : null;
const isSyntaxEnabled = (syntax: SqlParameterSyntax) => !enabledSyntaxes || enabledSyntaxes.has(syntax);
const complexTypeFieldSeparators = supportsNamedParameters && isSyntaxEnabled("named") ? collectComplexTypeFieldSeparators(sql, databaseType) : new Set<number>();
@@ -27,7 +27,7 @@ export function elasticsearchRestRequestRanges(sql: string, databaseType?: Datab
return requests.length > 0 && requests.every((request) => ELASTICSEARCH_REST_REQUEST.test(request.sql)) ? requests : [];
}
const NON_SQL_EXECUTION_TARGET_TYPES: ReadonlySet<DatabaseType> = new Set(["mongodb", "elasticsearch", "easysearch", "meilisearch", "solr", "qdrant", "milvus", "weaviate", "chromadb", "etcd", "zookeeper", "consul", "mq", "neo4j", "victoriametrics", "salesforce"]);
const NON_SQL_EXECUTION_TARGET_TYPES: ReadonlySet<DatabaseType> = new Set(["mongodb", "elasticsearch", "easysearch", "meilisearch", "solr", "qdrant", "milvus", "weaviate", "chromadb", "etcd", "zookeeper", "consul", "mq", "neo4j", "nebula", "victoriametrics", "salesforce"]);
export function supportsExecutionTargetPicker(databaseType?: DatabaseType): boolean {
return !!databaseType && (databaseType === "redis" || isHttpJsonRestDatabaseType(databaseType) || !NON_SQL_EXECUTION_TARGET_TYPES.has(databaseType));
@@ -72,7 +72,7 @@ export function resolveSqlVariableSyntaxToggles(overrides: SqlVariableSyntaxOver
const partial = dbType ? overrides?.[dbType] : undefined;
return {
positional: partial?.positional ?? true,
named: dbType === "neo4j" ? false : (partial?.named ?? true),
named: dbType === "neo4j" || dbType === "nebula" ? false : (partial?.named ?? true),
shell: partial?.shell ?? true,
mybatis: partial?.mybatis ?? true,
sqlserver: partial?.sqlserver ?? true,
+3 -1
View File
@@ -447,9 +447,11 @@ export function normalizeWhereInput(whereInput?: string): string {
* * Salesforce: SOQL has no `SELECT *`. With no known fields the backend builder
* falls back to the `FIELDS(ALL)` selector, which the org only accepts with
* `LIMIT 200` or less — awaiting the describe keeps every page size working.
* * NebulaGraph: without tag/edge properties, the grid can only show a single
* vertex/edge value instead of separate property columns.
*/
export function requiresEagerTableMetadataForDataOpen(databaseType: DatabaseType | undefined): boolean {
return databaseType === "mysql" || databaseType === "postgres" || databaseType === "salesforce";
return databaseType === "mysql" || databaseType === "postgres" || databaseType === "salesforce" || databaseType === "nebula";
}
export async function buildTableSelectSql(options: BuildTableSelectSqlOptions): Promise<string> {
@@ -1671,6 +1671,7 @@ export const useConnectionStore = defineStore("connection", () => {
informix: "Informix",
phoenix: "Apache Phoenix",
neo4j: "Neo4j",
nebula: "NebulaGraph",
cassandra: "Cassandra",
bigquery: "BigQuery",
spanner: "Cloud Spanner",
@@ -96,6 +96,7 @@ export const CONNECTION_PROFILES = {
dremio: { type: "jdbc", port: 31010, user: "", label: "Dremio", icon: "dremio" },
jdbcx: { type: "jdbc", port: 0, user: "", label: "JDBCX", icon: "jdbcx" },
neo4j: { type: "neo4j", port: 7687, user: "neo4j", label: "Neo4j", icon: "neo4j" },
nebula: { type: "nebula", port: 9669, user: "root", label: "NebulaGraph", icon: "nebula" },
cassandra: { type: "cassandra", port: 9042, user: "cassandra", label: "Cassandra", icon: "cassandra" },
bigquery: { type: "bigquery", port: 443, user: "", label: "BigQuery", icon: "bigquery", host: "https://www.googleapis.com/bigquery/v2" },
spanner: { type: "spanner", port: 443, user: "", label: "Cloud Spanner", icon: "spanner" },
@@ -206,6 +207,7 @@ export const CONNECTION_PROFILE_ICONS = {
dremio: "dremio",
jdbcx: "jdbcx",
neo4j: "neo4j",
nebula: "nebula",
cassandra: "cassandra",
bigquery: "bigquery",
spanner: "spanner",
@@ -311,6 +313,7 @@ export const CONNECTION_PICKER_OPTIONS = [
{ value: "dremio", label: "Dremio", category: "analytics" },
{ value: "jdbcx", label: "JDBCX", category: "sql" },
{ value: "neo4j", label: "Neo4j", category: "graph_ai" },
{ value: "nebula", label: "NebulaGraph", category: "graph_ai" },
{ value: "cassandra", label: "Cassandra", category: "document" },
{ value: "bigquery", label: "BigQuery", category: "analytics" },
{ value: "spanner", label: "Cloud Spanner", category: "sql" },
@@ -59,6 +59,7 @@ export const DATABASE_TYPES = [
"db2",
"informix",
"neo4j",
"nebula",
"cassandra",
"bigquery",
"spanner",
+1
View File
@@ -66,6 +66,7 @@ const BRIDGE_REQUIRED_TYPES: &[&str] = &[
"informix",
"iris",
"neo4j",
"nebula",
"cassandra",
"bigquery",
"spanner",
@@ -2123,6 +2123,46 @@
"driverManagement": true
}
},
{
"dbType": "nebula",
"label": "NebulaGraph",
"runtimeMode": "agent",
"mcpMode": "bridge",
"agentKey": "nebula",
"driverProfiles": [
{
"profile": "nebula-v3",
"agentKey": "nebula",
"label": "NebulaGraph 3.x",
"storeVisible": false
}
],
"driverStoreVisible": true,
"driverStoreOrder": 53,
"singleConnectionPool": false,
"metadataConnectionScoped": false,
"skipTcpProbe": true,
"defaultPort": 9669,
"supportLevel": "browse",
"capabilities": {
"queryExecution": true,
"metadataBrowse": true,
"objectBrowser": true,
"objectSource": true,
"schemaSearch": false,
"diagram": false,
"tableDataEdit": false,
"tableStructureEdit": false,
"tableImport": false,
"dataTransfer": false,
"sqlFileExecution": false,
"databaseCreate": false,
"fieldLineage": false,
"sqlExplain": false,
"userAdmin": false,
"driverManagement": true
}
},
{
"dbType": "cassandra",
"label": "Apache Cassandra",
+1
View File
@@ -330,6 +330,7 @@ macro_rules! agent_connection_pool_database_type {
| DatabaseType::Db2
| DatabaseType::Informix
| DatabaseType::Neo4j
| DatabaseType::Nebula
| DatabaseType::Cassandra
| DatabaseType::Bigquery
| DatabaseType::Spanner
@@ -82,6 +82,15 @@ mod tests {
assert!(!driver_store_entries().any(|(key, _)| key.starts_with("transwarp-")));
}
#[test]
fn nebula_v3_profile_reuses_the_legacy_agent_package() {
for profile in [None, Some("nebula"), Some("nebula-v3")] {
assert_eq!(agent_key(&DatabaseType::Nebula, profile), Some("nebula"));
}
assert_eq!(driver_store_entries().filter(|(key, _)| *key == "nebula").count(), 1);
assert!(!driver_store_entries().any(|(key, _)| key == "nebula-v3"));
}
#[test]
fn h2_profiles_share_the_same_agent() {
assert_eq!(agent_key(&DatabaseType::H2, None), Some("h2"));
@@ -268,6 +268,7 @@ pub fn supports_sql_query(database_type: DatabaseType) -> bool {
| DatabaseType::InfluxDb3
| DatabaseType::VictoriaMetrics
| DatabaseType::Neo4j
| DatabaseType::Nebula
| DatabaseType::Etcd
)
}
@@ -559,6 +559,7 @@ fn unsupported_pagination_type(database_type: Option<DatabaseType>) -> bool {
database_type,
Some(
DatabaseType::Neo4j
| DatabaseType::Nebula
| DatabaseType::MongoDb
| DatabaseType::Redis
| DatabaseType::Salesforce
@@ -905,7 +905,7 @@ pub fn profile_for(db_type: DatabaseType) -> DdlDialectProfile {
// Non-tabular / not applicable for relational CREATE TABLE
Redis | MongoDb | DynamoDb | Elasticsearch | Easysearch | Solr | Meilisearch | Qdrant | Milvus | Weaviate
| ChromaDb | Neo4j | Cassandra | Etcd | ZooKeeper | Nacos | Consul | InfluxDb3 | VictoriaMetrics
| ChromaDb | Neo4j | Nebula | Cassandra | Etcd | ZooKeeper | Nacos | Consul | InfluxDb3 | VictoriaMetrics
| MessageQueue | Mqtt | Hbase | Salesforce => conservative_ansi(db_type),
}
}
@@ -170,6 +170,7 @@ pub fn quote_table_identifier(database_type: Option<DatabaseType>, name: &str) -
}
Some(DatabaseType::Informix) if is_simple_informix_identifier(name) => name.to_string(),
Some(DatabaseType::Neo4j) => format!("`{}`", name.replace('`', "``")),
Some(DatabaseType::Nebula) => format!("`{}`", name.replace('\\', "\\\\").replace('`', "\\`")),
Some(DatabaseType::SqlServer) => format!("[{}]", name.replace(']', "]]")),
_ => format!("\"{}\"", name.replace('"', "\"\"")),
}
@@ -200,6 +200,9 @@ pub fn build_table_data_select_sql_with_database(
if database_type == Some(DatabaseType::Neo4j) {
return build_neo4j_table_select_sql(&options, limit);
}
if database_type == Some(DatabaseType::Nebula) {
return build_nebula_table_select_sql(&options, limit);
}
if database_type == Some(DatabaseType::Salesforce) {
return build_salesforce_table_select_sql(&options, limit);
}
@@ -939,6 +942,47 @@ pub(super) fn build_neo4j_table_select_sql(options: &TableDataSelectSqlOptions,
format!("MATCH (n:{label}){where_clause} RETURN {returns}{order}{skip} LIMIT {limit};")
}
fn build_nebula_table_select_sql(options: &TableDataSelectSqlOptions, limit: usize) -> String {
let quote = |name: &str| quote_table_identifier(Some(DatabaseType::Nebula), name);
let kind = quote(&options.table_name);
let is_edge = options.table_type.as_deref().is_some_and(|kind| kind.eq_ignore_ascii_case("VIEW"));
let (pattern, identity, projection) = if is_edge {
let projection = if options.columns.is_empty() {
"e AS `edge`".to_string()
} else {
options
.columns
.iter()
.map(|column| format!("e.{} AS {}", quote(column), quote(column)))
.collect::<Vec<_>>()
.join(", ")
};
(format!("MATCH ()-[e:{kind}]->()"), "src(e) AS `_src`, dst(e) AS `_dst`, rank(e) AS `_rank`", projection)
} else {
let projection = if options.columns.is_empty() {
"v AS `vertex`".to_string()
} else {
options
.columns
.iter()
.map(|column| format!("v.{kind}.{} AS {}", quote(column), quote(column)))
.collect::<Vec<_>>()
.join(", ")
};
(format!("MATCH (v:{kind})"), "id(v) AS `_vid`", projection)
};
let predicate = normalize_where_input(options.where_input.as_deref());
let where_clause = if predicate.is_empty() { String::new() } else { format!(" WHERE {predicate}") };
let order = options
.order_by
.as_deref()
.filter(|order| !order.trim().is_empty())
.map(|order| format!(" ORDER BY {order}"))
.unwrap_or_default();
let skip = options.offset.filter(|offset| *offset > 0).map(|offset| format!(" SKIP {offset}")).unwrap_or_default();
format!("{pattern}{where_clause} RETURN {identity}, {projection}{order}{skip} LIMIT {limit};")
}
/// Salesforce's `FIELDS(ALL)` selector is only legal with a LIMIT of 200 or less.
const SALESFORCE_FIELDS_ALL_MAX_LIMIT: usize = 200;
@@ -17,6 +17,7 @@ fn transfer_identifier_policy_preserves_legacy_output() {
#[test]
fn quotes_identifiers_by_database_type() {
assert_eq!(quote_table_identifier(Some(DatabaseType::Nebula), "tag`name"), "`tag\\`name`");
assert_eq!(quote_table_identifier(Some(DatabaseType::Mysql), "user`name"), "`user``name`");
assert_eq!(quote_table_identifier(Some(DatabaseType::ClickHouse), "user`name"), "`user``name`");
assert_eq!(quote_table_identifier(Some(DatabaseType::Doris), "user`name"), "`user``name`");
@@ -50,6 +51,29 @@ fn quotes_identifiers_by_database_type() {
assert!(is_schema_aware(DatabaseType::Argo));
}
#[test]
fn builds_nebula_tag_and_edge_queries() {
let tag = build_table_data_select_sql(TableDataSelectSqlOptions {
database_type: Some(DatabaseType::Nebula),
table_name: "player".into(),
table_type: Some("TABLE".into()),
columns: vec!["name".into()],
limit: Some(20),
..Default::default()
});
assert_eq!(tag, "MATCH (v:`player`) RETURN id(v) AS `_vid`, v.`player`.`name` AS `name` LIMIT 20;");
let edge = build_table_data_select_sql(TableDataSelectSqlOptions {
database_type: Some(DatabaseType::Nebula),
table_name: "serve".into(),
table_type: Some("VIEW".into()),
columns: vec!["start_year".into()],
limit: Some(10),
..Default::default()
});
assert_eq!(edge, "MATCH ()-[e:`serve`]->() RETURN src(e) AS `_src`, dst(e) AS `_dst`, rank(e) AS `_rank`, e.`start_year` AS `start_year` LIMIT 10;");
}
/// Spanner databases are created in one of two immutable dialects. The connected
/// agent reports the correct identifier quote; when it is missing the static mapping
/// must fall back to GoogleSQL (backticks), because GoogleSQL treats double quotes as
@@ -1185,6 +1185,7 @@ impl ConnectionConfig {
DatabaseType::Db2 => format!("db2://{host}:{port}{db_part}"),
DatabaseType::Informix => format!("informix://{host}:{port}{db_part}"),
DatabaseType::Neo4j => format!("neo4j://{host}:{port}{db_part}"),
DatabaseType::Nebula => format!("nebula://{host}:{port}{db_part}"),
DatabaseType::Cassandra => format!("cassandra://{host}:{port}{db_part}"),
DatabaseType::Bigquery => format!("bigquery://{host}/{db_part}"),
DatabaseType::Spanner => self.spanner_display_url(&host, port),
@@ -1442,6 +1443,9 @@ impl ConnectionConfig {
DatabaseType::Neo4j => {
format!("neo4j://{}:{}@{host}:{port}{db_part}", username, password)
}
DatabaseType::Nebula => {
format!("nebula://{}:{}@{host}:{port}{db_part}", username, password)
}
DatabaseType::Cassandra => {
format!("cassandra://{}:{}@{host}:{port}{db_part}", username, password)
}
Binary file not shown.

After

Width:  |  Height:  |  Size: 42 KiB

+37
View File
@@ -0,0 +1,37 @@
schemaVersion: 1
order: 565
dbType: nebula
rustVariant: Nebula
label: NebulaGraph
runtimeMode: agent
mcpMode: bridge
agentKey: nebula
driverProfiles:
- profile: nebula-v3
agentKey: nebula
label: NebulaGraph 3.x
storeVisible: false
driverStoreVisible: true
driverStoreOrder: 53
singleConnectionPool: false
metadataConnectionScoped: false
skipTcpProbe: true
defaultPort: 9669
supportLevel: browse
capabilities:
queryExecution: true
metadataBrowse: true
objectBrowser: true
objectSource: true
schemaSearch: false
diagram: false
tableDataEdit: false
tableStructureEdit: false
tableImport: false
dataTransfer: false
sqlFileExecution: false
databaseCreate: false
fieldLineage: false
sqlExplain: false
userAdmin: false
driverManagement: true
@@ -545,6 +545,13 @@ profiles:
port: 7687
user: neo4j
category: graph_ai
- id: nebula
dbType: nebula
label: NebulaGraph
icon: nebula
port: 9669
user: root
category: graph_ai
- id: cassandra
dbType: cassandra
label: Cassandra