diff --git a/.gitignore b/.gitignore index e32026c..1ee3778 100644 --- a/.gitignore +++ b/.gitignore @@ -1,3 +1,4 @@ **/*.iml /.idea **/target/ +out/ diff --git a/L7_threads/pom.xml b/L7_threads/pom.xml index c0d9e38..699959b 100644 --- a/L7_threads/pom.xml +++ b/L7_threads/pom.xml @@ -94,5 +94,20 @@ 1.7.5 + + + com.google.code.gson + gson + 2.8.2 + + + + + mysql + mysql-connector-java + 5.1.6 + + + \ No newline at end of file diff --git a/L7_threads/src/main/java/ru/track/homework/AdminThread.java b/L7_threads/src/main/java/ru/track/homework/AdminThread.java new file mode 100644 index 0000000..eb9785d --- /dev/null +++ b/L7_threads/src/main/java/ru/track/homework/AdminThread.java @@ -0,0 +1,52 @@ +package ru.track.homework; + +import org.jetbrains.annotations.NotNull; + +import java.util.List; +import java.util.Scanner; + +public class AdminThread extends Thread { + private List clients; + + AdminThread(@NotNull List clients) { + this.clients = clients; + } + + public void run() { + Scanner scanner = new Scanner(System.in); + + while (!isInterrupted()) { + String line = scanner.next(); + + if (line.equals("list")) { + synchronized (clients){ + System.out.println("Client list:"); + + if (clients.isEmpty()) { + System.out.println("empty"); + continue; + } + + for (ClientThread client: clients){ + if (!client.isInterrupted()) { + System.out.println(String.format("%s[%s]", client.getName(), client.getUsername())); + } + } + } + } + else if (line.equals("drop")) { + long clientId = scanner.nextLong(); + + synchronized (clients){ + for (ClientThread client: clients){ + if (!client.isInterrupted() && (client.getClientID() == clientId)) { + client.interrupt(); + break; + } + } + } + + } + } + } +} \ No newline at end of file diff --git a/L7_threads/src/main/java/ru/track/homework/BroadcosterThread.java b/L7_threads/src/main/java/ru/track/homework/BroadcosterThread.java new file mode 100644 index 0000000..7f2628a --- /dev/null +++ b/L7_threads/src/main/java/ru/track/homework/BroadcosterThread.java @@ -0,0 +1,47 @@ +package ru.track.homework; + +import org.jetbrains.annotations.NotNull; + +import java.io.IOException; +import java.util.List; + +public class BroadcosterThread extends Thread { + private static ConversationService messages; + private static List clients; + + private static long lastUpdateTS = System.currentTimeMillis(); + + BroadcosterThread(@NotNull List clients, @NotNull ConversationService messages) { + BroadcosterThread.messages = messages; + BroadcosterThread.clients = clients; + } + + public void run() + { + try { + while (!isInterrupted()) { + long currentTS = System.currentTimeMillis(); + List newMessages = messages.getHistory(lastUpdateTS, currentTS, 20); + lastUpdateTS = currentTS; + + for (Message newMessage: newMessages) { + synchronized (clients) { + for (ClientThread client: clients) { + if (!client.isInterrupted() && !client.getUsername().equals(newMessage.getUsername())) { + try { + client.send(newMessage); + } catch (IOException e) { + client.interrupt(); + } + } + } + } + } + + Thread.sleep(800); + } + } catch (InterruptedException e) { + //throw new RuntimeException(e); + } + } +} diff --git a/L7_threads/src/main/java/ru/track/homework/Client.java b/L7_threads/src/main/java/ru/track/homework/Client.java new file mode 100644 index 0000000..205d308 --- /dev/null +++ b/L7_threads/src/main/java/ru/track/homework/Client.java @@ -0,0 +1,102 @@ +package ru.track.homework; + +import com.google.gson.Gson; +import org.apache.commons.io.IOUtils; +import org.jetbrains.annotations.NotNull; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.net.Socket; +import java.util.Scanner; + +public class Client { + private final int port; + private final String host; + + private String username; + + private Thread writer; + private Thread reader; + + private static Gson gson = new Gson(); + + public Client(@NotNull String host, int port) { + this.port = port; + this.host = host; + } + + private void connect() + { + try { + Socket socket = new Socket(host, port); + final Socket finalSocket = socket; + + writer = new Thread(() -> { writeToSocket(finalSocket); }); + writer.start(); + + reader = new Thread(() -> { readFromSocket(finalSocket); }); + reader.start(); + + } catch (IOException e) { + //throw new RuntimeException(e); + } + } + + private void readFromSocket(@NotNull Socket socket) { + try { + final InputStream in = socket.getInputStream(); + final byte[] buffer = new byte[4096]; + + for (int nRead = 0; nRead != -1; nRead = in.read(buffer)) { + if (nRead == 0) continue; + + Message message = gson.fromJson(new String(buffer, 0, nRead), Message.class); + System.out.println(String.format("%s> %s", message.getUsername(), message.getText())); + } + socket.shutdownInput(); + + } catch (IOException e) { + //throw new RuntimeException(e); + } finally { + IOUtils.closeQuietly(socket); + } + } + + private void writeToSocket(@NotNull Socket socket) { + try{ + final OutputStream out = socket.getOutputStream(); + Scanner scanner = new Scanner(System.in); + + System.out.print("Введите ваш username: "); + username = scanner.nextLine(); + + out.write(username.getBytes()); + + while (true) { + String line = scanner.nextLine(); + if ("exit".equals(line)) { + System.out.println("bye!"); + break; + } + + Message message = new Message(username, line); + String result = gson.toJson(message); + + out.write(result.getBytes()); + out.flush(); + } + socket.shutdownOutput(); + + } catch (IOException e) { + //throw new RuntimeException(e); + } finally { + IOUtils.closeQuietly(socket); + } + } + + public static void main(String[] args) throws Exception { + Client client = new Client("localhost",9000); + client.connect(); + } +} diff --git a/L7_threads/src/main/java/ru/track/homework/ClientThread.java b/L7_threads/src/main/java/ru/track/homework/ClientThread.java new file mode 100644 index 0000000..7d82f62 --- /dev/null +++ b/L7_threads/src/main/java/ru/track/homework/ClientThread.java @@ -0,0 +1,88 @@ +package ru.track.homework; + +import com.google.gson.Gson; +import org.apache.commons.io.IOUtils; +import org.jetbrains.annotations.NotNull; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.net.InetAddress; +import java.net.ServerSocket; +import java.net.Socket; +import java.util.Date; +import java.util.List; +import java.util.Queue; +import java.util.concurrent.atomic.AtomicLong; + +public class ClientThread extends Thread { + private final long clientID; + private String username; + + private static AtomicCounter counter = new AtomicCounter(); + + private static ConversationService messages; + private static List clients; + + private final Socket socket; + + private InputStream in; + private OutputStream out; + + private static Gson gson = new Gson(); + + ClientThread(@NotNull Socket socket, @NotNull List clients, @NotNull ConversationService messages) { + this.clientID = counter.inc(); + this.socket = socket; + + ClientThread.clients = clients; + ClientThread.messages = messages; + + setName(String.format("Client[%d]@%s:%s", clientID, socket.getInetAddress(), socket.getPort())); + } + + public void run() { + try { + in = socket.getInputStream(); + out = socket.getOutputStream(); + + final byte[] buffer = new byte[4096]; + + int nUsername = in.read(buffer); + username = new String(buffer, 0, nUsername); + for (int nRead = 0; !isInterrupted() || (nRead != -1); nRead = in.read(buffer)) { + //while (in.available() == 0) { sleep(200); } + if (nRead == 0) continue; + + Message message = gson.fromJson(new String(buffer, 0, nRead), Message.class); + System.out.println(String.format("%s[%s]> %s", getName(), message.getUsername(), message.getText())); + messages.store(message); + } + + socket.shutdownInput(); + socket.shutdownOutput(); + + } catch (IOException e) { //| InterruptedException e + //throw new RuntimeException(e); + } finally { + clients.remove(this); + System.out.println(String.format("main> client disconnected: %s", getName())); + IOUtils.closeQuietly(socket); + } + } + + public void send(Message message) throws IOException { + out.write(gson.toJson(message).getBytes()); + out.flush(); + } + + public String getUsername() { return username; } + public long getClientID() { return clientID; } +} + +class AtomicCounter { + private AtomicLong val = new AtomicLong(0); + public long inc() { + return val.getAndIncrement(); + } +} diff --git a/L7_threads/src/main/java/ru/track/homework/ConversationService.java b/L7_threads/src/main/java/ru/track/homework/ConversationService.java new file mode 100644 index 0000000..4026f6b --- /dev/null +++ b/L7_threads/src/main/java/ru/track/homework/ConversationService.java @@ -0,0 +1,36 @@ +package ru.track.homework; +import java.util.List; + +public interface ConversationService { + /** + * В зависимости от Message.senderName нужно сохранять в разные базы + * + * @return вернуть ID, который был присвоен сообщению в базе (поле ID) + * + * + * 4 балла + */ + long store(Message msg); + + /** + * Получить историю сообщений за период времени. Важно учесть лимит, чтобы не свалить базы слишком большой выборкой. + * Ограничить LIMIT нужно именно при запросе в базу + * @param from - timestamp с какого времени + * @param to - timestamp до какого времени + * @param limit - максимальное кол-во ссобщений + * + * @return Список, отсротированный по timestamp + * + * 4 балла + */ + List getHistory(long from, long to, long limit); + + + /** + * Вернуть все сообщения от определенного пользователя + * + * 4 балла + */ + List getByUser(String username, long limit); + +} diff --git a/L7_threads/src/main/java/ru/track/homework/Database.java b/L7_threads/src/main/java/ru/track/homework/Database.java new file mode 100644 index 0000000..b889c2c --- /dev/null +++ b/L7_threads/src/main/java/ru/track/homework/Database.java @@ -0,0 +1,121 @@ +package ru.track.homework; + +import java.sql.*; +import java.util.ArrayList; +import java.util.List; + +public class Database implements ConversationService { + + private Connection[] connections = new Connection[3]; + + Database() { + try { + DriverManager.registerDriver((Driver) Class.forName("com.mysql.jdbc.Driver").newInstance()); + connections[0] = connect(1); + connections[1] = connect(2); + connections[2] = connect(3); + } catch (SQLException |IllegalAccessException | ClassNotFoundException | InstantiationException e) { + throw new RuntimeException(e); + } + } + + @Override + public synchronized long store(Message msg) { + Connection currentConnection = getConnection(msg.getUsername()); + + try { + String insert = "INSERT INTO messages (user_name, text, ts) VALUES (?, ?, ?)"; + PreparedStatement stmt = currentConnection.prepareStatement(insert, Statement.RETURN_GENERATED_KEYS); + stmt.setString(1, msg.getUsername()); + stmt.setString(2, msg.getText()); + stmt.setTimestamp(3, new java.sql.Timestamp(msg.getTimeStand())); + + int affectedRows = stmt.executeUpdate(); + + if (affectedRows == 0) { + throw new SQLException("Storing message failed, no rows affected."); + } + + try (ResultSet generatedKeys = stmt.getGeneratedKeys()) { + if (generatedKeys.next()) { + return generatedKeys.getLong(1); + } + else { + throw new SQLException("Storing message failed, no ID obtained."); + } + } + } catch (SQLException e) { + throw new RuntimeException(e); + } + } + + @Override + public synchronized List getHistory(long from, long to, long limit) { + List result = new ArrayList<>(); + + try { + for (Connection currentConnection : connections) { + String select = "SELECT * FROM messages WHERE ts >= ? AND ts < ? ORDER BY ts LIMIT ?"; + PreparedStatement stmt = currentConnection.prepareStatement(select); + stmt.setTimestamp(1, new java.sql.Timestamp(from)); + stmt.setTimestamp(2, new java.sql.Timestamp(to)); + stmt.setLong(3, limit); + + ResultSet resultSet = stmt.executeQuery(); + while (resultSet.next()) { + result.add(new Message(resultSet.getString(2), resultSet.getString(3), resultSet.getLong(4))); + } + } + + result.sort((a, b) -> (int) (a.getTimeStand() - b.getTimeStand())); + + } catch (SQLException e) { + throw new RuntimeException(e); + } + + return result; + } + + @Override + public synchronized List getByUser(String username, long limit) { + Connection currentConnection = getConnection(username); + List result = new ArrayList<>(); + + try { + String select = "SELECT * FROM messages WHERE user_name = '?' LIMIT ? ORDER BY ts"; + PreparedStatement stmt = currentConnection.prepareStatement(select); + stmt.setString(1, username); + stmt.setLong(2, limit); + + ResultSet resultSet = stmt.executeQuery(); + while(resultSet.next()) { + result.add(new Message(resultSet.getString(2), resultSet.getString(3), resultSet.getLong(4))); + } + + } catch (SQLException e) { + throw new RuntimeException(e); + } + + return result; + } + + private static Connection connect(int shardID) throws SQLException { + String url = String.format("jdbc:mysql://tdb-%d.trail5.net:3306/track17?user=track_student&password=7EsH.H6x", shardID); + return DriverManager.getConnection(url); + } + + private Connection getConnection(String username) { + char firstChar = Character.toUpperCase(username.charAt(0)); + + Connection currentConnection; + if (('A' <= firstChar) && (firstChar <= 'J')) + currentConnection = connections[0]; + else if (('K' <= firstChar) && (firstChar <= 'T')) + currentConnection = connections[1]; + else if (('U' <= firstChar) && (firstChar <= 'Z')) + currentConnection = connections[2]; + else currentConnection = null; + + return currentConnection; + } +} \ No newline at end of file diff --git a/L7_threads/src/main/java/ru/track/homework/Message.java b/L7_threads/src/main/java/ru/track/homework/Message.java new file mode 100644 index 0000000..116b87d --- /dev/null +++ b/L7_threads/src/main/java/ru/track/homework/Message.java @@ -0,0 +1,23 @@ +package ru.track.homework; + +import org.jetbrains.annotations.NotNull; + +class Message { + private String username; + private String text; + private long ts; + + Message(@NotNull String username, @NotNull String text) { + this(username, text, System.currentTimeMillis()); + } + + Message(@NotNull String username, @NotNull String text, long ts) { + this.username = username; + this.ts = ts; + this.text = text; + } + + public String getText() { return text; } + public String getUsername() { return username; } + public long getTimeStand() { return ts; } +} \ No newline at end of file diff --git a/L7_threads/src/main/java/ru/track/homework/Server.java b/L7_threads/src/main/java/ru/track/homework/Server.java new file mode 100644 index 0000000..576feef --- /dev/null +++ b/L7_threads/src/main/java/ru/track/homework/Server.java @@ -0,0 +1,69 @@ +package ru.track.homework; + +import org.apache.commons.io.IOUtils; +import org.jetbrains.annotations.NotNull; + +import java.io.IOException; +import java.net.InetAddress; +import java.net.ServerSocket; +import java.net.Socket; +import java.util.*; + +public class Server { + private final int port; + private final int backlog; + + private ConversationService messages = new Database(); + private List clients = Collections.synchronizedList(new ArrayList<>()); + + private BroadcosterThread broadcaster; + private AdminThread admin; + + + public Server(int port, int backlog) { + this.port = port; + this.backlog = backlog; + + admin = new AdminThread(clients); + broadcaster = new BroadcosterThread(clients, messages); + } + + public void serve() throws IOException { + ServerSocket serverSocket = null; + try { + final ServerSocket serverSocketFinal = new ServerSocket(port, backlog, InetAddress.getByName("localhost")); + serverSocket = serverSocketFinal; + + broadcaster.start(); + admin.start(); + + while (true){ + System.out.println("main> wait to accept new client..."); + handle(serverSocketFinal); + } + + } finally { + IOUtils.closeQuietly(serverSocket); + } + } + + private void handle(@NotNull ServerSocket serverSocket) { + try { + Socket socket = serverSocket.accept(); + + ClientThread client = new ClientThread(socket, clients, messages); + + clients.add(client); + client.start(); + + System.out.println(String.format("main> new client connected: %s", client.getName())); + } catch (IOException e) { + //throw new RuntimeException(e); + } + } + + public static void main(String[] args) throws Exception { + Server server = new Server(9000, 10); + server.serve(); + } +} diff --git a/OWNER.md b/OWNER.md new file mode 100644 index 0000000..8b2f6cb --- /dev/null +++ b/OWNER.md @@ -0,0 +1 @@ +Морев Максим Сергеевич