Class TcpClient

  • All Implemented Interfaces:
    AutoCloseable

    public class TcpClient
    extends Object
    implements AutoCloseable
    Клиент для TcpSession/TcpServer. Оперирует уже готовым байт-кодом для отправки.
    Для работы с лямбдами используйте TcpQuery
    Содержит в себе фоновый поток ос (socketReaderThread) для чтения входящих сообщений

    Протокол

    TcpClient оперирует понятием сообщения Message.

    Сообщение - это объект который серилизован стандартным механизмом сериализации Java ObjectOutputStream

    Сообщения передаются в обе стороны

    Клиент хранит счетчик сообщений (TcpProtocol.sid) на основании которого ведет свою уникальность идентификаторов. Данный счет увеличивает свои номера линейно вверх, без повторений.

    При посылке сообщения на сервер TcpProtocol.send(xyz.cofe.trambda.tcp.Message, java.util.function.Consumer, xyz.cofe.trambda.tcp.HeaderValue...) К каждому сообщению прикрепляется уникальный номер и дополнительная карта (ключ/значение) связанное с сообщением.

    Так же сервер отвечает на это сообщение. Сервер вместе с ответным сообщением может ссылаться на исходное сообщение - то на что ответ приходиться.

    Пример посылки сообщения

    Клиент Сервер
    1. увеличивает счетчик sid на 1
    sid++
    2. текущее значение sid сохраняет в msgId:
    msgId = sid
    3. сериализует сообщения
    byte[] payload = serialize(msg)
    4. формирует заголовок TcpHeader
    4.1 В заголовке указывает тип сообщения
    header.method = msg.getClass().getSimpleName()
    4.2 В заголовке указывает идентификатор сообщения
    header,msgId = msgId
    4.3 В заголовке указывает дополнительную информацию
    например md5, размер payload, ...
    5. сфорримированные данные склеивает и отсылает на сервер
    byte[] tcpPackage = header + payload
    6. Принимает заголовок, из которого узнает размер полезных данных
    payloadSize = payloadSizeOf(tcpPackage)
    headerSize = headerSizeOf(tcpPackage)
    7. Считывает недостающие данные в буфер
     while(readed < payloadSize+headerSize) {
       readed += read( buffer )
     }
     
    8. Восстаналивает из буфера сообщение и выполняет полезное действие связанное с сообщением
    9. Если в результате действий есть результат который надо отправить клиенту, то формируется ответное сообщение по аналогии как формирует клиент, и отправляется обратно клиенту.
    10. Фоновый поток принимает ответное сообщение и восстанавливает его TcpProtocol.readNow(), так же как сервер
    11. Если с данным ответом были связаны подписчики, то вызывает их TcpProtocol.getResponseConsumers()

    За связку конкретного запроса и ответа на него в синхронном/асинхронном режиме на клиенте отвечает класс ResultConsumer

    За структуру TCP сообщения следует смотреть класс TcpHeader

    • Field Detail

      • socket

        protected final Socket socket
        Клиентский сокет
      • proto

        protected final TcpProtocol proto
        Работа с клиентским сокетом
      • socketReaderThread

        protected final Thread socketReaderThread
        Поток читающих входящие сообщения
      • listeners

        protected final xyz.cofe.ecolls.ListenersHelper<TrListener,​TrEvent> listeners
    • Constructor Detail

      • TcpClient

        public TcpClient​(Socket socket)
        Конструктор
        создает (createThread(Runnable, Socket)) и запускает поток для чтения входящих сообщений Message.
        Parameters:
        socket - клиентский сокет
    • Method Detail

      • hasListener

        public boolean hasListener​(TrListener listener)
        Проверка наличия подписчика в списке обработки
        Parameters:
        listener - подписчик
        Returns:
        true - есть в списке обработки
      • getListeners

        public Set<TrListener> getListeners()
        Получение списка подписчиков
        Returns:
        подписчики
      • addListener

        public AutoCloseable addListener​(TrListener listener)
        Добавление подписчика.
        Parameters:
        listener - Подписчик.
        Returns:
        Интерфес для отсоединения подписчика
      • addListener

        public AutoCloseable addListener​(TrListener listener,
                                         boolean weakLink)
        Добавление подписчика.
        Parameters:
        listener - Подписчик.
        weakLink - true - добавить как weak ссылку / false - как hard ссылку
        Returns:
        Интерфес для отсоединения подписчика
      • removeListener

        public void removeListener​(TrListener listener)
        Удаление подписчика из списка обработки
        Parameters:
        listener - подписчик
      • removeAllListeners

        public void removeAllListeners()
        Удаление всех подписчиков
      • withQueue

        protected void withQueue​(Runnable run)
        Запустить выполнение кода в блоке, и не рассылать уведомления до завершения блока кода
        Parameters:
        run - блок кода
      • withQueue

        protected <T> T withQueue​(Supplier<T> run)
        Запустить выполнение кода в блоке, и не рассылать уведомления до завершения блока кода
        Parameters:
        run - блок кода
        Returns:
        возвращаемое значение
      • fireEvent

        protected void fireEvent​(TrEvent event)
        Рассылка уведомления подписчикам
        Parameters:
        event - уведомление
      • addEvent

        protected void addEvent​(TrEvent ev)
        Добавляет событие в очередь
        Parameters:
        ev - событие
      • runEventQueue

        protected void runEventQueue()
        Отправляет события из очереди подписчикам
      • createThread

        protected Thread createThread​(Runnable code,
                                      Socket socket)
        Создание потока для чтения входящих сообщений
        Parameters:
        code - код читающий входящие сообщения
        socket - сокет клиента
        Returns:
        поток ос
      • close

        public void close()
        Завершение работы клиента, вызывает shutdown()
        Specified by:
        close in interface AutoCloseable
      • shutdown

        public void shutdown()
        Завершение работы клиента
        Нельзя вызывать из того же потока, который осуществляет чтение входящих сообщений - socketReaderThread