root/src/java/src/com/omniti/reconnoiter/broker/RabbitBroker.java

Revision d0a64b649e4eac288431ac20d986830b57fb044d, 8.0 kB (checked in by Theo Schlossnagle <jesus@omniti.com>, 5 years ago)

Add support for riemann as the IEP subsystem.
Remove all traces of Esper.
Change the license on all our bits to simply match reconnoiter.
Cleanup copyrights and embelish auditing script.
Updated test 108 to check riemann iep results.

  • Property mode set to 100644
Line 
1 /*
2  * Copyright (c) 2013, Circonus, Inc. All rights reserved.
3  * Copyright (c) 2010, OmniTI Computer Consulting, Inc.
4  * All rights reserved.
5  *
6  * Redistribution and use in source and binary forms, with or without
7  * modification, are permitted provided that the following conditions are
8  * met:
9  *
10  *     * Redistributions of source code must retain the above copyright
11  *       notice, this list of conditions and the following disclaimer.
12  *     * Redistributions in binary form must reproduce the above
13  *       copyright notice, this list of conditions and the following
14  *       disclaimer in the documentation and/or other materials provided
15  *       with the distribution.
16  *     * Neither the name OmniTI Computer Consulting, Inc. nor the names
17  *       of its contributors may be used to endorse or promote products
18  *       derived from this software without specific prior written
19  *       permission.
20  *
21  * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
22  * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
23  * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
24  * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
25  * OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
26  * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT
27  * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,
28  * DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY
29  * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
30  * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
31  * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
32  */
33
34 package com.omniti.reconnoiter.broker;
35
36 import java.io.IOException;
37
38 import org.apache.log4j.Logger;
39 import com.omniti.reconnoiter.IEventHandler;
40 import com.omniti.reconnoiter.StratconConfig;
41 import com.rabbitmq.client.Connection;
42 import com.rabbitmq.client.Channel;
43 import com.rabbitmq.client.MessageProperties;
44 import com.rabbitmq.client.ConnectionFactory;
45 import com.rabbitmq.client.QueueingConsumer;
46
47
48 public class RabbitBroker implements IMQBroker  {
49   static Logger logger = Logger.getLogger(RabbitBroker.class.getName());
50   private int cidx;
51   private Connection conn;
52   private Channel channel;
53   private boolean noAck = true;
54   private String userName;
55   private String password;
56   private String virtualHost;
57   private String hostName[];
58   private ConnectionFactory factory[];
59   private int portNumber;
60   private String exchangeName;
61   private String exchangeType;
62   private String declaredQueueName;
63   private String queueName;
64   private String returnedQueueName;
65   private String routingKey;
66   private String alertRoutingKey;
67   private String alertExchangeName;
68   private Integer heartBeat;
69   private Integer connectTimeout;
70   private boolean exclusiveQueue;
71   private boolean durableQueue;
72   private boolean durableExchange;
73
74   public RabbitBroker(StratconConfig config) {
75     this.conn = null;
76     this.cidx = 0;
77     this.userName = config.getBrokerParameter("username", "guest");
78     this.password = config.getBrokerParameter("password", "guest");
79     this.virtualHost = config.getBrokerParameter("virtualhost", "/");
80     this.hostName = config.getBrokerParameter("hostname", "127.0.0.1").split(",");
81     this.portNumber = Integer.parseInt(config.getBrokerParameter("port", "5672"));
82     this.heartBeat = Integer.parseInt(config.getBrokerParameter("heartbeat", "5000"));
83     this.heartBeat = (this.heartBeat + 999) / 1000; // (ms -> seconds, rounding up)
84     this.connectTimeout = Integer.parseInt(config.getBrokerParameter("connect_timeout", "5000"));
85
86     this.exchangeType = config.getMQParameter("exchangetype", "fanout");
87     this.durableExchange = config.getMQParameter("durableexchange", "false").equals("true");
88     this.exchangeName = config.getMQParameter("exchange", "noit.firehose");
89     this.exclusiveQueue = config.getMQParameter("exclusivequeue", "false").equals("true");
90     this.durableQueue = config.getMQParameter("durablequeue", "false").equals("true");
91     this.declaredQueueName = config.getMQParameter("queue", "reconnoiter-{node}-{pid}");
92     this.routingKey = config.getMQParameter("routingkey", "");
93  
94     try {
95       this.queueName = this.declaredQueueName;
96       this.queueName = this.queueName.replace("{node}",
97         java.net.InetAddress.getLocalHost().getHostName());
98       this.queueName = this.queueName.replace("{pid}",
99         java.lang.management.ManagementFactory.getRuntimeMXBean()
100                                               .getName().replaceAll("@.+$", ""));
101     } catch(java.net.UnknownHostException uhe) {
102       System.err.println("Cannot self identify for queuename!");
103       System.exit(-1);
104     }     
105
106     this.alertRoutingKey = config.getBrokerParameter("routingkey", "noit.alerts.");
107     this.alertExchangeName = config.getBrokerParameter("exchange", "noit.alerts");
108
109     this.factory = new ConnectionFactory[this.hostName.length];
110     for(int i = 0; i<hostName.length; i++) {
111       this.factory[i] = new ConnectionFactory();
112       this.factory[i].setUsername(this.userName);
113       this.factory[i].setPassword(this.password);
114       this.factory[i].setVirtualHost(this.virtualHost);
115       this.factory[i].setRequestedHeartbeat(this.heartBeat);
116       this.factory[i].setConnectionTimeout(this.connectTimeout);
117       this.factory[i].setPort(this.portNumber);
118       this.factory[i].setHost(this.hostName[i]);
119     }
120   }
121  
122   //
123   public void disconnect() {
124     logger.info("AMQP disconnect.");
125     try {
126       channel.abort();
127     }
128     catch (Exception e) { }
129     channel = null;
130     try {
131       conn.abort();
132     }
133     catch (Exception e) { }
134     conn = null;
135   }
136   public void connect() throws Exception {
137     if(conn != null) disconnect();
138
139     int offset = ++cidx % factory.length;
140     logger.info("AMQP connect to " + this.hostName[offset]);
141     conn = factory[offset].newConnection();
142
143     if(conn == null) throw new Exception("connection failed");
144
145     channel = conn.createChannel();
146
147     boolean durable = durableExchange, exclusive = false,
148             internal = false, autoDelete = false;
149     channel.exchangeDeclare(exchangeName, exchangeType,
150                             durable, autoDelete, internal, null);
151
152     internal = false;
153     durable = true;
154     autoDelete = false;
155     channel.exchangeDeclare(getAlertExchangeName(), "topic", durable, autoDelete, internal, null);
156
157     autoDelete = true;
158     returnedQueueName = channel.queueDeclare(queueName, durableQueue,
159                                              exclusiveQueue, autoDelete, null).getQueue();
160     for (String rk : routingKey.split(",")) {
161         if ( rk.equalsIgnoreCase("null") ) rk = "";
162         channel.queueBind(returnedQueueName, exchangeName, rk);
163     }
164   }
165   public Channel getChannel() { return channel; }
166  
167   public void consume(IEventHandler eh) throws IOException {
168     QueueingConsumer consumer = new QueueingConsumer(channel);
169
170     channel.basicConsume(returnedQueueName, noAck, consumer);
171    
172     while (true) {
173       QueueingConsumer.Delivery delivery;
174       try {
175         delivery = consumer.nextDelivery();
176       } catch (InterruptedException ie) {
177         continue;
178       }
179       // (process the message components ...)
180    
181       String xml = new String(delivery.getBody());
182       try {
183         eh.processMessage(xml);
184       } catch (Exception e) {
185         // TODO Auto-generated catch block
186         e.printStackTrace();
187       }
188       if(!noAck)
189         channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
190     }
191   }
192
193   public String getAlertExchangeName() { return alertExchangeName; }
194   public String getAlertRoutingKey() { return alertRoutingKey; }
195   public void alert(String key, String json) {
196     String routingKey;
197     if(key == null) routingKey = getAlertRoutingKey();
198     else routingKey = getAlertRoutingKey() + key;
199     try {
200       byte[] messageBodyBytes = json.getBytes();
201       channel.basicPublish(getAlertExchangeName(), routingKey,
202                            MessageProperties.TEXT_PLAIN, messageBodyBytes);
203     } catch (Exception e) {
204       e.printStackTrace();
205     }
206   }
207 }
Note: See TracBrowser for help on using the browser.