Bressam icon

AlgoritmoBully

Bressam | PRO | 06/28/18 09:06:24 AM UTC | 0 ⭐ | 273 👁️ | Never ⏰ | []
Java |

22.74 KB

|

None

|

0 👍

/

0 👎

/*
 * To change this license header, choose License Headers in Project Properties.
 * To change this template file, choose Tools | Templates
 * and open the template in the editor.
 */
 
/**
 *
 * @author giovanne.gaspareto
 */
/*
 * To change this license header, choose License Headers in Project Properties.
 * To change this template file, choose Tools | Templates
 * and open the template in the editor.
 */
 
import java.io.BufferedOutputStream;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
//import java.io.EOFException;
//import java.io.IOException;
//import java.io.ObjectInputStream;
//import java.net.UnknownHostException;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.ObjectOutputStream;
import java.net.DatagramPacket;
import java.net.DatagramSocket;
import java.net.InetAddress;
import java.net.MulticastSocket;
import java.net.SocketException;
import java.net.SocketTimeoutException;
import java.util.Scanner;
import java.util.logging.Level;
import java.util.logging.Logger;
 
public class ProcessoSender {
    static DatagramPacket msg;
    static DatagramPacket recv;
    static DatagramPacket datagramLider;
    static InetAddress grupo;
    static MulticastSocket s;
    private static Scanner keyboard;
    
    static int contadorMensagens = 0;
    static final String GROUP_IP = "224.0.0.1";
    static final int MULTICAST_PORT = 2000;
    static int IDProcesso = 0;
    static int IDLider = 0;
    static InetAddress IPLider;
    static int PortLider;
    static int IDRecebido = 0;
    static boolean aguardarResposta = false;
    static boolean pararExecucao = false;
    static int tempoDeVerificacao = 15;
    /*
    // Leitura de teclado
    static String lerGrupo = "1";
    static String enviarMensagem = "2";
    static String verificarCoordenador = "3";
    */
    static boolean eleicaoEmAndamento = false;
    static boolean isLeader = false;
    static String mensagem = "";
    static final String OK_STRING = "ok";
    static final String DESCONECTAR_STRING = "q";
    static final String ELEICAO_STRING = "e";
    static final String VERIFICA_COORDENADOR_STRING = "v";
    static final String DIVISOR_STRING = ";";
    static final String OK_LIDER = "okL";
    static final String TIMEOUT_STRING = "";
    static final String LIDER_ELEITO_STRING = "L";
    static final int DESCONECTAR = 1;
    static final int REALIZAR_ELEICAO = 2;
    static final int VERIFICA_COORDENADOR = 3;
    static final int RESPONDER_ELEICAO = 4;
    static final int LIDER_ELEITO = 5;
    static final int OK = 6;
    static final int RESPONDER_VERIFICACAO_LIDER = 7;
    static final int IGNORAR = 0;
    static final int ERRO_SOCKET = -1;
    static final String MENU = "1-Ler grupo\n"
                                + "2-Enviar mensagem\n"
                                + "3-Verificar Coordenador\n"
                                + "q-Encerrar processo";
    
    
    //Unicast
    static DatagramSocket serverSocket;
    static final String LIDER_EXISTE_STRING = "lider";
    static final String PERGUNTA_LIDER_VIVO = "lider?";
    static final int TIMEOUT = 10;
    //
    
    /**
     *
     * @param msgStr
     */
    public static void sendMessageGroup(String msgStr) {
        msg=new DatagramPacket(msgStr.getBytes(), msgStr.length(),grupo, 2000);
       // System.out.println("Enviando mensagem para o grupo ...");
        try {
            s.send(msg);
        } catch (IOException e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
        }
    }
 
