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
33 changes: 15 additions & 18 deletions src/main/java/org/cache/Main.java
Original file line number Diff line number Diff line change
Expand Up @@ -3,27 +3,24 @@

import org.cache.core.LocalCache;
import org.cache.eviction.LruEvictionPolicy;
import org.cache.network.TcpCacheServer;
import org.cache.protocol.CommandParser;
import org.cache.protocol.CommandProcessor;
import org.cache.protocol.codec.StringCodec;

public class Main {
public static void main(String[] args) {

LruEvictionPolicy<String> lruEvictionPolicy = new LruEvictionPolicy<>();
LocalCache<String, String> localCache = new LocalCache<>(2, lruEvictionPolicy);

localCache.put("test", "first value", 100);

System.out.println(localCache.get("test").orElse(null));
import java.io.IOException;

localCache.put("test1", "second value", 100);
System.out.println(localCache.get("test").orElse(null));
localCache.put("test2", "third value", 100);


System.out.println(localCache.size());
System.out.println(localCache.get("test"));
System.out.println(localCache.get("test1").orElse(null));
public class Main {
public static void main(String[] args) throws IOException {
int port = args.length > 0 ? Integer.parseInt(args[0]) : 2020;
var cache = new LocalCache<String, String>(1_000, new LruEvictionPolicy<>());
var stringCodec = new StringCodec();
var commandParser = new CommandParser<>(stringCodec, stringCodec);
var commandProcessor = new CommandProcessor<>(cache, commandParser, stringCodec);


System.out.println(localCache.metrics().getEvictions());
try (cache; var server = new TcpCacheServer(port, commandProcessor)) {
server.start();
}
}
}
4 changes: 4 additions & 0 deletions src/main/java/org/cache/core/Cache.java
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
package org.cache.core;

import org.cache.core.metrics.Snapshot;

import java.util.Optional;

public interface Cache<K, V> {
Expand All @@ -13,4 +15,6 @@ public interface Cache<K, V> {
int size();

void clear();

Snapshot metrics();
}
1 change: 1 addition & 0 deletions src/main/java/org/cache/core/LocalCache.java
Original file line number Diff line number Diff line change
Expand Up @@ -163,6 +163,7 @@ public void close() {
cleanupScheduler.shutdownNow();
}

@Override
public Snapshot metrics() {
return metrics.snapshot();
}
Expand Down
40 changes: 40 additions & 0 deletions src/main/java/org/cache/network/ClientConnectionHandler.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
package org.cache.network;

import org.cache.protocol.CommandProcessor;

import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
import java.io.PrintWriter;
import java.net.Socket;
import java.nio.charset.StandardCharsets;

public class ClientConnectionHandler implements Runnable {

private final Socket socket;
private final CommandProcessor<?, ?> commandProcessor;

public ClientConnectionHandler(Socket socket, CommandProcessor<?, ?> commandProcessor) {
this.socket = socket;
this.commandProcessor = commandProcessor;
}

@Override
public void run() {
try (
socket;
var reader = new BufferedReader(
new InputStreamReader(socket.getInputStream(), StandardCharsets.UTF_8)
);
var writer = new PrintWriter(socket.getOutputStream(), true, StandardCharsets.UTF_8)
) {
String command;
while ((command = reader.readLine()) != null) {
String response = commandProcessor.process(command);
writer.println(response);
}
} catch (IOException exception) {
System.err.println("Client connection failed: " + exception.getMessage());
}
}
}
44 changes: 44 additions & 0 deletions src/main/java/org/cache/network/TcpCacheServer.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
package org.cache.network;

import org.cache.protocol.CommandProcessor;

import java.io.IOException;
import java.net.ServerSocket;
import java.net.Socket;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class TcpCacheServer implements AutoCloseable {

private final int port;
private final CommandProcessor<?, ?> commandProcessor;
private final ExecutorService executor;
private ServerSocket serverSocket;

private final static int NUMBER_OF_THREADS = 16;

public TcpCacheServer(int port, CommandProcessor<?, ?> commandProcessor) {
this.port = port;
this.commandProcessor = commandProcessor;
this.executor = Executors.newFixedThreadPool(NUMBER_OF_THREADS);
}

public void start() throws IOException {
serverSocket = new ServerSocket(port);
System.out.println("TCP cache server listening on port " + port);

while (!Thread.currentThread().isInterrupted() && !serverSocket.isClosed()) {
Socket socket = serverSocket.accept();
executor.submit(new ClientConnectionHandler(socket, commandProcessor));
}
}

@Override
public void close() throws IOException {
executor.shutdownNow();

if (serverSocket != null && !serverSocket.isClosed()) {
serverSocket.close();
}
}
}
102 changes: 102 additions & 0 deletions src/main/java/org/cache/protocol/CommandParser.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
package org.cache.protocol;

import org.cache.protocol.codec.Codec;
import org.cache.protocol.commands.*;

import java.util.concurrent.TimeUnit;

public class CommandParser<K, V> {

private final Codec<K> keyCodec;
private final Codec<V> valueCodec;


private final static int ONE = 1;
private final static int TWO = 2;
private final static int THREE = 3;
private final static int FOUR = 4;

public CommandParser(Codec<K> keyCodec, Codec<V> valueCodec) {
this.keyCodec = keyCodec;
this.valueCodec = valueCodec;
}

public CacheCommand<K, V> parse(String rawCommand) {
if (rawCommand == null || rawCommand.isBlank()) {
return new UnknownCommand<>();
}

String[] parts = rawCommand.trim().split("\\s+");

try {
CommandType type = CommandType.valueOf(parts[0].toUpperCase());
return switch (type) {
case PUT -> parsePut(parts);
case GET -> parseGet(parts);
case DELETE -> parseDelete(parts);
case SIZE -> parseSize(parts);
case CLEAR -> parseClear(parts);
case METRICS -> parseMetrics(parts);
case UNKNOWN -> new UnknownCommand<>();
};
} catch (IllegalArgumentException exception) {
return new UnknownCommand<>();
}
}

private CacheCommand<K, V> parsePut(String[] parts) {
if (parts.length != THREE && parts.length != FOUR) {
return new InvalidCommand<>("usage: PUT key value [ttlSeconds]");
}

try {
return new PutCommand<>(
keyCodec.decode(parts[ONE]),
valueCodec.decode(parts[TWO]),
Long.parseLong(parts[THREE])
);
} catch (NumberFormatException exception) {
return new InvalidCommand<>("ttl must be a number");
}
}

private CacheCommand<K, V> parseGet(String[] parts) {
if (parts.length != TWO) {
return new InvalidCommand<>("usage: GET key");
}

return new GetCommand<>(keyCodec.decode(parts[1]));
}

private CacheCommand<K, V> parseDelete(String[] parts) {
if (parts.length != TWO) {
return new InvalidCommand<>("usage: DELETE key");
}

return new DeleteCommand<>(keyCodec.decode(parts[1]));
}

private CacheCommand<K, V> parseSize(String[] parts) {
if (parts.length != ONE) {
return new InvalidCommand<>("usage: SIZE");
}

return new SizeCommand<>();
}

private CacheCommand<K, V> parseClear(String[] parts) {
if (parts.length != ONE) {
return new InvalidCommand<>("usage: CLEAR");
}

return new ClearCommand<>();
}

private CacheCommand<K, V> parseMetrics(String[] parts) {
if (parts.length != ONE) {
return new InvalidCommand<>("usage: METRICS");
}

return new MetricsCommand<>();
}
}
23 changes: 23 additions & 0 deletions src/main/java/org/cache/protocol/CommandProcessor.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
package org.cache.protocol;

import org.cache.core.Cache;
import org.cache.protocol.codec.Codec;
import org.cache.protocol.commands.CacheCommand;

public class CommandProcessor<K, V> {

private final Cache<K, V> cache;
private final CommandParser<K, V> parser;
private final Codec<V> valueCodec;

public CommandProcessor(Cache<K, V> cache, CommandParser<K, V> parser, Codec<V> valueCodec) {
this.cache = cache;
this.parser = parser;
this.valueCodec = valueCodec;
}

public String process(String rawCommand) {
CacheCommand<K, V> command = parser.parse(rawCommand);
return command.process(cache, valueCodec);
}
}
8 changes: 8 additions & 0 deletions src/main/java/org/cache/protocol/codec/Codec.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
package org.cache.protocol.codec;

public interface Codec<T> {

T decode(String value);

String encode(T value);
}
14 changes: 14 additions & 0 deletions src/main/java/org/cache/protocol/codec/StringCodec.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
package org.cache.protocol.codec;

public class StringCodec implements Codec<String> {

@Override
public String decode(String value) {
return value;
}

@Override
public String encode(String value) {
return value;
}
}
9 changes: 9 additions & 0 deletions src/main/java/org/cache/protocol/commands/CacheCommand.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
package org.cache.protocol.commands;

import org.cache.core.Cache;
import org.cache.protocol.codec.Codec;

public interface CacheCommand<K, V> {

String process(Cache<K, V> cache, Codec<V> valueCodec);
}
15 changes: 15 additions & 0 deletions src/main/java/org/cache/protocol/commands/ClearCommand.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
package org.cache.protocol.commands;

import org.cache.core.Cache;
import org.cache.protocol.codec.Codec;

import static org.cache.protocol.commands.ResponseConstants.OK;

public class ClearCommand<K, V> implements CacheCommand<K, V> {

@Override
public String process(Cache<K, V> cache, Codec<V> valueCodec) {
cache.clear();
return OK.name();
}
}
11 changes: 11 additions & 0 deletions src/main/java/org/cache/protocol/commands/CommandType.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
package org.cache.protocol.commands;

public enum CommandType {
PUT,
GET,
DELETE,
SIZE,
CLEAR,
METRICS,
UNKNOWN
}
21 changes: 21 additions & 0 deletions src/main/java/org/cache/protocol/commands/DeleteCommand.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
package org.cache.protocol.commands;

import org.cache.core.Cache;
import org.cache.protocol.codec.Codec;

import static org.cache.protocol.commands.ResponseConstants.OK;

public class DeleteCommand<K, V> implements CacheCommand<K, V> {

private final K key;

public DeleteCommand(K key) {
this.key = key;
}

@Override
public String process(Cache<K, V> cache, Codec<V> valueCodec) {
cache.delete(key);
return OK.name();
}
}
23 changes: 23 additions & 0 deletions src/main/java/org/cache/protocol/commands/GetCommand.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
package org.cache.protocol.commands;

import org.cache.core.Cache;
import org.cache.protocol.codec.Codec;

import static org.cache.protocol.commands.ResponseConstants.NOT_FOUND;
import static org.cache.protocol.commands.ResponseConstants.VALUE;

public class GetCommand<K, V> implements CacheCommand<K, V> {

private final K key;

public GetCommand(K key) {
this.key = key;
}

@Override
public String process(Cache<K, V> cache, Codec<V> valueCodec) {
return cache.get(key)
.map(value -> VALUE.name() + " " + valueCodec.encode(value))
.orElse(NOT_FOUND.name());
}
}
Loading