Class TcpClient
- java.lang.Object
-
- xyz.cofe.trambda.tcp.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 + payload6. Принимает заголовок, из которого узнает размер полезных данных
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 Summary
Fields Modifier and Type Field Description protected xyz.cofe.ecolls.ListenersHelper<TrListener,TrEvent>listenersprotected TcpProtocolprotoРабота с клиентским сокетомprotected SocketsocketКлиентский сокетprotected ThreadsocketReaderThreadПоток читающих входящие сообщенияprotected Map<String,List<xyz.cofe.fn.Tuple2<Consumer<ServerEvent>,AutoCloseable>>>subscribersСписок подписчиков на события сервера
-
Constructor Summary
Constructors Constructor Description TcpClient(Socket socket)Конструктор
создает (createThread(Runnable, Socket)) и запускает поток для чтения входящих сообщенийMessage.
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description protected voidaddEvent(TrEvent ev)Добавляет событие в очередьAutoCloseableaddListener(TrListener listener)Добавление подписчика.AutoCloseableaddListener(TrListener listener, boolean weakLink)Добавление подписчика.voidclose()Завершение работы клиента, вызываетshutdown()ResultConsumer<Compile,CompileResult>compile(xyz.cofe.trambda.LambdaDump methodDef)Компиляция лямбдыprotected ThreadcreateThread(Runnable code, Socket socket)Создание потока для чтения входящих сообщенийResultConsumer<Execute,ExecuteResult>execute(CompileResult cres)Выполнение ранее скомпилированной лямбды (compile(LambdaDump))protected voidfireEvent(TrEvent event)Рассылка уведомления подписчикамSet<TrListener>getListeners()Получение списка подписчиковbooleanhasListener(TrListener listener)Проверка наличия подписчика в списке обработкиvoidping(Consumer<Pong> consumer)voidremoveAllListeners()Удаление всех подписчиковvoidremoveListener(TrListener listener)Удаление подписчика из списка обработкиprotected voidrunEventQueue()Отправляет события из очереди подписчикамvoidshutdown()Завершение работы клиента
Нельзя вызывать из того же потока, который осуществляет чтение входящих сообщений -socketReaderThreadResultConsumer<Subscribe,SubscribeResult>subscribe(String publisher, Consumer<ServerEvent> listener)ResultConsumer<Subscribe,SubscribeResult>subscribe(Subscribe subscribe, Consumer<ServerEvent> listener)Подписка на события сервераResultConsumer<UnSubscribe,UnSubscribeResult>unsubscribe(String publisher)voidunsubscribe(Consumer<? super ServerEvent> listener)ResultConsumer<UnSubscribe,UnSubscribeResult>unsubscribe(UnSubscribe subscribe)protected voidwithQueue(Runnable run)Запустить выполнение кода в блоке, и не рассылать уведомления до завершения блока кодаprotected <T> TwithQueue(Supplier<T> run)Запустить выполнение кода в блоке, и не рассылать уведомления до завершения блока кода
-
-
-
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
-
subscribers
protected final Map<String,List<xyz.cofe.fn.Tuple2<Consumer<ServerEvent>,AutoCloseable>>> subscribers
Список подписчиков на события сервера
-
-
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:
closein interfaceAutoCloseable
-
shutdown
public void shutdown()
Завершение работы клиента
Нельзя вызывать из того же потока, который осуществляет чтение входящих сообщений -socketReaderThread
-
compile
public ResultConsumer<Compile,CompileResult> compile(xyz.cofe.trambda.LambdaDump methodDef)
Компиляция лямбдысм.
Compile,CompileResult,compile(LambdaDump),см call(Fn1, SerializedLambda, LambdaDump)()}- Parameters:
methodDef- лямбда- Returns:
- выполнение запроса
-
execute
public ResultConsumer<Execute,ExecuteResult> execute(CompileResult cres)
Выполнение ранее скомпилированной лямбды (compile(LambdaDump))см.
Compile,CompileResult,compile(LambdaDump),см call(Fn1, SerializedLambda, LambdaDump)()}- Parameters:
cres- результат компиляции- Returns:
- выполнение запроса
-
subscribe
public ResultConsumer<Subscribe,SubscribeResult> subscribe(Subscribe subscribe, Consumer<ServerEvent> listener)
Подписка на события сервера- Parameters:
subscribe- подписка (имя издателя)listener- подписчик- Returns:
- выполнение запроса
-
subscribe
public ResultConsumer<Subscribe,SubscribeResult> subscribe(String publisher, Consumer<ServerEvent> listener)
-
unsubscribe
public ResultConsumer<UnSubscribe,UnSubscribeResult> unsubscribe(UnSubscribe subscribe)
-
unsubscribe
public ResultConsumer<UnSubscribe,UnSubscribeResult> unsubscribe(String publisher)
-
unsubscribe
public void unsubscribe(Consumer<? super ServerEvent> listener)
-
-