Class 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
    • 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()

        Возможно два сценария закрытия сессии

        1. Нормальное закрытие сессии, сессия сама извещает о завершении sesListener
        2. Аварийное закрытие сессии, сессия не извещает о завершении 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
        Список именнованых издателей серверных событий
      • proxyPublishers

        protected final Map<Class<?>,​Object> proxyPublishers
    • 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:
        Сессии клиентов
      • run

        public void run()
        Specified by:
        run in interface Runnable
        Overrides:
        run in class Thread
      • sessionSoTimeout

        protected int sessionSoTimeout()
        Возвращает значение SoTimeout Socket.setSoTimeout(int) для сессии
        Returns:
        по умолчанию 3000 мс
      • 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:
        close in interface AutoCloseable
        Throws:
        Exception - Ошибки...
      • hasListener

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

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

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

        public void removeListener​(TrListener listener)
        Удаление подписчика из списка обработки
        Parameters:
        listener - подписчик
      • 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 SerializablePublisher<T> publisher​(String name)
        Создание нового или получение уже существующего издалетя серверных событий
        Type Parameters:
        T - Тип события
        Parameters:
        name - Имя издателя
        Returns:
        Издатель
      • createPublisher

        protected Publisher<?> createPublisher​(String name)
      • 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:
        Издатель