This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new eb332e20fe [Feature][Connector-V2] Add Neo4j source connectivity
dry-run (#12433)
eb332e20fe is described below
commit eb332e20fe5ccc1adee4c8b15fd6d5b189d7394b
Author: Goutam Adwant <[email protected]>
AuthorDate: Fri Oct 9 08:09:53 2026 +0000
[Feature][Connector-V2] Add Neo4j source connectivity dry-run (#12433)
Signed-off-by: Goutam Adwant <[email protected]>
---
docs/en/connectors/source/Neo4j.md | 12 ++
docs/en/engines/zeta/user-command.md | 1 +
docs/zh/connectors/source/Neo4j.md | 12 ++
docs/zh/engines/zeta/user-command.md | 1 +
.../neo4j/source/Neo4jSourceDryRunValidator.java | 95 +++++++++++++
.../seatunnel/neo4j/source/Neo4jSourceFactory.java | 67 ++++++++-
.../seatunnel/neo4j/Neo4jFactoryTest.java | 50 +++++++
.../source/Neo4jSourceDryRunValidatorTest.java | 150 +++++++++++++++++++++
.../seatunnel/e2e/connector/neo4j/Neo4jIT.java | 148 ++++++++++++++++++++
9 files changed, 531 insertions(+), 5 deletions(-)
diff --git a/docs/en/connectors/source/Neo4j.md
b/docs/en/connectors/source/Neo4j.md
index 4b98715b7c..04ee0c9cc4 100644
--- a/docs/en/connectors/source/Neo4j.md
+++ b/docs/en/connectors/source/Neo4j.md
@@ -43,6 +43,18 @@ the returned fields to a SeaTunnel schema.
| Map | MAP |
| Null | NULL |
+## Connectivity dry-run
+
+Run `bin/seatunnel.sh --config <your-job.conf> --dry-run connect` with your
own job configuration.
+
+This source uses the existing driver connection/authentication configuration
and returns the same configured single-table or `tables_configs` schemas as a
normal source.
+
+- The temporary driver calls `verifyConnectivityAsync()`. It does not create a
session, execute user Cypher, read graph records, or write data.
+- Success validates driver connectivity and the authentication performed by
the driver handshake. It does **not** validate the configured database,
database/query permissions, Cypher syntax, result columns/types, or schema
compatibility with stored data. In particular, an invalid query or a
nonexistent `database` can still pass this connectivity-only check.
+- Driver connection and verification waits are capped at 15 seconds, retaining
a smaller positive `max_connection_timeout`; zero is bounded for preflight.
Driver close is initiated even on failure and waited for at most 5 seconds. DNS
resolution follows the JVM resolver.
+- Retry/timeout limits apply only to the temporary driver. Normal source
configuration, source execution and sink behavior are unchanged.
+- Validation errors omit raw URIs, driver error text and secret-bearing
exception causes. Sink connectivity remains unsupported.
+
## Source Options
| Name | Type | Required | Default | Description
|
diff --git a/docs/en/engines/zeta/user-command.md
b/docs/en/engines/zeta/user-command.md
index c4fcbffeff..c117864c05 100644
--- a/docs/en/engines/zeta/user-command.md
+++ b/docs/en/engines/zeta/user-command.md
@@ -84,6 +84,7 @@ The `--dry-run connect` option runs the static checks first,
then uses connector
| Kafka | Yes ([topic metadata + runtime output
schema](../../connectors/source/Kafka.md#connectivity-dry-run), not
consumer/group permissions) | No |
| FakeSource | Yes (schema inference only, no external system) | - |
| MongoDB | No | Yes ([connectivity + configured
authentication](../../connectors/sink/MongoDB.md#connectivity-dry-run); no
collection, schema, write permission or transaction checks) |
+| Neo4j | Yes ([driver connectivity + configured schema; not database/query
validation](../../connectors/source/Neo4j.md#connectivity-dry-run)) | No |
| RabbitMQ | Yes ([existing-queue metadata + configured schema; not consumer
permissions](../../connectors/source/Rabbitmq.md#connectivity-dry-run)) | No |
| S3File | Yes (metadata connectivity + inline schema, single-table
text/csv/json/xml; see [supported
scope](../../connectors/source/S3File.md#connectivity-dry-run)) | - |
| Redis | Yes ([connectivity +
authentication](../../connectors/source/Redis.md#connectivity-dry-run) +
configured schema; no key access) | Yes ([connectivity +
authentication](../../connectors/sink/Redis.md#connectivity-dry-run); no field
compatibility or write permission check) |
diff --git a/docs/zh/connectors/source/Neo4j.md
b/docs/zh/connectors/source/Neo4j.md
index 1ebcef8d83..161d106cf4 100644
--- a/docs/zh/connectors/source/Neo4j.md
+++ b/docs/zh/connectors/source/Neo4j.md
@@ -43,6 +43,18 @@ Neo4j 源连接器通过执行 Cypher 查询从 Neo4j 读取数据,并把查
| Map | MAP |
| Null | NULL |
+## 连接预检查
+
+使用自己的作业配置执行 `bin/seatunnel.sh --config <your-job.conf> --dry-run connect`。
+
+Source 复用现有驱动的连接和认证配置,并返回与正常 Source 相同的单表或 `tables_configs` 配置 schema。
+
+- 临时驱动调用 `verifyConnectivityAsync()`,不创建 session,不执行用户 Cypher,不读取图数据,也不写入数据。
+- 成功表示驱动连接及握手时执行的认证成功,**不代表**配置的数据库存在、拥有数据库或查询权限、Cypher 语法正确、结果字段/类型正确,或
schema 与存储数据匹配。无效查询或不存在的 `database` 仍可能通过此连接检查。
+- 驱动连接和验证等待各自最多 15 秒,保留更小的正数
`max_connection_timeout`;零值在预检查中也受到限制。失败后仍会发起驱动关闭,最多等待 5 秒。DNS 解析由 JVM 解析器控制。
+- 重试和超时限制仅作用于临时驱动。正常 Source 配置、执行及 Sink 行为保持不变。
+- 校验错误不包含原始 URI、驱动原始错误或含凭据的异常原因。Sink 连通性仍不支持。
+
## 源选项
| 名称 | 类型 | 是否必填 | 默认值 | 描述
|
diff --git a/docs/zh/engines/zeta/user-command.md
b/docs/zh/engines/zeta/user-command.md
index 6dd61b4bbc..e2b056500b 100644
--- a/docs/zh/engines/zeta/user-command.md
+++ b/docs/zh/engines/zeta/user-command.md
@@ -100,6 +100,7 @@ bin/seatunnel.sh --config
$SEATUNNEL_HOME/config/v2.batch.config.template --dry-
| Kafka | 支持([主题元数据 + 运行时输出
schema](../../connectors/source/Kafka.md#连通性-dry-run),不含消费或消费组权限) | 不支持 |
| FakeSource | 支持(仅 schema 推断,无外部系统) | - |
| MongoDB | 不支持 | 支持([连通性 +
配置的认证](../../connectors/sink/MongoDB.md#连通性-dry-run),不检查集合、schema、写入权限或事务) |
+| Neo4j | 支持([驱动连通性 + 配置
schema,不含数据库/查询校验](../../connectors/source/Neo4j.md#连接预检查)) | 不支持 |
| RabbitMQ | 支持([已有队列元数据 + 配置
schema,不含消费权限](../../connectors/source/Rabbitmq.md#连接预检查)) | 不支持 |
| S3File | 支持(元数据连通性 + 内联 schema,仅单表
text/csv/json/xml;参见[支持范围](../../connectors/source/S3File.md#连接预检查)) | - |
| Redis | 支持([连通性 + 认证](../../connectors/source/Redis.md#连通性-dry-run) + 配置的
schema,不访问 key) | 支持([连通性 +
认证](../../connectors/sink/Redis.md#连通性-dry-run),不校验字段兼容性和写入权限) |
diff --git
a/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceDryRunValidator.java
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceDryRunValidator.java
new file mode 100644
index 0000000000..296f1949ae
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceDryRunValidator.java
@@ -0,0 +1,95 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.neo4j.source;
+
+import org.apache.seatunnel.connectors.seatunnel.neo4j.config.DriverBuilder;
+
+import org.neo4j.driver.Driver;
+
+import java.io.IOException;
+import java.util.concurrent.TimeUnit;
+
+/** Verifies the driver's connection only; never creates a session or executes
configured Cypher. */
+final class Neo4jSourceDryRunValidator {
+ private static final long MAX_TIMEOUT_SECONDS = 15;
+ private static final long CLOSE_TIMEOUT_SECONDS = 5;
+
+ private Neo4jSourceDryRunValidator() {}
+
+ static void validate(DriverBuilder builder) throws Exception {
+ if (Thread.currentThread().isInterrupted()) {
+ throw new InterruptedException("Neo4j connect dry-run
interrupted");
+ }
+ Driver driver = null;
+ Exception failure = null;
+ boolean interrupted = false;
+ try {
+ Long configuredTimeout = builder.getMaxConnectionTimeoutSeconds();
+ if (configuredTimeout != null && configuredTimeout < 0) {
+ throw new IllegalArgumentException("Invalid connection
timeout");
+ }
+ long timeout =
+ configuredTimeout == null || configuredTimeout == 0
+ ? MAX_TIMEOUT_SECONDS
+ : Math.min(configuredTimeout, MAX_TIMEOUT_SECONDS);
+ // This builder belongs only to the preflight; normal source
settings are unchanged.
+ builder.setMaxConnectionTimeoutSeconds(timeout);
+ builder.setMaxTransactionRetryTimeSeconds(0L);
+ driver = builder.build();
+ // Also bound routing/handshake waits, not just socket connection
establishment.
+
driver.verifyConnectivityAsync().toCompletableFuture().get(timeout,
TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ interrupted = true;
+ failure = new InterruptedException("Neo4j connect dry-run
interrupted");
+ } catch (Exception e) {
+ failure =
+ new IOException(
+ "Neo4j connect dry-run connectivity check failed;
check URI, authentication and TLS settings");
+ } finally {
+ interrupted |= Thread.interrupted();
+ if (driver != null) {
+ try {
+ // Close even after a failed or timed-out handshake; do
not wait indefinitely.
+ driver.closeAsync()
+ .toCompletableFuture()
+ .get(CLOSE_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ interrupted = true;
+ if (failure == null) {
+ failure =
+ new InterruptedException(
+ "Neo4j connect dry-run cleanup
interrupted");
+ }
+ } catch (Exception e) {
+ if (failure == null) {
+ failure = new IOException("Neo4j connect dry-run
driver cleanup failed");
+ }
+ }
+ }
+ if (interrupted) {
+ Thread.currentThread().interrupt();
+ if (failure == null) {
+ failure = new InterruptedException("Neo4j connect dry-run
interrupted");
+ }
+ }
+ }
+ if (failure != null) {
+ throw failure;
+ }
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceFactory.java
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceFactory.java
index b847b0e9a7..ab89e9aa03 100644
---
a/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceFactory.java
+++
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceFactory.java
@@ -23,6 +23,7 @@ import org.apache.seatunnel.api.common.SeaTunnelAPIErrorCode;
import org.apache.seatunnel.api.configuration.ReadonlyConfig;
import org.apache.seatunnel.api.configuration.util.ConditionExtension;
import org.apache.seatunnel.api.configuration.util.Conditions;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
import org.apache.seatunnel.api.configuration.util.OptionRule;
import org.apache.seatunnel.api.configuration.util.OptionValidationException;
import org.apache.seatunnel.api.options.ConnectorCommonOptions;
@@ -32,6 +33,7 @@ import org.apache.seatunnel.api.table.catalog.CatalogTable;
import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
import org.apache.seatunnel.api.table.connector.TableSource;
import org.apache.seatunnel.api.table.factory.Factory;
+import org.apache.seatunnel.api.table.factory.SupportSourceDryRunValidation;
import org.apache.seatunnel.api.table.factory.TableSourceFactory;
import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
import
org.apache.seatunnel.connectors.seatunnel.neo4j.config.Neo4jAuthenticationConditions;
@@ -41,15 +43,45 @@ import
org.apache.seatunnel.connectors.seatunnel.neo4j.exception.Neo4jConnectorE
import com.google.auto.service.AutoService;
+import java.io.IOException;
import java.io.Serializable;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
@AutoService(Factory.class)
-public class Neo4jSourceFactory implements TableSourceFactory {
+public class Neo4jSourceFactory implements TableSourceFactory,
SupportSourceDryRunValidation {
+
+ /** Reuses runtime catalog construction without opening a reader or a
network connection. */
+ @Override
+ public List<CatalogTable> inferSchemaForDryRun(TableSourceFactoryContext
context)
+ throws IOException {
+ return parseDryRunSourceConfig(context).catalogTables;
+ }
+
+ private SourceConfiguration
parseDryRunSourceConfig(TableSourceFactoryContext context)
+ throws IOException {
+ try {
+ ConfigValidator.of(context.getOptions()).validate(optionRule());
+ return parseSourceConfig(context.getOptions());
+ } catch (RuntimeException e) {
+ throw new IOException(
+ "Neo4j connect dry-run source configuration or schema is
invalid");
+ }
+ }
+
+ /** Checks connectivity only, without starting the source's data-reading
lifecycle. */
+ @Override
+ public void validateConnectionForDryRun(
+ TableSourceFactoryContext context, List<CatalogTable>
catalogTables) throws Exception {
+ // The connection hook can be called directly, without the preceding
schema hook.
+ Neo4jSourceDryRunValidator.validate(
+
parseDryRunSourceConfig(context).connectionInfo.getDriverBuilder());
+ }
+
@Override
public String factoryIdentifier() {
return Neo4jSourceOptions.PLUGIN_NAME;
@@ -102,10 +134,20 @@ public class Neo4jSourceFactory implements
TableSourceFactory {
}
private Neo4jSource createNeo4jSource(ReadonlyConfig config) {
+ SourceConfiguration parsed = parseSourceConfig(config);
+ return parsed.tableConfigs == null
+ ? new Neo4jSource(parsed.catalogTables.get(0),
parsed.connectionInfo)
+ : new Neo4jSource(parsed.catalogTables, parsed.connectionInfo,
parsed.tableConfigs);
+ }
+
+ // Shared pure configuration parsing keeps preflight schemas identical to
the runtime without
+ // creating a source, reader or driver.
+ private SourceConfiguration parseSourceConfig(ReadonlyConfig config) {
if
(!config.getOptional(ConnectorCommonOptions.TABLE_CONFIGS).isPresent()) {
- return new Neo4jSource(
- CatalogTableUtil.buildWithConfig(config),
- new Neo4jSourceQueryInfo(config.toConfig()));
+ return new SourceConfiguration(
+
Collections.singletonList(CatalogTableUtil.buildWithConfig(config)),
+ new Neo4jSourceQueryInfo(config.toConfig()),
+ null);
}
List<Map<String, Object>> entries =
config.get(ConnectorCommonOptions.TABLE_CONFIGS);
@@ -155,7 +197,22 @@ public class Neo4jSourceFactory implements
TableSourceFactory {
Neo4jSourceOptions.KEY_QUERY.key(),
ConfigValueFactory.fromAnyRef(
tableConfigs.get(0).getQuery())));
- return new Neo4jSource(catalogTables, connectionInfo, tableConfigs);
+ return new SourceConfiguration(catalogTables, connectionInfo,
tableConfigs);
+ }
+
+ private static final class SourceConfiguration {
+ private final List<CatalogTable> catalogTables;
+ private final Neo4jSourceQueryInfo connectionInfo;
+ private final List<Neo4jSourceTableConfig> tableConfigs;
+
+ private SourceConfiguration(
+ List<CatalogTable> catalogTables,
+ Neo4jSourceQueryInfo connectionInfo,
+ List<Neo4jSourceTableConfig> tableConfigs) {
+ this.catalogTables = catalogTables;
+ this.connectionInfo = connectionInfo;
+ this.tableConfigs = tableConfigs;
+ }
}
private static Neo4jConnectorException configError(String message) {
diff --git
a/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
b/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
index 0550d9bfca..682ae81d9a 100644
---
a/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
+++
b/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
@@ -21,6 +21,7 @@ import org.apache.seatunnel.api.configuration.ReadonlyConfig;
import org.apache.seatunnel.api.configuration.util.ConfigValidator;
import org.apache.seatunnel.api.configuration.util.OptionValidationException;
import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.factory.SupportSourceDryRunValidation;
import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
import org.apache.seatunnel.connectors.seatunnel.neo4j.config.Neo4jSinkOptions;
import
org.apache.seatunnel.connectors.seatunnel.neo4j.config.Neo4jSourceOptions;
@@ -40,6 +41,11 @@ import java.util.Map;
class Neo4jFactoryTest {
+ @Test
+ void supportsConnectivityDryRun() {
+ Assertions.assertTrue(new Neo4jSourceFactory() instanceof
SupportSourceDryRunValidation);
+ }
+
@Test
void optionRule() {
Assertions.assertNotNull((new Neo4jSourceFactory()).optionRule());
@@ -282,6 +288,50 @@ class Neo4jFactoryTest {
Assertions.assertEquals("companies",
tables.get(1).getTableId().toTablePath().toString());
}
+ @Test
+ void dryRunSchemaMatchesRuntimeForSingleAndMultipleTables() throws
Exception {
+ Map<String, Object> multi = sourceConnectionConfig();
+ multi.put(
+ "tables_configs",
+ Arrays.asList(
+ tableConfig("people", "CREATE (:MustNotExecute)"),
+ tableConfig("companies", "THIS IS NOT VALID CYPHER")));
+ Neo4jSourceFactory factory = new Neo4jSourceFactory();
+ for (Map<String, Object> options : Arrays.asList(validSourceConfig(),
multi)) {
+ TableSourceFactoryContext context =
+ new TableSourceFactoryContext(
+ ReadonlyConfig.fromMap(options),
getClass().getClassLoader());
+ List<CatalogTable> expected =
+ ((Neo4jSource) (Object)
factory.createSource(context).createSource())
+ .getProducedCatalogTables();
+ List<CatalogTable> actual = factory.inferSchemaForDryRun(context);
+ Assertions.assertEquals(expected.size(), actual.size());
+ for (int i = 0; i < expected.size(); i++) {
+ Assertions.assertEquals(expected.get(i).getTableId(),
actual.get(i).getTableId());
+ Assertions.assertEquals(
+ expected.get(i).getSeaTunnelRowType(),
actual.get(i).getSeaTunnelRowType());
+ Assertions.assertEquals(expected.get(i).getOptions(),
actual.get(i).getOptions());
+ }
+ }
+ }
+
+ @Test
+ void dryRunInvalidUriDoesNotEchoCredentials() {
+ Map<String, Object> options = validSourceConfig();
+ options.put("uri", "neo4j://private-secret@[invalid");
+ Exception failure =
+ Assertions.assertThrows(
+ java.io.IOException.class,
+ () ->
+ new Neo4jSourceFactory()
+ .inferSchemaForDryRun(
+ new TableSourceFactoryContext(
+
ReadonlyConfig.fromMap(options),
+
getClass().getClassLoader())));
+
Assertions.assertFalse(failure.getMessage().contains("private-secret"));
+ Assertions.assertNull(failure.getCause());
+ }
+
private static Map<String, Object> sourceConnectionConfig() {
Map<String, Object> config = new HashMap<>();
config.put("uri", "neo4j://localhost:7687");
diff --git
a/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceDryRunValidatorTest.java
b/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceDryRunValidatorTest.java
new file mode 100644
index 0000000000..775421a872
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceDryRunValidatorTest.java
@@ -0,0 +1,150 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.neo4j.source;
+
+import org.apache.seatunnel.connectors.seatunnel.neo4j.config.DriverBuilder;
+
+import org.junit.jupiter.api.Test;
+import org.neo4j.driver.Driver;
+
+import java.io.IOException;
+import java.util.concurrent.CompletableFuture;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
+import static org.mockito.Mockito.when;
+
+class Neo4jSourceDryRunValidatorTest {
+ @Test
+ void verifiesConnectionAndClosesWithoutCreatingSession() throws Exception {
+ DriverBuilder builder = mock(DriverBuilder.class);
+ Driver driver = mock(Driver.class);
+ when(builder.build()).thenReturn(driver);
+
when(driver.verifyConnectivityAsync()).thenReturn(CompletableFuture.completedFuture(null));
+
when(driver.closeAsync()).thenReturn(CompletableFuture.completedFuture(null));
+
+ Neo4jSourceDryRunValidator.validate(builder);
+
+ verify(driver).verifyConnectivityAsync();
+ verify(driver).closeAsync();
+ verifyNoMoreInteractions(driver);
+ verify(builder).setMaxConnectionTimeoutSeconds(15L);
+ verify(builder).setMaxTransactionRetryTimeSeconds(0L);
+ }
+
+ @Test
+ void closesAfterAuthenticationFailureAndDoesNotExposeDriverText() {
+ DriverBuilder builder = mock(DriverBuilder.class);
+ Driver driver = mock(Driver.class);
+ CompletableFuture<Void> failed = new CompletableFuture<>();
+ failed.completeExceptionally(new
IllegalArgumentException("neo4j://user:secret@host"));
+ when(builder.build()).thenReturn(driver);
+ when(driver.verifyConnectivityAsync()).thenReturn(failed);
+
when(driver.closeAsync()).thenReturn(CompletableFuture.completedFuture(null));
+
+ IOException error =
+ assertThrows(IOException.class, () ->
Neo4jSourceDryRunValidator.validate(builder));
+
+ assertFalse(error.getMessage().contains("secret"));
+ assertNull(error.getCause());
+ assertEquals(0, error.getSuppressed().length);
+ verify(driver).closeAsync();
+ }
+
+ @Test
+ void honorsSmallerTimeoutAndClosesStalledHandshake() {
+ DriverBuilder builder = mock(DriverBuilder.class);
+ Driver driver = mock(Driver.class);
+ when(builder.getMaxConnectionTimeoutSeconds()).thenReturn(1L);
+ when(builder.build()).thenReturn(driver);
+ when(driver.verifyConnectivityAsync()).thenReturn(new
CompletableFuture<>());
+
when(driver.closeAsync()).thenReturn(CompletableFuture.completedFuture(null));
+
+ assertThrows(IOException.class, () ->
Neo4jSourceDryRunValidator.validate(builder));
+ verify(builder).setMaxConnectionTimeoutSeconds(1L);
+ verify(driver).closeAsync();
+ }
+
+ @Test
+ void rejectsNegativeTimeoutBeforeDriverCreation() {
+ DriverBuilder builder = mock(DriverBuilder.class);
+ when(builder.getMaxConnectionTimeoutSeconds()).thenReturn(-1L);
+ assertThrows(IOException.class, () ->
Neo4jSourceDryRunValidator.validate(builder));
+ verify(builder).getMaxConnectionTimeoutSeconds();
+ verifyNoMoreInteractions(builder);
+ }
+
+ @Test
+ void rejectsPreInterruptedThreadWithoutCreatingDriver() {
+ DriverBuilder builder = mock(DriverBuilder.class);
+ Thread.currentThread().interrupt();
+ try {
+ assertThrows(
+ InterruptedException.class, () ->
Neo4jSourceDryRunValidator.validate(builder));
+ assertTrue(Thread.currentThread().isInterrupted());
+ verifyNoInteractions(builder);
+ } finally {
+ Thread.interrupted();
+ }
+ }
+
+ @Test
+ void preservesInterruptionDuringVerificationAndStillCloses() {
+ DriverBuilder builder = mock(DriverBuilder.class);
+ Driver driver = mock(Driver.class);
+ when(builder.build()).thenReturn(driver);
+ when(driver.verifyConnectivityAsync())
+ .thenAnswer(
+ ignored -> {
+ Thread.currentThread().interrupt();
+ return new CompletableFuture<>();
+ });
+
when(driver.closeAsync()).thenReturn(CompletableFuture.completedFuture(null));
+ try {
+ assertThrows(
+ InterruptedException.class, () ->
Neo4jSourceDryRunValidator.validate(builder));
+ assertTrue(Thread.currentThread().isInterrupted());
+ verify(driver).closeAsync();
+ } finally {
+ Thread.interrupted();
+ }
+ }
+
+ @Test
+ void cleanupFailureCannotTurnIntoSuccessOrExposeSecrets() {
+ DriverBuilder builder = mock(DriverBuilder.class);
+ Driver driver = mock(Driver.class);
+ when(builder.build()).thenReturn(driver);
+
when(driver.verifyConnectivityAsync()).thenReturn(CompletableFuture.completedFuture(null));
+ CompletableFuture<Void> failedClose = new CompletableFuture<>();
+ failedClose.completeExceptionally(new
IllegalArgumentException("private-token"));
+ when(driver.closeAsync()).thenReturn(failedClose);
+ IOException failure =
+ assertThrows(IOException.class, () ->
Neo4jSourceDryRunValidator.validate(builder));
+ assertTrue(failure.getMessage().contains("cleanup"));
+ assertFalse(failure.getMessage().contains("private-token"));
+ assertNull(failure.getCause());
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/java/org/apache/seatunnel/e2e/connector/neo4j/Neo4jIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/java/org/apache/seatunnel/e2e/connector/neo4j/Neo4jIT.java
index ad21d44672..40f7b3067a 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/java/org/apache/seatunnel/e2e/connector/neo4j/Neo4jIT.java
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/java/org/apache/seatunnel/e2e/connector/neo4j/Neo4jIT.java
@@ -17,6 +17,17 @@
package org.apache.seatunnel.e2e.connector.neo4j;
+import org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory;
+import org.apache.seatunnel.shade.com.typesafe.config.ConfigRenderOptions;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.table.factory.FactoryUtil;
+import org.apache.seatunnel.api.table.factory.SupportSourceDryRunValidation;
+import org.apache.seatunnel.api.table.factory.TableSourceFactory;
+import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import org.apache.seatunnel.core.starter.seatunnel.args.ClientCommandArgs;
+import
org.apache.seatunnel.core.starter.seatunnel.command.SeaTunnelConfValidateCommand;
+import org.apache.seatunnel.core.starter.utils.CommandLineUtils;
import org.apache.seatunnel.e2e.common.TestResource;
import org.apache.seatunnel.e2e.common.TestSuiteBase;
import org.apache.seatunnel.e2e.common.container.TestContainer;
@@ -24,7 +35,9 @@ import
org.apache.seatunnel.e2e.common.container.TestContainer;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.TestTemplate;
+import org.junit.jupiter.api.io.TempDir;
import org.neo4j.driver.AuthTokens;
import org.neo4j.driver.Driver;
import org.neo4j.driver.GraphDatabase;
@@ -46,9 +59,16 @@ import lombok.extern.slf4j.Slf4j;
import java.io.IOException;
import java.net.URI;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
import java.util.concurrent.TimeUnit;
import java.util.stream.Stream;
@@ -197,6 +217,134 @@ public class Neo4jIT extends TestSuiteBase implements
TestResource {
Assertions.assertEquals(0, execResult.getExitCode());
}
+ @Test
+ public void testDryRunDoesNotExecuteConfiguredCypher() throws Exception {
+ neo4jSession.run("MATCH (n:DryRunSentinel) DELETE n").consume();
+ neo4jSession.run("CREATE (:DryRunSentinel
{name:'unchanged'})").consume();
+ Map<String, Object> options = dryRunOptions();
+ options.put("query", "MATCH (n:DryRunSentinel) DELETE n");
+ validateDryRun(options);
+ Assertions.assertEquals(
+ "unchanged",
+ neo4jSession
+ .run("MATCH (n:DryRunSentinel) RETURN n.name AS name")
+ .single()
+ .get("name")
+ .asString());
+ options.put("query", "THIS IS DELIBERATELY INVALID CYPHER");
+ validateDryRun(options);
+ }
+
+ @Test
+ public void testDryRunMultiTableSchemaWithoutQueryExecution() throws
Exception {
+ Map<String, Object> options = dryRunOptions();
+ options.remove("query");
+ options.remove("schema");
+ Map<String, Object> first = new HashMap<>();
+ first.put("query", "CREATE (:DryRunMustNotExist)");
+ Map<String, Object> schema = new HashMap<>();
+ schema.put("table", "first");
+ schema.put("fields", Collections.singletonMap("name", "string"));
+ first.put("schema", schema);
+ Map<String, Object> second = new HashMap<>();
+ second.put("query", "INVALID CYPHER");
+ Map<String, Object> secondSchema = new HashMap<>(schema);
+ secondSchema.put("table", "second");
+ second.put("schema", secondSchema);
+ options.put("tables_configs", Arrays.asList(first, second));
+ validateDryRun(options);
+ Assertions.assertEquals(
+ 0,
+ neo4jSession
+ .run("MATCH (n:DryRunMustNotExist) RETURN count(n) AS
count")
+ .single()
+ .get("count")
+ .asInt());
+ }
+
+ @Test
+ public void testDryRunRejectsWrongPasswordWithoutEchoingIt() {
+ Map<String, Object> options = dryRunOptions();
+ options.put("password", "dry-run-private-password");
+ Exception error = Assertions.assertThrows(IOException.class, () ->
validateDryRun(options));
+
Assertions.assertFalse(error.toString().contains("dry-run-private-password"));
+ Assertions.assertNull(error.getCause());
+ }
+
+ @Test
+ public void testDryRunDoesNotClaimDatabaseValidation() throws Exception {
+ Map<String, Object> options = dryRunOptions();
+ options.put("database", "missing-preflight-database");
+ // Driver connectivity does not select a database or prove query
permissions.
+ validateDryRun(options);
+ }
+
+ private Map<String, Object> dryRunOptions() {
+ Map<String, Object> options = new HashMap<>();
+ options.put(
+ "uri",
+ String.format(
+ "bolt://%s:%s", container.getHost(),
container.getMappedPort(BOLT_PORT)));
+ options.put("database", "neo4j");
+ options.put("username", CONTAINER_NEO4J_USERNAME);
+ options.put("password", CONTAINER_NEO4J_PASSWORD);
+ options.put("query", "RETURN 'unused' AS name");
+ options.put(
+ "schema",
+ Collections.singletonMap("fields",
Collections.singletonMap("name", "string")));
+ return options;
+ }
+
+ @Test
+ public void testDryRunCommand(@TempDir Path directory) throws Exception {
+ Map<String, Object> options = dryRunOptions();
+ options.put("query", "CREATE (:DryRunCommandMustNotExist)");
+ checkDryRunCommand(options, directory);
+ Assertions.assertEquals(
+ 0,
+ neo4jSession
+ .run("MATCH (n:DryRunCommandMustNotExist) RETURN
count(n) AS count")
+ .single()
+ .get("count")
+ .asInt());
+ }
+
+ private void checkDryRunCommand(Map<String, Object> options, Path
directory) throws Exception {
+ options.put("plugin_name", "Neo4j");
+ options.put("plugin_output", "preflight");
+ Map<String, Object> sink = new HashMap<>();
+ sink.put("plugin_name", "Console");
+ sink.put("plugin_input", "preflight");
+ Map<String, Object> job = new HashMap<>();
+ job.put("source", Collections.singletonList(options));
+ job.put("sink", Collections.singletonList(sink));
+ Path file = directory.resolve("dry-run.json");
+ Files.write(
+ file,
+ ConfigFactory.parseMap(job)
+ .root()
+ .render(ConfigRenderOptions.concise())
+ .getBytes(StandardCharsets.UTF_8));
+ ClientCommandArgs args =
+ CommandLineUtils.parse(
+ new String[] {"-c", file.toString(), "--dry-run",
"connect"},
+ new ClientCommandArgs(),
+ "seatunnel.sh",
+ true);
+ new SeaTunnelConfValidateCommand(args).execute();
+ }
+
+ private void validateDryRun(Map<String, Object> options) throws Exception {
+ ClassLoader loader = getClass().getClassLoader();
+ TableSourceFactory factory =
+ FactoryUtil.discoverFactory(loader, TableSourceFactory.class,
"Neo4j");
+ Assertions.assertTrue(factory instanceof
SupportSourceDryRunValidation);
+ SupportSourceDryRunValidation validator =
(SupportSourceDryRunValidation) factory;
+ TableSourceFactoryContext context =
+ new TableSourceFactoryContext(ReadonlyConfig.fromMap(options),
loader);
+ validator.validateConnectionForDryRun(context,
validator.inferSchemaForDryRun(context));
+ }
+
@AfterAll
@Override
public void tearDown() {