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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
<awaitility.version>4.2.2</awaitility.version>
<build-resources.version>2024-12.1</build-resources.version>
<byte-buddy.version>1.17.0</byte-buddy.version>
<caniuse-neo4j-detection.version>1.0.0</caniuse-neo4j-detection.version>
<cdc.version>1.0.12</cdc.version>
<commons-collections4.version>4.4</commons-collections4.version>
<commons-lang3.version>3.17.0</commons-lang3.version>
Expand Down Expand Up @@ -279,6 +280,16 @@
<artifactId>mockito-kotlin</artifactId>
<version>${mockito-kotlin.version}</version>
</dependency>
<dependency>
<groupId>org.neo4j</groupId>
<artifactId>caniuse-core</artifactId>
<version>${caniuse-neo4j-detection.version}</version>
</dependency>
<dependency>
<groupId>org.neo4j</groupId>
<artifactId>caniuse-neo4j-detection</artifactId>
<version>${caniuse-neo4j-detection.version}</version>
</dependency>
<dependency>
<groupId>org.neo4j</groupId>
<artifactId>neo4j-cypher-dsl</artifactId>
Expand Down
8 changes: 8 additions & 0 deletions sink/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,14 @@
<artifactId>antlr4-runtime</artifactId>
<version>${antlr4.version}</version>
</dependency>
<dependency>
<groupId>org.neo4j</groupId>
<artifactId>caniuse-core</artifactId>
</dependency>
<dependency>
<groupId>org.neo4j</groupId>
<artifactId>caniuse-neo4j-detection</artifactId>
</dependency>
<dependency>
<groupId>org.neo4j</groupId>
<artifactId>neo4j-cypher-dsl</artifactId>
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
/*
* Copyright (c) "Neo4j"
* Neo4j Sweden AB [https://neo4j.com]
*
* Licensed 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.neo4j.connectors.kafka.sink

import org.neo4j.caniuse.CanIUse.canIUse
import org.neo4j.caniuse.Cypher
import org.neo4j.caniuse.Neo4j
import org.neo4j.caniuse.Neo4jVersion
import org.neo4j.cypherdsl.core.Statement
import org.neo4j.cypherdsl.core.renderer.Configuration
import org.neo4j.cypherdsl.core.renderer.Dialect
import org.neo4j.cypherdsl.core.renderer.Renderer

internal class Cypher5Renderer(neo4j: Neo4j) : Renderer {
companion object {
private val Neo4jVersion5 = Neo4jVersion(major = 5, minor = 0, patch = 0)
}

private val isCypher5PrefixSupported = canIUse(Cypher.explicitCypher5Selection()).withNeo4j(neo4j)
private val delegateRenderer =
Renderer.getRenderer(
Configuration.newConfig()
.withDialect(
if (neo4j.version < Neo4jVersion5) {
Dialect.DEFAULT
} else {
Dialect.NEO4J_5
})
.build())

override fun render(statement: Statement?): String? {
val rendered = delegateRenderer.render(statement)

if (isCypher5PrefixSupported) {
return "CYPHER 5 $rendered"
}

return rendered
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,9 @@ import org.apache.kafka.common.config.Config
import org.apache.kafka.common.config.ConfigDef
import org.apache.kafka.common.config.ConfigException
import org.apache.kafka.connect.sink.SinkTask
import org.jetbrains.annotations.TestOnly
import org.neo4j.caniuse.Neo4j
import org.neo4j.caniuse.detectedWith
import org.neo4j.connectors.kafka.configuration.ConnectorType
import org.neo4j.connectors.kafka.configuration.Groups
import org.neo4j.connectors.kafka.configuration.Neo4jConfiguration
Expand All @@ -30,16 +33,14 @@ import org.neo4j.connectors.kafka.configuration.helpers.Recommenders
import org.neo4j.connectors.kafka.configuration.helpers.SIMPLE_DURATION_PATTERN
import org.neo4j.connectors.kafka.configuration.helpers.Validators
import org.neo4j.connectors.kafka.configuration.helpers.toSimpleString
import org.neo4j.cypherdsl.core.Cypher
import org.neo4j.cypherdsl.core.renderer.Configuration
import org.neo4j.cypherdsl.core.renderer.Dialect
import org.neo4j.cypherdsl.core.renderer.Renderer

class SinkConfiguration : Neo4jConfiguration {
private var fixedRenderer: Renderer? = null

constructor(original: Map<String, *>) : this(original, null)

@TestOnly
constructor(
originals: Map<String, *>,
renderer: Renderer?
Expand Down Expand Up @@ -79,29 +80,9 @@ class SinkConfiguration : Neo4jConfiguration {
val patternBindValueAs
get(): String = getString(PATTERN_BIND_VALUE_AS)

val dialect: Dialect by lazy {
driver.session(sessionConfig()).use {
val name = Cypher.name("name")
val versions = Cypher.name("versions")
val stmt =
Cypher.call("dbms.components")
.yield(name, versions)
.where(name.eq(Cypher.anonParameter("Neo4j Kernel")))
.returning(Cypher.valueAt(versions, 0))
.build()
val neo4j: Neo4j by lazy { Neo4j.detectedWith(driver) }

val version = it.run(stmt.cypher, stmt.parameters).single().get(0).asString()
if (version.startsWith("4.")) {
return@lazy Dialect.DEFAULT
}

return@lazy Dialect.NEO4J_5
}
}

val renderer: Renderer by lazy {
fixedRenderer ?: Renderer.getRenderer(Configuration.newConfig().withDialect(dialect).build())
}
val renderer: Renderer by lazy { fixedRenderer ?: Cypher5Renderer(neo4j) }

val topicNames: List<String>
get() =
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,151 @@
/*
* Copyright (c) "Neo4j"
* Neo4j Sweden AB [https://neo4j.com]
*
* Licensed 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.neo4j.connectors.kafka.sink

import io.kotest.matchers.shouldBe
import java.util.stream.Stream
import org.junit.jupiter.api.extension.ExtensionContext
import org.junit.jupiter.params.ParameterizedTest
import org.junit.jupiter.params.provider.Arguments
import org.junit.jupiter.params.provider.ArgumentsProvider
import org.junit.jupiter.params.provider.ArgumentsSource
import org.neo4j.caniuse.Neo4j
import org.neo4j.caniuse.Neo4jDeploymentType
import org.neo4j.caniuse.Neo4jEdition
import org.neo4j.caniuse.Neo4jVersion
import org.neo4j.cypherdsl.core.Cypher

class Cypher5RendererTest {

@ParameterizedTest
@ArgumentsSource(Cypher5PrefixSupportedVersions::class)
fun `should prefix statements with cypher 5 on compatible neo4j versions`(neo4j: Neo4j) {
val person = Cypher.node("Person").named("p")
val stmt = Cypher.match(person).returning(person).build()

Cypher5Renderer(neo4j).render(stmt) shouldBe "CYPHER 5 MATCH (p:`Person`) RETURN p"
}

@ParameterizedTest
@ArgumentsSource(UnsupportedCypher5PrefixVersions::class)
fun `should not prefix statements on earlier neo4j versions`(neo4j: Neo4j) {
val person = Cypher.node("Person").named("p")
val stmt = Cypher.match(person).returning(person).build()

Cypher5Renderer(neo4j).render(stmt) shouldBe "MATCH (p:`Person`) RETURN p"
}

class Cypher5PrefixSupportedVersions : ArgumentsProvider {
override fun provideArguments(context: ExtensionContext?): Stream<out Arguments?>? {
return Stream.of(
Arguments.of(
Neo4j(
Neo4jVersion(5, 26, 0),
Neo4jEdition.ENTERPRISE,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(Neo4jVersion(5, 26, 0), Neo4jEdition.ENTERPRISE, Neo4jDeploymentType.AURA)),
Arguments.of(
Neo4j(
Neo4jVersion(5, 26, 0),
Neo4jEdition.COMMUNITY,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(
Neo4jVersion(5, 26, 2),
Neo4jEdition.ENTERPRISE,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(Neo4jVersion(5, 26, 2), Neo4jEdition.ENTERPRISE, Neo4jDeploymentType.AURA)),
Arguments.of(
Neo4j(
Neo4jVersion(5, 26, 2),
Neo4jEdition.COMMUNITY,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(
Neo4jVersion(5, 27, 0),
Neo4jEdition.ENTERPRISE,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(Neo4jVersion(5, 27, 0), Neo4jEdition.ENTERPRISE, Neo4jDeploymentType.AURA)),
Arguments.of(
Neo4j(
Neo4jVersion(5, 27, 0),
Neo4jEdition.COMMUNITY,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(
Neo4jVersion(2025, 1, 0),
Neo4jEdition.ENTERPRISE,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(Neo4jVersion(2025, 1, 0), Neo4jEdition.ENTERPRISE, Neo4jDeploymentType.AURA)),
Arguments.of(
Neo4j(
Neo4jVersion(2025, 1, 0),
Neo4jEdition.COMMUNITY,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(
Neo4jVersion.LATEST, Neo4jEdition.ENTERPRISE, Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(Neo4jVersion.LATEST, Neo4jEdition.ENTERPRISE, Neo4jDeploymentType.AURA)),
)
}
}

class UnsupportedCypher5PrefixVersions : ArgumentsProvider {
override fun provideArguments(context: ExtensionContext?): Stream<out Arguments?>? {
return Stream.of(
Arguments.of(
Neo4j(
Neo4jVersion(4, 4, 41),
Neo4jEdition.ENTERPRISE,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(Neo4jVersion(4, 4, 0), Neo4jEdition.ENTERPRISE, Neo4jDeploymentType.AURA)),
Arguments.of(
Neo4j(
Neo4jVersion(4, 4, 41),
Neo4jEdition.COMMUNITY,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(
Neo4jVersion(5, 8, 0),
Neo4jEdition.ENTERPRISE,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(Neo4jVersion(5, 8, 0), Neo4jEdition.ENTERPRISE, Neo4jDeploymentType.AURA)),
Arguments.of(
Neo4j(
Neo4jVersion(5, 8, 0), Neo4jEdition.COMMUNITY, Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(
Neo4jVersion(5, 21, 0),
Neo4jEdition.ENTERPRISE,
Neo4jDeploymentType.SELF_MANAGED)),
Arguments.of(
Neo4j(Neo4jVersion(5, 21, 0), Neo4jEdition.ENTERPRISE, Neo4jDeploymentType.AURA)),
Arguments.of(
Neo4j(
Neo4jVersion(5, 21, 0),
Neo4jEdition.COMMUNITY,
Neo4jDeploymentType.SELF_MANAGED)))
}
}
}