    /*
    public static void receiveMessage(int i) {
        byte[] buf=new byte[1000];
        DatagramPacket recv=new DatagramPacket(buf, buf.length);
        System.err.print("Aguardando mensagem "+i+" ...");
        try {
            s.receive(recv);
            byte[] dest=new byte[recv.getLength()];
            System.arraycopy(recv.getData(), 0, dest, 0, dest.length);
            System.err.println("Mensagen recebida: "+
                    new String(recv.getData(),recv.getOffset(),recv.getLength()));
        } catch (IOException e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
        }
    }
   */
    /*
    public static void sendObj() {
       
        try {
             //Prepare Data
            String message = "Hello there!";
            ByteArrayOutputStream baos = new ByteArrayOutputStream();
            ObjectOutputStream oos;
            oos = new ObjectOutputStream(baos);
            oos.writeObject(message);
            byte[] data = baos.toByteArray();
     
            //Send data
            s.send(new DatagramPacket(data, data.length, grupo, multicastPort));
        } catch (IOException e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
        }
    }
    */
    
    //Funcao de configuracao
    public static void setUp() {
        try {
            System.setProperty("java.net.preferIPv4Stack", "true"); // somente pra mac no wifi
            grupo=InetAddress.getByName(GROUP_IP);
            s=new MulticastSocket(MULTICAST_PORT);
            System.err.println("Entrando no grupo ...");
            s.joinGroup(grupo);
            System.err.println("Ok.");
           // s.setLoopbackMode(true);  //evita de ler a propria mensagem
        } catch (IOException e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
        }
       
    }
    
    public static void defineIDProcesso(){
        keyboard = new Scanner(System.in);
        System.out.println("Informe o ID deste processo: ");
        IDProcesso = keyboard.nextInt();
    }
    
 
    /*
    public static int opcoesTeclado(){
        keyboard = new Scanner(System.in);
        System.out.println(menu);
        mensagem = keyboard.nextLine();
        if(mensagem.equals(lerGrupo)){
            contadorMensagens++;
            receiveMessage(contadorMensagens);
        }else if(mensagem.equals(enviarMensagem)){
            //sendObj();
            keyboard = new Scanner(System.in);
            System.out.println("Mensagem a ser enviada: ");
            mensagem = keyboard.nextLine();
            sendMessage(mensagem);
        }else if(mensagem.equals(verificarCoordenador)){
            //System.out.println("Opcao 3 selecionada");
            if(verificaCoordenador()){
                //Achou coordenador
             }else{
                //Nao achou coordenador, faz eleicao
             }
        }else if(mensagem.equalsIgnoreCase(sair)){
            break;
        }
    }
    */
    public static String receiveMessageGroup(){
        byte[] buf=new byte[1000];
        recv=new DatagramPacket(buf, buf.length);
        try {
            s.setSoTimeout(10);
            s.setLoopbackMode(false);
           // System.out.println("timeout no recevei sem nada: " + s.getSoTimeout());
            s.receive(recv);
            byte[] dest=new byte[recv.getLength()];
            System.arraycopy(recv.getData(), 0, dest, 0, dest.length);
            String mensagemRecebida = new String(recv.getData(),recv.getOffset(),recv.getLength());
            //System.err.println("Mensagen recebida: "+ mensagemRecebida);
            return mensagemRecebida;
        } catch (SocketTimeoutException ex) {
          //  Logger.getLogger(ProcessoSender.class.getName()).log(Level.SEVERE, null, ex);
            return "x;0;0";
    } catch (IOException e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
            return "x;0;0";
        }
    }
    
