jgroupsmember.java

来自「JGRoups源码」· Java 代码 · 共 468 行 · 第 1/2 页

JAVA
468
字号
    /**
     * A MembershipListener callback. It is called when a member is suspected being crashed,
     * but has not yet been excluded from the group.
     * @param suspectedMbr is the Address of the suspected crashed member.
     */
    public void suspect(Address suspectedMbr) {
        if(logger.isDebugEnabled())
            logger.debug("[" + getLocalAddress() + "] MembershipListener.suspect(Address=" +
                    suspectedMbr + ") is called notifying a suspected crashed member...");
    }

    /**
     * A MembershipListener callback. It is called when a change in group membership has
     * occurred, either new member joins, existing member leaves or has been suspected crashed.
     * @param mbrView is a View containing Address of all members in the group.
     */
    public void viewAccepted(View mbrView) {
        String addr=getLocalAddress();
        System.out.println("[" + addr + "] MembershipListener.viewAccepted(View) is called...");

        if(null == mbrView) {
            logger.warn("[" + addr + "] a null View is received.");
            return;
        }

        Vector mbrs=mbrView.getMembers();
            System.out.println("[" + addr + "] View has " + mbrs.size() + " members.");
            Address mbr=null;
            for(Iterator it=mbrs.iterator(); it.hasNext();) {
                mbr=(Address)it.next();
                logger.debug("[" + addr + "] View a member address: " + mbr);
            }
    }

    /**
     * A MessageListener callback. This method is called when another member invokes its
     * Channel.getState() requesting for group state, and this member is chosen to respond,
     * which would indicate that this member is a group coordinator.
     * @return It returns a byte[] containing a serialized Address of the server known
     *         by this member. It returns null when failed to serialize the Address.
     */
    public byte[] getState() {
        String addr=getLocalAddress();
        logger.info("[" + addr + "] is responding to a Channel.getState() inquiry from another" +
                " member...");

        try {
            byte[] result=null;
            logger.debug("[" + addr + "] waiting to lock serverAddress to read...");
            synchronized(serverAddressMutex) {
                result=Util.objectToByteBuffer(serverAddress);
                logger.debug("[" + addr + "] has read serverAddress;" +
                        " releasing serverAddress lock...");
            }
            logger.info("[" + addr + "] replying Channel.getState() inquiry...");
            return result;
        }
        catch(Exception ex) {
            logger.error("[" + addr + "] failed to serialize reply to the Channel.getState()" +
                    " inquiry; Exception: ", ex);
            logger.error("[" + addr + "] replying null to Channel.getState() inquiry...");
            return null;
        }
    }

    /**
     * A MessageListener callback. This method is called to receive a group state after this
     * member has invoked a Channel.getState() call. The state received is a serialized
     * Address of server known by the current group coordinator.
     * @param state is a byte[] containing an Address of server known by the current group
     *              coordinator. It could be null if the group coordinator failed to
     *              serialize the server address, or the coordinator does not have the
     *              server address when responding to this member's inquiry.
     */
    public void setState(byte[] state) {
        String addr=getLocalAddress();
        logger.info("[" + addr + "] is receiving a group state...");

        try {
            if(null != state) {
                logger.debug("[" + addr + "] waiting to lock serverAddress to write...");
                synchronized(serverAddressMutex) {
                    serverAddress=(Address)Util.objectFromByteBuffer(state);
                    logger.debug("[" + addr + "] has written to serverAddress;" +
                            " releasing serverAddress lock...");
                }
                if(logger.isDebugEnabled())
                    logger.debug("[" + addr + "] Got a state from a group coordinator that " +
                            "the server address is " + serverAddress);
            }
            else
                logger.error("[" + addr + "] Got null state.");
        }
        catch(Exception ex) {
            logger.error("[" + addr + "] failed to de-serialize the state; Exception: ", ex);
        }
    }

    /**
     * A MessageListener callback. It is called to receive a Message sent from another member.
     * @param message the Message received.
     */
    public void receive(Message message) {
        String addr=getLocalAddress();
        if(logger.isDebugEnabled())
            logger.debug("[" + addr + "] is receiving a Message via MessageListener...");

        if(null == message) {
            logger.error("[" + addr + "] got a null Message.");
            return;
        }

        Address src=message.getSrc();
        if(message.getObject() instanceof String) {
            String contents=(String)message.getObject();
            if(logger.isDebugEnabled()) {
                logger.debug("[" + addr + "] got a Message \"" + contents + "\" from member at " +
                        src.toString());
            }
        }
        else {
            if(logger.isDebugEnabled()) {
                logger.debug("[" + addr + "] got a Message from member at " + src.toString());
            }
        }
    }

    /**
     * A RequestHandler callback. It is called to receive a Message multi-casted by another
     * member in the group, e.g. by a MessageDispatcher.castMessage(Message).
     * @param message the Message received.
     * @return If it is a server, it returns a String "ACK". If it is a client,
     *         it returns null.
     */
    public Object handle(Message message) {
        String addr=getLocalAddress();
        if(logger.isDebugEnabled())
            logger.debug("[" + addr + "] is receiving a Message via RequestHandler...");

        if(null == message) {
            logger.error("[" + addr + "] got a null Message.");
            return null;
        }

        Address src=message.getSrc();
        if(message.getObject() instanceof String) {
            String contents=(String)message.getObject();
            if(logger.isDebugEnabled()) {
                logger.debug("[" + addr + "] got a Message \"" + contents + "\" from member at " +
                        src.toString());
            }
            if(contents.equals("server")) {
                logger.debug("[" + addr + "] waiting to lock serverAddress to write...");
                synchronized(serverAddressMutex) {
                    serverAddress=src;
                    if(logger.isDebugEnabled())
                        logger.debug("[" + addr + "] has written to serverAddress = " + src +
                                "; releasing serverAddress lock..");
                }
            }
            else if(contents.equals("stop")) {
                logger.debug("[" + addr + "] received a \"stop\" message;" +
                        " waiting to lock serverEnds...");
                synchronized(serverEnds) {
                    logger.debug("[" + addr + "] waking up server main thread...");
                    serverEnds.notifyAll();
                }
            }
        }
        else {
            if(logger.isDebugEnabled()) {
                logger.debug("[" + addr + "] got a Message from member at " + src.toString());
            }
        }

        if(role.equals("server")) {
            logger.debug("[" + addr + "] replying \"ACK\"...");
            return "ACK";
        }
        else {
            logger.debug("[" + addr + "] replying null...");
            return null;
        }
    }

