Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions build.sbt
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@ import Tests._

name := "Rabbit Consumer"

organization := "com.paddypowerbetfair"

version := "0.1"

scalaVersion := "2.12.3"
Expand Down Expand Up @@ -55,6 +57,11 @@ val mockito = Seq (
"org.mockito" % "mockito-core" % mockitoV % "test"
)

// to use for local build only
publishTo := Some(Resolver.file("file", new File(Resolver.mavenLocal.root)))



libraryDependencies ++= logging ++ scalacheck ++ scalatest ++ amqpClient ++ scalaz ++ argonaut ++ typesafeConfig ++ mockito

lazy val root = (project in file("."))
Expand Down
132 changes: 95 additions & 37 deletions docker/resources/rabbitmq/definitions.json
Original file line number Diff line number Diff line change
@@ -1,39 +1,97 @@
{
"rabbit_version": "3.6.9",
"users": [
{
"name": "guest",
"password_hash": "FkYR4/EokTofabv5M6BrcganRjk/OjgbjT70+yIz0Ms05liE",
"hashing_algorithm": "rabbit_password_hashing_sha256",
"tags": "administrator"
}
],
"vhosts": [
{
"name": "/"
}
],
"permissions": [
{
"user": "guest",
"vhost": "/",
"configure": ".*",
"write": ".*",
"read": ".*"
}
],
"parameters": [],
"policies": [],
"queues": [],
"exchanges": [
{
"name": "myExchange",
"vhost": "/",
"type": "fanout",
"durable": true,
"auto_delete": false,
"internal": false,
"arguments": {}
}
]
"rabbit_version":"3.7.2",
"users":[
{
"name":"guest",
"password_hash":"FkYR4/EokTofabv5M6BrcganRjk/OjgbjT70+yIz0Ms05liE",
"hashing_algorithm":"rabbit_password_hashing_sha256",
"tags":"administrator"
}
],
"vhosts":[
{
"name":"/"
}
],
"permissions":[
{
"user":"guest",
"vhost":"/",
"configure":".*",
"write":".*",
"read":".*"
}
],
"topic_permissions":[

],
"parameters":[

],
"global_parameters":[
{
"name":"cluster_name",
"value":"rabbit@ppb-rabbit-server"
}
],
"policies":[

],
"queues":[
{
"name":"myTestQueue",
"vhost":"/",
"durable":false,
"auto_delete":false,
"arguments":{

}
},

{
"name":"myItTestSecondQueue",
"vhost":"/",
"durable":false,
"auto_delete":false,
"arguments":{

}
}
],
"exchanges":[
{
"name":"myExchange",
"vhost":"/",
"type":"fanout",
"durable":false,
"auto_delete":false,
"internal":false,
"arguments":{

}
}
],
"bindings":[
{
"source":"myExchange",
"vhost":"/",
"destination":"myTestQueue",
"destination_type":"queue",
"routing_key":"myItTest",
"arguments":{

}
},

{
"source":"myExchange",
"vhost":"/",
"destination":"myItTestSecondQueue",
"destination_type":"queue",
"routing_key":"myItTestSecond",
"arguments":{

}
}
]
}
24 changes: 15 additions & 9 deletions src/it/resources/local.conf
Original file line number Diff line number Diff line change
@@ -1,26 +1,32 @@
amqp {
connections = [
{
ip = "localhost"
name = "uat"
ip = "127.0.0.1"
port = 5672
user = "guest"
password = "guest"
useSSL = false
vhost = "/",
exchangeName = "myExchange"
queue = "dwimyQueue"
routingKey = ""
fileName = "~/output1.json"
},
routingKey = "myItTest"
queue = "myTestQueue"
fileName = "target/myTestQueue.log"
}
{
ip = "localhost"
name = "uat"
ip = "127.0.0.1"
port = 5672
user = "guest"
password = "guest"
useSSL = false
vhost = "/",
exchangeName = "myExchange"
queue = "dwimyQueue"
routingKey = ""
fileName = "~/output2.json"
routingKey = "myItTestSecond"
queue = "myItTestSecondQueue"
fileName = "target/myItTestSecondQueue.log"
}
]
}


71 changes: 63 additions & 8 deletions src/it/scala/com/ppb/rabbitconsumer/RabbitConsumerITSpec.scala
Original file line number Diff line number Diff line change
@@ -1,18 +1,73 @@
package com.ppb.rabbitconsumer

import java.io.File
import java.nio.file.{Files, Paths}
import java.util.UUID

import com.ppb.rabbitconsumer.RabbitConsumerAlgebra.Configurations
import com.typesafe.config.{Config, ConfigFactory}
import org.scalatest._

import scala.collection.JavaConverters._
import argonaut._
import Argonaut._
import org.scalatest.{FlatSpec, Matchers}
import org.slf4j.LoggerFactory

import scala.io.Source
import scala.util.{Failure, Success, Try}
import scala.collection.JavaConverters._
import com.typesafe.config.{Config, ConfigFactory}
//import scalaz.stream.{Sink, _}


