2 * Copyright 2011 GRNET S.A. All rights reserved.
4 * Redistribution and use in source and binary forms, with or
5 * without modification, are permitted provided that the following
8 * 1. Redistributions of source code must retain the above
9 * copyright notice, this list of conditions and the following
12 * 2. 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
15 * provided with the distribution.
17 * THIS SOFTWARE IS PROVIDED BY GRNET S.A. ``AS IS'' AND ANY EXPRESS
18 * OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
19 * WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR
20 * PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL GRNET S.A OR
21 * CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
22 * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT
23 * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF
24 * USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED
25 * AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT
26 * LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN
27 * ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE
28 * POSSIBILITY OF SUCH DAMAGE.
30 * The views and conclusions contained in the software and
31 * documentation are those of the authors and should not be
32 * interpreted as representing official policies, either expressed
33 * or implied, of GRNET S.A.
36 package gr.grnet.aquarium.messaging
39 import akka.amqp.{Topic, AMQP}
40 import akka.amqp.AMQP._
41 import gr.grnet.aquarium.MasterConf
42 import com.rabbitmq.client.Address
43 import gr.grnet.aquarium.util.Loggable
46 * Functionality for working with queues.
48 * @author Georgios Gousios <gousiosg@gmail.com>
50 trait AkkaAMQP extends Loggable {
52 class AMQPConnection {
53 private[messaging] lazy val connection = {
55 val servers = MasterConf.MasterConf.get(MasterConf.Keys.amqp_servers)
56 val port = MasterConf.MasterConf.get(MasterConf.Keys.amqp_port).toInt
58 val addresses = servers.split(",").foldLeft(Array[Address]()) {
59 (x, y) => x ++ Array(new Address(y, port))
65 MasterConf.MasterConf.get(MasterConf.Keys.amqp_username),
66 MasterConf.MasterConf.get(MasterConf.Keys.amqp_password),
67 MasterConf.MasterConf.get(MasterConf.Keys.amqp_vhost),
73 private lazy val exchanges = {
74 MasterConf.MasterConf.get(MasterConf.Keys.amqp_exchanges).split(",")
77 //Queues and exchnages are by default durable and persistent
78 val decl = ActiveDeclaration(durable = true, autoDelete = false)
80 def consumer(routekey: String, queue: String, exchange: String,
81 recipient: ActorRef, selfAck: Boolean) =
83 connection = (new AMQPConnection()).connection,
84 consumerParameters = ConsumerParameters(
85 routingKey = routekey,
86 exchangeParameters = Some(ExchangeParameters(exchange, Topic, decl)),
87 deliveryHandler = recipient,
88 queueName = Some(queue),
89 queueDeclaration = decl,
90 selfAcknowledging = selfAck
93 def producer(exchange: String) = {
95 if (!exchanges.contains(exchange))
96 logger.warn("Exchange %s is unknown".format(exchange))
99 connection = (new AMQPConnection()).connection,
100 producerParameters = ProducerParameters(
101 exchangeParameters = Some(ExchangeParameters(exchange, Topic, decl)),
102 channelParameters = Some(ChannelParameters(prefetchSize = 0))))