lastBookmarks = bookmarkManager.getBookmarks();
+ return Tuples.of(driver.session(ReactiveSession.class, Neo4jTransactionUtils.sessionConfig(false, lastBookmarks, targetDatabaseAndUser.getT1(), targetDatabaseAndUser.getT2())), lastBookmarks);
})))
- .map(t -> new DelegatingQueryRunner(t.getT1(), t.getT2(), (usedBookmarks, newBookmarks) -> {
- ReentrantReadWriteLock.WriteLock lock = bookmarksLock.writeLock();
- try {
- lock.lock();
- bookmarks.removeAll(usedBookmarks);
- bookmarks.addAll(newBookmarks);
- } finally {
- lock.unlock();
- }
- }));
+ .map(t -> new DelegatingQueryRunner(t.getT1(), t.getT2(), bookmarkManager::updateBookmarks));
}
private static class DelegatingQueryRunner implements ReactiveQueryRunner {
diff --git a/src/main/java/org/springframework/data/neo4j/core/Neo4jClient.java b/src/main/java/org/springframework/data/neo4j/core/Neo4jClient.java
index f3d330a28..45ee271ef 100644
--- a/src/main/java/org/springframework/data/neo4j/core/Neo4jClient.java
+++ b/src/main/java/org/springframework/data/neo4j/core/Neo4jClient.java
@@ -31,6 +31,7 @@ import org.neo4j.driver.summary.ResultSummary;
import org.neo4j.driver.types.TypeSystem;
import org.springframework.core.log.LogAccessor;
import org.springframework.data.neo4j.core.convert.Neo4jConversions;
+import org.springframework.data.neo4j.core.transaction.Neo4jBookmarkManager;
import org.springframework.lang.Nullable;
/**
@@ -79,6 +80,9 @@ public interface Neo4jClient {
@Nullable
Neo4jConversions neo4jConversions;
+ @Nullable
+ Neo4jBookmarkManager bookmarkManager;
+
private Builder(Driver driver) {
this.driver = driver;
}
@@ -121,6 +125,20 @@ public interface Neo4jClient {
return this;
}
+ /**
+ * Configures the {@link Neo4jBookmarkManager} to use.
+ * This should be the same instance as provided for the {@link org.springframework.data.neo4j.core.transaction.Neo4jTransactionManager}
+ * respectively the {@link org.springframework.data.neo4j.core.transaction.ReactiveNeo4jTransactionManager}.
+ *
+ * @param bookmarkManager Neo4jBookmarkManager instance that is shared with the transaction manager.
+ * @return The builder
+ * @since 7.1.2
+ */
+ public Builder withNeo4jBookmarkManager(Neo4jBookmarkManager bookmarkManager) {
+ this.bookmarkManager = bookmarkManager;
+ return this;
+ }
+
public Neo4jClient build() {
return new DefaultNeo4jClient(this);
}
diff --git a/src/main/java/org/springframework/data/neo4j/core/ReactiveNeo4jClient.java b/src/main/java/org/springframework/data/neo4j/core/ReactiveNeo4jClient.java
index 420635526..92dd6248e 100644
--- a/src/main/java/org/springframework/data/neo4j/core/ReactiveNeo4jClient.java
+++ b/src/main/java/org/springframework/data/neo4j/core/ReactiveNeo4jClient.java
@@ -15,6 +15,7 @@
*/
package org.springframework.data.neo4j.core;
+import org.springframework.data.neo4j.core.transaction.Neo4jBookmarkManager;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@@ -82,6 +83,9 @@ public interface ReactiveNeo4jClient {
@Nullable
Neo4jConversions neo4jConversions;
+ @Nullable
+ Neo4jBookmarkManager bookmarkManager;
+
private Builder(Driver driver) {
this.driver = driver;
}
@@ -124,6 +128,20 @@ public interface ReactiveNeo4jClient {
return this;
}
+ /**
+ * Configures the {@link Neo4jBookmarkManager} to use.
+ * This should be the same instance as provided for the {@link org.springframework.data.neo4j.core.transaction.Neo4jTransactionManager}
+ * respectively the {@link org.springframework.data.neo4j.core.transaction.ReactiveNeo4jTransactionManager}.
+ *
+ * @param bookmarkManager Neo4jBookmarkManager instance that is shared with the transaction manager.
+ * @return The builder
+ * @since 7.1.2
+ */
+ public Builder withNeo4jBookmarkManager(Neo4jBookmarkManager bookmarkManager) {
+ this.bookmarkManager = bookmarkManager;
+ return this;
+ }
+
public ReactiveNeo4jClient build() {
return new DefaultReactiveNeo4jClient(this);
}
diff --git a/src/main/java/org/springframework/data/neo4j/core/transaction/Neo4jBookmarkManager.java b/src/main/java/org/springframework/data/neo4j/core/transaction/Neo4jBookmarkManager.java
index e5b6c2e5f..23ca87d65 100644
--- a/src/main/java/org/springframework/data/neo4j/core/transaction/Neo4jBookmarkManager.java
+++ b/src/main/java/org/springframework/data/neo4j/core/transaction/Neo4jBookmarkManager.java
@@ -41,6 +41,13 @@ public sealed interface Neo4jBookmarkManager permits AbstractBookmarkManager, No
return new DefaultBookmarkManager(null);
}
+ /**
+ * @return default reactive version of bookmark manager
+ */
+ static Neo4jBookmarkManager createReactive() {
+ return new ReactiveDefaultBookmarkManager(null);
+ }
+
/**
* Use this factory method to add supplier of initial "seeding" bookmarks to the transaction managers
*
@@ -55,6 +62,20 @@ public sealed interface Neo4jBookmarkManager permits AbstractBookmarkManager, No
return new DefaultBookmarkManager(bookmarksSupplier);
}
+ /**
+ * Use this factory method to add supplier of initial "seeding" bookmarks to the transaction managers
+ *
+ * While this class will make sure that the supplier will be accessed in a thread-safe manner,
+ * it is the caller's duty to provide a thread safe supplier (not changing the seed during a call, etc.).
+ *
+ * @param bookmarksSupplier A supplier for seeding bookmarks, can be null. The supplier is free to provide different
+ * bookmarks on each call.
+ * @return A reactive bookmark manager
+ */
+ static Neo4jBookmarkManager createReactive(@Nullable Supplier> bookmarksSupplier) {
+ return new ReactiveDefaultBookmarkManager(bookmarksSupplier);
+ }
+
/**
* Use this bookmark manager at your own risk, it will effectively disable any bookmark management by dropping all
* bookmarks and never supplying any. In a cluster you will be at a high risk of experiencing stale reads. In a single
diff --git a/src/main/java/org/springframework/data/neo4j/core/transaction/ReactiveDefaultBookmarkManager.java b/src/main/java/org/springframework/data/neo4j/core/transaction/ReactiveDefaultBookmarkManager.java
new file mode 100644
index 000000000..d49ed050d
--- /dev/null
+++ b/src/main/java/org/springframework/data/neo4j/core/transaction/ReactiveDefaultBookmarkManager.java
@@ -0,0 +1,67 @@
+/*
+ * Copyright 2011-2023 the original author or authors.
+ *
+ * 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
+ *
+ * https://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.springframework.data.neo4j.core.transaction;
+
+import org.neo4j.driver.Bookmark;
+import org.springframework.context.ApplicationEventPublisher;
+import org.springframework.lang.Nullable;
+
+import java.util.Collection;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.Set;
+import java.util.function.Supplier;
+
+/**
+ * Default bookmark manager.
+ *
+ * @author Michael J. Simons
+ * @soundtrack Helge Schneider - The Last Jazz
+ * @since 7.0
+ */
+final class ReactiveDefaultBookmarkManager extends AbstractBookmarkManager {
+
+ private final Set bookmarks = Collections.synchronizedSet(new HashSet<>());
+
+ private final Supplier> bookmarksSupplier;
+
+ @Nullable
+ private ApplicationEventPublisher applicationEventPublisher;
+
+ ReactiveDefaultBookmarkManager(@Nullable Supplier> bookmarksSupplier) {
+ this.bookmarksSupplier = bookmarksSupplier == null ? Collections::emptySet : bookmarksSupplier;
+ }
+
+ @Override
+ public Collection getBookmarks() {
+ this.bookmarks.addAll(bookmarksSupplier.get());
+ return Collections.synchronizedSet(Collections.unmodifiableSet(this.bookmarks));
+ }
+
+ @Override
+ public void updateBookmarks(Collection usedBookmarks, Collection newBookmarks) {
+ bookmarks.removeAll(usedBookmarks);
+ bookmarks.addAll(newBookmarks);
+ if (applicationEventPublisher != null) {
+ applicationEventPublisher.publishEvent(new Neo4jBookmarksUpdatedEvent(new HashSet<>(bookmarks)));
+ }
+ }
+
+ @Override
+ public void setApplicationEventPublisher(@Nullable ApplicationEventPublisher applicationEventPublisher) {
+ this.applicationEventPublisher = applicationEventPublisher;
+ }
+}
diff --git a/src/main/java/org/springframework/data/neo4j/core/transaction/ReactiveNeo4jTransactionManager.java b/src/main/java/org/springframework/data/neo4j/core/transaction/ReactiveNeo4jTransactionManager.java
index e13ba5bf7..8e8fe6b64 100644
--- a/src/main/java/org/springframework/data/neo4j/core/transaction/ReactiveNeo4jTransactionManager.java
+++ b/src/main/java/org/springframework/data/neo4j/core/transaction/ReactiveNeo4jTransactionManager.java
@@ -179,7 +179,7 @@ public final class ReactiveNeo4jTransactionManager extends AbstractReactiveTrans
ReactiveUserSelectionProvider.getDefaultSelectionProvider() :
builder.userSelectionProvider;
this.bookmarkManager =
- builder.bookmarkManager == null ? Neo4jBookmarkManager.create() : builder.bookmarkManager;
+ builder.bookmarkManager == null ? Neo4jBookmarkManager.createReactive() : builder.bookmarkManager;
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/conversion_reactive/ReactiveCompositePropertiesIT.java b/src/test/java/org/springframework/data/neo4j/integration/conversion_reactive/ReactiveCompositePropertiesIT.java
index 4e7075729..47342db06 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/conversion_reactive/ReactiveCompositePropertiesIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/conversion_reactive/ReactiveCompositePropertiesIT.java
@@ -171,7 +171,7 @@ class ReactiveCompositePropertiesIT extends CompositePropertiesITBase {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/conversion_reactive/ReactiveCustomTypesIT.java b/src/test/java/org/springframework/data/neo4j/integration/conversion_reactive/ReactiveCustomTypesIT.java
index 13076c1d6..21130631a 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/conversion_reactive/ReactiveCustomTypesIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/conversion_reactive/ReactiveCustomTypesIT.java
@@ -221,7 +221,7 @@ public class ReactiveCustomTypesIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/imperative/Neo4jClientIT.java b/src/test/java/org/springframework/data/neo4j/integration/imperative/Neo4jClientIT.java
index 4d51b3658..791c78001 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/imperative/Neo4jClientIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/imperative/Neo4jClientIT.java
@@ -15,12 +15,6 @@
*/
package org.springframework.data.neo4j.integration.imperative;
-import static org.assertj.core.api.Assertions.assertThat;
-
-import java.util.Collection;
-import java.util.Collections;
-import java.util.List;
-
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.neo4j.cypherdsl.core.Cypher;
@@ -48,6 +42,12 @@ import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.transaction.annotation.EnableTransactionManagement;
import org.springframework.transaction.support.TransactionTemplate;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.List;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
/**
* @author Michael J. Simons
*/
diff --git a/src/test/java/org/springframework/data/neo4j/integration/issues/ReactiveIssuesIT.java b/src/test/java/org/springframework/data/neo4j/integration/issues/ReactiveIssuesIT.java
index ab5f84a20..690ad5469 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/issues/ReactiveIssuesIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/issues/ReactiveIssuesIT.java
@@ -447,7 +447,7 @@ class ReactiveIssuesIT extends TestBase {
BookmarkCapture bookmarkCapture = bookmarkCapture();
return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider,
- Neo4jBookmarkManager.create(bookmarkCapture));
+ Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/issues/gh2728/AbstractReactiveTestBase.java b/src/test/java/org/springframework/data/neo4j/integration/issues/gh2728/AbstractReactiveTestBase.java
index 1872b38d1..2c2f621f6 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/issues/gh2728/AbstractReactiveTestBase.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/issues/gh2728/AbstractReactiveTestBase.java
@@ -97,7 +97,7 @@ public abstract class AbstractReactiveTestBase {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
}
}
diff --git a/src/test/java/org/springframework/data/neo4j/integration/issues/pure_element_id/ReactiveElementIdIT.java b/src/test/java/org/springframework/data/neo4j/integration/issues/pure_element_id/ReactiveElementIdIT.java
index 6e8a129bc..13d31b438 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/issues/pure_element_id/ReactiveElementIdIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/issues/pure_element_id/ReactiveElementIdIT.java
@@ -358,7 +358,7 @@ public class ReactiveElementIdIT extends AbstractElementIdTestBase {
BookmarkCapture bookmarkCapture = bookmarkCapture();
return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider,
- Neo4jBookmarkManager.create(bookmarkCapture));
+ Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Bean
diff --git a/src/test/java/org/springframework/data/neo4j/integration/movies/reactive/ReactiveAdvancedMappingIT.java b/src/test/java/org/springframework/data/neo4j/integration/movies/reactive/ReactiveAdvancedMappingIT.java
index 7259fd74f..7e165705a 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/movies/reactive/ReactiveAdvancedMappingIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/movies/reactive/ReactiveAdvancedMappingIT.java
@@ -499,7 +499,7 @@ class ReactiveAdvancedMappingIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/properties/ReactivePropertyIT.java b/src/test/java/org/springframework/data/neo4j/integration/properties/ReactivePropertyIT.java
index 87155032e..8114fa139 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/properties/ReactivePropertyIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/properties/ReactivePropertyIT.java
@@ -304,7 +304,7 @@ class ReactivePropertyIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveAuditingIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveAuditingIT.java
index 1a1582246..b869ebb02 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveAuditingIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveAuditingIT.java
@@ -184,7 +184,7 @@ class ReactiveAuditingIT extends AuditingITBase {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveAuditingWithoutDatesIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveAuditingWithoutDatesIT.java
index 2a068f9c7..d78251245 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveAuditingWithoutDatesIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveAuditingWithoutDatesIT.java
@@ -119,7 +119,7 @@ class ReactiveAuditingWithoutDatesIT extends AuditingITBase {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCallbacksIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCallbacksIT.java
index 5c9cc4dbc..30faf5804 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCallbacksIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCallbacksIT.java
@@ -194,7 +194,7 @@ class ReactiveCallbacksIT extends CallbacksITBase {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCypherdslConditionExecutorIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCypherdslConditionExecutorIT.java
index 11237d135..4155e6745 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCypherdslConditionExecutorIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCypherdslConditionExecutorIT.java
@@ -181,7 +181,7 @@ class ReactiveCypherdslConditionExecutorIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCypherdslStatementExecutorIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCypherdslStatementExecutorIT.java
index 13233192d..a7d557cb1 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCypherdslStatementExecutorIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveCypherdslStatementExecutorIT.java
@@ -212,7 +212,7 @@ class ReactiveCypherdslStatementExecutorIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveDynamicLabelsIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveDynamicLabelsIT.java
index 292d8e4ae..3315dd57b 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveDynamicLabelsIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveDynamicLabelsIT.java
@@ -504,7 +504,7 @@ public class ReactiveDynamicLabelsIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Bean
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveDynamicRelationshipsIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveDynamicRelationshipsIT.java
index c60d76938..927e79ede 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveDynamicRelationshipsIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveDynamicRelationshipsIT.java
@@ -305,7 +305,7 @@ class ReactiveDynamicRelationshipsIT extends DynamicRelationshipsITBase
+ neo4jClient.query(cypher2)
+ .bindAll(Map.of("elementId", elementId))
+ .fetchAs(String.class).one()))
+ .expectNextCount(10000)
+ .verifyComplete();
+ }
+
+ @Test // GH-2755
+ public void testQueryExecutionPureDriver(@Autowired Driver driver, @Autowired BookmarkCapture bookmarkCapture) {
+
+ try (var session = driver.session(bookmarkCapture.createSessionConfig())) {
+ session.run("UNWIND range(1,10000) as count with count CREATE (u:VersionedExternalIdListBased) SET u.numberThing=count").consume();
+ bookmarkCapture.seedWith(session.lastBookmarks());
+ }
+
+ String cypher = "MATCH (n) RETURN elementId(n) as a";
+ String cypher2 = "MATCH (n) WHERE elementId(n) = $elementId RETURN elementId(n) as b";
+
+ StepVerifier.create(Flux.usingWhen(
+ Mono
+ .just(driver.session(ReactiveSession.class)),
+ session ->
+ Flux.from(session.run(cypher))
+ .flatMap(ReactiveResult::records)
+ .map(a -> a.get(0).asString())
+ .flatMap(elementId ->
+ Flux.usingWhen(
+ Mono.just(driver.session(ReactiveSession.class)),
+ innerSession ->
+ Flux.from(innerSession.run(cypher2, Map.of("elementId", elementId)))
+ .flatMap(ReactiveResult::records)
+ .map(result -> result.get(0).asString()),
+ innerSession -> Mono.fromDirect(innerSession.close())
+ )),
+ session -> Mono.fromDirect(session.close())))
+ .expectNextCount(10000)
+ .verifyComplete();
+ }
+
@Test // GH-2238
void clientShouldIntegrateWithCypherDSL(@Autowired TransactionalOperator transactionalOperator,
@Autowired ReactiveNeo4jClient client,
@@ -140,7 +201,7 @@ class ReactiveNeo4jClientIT {
BookmarkCapture bookmarkCapture = bookmarkCapture();
return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider,
- Neo4jBookmarkManager.create(bookmarkCapture));
+ Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Bean
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveNeo4jTemplateIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveNeo4jTemplateIT.java
index 60612e7ad..6af4260c7 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveNeo4jTemplateIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveNeo4jTemplateIT.java
@@ -1042,7 +1042,7 @@ class ReactiveNeo4jTemplateIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveNeo4jTransactionManagerTestIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveNeo4jTransactionManagerTestIT.java
index 0c526dc1b..41f09767c 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveNeo4jTransactionManagerTestIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveNeo4jTransactionManagerTestIT.java
@@ -106,7 +106,7 @@ class ReactiveNeo4jTransactionManagerTestIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveOptimisticLockingIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveOptimisticLockingIT.java
index 120d763cd..f8dc77974 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveOptimisticLockingIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveOptimisticLockingIT.java
@@ -380,7 +380,7 @@ class ReactiveOptimisticLockingIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveProjectionIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveProjectionIT.java
index 2ed3ea55e..df5a6f386 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveProjectionIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveProjectionIT.java
@@ -469,7 +469,7 @@ class ReactiveProjectionIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveQuerydslNeo4jPredicateExecutorIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveQuerydslNeo4jPredicateExecutorIT.java
index af1f461f1..31fc6e233 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveQuerydslNeo4jPredicateExecutorIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveQuerydslNeo4jPredicateExecutorIT.java
@@ -369,7 +369,7 @@ class ReactiveQuerydslNeo4jPredicateExecutorIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveRelationshipsIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveRelationshipsIT.java
index 3d88b5ac3..9913eaacf 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveRelationshipsIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveRelationshipsIT.java
@@ -244,7 +244,7 @@ class ReactiveRelationshipsIT extends RelationshipsITBase {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveRepositoryIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveRepositoryIT.java
index 8bf42da23..fa36c0b19 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveRepositoryIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveRepositoryIT.java
@@ -2882,7 +2882,7 @@ class ReactiveRepositoryIT {
return ReactiveNeo4jTransactionManager.with(driver)
.withDatabaseSelectionProvider(databaseSelectionProvider)
.withUserSelectionProvider(getUserSelectionProvider())
- .withBookmarkManager(Neo4jBookmarkManager.create(bookmarkCapture()))
+ .withBookmarkManager(Neo4jBookmarkManager.createReactive(bookmarkCapture()))
.build();
}
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveScrollingIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveScrollingIT.java
index 51b2e2cd6..68e13b535 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveScrollingIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveScrollingIT.java
@@ -292,7 +292,7 @@ class ReactiveScrollingIT {
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveStringlyTypeDynamicRelationshipsIT.java b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveStringlyTypeDynamicRelationshipsIT.java
index e7b2976c4..cd93e8e74 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveStringlyTypeDynamicRelationshipsIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/reactive/ReactiveStringlyTypeDynamicRelationshipsIT.java
@@ -301,7 +301,7 @@ class ReactiveStringlyTypeDynamicRelationshipsIT extends DynamicRelationshipsITB
public ReactiveTransactionManager reactiveTransactionManager(Driver driver, ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
- return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.create(bookmarkCapture));
+ return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider, Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override
diff --git a/src/test/java/org/springframework/data/neo4j/integration/versioned_self_references/ReactiveOptimisticLockingOfSelfReferencesIT.java b/src/test/java/org/springframework/data/neo4j/integration/versioned_self_references/ReactiveOptimisticLockingOfSelfReferencesIT.java
index 9844d3ebc..1e04a6326 100644
--- a/src/test/java/org/springframework/data/neo4j/integration/versioned_self_references/ReactiveOptimisticLockingOfSelfReferencesIT.java
+++ b/src/test/java/org/springframework/data/neo4j/integration/versioned_self_references/ReactiveOptimisticLockingOfSelfReferencesIT.java
@@ -316,7 +316,7 @@ class ReactiveOptimisticLockingOfSelfReferencesIT extends TestBase {
ReactiveDatabaseSelectionProvider databaseSelectionProvider) {
BookmarkCapture bookmarkCapture = bookmarkCapture();
return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider,
- Neo4jBookmarkManager.create(bookmarkCapture));
+ Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Bean