/*
* Returns a String of the Address this member is connected to the JGroups. It returns
* null if the channel is not connected.
*/

    private String getLocalAddress() {
        if(null == channel)
            return null;
        else {
            Address addr=channel.getLocalAddress();
            if(null == addr) // the channel is in closed state
                return null;
            else
                return addr.toString();
        }
    }

    public static void main(String args[]) {
        String role, propFile, message;
        int numOfRepeats;
        JGroupsMember member=null;

        if(4 != args.length) {
            System.out.println(
                    "Usage: JGroupsMember <role> <configXmlFile> <message> <repeats>");
            System.out.println(" role = client or server.");
            System.out.println(" configXmlFile = the JGroups protocol stack config file.");
            System.out.println(" For example: config/total-token.xml");
            System.out.println(" message = the message to multi-cast to the group.");
            System.out.println(" A message \"stop\" from a client will stop");
            System.out.println(" all servers.");
            System.out.println(" repeats = if the role is \"client\", then how many");
            System.out.println(" times to repeat the sending of the message.");
            return;
        }
        else {
            role=args[0];
            propFile=args[1];
            message=args[2];
            numOfRepeats=Integer.parseInt(args[3]);
        }

        member=new JGroupsMember(role, propFile, message, numOfRepeats);
        member.init();
member.run();
}
}

⌨️ 快捷键说明

复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?