GH-2755 - Use bookmark manager in Client.

Also make the `AbstractNeo4jConfig` and its reactive counter-part
use the same defined bookmark manager in transaction manager and
client.

Closes #2755

(cherry picked from commit 21c8a42b29)
This commit is contained in:
Gerrit Meier
2023-07-03 11:49:55 +02:00
parent aecf61ecfb
commit 4ea4905d17
40 changed files with 271 additions and 90 deletions

11
pom.xml
View File

@@ -72,6 +72,7 @@
<archunit.version>0.23.1</archunit.version>
<asciidoctor-maven-plugin.version>2.1.0</asciidoctor-maven-plugin.version>
<asciidoctorj-diagram.version>2.1.0</asciidoctorj-diagram.version>
<blockhound.version>1.0.8.RELEASE</blockhound.version>
<byte-buddy.version>1.14.3</byte-buddy.version>
<cdi>3.0.1</cdi>
<checkstyle.skip>${skipTests}</checkstyle.skip>
@@ -237,6 +238,11 @@
<type>pom</type>
<scope>import</scope>
</dependency>
<dependency>
<groupId>io.projectreactor.tools</groupId>
<artifactId>blockhound</artifactId>
<version>${blockhound.version}</version>
</dependency>
</dependencies>
</dependencyManagement>
@@ -449,6 +455,11 @@
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>io.projectreactor.tools</groupId>
<artifactId>blockhound</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<repositories>

View File

