Class AsyncNodeContext

java.lang.Object
org.jppf.nio.AbstractNioContext
org.jppf.server.nio.nodeserver.BaseNodeContext
org.jppf.server.nio.nodeserver.async.AsyncNodeContext
All Implemented Interfaces:
AutoCloseable, org.jppf.execute.ExecutorChannel<ServerTaskBundleNode>, org.jppf.nio.CloseableContext, org.jppf.nio.NioChannelHandler, org.jppf.nio.NioContext

public class AsyncNodeContext extends BaseNodeContext
Context or state information associated with a channel that exchanges heartbeat messages between the server and a node or client.
Author:
Laurent Cohen
  • Constructor Details

    • AsyncNodeContext

      public AsyncNodeContext(AsyncNodeNioServer server, SocketChannel socketChannel, boolean local)
      Parameters:
      server - the server that handles this context.
      socketChannel - the associated socket channel.
      local - whether this channel context is local.
  • Method Details

    • handleException

      public void handleException(Exception exception)
    • serializeBundle

      public AbstractTaskBundleMessage serializeBundle(ServerTaskBundleNode bundle) throws Exception
      Serialize specified bundle into a message to send.
      Parameters:
      bundle - the byndle to process.
      Returns:
      the created message.
      Throws:
      Exception - if any error occurs.
    • deserializeBundle

      public NodeBundleResults deserializeBundle(AbstractTaskBundleMessage message) throws Exception
      Deserialize a task bundle from the message read into this buffer.
      Parameters:
      message - the message to process.
      Returns:
      a pairing of the received result head and the serialized tasks.
      Throws:
      Exception - if an error occurs during the deserialization.
    • newMessage

      public AbstractTaskBundleMessage newMessage()
      Create a new message.
      Returns:
      an AbstractTaskBundleMessage instance.
    • readMessage

      public boolean readMessage() throws Exception
      Throws:
      Exception
    • writeMessage

      public boolean writeMessage() throws Exception
      Throws:
      Exception
    • addJobEntry

      public void addJobEntry(ServerTaskBundleNode bundle)
      Add a new job to the job map.
      Parameters:
      bundle - the job to add.
    • getJobEntry

      public ServerTaskBundleNode getJobEntry(String uuid, long bundleId)
      Retrieve the job entry with the specified id.
      Parameters:
      uuid - the job uuid.
      bundleId - the id of the bundle to remove.
      Returns:
      a ServerTaskBundleNode instance, or null if there is no entry with the specified id.
    • removeJobEntry

      public ServerTaskBundleNode removeJobEntry(String uuid, long bundleId)
      Remove the job entry with the specified id.
      Parameters:
      uuid - the job uuid.
      bundleId - the id of the bundle to remove.
      Returns:
      the removed ServerTaskBundleNode instance, or null if there is no entry with the specified id.
    • nextMessageToSend

      protected AbstractTaskBundleMessage nextMessageToSend()
      Overrides:
      nextMessageToSend in class org.jppf.nio.AbstractNioContext
    • takeNextMessageToSend

      public AbstractTaskBundleMessage takeNextMessageToSend() throws InterruptedException
      Take a message from the pending queue, waiting if necessary.
      Returns:
      a AbstractTaskBundleMessage instance.
      Throws:
      InterruptedException - if the thread is interrupted.
    • toString

      public String toString()
      Overrides:
      toString in class org.jppf.nio.AbstractNioContext
    • close

      public void close()
    • getMonitor

      public Object getMonitor()
    • getCurrentNbJobs

      public int getCurrentNbJobs()
    • submit

      public Future<?> submit(ServerTaskBundleNode nodeBundle) throws Exception
      Throws:
      Exception
    • getServer

      public AsyncNodeNioServer getServer()
      Returns:
      the server handling this channel context.
    • getLocalNodeReadLock

      public org.jppf.utils.concurrent.ThreadSynchronization getLocalNodeReadLock()
      Returns:
      a lock used to synchronize input I/O with a local node.
    • getLocalNodeWriteLock

      public org.jppf.utils.concurrent.ThreadSynchronization getLocalNodeWriteLock()
      Returns:
      a lock used to synchronize output I/O with a local node.
    • getMaxJobs

      public int getMaxJobs()
    • setMaxJobs

      public void setMaxJobs(int maxJobs)
      Set the maximum number of concurrent jobs for this channel.
      Parameters:
      maxJobs - the max number of jobs to set.
    • getNbBundlesForJob

      public int getNbBundlesForJob(String jobUuid)
      Get the number of dispatches to the node for the specified job.
      Parameters:
      jobUuid - the uuid of the job to check.
      Returns:
      the number of dispatches of the job.
    • isAcceptingNewJobs

      public boolean isAcceptingNewJobs()
      Returns:
      whether the job is accepting new jobs.
    • setAcceptingNewJobs

      public void setAcceptingNewJobs(boolean acceptingNewJobs)
      Parameters:
      acceptingNewJobs - whether the job is accepting new jobs.