de89d4013f16d23b65cc0010a286eeedcb8d9af5
[aquarium] / src / main / scala / gr / grnet / aquarium / service / ResourceEventProcessorService.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 import gr.grnet.aquarium.actor.DispatcherRole
39 import gr.grnet.aquarium.Configurator.Keys
40 import gr.grnet.aquarium.store.LocalFSEventStore
41 import com.ckkloverdos.maybe.{Maybe, Just, Failed, NoVal}
42 import gr.grnet.aquarium.actor.message.service.dispatcher.ProcessResourceEvent
43 import gr.grnet.aquarium.events.ResourceEvent
44 import gr.grnet.aquarium.util.date.TimeHelpers
45
46
47 /**
48  * An event processor service for resource events
49  *
50  * @author Georgios Gousios <gousiosg@gmail.com>
51  */
52 final class ResourceEventProcessorService extends EventProcessorService[ResourceEvent] {
53
54   override def parseJsonBytes(data: Array[Byte]) = ResourceEvent.fromBytes(data)
55
56   override def forward(event: ResourceEvent): Unit = {
57     if(event ne null) {
58       val businessLogicDispacther = _configurator.actorProvider.actorForRole(DispatcherRole)
59       businessLogicDispacther ! ProcessResourceEvent(event)
60     }
61   }
62
63   override def exists(event: ResourceEvent): Boolean =
64     _configurator.resourceEventStore.findResourceEventById(event.id).isJust
65
66   override def persist(event: ResourceEvent, initialPayload: Array[Byte]): Unit = {
67     // 1. Store to local FS for debugging purposes.
68     //    BUT be resilient to errors, since this is not critical
69     if(_configurator.eventsStoreFolder.isJust) {
70       Maybe {
71         LocalFSEventStore.storeResourceEvent(_configurator, event, initialPayload)
72       }
73     }
74
75     // 2. Store to DB
76     _configurator.resourceEventStore.storeResourceEvent(event)
77   }
78
79
80   protected def persistUnparsed(initialPayload: Array[Byte], exception: Throwable): Unit = {
81     // TODO: Also save to DB, just like we do for UserEvents
82     LocalFSEventStore.storeUnparsedResourceEvent(_configurator, initialPayload, exception)
83   }
84
85   override def queueReaderThreads: Int = 1
86
87   override def persisterThreads: Int = numCPUs + 4
88
89   override def numQueueActors: Int = 1 * queueReaderThreads
90
91   override def numPersisterActors: Int = 2 * persisterThreads
92
93   override def name = "resevtproc"
94
95   lazy val persister = new PersisterManager
96   lazy val queueReader = new QueueReaderManager
97
98   override def persisterManager = persister
99
100   override def queueReaderManager = queueReader
101
102   def start() {
103     logStarting()
104     val (ms0, ms1, _) = TimeHelpers.timed {
105       declareQueues(Keys.amqp_resevents_queues)
106     }
107     logStarted(ms0, ms1)
108   }
109
110   def stop() {
111     logStopping()
112     val (ms0, ms1, _) = TimeHelpers.timed {
113       queueReaderManager.stop()
114       persisterManager.stop()
115     }
116     logStopped(ms0, ms1)
117   }
118 }