diff --git a/src/main/java/io/github/winroot33/ClientConnectionHandler.java b/src/main/java/io/github/winroot33/ClientConnectionHandler.java new file mode 100644 index 0000000..bfda7ac --- /dev/null +++ b/src/main/java/io/github/winroot33/ClientConnectionHandler.java @@ -0,0 +1,53 @@ +package io.github.winroot33; + +import lombok.RequiredArgsConstructor; + +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; + +/** + * Класс для обработки соединений с клиентами, Runnable для передачи в пул потоков + */ +@RequiredArgsConstructor +public class ClientConnectionHandler implements Runnable { + private final Socket clientSocket; + + @Override + public void run() { + handleConnection(); + } + + /** + * Обработка нового соединения + */ + private void handleConnection() { + try (Socket socket = this.clientSocket; + PrintWriter out = new PrintWriter(socket.getOutputStream(), true, StandardCharsets.UTF_8); + BufferedReader in = new BufferedReader( + new InputStreamReader(socket.getInputStream(), StandardCharsets.UTF_8))) { + + processClientMessages(in, out); + + } catch (IOException e) { + System.err.println("IO Error: " + e.getMessage()); + } + } + + /** + * Метод для отправки эхо сообщений пользователю + * + * @param in входной поток данных + * @param out выходной поток + */ + private void processClientMessages(BufferedReader in, PrintWriter out) throws IOException { + String inputLine; + while ((inputLine = in.readLine()) != null) { + out.println(inputLine); + System.out.printf("Thread: %s\tMessage sent: %s\n", Thread.currentThread().getName(), inputLine); + } + } +} diff --git a/src/main/java/io/github/winroot33/ConnectionExecutor.java b/src/main/java/io/github/winroot33/ConnectionExecutor.java new file mode 100644 index 0000000..9fb7927 --- /dev/null +++ b/src/main/java/io/github/winroot33/ConnectionExecutor.java @@ -0,0 +1,46 @@ +package io.github.winroot33; + +import java.net.Socket; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +/** + * Обертка над ExecutorService для удобного создания новых соединений с клиентами + */ +public class ConnectionExecutor { + private final ExecutorService threadPool; + + /** + * Создание пула потоков с указанным количеством + * + * @param maxConnections количество потоков для FixedThreadPool + */ + public ConnectionExecutor(int maxConnections) { + this.threadPool = Executors.newFixedThreadPool(maxConnections); + } + + /** + * Обработка подключения нового клиента, создание и отправка задачи в ExecutorService + * + * @param clientSocket сокет с соединением нового клиента + */ + public void executeConnection(Socket clientSocket) { + Runnable connectionHandler = new ClientConnectionHandler(clientSocket); + threadPool.execute(connectionHandler); + } + + /** + * Корректная остановка пула потоков + */ + public void shutdownGracefully() { + threadPool.shutdown(); + try { + if (!threadPool.awaitTermination(30, TimeUnit.SECONDS)) { + threadPool.shutdownNow(); + } + } catch (InterruptedException e) { + threadPool.shutdownNow(); + } + } +} diff --git a/src/main/java/io/github/winroot33/EchoServer.java b/src/main/java/io/github/winroot33/EchoServer.java new file mode 100644 index 0000000..ac06306 --- /dev/null +++ b/src/main/java/io/github/winroot33/EchoServer.java @@ -0,0 +1,41 @@ +package io.github.winroot33; + +import java.io.IOException; +import java.net.ServerSocket; +import java.net.Socket; + +/** + * Класс для запуска эхо сервера + */ +public class EchoServer { + + private final ConnectionExecutor connectionExecutor; + private final int port; + + /** + * Создание эхо сервера + * + * @param port номер порта для сервера + * @param maxConnections количество потоков для FixedThreadPool + */ + public EchoServer(int port, int maxConnections) { + this.port = port; + this.connectionExecutor = new ConnectionExecutor(maxConnections); + } + + /** + * Запуск эхо сервера + */ + public void start() throws IOException { + + try (ServerSocket serverSocket = new ServerSocket(port)) { + System.out.println("Server started"); + while (true) { + Socket clientSocket = serverSocket.accept(); + connectionExecutor.executeConnection(clientSocket); + } + } finally { + connectionExecutor.shutdownGracefully(); + } + } +} diff --git a/src/main/java/io/github/winroot33/Main.java b/src/main/java/io/github/winroot33/Main.java index 94eaa94..8fc7d2d 100644 --- a/src/main/java/io/github/winroot33/Main.java +++ b/src/main/java/io/github/winroot33/Main.java @@ -1,6 +1,10 @@ package io.github.winroot33; +import java.io.IOException; + public class Main { - public static void main(String[] args) { + public static void main(String[] args) throws IOException { + EchoServer server = new EchoServer(7, 10); + server.start(); } } \ No newline at end of file