@@ -27,6 +27,7 @@ import org.springframework.data.neo4j.core.Neo4jOperations;
import org.springframework.data.neo4j.core.Neo4jTemplate;
import org.springframework.data.neo4j.core.UserSelectionProvider;
import org.springframework.data.neo4j.core.mapping.Neo4jMappingContext;
import org.springframework.data.neo4j.core.transaction.Neo4jBookmarkManager;
import org.springframework.data.neo4j.core.transaction.Neo4jTransactionManager;
import org.springframework.data.neo4j.repository.config.Neo4jRepositoryConfigurationExtension;
import org.springframework.lang.Nullable;
@@ -47,6 +48,9 @@ public abstract class AbstractNeo4jConfig extends Neo4jConfigurationSupport {
@Autowired
private ObjectProvider<UserSelectionProvider> userSelectionProviders;
@Autowired
private Neo4jBookmarkManager bookmarkManager;
/**
* The driver to be used for interacting with Neo4j.
*
@@ -66,6 +70,7 @@ public abstract class AbstractNeo4jConfig extends Neo4jConfigurationSupport {
return Neo4jClient.with(driver)
.withDatabaseSelectionProvider(databaseSelectionProvider)
.withUserSelectionProvider(getUserSelectionProvider())
.withNeo4jBookmarkManager(bookmarkManager)
.build();
}
@@ -94,9 +99,15 @@ public abstract class AbstractNeo4jConfig extends Neo4jConfigurationSupport {
.with(driver)
.withDatabaseSelectionProvider(databaseSelectionProvider)
.withUserSelectionProvider(getUserSelectionProvider())
.withBookmarkManager(bookmarkManager)
.build();
}
@Bean
public Neo4jBookmarkManager bookmarkManager() {
return Neo4jBookmarkManager.create();
}
/**
* Configures the database selection provider.
*

View File

@@ -26,6 +26,7 @@ import org.springframework.data.neo4j.core.ReactiveNeo4jClient;
import org.springframework.data.neo4j.core.ReactiveNeo4jTemplate;
import org.springframework.data.neo4j.core.ReactiveUserSelectionProvider;
import org.springframework.data.neo4j.core.mapping.Neo4jMappingContext;
import org.springframework.data.neo4j.core.transaction.Neo4jBookmarkManager;
import org.springframework.data.neo4j.core.transaction.ReactiveNeo4jTransactionManager;
import org.springframework.data.neo4j.repository.config.ReactiveNeo4jRepositoryConfigurationExtension;
import org.springframework.lang.Nullable;
@@ -47,6 +48,9 @@ public abstract class AbstractReactiveNeo4jConfig extends Neo4jConfigurationSupp
@Autowired
private ObjectProvider<ReactiveUserSelectionProvider> userSelectionProviders;
@Autowired
private Neo4jBookmarkManager bookmarkManager;
/**
* The driver to be used for interacting with Neo4j.
*
@@ -66,6 +70,7 @@ public abstract class AbstractReactiveNeo4jConfig extends Neo4jConfigurationSupp
return ReactiveNeo4jClient.with(driver)
.withDatabaseSelectionProvider(databaseSelectionProvider)
.withUserSelectionProvider(getUserSelectionProvider())
.withNeo4jBookmarkManager(bookmarkManager)
.build();
}
@@ -94,9 +99,15 @@ public abstract class AbstractReactiveNeo4jConfig extends Neo4jConfigurationSupp
return ReactiveNeo4jTransactionManager.with(driver)
.withDatabaseSelectionProvider(databaseSelectionProvider)
.withUserSelectionProvider(getUserSelectionProvider())
.withBookmarkManager(bookmarkManager)
.build();
}
@Bean
public Neo4jBookmarkManager bookmarkManager() {
return Neo4jBookmarkManager.createReactive();
}
/**
* Configures the database name provider.
*

View File

@@ -16,13 +16,9 @@
package org.springframework.data.neo4j.core;
import java.util.Collection;
import java.util.Collections;
import java.util.HashSet;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.function.BiConsumer;
import java.util.function.BiFunction;
import java.util.function.Function;
@@ -45,6 +41,7 @@ import org.springframework.core.convert.support.DefaultConversionService;
import org.springframework.dao.DataAccessException;
import org.springframework.dao.support.PersistenceExceptionTranslator;
import org.springframework.data.neo4j.core.convert.Neo4jConversions;
import org.springframework.data.neo4j.core.transaction.Neo4jBookmarkManager;
import org.springframework.data.neo4j.core.transaction.Neo4jTransactionManager;
import org.springframework.data.neo4j.core.transaction.Neo4jTransactionUtils;
import org.springframework.lang.Nullable;
@@ -67,15 +64,15 @@ final class DefaultNeo4jClient implements Neo4jClient {
private final ConversionService conversionService;
private final Neo4jPersistenceExceptionTranslator persistenceExceptionTranslator = new Neo4jPersistenceExceptionTranslator();
// Basically a local bookmark manager
private final Set<Bookmark> bookmarks = new HashSet<>();
private final ReentrantReadWriteLock bookmarksLock = new ReentrantReadWriteLock();
// Local bookmark manager when using outside managed transactions
private final Neo4jBookmarkManager bookmarkManager;
DefaultNeo4jClient(Builder builder) {
this.driver = builder.driver;
this.databaseSelectionProvider = builder.databaseSelectionProvider;
this.userSelectionProvider = builder.userSelectionProvider;
this.bookmarkManager = builder.bookmarkManager != null ? builder.bookmarkManager : Neo4jBookmarkManager.create();
this.conversionService = new DefaultConversionService();
Optional.ofNullable(builder.neo4jConversions).orElseGet(Neo4jConversions::new).registerConvertersIn((ConverterRegistry) conversionService);
@@ -85,29 +82,13 @@ final class DefaultNeo4jClient implements Neo4jClient {
public QueryRunner getQueryRunner(DatabaseSelection databaseSelection, UserSelection impersonatedUser) {
QueryRunner queryRunner = Neo4jTransactionManager.retrieveTransaction(driver, databaseSelection, impersonatedUser);
Collection<Bookmark> lastBookmarks = Collections.emptySet();
Collection<Bookmark> lastBookmarks = bookmarkManager.getBookmarks();
if (queryRunner == null) {
ReentrantReadWriteLock.ReadLock lock = bookmarksLock.readLock();
try {
lock.lock();
lastBookmarks = new HashSet<>(bookmarks);
queryRunner = driver.session(Neo4jTransactionUtils.sessionConfig(false, lastBookmarks, databaseSelection, impersonatedUser));
} finally {
lock.unlock();
}
queryRunner = driver.session(Neo4jTransactionUtils.sessionConfig(false, lastBookmarks, databaseSelection, impersonatedUser));
}
return new DelegatingQueryRunner(queryRunner, lastBookmarks, (usedBookmarks, newBookmarks) -> {
ReentrantReadWriteLock.WriteLock lock = bookmarksLock.writeLock();
try {
lock.lock();
bookmarks.removeAll(usedBookmarks);
bookmarks.addAll(newBookmarks);
} finally {
lock.unlock();
}
});
return new DelegatingQueryRunner(queryRunner, lastBookmarks, bookmarkManager::updateBookmarks);
}
private static class DelegatingQueryRunner implements QueryRunner {

View File

@@ -31,6 +31,7 @@ import org.springframework.core.convert.converter.ConverterRegistry;
import org.springframework.core.convert.support.DefaultConversionService;
import org.springframework.dao.DataAccessException;
import org.springframework.data.neo4j.core.convert.Neo4jConversions;
import org.springframework.data.neo4j.core.transaction.Neo4jBookmarkManager;
import org.springframework.data.neo4j.core.transaction.Neo4jTransactionUtils;
import org.springframework.data.neo4j.core.transaction.ReactiveNeo4jTransactionManager;
import org.springframework.lang.Nullable;
@@ -43,12 +44,8 @@ import reactor.util.function.Tuple2;
import reactor.util.function.Tuples;
import java.util.Collection;
import java.util.Collections;
import java.util.HashSet;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.function.BiConsumer;
import java.util.function.BiFunction;
import java.util.function.Function;
@@ -70,9 +67,8 @@ final class DefaultReactiveNeo4jClient implements ReactiveNeo4jClient {
private final ConversionService conversionService;
private final Neo4jPersistenceExceptionTranslator persistenceExceptionTranslator = new Neo4jPersistenceExceptionTranslator();
// Basically a local bookmark manager
private final Set<Bookmark> bookmarks = new HashSet<>();
private final ReentrantReadWriteLock bookmarksLock = new ReentrantReadWriteLock();
// Local bookmark manager when using outside managed transactions
private final Neo4jBookmarkManager bookmarkManager;
DefaultReactiveNeo4jClient(Builder builder) {
@@ -82,6 +78,7 @@ final class DefaultReactiveNeo4jClient implements ReactiveNeo4jClient {
this.conversionService = new DefaultConversionService();
Optional.ofNullable(builder.neo4jConversions).orElseGet(Neo4jConversions::new).registerConvertersIn((ConverterRegistry) conversionService);
this.bookmarkManager = builder.bookmarkManager != null ? builder.bookmarkManager : Neo4jBookmarkManager.createReactive();
}
@Override
@@ -91,27 +88,12 @@ final class DefaultReactiveNeo4jClient implements ReactiveNeo4jClient {
.flatMap(targetDatabaseAndUser ->
ReactiveNeo4jTransactionManager.retrieveReactiveTransaction(driver, targetDatabaseAndUser.getT1(), targetDatabaseAndUser.getT2())
.map(ReactiveQueryRunner.class::cast)
.zipWith(Mono.just(Collections.<Bookmark>emptySet()))
.zipWith(Mono.just(bookmarkManager.getBookmarks()))
.switchIfEmpty(Mono.fromSupplier(() -> {
ReentrantReadWriteLock.ReadLock lock = bookmarksLock.readLock();
try {
lock.lock();
Set<Bookmark> lastBookmarks = new HashSet<>(bookmarks);
return Tuples.of(driver.session(ReactiveSession.class, Neo4jTransactionUtils.sessionConfig(false, lastBookmarks, targetDatabaseAndUser.getT1(), targetDatabaseAndUser.getT2())), lastBookmarks);
} finally {
lock.unlock();
}
Collection<Bookmark> 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 {

View File

@@ -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);
}

View File

@@ -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);
}

View File

@@ -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
* <p>
@@ -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
* <p>
* 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<Set<Bookmark>> 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

View File

@@ -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<Bookmark> bookmarks = Collections.synchronizedSet(new HashSet<>());
private final Supplier<Set<Bookmark>> bookmarksSupplier;
@Nullable
private ApplicationEventPublisher applicationEventPublisher;
ReactiveDefaultBookmarkManager(@Nullable Supplier<Set<Bookmark>> bookmarksSupplier) {
this.bookmarksSupplier = bookmarksSupplier == null ? Collections::emptySet : bookmarksSupplier;
}
@Override
public Collection<Bookmark> getBookmarks() {
this.bookmarks.addAll(bookmarksSupplier.get());
return Collections.synchronizedSet(Collections.unmodifiableSet(this.bookmarks));
}
@Override
public void updateBookmarks(Collection<Bookmark> usedBookmarks, Collection<Bookmark> 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;
}
}

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -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
*/

View File

@@ -436,7 +436,7 @@ class ReactiveIssuesIT extends TestBase {
BookmarkCapture bookmarkCapture = bookmarkCapture();
return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider,
Neo4jBookmarkManager.create(bookmarkCapture));
Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Override

View File

@@ -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));
}
}
}

