Class Session

    • Constructor Detail

      • Session

        public Session​(Socket socket,
                       SSLSocket secureSocket)
                throws IOException
        Creates a new Session.
        Parameters:
        socket - The regular Socket on which the SocketSession will be based.
        secureSocket - The secure Socket on which the SocketSession will be based.
        Throws:
        IOException - When an IException happens on the socket.
    • Method Detail

      • close

        public void close()
        This method is called when the session with the remote must be closed. This object won't be used anymore after this method is called.

        A message which was published on this session but which its publisher thread had not sent yet is sent here rather than dropped, within the budget of DRAIN_BUDGET_MS, and its callback runs here once it is written. See sendWhatThePublisherLeftQueued(). What the close gives up on is not written, and the callbacks of those messages never run - see publish(ReplicationMsg, Runnable).

        Specified by:
        close in interface AutoCloseable
        Specified by:
        close in interface Closeable
      • closeInitiated

        public boolean closeInitiated()
        This methods allows to determine if the session close was initiated on this Session.
        Returns:
        A boolean allowing to determine if the session close was initiated on this Session.
      • getLastPublishTime

        public long getLastPublishTime()
        Gets the time the last replication message was published on this session.
        Returns:
        The timestamp in milliseconds of the last message published.
      • getLastReceiveTime

        public long getLastReceiveTime()
        Gets the time the last replication message was received on this session.
        Returns:
        The timestamp in milliseconds of the last message received.
      • getLocalUrl

        public HostPort getLocalUrl()
        Retrieve the local URL in the form host:port.
        Returns:
        The local URL.
      • getReadableRemoteAddress

        public String getReadableRemoteAddress()
        Retrieve the human readable address of the remote server.
        Returns:
        The human readable address of the remote server.
      • getRemoteAddress

        public HostPort getRemoteAddress()
        Retrieve the IP address and port of the remote server.
        Returns:
        The IP address and port of the remote server.
      • isEncrypted

        public boolean isEncrypted()
        Determine whether the session is using a security layer.
        Returns:
        true if the connection is encrypted, false otherwise.
      • publish

        public void publish​(ReplicationMsg msg)
                     throws IOException
        Sends a replication message to the remote peer.
        Parameters:
        msg - The message to be sent.
        Throws:
        IOException - If an IO error occurred.
      • publish

        public boolean publish​(ReplicationMsg msg,
                               Runnable whenWritten)
                        throws IOException
        Sends a replication message to the remote peer, and runs the provided callback once the message has been written to the socket.

        While the thread of this session runs, a message published is queued for it and written later, so the return of this method says only that the message is queued. The callback is the only word that the message has left this server: it runs once, on the thread which wrote the message, after the write returned - the thread of the session, or the one closing it for a message the close sends out of the queue - and never for a message which was not written, which is what becomes of a message the write of which fails, and of what a close gives up on (see close()). It must be short and must not block: the session writes nothing else until it returns.

        Parameters:
        msg - The message to be sent.
        whenWritten - What to run once the message has been written, or null.
        Returns:
        whether the message was written or queued to be written; false when it was neither, because it has no encoding for the protocol version of the peer or because the session is being closed - the callback then never runs. A message queued after a close drained the queue is taken back, and counts as neither.
        Throws:
        IOException - If an IO error occurred.
      • setProtocolVersion

        public void setProtocolVersion​(short version)
        This method is called at the establishment of the session and can be used to record the version of the protocol that is currently used.
        Parameters:
        version - The version of the protocol that is currently used.
      • getProtocolVersion

        public short getProtocolVersion()
        Returns the version of the protocol that is currently used.
        Returns:
        The version of the protocol that is currently used.
      • setSoTimeout

        public void setSoTimeout​(int timeout)
                          throws SocketException
        Set a timeout value. With this option set to a non-zero value, calls to the receive() method block for only this amount of time after which a java.net.SocketTimeoutException is raised. The Broker is valid and usable even after such an Exception is raised.
        Parameters:
        timeout - the specified timeout, in milliseconds.
        Throws:
        SocketException - if there is an error in the underlying protocol, such as a TCP error.
      • stopEncryption

        public void stopEncryption()
        Stop using the security layer, if there is any.
      • run

        public void run()
        Run method for the Session. Loops waiting for buffers from the queue and sends them when available.
        Specified by:
        run in interface Runnable
        Overrides:
        run in class Thread
      • waitForStartup

        public void waitForStartup()
                            throws InterruptedException
        This method can be called to wait until the session thread is properly started.
        Throws:
        InterruptedException - when interrupted