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
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
package org.cache.cluster.routing;

import org.cache.cluster.CacheNode;

import java.util.List;

public record ReplicationTargets(boolean includesCurrentNode, List<CacheNode> remoteNodes) {
}
187 changes: 134 additions & 53 deletions src/main/java/org/cache/cluster/routing/RoutedCacheService.java
Original file line number Diff line number Diff line change
Expand Up @@ -62,87 +62,89 @@ public RoutedCacheService(

@Override
public void putString(K key, String value, long ttlMillis) {
Optional<CacheNode> 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<String> getString(K key) {
Optional<CacheNode> 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<String> 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<CacheNode> 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<List<String>> lrange(K key, int from, int to) {
Optional<CacheNode> remoteOwner = remoteOwnerFor(key);
if (remoteOwner.isEmpty()) {
return localService.lrange(key, from, to);
}

List<String> 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<CacheNode> 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
Expand All @@ -160,18 +162,97 @@ public Snapshot metrics() {
return localService.metrics();
}

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();
private ReplicationTargets writeTargetsFor(K key) {
if (hashRing == null) {
return new ReplicationTargets(true, List.of());
}

List<CacheNode> 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<CacheNode> readOwnersFor(K key) {
if (hashRing == null) {
return List.of(currentNode);
}

return Optional.of(owner);
List<CacheNode> 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<String> remoteGetString(CacheNode owner, K key) {
List<String> 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<List<String>> remoteLrange(CacheNode owner, K key, int from, int to) {
List<String> 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<CacheNode> owners) {
return owners.stream()
.map(CacheNode::id)
.toList()
.toString();
}

private void expectOk(List<String> response) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,14 +32,6 @@ public Optional<String> value(List<String> response) {
return Optional.empty();
}

public Optional<String> errorMessage(List<String> response) {
if (isError(response) && response.size() > RESPONSE_VALUE_INDEX) {
return Optional.of(response.get(RESPONSE_VALUE_INDEX));
}

return Optional.empty();
}

public boolean isKnownResponse(List<String> response) {
return !response.isEmpty() && isKnownResponse(response.getFirst());
}
Expand Down
Loading
Loading