e3e738a2f2fa6efdec0920b31b109aed1a1b04d6
[aquarium] / src / main / scala / gr / grnet / aquarium / service / IMEventProcessorService.scala
1 /*
2  * Copyright 2011-2012 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.service
37
38
39 import gr.grnet.aquarium.actor.RouterRole
40 import gr.grnet.aquarium.store.LocalFSEventStore
41 import gr.grnet.aquarium.util.date.TimeHelpers
42 import gr.grnet.aquarium.util.makeString
43 import com.ckkloverdos.maybe._
44 import gr.grnet.aquarium.actor.message.event.ProcessIMEvent
45 import gr.grnet.aquarium.event.model.im.{StdIMEvent, IMEventModel}
46 import gr.grnet.aquarium.connector.rabbitmq.service.RabbitMQService
47
48 /**
49  * An event processor service for user events coming from the IM system
50  *
51  * @author Georgios Gousios <gousiosg@gmail.com>
52  */
53 class IMEventProcessorService extends EventProcessorService[IMEventModel] {
54
55   override def parseJsonBytes(data: Array[Byte]) = {
56     StdIMEvent.fromJsonBytes(data)
57   }
58
59   override def forward(event: IMEventModel) = {
60     if(event ne null) {
61       _configurator.actorProvider.actorForRole(RouterRole) ! ProcessIMEvent(event)
62     }
63   }
64
65   override def existsInStore(event: IMEventModel) =
66     _configurator.imEventStore.findIMEventById(event.id).isDefined
67
68   override def storeParsedEvent(event: IMEventModel, initialPayload: Array[Byte]) = {
69     // 1. Store to local FS for debugging purposes.
70     //    BUT be resilient to errors, since this is not critical
71     if(_configurator.eventsStoreFolder.isJust) {
72       Maybe {
73         LocalFSEventStore.storeIMEvent(_configurator, event, initialPayload)
74       }
75     }
76
77     // 2. Store to DB
78     _configurator.imEventStore.insertIMEvent(event)
79   }
80
81   protected def storeUnparsedEvent(initialPayload: Array[Byte], exception: Throwable): Unit = {
82     val json = makeString(initialPayload)
83
84     LocalFSEventStore.storeUnparsedIMEvent(_configurator, initialPayload, exception)
85   }
86
87   override def queueReaderThreads: Int = 1
88
89   override def persisterThreads: Int = numCPUs
90
91   protected def numQueueActors = 2 * queueReaderThreads
92
93   protected def numPersisterActors = 2 * persisterThreads
94
95   override def name = "usrevtproc"
96
97   lazy val persister = new PersisterManager
98   lazy val queueReader = new QueueReaderManager
99
100   override def persisterManager = persister
101
102   override def queueReaderManager = queueReader
103
104   def start() {
105     logStarting()
106     val (ms0, ms1, _) = TimeHelpers.timed {
107       declareQueues(RabbitMQService.RabbitMQConfKeys.amqp_imevents_queues)
108     }
109     logStarted(ms0, ms1)
110   }
111
112   def stop() {
113     logStopping()
114     val (ms0, ms1, _) = TimeHelpers.timed {
115       queueReaderManager.stop()
116       persisterManager.stop()
117     }
118     logStopped(ms0, ms1)
119   }
120 }