Skip to content
Open
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
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
**/*.iml
/.idea
**/target/
out/
15 changes: 15 additions & 0 deletions L7_threads/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -94,5 +94,20 @@
<version>1.7.5</version>
</dependency>

<!-- https://mvnrepository.com/artifact/com.google.code.gson/gson -->
<dependency>
<groupId>com.google.code.gson</groupId>
<artifactId>gson</artifactId>
<version>2.8.2</version>
</dependency>

<!-- https://mvnrepository.com/artifact/mysql/mysql-connector-java -->
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>5.1.6</version>
</dependency>


</dependencies>
</project>
52 changes: 52 additions & 0 deletions L7_threads/src/main/java/ru/track/homework/AdminThread.java
Original file line number Diff line number Diff line change
@@ -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<ClientThread> clients;

AdminThread(@NotNull List<ClientThread> 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;
}
}
}

}
}
}
}
47 changes: 47 additions & 0 deletions L7_threads/src/main/java/ru/track/homework/BroadcosterThread.java
Original file line number Diff line number Diff line change
@@ -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<ClientThread> clients;

private static long lastUpdateTS = System.currentTimeMillis();

BroadcosterThread(@NotNull List<ClientThread> clients, @NotNull ConversationService messages) {
BroadcosterThread.messages = messages;
BroadcosterThread.clients = clients;
}

public void run()
{
try {
while (!isInterrupted()) {
long currentTS = System.currentTimeMillis();
List<Message> 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);
}
}
}
102 changes: 102 additions & 0 deletions L7_threads/src/main/java/ru/track/homework/Client.java
Original file line number Diff line number Diff line change
@@ -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();
}
}
88 changes: 88 additions & 0 deletions L7_threads/src/main/java/ru/track/homework/ClientThread.java
Original file line number Diff line number Diff line change
@@ -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<ClientThread> clients;

private final Socket socket;

private InputStream in;
private OutputStream out;

private static Gson gson = new Gson();

ClientThread(@NotNull Socket socket, @NotNull List<ClientThread> 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();
}
}
Original file line number Diff line number Diff line change
@@ -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<Message> getHistory(long from, long to, long limit);


/**
* Вернуть все сообщения от определенного пользователя
*
* 4 балла
*/
List<Message> getByUser(String username, long limit);

}
Loading