    public static int trataMensagem(String mensagemRecebida, DatagramPacket PacketLider){
        String[] partes = mensagemRecebida.split(DIVISOR_STRING);
        String operacao = partes[0];
        String mensagem = partes[1];
        IDRecebido = Integer.parseInt(mensagem);
        if(operacao.equals(LIDER_ELEITO_STRING)){
            //IDRecebido = Integer.parseInt(mensagem);
            setPortIPLider(recv,mensagem);
            return LIDER_ELEITO;
        }
        else if(operacao.equals(VERIFICA_COORDENADOR_STRING)){
            return RESPONDER_VERIFICACAO_LIDER;
        }else if(operacao.equals(ELEICAO_STRING)){
            if(Integer.parseInt(mensagem) == IDProcesso){ //foi o proprio processoq que iniciou a eleicao, ignora sua propria mensagem
                return IGNORAR;
            }
            else{
                IDRecebido = Integer.parseInt(mensagem);
                return RESPONDER_ELEICAO;
            }
        }else if(operacao.equals(DESCONECTAR_STRING)){
            if(Integer.parseInt(mensagem) == IDProcesso){
                return DESCONECTAR;
            }
            else return IGNORAR;
        }else if(operacao.equals(OK_STRING)){
            if(Integer.parseInt(mensagem) == IDProcesso){ //se foi ele mesmo que enviou o OK ou se nao esta mais na eleicao, ignora
                return IGNORAR;
            }
            else{
               // eleicaoEmAndamento = false;  //processo para de participar da eleicao, tem alguem maior
               //setPortIPLider(PacketLider, mensagem);
               String idRecebido = partes[2];
               IDRecebido = Integer.parseInt(idRecebido);
               
               if(Integer.parseInt(idRecebido) == IDProcesso){
                 return OK;
               }
               else{
                return IGNORAR;
               }
            }
        }else if(operacao.equals("")){
            return ERRO_SOCKET;
        }else{
            // mensagem ignorada
            return IGNORAR;
        }
    }
    
    public static void setPortIPLider(DatagramPacket PacketLider, String mensagemIDLider){
        IDLider = Integer.parseInt(mensagemIDLider);
        IPLider = PacketLider.getAddress();
        PortLider = PacketLider.getPort();
        
    }
    
    public static void realizarEleicao(){
        //Quando coordenador for eleito, envia mensagem para o grupo,
        //dessa mensagem todos pegaram o IP e porta para poder realizar a conexão unicast
        
        //Etapas: Enviar mensagem de eleicao = e;ID. Todos recebem, se tiverem ID menor ignora
        // Quem tiver ID maior ou por algum motivo tiver flag de lider ligada, responde
        mensagem = ELEICAO_STRING+DIVISOR_STRING+IDProcesso;
        sendMessageGroup(mensagem);
        eleicaoEmAndamento = true;
    }
    
    public static boolean responderEleicao(){
        //IDRecebido = Integer.parseInt(mensagemRecebida);
        if(IDProcesso >= IDRecebido){
           // isLeader = true;
            mensagem = OK_STRING+DIVISOR_STRING+IDProcesso+DIVISOR_STRING+IDRecebido;
            sendMessageGroup(mensagem);
            return true;
        }
        else{
            //ignora mensagem
            eleicaoEmAndamento = false; //sai da eleicao
            isLeader = false;           //nao é mais o lider
            return false;
        }
    }
    
    public static boolean verificaCoordenador() throws SocketException{
        //envia mensagem para coordenador, se nao receber resposta dentro do timeout faz eleicao
        /*
        //Se receber resposta do coordenador
        if(liderVivo()){
            return true;
        }
        else {
            realizarEleicao();
            return false;
        }
        */
        mensagem = VERIFICA_COORDENADOR_STRING+DIVISOR_STRING+IDProcesso;
        sendMessageGroup(mensagem); //envia mensagem de verificaçao para o lider responder caso esteja vivo
        TimerProcessos tempoEsperaPorOkLider = new TimerProcessos(1);
        while(!tempoEsperaPorOkLider.shouldCreate)
        {
           mensagem = receiveMessageGroup();  // espera resposta do lider, pelo delay/timeout definido
           if (!mensagem.equals("")){ //teve resposta, verificar qual foi
               String[] partes = mensagem.split(DIVISOR_STRING);
               String operacao = partes[0];
               String conteudo = partes[1];
               if(operacao.equals(OK_LIDER)){ //lider respondeu
                   tempoEsperaPorOkLider.cancelTimer();
                   setPortIPLider(recv,conteudo);
                   return true;
               }
           }
           else{
               tempoEsperaPorOkLider.cancelTimer();
               return false; //teve timeout, nao existe lider
           }
        }
        tempoEsperaPorOkLider.cancelTimer();
        return false;
    }
 
