/* * 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); } }