Package xyz.cofe.trambda.tcp
Class TcpProtocol
- java.lang.Object
-
- xyz.cofe.trambda.tcp.TcpProtocol
-
public class TcpProtocol extends Object
Поддержка TCP протокола
-
-
Field Summary
Fields Modifier and Type Field Description protected Map<String,Set<Consumer<ServerEvent>>>serverEventListeners
-
Constructor Summary
Constructors Constructor Description TcpProtocol(OutputStream outputStream, InputStream inputStream, org.slf4j.Logger logger)КонструкторTcpProtocol(Socket socket)КонструкторTcpProtocol(Socket socket, org.slf4j.Logger logger)Конструктор
-
Method Summary
-
-
-
Field Detail
-
serverEventListeners
protected final Map<String,Set<Consumer<ServerEvent>>> serverEventListeners
-
-
Constructor Detail
-
TcpProtocol
public TcpProtocol(Socket socket)
Конструктор- Parameters:
socket- сокет
-
TcpProtocol
public TcpProtocol(Socket socket, org.slf4j.Logger logger)
Конструктор- Parameters:
socket- сокетlogger- логгер
-
TcpProtocol
public TcpProtocol(OutputStream outputStream, InputStream inputStream, org.slf4j.Logger logger)
Конструктор- Parameters:
outputStream- исходящий поток данныхinputStream- входящий поток данныхlogger- логгер
-
-
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:
- идентификатор сообщения
- Throws:
IOException- ошибка сети
-
send
@SafeVarargs public final int send(Message message, HeaderValue<? extends Object>... headerValues) throws IOException
Отправка сообщения- Parameters:
message- сообщениеheaderValues- дополнительные заголовки сообщения- Returns:
- идентификатор сообщения
- Throws:
IOException- ошибка сети
-
send
@SafeVarargs public final int send(Message message, Consumer<Integer> sid, HeaderValue<? extends Object>... headerValues) throws IOException
Отправка сообщения- Parameters:
message- сообщениеsid- идентификатор сообщенияheaderValues- дополнительные заголовки сообщения- Returns:
- идентификатор сообщения
- 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()
-
unbindedError
public BiConsumer<ErrMessage,TcpHeader> unbindedError()
-
unbindedError
public TcpProtocol unbindedError(BiConsumer<ErrMessage,TcpHeader> cons)
-
unbindedMessage
public BiConsumer<Message,TcpHeader> unbindedMessage()
-
unbindedMessage
public TcpProtocol unbindedMessage(BiConsumer<Message,TcpHeader> cons)
-
readNow
public boolean readNow() throws IOException- Throws:
IOException
-
listenServerEvent
public AutoCloseable listenServerEvent(String publisher, Consumer<ServerEvent> listener)
-
removeServerEventListener
public void removeServerEventListener(Consumer<ServerEvent> listener)
-
process
protected void process(ServerEvent sevent, TcpHeader header)
-
compile
public ResultConsumer<Compile,CompileResult> compile(LambdaDump methodDef)
-
execute
public ResultConsumer<Execute,ExecuteResult> execute(CompileResult cres)
-
subscribe
public ResultConsumer<Subscribe,SubscribeResult> subscribe(Subscribe subscribe)
-
unsubscribe
public ResultConsumer<UnSubscribe,UnSubscribeResult> unsubscribe(UnSubscribe subscribe)
-
-