From 222147f6d8544a335884138ac62dd9bfe6280f1f Mon Sep 17 00:00:00 2001 From: ferid333 <135500346+ferid333@users.noreply.github.com> Date: Wed, 22 Jul 2026 00:54:09 +1000 Subject: [PATCH] Phase-4: Adding TCP server, Command parser and processor --- src/main/java/org/cache/Main.java | 33 +++--- src/main/java/org/cache/core/Cache.java | 4 + src/main/java/org/cache/core/LocalCache.java | 1 + .../network/ClientConnectionHandler.java | 40 +++++++ .../org/cache/network/TcpCacheServer.java | 44 ++++++++ .../org/cache/protocol/CommandParser.java | 102 ++++++++++++++++++ .../org/cache/protocol/CommandProcessor.java | 23 ++++ .../java/org/cache/protocol/codec/Codec.java | 8 ++ .../org/cache/protocol/codec/StringCodec.java | 14 +++ .../cache/protocol/commands/CacheCommand.java | 9 ++ .../cache/protocol/commands/ClearCommand.java | 15 +++ .../cache/protocol/commands/CommandType.java | 11 ++ .../protocol/commands/DeleteCommand.java | 21 ++++ .../cache/protocol/commands/GetCommand.java | 23 ++++ .../protocol/commands/InvalidCommand.java | 20 ++++ .../protocol/commands/MetricsCommand.java | 20 ++++ .../cache/protocol/commands/PutCommand.java | 29 +++++ .../protocol/commands/ResponseConstants.java | 11 ++ .../cache/protocol/commands/SizeCommand.java | 14 +++ .../protocol/commands/UnknownCommand.java | 14 +++ 20 files changed, 438 insertions(+), 18 deletions(-) create mode 100644 src/main/java/org/cache/network/ClientConnectionHandler.java create mode 100644 src/main/java/org/cache/network/TcpCacheServer.java create mode 100644 src/main/java/org/cache/protocol/CommandParser.java create mode 100644 src/main/java/org/cache/protocol/CommandProcessor.java create mode 100644 src/main/java/org/cache/protocol/codec/Codec.java create mode 100644 src/main/java/org/cache/protocol/codec/StringCodec.java create mode 100644 src/main/java/org/cache/protocol/commands/CacheCommand.java create mode 100644 src/main/java/org/cache/protocol/commands/ClearCommand.java create mode 100644 src/main/java/org/cache/protocol/commands/CommandType.java create mode 100644 src/main/java/org/cache/protocol/commands/DeleteCommand.java create mode 100644 src/main/java/org/cache/protocol/commands/GetCommand.java create mode 100644 src/main/java/org/cache/protocol/commands/InvalidCommand.java create mode 100644 src/main/java/org/cache/protocol/commands/MetricsCommand.java create mode 100644 src/main/java/org/cache/protocol/commands/PutCommand.java create mode 100644 src/main/java/org/cache/protocol/commands/ResponseConstants.java create mode 100644 src/main/java/org/cache/protocol/commands/SizeCommand.java create mode 100644 src/main/java/org/cache/protocol/commands/UnknownCommand.java diff --git a/src/main/java/org/cache/Main.java b/src/main/java/org/cache/Main.java index f13ccb3..d32c17e 100644 --- a/src/main/java/org/cache/Main.java +++ b/src/main/java/org/cache/Main.java @@ -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 lruEvictionPolicy = new LruEvictionPolicy<>(); - LocalCache 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(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(); + } } } \ No newline at end of file diff --git a/src/main/java/org/cache/core/Cache.java b/src/main/java/org/cache/core/Cache.java index f60f2c8..845b83d 100644 --- a/src/main/java/org/cache/core/Cache.java +++ b/src/main/java/org/cache/core/Cache.java @@ -1,5 +1,7 @@ package org.cache.core; +import org.cache.core.metrics.Snapshot; + import java.util.Optional; public interface Cache { @@ -13,4 +15,6 @@ public interface Cache { int size(); void clear(); + + Snapshot metrics(); } diff --git a/src/main/java/org/cache/core/LocalCache.java b/src/main/java/org/cache/core/LocalCache.java index 0be6cb6..2129abd 100644 --- a/src/main/java/org/cache/core/LocalCache.java +++ b/src/main/java/org/cache/core/LocalCache.java @@ -163,6 +163,7 @@ public void close() { cleanupScheduler.shutdownNow(); } + @Override public Snapshot metrics() { return metrics.snapshot(); } diff --git a/src/main/java/org/cache/network/ClientConnectionHandler.java b/src/main/java/org/cache/network/ClientConnectionHandler.java new file mode 100644 index 0000000..b48a85b --- /dev/null +++ b/src/main/java/org/cache/network/ClientConnectionHandler.java @@ -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()); + } + } +} diff --git a/src/main/java/org/cache/network/TcpCacheServer.java b/src/main/java/org/cache/network/TcpCacheServer.java new file mode 100644 index 0000000..c805f83 --- /dev/null +++ b/src/main/java/org/cache/network/TcpCacheServer.java @@ -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(); + } + } +} diff --git a/src/main/java/org/cache/protocol/CommandParser.java b/src/main/java/org/cache/protocol/CommandParser.java new file mode 100644 index 0000000..239f6f5 --- /dev/null +++ b/src/main/java/org/cache/protocol/CommandParser.java @@ -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 { + + private final Codec keyCodec; + private final Codec 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 keyCodec, Codec valueCodec) { + this.keyCodec = keyCodec; + this.valueCodec = valueCodec; + } + + public CacheCommand 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 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 parseGet(String[] parts) { + if (parts.length != TWO) { + return new InvalidCommand<>("usage: GET key"); + } + + return new GetCommand<>(keyCodec.decode(parts[1])); + } + + private CacheCommand parseDelete(String[] parts) { + if (parts.length != TWO) { + return new InvalidCommand<>("usage: DELETE key"); + } + + return new DeleteCommand<>(keyCodec.decode(parts[1])); + } + + private CacheCommand parseSize(String[] parts) { + if (parts.length != ONE) { + return new InvalidCommand<>("usage: SIZE"); + } + + return new SizeCommand<>(); + } + + private CacheCommand parseClear(String[] parts) { + if (parts.length != ONE) { + return new InvalidCommand<>("usage: CLEAR"); + } + + return new ClearCommand<>(); + } + + private CacheCommand parseMetrics(String[] parts) { + if (parts.length != ONE) { + return new InvalidCommand<>("usage: METRICS"); + } + + return new MetricsCommand<>(); + } +} diff --git a/src/main/java/org/cache/protocol/CommandProcessor.java b/src/main/java/org/cache/protocol/CommandProcessor.java new file mode 100644 index 0000000..0bdce98 --- /dev/null +++ b/src/main/java/org/cache/protocol/CommandProcessor.java @@ -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 { + + private final Cache cache; + private final CommandParser parser; + private final Codec valueCodec; + + public CommandProcessor(Cache cache, CommandParser parser, Codec valueCodec) { + this.cache = cache; + this.parser = parser; + this.valueCodec = valueCodec; + } + + public String process(String rawCommand) { + CacheCommand command = parser.parse(rawCommand); + return command.process(cache, valueCodec); + } +} diff --git a/src/main/java/org/cache/protocol/codec/Codec.java b/src/main/java/org/cache/protocol/codec/Codec.java new file mode 100644 index 0000000..e6fe5f9 --- /dev/null +++ b/src/main/java/org/cache/protocol/codec/Codec.java @@ -0,0 +1,8 @@ +package org.cache.protocol.codec; + +public interface Codec { + + T decode(String value); + + String encode(T value); +} diff --git a/src/main/java/org/cache/protocol/codec/StringCodec.java b/src/main/java/org/cache/protocol/codec/StringCodec.java new file mode 100644 index 0000000..c75d45a --- /dev/null +++ b/src/main/java/org/cache/protocol/codec/StringCodec.java @@ -0,0 +1,14 @@ +package org.cache.protocol.codec; + +public class StringCodec implements Codec { + + @Override + public String decode(String value) { + return value; + } + + @Override + public String encode(String value) { + return value; + } +} diff --git a/src/main/java/org/cache/protocol/commands/CacheCommand.java b/src/main/java/org/cache/protocol/commands/CacheCommand.java new file mode 100644 index 0000000..8e26a69 --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/CacheCommand.java @@ -0,0 +1,9 @@ +package org.cache.protocol.commands; + +import org.cache.core.Cache; +import org.cache.protocol.codec.Codec; + +public interface CacheCommand { + + String process(Cache cache, Codec valueCodec); +} diff --git a/src/main/java/org/cache/protocol/commands/ClearCommand.java b/src/main/java/org/cache/protocol/commands/ClearCommand.java new file mode 100644 index 0000000..2e352ac --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/ClearCommand.java @@ -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 implements CacheCommand { + + @Override + public String process(Cache cache, Codec valueCodec) { + cache.clear(); + return OK.name(); + } +} diff --git a/src/main/java/org/cache/protocol/commands/CommandType.java b/src/main/java/org/cache/protocol/commands/CommandType.java new file mode 100644 index 0000000..a5c2440 --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/CommandType.java @@ -0,0 +1,11 @@ +package org.cache.protocol.commands; + +public enum CommandType { + PUT, + GET, + DELETE, + SIZE, + CLEAR, + METRICS, + UNKNOWN +} diff --git a/src/main/java/org/cache/protocol/commands/DeleteCommand.java b/src/main/java/org/cache/protocol/commands/DeleteCommand.java new file mode 100644 index 0000000..f63cef7 --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/DeleteCommand.java @@ -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 implements CacheCommand { + + private final K key; + + public DeleteCommand(K key) { + this.key = key; + } + + @Override + public String process(Cache cache, Codec valueCodec) { + cache.delete(key); + return OK.name(); + } +} diff --git a/src/main/java/org/cache/protocol/commands/GetCommand.java b/src/main/java/org/cache/protocol/commands/GetCommand.java new file mode 100644 index 0000000..917417e --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/GetCommand.java @@ -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 implements CacheCommand { + + private final K key; + + public GetCommand(K key) { + this.key = key; + } + + @Override + public String process(Cache cache, Codec valueCodec) { + return cache.get(key) + .map(value -> VALUE.name() + " " + valueCodec.encode(value)) + .orElse(NOT_FOUND.name()); + } +} diff --git a/src/main/java/org/cache/protocol/commands/InvalidCommand.java b/src/main/java/org/cache/protocol/commands/InvalidCommand.java new file mode 100644 index 0000000..8229d55 --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/InvalidCommand.java @@ -0,0 +1,20 @@ +package org.cache.protocol.commands; + +import org.cache.core.Cache; +import org.cache.protocol.codec.Codec; + +import static org.cache.protocol.commands.ResponseConstants.ERROR; + +public class InvalidCommand implements CacheCommand { + + private final String message; + + public InvalidCommand(String message) { + this.message = message; + } + + @Override + public String process(Cache cache, Codec valueCodec) { + return ERROR.name() + " " + message; + } +} diff --git a/src/main/java/org/cache/protocol/commands/MetricsCommand.java b/src/main/java/org/cache/protocol/commands/MetricsCommand.java new file mode 100644 index 0000000..8509ff7 --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/MetricsCommand.java @@ -0,0 +1,20 @@ +package org.cache.protocol.commands; + +import org.cache.core.Cache; +import org.cache.core.metrics.Snapshot; +import org.cache.protocol.codec.Codec; + +import static org.cache.protocol.commands.ResponseConstants.METRICS; + +public class MetricsCommand implements CacheCommand { + + @Override + public String process(Cache cache, Codec valueCodec) { + Snapshot metrics = cache.metrics(); + return METRICS.name() + " hits=" + metrics.getHits() + + " misses=" + metrics.getMisses() + + " evictions=" + metrics.getEvictions() + + " expirations=" + metrics.getExpirations() + + " hitRate=" + metrics.getHitRate(); + } +} diff --git a/src/main/java/org/cache/protocol/commands/PutCommand.java b/src/main/java/org/cache/protocol/commands/PutCommand.java new file mode 100644 index 0000000..a3b7d41 --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/PutCommand.java @@ -0,0 +1,29 @@ +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 PutCommand implements CacheCommand { + + private final K key; + private final V value; + private final long ttlMillis; + + public PutCommand(K key, V value, long ttlMillis) { + this.key = key; + this.value = value; + this.ttlMillis = ttlMillis; + } + + @Override + public String process(Cache cache, Codec valueCodec) { + cache.put(key, value, ttlMillis); + return OK.name(); + } + + public V getValue() { + return value; + } +} diff --git a/src/main/java/org/cache/protocol/commands/ResponseConstants.java b/src/main/java/org/cache/protocol/commands/ResponseConstants.java new file mode 100644 index 0000000..192b3f7 --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/ResponseConstants.java @@ -0,0 +1,11 @@ +package org.cache.protocol.commands; + +public enum ResponseConstants { + + OK, + ERROR, + NOT_FOUND, + VALUE, + SIZE, + METRICS; +} diff --git a/src/main/java/org/cache/protocol/commands/SizeCommand.java b/src/main/java/org/cache/protocol/commands/SizeCommand.java new file mode 100644 index 0000000..7b06fb5 --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/SizeCommand.java @@ -0,0 +1,14 @@ +package org.cache.protocol.commands; + +import org.cache.core.Cache; +import org.cache.protocol.codec.Codec; + +import static org.cache.protocol.commands.ResponseConstants.SIZE; + +public class SizeCommand implements CacheCommand { + + @Override + public String process(Cache cache, Codec valueCodec) { + return SIZE.name() + " " + cache.size(); + } +} diff --git a/src/main/java/org/cache/protocol/commands/UnknownCommand.java b/src/main/java/org/cache/protocol/commands/UnknownCommand.java new file mode 100644 index 0000000..f0e82cd --- /dev/null +++ b/src/main/java/org/cache/protocol/commands/UnknownCommand.java @@ -0,0 +1,14 @@ +package org.cache.protocol.commands; + +import org.cache.core.Cache; +import org.cache.protocol.codec.Codec; + +import static org.cache.protocol.commands.ResponseConstants.ERROR; + +public class UnknownCommand implements CacheCommand { + + @Override + public String process(Cache cache, Codec valueCodec) { + return ERROR.name() + " unknown command"; + } +}