diff --git a/OWNER.md b/OWNER.md
new file mode 100644
index 0000000..708fa38
--- /dev/null
+++ b/OWNER.md
@@ -0,0 +1 @@
+MAXIM
\ No newline at end of file
diff --git a/client.xml b/client.xml
index 10873de..59e05d5 100644
--- a/client.xml
+++ b/client.xml
@@ -1,11 +1,9 @@
-
-
-
+
-
+
diff --git a/logs/application.log b/logs/application.log
new file mode 100644
index 0000000..2419877
Binary files /dev/null and b/logs/application.log differ
diff --git a/pom.xml b/pom.xml
index f822ba6..b079c12 100644
--- a/pom.xml
+++ b/pom.xml
@@ -28,7 +28,7 @@
exec-maven-plugin
1.4.0
- arhangel.dim.lections.jdbc.JdbcExample
+ arhangel.dim.container.Application
@@ -84,17 +84,10 @@
slf4j-log4j12
1.7.5
-
- com.oracle
- ojdbc
- 14
- system
- ${basedir}/lib/ojdbc6.jar
-
org.postgresql
postgresql
- 9.4-1204-jdbc42
+ 9.4-1201-jdbc41
diff --git a/src/main/java/arhangel/dim/client/Client.java b/src/main/java/arhangel/dim/client/Client.java
index f9daeea..47ea7cb 100644
--- a/src/main/java/arhangel/dim/client/Client.java
+++ b/src/main/java/arhangel/dim/client/Client.java
@@ -1,23 +1,31 @@
package arhangel.dim.client;
-import java.io.IOException;
-import java.io.InputStream;
-import java.io.OutputStream;
-import java.net.Socket;
-import java.util.Arrays;
-import java.util.Scanner;
-
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
import arhangel.dim.container.Container;
-import arhangel.dim.container.InvalidConfigurationException;
+import arhangel.dim.container.exceptions.InvalidConfigurationException;
+import arhangel.dim.core.messages.InfoMessage;
+import arhangel.dim.core.messages.ListChatResultMessage;
import arhangel.dim.core.messages.Message;
+import arhangel.dim.core.messages.StatusMessage;
import arhangel.dim.core.messages.TextMessage;
-import arhangel.dim.core.messages.Type;
+import arhangel.dim.core.messages.CreateChatMessage;
+import arhangel.dim.core.messages.LoginMessage;
import arhangel.dim.core.net.ConnectionHandler;
import arhangel.dim.core.net.Protocol;
import arhangel.dim.core.net.ProtocolException;
+import arhangel.dim.core.net.StringProtocol;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.net.ConnectException;
+import java.net.Socket;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.Scanner;
+import java.util.stream.Collectors;
/**
* Клиент для тестирования серверного приложения
@@ -26,13 +34,12 @@ public class Client implements ConnectionHandler {
/**
* Механизм логирования позволяет более гибко управлять записью данных в лог (консоль, файл и тд)
- * */
+ */
static Logger log = LoggerFactory.getLogger(Client.class);
/**
* Протокол, хост и порт инициализируются из конфига
- *
- * */
+ */
private Protocol protocol;
private int port;
private String host;
@@ -41,6 +48,7 @@ public class Client implements ConnectionHandler {
* Тред "слушает" сокет на наличие входящих сообщений от сервера
*/
private Thread socketThread;
+ private Socket socket;
/**
* С каждым сокетом связано 2 канала in/out
@@ -48,6 +56,8 @@ public class Client implements ConnectionHandler {
private InputStream in;
private OutputStream out;
+ private Long userId;
+
public Protocol getProtocol() {
return protocol;
}
@@ -72,10 +82,26 @@ public void setHost(String host) {
this.host = host;
}
+ public static int CONNECTION_TIMEOUT_SECONDS = 10;
+
public void initSocket() throws IOException {
- Socket socket = new Socket(host, port);
- in = socket.getInputStream();
- out = socket.getOutputStream();
+ log.info("Создание сетевого подключения: " + host + ":" + port);
+ try {
+ Socket socket = new Socket(host, port);
+ in = socket.getInputStream();
+ out = socket.getOutputStream();
+ } catch (ConnectException error) {
+ log.error("Невозможно установить подключение, " +
+ "следующая попытка через " + CONNECTION_TIMEOUT_SECONDS + " секунд");
+ try {
+ Thread.sleep(CONNECTION_TIMEOUT_SECONDS * 1000L);
+ } catch (InterruptedException e) {
+ throw new IOException(e);
+ }
+
+ initSocket();
+ return;
+ }
/**
* Инициализируем поток-слушатель. Синтаксис лямбды скрывает создание анонимного класса Runnable
@@ -87,15 +113,14 @@ public void initSocket() throws IOException {
try {
// Здесь поток блокируется на ожидании данных
int read = in.read(buf);
- if (read > 0) {
- // По сети передается поток байт, его нужно раскодировать с помощью протокола
+ if (read > 0) {
+ log.info("< New Message received >");
Message msg = protocol.decode(Arrays.copyOf(buf, read));
onMessage(msg);
}
} catch (Exception e) {
log.error("Failed to process connection: {}", e);
- e.printStackTrace();
Thread.currentThread().interrupt();
}
}
@@ -109,7 +134,38 @@ public void initSocket() throws IOException {
*/
@Override
public void onMessage(Message msg) {
- log.info("Message received: {}", msg);
+ switch (msg.getType()) {
+ case MSG_STATUS:
+ StatusMessage msgStatus = (StatusMessage) msg;
+ log.info(msgStatus.getStatus());
+ System.out.println(msgStatus.getStatus());
+ break;
+ case MSG_CHAT_LIST_RESULT:
+ ListChatResultMessage msgChatListResult = (ListChatResultMessage) msg;
+ if (msgChatListResult.getChatIds().size() == 0) {
+ log.info("У вас нет активных чатов в данный момент");
+ } else {
+ log.info("Ваши активные чаты: " + String.join(",", msgChatListResult.getChatIds().stream()
+ .map(Object::toString)
+ .collect(Collectors.toList())));
+ }
+ System.out.println("Ваши активные чаты: " + String.join(",", msgChatListResult.getChatIds().stream()
+ .map(Object::toString)
+ .collect(Collectors.toList())));
+ break;
+ case MSG_INFO:
+ InfoMessage infoMessage = (InfoMessage) msg;
+ log.info(infoMessage.getInfo());
+ break;
+ case MSG_TEXT:
+ log.info(msg.toString());
+ System.out.println(msg.toString());
+ break;
+ default:
+ log.error("Данный тип сообщений не поддерживается");
+ System.err.println("Полученный тип сообщений не поддерживается");
+ break;
+ }
}
/**
@@ -122,22 +178,62 @@ public void processInput(String line) throws IOException, ProtocolException {
String cmdType = tokens[0];
switch (cmdType) {
case "/login":
- // TODO: реализация
+ if (tokens.length != 3) {
+ log.error("Неправильное использование команды");
+ break;
+ }
+ LoginMessage loginMessage = new LoginMessage();
+ loginMessage.setLogin(tokens[1]);
+ loginMessage.setPassword(tokens[2]);
+ send(loginMessage);
break;
case "/help":
- // TODO: реализация
+ log.info(
+ "login <имя пользователя> <пароль> - выполнить вход\n" +
+ "help - справка по месседжеру\n" +
+ "text <сообщение> - отправить сообщение в чат с заданным id\n" +
+ "info - получить информациию о пользователе\n" +
+ "info - получить информациию о себе\n" +
+ "chat_create - создать новый чат " +
+ "(использование: chat_create 1, 2, 3, 66)\n"
+ );
break;
case "/text":
- // FIXME: пример реализации для простого текстового сообщения
+ if (tokens.length != 3) {
+ log.error("Неправильное использование команды");
+ break;
+ }
TextMessage sendMessage = new TextMessage();
- sendMessage.setType(Type.MSG_TEXT);
- sendMessage.setText(tokens[1]);
+ sendMessage.setChatId(Long.parseLong(tokens[1]));
+ sendMessage.setText(tokens[2]);
send(sendMessage);
break;
- // TODO: implement another types from wiki
-
+ case "/info":
+ InfoMessage infoMessage = new InfoMessage();
+ if (tokens.length > 1) {
+ infoMessage.setUserId(Long.parseLong(tokens[1]));
+ } else {
+ infoMessage.setUserId(-1);
+ }
+ infoMessage.setInfo(null);
+ send(infoMessage);
+ break;
+ case "/chat_create":
+ CreateChatMessage createChatMessage = new CreateChatMessage();
+ String[] userIdsStr = tokens[1].split(",");
+ List userIds = new ArrayList<>();
+ for (String anUserIdsStr : userIdsStr) {
+ userIds.add(Long.parseLong(anUserIdsStr));
+ }
+ createChatMessage.setUsersIds(userIds);
+ send(createChatMessage);
+ break;
+ case "/chat_hist":
+ //TODO: Доделать
+ case "/chat_list":
+ //TODO: Доделать
default:
- log.error("Invalid input: " + line);
+ log.error("Команда не найдена: " + line);
}
}
@@ -146,23 +242,33 @@ public void processInput(String line) throws IOException, ProtocolException {
*/
@Override
public void send(Message msg) throws IOException, ProtocolException {
- log.info(msg.toString());
+ log.info("< Sending new message >");
out.write(protocol.encode(msg));
- out.flush(); // принудительно проталкиваем буфер с данными
+ out.flush();
}
@Override
- public void close() {
- // TODO: написать реализацию. Закройте ресурсы и остановите поток-слушатель
+ public void close() throws IOException {
+ log.error("Closing socket...");
+ if (!socketThread.isInterrupted()) {
+ socketThread.interrupt();
+ }
+
+ if (!socket.isClosed()) {
+ socket.close();
+ }
}
public static void main(String[] args) throws Exception {
- Client client = null;
+ Client client;
// Пользуемся механизмом контейнера
+ //TODO: or not o:
try {
Container context = new Container("client.xml");
client = (Client) context.getByName("client");
+ //TODO FIX
+ client.setProtocol(new StringProtocol());
} catch (InvalidConfigurationException e) {
log.error("Failed to create client", e);
return;
@@ -175,7 +281,7 @@ public static void main(String[] args) throws Exception {
System.out.println("$");
while (true) {
String input = scanner.nextLine();
- if ("q".equals(input)) {
+ if ("exit".equals(input)) {
return;
}
try {
diff --git a/src/main/java/arhangel/dim/container/Application.java b/src/main/java/arhangel/dim/container/Application.java
new file mode 100644
index 0000000..e14189a
--- /dev/null
+++ b/src/main/java/arhangel/dim/container/Application.java
@@ -0,0 +1,10 @@
+package arhangel.dim.container;
+
+/**
+ * Created by spec45as on 3/23/2016.
+ */
+public class Application {
+ public static void main(String[] args) {
+ System.out.println("Container Application");
+ }
+}
diff --git a/src/main/java/arhangel/dim/container/BeanGraph.java b/src/main/java/arhangel/dim/container/BeanGraph.java
index 306c0e6..91b0c41 100644
--- a/src/main/java/arhangel/dim/container/BeanGraph.java
+++ b/src/main/java/arhangel/dim/container/BeanGraph.java
@@ -1,8 +1,13 @@
package arhangel.dim.container;
-import java.util.HashMap;
+import arhangel.dim.container.exceptions.CycleReferenceException;
+
import java.util.List;
+import java.util.ArrayList;
import java.util.Map;
+import java.util.HashMap;
+
+import static java.util.Collections.reverse;
/**
*
@@ -10,13 +15,19 @@
public class BeanGraph {
// Граф представлен в виде списка связности для каждой вершины
private Map> vertices = new HashMap<>();
+ private Map vertexByName = new HashMap<>();
/**
* Добавить вершину в граф
* @param value - объект, привязанный к вершине
*/
public BeanVertex addVertex(Bean value) {
- return null;
+ BeanVertex newVertex = new BeanVertex(value);
+
+ vertices.put(newVertex, new ArrayList());
+ vertexByName.put(value.getName(), newVertex);
+
+ return newVertex;
}
/**
@@ -25,26 +36,103 @@ public BeanVertex addVertex(Bean value) {
* @param to в какую вершину
*/
public void addEdge(BeanVertex from ,BeanVertex to) {
+ List incidentVertices = vertices.get(from);
+ incidentVertices.add(to);
+ }
+
+ public BeanGraph() {}
+
+ public BeanGraph(List beans) {
+ // adding all the vertices
+ for (Bean bean : beans) {
+ addVertex(bean);
+ }
+
+ // adding edges between vertices
+ for (Bean bean : beans) {
+ BeanVertex from = vertexByName.get(bean.getName());
+
+ HashMap properties = (HashMap) bean.getProperties();
+ for (Property property : properties.values()) {
+
+ if (property.getType() == ValueType.VAL) {
+ continue;
+ }
+
+ BeanVertex to = vertexByName.get(property.getName());
+ addEdge(from, to);
+ }
+ }
}
/**
* Проверяем, связаны ли вершины
*/
public boolean isConnected(BeanVertex v1, BeanVertex v2) {
- return false;
+ List incidentVertices = vertices.get(v1);
+ return incidentVertices.contains(v2);
}
/**
* Получить список вершин, с которыми связана vertex
*/
public List getLinked(BeanVertex vertex) {
- return null;
+ return vertices.get(vertex);
}
/**
* Количество вершин в графе
*/
public int size() {
- return 0;
+ return vertices.size();
+ }
+
+ private enum VertexType {
+ NOT_PROCESSED, // dfs в вершину еще не заходил
+ STARTED_PROCESSING, // dfs зашел в вершину
+ FINISHED_PROCESSING // dfs вышел из вершины
+ }
+
+ /**
+ * Проверить граф на наличие циклов
+ */
+ private boolean isCircle(BeanVertex vertex, List sortedVertices,
+ Map usedVertices) {
+ usedVertices.put(vertex, VertexType.STARTED_PROCESSING);
+
+ for (BeanVertex incidentVertex : vertices.get(vertex)) {
+ if (usedVertices.get(incidentVertex).equals(VertexType.NOT_PROCESSED)) {
+ isCircle(incidentVertex, sortedVertices, usedVertices);
+ } else if (usedVertices.get(incidentVertex).equals(VertexType.STARTED_PROCESSING)) {
+ return true;
+ }
+ }
+ sortedVertices.add(vertex);
+ usedVertices.put(vertex, VertexType.FINISHED_PROCESSING);
+ return false;
+ }
+
+ /**
+ * Отсортировать вершины графа в топологическом порядке
+ */
+ public List sortTopologically() throws CycleReferenceException {
+ Map usedVertices = new HashMap<>();
+ for (BeanVertex vertex : vertices.keySet()) {
+ usedVertices.put(vertex, VertexType.NOT_PROCESSED);
+ }
+ List sortedVertices = new ArrayList<>();
+
+ for (BeanVertex vertex : vertices.keySet()) {
+ if (!usedVertices.get(vertex).equals(0)) {
+ continue;
+ }
+
+ boolean foundCircle = isCircle(vertex, sortedVertices, usedVertices);
+ if (foundCircle) {
+ throw new CycleReferenceException("circle reference found");
+ }
+
+ }
+ return sortedVertices;
}
}
diff --git a/src/main/java/arhangel/dim/container/BeanXmlReader.java b/src/main/java/arhangel/dim/container/BeanXmlReader.java
index 3350f1e..75b3edf 100644
--- a/src/main/java/arhangel/dim/container/BeanXmlReader.java
+++ b/src/main/java/arhangel/dim/container/BeanXmlReader.java
@@ -1,10 +1,22 @@
package arhangel.dim.container;
+import javax.xml.parsers.DocumentBuilderFactory;
+import javax.xml.parsers.DocumentBuilder;
+import javax.xml.parsers.ParserConfigurationException;
+
+import arhangel.dim.container.exceptions.InvalidConfigurationException;
+import org.w3c.dom.Document;
+import org.w3c.dom.NodeList;
+import org.w3c.dom.Node;
+import org.w3c.dom.Element;
+import org.xml.sax.SAXException;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.HashMap;
import java.util.List;
-/**
- *
- */
public class BeanXmlReader {
private static final String TAG_BEAN = "bean";
private static final String TAG_PROPERTY = "property";
@@ -14,8 +26,55 @@ public class BeanXmlReader {
private static final String ATTR_BEAN_ID = "id";
private static final String ATTR_BEAN_CLASS = "class";
- public List parseBeans(String pathToFile) {
- return null;
- }
+ public List parseBeans(String pathToFile) throws InvalidConfigurationException {
+ try {
+ File xmlFile = new File(pathToFile);
+ DocumentBuilderFactory dbFactory = DocumentBuilderFactory.newInstance();
+ DocumentBuilder dbBuilder = dbFactory.newDocumentBuilder();
+ Document document = dbBuilder.parse(xmlFile);
+
+ NodeList beanNodeList = document.getElementsByTagName(TAG_BEAN);
+ List beanList = new ArrayList();
+ for (int i = 0; i < beanNodeList.getLength(); i++) {
+ Node beanNode = beanNodeList.item(i);
+
+ if (beanNode.getNodeType() != Node.ELEMENT_NODE) {
+ continue;
+ }
+
+ Element beanElement = (Element) beanNode;
+
+ String id = beanElement.getAttribute(ATTR_BEAN_ID);
+ String clazz = beanElement.getAttribute(ATTR_BEAN_CLASS);
+
+ NodeList propertyList = beanElement.getElementsByTagName(TAG_PROPERTY);
+ HashMap properties = new HashMap();
+ for (int j = 0; j < propertyList.getLength(); j++) {
+ Node propertyNode = propertyList.item(j);
+ if (propertyNode.getNodeType() != Node.ELEMENT_NODE) {
+ continue;
+ }
+
+ Element propertyElement = (Element) propertyNode;
+
+ String name = propertyElement.getAttribute(ATTR_NAME);
+ String refValue = propertyElement.getAttribute(ATTR_REF);
+
+ ValueType ref = refValue.equals("") ? ValueType.VAL : ValueType.REF;
+ String val = refValue.equals("") ? propertyElement.getAttribute(ATTR_VALUE) : refValue;
+
+ Property property = new Property(name, val, ref);
+ properties.put(name, property);
+ }
+ Bean bean = new Bean(id, clazz, properties);
+ beanList.add(bean);
+ }
+
+ return beanList;
+
+ } catch (ParserConfigurationException | IOException | SAXException e) {
+ throw new InvalidConfigurationException("invalid config file format");
+ }
+ }
}
diff --git a/src/main/java/arhangel/dim/container/Container.java b/src/main/java/arhangel/dim/container/Container.java
index 4213396..86a6f8c 100644
--- a/src/main/java/arhangel/dim/container/Container.java
+++ b/src/main/java/arhangel/dim/container/Container.java
@@ -1,12 +1,25 @@
package arhangel.dim.container;
+import arhangel.dim.container.beans.Car;
+import arhangel.dim.container.beans.Engine;
+import arhangel.dim.container.beans.Gear;
+import arhangel.dim.container.exceptions.InvalidConfigurationException;
+
+import java.lang.reflect.Field;
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
/**
* Используйте ваш xml reader чтобы прочитать конфиг и получить список бинов
*/
public class Container {
+
private List beans;
+ private Map objByName = new HashMap<>();
+ private Map objByClassName = new HashMap<>();
/**
* Если не получается считать конфиг, то бросьте исключение
@@ -14,7 +27,16 @@ public class Container {
*/
public Container(String pathToConfig) throws InvalidConfigurationException {
- // вызываем BeanXmlReader
+ try {
+ BeanXmlReader beanXmlReader = new BeanXmlReader();
+ beans = beanXmlReader.parseBeans(pathToConfig);
+ for (Bean bean: beans) {
+ instantiateBean(bean);
+ }
+ } catch (InvalidConfigurationException e) {
+ System.err.println(e.getMessage());
+ e.printStackTrace();
+ }
}
/**
@@ -22,7 +44,7 @@ public Container(String pathToConfig) throws InvalidConfigurationException {
* Например, Car car = (Car) container.getByName("carBean")
*/
public Object getByName(String name) {
- return null;
+ return objByName.get(name);
}
/**
@@ -30,35 +52,71 @@ public Object getByName(String name) {
* Например, Car car = (Car) container.getByClass("arhangel.dim.container.Car")
*/
public Object getByClass(String className) {
- return null;
+ return objByClassName.get(className);
}
- private void instantiateBean(Bean bean) {
-
- /*
- // Примерный ход работы
-
- String className = bean.getClassName();
- Class clazz = Class.forName(className);
- // ищем дефолтный конструктор
- Object ob = clazz.newInstance();
-
-
- for (String name : bean.getProperties().keySet()) {
- // ищем поле с таким именен внутри класса
- // учитывая приватные
- Field field = clazz.getDeclaredField(name);
- // проверяем, если такого поля нет, то кидаем InvalidConfigurationException с описание ошибки
-
- // Делаем приватные поля доступными
- field.setAccessible(true);
-
- // Далее определяем тип поля и заполняем его
- // Если поле - примитив, то все просто
- // Если поле ссылка, то эта ссылка должа была быть инициализирована ранее
+ public List getBeans() {
+ return beans;
+ }
- */
+ // метод, кастующий строку аргумента value к его типу clazz
+ private Object parseValue(String clazz, String value) {
+ switch (clazz) {
+ case "boolean":
+ return Boolean.valueOf(value);
+ case "byte":
+ return Byte.valueOf(value);
+ case "short":
+ return Short.valueOf(value);
+ case "int":
+ return Integer.valueOf(value);
+ case "long":
+ return Long.valueOf(value);
+ case "double":
+ return Double.valueOf(value);
+ case "float":
+ return Float.valueOf(value);
+ case "java.lang.String":
+ return value;
+ default:
+ return objByName.get(value);
+ }
+ }
+ private void instantiateBean(Bean bean) throws InvalidConfigurationException {
+ try {
+ String className = bean.getClassName();
+ Class clazz = Class.forName(className);
+ // ищем дефолтный конструктор
+ Object object = clazz.newInstance();
+
+ for (String name : bean.getProperties().keySet()) {
+ // ищем поле с таким именен внутри класса
+ // учитывая приватные
+ Field field = clazz.getDeclaredField(name);
+ // проверяем, если такого поля нет, то кидаем InvalidConfigurationException с описанием ошибки
+
+ // определяем тип аргумента для передачи методу set
+ Class[] argTypes = new Class[] { field.getType() };
+ // определяем имя сеттера
+ String methodName = "set" + name.substring(0, 1).toUpperCase() + name.substring(1);
+ // по имени и типу аргумента получаем метод
+ Method setFieldValue = clazz.getDeclaredMethod(methodName, argTypes);
+
+ // определем value to be set
+ Object valueToSet = parseValue(field.getType().getName(), bean.getProperties().get(name).getValue());
+ // вызываем setter
+ setFieldValue.invoke(object, valueToSet);
+
+ }
+
+ objByName.put(bean.getName(), object);
+ objByClassName.put(className, object);
+
+ } catch (ClassNotFoundException | IllegalAccessException | InstantiationException |
+ NoSuchFieldException | NoSuchMethodException | InvocationTargetException e) {
+ throw new InvalidConfigurationException(e.getMessage());
+ }
}
}
diff --git a/src/main/java/arhangel/dim/container/Main.java b/src/main/java/arhangel/dim/container/Main.java
new file mode 100644
index 0000000..ee50711
--- /dev/null
+++ b/src/main/java/arhangel/dim/container/Main.java
@@ -0,0 +1,37 @@
+package arhangel.dim.container;
+
+import arhangel.dim.container.beans.Car;
+import arhangel.dim.core.net.Protocol;
+import arhangel.dim.core.net.SerialProtocol;
+
+import java.util.ArrayList;
+import java.util.List;
+
+public class Main {
+ public static void main(String[] args) {
+ BeanXmlReader reader = new BeanXmlReader();
+ try {
+ List beans = reader.parseBeans("/home/spec45as/technotrack/messenger/config.xml");
+ for (Bean bean : beans) {
+ System.out.println(bean.toString());
+ }
+ BeanGraph graph = new BeanGraph(reader.parseBeans("/home/spec45as/technotrack/messenger/config.xml"));
+
+ List vertices = graph.sortTopologically();
+ for (BeanVertex vertex : vertices) {
+ System.out.println(vertex.getBean());
+ }
+ } catch (Exception e) {
+ e.printStackTrace();
+ }
+
+ try {
+ Container container = new Container("/home/spec45as/technotrack/messenger/config.xml");
+ Car car = (Car) container.getByName("carBean");
+ System.out.println(car.getEngine());
+ } catch (Exception e) {
+ e.printStackTrace();
+ }
+ }
+
+}
diff --git a/src/main/java/arhangel/dim/container/exceptions/CycleReferenceException.java b/src/main/java/arhangel/dim/container/exceptions/CycleReferenceException.java
new file mode 100644
index 0000000..9f5a776
--- /dev/null
+++ b/src/main/java/arhangel/dim/container/exceptions/CycleReferenceException.java
@@ -0,0 +1,14 @@
+package arhangel.dim.container.exceptions;
+
+/**
+ * Обнаружена циклическая зависимость
+ */
+public class CycleReferenceException extends Exception {
+ public CycleReferenceException(String message) {
+ super(message);
+ }
+
+ public String getMessage() {
+ return super.getMessage();
+ }
+}
diff --git a/src/main/java/arhangel/dim/container/exceptions/InvalidConfigurationException.java b/src/main/java/arhangel/dim/container/exceptions/InvalidConfigurationException.java
new file mode 100644
index 0000000..be5e1b2
--- /dev/null
+++ b/src/main/java/arhangel/dim/container/exceptions/InvalidConfigurationException.java
@@ -0,0 +1,11 @@
+package arhangel.dim.container.exceptions;
+
+public class InvalidConfigurationException extends Exception {
+ public InvalidConfigurationException(String message) {
+ super(message);
+ }
+
+ public String getMessage() {
+ return super.getMessage();
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/Chat.java b/src/main/java/arhangel/dim/core/Chat.java
index 4363d4b..b12952b 100644
--- a/src/main/java/arhangel/dim/core/Chat.java
+++ b/src/main/java/arhangel/dim/core/Chat.java
@@ -1,8 +1,54 @@
package arhangel.dim.core;
+import java.util.ArrayList;
+import java.util.List;
+
/**
* А над этим классом надо еще поработать
*/
public class Chat {
private Long id;
+ private Long creatorId;
+ private List messageIds = new ArrayList<>();
+ private List participantIds = new ArrayList<>();
+
+ public Long getId() {
+ return id;
+ }
+
+ public void setId(Long id) {
+ this.id = id;
+ }
+
+ public void setCreatorId(Long creatorId) {
+ this.creatorId = creatorId;
+ }
+
+ public Long getCreatorId() {
+ return creatorId;
+ }
+
+ public List getMessageIds() {
+ return messageIds;
+ }
+
+ public void setMessageIds(List messageIds) {
+ this.messageIds = messageIds;
+ }
+
+ public List getParticipantIds() {
+ return participantIds;
+ }
+
+ public void setParticipantIds(List participantIds) {
+ this.participantIds = participantIds;
+ }
+
+ public void addParticipant(Long id) {
+ participantIds.add(id);
+ }
+
+ public void addMessage(Long id) {
+ messageIds.add(id);
+ }
}
diff --git a/src/main/java/arhangel/dim/core/User.java b/src/main/java/arhangel/dim/core/User.java
index 0f21c0c..3757ba7 100644
--- a/src/main/java/arhangel/dim/core/User.java
+++ b/src/main/java/arhangel/dim/core/User.java
@@ -1,11 +1,27 @@
package arhangel.dim.core;
+import arhangel.dim.core.service.AuthorizationService;
+
+import java.math.BigInteger;
+
/**
* Представление пользователя
*/
public class User {
private Long id;
private String name;
+ private String hash;
+
+ public User() {}
+
+ public User(final String name) {
+ this.name = name;
+ }
+
+ public User(final String name, final String password) {
+ this.name = name;
+ hash = password;
+ }
public Long getId() {
return id;
@@ -22,4 +38,25 @@ public String getName() {
public void setName(String name) {
this.name = name;
}
+
+ public void setHash(final String hash) {
+ this.hash = hash;
+ }
+
+ public String getHash() {
+ return hash;
+ }
+
+ public void setPass(final String password) {
+ hash = password;
+ }
+
+ @Override
+ public String toString() {
+ return "User{" +
+ "userName='" + name + '\'' +
+ ", hash=" + hash +
+ ", userID=" + id +
+ '}';
+ }
}
diff --git a/src/main/java/arhangel/dim/core/command/CreateChatCommand.java b/src/main/java/arhangel/dim/core/command/CreateChatCommand.java
new file mode 100644
index 0000000..4748943
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/command/CreateChatCommand.java
@@ -0,0 +1,75 @@
+package arhangel.dim.core.command;
+
+import arhangel.dim.core.Chat;
+import arhangel.dim.core.messages.Command;
+import arhangel.dim.core.messages.CommandException;
+import arhangel.dim.core.messages.CreateChatMessage;
+import arhangel.dim.core.messages.Message;
+import arhangel.dim.core.messages.StatusMessage;
+import arhangel.dim.core.net.ProtocolException;
+import arhangel.dim.core.net.Session;
+import arhangel.dim.core.store.MessageStore;
+import arhangel.dim.server.Server;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.util.List;
+
+public class CreateChatCommand implements Command {
+
+ private Server server;
+
+ static Logger log = LoggerFactory.getLogger(CreateChatCommand.class);
+
+ public CreateChatCommand(Server server) {
+ this.server = server;
+ }
+
+ @Override
+ public void execute(Session session, Message message) throws CommandException {
+ try {
+ if (session.getUser() == null) {
+ StatusMessage errorMessage = new StatusMessage();
+ errorMessage.setStatus("Ony authorised person can create chat");
+ session.send(errorMessage);
+ return;
+ }
+
+ MessageStore messageStore = session.getMessageStore();
+
+ List participants = ((CreateChatMessage) message).getUsersIds();
+ if (participants.size() == 1) {
+ List chatsByUser = messageStore.getChatsByUserId(session.getUser().getId());
+
+ for (Long chatId : chatsByUser) {
+ Chat chat = messageStore.getChatById(chatId);
+ log.info("CHAT " + chat.toString());
+
+ if (chat.getParticipantIds().size() == 1 &&
+ chat.getParticipantIds().get(0).equals(participants.get(0))) {
+ StatusMessage response = new StatusMessage();
+ response.setStatus(String.format(
+ "Сhat with user %d already exists. Chat id: %d",
+ participants.get(0), chat.getId()));
+ session.send(response);
+ return;
+ }
+
+ }
+
+ }
+
+ Chat chat = messageStore.createChat(session.getUser().getId(), participants);
+
+ StatusMessage response = new StatusMessage();
+ response.setStatus(String.format("Chat created. Chat id: %d", chat.getId()));
+ session.send(response);
+ return;
+
+
+ } catch (ProtocolException | IOException e) {
+ throw new CommandException(e);
+ }
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/command/HistChatCommand.java b/src/main/java/arhangel/dim/core/command/HistChatCommand.java
new file mode 100644
index 0000000..36bff23
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/command/HistChatCommand.java
@@ -0,0 +1,79 @@
+package arhangel.dim.core.command;
+
+import arhangel.dim.core.User;
+import arhangel.dim.core.messages.Command;
+import arhangel.dim.core.messages.CommandException;
+import arhangel.dim.core.messages.HistChatMessage;
+import arhangel.dim.core.messages.HistChatResultMessage;
+import arhangel.dim.core.messages.Message;
+import arhangel.dim.core.messages.StatusMessage;
+import arhangel.dim.core.net.ProtocolException;
+import arhangel.dim.core.net.Session;
+import arhangel.dim.core.store.MessageStore;
+import arhangel.dim.server.Server;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.util.List;
+
+public class HistChatCommand implements Command {
+
+ private Server server;
+
+ static Logger log = LoggerFactory.getLogger(CreateChatCommand.class);
+
+ public HistChatCommand(Server server) {
+ this.server = server;
+ }
+
+ @Override
+ public void execute(Session session, Message message) throws CommandException {
+ try {
+ HistChatMessage chatHistoryMessage = (HistChatMessage) message;
+
+ Long chatId = chatHistoryMessage.getChatId();
+
+ StatusMessage errorMessage = new StatusMessage();
+
+ User user = session.getUser();
+ if (user == null) {
+ log.info("HistChat: user not authorised ");
+ errorMessage.setStatus("You are not authorised");
+ session.send(errorMessage);
+ return;
+ }
+
+ MessageStore messageStore = server.getMessageStore();
+
+ if (!messageStore.getChatById(chatId).getParticipantIds().contains(user.getId())) {
+ log.info(String.format("HistChat: user %d not invited to chat %d", user.getId(), chatId));
+ errorMessage.setStatus("user with id=" + user.getId() + " wasn't invited to chat with id=" + chatId);
+ session.send(errorMessage);
+ return;
+ }
+
+ log.info(String.format("User %d: chat %d history", user.getId(), chatId));
+ List messagesFromChat = messageStore.getMessagesFromChat(chatId);
+
+ StringBuilder chatHistory = new StringBuilder();
+
+ for (long i = 0; i < messagesFromChat.size(); ++i) {
+ Long id = messagesFromChat.get((int) i);
+ chatHistory.append(messageStore.getMessageById(messagesFromChat.get((int) i)).toString() + "\n");
+ }
+
+ HistChatResultMessage histChatResultMessage = new HistChatResultMessage();
+ histChatResultMessage.setChatId(chatId);
+ histChatResultMessage.setHistory(chatHistory.toString());
+ session.send(histChatResultMessage);
+ return;
+
+ } catch (IOException | ProtocolException e) {
+ throw new CommandException(e);
+ }
+ }
+
+
+
+}
diff --git a/src/main/java/arhangel/dim/core/command/InfoCommand.java b/src/main/java/arhangel/dim/core/command/InfoCommand.java
new file mode 100644
index 0000000..0ec7ba8
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/command/InfoCommand.java
@@ -0,0 +1,64 @@
+package arhangel.dim.core.command;
+
+import arhangel.dim.core.User;
+import arhangel.dim.core.messages.Command;
+import arhangel.dim.core.messages.CommandException;
+import arhangel.dim.core.messages.InfoMessage;
+import arhangel.dim.core.messages.Message;
+import arhangel.dim.core.messages.StatusMessage;
+import arhangel.dim.core.net.ProtocolException;
+import arhangel.dim.core.net.Session;
+import arhangel.dim.core.store.UserStore;
+import arhangel.dim.server.Server;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+
+public class InfoCommand implements Command {
+
+ private Server server;
+
+ static Logger log = LoggerFactory.getLogger(CreateChatCommand.class);
+
+ public InfoCommand(Server server) {
+ this.server = server;
+ }
+
+ @Override
+ public void execute(Session session, Message message) throws CommandException {
+ try {
+ InfoMessage infoMessage = (InfoMessage) message;
+ StatusMessage errorMessage = new StatusMessage();
+
+ UserStore userStore = server.getUserStore();
+
+ if (session.getUser() == null) {
+ errorMessage.setStatus("Ony authorised person can get info");
+ session.send(errorMessage);
+ return;
+ }
+
+ long userId = session.getUser().getId();
+ if (infoMessage.getUserId() != -1) {
+ userId = infoMessage.getUserId();
+ }
+
+ User user = server.getUserStore().getUserById(userId);
+
+ if (user == null) {
+ errorMessage.setStatus("No user with id " + userId);
+ session.send(errorMessage);
+ return;
+ }
+
+ InfoMessage response = new InfoMessage();
+ response.setUserId(userId);
+ response.setInfo("User Login: " + user.getName());
+ session.send(response);
+
+ } catch (IOException | ProtocolException e) {
+ throw new CommandException(e);
+ }
+ }
+}
\ No newline at end of file
diff --git a/src/main/java/arhangel/dim/core/command/ListChatCommand.java b/src/main/java/arhangel/dim/core/command/ListChatCommand.java
new file mode 100644
index 0000000..cbc1dac
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/command/ListChatCommand.java
@@ -0,0 +1,61 @@
+package arhangel.dim.core.command;
+
+import arhangel.dim.core.User;
+import arhangel.dim.core.messages.Command;
+import arhangel.dim.core.messages.CommandException;
+import arhangel.dim.core.messages.ListChatMessage;
+import arhangel.dim.core.messages.ListChatResultMessage;
+import arhangel.dim.core.messages.Message;
+import arhangel.dim.core.messages.StatusMessage;
+import arhangel.dim.core.net.ProtocolException;
+import arhangel.dim.core.net.Session;
+import arhangel.dim.core.store.MessageStore;
+import arhangel.dim.core.store.MessageStoreImpl;
+import arhangel.dim.server.Server;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.util.List;
+
+public class ListChatCommand implements Command {
+
+ private Server server;
+
+ static Logger log = LoggerFactory.getLogger(CreateChatCommand.class);
+
+ public ListChatCommand(Server server) {
+ this.server = server;
+ }
+
+ @Override
+ public void execute(Session session, Message message) throws CommandException {
+ try {
+ ListChatMessage chatListMessage = (ListChatMessage)message;
+ StatusMessage errorMessage = new StatusMessage();
+
+ User user = session.getUser();
+ if (user == null) {
+ errorMessage.setStatus("Ony authorised person can see chats");
+ session.send(errorMessage);
+ return;
+ }
+
+ MessageStoreImpl messageStore = (MessageStoreImpl)server.getMessageStore();
+ List chatIds = messageStore.getChatsByUserId(user.getId());
+ if (chatIds == null) {
+ errorMessage.setStatus("No chats for user");
+ session.send(errorMessage);
+ return;
+ }
+
+ ListChatResultMessage listChatResultMessage = new ListChatResultMessage();
+ listChatResultMessage.setChatIds(chatIds);
+ session.send(listChatResultMessage);
+
+ } catch (IOException | ProtocolException e) {
+ throw new CommandException(e);
+ }
+
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/command/LoginCommand.java b/src/main/java/arhangel/dim/core/command/LoginCommand.java
new file mode 100644
index 0000000..f4b738b
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/command/LoginCommand.java
@@ -0,0 +1,64 @@
+package arhangel.dim.core.command;
+
+import arhangel.dim.core.User;
+import arhangel.dim.core.messages.Command;
+import arhangel.dim.core.messages.CommandException;
+import arhangel.dim.core.messages.LoginMessage;
+import arhangel.dim.core.messages.Message;
+import arhangel.dim.core.messages.TextMessage;
+import arhangel.dim.core.net.ProtocolException;
+import arhangel.dim.core.net.Session;
+import arhangel.dim.core.store.MessageStoreImpl;
+import arhangel.dim.core.store.UserStoreImpl;
+import arhangel.dim.server.Server;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+
+public class LoginCommand implements Command {
+
+ private Server server;
+
+ static Logger log = LoggerFactory.getLogger(CreateChatCommand.class);
+
+ public LoginCommand(Server server) {
+ this.server = server;
+ }
+
+ @Override
+ public void execute(Session session, Message message) throws CommandException {
+ LoginMessage loginMsg = (LoginMessage) message;
+
+ String userName = loginMsg.getLogin();
+ String password = loginMsg.getPassword();
+
+ UserStoreImpl userStore = session.getUserStore();
+ MessageStoreImpl messageStore = session.getMessageStore();
+
+ TextMessage result = new TextMessage();
+ if (!userStore.isUserExist(userName)) {
+ result.setText("User not found");
+ try {
+ session.send(result);
+ } catch (IOException | ProtocolException e) {
+ e.printStackTrace();
+ }
+ return;
+ }
+
+ User user = userStore.getUser(userName, password);
+ if (user == null) {
+ result.setText("Login or password is incorrect");
+ } else {
+ session.setUser(user);
+ result.setText(String.format("Success: userid: %d, name: %s", user.getId(), user.getName()));
+ }
+
+ try {
+ session.send(result);
+ } catch (IOException | ProtocolException e) {
+ e.printStackTrace();
+ }
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/command/TextCommand.java b/src/main/java/arhangel/dim/core/command/TextCommand.java
new file mode 100644
index 0000000..0dbf584
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/command/TextCommand.java
@@ -0,0 +1,65 @@
+package arhangel.dim.core.command;
+
+import arhangel.dim.core.Chat;
+import arhangel.dim.core.User;
+import arhangel.dim.core.messages.Command;
+import arhangel.dim.core.messages.CommandException;
+import arhangel.dim.core.messages.Message;
+import arhangel.dim.core.messages.StatusMessage;
+import arhangel.dim.core.messages.TextMessage;
+import arhangel.dim.core.messages.Type;
+import arhangel.dim.core.net.Session;
+import arhangel.dim.core.store.MessageStoreImpl;
+import arhangel.dim.core.store.UserStoreImpl;
+import arhangel.dim.server.Server;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class TextCommand implements Command {
+
+ private Server server;
+
+ static Logger log = LoggerFactory.getLogger(CreateChatCommand.class);
+
+ public TextCommand(Server server) {
+ this.server = server;
+ }
+
+ @Override
+ public void execute(Session session, Message message) throws CommandException {
+ try {
+
+ StatusMessage errorMessage = new StatusMessage();
+ if (session.getUser() == null) {
+ errorMessage.setStatus("Only authorised person can send messages");
+ session.send(errorMessage);
+ return;
+ }
+
+ MessageStoreImpl messageStore = (MessageStoreImpl)server.getMessageStore();
+
+ TextMessage textMessage = (TextMessage)message;
+ textMessage.setSenderId(session.getUser().getId());
+
+ Chat chat = messageStore.getChatById(textMessage.getChatId());
+
+ StatusMessage response = new StatusMessage();
+ response.setStatus(String.format("User %s wrote to chat %d: %s",
+ session.getUser().getName(), chat.getId(), textMessage.getText()));
+ response.setSenderId(textMessage.getSenderId());
+
+ for (Long chatUsersId : chat.getParticipantIds()) {
+ for (Session s : server.getSessions()) {
+ if (s.getUser() != null && chatUsersId.equals(s.getUser().getId())) {
+ s.send(response);
+ }
+ }
+
+ }
+
+ } catch (Exception e) {
+ throw new CommandException(e);
+ }
+ }
+
+}
diff --git a/src/main/java/arhangel/dim/core/dbservice/dao/UsersDao.java b/src/main/java/arhangel/dim/core/dbservice/dao/UsersDao.java
new file mode 100644
index 0000000..6d06de4
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/dbservice/dao/UsersDao.java
@@ -0,0 +1,347 @@
+package arhangel.dim.core.dbservice.dao;
+
+import arhangel.dim.core.Chat;
+import arhangel.dim.core.User;
+import arhangel.dim.core.dbservice.executor.QueryExecutor;
+import arhangel.dim.core.messages.Message;
+import arhangel.dim.core.messages.Type;
+import org.postgresql.ds.PGPoolingDataSource;
+import org.slf4j.LoggerFactory;
+
+import java.sql.Connection;
+import java.sql.SQLException;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+public class UsersDao {
+ static org.slf4j.Logger log = LoggerFactory.getLogger(UsersDao.class);
+ QueryExecutor queryExecutor;
+
+ public UsersDao() {
+ }
+
+ // database access test
+ public static void main(String[] args) throws Exception {
+ UsersDao usersDao = new UsersDao();
+ usersDao.init();
+
+ //Long num = new Long(0);
+ usersDao.createTables();
+ //usersDao.addUserToChat(new Long(0),new Long(1));
+ usersDao.addUser("spec45as", "alphatest");
+ usersDao.addUser("spec45as2", "alphatest2");
+
+ }
+
+ public void init() throws SQLException, ClassNotFoundException {
+
+ Class.forName("org.postgresql.Driver");
+
+ PGPoolingDataSource source = new PGPoolingDataSource();
+ source.setDataSourceName("jdbc:postgresql");
+ source.setServerName("178.62.140.149");
+
+ //TODO: актуальный когда допилишь для проекта
+ source.setDatabaseName("spec45as");
+ source.setUser("trackuser");
+ source.setPassword("trackuser");
+ source.setMaxConnections(10);
+
+ Connection connection = source.getConnection();
+
+ queryExecutor = new QueryExecutor();
+ queryExecutor.setConnection(connection);
+ createTables();
+
+ }
+
+ public void createTables() throws SQLException {
+ try {
+ createUsers();
+ createChats();
+ createMessages();
+ } catch (SQLException e) {
+ throw new SQLException("Ошибка создания таблиц", e);
+ }
+ }
+
+ public void createUsers() throws SQLException {
+ String sql = "CREATE TABLE IF NOT EXISTS users (" +
+ " user_id SERIAL PRIMARY KEY, " +
+ " user_login character varying, " +
+ " user_pass character varying " +
+ ");";
+ queryExecutor.updateQuery(sql);
+ }
+
+ public void createMessages() throws SQLException {
+ String sql = "CREATE TABLE IF NOT EXISTS messages (" +
+ " msg_id SERIAL PRIMARY KEY, " +
+ " msg_text character varying, " +
+ " chat_id integer, " +
+ " timestamp timestamp, " +
+ " user_id integer " +
+ ");";
+ queryExecutor.updateQuery(sql);
+ }
+
+ public void createChats() throws SQLException {
+ String sql = "CREATE TABLE IF NOT EXISTS chats (" +
+ " user_id bigint, " +
+ " chat_id bigint " +
+ ");";
+ queryExecutor.updateQuery(sql);
+ }
+
+
+ public User getUser(String userName) throws Exception {
+
+ String sql = "SELECT user_id, user_login, user_pass " +
+ " FROM users " +
+ " WHERE user_login = ? ";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, userName);
+
+ return queryExecutor.execQuery(sql, prepared, resultSet -> {
+ if (resultSet.next()) {
+ String login = resultSet.getString("user_login");
+ String hash = resultSet.getString("user_pass");
+ Long id = resultSet.getLong("user_id");
+ User user1 = new User(login);
+ user1.setPass(hash);
+ user1.setId(id);
+ return user1;
+ }
+ return null;
+ });
+
+ }
+
+ public User getUserById(Long id) throws Exception {
+ String sql = "SELECT user_login, user_pass " +
+ " FROM users " +
+ " WHERE user_id = ? ";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, id);
+
+ User user = queryExecutor.execQuery(sql, prepared, resultSet -> {
+ if (resultSet.next()) {
+ String login = resultSet.getString("user_login");
+ String hash = resultSet.getString("user_pass");
+ User user1 = new User(login);
+ user1.setHash(hash);
+ return user1;
+ }
+ return null;
+ });
+ if (user != null) {
+ user.setId(id);
+ }
+ return user;
+ }
+
+ public void setNewPass(String login, String password) throws Exception {
+
+ String sqlUpdate = "UPDATE users " +
+ " SET user_pass = ? " +
+ " WHERE user_login = ? ;";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, login);
+ prepared.put(2, password);
+
+ queryExecutor.updateQuery(sqlUpdate, prepared);
+ }
+
+ public void addChat(Chat chat) {
+
+ String sqlInsert = "INSERT INTO chats (chat_id) VALUES " +
+ "(?) ;";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, chat.getCreatorId());
+
+ Long key = null;
+ try {
+ key = queryExecutor.updateQueryWithGeneratedKey(sqlInsert, prepared, "chat_id");
+ log.info("key = " + key);
+ } catch (SQLException sqlExc) {
+ log.error("Проблемы с sql запросом: " + sqlInsert + ", where ?=" + chat.getCreatorId(), sqlExc);
+ }
+
+ chat.setId(key);
+
+ }
+
+ public void addUserToChat(Long userId, Long chatId) {
+
+ String sqlInsert = "INSERT INTO chats (\"chat_id\", \"user_id\") " +
+ "VALUES " +
+ "(?, ?); ";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, chatId);
+ prepared.put(2, userId);
+
+ try {
+ queryExecutor.updateQuery(sqlInsert, prepared);
+ } catch (SQLException sqlExc) {
+ log.error("Ошибка запроса:\n" + sqlInsert + "\n" + "where ?=" + chatId + ", ?=" + userId);
+ log.error(String.valueOf(sqlExc));
+ }
+ }
+
+ public List getChatsByUserId(Long userId) {
+
+ String sql = "SELECT chat_id " +
+ "FROM chats " +
+ "WHERE user_id = ?";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, userId);
+
+ List chatList = null;
+ try {
+ chatList = queryExecutor.execQuery(sql, prepared, resultSet -> {
+ List chatList1 = new ArrayList<>();
+ while (resultSet.next()) {
+ chatList1.add(resultSet.getLong("chat_id"));
+ }
+ return chatList1;
+ });
+ } catch (SQLException sqlExc) {
+ log.error("Проблемы с sql запросом: " + sql + "?=" + userId, sqlExc);
+ }
+ return chatList;
+ }
+
+ public List getUsersByChatId(Long chatId) {
+ String sql = "SELECT user_id " +
+ "FROM chats " +
+ "WHERE chat_id = ?";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, chatId);
+
+ List userList = null;
+ try {
+ userList = queryExecutor.execQuery(sql, prepared, resultSet -> {
+ List userList1 = new ArrayList<>();
+ while (resultSet.next()) {
+ userList1.add(resultSet.getLong("user_id"));
+ }
+ return userList1;
+ });
+ } catch (SQLException sqlExc) {
+ log.error("Проблемы с sql запросом: " + sql + "?=" + chatId, sqlExc);
+ }
+ return userList;
+ }
+
+ public void addMessage(Long chatId, Message msg) {
+
+ String sqlInsert = "INSERT INTO messages " +
+ "(msg_text, chat_id, timestamp, user_id) VALUES " +
+ "(?, ?, ?, ?) ";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, msg.getMessage());
+ prepared.put(2, chatId);
+ prepared.put(3, msg.getTime());
+ prepared.put(4, msg.getSenderId());
+
+ try {
+ Long id = queryExecutor.updateQueryWithGeneratedKey(sqlInsert, prepared, "msg_id");
+ msg.setId(id);
+ } catch (SQLException sqlExc) {
+ log.error("Проблемы с sql запросом: " + sqlInsert +
+ prepared.get(1) +
+ prepared.get(2) +
+ prepared.get(3) +
+ prepared.get(4),
+ sqlExc);
+ }
+ }
+
+ public List getMessagesByChatId(Long chatId) {
+
+ String sql = "SELECT msg_id " +
+ "FROM messages " +
+ "WHERE chat_id = ?";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, chatId);
+
+ List msgList = null;
+ try {
+ msgList = queryExecutor.execQuery(sql, prepared, resultSet -> {
+ List msgList1 = new ArrayList<>();
+ while (resultSet.next()) {
+ msgList1.add(resultSet.getLong("msg_id"));
+ }
+ return msgList1;
+ });
+ } catch (SQLException sqlExc) {
+ log.error("Проблемы с sql запросом: " + sql + "?=" + chatId, sqlExc);
+ }
+ return msgList;
+
+ }
+
+ public Message getMessageById(Long messageId) {
+
+ String sql = "" +
+ "SELECT * " +
+ "FROM messages " +
+ "WHERE msg_id = ?";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, messageId);
+ Message msg = null;
+ try {
+ msg = queryExecutor.execQuery(sql, prepared, resultSet -> {
+ Message msg1 = new Message();
+ if (resultSet.next()) {
+ msg1.setId(resultSet.getLong("msg_id"));
+ msg1.setMessage(resultSet.getString("msg_text"));
+ msg1.setType(Type.MSG_TEXT);
+ msg1.setSenderId(resultSet.getLong("user_id"));
+ msg1.setTime(resultSet.getString("timestamp"));
+ return msg1;
+ }
+ return null;
+ });
+ } catch (SQLException sqlExc) {
+ log.error("Проблемы с sql запросом: " + sql + "where ?=" + messageId, sqlExc);
+ }
+ return msg;
+ }
+
+ public User addUser(String userName, String password) throws Exception {
+
+ User user;
+
+ String sqlInsert = "INSERT INTO users (\"user_login\", \"user_pass\") VAlUES " +
+ "(?, ?);";
+
+ Map prepared = new HashMap<>();
+ prepared.put(1, userName);
+ prepared.put(2, password);
+
+ Long id = queryExecutor.updateQueryWithGeneratedKey(sqlInsert, prepared, "user_id");
+
+ user = new User(userName);
+ user.setHash(password);
+ user.setId(id);
+ return user;
+
+ }
+
+ public void close() {
+ queryExecutor.close();
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/dbservice/executor/QueryExecutor.java b/src/main/java/arhangel/dim/core/dbservice/executor/QueryExecutor.java
new file mode 100644
index 0000000..c430de4
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/dbservice/executor/QueryExecutor.java
@@ -0,0 +1,115 @@
+package arhangel.dim.core.dbservice.executor;
+
+import org.slf4j.LoggerFactory;
+
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.sql.SQLException;
+import java.sql.Statement;
+import java.sql.ResultSet;
+import java.util.HashMap;
+import java.util.Map;
+
+public class QueryExecutor {
+ static org.slf4j.Logger log = LoggerFactory.getLogger(QueryExecutor.class);
+ Connection connection;
+
+ private Map execPrepareds = new HashMap<>();
+ private Map updatePrepareds = new HashMap<>();
+ private Map updateWithKeysPrepareds = new HashMap<>();
+
+ public void setConnection(Connection connection) {
+ this.connection = connection;
+ }
+
+ public T execQuery(String query, Map args, ResultHandler handler) throws SQLException {
+
+ PreparedStatement stmt;
+ if (execPrepareds.containsKey(query)) {
+ stmt = execPrepareds.get(query);
+ } else {
+ stmt = connection.prepareStatement(query);
+ execPrepareds.put(query, stmt);
+ }
+
+ stmt.clearParameters();
+
+ for (Map.Entry entry : args.entrySet()) {
+ stmt.setObject(entry.getKey(), entry.getValue());
+ }
+ ResultSet resultset = stmt.executeQuery();
+ T value = handler.handle(resultset);
+ resultset.close();
+ return value;
+ }
+
+ public void updateQuery(String query) throws SQLException {
+
+ Statement stmt = connection.createStatement();
+ stmt.execute(query);
+ stmt.close();
+
+ }
+
+ public void updateQuery(String query, Map args) throws SQLException {
+
+ PreparedStatement stmt;
+ if (updatePrepareds.containsKey(query)) {
+ stmt = updatePrepareds.get(query);
+ } else {
+ stmt = connection.prepareStatement(query);
+ updatePrepareds.put(query, stmt);
+ }
+
+ stmt.clearParameters();
+
+ for (Map.Entry entry : args.entrySet()) {
+ stmt.setObject(entry.getKey(), entry.getValue());
+ }
+ stmt.executeQuery();
+
+ }
+
+ public Long updateQueryWithGeneratedKey(String query, Map args, String keyName) throws SQLException {
+
+ PreparedStatement stmt;
+ if (updateWithKeysPrepareds.containsKey(query)) {
+ stmt = updateWithKeysPrepareds.get(query);
+ } else {
+ stmt = connection.prepareStatement(query, Statement.RETURN_GENERATED_KEYS);
+ updateWithKeysPrepareds.put(query, stmt);
+ }
+
+ stmt.clearParameters();
+
+ for (Map.Entry entry : args.entrySet()) {
+ stmt.setObject(entry.getKey(), entry.getValue());
+ }
+
+ stmt.executeUpdate();
+
+ ResultSet resultset = stmt.getGeneratedKeys();
+ resultset.next();
+ Long key = resultset.getLong(keyName);
+
+ resultset.close();
+
+ return key;
+ }
+
+ public void close() {
+ try {
+ for (Map.Entry entry : execPrepareds.entrySet()) {
+ entry.getValue().close();
+ }
+ for (Map.Entry entry : updatePrepareds.entrySet()) {
+ entry.getValue().close();
+ }
+ for (Map.Entry entry : updateWithKeysPrepareds.entrySet()) {
+ entry.getValue().close();
+ }
+ } catch (SQLException sqlExc) {
+ log.error("Ошибка закрытия sql", sqlExc);
+ }
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/dbservice/executor/ResultHandler.java b/src/main/java/arhangel/dim/core/dbservice/executor/ResultHandler.java
new file mode 100644
index 0000000..41efa34
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/dbservice/executor/ResultHandler.java
@@ -0,0 +1,8 @@
+package arhangel.dim.core.dbservice.executor;
+
+import java.sql.ResultSet;
+import java.sql.SQLException;
+
+public interface ResultHandler {
+ T handle(ResultSet resultSet) throws SQLException;
+}
diff --git a/src/main/java/arhangel/dim/core/messages/CommandExecutor.java b/src/main/java/arhangel/dim/core/messages/CommandExecutor.java
new file mode 100644
index 0000000..7517e03
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/messages/CommandExecutor.java
@@ -0,0 +1,30 @@
+package arhangel.dim.core.messages;
+
+import arhangel.dim.core.net.Session;
+
+import java.util.HashMap;
+import java.util.Map;
+
+public class CommandExecutor {
+ private Map commands;
+
+ public CommandExecutor() {
+ commands = new HashMap<>();
+ }
+
+ public CommandExecutor addCommand(Type type, Command command) {
+ commands.put(type, command);
+ return this;
+ }
+
+ public void handleMessage(Message message, Session session) throws CommandException {
+ Type messageType = message.getType();
+ if (messageType == null) {
+ throw new CommandException("Message type is null");
+ } else if (commands.containsKey(messageType)) {
+ commands.get(messageType).execute(session, message);
+ } else {
+ throw new CommandException("Undefined message type");
+ }
+ }
+}
\ No newline at end of file
diff --git a/src/main/java/arhangel/dim/core/messages/CreateChatMessage.java b/src/main/java/arhangel/dim/core/messages/CreateChatMessage.java
new file mode 100644
index 0000000..19fef8a
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/messages/CreateChatMessage.java
@@ -0,0 +1,47 @@
+package arhangel.dim.core.messages;
+
+import java.util.List;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+public class CreateChatMessage extends Message {
+
+ private List usersIds;
+
+ public CreateChatMessage() {
+ this.setType(Type.MSG_CHAT_CREATE);
+ }
+
+ public void setUsersIds(List userIds) {
+ this.usersIds = userIds;
+ }
+
+ public List getUsersIds() {
+ return usersIds;
+ }
+
+ @Override
+ public boolean equals(Object other) {
+ if (this == other) {
+ return true;
+ }
+ if (other == null || getClass() != other.getClass()) {
+ return false;
+ }
+
+ CreateChatMessage createChatMessage = (CreateChatMessage) other;
+
+ return Objects.equals(usersIds, createChatMessage.usersIds);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(super.hashCode(), getUsersIds());
+ }
+
+ @Override
+ public String toString() {
+ return "ChatCreateMessage: " +
+ String.join(",", usersIds.stream().map(Object::toString).collect(Collectors.toList()));
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/messages/HistChatMessage.java b/src/main/java/arhangel/dim/core/messages/HistChatMessage.java
new file mode 100644
index 0000000..486b91f
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/messages/HistChatMessage.java
@@ -0,0 +1,44 @@
+package arhangel.dim.core.messages;
+
+import java.util.Objects;
+
+public class HistChatMessage extends Message {
+ private Long chatId;
+
+ public HistChatMessage() {
+ this.setType(Type.MSG_CHAT_HIST);
+ }
+
+ public void setChatId(Long chatId) {
+ this.chatId = chatId;
+ }
+
+ public Long getChatId() {
+ return chatId;
+ }
+
+ @Override
+ public boolean equals(Object other) {
+ if (this == other) {
+ return true;
+ }
+ if (other == null || getClass() != other.getClass()) {
+ return false;
+ }
+
+ HistChatMessage histChatMessage = (HistChatMessage) other;
+
+ return Objects.equals(chatId, histChatMessage.chatId);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(chatId);
+ }
+
+ @Override
+ public String toString() {
+ return "Chat list: " +
+ chatId.toString();
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/messages/HistChatResultMessage.java b/src/main/java/arhangel/dim/core/messages/HistChatResultMessage.java
new file mode 100644
index 0000000..0f95526
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/messages/HistChatResultMessage.java
@@ -0,0 +1,58 @@
+package arhangel.dim.core.messages;
+
+import java.util.Objects;
+
+public class HistChatResultMessage extends Message {
+ private String history = "";
+ private long chatId;
+
+ public HistChatResultMessage() {
+ this.setType(Type.MSG_CHAT_HIST_RESULT);
+ }
+
+ public void setChatId(long id) {
+ this.chatId = id;
+ }
+
+ public long getChatId() {
+ return chatId;
+ }
+
+ public void setHistory(String history) {
+ this.history = history;
+ }
+
+ public String getHistory() {
+ return history;
+ }
+
+ public void addMsg(String msg) {
+ history = history + "\n" + msg;
+ }
+
+ @Override
+ public boolean equals(Object other) {
+ if (this == other) {
+ return true;
+ }
+ if (other == null || getClass() != other.getClass()) {
+ return false;
+ }
+ HistChatResultMessage histChatResultMessage = (HistChatResultMessage) other;
+ return Objects.equals(history, histChatResultMessage.history) &&
+ Objects.equals(chatId, histChatResultMessage.chatId);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(super.hashCode(), getHistory(), getChatId());
+ }
+
+ @Override
+ public String toString() {
+ return "Chat id =" +
+ chatId +
+ " history: " +
+ history;
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/messages/InfoMessage.java b/src/main/java/arhangel/dim/core/messages/InfoMessage.java
new file mode 100644
index 0000000..0fdf504
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/messages/InfoMessage.java
@@ -0,0 +1,54 @@
+package arhangel.dim.core.messages;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+public class InfoMessage extends Message {
+ public InfoMessage() {
+ this.setType(Type.MSG_INFO);
+ }
+
+ private String info = "";
+ private long userId;
+
+ public void setInfo(String info) {
+ this.info = info;
+ }
+
+ public String getInfo() {
+ return info;
+ }
+
+ public void setUserId(long userId) {
+ this.userId = userId;
+ }
+
+ public long getUserId() {
+ return userId;
+ }
+
+ @Override
+ public boolean equals(Object object) {
+ if (this == object) {
+ return true;
+ }
+ if (object == null || getClass() != object.getClass()) {
+ return false;
+ }
+
+ return Objects.equals(info, ((InfoMessage)object).info) &&
+ Objects.equals(userId, ((InfoMessage)object).userId);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(super.hashCode(), userId, info);
+ }
+
+ @Override
+ public String toString() {
+ return "User id: " + userId + " info: " + info;
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/messages/ListChatMessage.java b/src/main/java/arhangel/dim/core/messages/ListChatMessage.java
new file mode 100644
index 0000000..0dafa22
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/messages/ListChatMessage.java
@@ -0,0 +1,45 @@
+package arhangel.dim.core.messages;
+
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+public class ListChatMessage extends Message {
+ String chatsList;
+
+ public ListChatMessage() {
+ this.setType(Type.MSG_CHAT_LIST);
+ }
+
+ public void setChatsList(String chatsList) {
+ this.chatsList = chatsList;
+ }
+
+ public String getChatsList() {
+ return chatsList;
+ }
+
+ @Override
+ public boolean equals(Object other) {
+ if (this == other) {
+ return true;
+ }
+ if (other == null || getClass() != other.getClass()) {
+ return false;
+ }
+
+ ListChatMessage listChatMessage = (ListChatMessage) other;
+
+ return Objects.equals(chatsList, listChatMessage.chatsList);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(super.hashCode(), getChatsList());
+ }
+
+ @Override
+ public String toString() {
+ return "Chat list: " +
+ chatsList;
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/messages/ListChatResultMessage.java b/src/main/java/arhangel/dim/core/messages/ListChatResultMessage.java
new file mode 100644
index 0000000..d10bc9d
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/messages/ListChatResultMessage.java
@@ -0,0 +1,48 @@
+package arhangel.dim.core.messages;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+public class ListChatResultMessage extends Message {
+
+ public ListChatResultMessage() {
+ this.setType(Type.MSG_CHAT_LIST_RESULT);
+ }
+
+ private List chatsIds = new ArrayList<>();
+
+ public List getChatIds() {
+ return chatsIds;
+ }
+
+ public void setChatIds(List chatIds) {
+ this.chatsIds = chatIds;
+ }
+
+ @Override
+ public boolean equals(Object object) {
+ if (this == object) {
+ return true;
+ }
+ if (object == null || getClass() != object.getClass()) {
+ return false;
+ }
+
+ ListChatResultMessage listChatResultMessage = (ListChatResultMessage) object;
+
+ return Objects.equals(chatsIds, listChatResultMessage.chatsIds);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(super.hashCode(), chatsIds);
+ }
+
+ @Override
+ public String toString() {
+ return "Chat list: " +
+ String.join(",", chatsIds.stream().map(Object::toString).collect(Collectors.toList()));
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/messages/LoginMessage.java b/src/main/java/arhangel/dim/core/messages/LoginMessage.java
new file mode 100644
index 0000000..b8b8048
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/messages/LoginMessage.java
@@ -0,0 +1,53 @@
+package arhangel.dim.core.messages;
+
+import java.util.Objects;
+
+public class LoginMessage extends Message {
+ private String login;
+ private String password;
+
+ public LoginMessage() {
+ this.setType(Type.MSG_LOGIN);
+ }
+
+ public String getLogin() {
+ return this.login;
+ }
+
+ public String getPassword() {
+ return this.password;
+ }
+
+ public void setLogin(String login) {
+ this.login = login;
+ }
+
+ public void setPassword(String password) {
+ this.password = password;
+ }
+
+ @Override
+ public boolean equals(Object other) {
+ if (this == other) {
+ return true;
+ }
+ if (other == null || getClass() != other.getClass()) {
+ return false;
+ }
+ LoginMessage loginMessage = (LoginMessage) other;
+ return Objects.equals(login, loginMessage.login) &&
+ Objects.equals(password, loginMessage.password);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(super.hashCode(), getLogin(), getPassword());
+ }
+
+ @Override
+ public String toString() {
+ return "LoginMessage: " +
+ "login: " + getLogin() + '\n' +
+ "password: " + getPassword() + '\n';
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/messages/Message.java b/src/main/java/arhangel/dim/core/messages/Message.java
index 65486c8..df89a9b 100644
--- a/src/main/java/arhangel/dim/core/messages/Message.java
+++ b/src/main/java/arhangel/dim/core/messages/Message.java
@@ -1,16 +1,24 @@
package arhangel.dim.core.messages;
import java.io.Serializable;
+import java.util.Date;
import java.util.Objects;
/**
* Базовый класс для всех сообщений
*/
-public abstract class Message implements Serializable {
+public class Message implements Serializable {
private Long id;
private Long senderId;
private Type type;
+ private String message;
+ private String messageId;
+ private String time;
+
+ public Message() {
+ time = new Date().toString();
+ }
public Long getId() {
return id;
@@ -36,6 +44,22 @@ public void setType(Type type) {
this.type = type;
}
+ public void setMessage(String message) {
+ this.message = message;
+ }
+
+ public String getMessage() {
+ return message;
+ }
+
+ public void setTime(String time) {
+ this.time = time;
+ }
+
+ public String getTime() {
+ return time;
+ }
+
@Override
public boolean equals(Object other) {
if (this == other) {
diff --git a/src/main/java/arhangel/dim/core/messages/StatusMessage.java b/src/main/java/arhangel/dim/core/messages/StatusMessage.java
new file mode 100644
index 0000000..3eb2e19
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/messages/StatusMessage.java
@@ -0,0 +1,43 @@
+package arhangel.dim.core.messages;
+
+import java.util.Objects;
+
+public class StatusMessage extends Message {
+ private String status;
+
+ public StatusMessage() {
+ this.setType(Type.MSG_STATUS);
+ }
+
+ public String getStatus() {
+ return status;
+ }
+
+ public void setStatus(String status) {
+ this.status = status;
+ }
+
+ @Override
+ public boolean equals(Object other) {
+ if (this == other) {
+ return true;
+ }
+ if (other == null || getClass() != other.getClass()) {
+ return false;
+ }
+
+ StatusMessage statusMessage = (StatusMessage) other;
+
+ return Objects.equals(status, statusMessage.status);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(super.hashCode(), getStatus());
+ }
+
+ @Override
+ public String toString() {
+ return "Status: " + status;
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/messages/TextMessage.java b/src/main/java/arhangel/dim/core/messages/TextMessage.java
index 59f06b5..25c52e4 100644
--- a/src/main/java/arhangel/dim/core/messages/TextMessage.java
+++ b/src/main/java/arhangel/dim/core/messages/TextMessage.java
@@ -7,15 +7,28 @@
*/
public class TextMessage extends Message {
private String text;
+ private long chatId;
- public String getText() {
- return text;
+ public TextMessage() {
+ this.setType(Type.MSG_TEXT);
}
public void setText(String text) {
this.text = text;
}
+ public String getText() {
+ return text;
+ }
+
+ public void setChatId(long chatId) {
+ this.chatId = chatId;
+ }
+
+ public long getChatId() {
+ return chatId;
+ }
+
@Override
public boolean equals(Object other) {
if (this == other) {
@@ -38,8 +51,6 @@ public int hashCode() {
@Override
public String toString() {
- return "TextMessage{" +
- "text='" + text + '\'' +
- '}';
+ return "Message: " + text;
}
}
diff --git a/src/main/java/arhangel/dim/core/net/ConnectionHandler.java b/src/main/java/arhangel/dim/core/net/ConnectionHandler.java
index c5df653..3ce84d6 100644
--- a/src/main/java/arhangel/dim/core/net/ConnectionHandler.java
+++ b/src/main/java/arhangel/dim/core/net/ConnectionHandler.java
@@ -26,6 +26,6 @@ public interface ConnectionHandler {
/**
* Молча (без проброса ошибок) закрываем соединение и освобождаем ресурсы
*/
- void close();
+ void close() throws IOException;
}
diff --git a/src/main/java/arhangel/dim/core/net/MessageListener.java b/src/main/java/arhangel/dim/core/net/MessageListener.java
new file mode 100644
index 0000000..1e89169
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/net/MessageListener.java
@@ -0,0 +1,7 @@
+package arhangel.dim.core.net;
+
+import arhangel.dim.core.messages.Message;
+
+public interface MessageListener {
+ void onMessage(Message message, long id);
+}
diff --git a/src/main/java/arhangel/dim/core/net/SerialProtocol.java b/src/main/java/arhangel/dim/core/net/SerialProtocol.java
new file mode 100644
index 0000000..1741807
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/net/SerialProtocol.java
@@ -0,0 +1,60 @@
+package arhangel.dim.core.net;
+
+import arhangel.dim.core.messages.Message;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
+import java.io.ObjectInput;
+import java.io.ObjectInputStream;
+import java.io.ObjectOutput;
+import java.io.ObjectOutputStream;
+
+public class SerialProtocol implements Protocol {
+
+ static Logger log = LoggerFactory.getLogger(StringProtocol.class);
+
+ @Override
+ public Message decode(byte[] bytes) throws ProtocolException {
+
+ ByteArrayInputStream inputStream = new ByteArrayInputStream(bytes);
+ ObjectInput in = null;
+ Message msg = null;
+ try {
+ in = new ObjectInputStream(inputStream);
+ msg = (Message) in.readObject();
+ } catch (IOException | ClassNotFoundException e) {
+ e.printStackTrace();
+ throw new ProtocolException(e);
+ }
+ return msg;
+ }
+
+ @Override
+ public byte[] encode(Message msg) throws ProtocolException {
+ ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
+ ObjectOutput out = null;
+ byte[] bytes = null;
+ try {
+ out = new ObjectOutputStream(outputStream);
+ out.writeObject(msg);
+ bytes = outputStream.toByteArray();
+ } catch (IOException e) {
+ e.printStackTrace();
+ throw new ProtocolException(e);
+ }
+ return bytes;
+ }
+
+ private Long parseLong(String str) {
+ try {
+ return Long.parseLong(str);
+ } catch (Exception e) {
+ // who care
+ }
+ return null;
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/net/Session.java b/src/main/java/arhangel/dim/core/net/Session.java
index f9a161f..4ac5a08 100644
--- a/src/main/java/arhangel/dim/core/net/Session.java
+++ b/src/main/java/arhangel/dim/core/net/Session.java
@@ -4,25 +4,43 @@
import java.io.InputStream;
import java.io.OutputStream;
import java.net.Socket;
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.SQLException;
+import java.util.Arrays;
import arhangel.dim.core.User;
+import arhangel.dim.core.command.CreateChatCommand;
+import arhangel.dim.core.command.HistChatCommand;
+import arhangel.dim.core.command.InfoCommand;
+import arhangel.dim.core.command.ListChatCommand;
+import arhangel.dim.core.command.LoginCommand;
+import arhangel.dim.core.command.TextCommand;
+import arhangel.dim.core.dbservice.dao.UsersDao;
+import arhangel.dim.core.messages.CommandException;
+import arhangel.dim.core.messages.CommandExecutor;
import arhangel.dim.core.messages.Message;
+import arhangel.dim.core.messages.Type;
+import arhangel.dim.core.store.MessageStore;
+import arhangel.dim.core.store.MessageStoreImpl;
+import arhangel.dim.core.store.UserStore;
+import arhangel.dim.core.store.UserStoreImpl;
+import arhangel.dim.server.Server;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
/**
* Здесь храним всю информацию, связанную с отдельным клиентом.
* - объект User - описание пользователя
* - сокеты на чтение/запись данных в канал пользователя
*/
-public class Session implements ConnectionHandler {
+public class Session implements ConnectionHandler, Runnable, AutoCloseable {
- /**
- * Пользователь сессии, пока не прошел логин, user == null
- * После логина устанавливается реальный пользователь
- */
private User user;
// сокет на клиента
private Socket socket;
+ private Server server;
/**
* С каждым сокетом связано 2 канала in/out
@@ -30,18 +48,107 @@ public class Session implements ConnectionHandler {
private InputStream in;
private OutputStream out;
+ private Connection connection;
+ private MessageStoreImpl messageStore;
+ private UserStoreImpl userStore;
+ private UsersDao usersDao;
+ private Protocol protocol;
+
+ private CommandExecutor commandExecutor;
+
+ private static Logger log = LoggerFactory.getLogger(Session.class);
+
+ public Session(Socket socket, Protocol protocol, CommandExecutor commandExecutor, UsersDao usersDao) throws IOException, SQLException, ClassNotFoundException {
+ this.socket = socket;
+ in = socket.getInputStream();
+ out = socket.getOutputStream();
+ this.protocol = protocol;
+ this.server = server;
+
+ this.usersDao = usersDao;
+
+ this.messageStore = new MessageStoreImpl(usersDao);
+ this.userStore = new UserStoreImpl(usersDao);
+
+ this.commandExecutor = commandExecutor;
+
+ }
+
@Override
public void send(Message msg) throws ProtocolException, IOException {
- // TODO: Отправить клиенту сообщение
+ log.info("< Sending new message >");
+ out.write(protocol.encode(msg));
+ out.flush();
}
@Override
public void onMessage(Message msg) {
- // TODO: Пришло некое сообщение от клиента, его нужно обработать
+ try {
+ commandExecutor.handleMessage(msg, this);
+ } catch (CommandException e) {
+ e.printStackTrace();
+ }
}
@Override
public void close() {
- // TODO: закрыть in/out каналы и сокет. Освободить другие ресурсы, если необходимо
+ try {
+ in.close();
+ out.close();
+ connection.close();
+ Thread.currentThread().interrupt();
+ log.info("Session closed");
+ } catch (IOException | SQLException e) {
+ log.error(e.getMessage());
+ }
+ }
+
+ @Override
+ public void run() {
+ final byte[] buf = new byte[1024 * 64];
+ log.info("Running client section");
+ while (!Thread.currentThread().isInterrupted()) {
+ try {
+ int read = in.read(buf);
+ log.info("Read data: " + read);
+ if (read > 0) {
+ log.info("< New Message received >");
+ Message message = protocol.decode(Arrays.copyOf(buf, read));
+ onMessage(message);
+ } else {
+ if (read == -1) {
+ close();
+ }
+ }
+ } catch (Exception e) {
+ log.error("Failed to process client connection:", e);
+ Thread.currentThread().interrupt();
+ }
+ }
}
+
+ public User getUser() {
+ return user;
+ }
+
+ public void setUser(User user) {
+ this.user = user;
+ }
+
+ public MessageStoreImpl getMessageStore() {
+ return messageStore;
+ }
+
+ public void setMessageStore(MessageStoreImpl messageStore) {
+ this.messageStore = messageStore;
+ }
+
+ public UserStoreImpl getUserStore() {
+ return userStore;
+ }
+
+ public void setUserStore(UserStoreImpl userStore) {
+ this.userStore = userStore;
+ }
+
}
diff --git a/src/main/java/arhangel/dim/core/net/StringProtocol.java b/src/main/java/arhangel/dim/core/net/StringProtocol.java
index d4cd40a..5d0fe74 100644
--- a/src/main/java/arhangel/dim/core/net/StringProtocol.java
+++ b/src/main/java/arhangel/dim/core/net/StringProtocol.java
@@ -1,11 +1,13 @@
package arhangel.dim.core.net;
+import arhangel.dim.core.messages.*;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import arhangel.dim.core.messages.Message;
-import arhangel.dim.core.messages.TextMessage;
-import arhangel.dim.core.messages.Type;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.stream.Collectors;
/**
* Простейший протокол передачи данных
@@ -19,16 +21,63 @@ public class StringProtocol implements Protocol {
@Override
public Message decode(byte[] bytes) throws ProtocolException {
String str = new String(bytes);
- log.info("decoded: {}", str);
+ log.info("DECODED MESSAGE: {}", str);
String[] tokens = str.split(DELIMITER);
Type type = Type.valueOf(tokens[0]);
switch (type) {
case MSG_TEXT:
TextMessage textMsg = new TextMessage();
textMsg.setSenderId(parseLong(tokens[1]));
- textMsg.setText(tokens[2]);
- textMsg.setType(type);
+ textMsg.setChatId(parseLong(tokens[2]));
+ textMsg.setText(tokens[3]);
return textMsg;
+ case MSG_INFO:
+ InfoMessage infoMessage = new InfoMessage();
+ infoMessage.setSenderId(parseLong(tokens[1]));
+ infoMessage.setUserId(parseLong(tokens[2]));
+ infoMessage.setInfo(tokens[3]);
+ return infoMessage;
+ case MSG_STATUS:
+ StatusMessage statusMessage = new StatusMessage();
+ statusMessage.setSenderId(parseLong(tokens[1]));
+ statusMessage.setStatus(tokens[2]);
+ return statusMessage;
+ case MSG_LOGIN:
+ LoginMessage loginMessage = new LoginMessage();
+ loginMessage.setSenderId(parseLong(tokens[1]));
+ loginMessage.setLogin(tokens[2]);
+ loginMessage.setPassword(tokens[3]);
+ return loginMessage;
+ case MSG_CHAT_LIST:
+ ListChatMessage listChatMessage = new ListChatMessage();
+ listChatMessage.setSenderId(parseLong(tokens[1]));
+ return listChatMessage;
+ case MSG_CHAT_CREATE:
+ CreateChatMessage chatCreateMessage = new CreateChatMessage();
+ chatCreateMessage.setSenderId(parseLong(tokens[1]));
+ String[] userIdsStr = tokens[2].split(",");
+ List userIds = new ArrayList();
+ for (int i = 0; i < userIdsStr.length; ++i) {
+ userIds.add(Long.parseLong(userIdsStr[i]));
+ }
+ chatCreateMessage.setUsersIds(userIds);
+ return chatCreateMessage;
+ case MSG_CHAT_LIST_RESULT:
+ ListChatResultMessage chatListResultMessage = new ListChatResultMessage();
+ chatListResultMessage.setSenderId(parseLong(tokens[1]));
+ if (tokens.length >= 3) {
+ chatListResultMessage.setChatIds(Arrays.asList(tokens[2].split(",")).stream()
+ .map(this::parseLong)
+ .collect(Collectors.toList()));
+ }
+ return chatListResultMessage;
+ case MSG_CHAT_HIST_RESULT:
+ HistChatResultMessage chatHistResultMessage = new HistChatResultMessage();
+ chatHistResultMessage.setSenderId(parseLong(tokens[1]));
+ if (tokens.length >= 3) {
+ chatHistResultMessage.setHistory(tokens[2]);
+ }
+ return chatHistResultMessage;
default:
throw new ProtocolException("Invalid type: " + type);
}
@@ -43,14 +92,56 @@ public byte[] encode(Message msg) throws ProtocolException {
case MSG_TEXT:
TextMessage sendMessage = (TextMessage) msg;
builder.append(String.valueOf(sendMessage.getSenderId())).append(DELIMITER);
+ builder.append(sendMessage.getChatId()).append(DELIMITER);
builder.append(sendMessage.getText()).append(DELIMITER);
break;
+ case MSG_LOGIN:
+ LoginMessage loginMessage = (LoginMessage) msg;
+ builder.append(String.valueOf(loginMessage.getSenderId())).append(DELIMITER);
+ builder.append(loginMessage.getLogin()).append(DELIMITER);
+ builder.append(loginMessage.getPassword()).append(DELIMITER);
+ break;
+ case MSG_CHAT_CREATE:
+ CreateChatMessage chatCreateMessage = (CreateChatMessage) msg;
+ List userIds = chatCreateMessage.getUsersIds();
+ List userIdsStr = new ArrayList<>();
+ for (int i = 0; i < userIds.size(); ++i) {
+ userIdsStr.add(String.valueOf(userIds.get(i)));
+ }
+ builder.append(String.valueOf(chatCreateMessage.getSenderId())).append(DELIMITER);
+ builder.append(String.join(",",
+ userIdsStr)).append(DELIMITER);
+ break;
+ case MSG_STATUS:
+ StatusMessage statusMessage = (StatusMessage) msg;
+ builder.append(String.valueOf(statusMessage.getSenderId())).append(DELIMITER);
+ builder.append(statusMessage.getStatus()).append(DELIMITER);
+ break;
+ case MSG_CHAT_LIST_RESULT:
+ ListChatResultMessage listChatResultMessage = (ListChatResultMessage) msg;
+ builder.append(String.valueOf(listChatResultMessage.getSenderId())).append(DELIMITER);
+ builder.append(String.join(",",
+ listChatResultMessage.getChatIds().stream()
+ .map(Object::toString)
+ .collect(Collectors.toList()))).append(DELIMITER);
+ break;
+ case MSG_INFO:
+ InfoMessage infoMessage = (InfoMessage) msg;
+ builder.append(String.valueOf(infoMessage.getSenderId())).append(DELIMITER);
+ builder.append(infoMessage.getUserId()).append(DELIMITER);
+ builder.append(infoMessage.getInfo()).append(DELIMITER);
+ break;
+ case MSG_CHAT_LIST:
+ ListChatMessage listChatMessage = new ListChatMessage();
+ builder.append(String.valueOf(listChatMessage.getSenderId())).append(DELIMITER);
+ builder.append(listChatMessage.getChatsList()).append(DELIMITER);
+ break;
default:
throw new ProtocolException("Invalid type: " + type);
}
- log.info("encoded: {}", builder.toString());
+ log.info("ENCODED MESSAGE: {}", builder.toString());
return builder.toString().getBytes();
}
diff --git a/src/main/java/arhangel/dim/core/service/AuthorizationService.java b/src/main/java/arhangel/dim/core/service/AuthorizationService.java
new file mode 100644
index 0000000..6278ddd
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/service/AuthorizationService.java
@@ -0,0 +1,52 @@
+package arhangel.dim.core.service;
+
+import arhangel.dim.core.User;
+import arhangel.dim.core.store.UserStoreImpl;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.math.BigInteger;
+import java.security.MessageDigest;
+
+public class AuthorizationService {
+ private UserStoreImpl store;
+ private static MessageDigest md;
+ private boolean success = false;
+
+ User user = null;
+ private static Logger log = LoggerFactory.getLogger(AuthorizationService.class);
+
+ public AuthorizationService(UserStoreImpl store, InputStream is, OutputStream os) {
+ this.store = store;
+ }
+
+ public static String calculateHash(final String str) {
+
+ /* try {
+ if (md == null) {
+ md = MessageDigest.getInstance("SHA-256");
+ }
+ md.update(str.getBytes("UTF-8"));
+ } catch (Exception exc) {
+ System.err.println("can't get instance of hashing algo SHA-256");
+ }
+ return md.digest();
+ */
+ return str;
+ }
+
+ public static boolean isCorrect(User user, String word) {
+
+ if (word == null) {
+ return false;
+ }
+
+ String hash = user.getHash();
+ String newHash = AuthorizationService.calculateHash(word);
+ log.info(hash);
+ log.info(newHash);
+ return hash.equals(newHash);
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/store/MessageStore.java b/src/main/java/arhangel/dim/core/store/MessageStore.java
index e16bed8..e814ec8 100644
--- a/src/main/java/arhangel/dim/core/store/MessageStore.java
+++ b/src/main/java/arhangel/dim/core/store/MessageStore.java
@@ -40,5 +40,7 @@ public interface MessageStore {
*/
void addUserToChat(Long userId, Long chatId);
+ public Chat createChat(Long chatCreator, List userIds);
+
}
diff --git a/src/main/java/arhangel/dim/core/store/MessageStoreImpl.java b/src/main/java/arhangel/dim/core/store/MessageStoreImpl.java
new file mode 100644
index 0000000..b77011b
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/store/MessageStoreImpl.java
@@ -0,0 +1,87 @@
+package arhangel.dim.core.store;
+
+import arhangel.dim.core.Chat;
+import arhangel.dim.core.dbservice.dao.UsersDao;
+import arhangel.dim.core.messages.Message;
+
+import java.util.List;
+
+public class MessageStoreImpl implements MessageStore {
+ UsersDao usersDao;
+
+ public MessageStoreImpl(UsersDao usersDao) {
+ this.usersDao = usersDao;
+ }
+
+ public List getUsersByChatId(Long chatId) {
+ List users = usersDao.getUsersByChatId(chatId);
+ return users;
+ }
+
+ @Override
+ public List getChatsByUserId(Long userId) {
+
+ List chats = usersDao.getChatsByUserId(userId);
+ return chats;
+
+ }
+
+ @Override
+ public Chat getChatById(Long chatId) {
+
+ Chat chat = new Chat();
+ List userIds = getUsersByChatId(chatId);
+ List messagesIds = getMessagesFromChat(chatId);
+
+ for (Long userId : userIds) {
+ chat.addParticipant(userId);
+ }
+
+ /*for (Long messageId : messagesIds) {
+ chat.addMessage(messageId);
+ }*/
+
+ chat.setId(chatId);
+
+ return chat;
+
+ }
+
+ @Override
+ public List getMessagesFromChat(Long chatId) {
+ List messageList = usersDao.getMessagesByChatId(chatId);
+ return messageList;
+ }
+
+ @Override
+ public Message getMessageById(Long messageId) {
+
+ return usersDao.getMessageById(messageId);
+
+ }
+
+ @Override
+ public void addMessage(Long chatId, Message message) {
+ usersDao.addMessage(chatId, message);
+ }
+
+ @Override
+ public void addUserToChat(Long userId, Long chatId) {
+ usersDao.addUserToChat(userId, chatId);
+ }
+
+ @Override
+ public Chat createChat(Long chatCreator, List userIds) {
+ Chat chat = new Chat();
+ chat.setCreatorId(chatCreator);
+
+ usersDao.addChat(chat);
+
+ for (Long id : userIds) {
+ chat.addParticipant(id);
+ addUserToChat(id, chat.getId());
+ }
+
+ return chat;
+ }
+}
diff --git a/src/main/java/arhangel/dim/core/store/UserStoreImpl.java b/src/main/java/arhangel/dim/core/store/UserStoreImpl.java
new file mode 100644
index 0000000..fad6f2c
--- /dev/null
+++ b/src/main/java/arhangel/dim/core/store/UserStoreImpl.java
@@ -0,0 +1,95 @@
+package arhangel.dim.core.store;
+
+import arhangel.dim.core.User;
+import arhangel.dim.core.dbservice.dao.UsersDao;
+import arhangel.dim.core.service.AuthorizationService;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.HashMap;
+import java.util.Map;
+
+public class UserStoreImpl implements UserStore {
+ private static Logger log = LoggerFactory.getLogger(UserStore.class);
+
+ private Map users;
+ private UsersDao usersDao;
+ private Map userLogins = new HashMap<>();
+
+ public UserStoreImpl(UsersDao usersDao) {
+ this.usersDao = usersDao;
+ }
+
+ public boolean isUserExist(String login) {
+ try {
+ return (usersDao.getUser(login) != null);
+ } catch (Exception e) {
+ log.error("Невозможно получить данные пользователя:" + login, e);
+ }
+ return false;
+ }
+
+ // Добавить пользователя в хранилище
+ @Override
+ public User addUser(User user) {
+
+
+ try {
+ usersDao.addUser(user.getName(), user.getHash());
+ } catch (Exception ioExc) {
+ log.error(ioExc.getMessage());
+ return null;
+ }
+ return user;
+ }
+
+ // Получить пользователя по имени и паролю
+ @Override
+ public User getUser(String name, String pass) {
+
+ User user;
+
+ try {
+ user = usersDao.getUser(name);
+ } catch (Exception e) {
+ log.error("Проблема с получением данных пользователя: " + name, e);
+ return null;
+ }
+
+ if (user == null) {
+ return null;
+ }
+
+ if (AuthorizationService.isCorrect(user, pass)) {
+ return user;
+ }
+ return null;
+ }
+
+ @Override
+ public User getUserById(Long id) {
+
+ try {
+ return usersDao.getUserById(id);
+ } catch (Exception e) {
+ log.error("Проблема с получением данных пользователя с id: " + id, e);
+ }
+ return null;
+ }
+
+ @Override
+ public User updateUser(User user) {
+
+ try {
+ usersDao.setNewPass(user.getName(), user.getHash());
+ } catch (Exception e) {
+ log.error("Проблема с изменением данных пользователя: " + user.getName(), e);
+ }
+ return user;
+ }
+
+ public void close() {
+ usersDao.close();
+ }
+
+}
diff --git a/src/main/java/arhangel/dim/lections/threads/queueu/BlockingQueue.java b/src/main/java/arhangel/dim/lections/threads/queueu/BlockingQueue.java
new file mode 100644
index 0000000..27eb325
--- /dev/null
+++ b/src/main/java/arhangel/dim/lections/threads/queueu/BlockingQueue.java
@@ -0,0 +1,19 @@
+package arhangel.dim.lections.threads.queueu;
+
+/**
+ *
+ */
+public interface BlockingQueue {
+
+ /**
+ *
+ * @param e the element to add
+ */
+ void put(E elem) throws InterruptedException;
+
+ /**
+ *
+ * @return the head element
+ */
+ E take() throws InterruptedException;
+}
diff --git a/src/main/java/arhangel/dim/lections/threads/queueu/ListBlockingQueue.java b/src/main/java/arhangel/dim/lections/threads/queueu/ListBlockingQueue.java
index ab3d557..6b8b083 100644
--- a/src/main/java/arhangel/dim/lections/threads/queueu/ListBlockingQueue.java
+++ b/src/main/java/arhangel/dim/lections/threads/queueu/ListBlockingQueue.java
@@ -3,9 +3,9 @@
import java.util.LinkedList;
/**
- * TODO: implements on lesson
+ *
*/
-public class ListBlockingQueue {
+public class ListBlockingQueue implements BlockingQueue {
public static final int DEFAULT_CAPACITY = 10;
@@ -20,12 +20,22 @@ public ListBlockingQueue(int capacity) {
this.capacity = capacity;
}
+ @Override
public void put(E elem) throws InterruptedException {
+ /*
+ Внутри критической секции пытаемся поместить элемент в очередь
+ Если очередь полная isFull==true то блокируемся на wait() пока условие не будет выполнено
+ */
}
+ @Override
public E take() throws InterruptedException {
+ /*
+ Внутри критической секции пытаемся достать элемент из очереди
+ Если очередь пустая isEmpty == true то блокируемся на wait() пока условие не будет выполнено
+ */
return null;
}
diff --git a/src/main/java/arhangel/dim/lections/threads/queueu/ProducerConsumer.java b/src/main/java/arhangel/dim/lections/threads/queueu/ProducerConsumer.java
index 5147ea3..1b60c9c 100644
--- a/src/main/java/arhangel/dim/lections/threads/queueu/ProducerConsumer.java
+++ b/src/main/java/arhangel/dim/lections/threads/queueu/ProducerConsumer.java
@@ -1,8 +1,5 @@
package arhangel.dim.lections.threads.queueu;
-import java.util.ArrayList;
-import java.util.List;
-
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -10,7 +7,7 @@ public class ProducerConsumer {
static Logger log = LoggerFactory.getLogger(ProducerConsumer.class);
- static List list = new ArrayList<>();
+ static boolean isReady = false;
static class Producer extends Thread {
private final Object lock;
@@ -27,17 +24,14 @@ public void run() {
} catch (InterruptedException e) {
e.printStackTrace();
}
- log.info("[PRODUCER] Data prepared.");
+ log.info("[PRODUCER] Data prepared. Notify All!");
+ isReady = true;
synchronized (lock) {
- list.add("READY");
+ lock.notifyAll();
}
-
-// synchronized (lock) {
-// list.add("READY");
-// lock.notifyAll();
-// }
}
+
}
static class Consumer extends Thread {
@@ -50,39 +44,23 @@ public Consumer(Object lock) {
@Override
public void run() {
- // 1) Ждем данные, busy loop
- while (list.isEmpty()) {
- // 2) Данные готовы, захватываем КС
- synchronized (lock) {
- // 3) Проверяем, что данные никто не поменял, иначе отпускаем блокировку
- if (!list.isEmpty()) {
- // 4) Мы владеем данными в контексте КС, можно изменять
- System.out.println("Processing data: " + list.get(0));
+ synchronized (lock) {
+ log.info("[CONSUMER] Waiting for data...");
+
+ // Если данные еще не готовы
+ while (!isReady) {
+ try {
+ // ждем
+ lock.wait();
+ // как только пробудились, заново проверяем состояние данных
+ // если они не готовы (или кто-то уже их поменял), то снова ждем
+ } catch (InterruptedException e) {
+ e.printStackTrace();
}
}
- }
-// // 1) Захватываем блокировку
-// synchronized (lock) {
-// try {
-// log.info("[CONSUMER] Waiting for data...");
-//
-// // 1) Если данные еще не готовы - ждем. Но wait() отпускает блокировку, при этом текущий поток
-// // переходит в состояние WAITING и помещается в wait set монитора
-// while (!isReady) {
-// lock.wait();
-// // как только пробудились, заново проверяем состояние данных
-// // если они не готовы (или кто-то уже их поменял), то снова ждем
-// }
-//
-// // 2) Данные готовы, мы внутри КС - можно изменять данные
-// System.out.println("Processing data.");
-// } catch (InterruptedException e) {
-// e.printStackTrace();
-// }
-//
-// log.info("[CONSUMER] Data received");
-// }
+ log.info("[CONSUMER] Data received");
+ }
}
}
diff --git a/src/main/java/arhangel/dim/server/Server.java b/src/main/java/arhangel/dim/server/Server.java
index 92452b8..384a7dd 100644
--- a/src/main/java/arhangel/dim/server/Server.java
+++ b/src/main/java/arhangel/dim/server/Server.java
@@ -1,21 +1,177 @@
package arhangel.dim.server;
+import arhangel.dim.core.command.CreateChatCommand;
+import arhangel.dim.core.command.HistChatCommand;
+import arhangel.dim.core.command.InfoCommand;
+import arhangel.dim.core.command.ListChatCommand;
+import arhangel.dim.core.command.LoginCommand;
+import arhangel.dim.core.command.TextCommand;
+import arhangel.dim.core.dbservice.dao.UsersDao;
+import arhangel.dim.core.messages.CommandExecutor;
+import arhangel.dim.core.messages.Type;
import arhangel.dim.core.net.Protocol;
+import arhangel.dim.core.net.Session;
+import arhangel.dim.core.net.StringProtocol;
+import arhangel.dim.core.store.MessageStore;
+import arhangel.dim.core.store.MessageStoreImpl;
+import arhangel.dim.core.store.UserStore;
+import arhangel.dim.core.store.UserStoreImpl;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.net.ServerSocket;
+import java.net.Socket;
+import java.net.SocketException;
+import java.sql.SQLException;
+import java.util.Set;
+import java.util.HashSet;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
/**
* Основной класс для сервера сообщений
*/
public class Server {
+ private static Logger log = LoggerFactory.getLogger(Server.class);
+
public static final int DEFAULT_MAX_CONNECT = 16;
- // Засетить из конфига
private int port;
private Protocol protocol;
private int maxConnection = DEFAULT_MAX_CONNECT;
+ private final ExecutorService threadPool = Executors.newCachedThreadPool();
+
+ private ServerSocket serverSocket;
+ private Set sessions = new HashSet<>();
+ private CommandExecutor commandExecutor = new CommandExecutor();
+
+ private UsersDao usersDao = new UsersDao();
+
+ private MessageStore messageStore = new MessageStoreImpl(usersDao);
+ private UserStore userStore = new UserStoreImpl(usersDao);
+
+ public Server(int port) throws IOException {
+ try {
+ this.port = port;
+ serverSocket = new ServerSocket(port);
+ usersDao.init();
+
+ this.commandExecutor.addCommand(Type.MSG_LOGIN, new LoginCommand(this));
+ this.commandExecutor.addCommand(Type.MSG_INFO, new InfoCommand(this));
+ this.commandExecutor.addCommand(Type.MSG_CHAT_CREATE, new CreateChatCommand(this));
+ this.commandExecutor.addCommand(Type.MSG_CHAT_HIST, new HistChatCommand(this));
+ this.commandExecutor.addCommand(Type.MSG_TEXT, new TextCommand(this));
+ this.commandExecutor.addCommand(Type.MSG_CHAT_LIST, new ListChatCommand(this));
+ } catch (ClassNotFoundException | SQLException e) {
+ e.printStackTrace();
+ }
+ }
+
+ public void start() {
+ threadPool.submit(listenForClients);
+ }
+
public void stop() {
- // TODO: закрыть все сетевые подключения, остановить потоки-обработчики, закрыть ресурсы, если есть.
+ try {
+ serverSocket.close();
+ } catch (IOException e) {
+ log.error(e.getMessage());
+ }
+ threadPool.shutdown();
+ }
+
+ // Сервер слушает клиентов. Каждому новому клиенту соответствует своя сессия, которая работает в новом потоке.
+ private Runnable listenForClients = new Runnable() {
+ @Override
+ public void run() {
+ try {
+ while (!serverSocket.isClosed()) {
+ Socket newClientSocket = null;
+ try {
+ log.info("Waiting for connect...");
+ newClientSocket = serverSocket.accept();
+ log.info("Connected: " + newClientSocket.getInetAddress());
+ } catch (SocketException e) {
+ log.error(e.getMessage());
+ }
+ Session session = new Session(newClientSocket, protocol, commandExecutor, usersDao);
+ sessions.add(session);
+ threadPool.submit(session);
+ }
+ } catch (Exception e) {
+ log.error(e.getMessage());
+ }
+ if (!serverSocket.isClosed()) {
+ threadPool.submit(listenForClients);
+ }
+ }
+ };
+
+ public int getPort() {
+ return port;
+ }
+
+ public void setPort(int port) {
+ this.port = port;
+ }
+
+ public Protocol getIProtocol() {
+ return protocol;
+ }
+
+ public void setIProtocol(Protocol protocol) {
+ this.protocol = protocol;
+ }
+
+ public void setExecutor(CommandExecutor executor) {
+ this.commandExecutor = executor;
}
-}
+ public CommandExecutor getExecutor() {
+ return commandExecutor;
+ }
+
+ public void setUsersDao(UsersDao usersDao) {
+ this.usersDao = usersDao;
+ }
+
+ public UsersDao getUsersDao() {
+ return usersDao;
+ }
+
+ public void setMessageStore(MessageStore messageStore) {
+ this.messageStore = messageStore;
+ }
+
+ public MessageStore getMessageStore() {
+ return messageStore;
+ }
+
+ public void setUserStore(UserStore userStore) {
+ this.userStore = userStore;
+ }
+
+ public UserStore getUserStore() {
+ return userStore;
+ }
+
+ public void setSessions(Set sessions) {
+ this.sessions = sessions;
+ }
+
+ public Set getSessions() {
+ return sessions;
+ }
+
+ public static void main(String[] args) throws Exception {
+ Protocol stringProtocol = new StringProtocol();
+ Server server = new Server(8000);
+ server.setIProtocol(stringProtocol);
+ server.start();
+ }
+
+
+}
\ No newline at end of file
diff --git a/src/main/resources/log4j.xml b/src/main/resources/log4j.xml
index 4331de1..809bb39 100644
--- a/src/main/resources/log4j.xml
+++ b/src/main/resources/log4j.xml
@@ -10,20 +10,24 @@
-
-
-
-
+
+
+
+
+ value="%d{yyyy-MM-dd HH:mm:ss} %-5p %c{1}:%L - %m%n"/>
+
+
-
-
-
+
+
diff --git a/src/test/java/arhangel/dim/container/ContainerTest.java b/src/test/java/arhangel/dim/container/ContainerTest.java
index 3152f67..8004877 100644
--- a/src/test/java/arhangel/dim/container/ContainerTest.java
+++ b/src/test/java/arhangel/dim/container/ContainerTest.java
@@ -8,7 +8,7 @@
import arhangel.dim.container.beans.Car;
import arhangel.dim.container.beans.Engine;
import arhangel.dim.container.beans.Gear;
-
+import arhangel.dim.container.exceptions.InvalidConfigurationException;
/**
*
*/