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
20 changes: 19 additions & 1 deletion src/main/java/org/cache/Main.java
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,24 @@ public CacheOperations<Object> routedCacheService(
return new RoutedCacheService<>(cacheService, cacheNode, cacheConfig.clusterInfo(), forwardingClient, keyCodec);
}

@Bean
public CacheOperations<Object> clusterRoutedCacheService(
CacheService<Object> cacheService,
CacheNode cacheNode,
CacheConfig cacheConfig,
ClusterForwardingClient forwardingClient,
KeyCodec<Object> keyCodec
) {
return new RoutedCacheService<>(
cacheService,
cacheNode,
cacheConfig.clusterInfo(),
forwardingClient,
keyCodec,
false
);
}

@Bean
public CommandProcessor<Object> commandProcessor(KeyCodec<Object> keyCodec, CacheOperations<Object> cacheService) {
return new CommandProcessor<>(keyCodec, cacheService);
Expand All @@ -111,7 +129,7 @@ public CommandProcessor<Object> commandProcessor(KeyCodec<Object> keyCodec, Cach
@Bean
public CommandProcessor<Object> clusterCommandProcessor(
KeyCodec<Object> keyCodec,
CacheService<Object> cacheService
@Qualifier("clusterRoutedCacheService") CacheOperations<Object> cacheService
) {
return new CommandProcessor<>(keyCodec, cacheService);
}
Expand Down
64 changes: 42 additions & 22 deletions src/main/java/org/cache/cluster/routing/RoutedCacheService.java
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ public class RoutedCacheService<K> implements CacheOperations<K> {
private final ConsistentHashRing hashRing;
private final KeyCodec<K> keyCodec;
private final CacheResponseParser responseParser;
private final boolean forwardingAllowed;

public RoutedCacheService(
CacheOperations<K> localService,
Expand All @@ -27,7 +28,18 @@ public RoutedCacheService(
ClusterForwardingClient forwardingClient,
KeyCodec<K> keyCodec
) {
this(localService, currentNode, clusterInfo, forwardingClient, keyCodec, new CacheResponseParser());
this(localService, currentNode, clusterInfo, forwardingClient, keyCodec, true);
}

public RoutedCacheService(
CacheOperations<K> localService,
CacheNode currentNode,
ClusterInfo clusterInfo,
ClusterForwardingClient forwardingClient,
KeyCodec<K> keyCodec,
boolean forwardingAllowed
) {
this(localService, currentNode, clusterInfo, forwardingClient, keyCodec, new CacheResponseParser(), forwardingAllowed);
}

RoutedCacheService(
Expand All @@ -36,25 +48,27 @@ public RoutedCacheService(
ClusterInfo clusterInfo,
ClusterForwardingClient forwardingClient,
KeyCodec<K> keyCodec,
CacheResponseParser responseParser
CacheResponseParser responseParser,
boolean forwardingAllowed
) {
this.localService = localService;
this.currentNode = currentNode;
this.forwardingClient = forwardingClient;
this.hashRing = clusterInfo == null ? null : new ConsistentHashRing(clusterInfo);
this.keyCodec = keyCodec;
this.responseParser = responseParser;
this.forwardingAllowed = forwardingAllowed;
}

@Override
public void putString(K key, String value, long ttlMillis) {
CacheNode owner = ownerFor(key);
if (isLocal(owner)) {
Optional<CacheNode> remoteOwner = remoteOwnerFor(key);
if (remoteOwner.isEmpty()) {
localService.putString(key, value, ttlMillis);
return;
}

expectOk(forwardingClient.forward(owner, List.of(
expectOk(forwardingClient.forward(remoteOwner.get(), List.of(
"PUT",
keyCodec.encode(key),
value,
Expand All @@ -64,12 +78,12 @@ public void putString(K key, String value, long ttlMillis) {

@Override
public Optional<String> getString(K key) {
CacheNode owner = ownerFor(key);
if (isLocal(owner)) {
Optional<CacheNode> remoteOwner = remoteOwnerFor(key);
if (remoteOwner.isEmpty()) {
return localService.getString(key);
}

List<String> response = forwardingClient.forward(owner, List.of("GET", keyCodec.encode(key)));
List<String> response = forwardingClient.forward(remoteOwner.get(), List.of("GET", keyCodec.encode(key)));
if (responseParser.isNotFound(response)) {
return Optional.empty();
}
Expand All @@ -82,23 +96,23 @@ public Optional<String> getString(K key) {

@Override
public void push(K key, String value) {
CacheNode owner = ownerFor(key);
if (isLocal(owner)) {
Optional<CacheNode> remoteOwner = remoteOwnerFor(key);
if (remoteOwner.isEmpty()) {
localService.push(key, value);
return;
}

expectOk(forwardingClient.forward(owner, List.of("PUSH", keyCodec.encode(key), value)));
expectOk(forwardingClient.forward(remoteOwner.get(), List.of("PUSH", keyCodec.encode(key), value)));
}

@Override
public Optional<List<String>> lrange(K key, int from, int to) {
CacheNode owner = ownerFor(key);
if (isLocal(owner)) {
Optional<CacheNode> remoteOwner = remoteOwnerFor(key);
if (remoteOwner.isEmpty()) {
return localService.lrange(key, from, to);
}

List<String> response = forwardingClient.forward(owner, List.of(
List<String> response = forwardingClient.forward(remoteOwner.get(), List.of(
"LRANGE",
keyCodec.encode(key),
Integer.toString(from),
Expand All @@ -122,13 +136,13 @@ public Optional<List<String>> lrange(K key, int from, int to) {

@Override
public void delete(K key) {
CacheNode owner = ownerFor(key);
if (isLocal(owner)) {
Optional<CacheNode> remoteOwner = remoteOwnerFor(key);
if (remoteOwner.isEmpty()) {
localService.delete(key);
return;
}

expectOk(forwardingClient.forward(owner, List.of("DELETE", keyCodec.encode(key))));
expectOk(forwardingClient.forward(remoteOwner.get(), List.of("DELETE", keyCodec.encode(key))));
}

@Override
Expand All @@ -146,12 +160,18 @@ public Snapshot metrics() {
return localService.metrics();
}

private CacheNode ownerFor(K key) {
return hashRing == null ? currentNode : hashRing.nodeFor(keyCodec.encode(key));
}
private Optional<CacheNode> remoteOwnerFor(K key) {
CacheNode owner = hashRing == null ? currentNode : hashRing.nodeFor(keyCodec.encode(key));
if (owner.id().equals(currentNode.id())) {
return Optional.empty();
}

if (!forwardingAllowed) {
throw new ClusterForwardingException("Request routed to wrong node. Expected owner: " + owner.id()
+ ", current node: " + currentNode.id());
}

private boolean isLocal(CacheNode owner) {
return owner.id().equals(currentNode.id());
return Optional.of(owner);
}

private void expectOk(List<String> response) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
import java.util.Optional;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
Expand Down Expand Up @@ -63,6 +64,26 @@ void getStringForwardsWhenAnotherNodeOwnsKey() {
verifyNoInteractions(localService);
}

@Test
void getStringRejectsRemoteOwnerWhenForwardingIsDisabled() {
CacheOperations<String> localService = localService();
ClusterForwardingClient forwardingClient = mock(ClusterForwardingClient.class);
String key = keyOwnedBy(nodeB);
RoutedCacheService<String> service = new RoutedCacheService<>(
localService,
nodeA,
clusterInfo,
forwardingClient,
keyCodec,
false
);

var exception = assertThrows(ClusterForwardingException.class, () -> service.getString(key));

assertEquals("Request routed to wrong node. Expected owner: node-b, current node: node-a", exception.getMessage());
verifyNoInteractions(localService, forwardingClient);
}

@Test
void lrangeForwardsAndReturnsListValues() {
CacheOperations<String> localService = localService();
Expand Down
Loading