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
5 changes: 5 additions & 0 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@ plugins {
id 'java'
id 'application'
id 'checkstyle'
id 'org.springframework.boot' version '3.3.4'
id 'io.spring.dependency-management' version '1.1.6'
}

group = 'org.example'
Expand All @@ -16,6 +18,9 @@ dependencies {
testImplementation 'org.junit.jupiter:junit-jupiter'
testImplementation 'org.mockito:mockito-junit-jupiter:5.12.0'
testRuntimeOnly 'org.junit.platform:junit-platform-launcher'

implementation 'org.springframework.boot:spring-boot-starter'
implementation 'org.springframework.boot:spring-boot-starter-web'
}

application {
Expand Down
74 changes: 52 additions & 22 deletions src/main/java/org/cache/Main.java
Original file line number Diff line number Diff line change
@@ -1,43 +1,73 @@
package org.cache;

import org.cache.core.Cache;
import org.cache.core.CacheService;
import org.cache.core.LocalCache;
import org.cache.core.ValueType;
import org.cache.eviction.LruEvictionPolicy;
import org.cache.network.connection.ClientConnectionHandler;
import org.cache.network.TcpCacheServer;
import org.cache.protocol.CommandParser;
import org.cache.network.connection.ClientConnectionHandler;
import org.cache.protocol.CommandProcessor;
import org.cache.protocol.codec.ListValueCodec;
import org.cache.protocol.codec.StringKeyCodec;
import org.cache.protocol.codec.StringValueCodec;
import org.cache.protocol.codec.ValueCodecRegistry;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;

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

@SpringBootApplication
public class Main {
private static final int SERVER_THREAD_COUNT = 16;
private static final int DEFAULT_CACHE_CAPACITY = 1_000;

public static void main(String[] args) {
SpringApplication.run(Main.class, args);
}

@Bean(destroyMethod = "close")
public LocalCache<String> cache() {
return new LocalCache<>(DEFAULT_CACHE_CAPACITY, new LruEvictionPolicy<>());
}

@Bean
public ValueCodecRegistry valueCodecs() {
return new ValueCodecRegistry()
.register(ValueType.STRING, new StringValueCodec())
.register(ValueType.LIST, new ListValueCodec());
}

@Bean
public CacheService<String> cacheService(Cache<String> cache, ValueCodecRegistry valueCodecs) {
return new CacheService<>(cache, valueCodecs);
}

@Bean
public CommandProcessor<String> commandProcessor(CacheService<String> cacheService) {
return new CommandProcessor<>(new StringKeyCodec(), cacheService);
}

@Bean(destroyMethod = "shutdownNow")
public ExecutorService tcpClientExecutor() {
return Executors.newFixedThreadPool(SERVER_THREAD_COUNT);
}

public static void main(String[] args) throws IOException {
int port = args.length > 0 ? Integer.parseInt(args[0]) : 2020;
var cache = new LocalCache<String>(1_000, new LruEvictionPolicy<>());
var keyCodec = new StringKeyCodec();
var valueCodecs = new ValueCodecRegistry();
valueCodecs.register(ValueType.STRING, new StringValueCodec()).register(ValueType.LIST, new ListValueCodec());
var commandParser = new CommandParser<>(keyCodec);
var commandProcessor = new CommandProcessor<>(cache, commandParser, valueCodecs);


try (cache;
var serverSocket = new ServerSocket(port);
var executor = Executors.newFixedThreadPool(SERVER_THREAD_COUNT);
var server = new TcpCacheServer(
serverSocket,
executor,
socket -> new ClientConnectionHandler(socket, commandProcessor)
)) {
server.start();
}
@Bean
public TcpCacheServer tcpCacheServer(
@Value("${cache.tcp.port:2020}") int port,
ExecutorService tcpClientExecutor,
CommandProcessor<String> commandProcessor
) throws IOException {
return new TcpCacheServer(
new ServerSocket(port),
tcpClientExecutor,
socket -> new ClientConnectionHandler(socket, commandProcessor)
);
}
}
6 changes: 3 additions & 3 deletions src/main/java/org/cache/client/TcpCacheClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -4,16 +4,16 @@
import org.cache.core.metrics.Snapshot;
import org.cache.network.connection.RespConnection;
import org.cache.protocol.codec.KeyCodec;
import org.cache.protocol.commands.ResponseConstants;
import org.cache.protocol.handlers.ResponseConstants;

import java.io.IOException;
import java.net.InetSocketAddress;
import java.net.Socket;
import java.util.List;
import java.util.Optional;

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

public class TcpCacheClient<K, V> implements CacheClient<K, V>, AutoCloseable {

Expand Down
108 changes: 108 additions & 0 deletions src/main/java/org/cache/core/CacheService.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
package org.cache.core;

import org.cache.core.metrics.Snapshot;
import org.cache.protocol.codec.ValueCodec;
import org.cache.protocol.codec.ValueCodecRegistry;
import org.cache.protocol.handlers.WrongValueTypeException;

import java.util.ArrayList;
import java.util.List;
import java.util.Optional;

public class CacheService<K> {

private final Cache<K> cache;
private final ValueCodecRegistry valueCodecs;

public CacheService(Cache<K> cache, ValueCodecRegistry valueCodecs) {
this.cache = cache;
this.valueCodecs = valueCodecs;
}

@SuppressWarnings("unchecked")
public void putString(K key, String value, long ttlMillis) {
ValueCodec<String> codec = (ValueCodec<String>) valueCodecs.get(ValueType.STRING);
cache.put(key, codec.encode(value), ValueType.STRING, ttlMillis);
}

public Optional<String> getString(K key) {
Optional<CacheEntry> entry = cache.get(key);

if (entry.isEmpty()) {
return Optional.empty();
}

CacheEntry cacheEntry = entry.get();
if (cacheEntry.getType() != ValueType.STRING) {
throw new WrongValueTypeException(ValueType.STRING, cacheEntry.getType());
}

return Optional.of(valueCodecs.get(ValueType.STRING).toString(cacheEntry.getValue()));
}

@SuppressWarnings("unchecked")
public void push(K key, String value) {
Optional<CacheEntry> existingList = cache.get(key);
ValueCodec<List<String>> listCodec = (ValueCodec<List<String>>) valueCodecs.get(ValueType.LIST);

List<String> list;

if (existingList.isEmpty()) {
list = new ArrayList<>();
} else {
CacheEntry entry = existingList.get();
if (entry.getType() != ValueType.LIST) {
throw new WrongValueTypeException(ValueType.LIST, entry.getType());
}

list = listCodec.decode(entry.getValue());
}

list.add(value);
cache.put(key, listCodec.encode(list), ValueType.LIST, 0);
}

@SuppressWarnings("unchecked")
public Optional<List<String>> lrange(K key, int from, int to) {
if (from < 0 || to < from) {
throw new IllegalArgumentException("from must be >= 0 and to must be >= from");
}

Optional<CacheEntry> entry = cache.get(key);

if (entry.isEmpty()) {
return Optional.empty();
}

CacheEntry cacheEntry = entry.get();
if (cacheEntry.getType() != ValueType.LIST) {
throw new WrongValueTypeException(ValueType.LIST, cacheEntry.getType());
}

ValueCodec<List<String>> listCodec = (ValueCodec<List<String>>) valueCodecs.get(ValueType.LIST);
List<String> list = listCodec.decode(cacheEntry.getValue());

if (from >= list.size()) {
return Optional.of(List.of());
}

int boundedTo = Math.min(to, list.size());
return Optional.of(List.copyOf(list.subList(from, boundedTo)));
}

public void delete(K key) {
cache.delete(key);
}

public int size() {
return cache.size();
}

public void clear() {
cache.clear();
}

public Snapshot metrics() {
return cache.metrics();
}
}
2 changes: 1 addition & 1 deletion src/main/java/org/cache/network/TcpCacheServer.java
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ public class TcpCacheServer implements AutoCloseable {

private final ExecutorService executor;
private final ClientConnectionHandlerFactory handlerFactory;
private ServerSocket serverSocket;
private final ServerSocket serverSocket;

public TcpCacheServer(
ServerSocket serverSocket,
Expand Down
61 changes: 61 additions & 0 deletions src/main/java/org/cache/network/TcpCacheServerLifecycle.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
package org.cache.network;

import org.springframework.context.SmartLifecycle;
import org.springframework.stereotype.Component;

import java.io.IOException;

@Component
public class TcpCacheServerLifecycle implements SmartLifecycle {

private final TcpCacheServer server;
private volatile boolean running;
private Thread serverThread;

public TcpCacheServerLifecycle(TcpCacheServer server) {
this.server = server;
}

@Override
public void start() {
if (running) {
return;
}

running = true;
serverThread = new Thread(this::runServer, "tcp-cache-server");
serverThread.start();
}

@Override
public void stop() {
running = false;

try {
server.close();
} catch (IOException exception) {
System.err.println("Failed to stop TCP cache server: " + exception.getMessage());
}

if (serverThread != null) {
serverThread.interrupt();
}
}

@Override
public boolean isRunning() {
return running;
}

private void runServer() {
try {
server.start();
} catch (IOException exception) {
if (running) {
System.err.println("TCP cache server stopped unexpectedly: " + exception.getMessage());
}
} finally {
running = false;
}
}
}
10 changes: 5 additions & 5 deletions src/main/java/org/cache/network/connection/RespConnection.java
Original file line number Diff line number Diff line change
Expand Up @@ -23,11 +23,11 @@
import static org.cache.protocol.RegexConstants.KEY_VALUE_SEPARATOR;
import static org.cache.protocol.RegexConstants.SPACE;
import static org.cache.protocol.RegexConstants.WHITESPACE;
import static org.cache.protocol.commands.ResponseConstants.ERROR;
import static org.cache.protocol.commands.ResponseConstants.LIST;
import static org.cache.protocol.commands.ResponseConstants.METRICS;
import static org.cache.protocol.commands.ResponseConstants.SIZE;
import static org.cache.protocol.commands.ResponseConstants.VALUE;
import static org.cache.protocol.handlers.ResponseConstants.ERROR;
import static org.cache.protocol.handlers.ResponseConstants.LIST;
import static org.cache.protocol.handlers.ResponseConstants.METRICS;
import static org.cache.protocol.handlers.ResponseConstants.SIZE;
import static org.cache.protocol.handlers.ResponseConstants.VALUE;

public class RespConnection implements ProtocolConnection, AutoCloseable {

Expand Down
Loading
Loading