Package xyz.cofe.trambda.tcp
Class TcpServer<ENV>
- java.lang.Object
-
- java.lang.Thread
-
- xyz.cofe.trambda.tcp.TcpServer<ENV>
-
- Type Parameters:
ENV- Класс сервиса
- All Implemented Interfaces:
AutoCloseable,Runnable,TrEventPublisher
public class TcpServer<ENV> extends Thread implements AutoCloseable, TrEventPublisher
TCP Сервер для предоставления сервиса.Создание простого клиент/сервера
В первую очередь необходимо- созданный
ServerSocket - сам сервис или функция возвращающая ссылку на сервис
// Содержит описание процесса ОС public class OsProc implements Serializable { private final Integer ppid; private final Integer pid; private final String name; private final String cmdLine; public OsProc(int ppid,int pid,String name,String cmdLine){ this.ppid = ppid; this.pid = pid; this.name = name; this.cmdLine = cmdLine; } public Optional<Integer> getPpid(){ return ppid!=null ? Optional.of(ppid) : Optional.empty(); } public int getPid(){ return pid; } public String getName(){ return name; } public Optional<String> getCmdline(){ return cmdLine!=null ? Optional.of(cmdLine) : Optional.empty(); } public static OsProc linuxProc(File file){ if( file.isDir() ){ try{ var statusFile = file.resolve("status"); var status = statusFile.isFile() ? statusFile.readText(Charset.defaultCharset()) : ""; var cmdLineFile = file.resolve("cmdline"); var cmdLine = cmdLineFile.isFile() ? cmdLineFile.readText(Charset.defaultCharset()) : ""; var keyVals = Text.splitNewLinesIterable(status) .map(line -> line.split("\\s*:\\s*",2)) .filter(kv -> kv.length==2) .map(kv -> Tuple2.of(kv[0].toLowerCase(), kv[1])); String name = keyVals.filter( kv->kv.a().equals("name")) .map(Tuple2::b).first().orElse("?"); String pid = keyVals.filter( kv->kv.a().equals("pid")) .map(Tuple2::b).first().orElse("-1"); String ppid = keyVals.filter( kv->kv.a().equals("ppid")) .map(Tuple2::b).first().orElse("-1"); return new OsProc( Integer.parseInt(ppid), Integer.parseInt(pid), name, cmdLine.replace((char)0,' ') ); } catch( Throwable err ){ System.err.println("err for "+file+": "+err.getMessage()); } } return new OsProc(-1,"?"); } } // Это будет интерфейсом нашого сервиса, // через который будет происходить общение public interface IEnv { // Возвращает список процессов public List<OsProc> processes(); } // Это реализация нашего сервиса для Linux public class LinuxEnv implements IEnv { @Override public List<OsProc> processes(){ ArrayList<OsProc> procs = new ArrayList<>(); File procDir = new File("/proc"); procDir.dirList().stream() .filter( d -> d.getName().matches("\\d+") && d.isDir() ) .map(OsProc::linuxProc) .forEach(procs::add); return procs; } }Далее поднимим Tcp серверimport java.io.IOException; import java.net.ServerSocket; import xyz.cofe.trambda.tcp.TcpServer; ... TcpServer<IEnv> createServer(IEnv myService, int port){ ... ServerSocket ssocket = null; TcpServer<IEnv> server = null; try{ // Создаем сокет ssocket = new ServerSocket(port); // Указываем SoTimeout что бы // сервер не повис в ожидании пакета ssocket.setSoTimeout(1000*5); // Создаем сам сервер server = new TcpServer<IEnv>( // Передаем сокет, который уже привязан к порту ssocket, // ссылку на сервис session -> myService ); // Указывем характеристики Thread server.setDaemon(true); server.setName("server"); // Добавляем подписчиков на сообытия сервера server.addListener(System.out::println); // Запускаем сервер server.start(); return server; } catch( IOException error ) { log.error( "can't start server", error ); return; } }Сервис готов и поднят, осталось написать клиента:
Query<IEnv> query = TcpQuery // Указываем интерфейс сервиса .create(IEnv.class) // Указывам адрес и порт, на котором располагается сервис .host("localhost") .port(port) // Создаем клиента .build(); // Указываем какие процессы нас интересуют на сервере var qRegex = "chrome|java"; // Выполняем запрос к серверу query.apply( env -> // Данный код выполняется на сервере env.processes().stream() .filter( p -> p.getName().matches("(?is).*("+qRegex+").*") ) .collect(Collectors.toList()) ) // Получаем результат выполнения в клиенте // И отображаем его в логе .stream().map(OsProc::toString).forEach(log::info);Дополнительные темы
- Настройка безопасности, см
SecurError - Подписка на события сервера, см
Subscribe
-
-
Nested Class Summary
Nested Classes Modifier and Type Class Description static classTcpServer.SessionClosed<ENV>Событие о завершении сессииstatic classTcpServer.SessionCreated<ENV>Событие о создании сессии-
Nested classes/interfaces inherited from class java.lang.Thread
Thread.State, Thread.UncaughtExceptionHandler
-
-
Field Summary
Fields Modifier and Type Field Description protected Function<TcpSession<ENV>,ENV>envBuilderФункция получения сервиса для новой сессииprotected Map<Integer,Long>fireClosedИнформация когда было уведомление о закрытии сессии: ses.id / System.currentTimeMillis()protected xyz.cofe.ecolls.ListenersHelper<TrListener,TrEvent>listenersprotected Map<Class<?>,Object>proxyPublishersprotected Map<String,Publisher<?>>publishersСписок именнованых издателей серверных событийprotected xyz.cofe.trambda.sec.SecurityFilter<String,xyz.cofe.fn.Tuple2<xyz.cofe.trambda.LambdaDump,xyz.cofe.trambda.LambdaNode>>securityFilterФункция фильтрации байт-кодаprotected Set<TcpSession<ENV>>sessionsСессии клиентовprotected ServerSocketsocketСокет через который осуществляется общение-
Fields inherited from class java.lang.Thread
MAX_PRIORITY, MIN_PRIORITY, NORM_PRIORITY
-
-
Constructor Summary
Constructors Constructor Description TcpServer(ServerSocket socket, Function<TcpSession<ENV>,ENV> envBuilder)Создание сервераTcpServer(ServerSocket socket, Function<TcpSession<ENV>,ENV> envBuilder, xyz.cofe.trambda.sec.SecurityFilter<String,xyz.cofe.fn.Tuple2<xyz.cofe.trambda.LambdaDump,xyz.cofe.trambda.LambdaNode>> securityFilter)Создание сервера
-
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)Добавление подписчика.protected voidaddSesListener(TcpSession<ENV> ses)Добавление подписчика sesListener на события сессииvoidclose()Завершение работы сервераprotected voidcloseSessions()Завершение всех сессийprotected voidcloseSocket()Закрытие сокетаprotected TcpSession<ENV>create(Socket sock)Создание сессииprotected <T> TcreateProxyPublisher(Class<T> cls)Создание прокси для класса событийprotected Publisher<?>createPublisher(String name)protected voidfireEvent(TrEvent event)Рассылка уведомления подписчикамSet<TrListener>getListeners()Получение списка подписчиковSet<TcpSession<ENV>>getSessions()Возвращает сессииbooleanhasListener(TrListener listener)Проверка наличия подписчика в списке обработки<T extends Serializable>
Publisher<T>publisher(String name)Создание нового или получение уже существующего издалетя серверных событий<T> Tpublishers(Class<T> cls)Создание нового или получение уже существующего издалетя класса серверных событий.voidremoveAllListeners()Удаление всех подписчиковvoidremoveListener(TrListener listener)Удаление подписчика из списка обработкиvoidrun()protected voidrunEventQueue()Отправляет события из очереди подписчикамprotected longsessionCloseTimeout()Таймаут согласно которому сессия должна быть завершенаprotected intsessionSoTimeout()Возвращает значение SoTimeoutSocket.setSoTimeout(int)для сессииvoidshutdown()Завершение всех сессий и остановка сервераprotected voidwithQueue(Runnable run)Запустить выполнение кода в блоке, и не рассылать уведомления до завершения блока кодаprotected <T> TwithQueue(Supplier<T> run)Запустить выполнение кода в блоке, и не рассылать уведомления до завершения блока кода-
Methods inherited from class java.lang.Thread
activeCount, checkAccess, clone, countStackFrames, currentThread, dumpStack, enumerate, getAllStackTraces, getContextClassLoader, getDefaultUncaughtExceptionHandler, getId, getName, getPriority, getStackTrace, getState, getThreadGroup, getUncaughtExceptionHandler, holdsLock, interrupt, interrupted, isAlive, isDaemon, isInterrupted, join, join, join, onSpinWait, resume, setContextClassLoader, setDaemon, setDefaultUncaughtExceptionHandler, setName, setPriority, setUncaughtExceptionHandler, sleep, sleep, start, stop, suspend, toString, yield
-
-
-
-
Field Detail
-
socket
protected final ServerSocket socket
Сокет через который осуществляется общение
-
sessions
protected final Set<TcpSession<ENV>> sessions
Сессии клиентов
-
fireClosed
protected final Map<Integer,Long> fireClosed
Информация когда было уведомление о закрытии сессии: ses.id / System.currentTimeMillis()Возможно два сценария закрытия сессии
-
Нормальное закрытие сессии, сессия сама извещает о завершении
sesListener -
Аварийное закрытие сессии, сессия не извещает о завершении
checkTerminatedSessions()
-
Нормальное закрытие сессии, сессия сама извещает о завершении
-
envBuilder
protected final Function<TcpSession<ENV>,ENV> envBuilder
Функция получения сервиса для новой сессии
-
securityFilter
protected final xyz.cofe.trambda.sec.SecurityFilter<String,xyz.cofe.fn.Tuple2<xyz.cofe.trambda.LambdaDump,xyz.cofe.trambda.LambdaNode>> securityFilter
Функция фильтрации байт-кода
-
listeners
protected final xyz.cofe.ecolls.ListenersHelper<TrListener,TrEvent> listeners
-
publishers
protected final Map<String,Publisher<?>> publishers
Список именнованых издателей серверных событий
-
-
Constructor Detail
-
TcpServer
public TcpServer(ServerSocket socket, Function<TcpSession<ENV>,ENV> envBuilder, xyz.cofe.trambda.sec.SecurityFilter<String,xyz.cofe.fn.Tuple2<xyz.cofe.trambda.LambdaDump,xyz.cofe.trambda.LambdaNode>> securityFilter)
Создание сервера- Parameters:
socket- сокетenvBuilder- Функция получения сервиса для новой сессииsecurityFilter- Функция фильтрации байт-кода
-
TcpServer
public TcpServer(ServerSocket socket, Function<TcpSession<ENV>,ENV> envBuilder)
Создание сервера- Parameters:
socket- сокетenvBuilder- Функция получения сервиса для новой сессии
-
-
Method Detail
-
getSessions
public Set<TcpSession<ENV>> getSessions()
Возвращает сессии- Returns:
- Сессии клиентов
-
sessionSoTimeout
protected int sessionSoTimeout()
Возвращает значение SoTimeoutSocket.setSoTimeout(int)для сессии- Returns:
- по умолчанию 3000 мс
-
create
protected TcpSession<ENV> create(Socket sock)
Создание сессии- Parameters:
sock- сокет- Returns:
- сессия
- See Also:
sessionSoTimeout(),addSesListener(TcpSession),TcpServer.SessionCreated
-
shutdown
public void shutdown()
Завершение всех сессий и остановка сервера- See Also:
closeSocket(),closeSessions()
-
closeSocket
protected void closeSocket()
Закрытие сокета
-
sessionCloseTimeout
protected long sessionCloseTimeout()
Таймаут согласно которому сессия должна быть завершена- Returns:
- 5000 мс
-
closeSessions
protected void closeSessions()
Завершение всех сессий- See Also:
sessionCloseTimeout()
-
addSesListener
protected void addSesListener(TcpSession<ENV> ses)
Добавление подписчика sesListener на события сессии- Parameters:
ses- сессия
-
close
public void close() throws ExceptionЗавершение работы сервера- Specified by:
closein interfaceAutoCloseable- Throws:
Exception- Ошибки...
-
hasListener
public boolean hasListener(TrListener listener)
Проверка наличия подписчика в списке обработки- Specified by:
hasListenerin interfaceTrEventPublisher- Parameters:
listener- подписчик- Returns:
- true - есть в списке обработки
-
getListeners
public Set<TrListener> getListeners()
Получение списка подписчиков- Specified by:
getListenersin interfaceTrEventPublisher- Returns:
- подписчики
-
addListener
public AutoCloseable addListener(TrListener listener)
Добавление подписчика.- Specified by:
addListenerin interfaceTrEventPublisher- Parameters:
listener- Подписчик.- Returns:
- Интерфейс для отсоединения подписчика
-
addListener
public AutoCloseable addListener(TrListener listener, boolean weakLink)
Добавление подписчика.- Specified by:
addListenerin interfaceTrEventPublisher- Parameters:
listener- Подписчик.weakLink- true - добавить как weak ссылку / false - как hard ссылку- Returns:
- Интерфейс для отсоединения подписчика
-
removeListener
public void removeListener(TrListener listener)
Удаление подписчика из списка обработки- Parameters:
listener- подписчик
-
removeAllListeners
public void removeAllListeners()
Удаление всех подписчиков- Specified by:
removeAllListenersin interfaceTrEventPublisher
-
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()
Отправляет события из очереди подписчикам
-
publisher
public <T extends Serializable> Publisher<T> publisher(String name)
Создание нового или получение уже существующего издалетя серверных событий- Type Parameters:
T- Тип события- Parameters:
name- Имя издателя- Returns:
- Издатель
-
publishers
public <T> T publishers(Class<T> cls)
Создание нового или получение уже существующего издалетя класса серверных событий.Издатель заданный классом - это интерфейс, например такой
public interface Events { Publisher<ServerDemoEvent> defaultPublisher(); Publisher<ServerDemoEvent2> timedEvents(); }Для каждого метода (в данном примере для defaultPublisher и timedEvents) будет зарегистриован (publisher(java.lang.String)) свой издатель, имя издателя будет совпадать с именем метода.Грубо будет так:
proxy = new
PubProxy(){ @Override protected Publisher<?> publisher(Method method){ String publisherName = method.getName(); return TcpServer.this.publisher(publisherName); } };- Type Parameters:
T- Класс серверных событий - интерфейс c методами без параметра- Parameters:
cls- Класс серверных событий- Returns:
- Издатель
-
createProxyPublisher
protected <T> T createProxyPublisher(Class<T> cls)
Создание прокси для класса событий- Type Parameters:
T- Класс серверных событий - интерфейс c методами без параметра- Parameters:
cls- Класс серверных событий- Returns:
- Издатель
-
-