    public static boolean liderVivo(){
        try {
            serverSocket = new DatagramSocket();
            String toSend= PERGUNTA_LIDER_VIVO;
            byte[] buffer=toSend.getBytes();
            DatagramPacket pergunta=new DatagramPacket(buffer, buffer.length,
                        IPLider, PortLider);
            serverSocket.send(pergunta);
            s.setSoTimeout(TIMEOUT);
            while(true){
            s.receive(pergunta);
            String mensagemRecebida = new String(pergunta.getData(),pergunta.getOffset(),pergunta.getLength());
                if(mensagemRecebida.equals(LIDER_EXISTE_STRING)){
                    serverSocket.close();
                  return true; //Se a mensagem for a confirmacao do lider, sai da espera
                }
            }
        } catch (SocketException ex) {
            Logger.getLogger(ProcessoSender.class.getName()).log(Level.SEVERE, null, ex);
        } catch (IOException ex) {
            Logger.getLogger(ProcessoSender.class.getName()).log(Level.SEVERE, null, ex);
        }
        serverSocket.close();
        return true;
    }
    
    public static String receiveMessageGroupWithDelay() throws SocketException{
        byte[] buf=new byte[1000];
        recv=new DatagramPacket(buf, buf.length);
        try {
            s.setLoopbackMode(true);  //evita de ler a propria mensagem
            s.setSoTimeout(TIMEOUT);
           // System.out.println("Timeout no receive com timeout: " + s.getSoTimeout());
            s.receive(recv);
            byte[] dest=new byte[recv.getLength()];
            System.arraycopy(recv.getData(), 0, dest, 0, dest.length);
            String mensagemRecebida = new String(recv.getData(),recv.getOffset(),recv.getLength());
            //System.err.println("Mensagen recebida do delay: "+ mensagemRecebida);
            s.setLoopbackMode(false);  //habilita ler a propria mensagem
            return mensagemRecebida;
        }catch (SocketTimeoutException e) {
            // timeout exception.
            //System.out.println("Timeout reached!!! " + e);
            //System.out.println("TIMEOUUTTT");
            s.setSoTimeout(0);
            s.setLoopbackMode(false);
            return "x;0;0";
        } catch (IOException ex) {
            Logger.getLogger(ProcessoSender.class.getName()).log(Level.SEVERE, null, ex);
            s.setSoTimeout(0);
            s.setLoopbackMode(false);
            return "x;0;0";
        }
    }
    
