diff --git a/src/main/java/org/cache/cluster/routing/ReplicationTargets.java b/src/main/java/org/cache/cluster/routing/ReplicationTargets.java new file mode 100644 index 0000000..1e10656 --- /dev/null +++ b/src/main/java/org/cache/cluster/routing/ReplicationTargets.java @@ -0,0 +1,8 @@ +package org.cache.cluster.routing; + +import org.cache.cluster.CacheNode; + +import java.util.List; + +public record ReplicationTargets(boolean includesCurrentNode, List remoteNodes) { +} \ No newline at end of file diff --git a/src/main/java/org/cache/cluster/routing/RoutedCacheService.java b/src/main/java/org/cache/cluster/routing/RoutedCacheService.java index ea0e24c..837b3d7 100644 --- a/src/main/java/org/cache/cluster/routing/RoutedCacheService.java +++ b/src/main/java/org/cache/cluster/routing/RoutedCacheService.java @@ -62,87 +62,89 @@ public RoutedCacheService( @Override public void putString(K key, String value, long ttlMillis) { - Optional remoteOwner = remoteOwnerFor(key); - if (remoteOwner.isEmpty()) { + ReplicationTargets targets = writeTargetsFor(key); + if (targets.includesCurrentNode()) { localService.putString(key, value, ttlMillis); - return; } - expectOk(forwardingClient.forward(remoteOwner.get(), List.of( - "PUT", - keyCodec.encode(key), - value, - Long.toString(ttlMillis) - ))); + for (CacheNode replica : targets.remoteNodes()) { + expectOk(forwardingClient.forward(replica, List.of( + "PUT", + keyCodec.encode(key), + value, + Long.toString(ttlMillis) + ))); + } } @Override public Optional getString(K key) { - Optional remoteOwner = remoteOwnerFor(key); - if (remoteOwner.isEmpty()) { - return localService.getString(key); + ClusterForwardingException failure = null; + + for (CacheNode owner : readOwnersFor(key)) { + try { + if (isCurrentNode(owner)) { + return localService.getString(key); + } + + return remoteGetString(owner, key); + } catch (ClusterForwardingException exception) { + failure = exception; + } } - List response = forwardingClient.forward(remoteOwner.get(), List.of("GET", keyCodec.encode(key))); - if (responseParser.isNotFound(response)) { - return Optional.empty(); + if (failure == null) { + throw new ClusterForwardingException("No replica owner found for key: " + keyCodec.encode(key)); } - return responseParser.value(response) - .or(() -> { - throw new ClusterForwardingException("Unexpected cluster response: " + response); - }); + throw failure; } @Override public void push(K key, String value) { - Optional remoteOwner = remoteOwnerFor(key); - if (remoteOwner.isEmpty()) { + ReplicationTargets targets = writeTargetsFor(key); + if (targets.includesCurrentNode()) { localService.push(key, value); - return; } - expectOk(forwardingClient.forward(remoteOwner.get(), List.of("PUSH", keyCodec.encode(key), value))); + for (CacheNode replica : targets.remoteNodes()) { + expectOk(forwardingClient.forward(replica, List.of("PUSH", keyCodec.encode(key), value))); + } } @Override public Optional> lrange(K key, int from, int to) { - Optional remoteOwner = remoteOwnerFor(key); - if (remoteOwner.isEmpty()) { - return localService.lrange(key, from, to); - } - - List response = forwardingClient.forward(remoteOwner.get(), List.of( - "LRANGE", - keyCodec.encode(key), - Integer.toString(from), - Integer.toString(to) - )); + ClusterForwardingException failure = null; - if (responseParser.isNotFound(response)) { - return Optional.empty(); - } + for (CacheNode owner : readOwnersFor(key)) { + try { + if (isCurrentNode(owner)) { + return localService.lrange(key, from, to); + } - if (response.isEmpty()) { - return Optional.of(List.of()); + return remoteLrange(owner, key, from, to); + } catch (ClusterForwardingException exception) { + failure = exception; + } } - if (!responseParser.isKnownResponse(response)) { - return Optional.of(List.copyOf(response)); + if (failure == null) { + throw new ClusterForwardingException("No replica owner found for key: " + keyCodec.encode(key)); } - throw new ClusterForwardingException("Unexpected cluster response: " + response); + throw failure; } @Override public void delete(K key) { - Optional remoteOwner = remoteOwnerFor(key); - if (remoteOwner.isEmpty()) { + ReplicationTargets targets = writeTargetsFor(key); + if (targets.includesCurrentNode()) { localService.delete(key); - return; } - expectOk(forwardingClient.forward(remoteOwner.get(), List.of("DELETE", keyCodec.encode(key)))); + for (CacheNode replica : targets.remoteNodes()) { + expectOk(forwardingClient.forward(replica, List.of("DELETE", keyCodec.encode(key)))); + } } @Override @@ -160,18 +162,97 @@ public Snapshot metrics() { return localService.metrics(); } - private Optional remoteOwnerFor(K key) { - CacheNode owner = hashRing == null ? currentNode : hashRing.nodeFor(keyCodec.encode(key)); - if (owner.id().equals(currentNode.id())) { - return Optional.empty(); + private ReplicationTargets writeTargetsFor(K key) { + if (hashRing == null) { + return new ReplicationTargets(true, List.of()); + } + + List owners = hashRing.nodesFor(keyCodec.encode(key)); + boolean includesCurrentNode = owners.stream() + .anyMatch(this::isCurrentNode); + + if (!includesCurrentNode && !forwardingAllowed) { + throw new ClusterForwardingException("Request routed to wrong node. Expected one of owners: " + + ownerIds(owners) + ", current node: " + currentNode.id()); } if (!forwardingAllowed) { - throw new ClusterForwardingException("Request routed to wrong node. Expected owner: " + owner.id() - + ", current node: " + currentNode.id()); + return new ReplicationTargets(true, List.of()); + } + + return new ReplicationTargets( + includesCurrentNode, + owners.stream() + .filter(owner -> !isCurrentNode(owner)) + .toList() + ); + } + + private List readOwnersFor(K key) { + if (hashRing == null) { + return List.of(currentNode); } - return Optional.of(owner); + List owners = hashRing.nodesFor(keyCodec.encode(key)); + boolean includesCurrentNode = owners.stream() + .anyMatch(this::isCurrentNode); + + if (!includesCurrentNode && !forwardingAllowed) { + throw new ClusterForwardingException("Request routed to wrong node. Expected one of owners: " + + ownerIds(owners) + ", current node: " + currentNode.id()); + } + + if (!forwardingAllowed) { + return List.of(currentNode); + } + + return owners; + } + + private Optional remoteGetString(CacheNode owner, K key) { + List response = forwardingClient.forward(owner, List.of("GET", keyCodec.encode(key))); + if (responseParser.isNotFound(response)) { + return Optional.empty(); + } + + return responseParser.value(response) + .or(() -> { + throw new ClusterForwardingException("Unexpected cluster response: " + response); + }); + } + + private Optional> remoteLrange(CacheNode owner, K key, int from, int to) { + List response = forwardingClient.forward(owner, List.of( + "LRANGE", + keyCodec.encode(key), + Integer.toString(from), + Integer.toString(to) + )); + + if (responseParser.isNotFound(response)) { + return Optional.empty(); + } + + if (response.isEmpty()) { + return Optional.of(List.of()); + } + + if (!responseParser.isKnownResponse(response)) { + return Optional.of(List.copyOf(response)); + } + + throw new ClusterForwardingException("Unexpected cluster response: " + response); + } + + private boolean isCurrentNode(CacheNode owner) { + return owner.id().equals(currentNode.id()); + } + + private String ownerIds(List owners) { + return owners.stream() + .map(CacheNode::id) + .toList() + .toString(); } private void expectOk(List response) { diff --git a/src/main/java/org/cache/network/tcp/TcpCacheServerLifecycle.java b/src/main/java/org/cache/network/tcp/TcpCacheServerLifecycle.java index ca70903..78c0425 100644 --- a/src/main/java/org/cache/network/tcp/TcpCacheServerLifecycle.java +++ b/src/main/java/org/cache/network/tcp/TcpCacheServerLifecycle.java @@ -11,10 +11,6 @@ public class TcpCacheServerLifecycle implements SmartLifecycle { private volatile boolean running; private Thread serverThread; - public TcpCacheServerLifecycle(TcpCacheServer server) { - this(server, "tcp-cache-server"); - } - public TcpCacheServerLifecycle(TcpCacheServer server, String threadName) { this.server = server; this.threadName = threadName; diff --git a/src/main/java/org/cache/network/tcp/connection/CacheResponseParser.java b/src/main/java/org/cache/network/tcp/connection/CacheResponseParser.java index e4b3f57..bf1feb5 100644 --- a/src/main/java/org/cache/network/tcp/connection/CacheResponseParser.java +++ b/src/main/java/org/cache/network/tcp/connection/CacheResponseParser.java @@ -32,14 +32,6 @@ public Optional value(List response) { return Optional.empty(); } - public Optional errorMessage(List response) { - if (isError(response) && response.size() > RESPONSE_VALUE_INDEX) { - return Optional.of(response.get(RESPONSE_VALUE_INDEX)); - } - - return Optional.empty(); - } - public boolean isKnownResponse(List response) { return !response.isEmpty() && isKnownResponse(response.getFirst()); } diff --git a/src/test/java/org/cache/cluster/routing/RoutedCacheServiceTest.java b/src/test/java/org/cache/cluster/routing/RoutedCacheServiceTest.java index 68dd4ee..8ba76f3 100644 --- a/src/test/java/org/cache/cluster/routing/RoutedCacheServiceTest.java +++ b/src/test/java/org/cache/cluster/routing/RoutedCacheServiceTest.java @@ -21,7 +21,9 @@ class RoutedCacheServiceTest { private final CacheNode nodeA = node("node-a", 8080, 2020, 10001); private final CacheNode nodeB = node("node-b", 8081, 2021, 10002); + private final CacheNode nodeC = node("node-c", 8082, 2022, 10003); private final ClusterInfo clusterInfo = new ClusterInfo(1, List.of(nodeA, nodeB)); + private final ClusterInfo replicatedClusterInfo = new ClusterInfo(2, List.of(nodeA, nodeB, nodeC)); private final StringKeyCodec keyCodec = new StringKeyCodec(); @Test @@ -64,15 +66,65 @@ void getStringForwardsWhenAnotherNodeOwnsKey() { verifyNoInteractions(localService); } + @Test + void getStringFallsBackToRemoteReplicaWhenPrimaryForwardFails() { + CacheOperations localService = localService(); + ClusterForwardingClient forwardingClient = mock(ClusterForwardingClient.class); + String key = keyReplicatedTo(nodeB, nodeC); + RoutedCacheService service = new RoutedCacheService<>( + localService, + nodeA, + replicatedClusterInfo, + forwardingClient, + keyCodec + ); + + when(forwardingClient.forward(nodeB, List.of("GET", key))) + .thenThrow(new ClusterForwardingException("primary unavailable")); + when(forwardingClient.forward(nodeC, List.of("GET", key))).thenReturn(List.of("VALUE", "apple")); + + assertEquals(Optional.of("apple"), service.getString(key)); + verify(forwardingClient).forward(nodeB, List.of("GET", key)); + verify(forwardingClient).forward(nodeC, List.of("GET", key)); + verifyNoInteractions(localService); + } + + @Test + void getStringFallsBackToLocalReplicaWhenPrimaryForwardFails() { + CacheOperations localService = localService(); + ClusterForwardingClient forwardingClient = mock(ClusterForwardingClient.class); + String key = keyReplicatedToRemotePrimaryAndLocalReplica(); + CacheNode primaryOwner = replicaOwnersFor(key).getFirst(); + RoutedCacheService service = new RoutedCacheService<>( + localService, + nodeA, + replicatedClusterInfo, + forwardingClient, + keyCodec + ); + + when(forwardingClient.forward(primaryOwner, List.of("GET", key))) + .thenThrow(new ClusterForwardingException("primary unavailable")); + when(localService.getString(key)).thenReturn(Optional.of("apple")); + + assertEquals(Optional.of("apple"), service.getString(key)); + verify(forwardingClient).forward(primaryOwner, List.of("GET", key)); + verify(localService).getString(key); + } + @Test void getStringRejectsRemoteOwnerWhenForwardingIsDisabled() { CacheOperations localService = localService(); ClusterForwardingClient forwardingClient = mock(ClusterForwardingClient.class); - String key = keyOwnedBy(nodeB); + String key = keyNotReplicatedTo(nodeA); + List ownerIds = replicaOwnersFor(key) + .stream() + .map(CacheNode::id) + .toList(); RoutedCacheService service = new RoutedCacheService<>( localService, nodeA, - clusterInfo, + replicatedClusterInfo, forwardingClient, keyCodec, false @@ -80,7 +132,105 @@ void getStringRejectsRemoteOwnerWhenForwardingIsDisabled() { var exception = assertThrows(ClusterForwardingException.class, () -> service.getString(key)); - assertEquals("Request routed to wrong node. Expected owner: node-b, current node: node-a", exception.getMessage()); + assertEquals("Request routed to wrong node. Expected one of owners: " + ownerIds + ", current node: node-a", + exception.getMessage()); + verifyNoInteractions(localService, forwardingClient); + } + + @Test + void putStringWritesToLocalAndRemoteReplicas() { + CacheOperations localService = localService(); + ClusterForwardingClient forwardingClient = mock(ClusterForwardingClient.class); + String key = keyReplicatedIncluding(nodeA); + List remoteOwners = replicaOwnersFor(key) + .stream() + .filter(owner -> !owner.id().equals(nodeA.id())) + .toList(); + RoutedCacheService service = new RoutedCacheService<>( + localService, + nodeA, + replicatedClusterInfo, + forwardingClient, + keyCodec + ); + + for (CacheNode owner : remoteOwners) { + when(forwardingClient.forward(owner, List.of("PUT", key, "apple", "1000"))).thenReturn(List.of("OK")); + } + + service.putString(key, "apple", 1_000); + + verify(localService).putString(key, "apple", 1_000); + for (CacheNode owner : remoteOwners) { + verify(forwardingClient).forward(owner, List.of("PUT", key, "apple", "1000")); + } + } + + @Test + void putStringWritesToAllRemoteReplicasWhenCurrentNodeIsNotReplicaOwner() { + CacheOperations localService = localService(); + ClusterForwardingClient forwardingClient = mock(ClusterForwardingClient.class); + String key = keyReplicatedTo(nodeB, nodeC); + RoutedCacheService service = new RoutedCacheService<>( + localService, + nodeA, + replicatedClusterInfo, + forwardingClient, + keyCodec + ); + + when(forwardingClient.forward(nodeB, List.of("PUT", key, "apple", "1000"))).thenReturn(List.of("OK")); + when(forwardingClient.forward(nodeC, List.of("PUT", key, "apple", "1000"))).thenReturn(List.of("OK")); + + service.putString(key, "apple", 1_000); + + verify(forwardingClient).forward(nodeB, List.of("PUT", key, "apple", "1000")); + verify(forwardingClient).forward(nodeC, List.of("PUT", key, "apple", "1000")); + verifyNoInteractions(localService); + } + + @Test + void putStringExecutesLocallyOnReplicaOwnerWhenForwardingIsDisabled() { + CacheOperations localService = localService(); + ClusterForwardingClient forwardingClient = mock(ClusterForwardingClient.class); + String key = keyReplicatedIncluding(nodeA); + RoutedCacheService service = new RoutedCacheService<>( + localService, + nodeA, + replicatedClusterInfo, + forwardingClient, + keyCodec, + false + ); + + service.putString(key, "apple", 1_000); + + verify(localService).putString(key, "apple", 1_000); + verifyNoInteractions(forwardingClient); + } + + @Test + void putStringRejectsNonReplicaOwnerWhenForwardingIsDisabled() { + CacheOperations localService = localService(); + ClusterForwardingClient forwardingClient = mock(ClusterForwardingClient.class); + String key = keyNotReplicatedTo(nodeA); + List ownerIds = replicaOwnersFor(key) + .stream() + .map(CacheNode::id) + .toList(); + RoutedCacheService service = new RoutedCacheService<>( + localService, + nodeA, + replicatedClusterInfo, + forwardingClient, + keyCodec, + false + ); + + var exception = assertThrows(ClusterForwardingException.class, () -> service.putString(key, "apple", 1_000)); + + assertEquals("Request routed to wrong node. Expected one of owners: " + ownerIds + ", current node: node-a", + exception.getMessage()); verifyNoInteractions(localService, forwardingClient); } @@ -104,6 +254,29 @@ void lrangeForwardsAndReturnsListValues() { verifyNoInteractions(localService); } + @Test + void lrangeFallsBackToRemoteReplicaWhenPrimaryForwardFails() { + CacheOperations localService = localService(); + ClusterForwardingClient forwardingClient = mock(ClusterForwardingClient.class); + String key = keyReplicatedTo(nodeB, nodeC); + RoutedCacheService service = new RoutedCacheService<>( + localService, + nodeA, + replicatedClusterInfo, + forwardingClient, + keyCodec + ); + + when(forwardingClient.forward(nodeB, List.of("LRANGE", key, "0", "2"))) + .thenThrow(new ClusterForwardingException("primary unavailable")); + when(forwardingClient.forward(nodeC, List.of("LRANGE", key, "0", "2"))).thenReturn(List.of("one", "two")); + + assertEquals(Optional.of(List.of("one", "two")), service.lrange(key, 0, 2)); + verify(forwardingClient).forward(nodeB, List.of("LRANGE", key, "0", "2")); + verify(forwardingClient).forward(nodeC, List.of("LRANGE", key, "0", "2")); + verifyNoInteractions(localService); + } + @Test void clusterDisabledUsesLocalService() { CacheOperations localService = localService(); @@ -135,6 +308,74 @@ private String keyOwnedBy(CacheNode owner) { throw new IllegalStateException("Could not find key owned by node: " + owner.id()); } + private String keyReplicatedTo(CacheNode firstOwner, CacheNode secondOwner) { + ConsistentHashRing ring = new ConsistentHashRing(replicatedClusterInfo); + + for (int index = 0; index < 10_000; index++) { + String key = "replica-key-" + index; + List ownerIds = ring.nodesFor(key) + .stream() + .map(CacheNode::id) + .toList(); + + if (ownerIds.equals(List.of(firstOwner.id(), secondOwner.id()))) { + return key; + } + } + + throw new IllegalStateException("Could not find key replicated to owners: " + + firstOwner.id() + ", " + secondOwner.id()); + } + + private String keyReplicatedIncluding(CacheNode owner) { + for (int index = 0; index < 10_000; index++) { + String key = "replica-key-" + index; + boolean containsOwner = replicaOwnersFor(key).stream() + .anyMatch(replicaOwner -> replicaOwner.id().equals(owner.id())); + + if (containsOwner) { + return key; + } + } + + throw new IllegalStateException("Could not find key replicated to owner: " + owner.id()); + } + + private String keyNotReplicatedTo(CacheNode owner) { + for (int index = 0; index < 10_000; index++) { + String key = "replica-key-" + index; + boolean containsOwner = replicaOwnersFor(key).stream() + .anyMatch(replicaOwner -> replicaOwner.id().equals(owner.id())); + + if (!containsOwner) { + return key; + } + } + + throw new IllegalStateException("Could not find key not replicated to owner: " + owner.id()); + } + + private String keyReplicatedToRemotePrimaryAndLocalReplica() { + for (int index = 0; index < 10_000; index++) { + String key = "replica-key-" + index; + List owners = replicaOwnersFor(key); + boolean hasRemotePrimary = !owners.getFirst().id().equals(nodeA.id()); + boolean hasLocalReplica = owners.stream() + .skip(1) + .anyMatch(owner -> owner.id().equals(nodeA.id())); + + if (hasRemotePrimary && hasLocalReplica) { + return key; + } + } + + throw new IllegalStateException("Could not find key with remote primary and local replica"); + } + + private List replicaOwnersFor(String key) { + return new ConsistentHashRing(replicatedClusterInfo).nodesFor(key); + } + private CacheNode node(String id, int httpPort, int tcpPort, int clusterPort) { return new CacheNode(id, "localhost", httpPort, tcpPort, clusterPort); }