From ef05a6a19b0cdb1fd8ee67cb9a60aac2c2484618 Mon Sep 17 00:00:00 2001 From: Hand Sonic <8078023+HandSonic@users.noreply.github.com> Date: Tue, 29 Sep 2026 21:03:12 +0800 Subject: [PATCH] feat(nebula): add NebulaGraph 3.x browser support --- .github/scripts/bump-agent-versions.mjs | 3 +- .github/scripts/ci-config.mjs | 1 + .github/scripts/ci-plan.test.mjs | 6 +- .github/scripts/label-pull-request.mjs | 1 + .../scripts/reuse-agent-release-assets.mjs | 2 +- .github/workflows/agents-release.yml | 46 ++- agents/README.md | 1 + agents/README.zh-CN.md | 1 + agents/drivers/nebula-go/go.mod | 11 + agents/drivers/nebula-go/go.sum | 55 +++ agents/drivers/nebula-go/main.go | 347 ++++++++++++++++++ agents/drivers/nebula-go/main_test.go | 135 +++++++ agents/drivers/nebula-go/metadata.go | 282 ++++++++++++++ agents/drivers/nebula-go/protocol_error.go | 79 ++++ agents/drivers/nebula-go/query.go | 262 +++++++++++++ agents/metadata-constraint-coverage.tsv | 1 + .../scripts/driver_release_packages_test.py | 20 + agents/scripts/validate_agents.py | 1 + agents/scripts/version_agent_artifacts.py | 2 +- agents/versions.json | 1 + apps/desktop/public/icons/database/nebula.png | Bin 0 -> 9347 bytes .../connection/ConnectionDialog.vue | 36 +- .../editor/EditorSettingsDialog.vue | 4 +- apps/desktop/src/components/grid/DataGrid.vue | 2 +- .../src/components/icons/DatabaseIcon.vue | 1 + .../src/components/objects/ObjectBrowser.vue | 38 +- ...owser.tableCopyClipboard.component.spec.ts | 15 + .../sidebar/SidebarTreeRuntimeHost.vue | 35 +- ...rTreeRuntimeHost.nebulaContextMenu.spec.ts | 80 ++++ .../useSidebarDataOpenRuntime.spec.ts | 34 ++ .../connection/connectionUrl.nebula.spec.ts | 21 ++ .../database/databaseDriverManifest.spec.ts | 4 + .../lib/__tests__/sql/sqlParameters.spec.ts | 9 + .../__tests__/sql/sqlVariableSyntax.spec.ts | 7 + .../__tests__/table/tableSelectSql.spec.ts | 1 + .../__tests__/agentDriverInstallHint.spec.ts | 9 + .../driverCategoryDefinitions.spec.ts | 1 + .../lib/connection/agentDriverInstallHint.ts | 5 + .../connection/connectionAttemptTimeout.ts | 1 + .../lib/connection/connectionPresentation.ts | 3 + .../src/lib/connection/connectionUrl.ts | 1 + .../connection/driver-category-definitions.ts | 1 + .../__tests__/dataDictionarySupport.spec.ts | 1 + .../lib/database/databaseFeatureSupport.ts | 20 +- .../lib/database/databaseNamespaceCreation.ts | 1 + .../database/databaseObjectCapabilities.ts | 1 + .../lib/database/databasePropertyEditing.ts | 1 + apps/desktop/src/lib/sql/newQueryContext.ts | 2 +- apps/desktop/src/lib/sql/sqlParameters.ts | 2 +- .../desktop/src/lib/sql/sqlStatementRanges.ts | 2 +- apps/desktop/src/lib/sql/sqlVariableSyntax.ts | 2 +- apps/desktop/src/lib/table/tableSelectSql.ts | 4 +- apps/desktop/src/stores/connectionStore.ts | 1 + .../src/types/generated/connectionProfiles.ts | 3 + .../src/types/generated/databaseTypes.ts | 1 + crates/dbx-cli/src/main.rs | 1 + .../assets/database-drivers.manifest.json | 40 ++ crates/dbx-core/src/connection/mod.rs | 1 + crates/dbx-driver-agent/src/agent_catalog.rs | 9 + .../dbx-sql-core/src/query_execution_sql.rs | 1 + crates/dbx-sql-data/src/query_result_sql.rs | 1 + .../src/sql_dialect/ddl_profile.rs | 2 +- .../src/sql_dialect/identifiers.rs | 1 + .../src/sql_dialect/table_select.rs | 44 +++ .../dbx-sql-dialect/src/sql_dialect/tests.rs | 24 ++ crates/dbx-types/src/models/connection.rs | 4 + docs/screenshots/nebula-graph-connection.png | Bin 0 -> 43197 bytes plugins/connection-types/nebula.yaml | 37 ++ .../connection-types/profiles/catalog.yaml | 7 + 69 files changed, 1737 insertions(+), 41 deletions(-) create mode 100644 agents/drivers/nebula-go/go.mod create mode 100644 agents/drivers/nebula-go/go.sum create mode 100644 agents/drivers/nebula-go/main.go create mode 100644 agents/drivers/nebula-go/main_test.go create mode 100644 agents/drivers/nebula-go/metadata.go create mode 100644 agents/drivers/nebula-go/protocol_error.go create mode 100644 agents/drivers/nebula-go/query.go create mode 100644 apps/desktop/public/icons/database/nebula.png create mode 100644 apps/desktop/src/components/sidebar/__tests__/SidebarTreeRuntimeHost.nebulaContextMenu.spec.ts create mode 100644 apps/desktop/src/lib/__tests__/connection/connectionUrl.nebula.spec.ts create mode 100644 docs/screenshots/nebula-graph-connection.png create mode 100644 plugins/connection-types/nebula.yaml diff --git a/.github/scripts/bump-agent-versions.mjs b/.github/scripts/bump-agent-versions.mjs index 6b6bb6f81..310520a0a 100644 --- a/.github/scripts/bump-agent-versions.mjs +++ b/.github/scripts/bump-agent-versions.mjs @@ -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", diff --git a/.github/scripts/ci-config.mjs b/.github/scripts/ci-config.mjs index 91a256e41..0aa5dd3bd 100644 --- a/.github/scripts/ci-config.mjs +++ b/.github/scripts/ci-config.mjs @@ -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 }, ]; diff --git a/.github/scripts/ci-plan.test.mjs b/.github/scripts/ci-plan.test.mjs index 57e3b98b7..e9e69ac78 100644 --- a/.github/scripts/ci-plan.test.mjs +++ b/.github/scripts/ci-plan.test.mjs @@ -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); }); diff --git a/.github/scripts/label-pull-request.mjs b/.github/scripts/label-pull-request.mjs index a81e06b25..223b3afbd 100644 --- a/.github/scripts/label-pull-request.mjs +++ b/.github/scripts/label-pull-request.mjs @@ -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", diff --git a/.github/scripts/reuse-agent-release-assets.mjs b/.github/scripts/reuse-agent-release-assets.mjs index 110c04226..ae27a018b 100644 --- a/.github/scripts/reuse-agent-release-assets.mjs +++ b/.github/scripts/reuse-agent-release-assets.mjs @@ -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", diff --git a/.github/workflows/agents-release.yml b/.github/workflows/agents-release.yml index db8143d45..b43ce6a24 100644 --- a/.github/workflows/agents-release.yml +++ b/.github/workflows/agents-release.yml @@ -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 diff --git a/agents/README.md b/agents/README.md index 3b9147e87..add479359 100644 --- a/agents/README.md +++ b/agents/README.md @@ -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 | diff --git a/agents/README.zh-CN.md b/agents/README.zh-CN.md index cead6e6dd..f8c66ff57 100644 --- a/agents/README.zh-CN.md +++ b/agents/README.zh-CN.md @@ -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 | diff --git a/agents/drivers/nebula-go/go.mod b/agents/drivers/nebula-go/go.mod new file mode 100644 index 000000000..a55b46a8f --- /dev/null +++ b/agents/drivers/nebula-go/go.mod @@ -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 +) diff --git a/agents/drivers/nebula-go/go.sum b/agents/drivers/nebula-go/go.sum new file mode 100644 index 000000000..d44cf94e7 --- /dev/null +++ b/agents/drivers/nebula-go/go.sum @@ -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= diff --git a/agents/drivers/nebula-go/main.go b/agents/drivers/nebula-go/main.go new file mode 100644 index 000000000..c47d3151a --- /dev/null +++ b/agents/drivers/nebula-go/main.go @@ -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 +} diff --git a/agents/drivers/nebula-go/main_test.go b/agents/drivers/nebula-go/main_test.go new file mode 100644 index 000000000..c4806444c --- /dev/null +++ b/agents/drivers/nebula-go/main_test.go @@ -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") + } +} diff --git a/agents/drivers/nebula-go/metadata.go b/agents/drivers/nebula-go/metadata.go new file mode 100644 index 000000000..47f4eeeb9 --- /dev/null +++ b/agents/drivers/nebula-go/metadata.go @@ -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 +} diff --git a/agents/drivers/nebula-go/protocol_error.go b/agents/drivers/nebula-go/protocol_error.go new file mode 100644 index 000000000..45406acb4 --- /dev/null +++ b/agents/drivers/nebula-go/protocol_error.go @@ -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) +} diff --git a/agents/drivers/nebula-go/query.go b/agents/drivers/nebula-go/query.go new file mode 100644 index 000000000..066d27b76 --- /dev/null +++ b/agents/drivers/nebula-go/query.go @@ -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 +} diff --git a/agents/metadata-constraint-coverage.tsv b/agents/metadata-constraint-coverage.tsv index b3345f91a..1e695d626 100644 --- a/agents/metadata-constraint-coverage.tsv +++ b/agents/metadata-constraint-coverage.tsv @@ -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. diff --git a/agents/scripts/driver_release_packages_test.py b/agents/scripts/driver_release_packages_test.py index da3da6e02..7773a9cde 100644 --- a/agents/scripts/driver_release_packages_test.py +++ b/agents/scripts/driver_release_packages_test.py @@ -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) diff --git a/agents/scripts/validate_agents.py b/agents/scripts/validate_agents.py index 3a8e6c7f6..cad121643 100644 --- a/agents/scripts/validate_agents.py +++ b/agents/scripts/validate_agents.py @@ -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", diff --git a/agents/scripts/version_agent_artifacts.py b/agents/scripts/version_agent_artifacts.py index 038f5bb6f..238f7abde 100644 --- a/agents/scripts/version_agent_artifacts.py +++ b/agents/scripts/version_agent_artifacts.py @@ -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", diff --git a/agents/versions.json b/agents/versions.json index 5d8c15ab5..a3381da6f 100644 --- a/agents/versions.json +++ b/agents/versions.json @@ -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", diff --git a/apps/desktop/public/icons/database/nebula.png b/apps/desktop/public/icons/database/nebula.png new file mode 100644 index 0000000000000000000000000000000000000000..45df6ec4ae6153b423dc7c91399309608df2bd8d GIT binary patch literal 9347 zcmXAPby(ER_x5L(U3TfE7wHa>l9pJyQxs4d1Rm*4ul@`}bi!=m~US*tim&#ap7vVN-E+ZM7vhbswj$?vz(*ged*J>K-u}Di(HTS4Ne-` zzPCC4`$w|BCw+6{(@SY&Qxf_4{wT7R;0@SI?YI7pu8%z<&*9J7yiAK>4qoeuO!Ros zALAT-HpolPxS|yxC8|s?Xl!X6X#0u!%~abh0#pEjqhyXw4D3+i z&3>C3u|!xF=T~E^cni?5&G+vv!%rQLQ)`&Y@%V$qdYMD{c!5}ol0a>pSQY2Y^dxc4 za=a=JoOcg+0d`}zBZ+#F92P^@qpg2lR*_v`f$%~8b(Nk%NOv9+2M1Npk!ZKP0FnwV zQ@zBAL(mBr0TZ2(dAz(3X$#OMMli_9f_d}fDcJQYBf#4_W(YwW)Ykvs8@Kqk<81*( zFF{<}=NiY^T=ddPnmG~;;_4;aJdipH^dtT?9sKURC0vGL1*h)N7kFs#6V9{`n<7@; zd0T^_SX_qE6tUs_%<8r8*AkZ3d8D+mqj@oc2M;cds_ z-XA`I$)sHX6XE(=FISHGE1&3;dt1; zh5a!v0L9ijx;4+YPcGB(k$n#Qr#y1+nNc-bb{S@c*?sPUprp}%{YYt2YcnQE(4tQs z-g~YiuwDePxoy2-2Zs8Fcz5s=b?JUg)l|EPgSG$)9r~Vvf zLI=$^+lp^u!5sD=_}qDe6Z78-oR7O*{oeD;^FdrP^Sp#NS2hVhgB=LEQ#KD!P{>Y#~@@X6o#&|4TsySPw0ZLeAf1+uwSeBIVg_Na6Z`ywL5B zmcPC2I>_78M-;*lOaW@gL*~z=aXAR(DMJ<^Opuii%E6m`rc%NJH&-k4&9N&$)$#KfIH!+XT6EPuB&d!pTd8==c37GpdPfF`q$J^9P_+>K0PG^ z4xruTq;oF#bl$J!s=VLizP%gA&cH9OnheCXcf@~@u^eV}*X3DPa$EQIp#S;8Nxv@J zstiszjEmFknn3W>7)etk{p=TcsUjfS@xEitIvpl?I-h+?zEh|~Jt$V%MqSyQ^Tu?1 zuXUNYGp)WPbiv=hWDc1(bAnS&7d>@l{_O~IxbTZUnyUZ5)Ag8tfNyULaL#kqxpsGg zvl`CfcXU_(*KK1IzM2!(=3O7bOreLyg5Y@&7F@aR4XlrxfOzLr9?w9S8Pwt^?t=7y z+?=&f>gxBdSQch;>6S&1>Q5OACk&DAK{}p*fLJ23E$Biw4)eu8*wh1+G)DTsdK4KU z=S@U~TVpevlB1yx3UVm`GoP;oSE~$%6|3U!9i&nmvjFYyfJP^OjuJU}Zf!MvX$K`%=R0zxXxapyc=??$A7zTL;3H;1Qwy9C!@jT9 ztH=%TV&Ne54;i)_%|;HBKZPQ0YNQ+(R`K(dQ-PFnIr|gr{!A04ZR%&GzIMqBi8}fc z@FMG8Lo&AFD{bsdBC`(^+j@GFn_egW8zxSP9xOTf-ETST@?yKYEZDW)^kexGEI!kx zXY2TjP*;7|M_nlX7+nhziLPmbMJuNd3M>3()o1TSxfHb8JA3)%M#<#!ugqKstSP7X z=8A!dSrX=&WXP~SI7VAIx{ZUFEKK(+o3Gp+`dg;5>`>(a;Ei$ zD*pg?o!j-hE$fV*A6WUV>ljD~cp`~|f|33;8oC~bI%l~kiSLseF%k<+G+WzgOHks` z0ipR9ENn9g`}jzI;ETiU;`Z-x>pPy70gGd}p)2`sueeIO-VkG`46d^U(+^nG56Scf zbT4nNv=IYs8>p87bNEqJUEGnY_rJo|*U=Xh0zuiwss@+;D*DE9+Po`oHQ}Kq(8#lT z1(>Jy@#oGt|41eoKd1+^rz1XYKIbxdN_coZ_63RQ z&Iq2J*OK?)6ThMH()1=s$3w6DlP;3TL1Q;^EaAJC0@WLcQ zFTsoF$5F}57bWgpQldLS&eCSbToFaXwXdxYqV`0ru1((C|HRil!#vI9YYl-eL1(od z*yEuSM;YBjdIM6w_iZIoFJ&cuYBGh#$!dlJVM^8ZUQ-5E2E4g<`K`KVGb``|zx1}K zXToiOG{`mJ5@&ohZ)eKwrwO^VSfIpR)r(LDl$+=n5}q?9D6)1DT9kE@E@C9ek?n)G3s3X5w^6~a<%csQ&< zfKo9o?iW;LDMRK|_}_9^_vu5s{ypR!ru3o$W$GV*)d|qLX{@UuZ8Kx-LAR2|uo9Ec zS~4<8mLVE0i-BR1vOyhTxe94X{qBCX|0dJOfZ22M{np026z{aV7dCgMjdy(ettH@A z)(`W*Kmq7Lo%*ECMv0 z(nQp}frif1cHAI@xxwEe0lB?q{-00oG$xibYbF2Cx{``!ymbrqc(E&zif>$)JxF|S zlK9ueZ#esL%5}!Pem>szv?UR!U~-GCEpu|g*d!AWr^-Ui+FW~lGHe!INtn(yZuj*O zz_i)x6iLpvjv=@#)VO^t`Pv!p7EvH<#x)L)&hV?By@`2kcMMj%i#%)sq@ZhUYU0oH zI;D6mQS3M{{QL8B)f_Cw1Gcz2PP8mlacs075g`pR^&b0D=^t>=(z$lk$&6axKZfxT zq3VAI!ycN%hN7gj*ytN+W4J943>TnPZtZ`)@Ps_ zBOgX$VK}C+nNkir$28#L_$-VfB#??a4h8Wf?@He@h#EjmIg80&yayKW5IyzCys z&;+XSj7-P&Lp5umg>>UcJi4Jto5O!$&_IiarixNyuih7k# ztkxMcj-F>mZ~T+hJpdeCTFR^|?BP?b%>;?n=|Z2AX+y={`H>rSyvE>yHXHGuQV58W zqh~K)J`z;|9DZsKB~$TV^Iim?ei!Exfva2T>JMg$Bc+!+EuTX8OvW83xFoufC0FsnRf9(oog+`ZDW%&Z1Dv5E{|7unJ_bNC2YGJA_sYeLsr3 zeGt8GFA%~`&3qz^mmSwoh>LURMvH)x8p<_*`SZ6(Td2*N(|fxO*+~ujfSir{jFI$Am zoD4JCIx*G)+Xg8IlgP6VeND~o{H>`6StV+Ef9oU~>9go2QJ(x#tz3@lW#wt59{>sm zc=j3xX#>jE$*}u>)Yu{kd2ovN8lAb4lBlcyAPnH>ZEZKhQsl?MC)f2V@llX>h8$U5 z@(LVKtr0wEQk?3&t#GT^J)j0By)t;6obO=Z4v0e?BS4MHnL^E8Slt@S=;i{T6|G1X z8qsGCN43)_-*uEfY(I*!iP*to|9KhGe*XSlsOTX`vH*yeh4M4)Pyqv7jUf25=lD3h zz5sX5+Z%xgFu*0iRG99CB!4BMJQN~ov=DJ^%-TZkzDD~esHl3N(Lh=7`4Ih z1lK^dfS0T5_8AFNp(_OS9pP=9y?%A;RSvZ$`!iQg<SCq58huFT^PHfrqic{!~EW zQ(`y^9(*)&RXOW^OwN>(TJ?*Apk|OGc;ttN}nQ`LV!<) z`}HGC;R5-5W&rUWE|dxkj8m`P`uF1--e&IYV&p73ZF#ihi*x+2V1o~FiaSvFj0zl) z4KOhn3M#&wu{~Iw93fEMTF+0=mL_~M@*i<5kuZ)(soV@V%7Bm*-?2jWNs-(G9wzS6 zmgf(J9e*AAo7Gt2D=1pGuqN5TLEu-n1V#M(xxho&v!g~G`4nnq?RJCA@SvN=LLcwv z1i9Z0MG@9uSD3Ma2ss^De^mOfF_o8&_};9pOvn#V$~FLw(li4k(P1Y5;zM~W!utj@ zAI8`MN_TG7V&}Hc%;`8U%M^>ATwsoaC6k0w0{KEQ#W%G_t5&LQRQ4>P&1Z;2xW_gGgNqes< z0gSB(Ft76Om7tFzM@$(-Z=B&0^j~5>0sIAk-j>JB%9u(|2_f?Vao{F(`5Y)j#wkA| z|AUw}dKc13fPRg%cm~L57E0hwzVb>5$2OsOFx-~$56@g%Z7Jid450@6R&xTAG>*LzCKfKt)u>0c^M zZ1|vIU+gr&RP>f#WqMY)qfcacRT#1h%>@CJOEVMZU7Fa^2Oi1!sBo-&6IA1x=;SWW zEhs<#j=)6$R9s*Km#yhvJ;kfgML|KD5>w5Q@$@`uaItgK>o6mTB}Sft0niP^I`O-z z{vbf! zjOKs<mR@g0ZV5!lHzUB$?r4(>^cppEk4(#J<<6HSqy_m9=8w zIq|Jt{EIsO5s@gR=;j=PiY?LG#4)UZAdSxmZhHu;xr#NWPBFo=2PyF@Fax2CHm<0D68* zeOBCvF?R&`e;CW{w1@lE!W4AWi&&R$FAjVGR3AZg+aHPVy7`R~;#xnfx6@Wmyp;;r z=swz|vOp*kKPDVBN+wNqkFkFeP*uK~m zy7m?QAt`!qEhj$eaOr!t0@Qe&|ATdTw)$+Q=k}LB;Z-;U6tB&NfqH3le`d%8PWPzs zicO0g&WBX}=7)f0UtH9$(@LB8SnBG*uHRoH5NHjbg;SvPMz%(OR&R21L~M?6a_0<# zi66Vy^U8MxXl}xP(S?6Vkp^%9WyY02W}vbk3Sgj(Pr=skzRy0LcczsY_ zysjFg3kS}gJ@BCnn?V4IZlah<*lq zaJgIflpa+yz7e9R4t%{b!KAL$d}5o!wx-?Z#&OAoAl@{IA?C*85WLArI+GhJ`M2Zb z4=}v{44^4N_9C0_Tw%mdR~PZ*v~IP&zj+P^+~G%5EtfX>*yTb$8oT$4@}Hc=U{N4` z3p+IsoH_jUnlPuVSq7ZKpFN*8C!7g`|7b?-L92N6i-)H-ewx7sxig4&Xk+W~2}6j5 zr!YMTklbx@g_#~>6d?M%BtlLG#cWXsT(6mXpUp&?&bbQRoFDU8oBwF?-pV&gPXqFN z7lawzXS*SoAchVqu?R4Gx@>n6xRO>dcO)(6gftf}gW?T&fl4pYc4a_+07unR9`r#P zOZ}x5{z+LgPQO;~#pvS}Wk9*tK8;L4i~{v}ooAGYR6JwK=fJ3GPeS6Q7qP{`z61Yl zB`MnALFh!7E~>=Y$v@Lo>^w*RA60yU0~J+Hul`a<)ZOmaY5B7aQC{319_m8RB6LwF zG49emATVYLRoBzHp_^)^cAJnDQhVG5pgAu~BWp6usAzCEAmwZ7R+9}dq1uspF*vG3 z8F>i8SwQ*| z+&1?)`D>DLzTH&0#ygL>lVG*qYx$4VACoZ28Jf{3_|G5py9g7w%~Z=e&Ryx^kyi{z z*wEzcoLuCB1w{W|#XSY-bO|Vgl>n2oVuyREx!F;kYLOw+B7K8r=?#rjEkCMC08f($ z33cRu_DFV)1S0l=Bg2_5jm|OJZ!?6p*n0LwFO>Zhloh&ceM(q>L3}Bo@owj5@VgPv z!f|dXt+rklX4t}pBQo;IRU>QSS^C53dHPn~*~_nYb23f_*4okJt*MeudnXZt4h!lG zC`$0~HshZ&NGpcnCW7eqW3O_?>C zY0%on1F>it@A$T`c?R=><9LP$E`Pz8MZ`$EQ+Hllo?rn=ur~2uoWw#c5|mPUU;_h* z`<k}^=rPnSG;GDI+i@=E?CV34cvvZ$~2_~)$#Elw?Y)OJ|sz0{RH_QNm-G{P3PRkI8mYMG}uf#(B_Am4}=uxv8+%7Ouk-L z1I`lOrt_?v?_s-Gqd_T7xs<4m<#-EAt^(R4Y&KXx7eGfdYm=_9k(0=cv^lbALzR`zO#9N+i z*Q|YWJu?4VAkkU?9=>haq%0SNc_%$~4NV`SOTv?n25sQ_;e>dOOqVT#&NTsBUx}Du*ho2qz-m@q* zC;{r`qLP%%#k26xlF+94C5K|A8*~l4q?TRrf76C` zwg6m%caT(z3w!12q3evP=L?fti9@nY7Y&4K(m=VL+OXkr?{kKHY9U@!T|S@u;-B-7 z*B3f=T$RM#l6Mn@Y(PZq>XsG}LwMEa%L4knhTbtHxCyPS?k+(_hxVS_*8bM#Z&xJm zA^8fGtw4{n%&Zu)qbi|gCxR)J4pe)SUb*@Aob*fCaYX*@{`VBeVM6YYynvXF&BpdR z^qgm*aaLiNfE%Jp)oUh#RxNyd_z`AEn_$UfQtfln$p4^3cK855-H$wrIU>bBKL$^* zR^XC&lSK}yK8VHoFNac)<;Z9ov)vJMi$=Ct9~|nL3u^}_C39q%|B@wNkd_kx3FamK zkb+C>`On?ejhvw3Pv=IpwVO3{x53@AS9&Vh(!|r^SNWL$=CRvI_&40_Th7`RH5_t~ zI7|DLNI5ArW1R=R53*4(@@>0r{e)*+=W34P3WIv#4JKN|WaGr?94NUPH|Gk}Ablah z$d@%jN`KaIQ;F*+lVkmTjba?q3NV3Xls!6$CZ}yy!}cVi?W4RN!5bAw3RSoV?VjxV z8R!JB8X%rad|y0<71{uoV0}iqzX{oF0hXEeAG&iwQ0>F=6d8+x5(G5l4duN?9ZHSGl0@ z#ifO>VvFxc7!32KH`$L|bWdl_4UIp2!Rr*lu%Zacb^}0BP-l!N=Xjd=V(XrZ(do-X z=ZCI3%4z-L9BP#B2Tf*s!tya%Yd0aMTx2Z|dJ{Lsy+%n)f@}NX#;Y@TQsb?==C-wx zC%C1d`qX@%^P9duCnr*O^nb#jEa08D#MI(`gxEp>5EAJ$-9Xkj*)bW zng9?5-nPJtiK}I{U^XJq4v?B?4l{2Pz7}`CVb_LQTwerO`F^t}{}LFMgpb8dqIMOb zgUF-(x4j-QE*~gG+emHz9<+3vCRuSB-};ohb!4;gfkj;AzyK53NCZ<*Wt^df-?qlv z+mRdV2^PuIlY?Jd&-pgMz;=)!=6&I^(p5!XGbMP7GJo(~J=CFZB>H4K>`hwhWx8bc z#_>QU`SjCOBtKYpmePo#0F|H|J>5GwsKB&Oi$YjY;~-@i#9pRk2OwpaZQN5Ql7i)(BoM~iw`{EC>AJ4 zo~avZ9wuptZ7hq@l$mcijMm+LBED17_V-m=2je^xfa@$iI!GyrsZmaqA0|?i^Mi2- zn(ak{k>U-mr-_Rxq5^8oS42qz>pHOQ>F-@k}2p4yp18Vh-O51Z0q>Oqox1cniGjr z@qkT2iijciR~XVSxXGX3+JW!SrTFL$G7t8Ziok%PzN`8*S$lTn-(?5a|Kj`W?e|{| zQPDU0{@J?;dd5QyVAkY1n558#ixnT33>NoTht6vBd~hVB^&J0xNoRrpuD-*v?)Kad zY)xDw@$rziN5QAnU#m!-y}!C7jlX9oglwac79cLEFoW;;LWUi2N%4#jb^xuE@AnD$ z?@sU9GSr0i?eq2WC)5g$gYIAU_dQ%utlDRyi^x_Sq%@UBq-QX*;tEU-AFv>WlTU~4 zcu1cMSYVILznNZN>?~<15s#nCRf7Q|$TG?ql(y5^5cE4oz+~q$u?i>ii?^0oW*rgy z>2tZ^{|~{HXEF>^U|aRWI^p~caPgdo{CTn@aJu}c_TW_ktt{LO0M76PMo3k2f|2*} z+o@}PnMVu;(Le*yW9HFcF~qZyxr2HTMOFaKj}Vkg8I{*S-c54WYEaUtPJxFu!f{vEH8ir(LpYm}4F?s6-c20%e7){g zeE|UeIPt{dCm&*RqsuCEqE?}(u6f{T+L+Fb^XvPSwSjCAUD%DZw9Oa_>zkl4b%0Ij z2+<0nx_?@}qq!Ag8-eBHH#aD!ZGjF78bp3+rc-#`RBi1^K9DUff(;XOg%*)YrX18S zzW{lgBkyX)QUbb5+W=K@-cp!Az-!Kn^=M0w5eb3Z$JK94;Q5R!9Xq;HuTJnVj4D*o zqMH#;fI%3+St5u!DDF88A&6BPdMc7B3c=1U!S;KWL+t-?C3d6Fu>yjfN!E%c%2m}% z6`kOGyn(2=LgBgoEj3~i9vs$2bDEn>Gv^^2EvxsIplXzz9q1B1_0W|ho}teF&Cw$a z*dlqzG>m(Qa;n6B_J{qOe^hNlC2i710G(bq-1=CRLkfoqM|x^=t{c1vIwe$p+IZ%8 zd}{&1{yNu`X>V1FmeYXCDWAo0T8NEjL6bX!X@80+3hk58wwW6NJjGbzUVCp-NP+6& z+Hvl=x7mY>NZaYfiU=r^z#=p693OsA z9(7NZF{iR(>M39)BsaJ@_vqOrNH5@=U%>pBCv&l07zWe&pH~Bl#8Q#pgAw;QIpk7! zbO|S%k=$wU36VN|L6Ron#Y5^@jGYc@ = { ...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([ "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) { +
+ +
+ +
+
+
@@ -9322,7 +9348,7 @@ function openExternalUrl(url: string) {
-