    public static void main(String[] args) {
        try {
            defineIDProcesso();
            setUp();
            TimerProcessos timerVerificacaoLider = new TimerProcessos(tempoDeVerificacao);
            //sendMessage("Novo membro no grupo");
            while(!pararExecucao) {
                //If passou x segundos - Verificar lider
                if (eleicaoEmAndamento){
                   // timerVerificacaoLider.cancelTimer();
                    boolean idMaior = responderEleicao();
                    if(idMaior){
                        realizarEleicao();
                        // timeout nao funciona ja que esta tud9o nomesmo ip. Usar classe de timer e, enquanto o timer nao der o valor de timeout, aguardar algo que seja "OK"
                       TimerProcessos tProcesso = new TimerProcessos(1);
                       aguardarResposta = true;
                        while(!tProcesso.shouldCreate && aguardarResposta){
                            //System.out.println("Aguardando resposta da eleicao");
                            switch (trataMensagem(receiveMessageGroupWithDelay(),recv)) {
                                case OK: //teve Ok de outro processo
                                    isLeader = false;
                                    aguardarResposta = false;
                                  //  System.out.println("OK da eleicao recebida");
                                    break;
                                case RESPONDER_ELEICAO:
                                    idMaior = responderEleicao();
                                   // System.out.println("Respondendo eleicao");
                                    if(!idMaior){
                                        aguardarResposta = false;
                                        isLeader = false;
                                    }
                                    break;
                                case LIDER_ELEITO:
                                    if(IDRecebido>IDProcesso){
                                        aguardarResposta = false;
                                        isLeader = false;
                                    }
                                    break;
                                default:
                                  //  System.out.println("default da eleicao");
                                    //sendMessageGroup("a;0");
                                    break;
                            }
                        }
                        //if(tProcesso.shouldCreate){
                        if(aguardarResposta){ //se saiu enquanto aguardava resposta, é lider
                            System.out.println("Processo "+IDProcesso+" é novo lider");
                            isLeader=true;
                            mensagem = LIDER_ELEITO_STRING+DIVISOR_STRING+IDProcesso;
                            sendMessageGroup(mensagem);
                        }else{
                            System.out.println("Nao é lider");
                            isLeader = false;
                        }
                        tProcesso.cancelTimer();
                    }
                    eleicaoEmAndamento = false;
                    System.out.println("Eleicao acabou");
                    timerVerificacaoLider = new TimerProcessos(tempoDeVerificacao);
 
                }
                else{
                  //  System.out.println("Voltou a ouvir mensagens");
                  //  System.out.println("Timeout no main: " + s.getSoTimeout());
                  if (timerVerificacaoLider.shouldCreate && !isLeader && !eleicaoEmAndamento){
                      //acabou tempo, verificar lider
                    boolean liderVivo = verificaCoordenador();
                    
                    if(!liderVivo){
                        System.out.println("Nao ha lider!");
                        eleicaoEmAndamento = true; //inicia eleicao
                        mensagem = ELEICAO_STRING+DIVISOR_STRING+IDProcesso;
                        sendMessageGroup(mensagem);
                    }else{
                        System.out.println("Lider existe! Lider: [Processo "+IDLider+"]");
                    }
                    
                    timerVerificacaoLider = new TimerProcessos(tempoDeVerificacao); // novo timer
                  }
                    switch(trataMensagem(receiveMessageGroup(), recv)){
                    case VERIFICA_COORDENADOR:
                            //Verificar Coordenador
                        /*
                            boolean liderVivo = verificaCoordenador();
                            if(!liderVivo){
                                System.out.println("Nao ha lider!");
                                eleicaoEmAndamento = true; //inicia eleicao
                                mensagem = ELEICAO_STRING+DIVISOR_STRING+IDProcesso;
                                sendMessageGroup(mensagem);
                            }
*/
                            break;
                    case RESPONDER_ELEICAO:
                            //responderEleicao();
                            if(IDRecebido < IDProcesso){
                                System.out.println("Eleicao em andamento");
                                eleicaoEmAndamento = true;
                            }
                            break;
                    case LIDER_ELEITO:
                        //System.out.println("ID recebido: "+IDRecebido);
                        if(IDRecebido < IDProcesso)
                        {
                            eleicaoEmAndamento = true;
                        }
                        else{
                            eleicaoEmAndamento = false;
                            if(IDLider != IDProcesso){
                                isLeader = false;
                                System.out.println("[Processo "+IDProcesso+"] Nao sou lider!");
                            }
                            System.out.println("Novo lider eleito, ID do lider: "+IDLider
                                                + ", IP do lider: " +IPLider
                                                + ", Porta do lider: " +PortLider);
                        }
                            break;
                    case RESPONDER_VERIFICACAO_LIDER:
                            if(isLeader){
                                mensagem = OK_LIDER+DIVISOR_STRING+IDProcesso;
                                sendMessageGroup(mensagem);
                                System.out.println("Respondendo verificacao de lider vivo");
                                System.out.println("[Processo "+IDProcesso+"] Sou lider!");
                            }
                            else{
                                System.out.println("[Processo "+IDProcesso+"] Nao sou lider!");
                            }
                            break;
                    case IGNORAR:
                    //  System.err.print("mensagem ignorada");
                            break;
                    case ERRO_SOCKET:
                            System.err.print("Mensagem nao recebida");
                            break;
                    case DESCONECTAR:
                            System.err.println("Processo "+IDProcesso+" saindo do grupo...");
                            s.leaveGroup(grupo);
                            System.err.println("Desconectado do grupo.");
                            pararExecucao = true;
                            break;
                    default:
                           // System.err.print("Case default");
                            break;
                    }
                }
            }
        } catch(IOException exc) {
            exc.printStackTrace();
        }
        System.out.println("Processo encerrado");
        System.exit(0);
    }
}

Comments

  •  icon
    01/01/70 12:00:00 AM UTC
    Plain Text |

    0 B

    |

    👍

    /

    👎