View File

@@ -358,7 +358,7 @@ public class ReactiveElementIdIT extends AbstractElementIdTestBase {
BookmarkCapture bookmarkCapture = bookmarkCapture();
return new ReactiveNeo4jTransactionManager(driver, databaseSelectionProvider,
Neo4jBookmarkManager.create(bookmarkCapture));
Neo4jBookmarkManager.createReactive(bookmarkCapture));
}
@Bean

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -305,7 +305,7 @@ class ReactiveDynamicRelationshipsIT extends DynamicRelationshipsITBase<PersonWi
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

View File

@@ -159,7 +159,7 @@ class ReactiveIdGeneratorsIT extends IdGeneratorsITBase {
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

View File

@@ -429,7 +429,7 @@ public class ReactiveImmutableAssignedIdsIT {
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

View File

@@ -401,7 +401,7 @@ public class ReactiveImmutableExternallyGeneratedIdsIT {
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

View File

@@ -377,7 +377,7 @@ public class ReactiveImmutableGeneratedIdsIT {
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

View File

@@ -27,7 +27,9 @@ import java.util.function.Consumer;
import java.util.function.Function;
import java.util.stream.Stream;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.Test;
import org.neo4j.cypherdsl.core.Cypher;
import org.neo4j.cypherdsl.core.Node;
@@ -42,6 +44,7 @@ import org.neo4j.driver.Transaction;
import org.neo4j.driver.async.AsyncQueryRunner;
import org.neo4j.driver.reactivestreams.ReactiveQueryRunner;
import org.neo4j.driver.reactivestreams.ReactiveResult;
import org.neo4j.driver.reactivestreams.ReactiveSession;
import org.neo4j.driver.summary.ResultSummary;
import org.reactivestreams.Publisher;
import org.springframework.beans.factory.annotation.Autowired;
@@ -61,6 +64,7 @@ import org.springframework.transaction.ReactiveTransactionManager;
import org.springframework.transaction.annotation.EnableTransactionManagement;
import org.springframework.transaction.reactive.TransactionalOperator;
import reactor.blockhound.BlockHound;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
@@ -68,14 +72,19 @@ import reactor.test.StepVerifier;
/**
* @author Michael J. Simons
*/
@Disabled // we cannot use BlockHound right now in the Maven build
@Neo4jIntegrationTest
class ReactiveNeo4jClientIT {
protected static Neo4jConnectionSupport neo4jConnectionSupport;
@BeforeAll
static void setupBlockHound() {
BlockHound.install();
}
@BeforeEach
void setupData(@Autowired BookmarkCapture bookmarkCapture, @Autowired Driver driver) {
try (
Session session = driver.session(bookmarkCapture.createSessionConfig());
Transaction transaction = session.beginTransaction()
@@ -85,6 +94,58 @@ class ReactiveNeo4jClientIT {
}
}
@Test // GH-2755
public void testQueryExecutionNeo4jClient(@Autowired ReactiveNeo4jClient neo4jClient, @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)";
String cypher2 = "MATCH (n) WHERE elementId(n) = $elementId RETURN elementId(n)";
StepVerifier.create(neo4jClient.query(cypher).fetchAs(String.class).all()
.flatMap(elementId ->
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

View File

@@ -957,7 +957,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

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -2882,7 +2882,7 @@ class ReactiveRepositoryIT {
return ReactiveNeo4jTransactionManager.with(driver)
.withDatabaseSelectionProvider(databaseSelectionProvider)
.withUserSelectionProvider(getUserSelectionProvider())
.withBookmarkManager(Neo4jBookmarkManager.create(bookmarkCapture()))
.withBookmarkManager(Neo4jBookmarkManager.createReactive(bookmarkCapture()))
.build();
}

View File

@@ -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

View File

@@ -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

View File

@@ -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