Package xyz.cofe.trambda.tcp
Class TcpProtocol
- java.lang.Object
-
- xyz.cofe.trambda.tcp.TcpProtocol
-
public class TcpProtocol extends Object
Поддержка TCP протокола.
Используется в клиентеTcpClient, взаимодействует с сессиейTcpSession.
Умеет- отправку сообщений
Message, через методыsend(Message, HeaderValue[]),send(Message, Consumer, HeaderValue[]),sendRaw(String, byte[], Consumer, HeaderValue[]) - читает очередное сообщение
readNow() - обрабатывает входящие сообщения
process(Message, TcpHeader)
- отправку сообщений
-
-
Nested Class Summary
Nested Classes Modifier and Type Class Description static classTcpProtocol.SentRawData
-
Field Summary
Fields Modifier and Type Field Description protected Supplier<InputStream>getInputКанал/поток входящих сообщенийprotected Supplier<xyz.cofe.trambda.log.api.Logger>getLoggerЛогированиеprotected Supplier<OutputStream>getOutputКанал/поток исходящих сообщенийprotected ThreadLocal<Message>sendMessageConsumer<TcpProtocol.SentRawData>sentRawDataConsumerprotected Map<String,Set<Consumer<ServerEvent>>>serverEventListenersПодписчики на серверные событияprotected AtomicIntegersidСчетчик/генератор идентификаторов сообщенийprotected SocketsocketСокет с которым производится работа
-
Constructor Summary
Constructors Constructor Description TcpProtocol(OutputStream outputStream, InputStream inputStream, xyz.cofe.trambda.log.api.Logger logger)КонструкторTcpProtocol(OutputStream outputStream, InputStream inputStream, xyz.cofe.trambda.log.api.Logger logger, Consumer<TcpProtocol.SentRawData> sentRawDataConsumer)КонструкторTcpProtocol(Socket socket)КонструкторTcpProtocol(Socket socket, Consumer<TcpProtocol.SentRawData> sentRawDataConsumer)КонструкторTcpProtocol(Socket socket, xyz.cofe.trambda.log.api.Logger logger)Конструктор
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description ResultConsumer<Compile,CompileResult>compile(xyz.cofe.trambda.LambdaDump methodDef)Компиляция лямбдыResultConsumer<Execute,ExecuteResult>execute(CompileResult cres)Выполнение ранее скомпилированной лямбдыMap<Integer,Consumer<ErrMessage>>getErrorConsumers()Возвращает обработчики ошибок на конкретные запросыTcpHeader.getSid()Map<Integer,Consumer<? extends Message>>getResponseConsumers()Возвращает обработчики на конкретные запросыTcpHeader.getSid()AutoCloseablelistenServerEvent(String publisher, Consumer<ServerEvent> listener)Добавляет подписчика на событие сервера.intping(Consumer<Pong> consumer)ОтправляетPingзапросprotected voidprocess(Message msg, TcpHeader header)Обработка входящего сообщения.protected voidprocess(Ping ping, TcpHeader header)При получении Ping сообщения, отсылает ответPongprotected voidprocess(Pong pong, TcpHeader header)При получении ответаPongуведомляет подписанотовping(Consumer), тех кто отправил запрос Pingprotected voidprocess(ServerEvent sevent, TcpHeader header)Обработка входящего серверного события Уведомляет подписчиковlistenServerEvent(String, Consumer)о серверном событииbooleanreadNow()Optional<xyz.cofe.fn.Tuple2<RawPack,Optional<BuffInputStream>>>receiveRaw(Optional<BuffInputStream> prevStream)Получение сообщения из сетиvoidremoveServerEventListener(Consumer<ServerEvent> listener)Отписка от событий сервераintsend(Message message, Consumer<Integer> sid, HeaderValue<? extends Object>... headerValues)Отправка сообщенияintsend(Message message, HeaderValue<? extends Object>... headerValues)Отправка сообщенияintsendRaw(String method, byte[] payload, Consumer<Integer> sendId, HeaderValue<? extends Object>... headerValues)Отправка сообщенияResultConsumer<Subscribe,SubscribeResult>subscribe(Subscribe subscribe)Подписка на события сервера
Используйте возможности клиентаTcpClient.subscribe(String, Consumer)BiConsumer<ErrMessage,TcpHeader>unbindedError()Возвращает обработчик общих ошибок, смgetErrorConsumers()TcpProtocolunbindedError(BiConsumer<ErrMessage,TcpHeader> cons)Указывает обработчик общих ошибок, смgetErrorConsumers()BiConsumer<Message,TcpHeader>unbindedMessage()Возвращает обработчик сообщений, смgetResponseConsumers()TcpProtocolunbindedMessage(BiConsumer<Message,TcpHeader> cons)Указывает обработчик сообщений, смgetResponseConsumers()ResultConsumer<UnSubscribe,UnSubscribeResult>unsubscribe(UnSubscribe subscribe)Отписка от событий сервера
Используйте возможности клиентаTcpClient.unsubscribe(String),TcpClient.unsubscribe(Consumer)
-
-
-
Field Detail
-
sentRawDataConsumer
public final Consumer<TcpProtocol.SentRawData> sentRawDataConsumer
-
socket
protected final Socket socket
Сокет с которым производится работа
-
getOutput
protected final Supplier<OutputStream> getOutput
Канал/поток исходящих сообщений
-
getInput
protected final Supplier<InputStream> getInput
Канал/поток входящих сообщений
-
getLogger
protected final Supplier<xyz.cofe.trambda.log.api.Logger> getLogger
Логирование
-
sid
protected final AtomicInteger sid
Счетчик/генератор идентификаторов сообщений
-
sendMessage
protected final ThreadLocal<Message> sendMessage
-
serverEventListeners
protected final Map<String,Set<Consumer<ServerEvent>>> serverEventListeners
Подписчики на серверные события
-
-
Constructor Detail
-
TcpProtocol
public TcpProtocol(Socket socket)
Конструктор- Parameters:
socket- сокет
-
TcpProtocol
public TcpProtocol(Socket socket, Consumer<TcpProtocol.SentRawData> sentRawDataConsumer)
Конструктор- Parameters:
socket- сокетsentRawDataConsumer- подписчик на отправленные сообщения
-
TcpProtocol
public TcpProtocol(Socket socket, xyz.cofe.trambda.log.api.Logger logger)
Конструктор- Parameters:
socket- сокетlogger- логгер
-
TcpProtocol
public TcpProtocol(OutputStream outputStream, InputStream inputStream, xyz.cofe.trambda.log.api.Logger logger)
Конструктор- Parameters:
outputStream- исходящий поток данныхinputStream- входящий поток данныхlogger- логгер
-
TcpProtocol
public TcpProtocol(OutputStream outputStream, InputStream inputStream, xyz.cofe.trambda.log.api.Logger logger, Consumer<TcpProtocol.SentRawData> sentRawDataConsumer)
Конструктор- Parameters:
outputStream- исходящий поток данныхinputStream- входящий поток данныхlogger- логгерsentRawDataConsumer- подписчик на отправленные сообщения
-
-
Method Detail
-
sendRaw
@SafeVarargs public final int sendRaw(String method, byte[] payload, Consumer<Integer> sendId, HeaderValue<? extends Object>... headerValues) throws IOException
Отправка сообщения- Parameters:
method- методpayload- полезная нагрузкаsendId- идентификатор сообщенияheaderValues- дополнительные заголовки сообщения- Returns:
- идентификатор сообщения, см
TcpHeader.getSid() - Throws:
IOException- ошибка сети
-
send
@SafeVarargs public final int send(Message message, HeaderValue<? extends Object>... headerValues) throws IOException
Отправка сообщения- Parameters:
message- сообщениеheaderValues- дополнительные заголовки сообщения- Returns:
- идентификатор сообщения, см
TcpHeader.getSid() - Throws:
IOException- ошибка сети
-
send
@SafeVarargs public final int send(Message message, Consumer<Integer> sid, HeaderValue<? extends Object>... headerValues) throws IOException
Отправка сообщения- Parameters:
message- сообщениеsid- идентификатор сообщенияheaderValues- дополнительные заголовки сообщения- Returns:
- идентификатор сообщения, см
TcpHeader.getSid() - Throws:
IOException- ошибка сети
-
receiveRaw
public Optional<xyz.cofe.fn.Tuple2<RawPack,Optional<BuffInputStream>>> receiveRaw(Optional<BuffInputStream> prevStream) throws IOException
Получение сообщения из сети- Parameters:
prevStream- предыдущая часть сообщения- Returns:
- сообщение или входящий поток закрыт
- Throws:
IOException- ошибка сети
-
getErrorConsumers
public Map<Integer,Consumer<ErrMessage>> getErrorConsumers()
Возвращает обработчики ошибок на конкретные запросыTcpHeader.getSid()- Returns:
- Карта ключ - идентификатор сообщения
TcpHeader.getSid()/ обработчик
-
getResponseConsumers
public Map<Integer,Consumer<? extends Message>> getResponseConsumers()
Возвращает обработчики на конкретные запросыTcpHeader.getSid()- Returns:
- Карта ключ - идентификатор сообщения
TcpHeader.getSid()/ обработчик
-
unbindedError
public BiConsumer<ErrMessage,TcpHeader> unbindedError()
Возвращает обработчик общих ошибок, смgetErrorConsumers()- Returns:
- обработчик общих ошибок
-
unbindedError
public TcpProtocol unbindedError(BiConsumer<ErrMessage,TcpHeader> cons)
Указывает обработчик общих ошибок, смgetErrorConsumers()- Parameters:
cons- обработчик общих ошибок- Returns:
- SELF ссылка
-
unbindedMessage
public BiConsumer<Message,TcpHeader> unbindedMessage()
Возвращает обработчик сообщений, смgetResponseConsumers()- Returns:
- обработчик сообщений
-
unbindedMessage
public TcpProtocol unbindedMessage(BiConsumer<Message,TcpHeader> cons)
Указывает обработчик сообщений, смgetResponseConsumers()- Parameters:
cons- обработчик сообщений- Returns:
- SELF ссылка
-
readNow
public boolean readNow() throws IOExceptionЧитает сообщениеMessageиз входного потока (getInput).- Делает попытку чтения данных из входящего потока
receiveRaw(Optional) - Проверяет контрольную сумму сообщения
RawPackReadonly.isPayloadChecksumMatched() - Вычленяет из пакета байтов само сообщение
RawPackReadonly.payloadMessage() - Отправляет сообщение на последующую обработку
process(Message, TcpHeader)
- Returns:
- true - сообщение прочитано
- Throws:
IOException- Ошибка чтения данных
- Делает попытку чтения данных из входящего потока
-
process
protected void process(Message msg, TcpHeader header)
Обработка входящего сообщения.- Отрабатывает сообщения
Ping- вызываетprocess(Ping, TcpHeader) - Отрабатывает сообщения
Pong- вызываетprocess(Pong, TcpHeader) - Отрабатывает сообщения
ErrMessage- Если указан на какое сообщение
TcpHeader.getReferrer()ответ- то вызывает обработчик конкретный
getErrorConsumers() - или общий обработчик
unbindedError()
- то вызывает обработчик конкретный
- Если указан на какое сообщение
-
Обрабатывает сообщение
ServerEventи вызываетprocess(ServerEvent, TcpHeader) -
Для всех остальных сообщений
- Если есть зарегистрированный обработчик
getResponseConsumers(), то вызывается он - Либо общий обработчик
unbindedMessage()
- Если есть зарегистрированный обработчик
- Parameters:
msg- входящие сообщениеheader- заголовки сообщения
- Отрабатывает сообщения
-
process
protected void process(Pong pong, TcpHeader header)
При получении ответаPongуведомляет подписанотовping(Consumer), тех кто отправил запрос Ping- Parameters:
pong- Ответheader- заголовки
-
process
protected void process(Ping ping, TcpHeader header)
При получении Ping сообщения, отсылает ответPong- Parameters:
ping- Входящий Pingheader- Заголовки
-
listenServerEvent
public AutoCloseable listenServerEvent(String publisher, Consumer<ServerEvent> listener)
Добавляет подписчика на событие сервера.
Используйте возможности клиентаTcpClient.subscribe(String, Consumer)- Parameters:
publisher- издатель событияlistener- подписчик- Returns:
- отписка от событий
-
removeServerEventListener
public void removeServerEventListener(Consumer<ServerEvent> listener)
Отписка от событий сервера- Parameters:
listener- подписчик
-
process
protected void process(ServerEvent sevent, TcpHeader header)
Обработка входящего серверного события-
Уведомляет подписчиков
listenServerEvent(String, Consumer)о серверном событии
- Parameters:
sevent- серверное событиеheader- заголовки
-
Уведомляет подписчиков
-
ping
public int ping(Consumer<Pong> consumer)
ОтправляетPingзапрос- Parameters:
consumer- приемник ответа- Returns:
- Идентификатор
TcpHeader.getSid()
-
compile
public ResultConsumer<Compile,CompileResult> compile(xyz.cofe.trambda.LambdaDump methodDef)
Компиляция лямбды- Parameters:
methodDef- лямбда- Returns:
- Отправка запроса
-
execute
public ResultConsumer<Execute,ExecuteResult> execute(CompileResult cres)
Выполнение ранее скомпилированной лямбды- Parameters:
cres- ранее скомпилированная лямбдаcompile(LambdaDump)- Returns:
- Результат выполнения
-
subscribe
public ResultConsumer<Subscribe,SubscribeResult> subscribe(Subscribe subscribe)
Подписка на события сервера
Используйте возможности клиентаTcpClient.subscribe(String, Consumer)- Parameters:
subscribe- имя издателя сервера- Returns:
- Результат подписки
-
unsubscribe
public ResultConsumer<UnSubscribe,UnSubscribeResult> unsubscribe(UnSubscribe subscribe)
Отписка от событий сервера
Используйте возможности клиентаTcpClient.unsubscribe(String),TcpClient.unsubscribe(Consumer)- Parameters:
subscribe- подписчик- Returns:
- Результат отписки
-
-