Remove dead code before migrating to single project setup
[aquarium] / logic / src / main / scala / gr / grnet / aquarium / messaging / AkkaAMQP.scala
1 /*
2  * Copyright 2011 GRNET S.A. All rights reserved.
3  *
4  * Redistribution and use in source and binary forms, with or
5  * without modification, are permitted provided that the following
6  * conditions are met:
7  *
8  *   1. Redistributions of source code must retain the above
9  *      copyright notice, this list of conditions and the following
10  *      disclaimer.
11  *
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.
16  *
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.
29  *
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.
34  */
35
36 package gr.grnet.aquarium.messaging
37
38 import akka.actor._
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
44
45 /**
46  * Functionality for working with queues.
47  *
48  * @author Georgios Gousios <gousiosg@gmail.com>
49  */
50 trait AkkaAMQP extends Loggable {
51
52   class AMQPConnection {
53     private[messaging] lazy val connection = {
54
55       val servers = MasterConf.MasterConf.get(MasterConf.Keys.amqp_servers)
56       val port = MasterConf.MasterConf.get(MasterConf.Keys.amqp_port).toInt
57
58       val addresses = servers.split(",").foldLeft(Array[Address]()) {
59         (x, y) => x ++ Array(new Address(y, port))
60       }
61
62       AMQP.newConnection(
63         ConnectionParameters(
64           addresses,
65           MasterConf.MasterConf.get(MasterConf.Keys.amqp_username),
66           MasterConf.MasterConf.get(MasterConf.Keys.amqp_password),
67           MasterConf.MasterConf.get(MasterConf.Keys.amqp_vhost),
68           1000,
69           None))
70     }
71   }
72
73   private lazy val exchanges = {
74     MasterConf.MasterConf.get(MasterConf.Keys.amqp_exchanges).split(",")
75   }
76
77   //Queues and exchnages are by default durable and persistent
78   val decl = ActiveDeclaration(durable = true, autoDelete = false)
79
80   def consumer(routekey: String, queue: String, exchange: String,
81                recipient: ActorRef, selfAck: Boolean) =
82     AMQP.newConsumer(
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
91         ))
92
93   def producer(exchange: String) = {
94     
95     if (!exchanges.contains(exchange))
96       logger.warn("Exchange %s is unknown".format(exchange))
97
98     AMQP.newProducer(
99       connection = (new AMQPConnection()).connection,
100       producerParameters = ProducerParameters(
101         exchangeParameters = Some(ExchangeParameters(exchange, Topic, decl)),
102         channelParameters = Some(ChannelParameters(prefetchSize = 0))))
103   }
104 }