class RabbitConsumerITSpec extends FlatSpec {

val configFile:String = "src/it/resources/local.conf"
val configurations:Configurations = Configurations(configFile, ConfigFactory.parseFile(new File(configFile)).getConfigList("amqp.connections").asScala.toList)

val setUp:Config => Unit = (config) => {
println("Inside setup "+ config)
val rabbitConn:RabbitConnection = ConnectionService.rabbitConnection(config)
val queueName = ConfigService.readQueue(config)
rabbitConn.channel.queuePurge(queueName)
println(" Publishing dummy to queue "+ queueName)
for (i <- 1 to 3) {
val map:Map[String, String] = Map("QueueName"->queueName, "message" -> UUID.randomUUID().toString, "ID"-> i.toString)
rabbitConn.channel.basicPublish("", queueName, null, map.asJson.toString().getBytes())
}

class RabbitConsumerITSpec extends FlatSpec with Matchers {
RabbitConsumer.readFromFile(configFile)

println("closing connection")
rabbitConn.connection.close(1000)
Option(ConfigService.getFilename(config)) map(
outputfile => { Files.deleteIfExists(Paths.get(outputfile)) }
)getOrElse( println (" To Stdout") )

it should "read docker local configuration files and process it" in {
RabbitConsumer.local
}

val verifyOrFail:Config => Unit = (config) => {
Option(ConfigService.getFilename(config)) map(
outputfile => {
assert(Files.exists(Paths.get(outputfile)))
assert(Files.size(Paths.get(outputfile)) > 0)
println("============== "+ outputfile +"============") ;
scala.io.Source.fromFile(outputfile).getLines foreach println _
}
)getOrElse( )
}



configurations.configs.map(setUp)

it should "read stream message to local file it" in {
RabbitConsumer.readFromFile(configFile)

configurations.configs.map(verifyOrFail)

// configurations.configs.foreach(
// (x) => {
// println("============== "+x +"============") ;
// scala.io.Source.fromFile(ConfigService.getFilename(x)).getLines foreach println _
// }
// )
}

}
56 changes: 46 additions & 10 deletions src/main/scala/com/ppb/rabbitconsumer/ConnectionService.scala
Original file line number Diff line number Diff line change
@@ -1,16 +1,18 @@
package com.ppb.rabbitconsumer

import com.rabbitmq.client.ConnectionFactory
import com.typesafe.config.Config

import argonaut._
import com.ppb.rabbitconsumer.ConfigService.{getFilename, readExchange, readQueue, readRoutingKey}
import com.ppb.rabbitconsumer.RabbitConnection._
import com.ppb.rabbitconsumer.RabbitConsumerAlgebra.{SingleResponseConnection, SubscriptionConnection}
import com.rabbitmq.client.ConnectionFactory
import com.typesafe.config.Config
import org.slf4j.LoggerFactory
import scodec.bits.ByteVector

import scalaz.Sink
import scalaz.concurrent.Task
import scalaz.stream.{Sink, io, sink}


sealed trait RabbitResponse extends Product with Serializable
case object NoMoreMessages extends RabbitResponse
case class RabbitMessage(payload: Json) extends RabbitResponse

object ConnectionService {

Expand Down Expand Up @@ -40,19 +42,53 @@ object ConnectionService {
RabbitConnection(connection, channel, nextMessage)
}

def init(config: Config): Cxn = {

def createConnection(config: Config) : SingleResponseConnection = {
implicit val rabbitConnnection: RabbitConnection = rabbitConnection(config)
val queueName = readQueue(config)
val exchangeName = readExchange(config)
val routingKey = readRoutingKey(config)
bindQueueToExchange(queueName, exchangeName, routingKey)

val sink = Option(getFilename(config)) map {
fileNameDefined => io.fileChunkW(fileNameDefined)
} getOrElse io.stdOutBytes

SingleResponseConnection(() => RabbitConnection.readNextPayload(queueName), () => RabbitConnection.disconnect, sink)
}






// def init(config: Config, createQueue:Boolean = false, sink:Sink[Task, ByteVector]): SingleResponseConnection = {
// createConnection(config, createQueue, sink)
// }


def subscribe(config: Config, messageParser:MessageParser = PlainMessageParser): SubscriptionConnection = {
implicit val rabbitConnnection: RabbitConnection = rabbitConnection(config)

val queueName = readQueue(config)
val exchangeName = readExchange(config)
val routingKey = readRoutingKey(config)
val fileName = getFilename(config)

createQueue(queueName)
bindQueueToExchange(queueName, exchangeName, routingKey)

Cxn(getFilename(config), () => RabbitConnection.nextPayload(queueName), () => RabbitConnection.disconnect)


println(" Created a subscription with "+ queueName + " " +exchangeName+ " "+ routingKey )

SubscriptionConnection(RabbitConnection.registerListener(queueName, messageParser), ()=>RabbitConnection.disconnect, Option(fileName))


}




def done(config: Config): Unit = {
val queueName = readQueue(config)
logger.info(s"Deleting $queueName")
Expand Down
Loading