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; /** * */