From 2b0beb989a6f15380df1aa912562c719422b8089 Mon Sep 17 00:00:00 2001 From: Ilya Brin Date: Mon, 29 Jan 2018 22:09:57 -0800 Subject: [PATCH 1/7] save the sender, and use it when answering --- project/build.properties | 1 + .../scala/com/github/sstone/amqp/ChannelOwner.scala | 12 ++++++------ src/main/scala/com/github/sstone/amqp/Consumer.scala | 6 +++--- .../scala/com/github/sstone/amqp/RpcServer.scala | 8 ++++---- src/main/scala/com/github/sstone/amqp/package.scala | 8 ++++++++ .../com/github/sstone/amqp/samples/Consumer2.scala | 2 +- 6 files changed, 23 insertions(+), 14 deletions(-) create mode 100644 project/build.properties create mode 100644 src/main/scala/com/github/sstone/amqp/package.scala diff --git a/project/build.properties b/project/build.properties new file mode 100644 index 0000000..c091b86 --- /dev/null +++ b/project/build.properties @@ -0,0 +1 @@ +sbt.version=0.13.16 diff --git a/src/main/scala/com/github/sstone/amqp/ChannelOwner.scala b/src/main/scala/com/github/sstone/amqp/ChannelOwner.scala index 807a00b..87ddd7e 100644 --- a/src/main/scala/com/github/sstone/amqp/ChannelOwner.scala +++ b/src/main/scala/com/github/sstone/amqp/ChannelOwner.scala @@ -21,7 +21,7 @@ object ChannelOwner { case class NotConnectedError(request: Request) - def props(init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None): Props = Props(new ChannelOwner(init, channelParams)) + def props(init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None): Props = Props(new ChannelOwner(init, channelParams)) private[amqp] class Forwarder(channel: Channel) extends Actor with ActorLogging { @@ -153,11 +153,11 @@ object ChannelOwner { } } -class ChannelOwner(init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None) extends Actor with ActorLogging { +class ChannelOwner(init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None) extends Actor with ActorLogging { import ChannelOwner._ - var requestLog: Vector[Request] = init.toVector + var requestLog: Vector[RequestAndSender] = init.toVector val statusListeners = mutable.HashSet.empty[ActorRef] override def preStart() = context.parent ! ConnectionOwner.CreateChannel @@ -185,13 +185,13 @@ class ChannelOwner(init: Seq[Request] = Seq.empty[Request], channelParams: Optio forwarder ! AddShutdownListener(self) forwarder ! AddReturnListener(self) onChannel(channel, forwarder) - requestLog.map(r => self forward r) + requestLog.foreach { case (request, sender) => self.tell(request, sender.getOrElse(context.sender())) } log.info(s"got channel $channel") statusListeners.map(a => a ! Connected) context.become(connected(channel, forwarder)) } case Record(request: Request) => { - requestLog :+= request + requestLog :+= request -> Some(sender()) } case AddStatusListener(actor) => addStatusListener(actor) @@ -203,7 +203,7 @@ class ChannelOwner(init: Seq[Request] = Seq.empty[Request], channelParams: Optio def connected(channel: Channel, forwarder: ActorRef): Receive = LoggingReceive { case Amqp.Ok(_, _) => () case Record(request: Request) => { - requestLog :+= request + requestLog :+= request -> Some(sender()) self forward request } case AddStatusListener(listener) => { diff --git a/src/main/scala/com/github/sstone/amqp/Consumer.scala b/src/main/scala/com/github/sstone/amqp/Consumer.scala index f896f6b..abb1ee5 100644 --- a/src/main/scala/com/github/sstone/amqp/Consumer.scala +++ b/src/main/scala/com/github/sstone/amqp/Consumer.scala @@ -9,12 +9,12 @@ import akka.event.LoggingReceive import scala.collection.JavaConversions._ object Consumer { - def props(listener: Option[ActorRef], autoack: Boolean = false, init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None, + def props(listener: Option[ActorRef], autoack: Boolean = false, init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None, consumerTag: String = "", noLocal: Boolean = false, exclusive: Boolean = false, arguments: Map[String, AnyRef] = Map.empty): Props = Props(new Consumer(listener, autoack, init, channelParams, consumerTag, noLocal, exclusive, arguments)) def props(listener: ActorRef, exchange: ExchangeParameters, queue: QueueParameters, routingKey: String, channelParams: Option[ChannelParameters], autoack: Boolean): Props = - props(Some(listener), init = List(AddBinding(Binding(exchange, queue, routingKey))), channelParams = channelParams, autoack = autoack) + props(Some(listener), init = List(AddBinding(Binding(exchange, queue, routingKey)) -> None), channelParams = channelParams, autoack = autoack) def props(listener: ActorRef, channelParams: Option[ChannelParameters], autoack: Boolean): Props = props(Some(listener), channelParams = channelParams, autoack = autoack) } @@ -32,7 +32,7 @@ object Consumer { */ class Consumer(listener: Option[ActorRef], autoack: Boolean = false, - init: Seq[Request] = Seq.empty[Request], + init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None, consumerTag: String = "", noLocal: Boolean = false, diff --git a/src/main/scala/com/github/sstone/amqp/RpcServer.scala b/src/main/scala/com/github/sstone/amqp/RpcServer.scala index 865ce30..b713d4b 100644 --- a/src/main/scala/com/github/sstone/amqp/RpcServer.scala +++ b/src/main/scala/com/github/sstone/amqp/RpcServer.scala @@ -39,14 +39,14 @@ object RpcServer { def onFailure(delivery: Delivery, e: Throwable): ProcessResult } - def props(processor: RpcServer.IProcessor, init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None)(implicit ctx: ExecutionContext): Props = + def props(processor: RpcServer.IProcessor, init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None)(implicit ctx: ExecutionContext): Props = Props(new RpcServer(processor, init, channelParams)) def props(queue: QueueParameters, exchange: ExchangeParameters, routingKey: String, proc: RpcServer.IProcessor, channelParams: ChannelParameters)(implicit ctx: ExecutionContext): Props = - props(processor = proc, init = List(AddBinding(Binding(exchange, queue, routingKey))), channelParams = Some(channelParams)) + props(processor = proc, init = List(AddBinding(Binding(exchange, queue, routingKey)) -> None), channelParams = Some(channelParams)) def props(queue: QueueParameters, exchange: ExchangeParameters, routingKey: String, proc: RpcServer.IProcessor)(implicit ctx: ExecutionContext): Props = - props(processor = proc, init = List(AddBinding(Binding(exchange, queue, routingKey)))) + props(processor = proc, init = List(AddBinding(Binding(exchange, queue, routingKey)) -> None)) } @@ -60,7 +60,7 @@ object RpcServer { * @param processor [[com.github.sstone.amqp.RpcServer.IProcessor]] implementation * @param channelParams optional channel parameters */ -class RpcServer(processor: RpcServer.IProcessor, init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None)(implicit ctx: ExecutionContext = ExecutionContext.Implicits.global) extends Consumer(listener = None, autoack = false, init = init, channelParams = channelParams) { +class RpcServer(processor: RpcServer.IProcessor, init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None)(implicit ctx: ExecutionContext = ExecutionContext.Implicits.global) extends Consumer(listener = None, autoack = false, init = init, channelParams = channelParams) { import RpcServer._ private def sendResponse(result: ProcessResult, properties: BasicProperties, channel: Channel) { diff --git a/src/main/scala/com/github/sstone/amqp/package.scala b/src/main/scala/com/github/sstone/amqp/package.scala new file mode 100644 index 0000000..58ea0de --- /dev/null +++ b/src/main/scala/com/github/sstone/amqp/package.scala @@ -0,0 +1,8 @@ +package com.github.sstone + +import akka.actor.ActorRef +import com.github.sstone.amqp.Amqp.Request + +package object amqp { + type RequestAndSender = (Request, Option[ActorRef]) +} diff --git a/src/main/scala/com/github/sstone/amqp/samples/Consumer2.scala b/src/main/scala/com/github/sstone/amqp/samples/Consumer2.scala index d9586e4..642b642 100644 --- a/src/main/scala/com/github/sstone/amqp/samples/Consumer2.scala +++ b/src/main/scala/com/github/sstone/amqp/samples/Consumer2.scala @@ -31,7 +31,7 @@ object Consumer2 extends App { // to the broker is lost and restored val consumer = ConnectionOwner.createChildActor(conn, Consumer.props( listener = Some(listener), - init = List(AddBinding(Binding(StandardExchanges.amqDirect, queueParams, "my_key"))) + init = List(AddBinding(Binding(StandardExchanges.amqDirect, queueParams, "my_key")) -> None) ), name = Some("consumer")) // wait till everyone is actually connected to the broker From 21bff16168d49af28232b2bbd2d4526ce92e973d Mon Sep 17 00:00:00 2001 From: Ilya Brin Date: Mon, 29 Jan 2018 22:24:59 -0800 Subject: [PATCH 2/7] tests compile --- build.sbt | 2 +- src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/build.sbt b/build.sbt index bf7803a..39f5d5c 100644 --- a/build.sbt +++ b/build.sbt @@ -6,7 +6,7 @@ version := "1.5-SNAPSHOT" scalaVersion := "2.11.7" -scalacOptions ++= Seq("-feature", "-language:postfixOps") +scalacOptions ++= Seq("-feature", "-language:postfixOps") resolvers += "Typesafe Repository" at "http://repo.typesafe.com/typesafe/releases/" diff --git a/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala b/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala index 37b9c1d..866b7de 100644 --- a/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala +++ b/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala @@ -168,7 +168,7 @@ class ConsumerSpec extends ChannelSpec { val probe = TestProbe() val queue = randomQueue val consumer = ConnectionOwner.createChildActor(conn, Consumer.props( - listener = Some(probe.ref), autoack = false, init = Seq.empty[Request], channelParams = None, + listener = Some(probe.ref), autoack = false, init = Seq.empty[RequestAndSender], channelParams = None, consumerTag = "", noLocal = false, exclusive = true, arguments = Map.empty), timeout = 5000 millis) val producer = ConnectionOwner.createChildActor(conn, ChannelOwner.props()) @@ -186,7 +186,7 @@ class ConsumerSpec extends ChannelSpec { assert(delivery.consumerTag === consumerTag) val consumer1 = ConnectionOwner.createChildActor(conn, Consumer.props( - listener = Some(probe.ref), autoack = false, init = Seq.empty[Request], channelParams = None, + listener = Some(probe.ref), autoack = false, init = Seq.empty[RequestAndSender], channelParams = None, consumerTag = "", noLocal = false, exclusive = true, arguments = Map.empty), timeout = 5000 millis) consumer1 ! AddStatusListener(probe.ref) From 248bb8afa6e32235be0d0fa5e3667ec36ef78be9 Mon Sep 17 00:00:00 2001 From: Ilya Brin Date: Wed, 31 Jan 2018 14:11:43 -0800 Subject: [PATCH 3/7] ignore publish oks --- src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala b/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala index 866b7de..07410bc 100644 --- a/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala +++ b/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala @@ -144,6 +144,9 @@ class ConsumerSpec extends ChannelSpec { probe.expectNoMsg() } "send consumer cancellation notifications" in { + ignoreMsg { + case Amqp.Ok(p:Publish, _) => true + } val probe = TestProbe() val queue = randomQueue val consumer = ConnectionOwner.createChildActor(conn, Consumer.props(listener = Some(probe.ref), autoack = false), timeout = 5000 millis) @@ -165,6 +168,9 @@ class ConsumerSpec extends ChannelSpec { probe.expectMsg(1 second, ConsumerCancelled(consumerTag)) } "create exclusive consumers" in { + ignoreMsg { + case Amqp.Ok(p:Publish, _) => true + } val probe = TestProbe() val queue = randomQueue val consumer = ConnectionOwner.createChildActor(conn, Consumer.props( From 21db6310cff0bbe4b7dc6cc1ddde61d71b40b059 Mon Sep 17 00:00:00 2001 From: Ilya Brin Date: Wed, 31 Jan 2018 14:15:29 -0800 Subject: [PATCH 4/7] inits dont need the sender anyways, so only store it internally --- src/main/scala/com/github/sstone/amqp/ChannelOwner.scala | 6 +++--- src/main/scala/com/github/sstone/amqp/Consumer.scala | 6 +++--- src/main/scala/com/github/sstone/amqp/RpcServer.scala | 8 ++++---- src/main/scala/com/github/sstone/amqp/package.scala | 8 -------- .../scala/com/github/sstone/amqp/samples/Consumer2.scala | 2 +- src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala | 4 ++-- 6 files changed, 13 insertions(+), 21 deletions(-) delete mode 100644 src/main/scala/com/github/sstone/amqp/package.scala diff --git a/src/main/scala/com/github/sstone/amqp/ChannelOwner.scala b/src/main/scala/com/github/sstone/amqp/ChannelOwner.scala index 87ddd7e..d4d528b 100644 --- a/src/main/scala/com/github/sstone/amqp/ChannelOwner.scala +++ b/src/main/scala/com/github/sstone/amqp/ChannelOwner.scala @@ -21,7 +21,7 @@ object ChannelOwner { case class NotConnectedError(request: Request) - def props(init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None): Props = Props(new ChannelOwner(init, channelParams)) + def props(init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None): Props = Props(new ChannelOwner(init, channelParams)) private[amqp] class Forwarder(channel: Channel) extends Actor with ActorLogging { @@ -153,11 +153,11 @@ object ChannelOwner { } } -class ChannelOwner(init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None) extends Actor with ActorLogging { +class ChannelOwner(init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None) extends Actor with ActorLogging { import ChannelOwner._ - var requestLog: Vector[RequestAndSender] = init.toVector + var requestLog: Vector[(Request, Option[ActorRef])] = init.map(_ -> None).toVector val statusListeners = mutable.HashSet.empty[ActorRef] override def preStart() = context.parent ! ConnectionOwner.CreateChannel diff --git a/src/main/scala/com/github/sstone/amqp/Consumer.scala b/src/main/scala/com/github/sstone/amqp/Consumer.scala index abb1ee5..f896f6b 100644 --- a/src/main/scala/com/github/sstone/amqp/Consumer.scala +++ b/src/main/scala/com/github/sstone/amqp/Consumer.scala @@ -9,12 +9,12 @@ import akka.event.LoggingReceive import scala.collection.JavaConversions._ object Consumer { - def props(listener: Option[ActorRef], autoack: Boolean = false, init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None, + def props(listener: Option[ActorRef], autoack: Boolean = false, init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None, consumerTag: String = "", noLocal: Boolean = false, exclusive: Boolean = false, arguments: Map[String, AnyRef] = Map.empty): Props = Props(new Consumer(listener, autoack, init, channelParams, consumerTag, noLocal, exclusive, arguments)) def props(listener: ActorRef, exchange: ExchangeParameters, queue: QueueParameters, routingKey: String, channelParams: Option[ChannelParameters], autoack: Boolean): Props = - props(Some(listener), init = List(AddBinding(Binding(exchange, queue, routingKey)) -> None), channelParams = channelParams, autoack = autoack) + props(Some(listener), init = List(AddBinding(Binding(exchange, queue, routingKey))), channelParams = channelParams, autoack = autoack) def props(listener: ActorRef, channelParams: Option[ChannelParameters], autoack: Boolean): Props = props(Some(listener), channelParams = channelParams, autoack = autoack) } @@ -32,7 +32,7 @@ object Consumer { */ class Consumer(listener: Option[ActorRef], autoack: Boolean = false, - init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], + init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None, consumerTag: String = "", noLocal: Boolean = false, diff --git a/src/main/scala/com/github/sstone/amqp/RpcServer.scala b/src/main/scala/com/github/sstone/amqp/RpcServer.scala index b713d4b..865ce30 100644 --- a/src/main/scala/com/github/sstone/amqp/RpcServer.scala +++ b/src/main/scala/com/github/sstone/amqp/RpcServer.scala @@ -39,14 +39,14 @@ object RpcServer { def onFailure(delivery: Delivery, e: Throwable): ProcessResult } - def props(processor: RpcServer.IProcessor, init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None)(implicit ctx: ExecutionContext): Props = + def props(processor: RpcServer.IProcessor, init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None)(implicit ctx: ExecutionContext): Props = Props(new RpcServer(processor, init, channelParams)) def props(queue: QueueParameters, exchange: ExchangeParameters, routingKey: String, proc: RpcServer.IProcessor, channelParams: ChannelParameters)(implicit ctx: ExecutionContext): Props = - props(processor = proc, init = List(AddBinding(Binding(exchange, queue, routingKey)) -> None), channelParams = Some(channelParams)) + props(processor = proc, init = List(AddBinding(Binding(exchange, queue, routingKey))), channelParams = Some(channelParams)) def props(queue: QueueParameters, exchange: ExchangeParameters, routingKey: String, proc: RpcServer.IProcessor)(implicit ctx: ExecutionContext): Props = - props(processor = proc, init = List(AddBinding(Binding(exchange, queue, routingKey)) -> None)) + props(processor = proc, init = List(AddBinding(Binding(exchange, queue, routingKey)))) } @@ -60,7 +60,7 @@ object RpcServer { * @param processor [[com.github.sstone.amqp.RpcServer.IProcessor]] implementation * @param channelParams optional channel parameters */ -class RpcServer(processor: RpcServer.IProcessor, init: Seq[RequestAndSender] = Seq.empty[RequestAndSender], channelParams: Option[ChannelParameters] = None)(implicit ctx: ExecutionContext = ExecutionContext.Implicits.global) extends Consumer(listener = None, autoack = false, init = init, channelParams = channelParams) { +class RpcServer(processor: RpcServer.IProcessor, init: Seq[Request] = Seq.empty[Request], channelParams: Option[ChannelParameters] = None)(implicit ctx: ExecutionContext = ExecutionContext.Implicits.global) extends Consumer(listener = None, autoack = false, init = init, channelParams = channelParams) { import RpcServer._ private def sendResponse(result: ProcessResult, properties: BasicProperties, channel: Channel) { diff --git a/src/main/scala/com/github/sstone/amqp/package.scala b/src/main/scala/com/github/sstone/amqp/package.scala deleted file mode 100644 index 58ea0de..0000000 --- a/src/main/scala/com/github/sstone/amqp/package.scala +++ /dev/null @@ -1,8 +0,0 @@ -package com.github.sstone - -import akka.actor.ActorRef -import com.github.sstone.amqp.Amqp.Request - -package object amqp { - type RequestAndSender = (Request, Option[ActorRef]) -} diff --git a/src/main/scala/com/github/sstone/amqp/samples/Consumer2.scala b/src/main/scala/com/github/sstone/amqp/samples/Consumer2.scala index 642b642..d9586e4 100644 --- a/src/main/scala/com/github/sstone/amqp/samples/Consumer2.scala +++ b/src/main/scala/com/github/sstone/amqp/samples/Consumer2.scala @@ -31,7 +31,7 @@ object Consumer2 extends App { // to the broker is lost and restored val consumer = ConnectionOwner.createChildActor(conn, Consumer.props( listener = Some(listener), - init = List(AddBinding(Binding(StandardExchanges.amqDirect, queueParams, "my_key")) -> None) + init = List(AddBinding(Binding(StandardExchanges.amqDirect, queueParams, "my_key"))) ), name = Some("consumer")) // wait till everyone is actually connected to the broker diff --git a/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala b/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala index 07410bc..553ca31 100644 --- a/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala +++ b/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala @@ -174,7 +174,7 @@ class ConsumerSpec extends ChannelSpec { val probe = TestProbe() val queue = randomQueue val consumer = ConnectionOwner.createChildActor(conn, Consumer.props( - listener = Some(probe.ref), autoack = false, init = Seq.empty[RequestAndSender], channelParams = None, + listener = Some(probe.ref), autoack = false, init = Seq.empty[Request], channelParams = None, consumerTag = "", noLocal = false, exclusive = true, arguments = Map.empty), timeout = 5000 millis) val producer = ConnectionOwner.createChildActor(conn, ChannelOwner.props()) @@ -192,7 +192,7 @@ class ConsumerSpec extends ChannelSpec { assert(delivery.consumerTag === consumerTag) val consumer1 = ConnectionOwner.createChildActor(conn, Consumer.props( - listener = Some(probe.ref), autoack = false, init = Seq.empty[RequestAndSender], channelParams = None, + listener = Some(probe.ref), autoack = false, init = Seq.empty[Request], channelParams = None, consumerTag = "", noLocal = false, exclusive = true, arguments = Map.empty), timeout = 5000 millis) consumer1 ! AddStatusListener(probe.ref) From b2b5a277f184461c5d6af9d9c325f2fbf51dc34d Mon Sep 17 00:00:00 2001 From: Ilya Brin Date: Wed, 31 Jan 2018 14:36:14 -0800 Subject: [PATCH 5/7] fix tests --- .../com/github/sstone/amqp/ChannelSpec.scala | 12 +++- .../com/github/sstone/amqp/ConsumerSpec.scala | 62 ++++++++++++------- 2 files changed, 47 insertions(+), 27 deletions(-) diff --git a/src/test/scala/com/github/sstone/amqp/ChannelSpec.scala b/src/test/scala/com/github/sstone/amqp/ChannelSpec.scala index f179d87..3c2016f 100644 --- a/src/test/scala/com/github/sstone/amqp/ChannelSpec.scala +++ b/src/test/scala/com/github/sstone/amqp/ChannelSpec.scala @@ -13,11 +13,11 @@ import com.rabbitmq.client.ConnectionFactory import com.github.sstone.amqp.Amqp._ import scala.util.Random -class ChannelSpec extends TestKit(ActorSystem("TestSystem")) with WordSpecLike with ShouldMatchers with BeforeAndAfter with ImplicitSender { +trait ChannelSpecNoTestKit extends WordSpecLike with ShouldMatchers with BeforeAndAfter { + implicit val system: ActorSystem + implicit val timeout = Timeout(5 seconds) val connFactory = new ConnectionFactory() - val uri = system.settings.config.getString("amqp-client-test.rabbitmq.uri") - connFactory.setUri(uri) var conn: ActorRef = _ var channelOwner: ActorRef = _ val random = new Random() @@ -32,6 +32,8 @@ class ChannelSpec extends TestKit(ActorSystem("TestSystem")) with WordSpecLike w before { println("before") + val uri = system.settings.config.getString("amqp-client-test.rabbitmq.uri") + connFactory.setUri(uri) conn = system.actorOf(ConnectionOwner.props(connFactory, 1 second)) channelOwner = ConnectionOwner.createChildActor(conn, ChannelOwner.props()) waitForConnection(system, conn, channelOwner).await(5, TimeUnit.SECONDS) @@ -42,3 +44,7 @@ class ChannelSpec extends TestKit(ActorSystem("TestSystem")) with WordSpecLike w Await.result(gracefulStop(conn, 5 seconds), 6 seconds) } } + +class ChannelSpec extends TestKit(ActorSystem("TestSystem")) with ChannelSpecNoTestKit with ImplicitSender { + +} \ No newline at end of file diff --git a/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala b/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala index 553ca31..51ebfbb 100644 --- a/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala +++ b/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala @@ -1,21 +1,25 @@ package com.github.sstone.amqp +import akka.actor.ActorSystem import akka.testkit.TestProbe import com.github.sstone.amqp.Amqp._ + import concurrent.duration._ import org.junit.runner.RunWith import org.scalatest.junit.JUnitRunner @RunWith(classOf[JUnitRunner]) -class ConsumerSpec extends ChannelSpec { +class ConsumerSpec extends ChannelSpecNoTestKit { + override implicit val system: ActorSystem = ActorSystem("ConsumerSpec") "Consumers" should { "receive messages sent by producers" in { val exchange = ExchangeParameters(name = "amq.direct", exchangeType = "", passive = true) val queue = QueueParameters(name = "", passive = false, exclusive = true) - ignoreMsg { + val probe = TestProbe() + implicit val sender = probe.ref + probe.ignoreMsg { case Amqp.Ok(p:Publish, _) => true } - val probe = TestProbe() val consumer = ConnectionOwner.createChildActor(conn, Consumer.props(listener = Some(probe.ref)), timeout = 5000 millis) val producer = ConnectionOwner.createChildActor(conn, ChannelOwner.props()) consumer ! AddStatusListener(probe.ref) @@ -23,7 +27,7 @@ class ConsumerSpec extends ChannelSpec { probe.expectMsg(1 second, ChannelOwner.Connected) probe.expectMsg(1 second, ChannelOwner.Connected) consumer ! AddBinding(Binding(exchange, queue, "my_key")) - val check = receiveOne(1 second) + val check = probe.receiveOne(1 second) println(check) val message = "yo!".getBytes producer ! Publish(exchange.name, "my_key", message) @@ -32,12 +36,16 @@ class ConsumerSpec extends ChannelSpec { "be able to set their channel's prefetch size" in { val queue = randomQueue val probe = TestProbe() + implicit val sender = probe.ref + probe.ignoreMsg { + case Amqp.Ok(p:Publish, _) => true + } val consumer = ConnectionOwner.createChildActor(conn, Consumer.props(listener = probe.ref, autoack = false, channelParams = Some(ChannelParameters(qos = 3))), timeout = 5000 millis) consumer ! AddStatusListener(probe.ref) probe.expectMsg(1 second, ChannelOwner.Connected) consumer ! AddQueue(queue) - val Amqp.Ok(AddQueue(_), _) = receiveOne(1 second) + val Amqp.Ok(AddQueue(_), _) = probe.receiveOne(1 second) consumer ! Publish("", queue.name, "test".getBytes("UTF-8")) val delivery1 = probe.expectMsgClass(200 milliseconds, classOf[Delivery]) @@ -52,14 +60,15 @@ class ConsumerSpec extends ChannelSpec { // but if we ack one our our messages we shoule get the 4th delivery consumer ! Ack(deliveryTag = delivery1.envelope.getDeliveryTag) - val Amqp.Ok(Ack(_), _) = receiveOne(1 second) + val Amqp.Ok(Ack(_), _) = probe.receiveOne(1 second) val delivery4 = probe.expectMsgClass(200 milliseconds, classOf[Delivery]) } "be restarted if their channel crashes" in { val exchange = ExchangeParameters(name = "amq.direct", exchangeType = "", passive = true) val queue = randomQueue val probe = TestProbe() - ignoreMsg { + implicit val sender = probe.ref + probe.ignoreMsg { case Amqp.Ok(p:Publish, _) => true } val consumer = ConnectionOwner.createChildActor(conn, Consumer.props(listener = Some(probe.ref)), timeout = 5000 millis) @@ -69,7 +78,7 @@ class ConsumerSpec extends ChannelSpec { probe.expectMsg(1 second, ChannelOwner.Connected) probe.expectMsg(1 second, ChannelOwner.Connected) consumer ! Record(AddBinding(Binding(exchange, queue, "my_key"))) - val Amqp.Ok(AddBinding(_), _) = receiveOne(1 second) + val Amqp.Ok(AddBinding(_), _) = probe.receiveOne(1 second) val message = "yo!".getBytes producer ! Publish(exchange.name, "my_key", message) @@ -77,8 +86,9 @@ class ConsumerSpec extends ChannelSpec { // crash the consumer's channel consumer ! DeclareExchange(ExchangeParameters(name = "foo", passive = true, exchangeType ="")) - receiveOne(1 second) + val Amqp.Error(DeclareExchange(_), _) = probe.receiveOne(1 second) probe.expectMsgAllOf(1 second, ChannelOwner.Disconnected, ChannelOwner.Connected) + val Ok(AddBinding(Binding(`exchange`, `queue`, "my_key")), Some(_)) = probe.receiveOne(1 second) Thread.sleep(100) producer ! Publish(exchange.name, "my_key", message) @@ -89,7 +99,8 @@ class ConsumerSpec extends ChannelSpec { val exchange = ExchangeParameters(name = randomExchangeName, exchangeType = "direct", passive = false, durable = false, autodelete = true) val queue = randomQueue val probe = TestProbe() - ignoreMsg { + implicit val sender = probe.ref + probe.ignoreMsg { case Amqp.Ok(p:Publish, _) => true } val consumer = ConnectionOwner.createChildActor(conn, Consumer.props(listener = Some(probe.ref)), timeout = 5000 millis) @@ -99,12 +110,12 @@ class ConsumerSpec extends ChannelSpec { probe.expectMsg(1 second, ChannelOwner.Connected) probe.expectMsg(1 second, ChannelOwner.Connected) consumer ! AddBinding(Binding(exchange, queue, "test_key")) - val Amqp.Ok(AddBinding(_), _) = receiveOne(1 second) + val Amqp.Ok(AddBinding(_), _) = probe.receiveOne(1 second) // check that our exchange was created val exchange1 = exchange.copy(passive = true) consumer ! DeclareExchange(exchange1) - val Amqp.Ok(DeclareExchange(_), _) = receiveOne(1 second) + val Amqp.Ok(DeclareExchange(_), _) = probe.receiveOne(1 second) // check that publishing works producer ! Publish(exchange.name, "test_key", "test message".getBytes("UTF-8")) @@ -114,7 +125,8 @@ class ConsumerSpec extends ChannelSpec { val queue1 = randomQueue val queue2 = randomQueue val probe = TestProbe() - ignoreMsg { + implicit val sender = probe.ref + probe.ignoreMsg { case Amqp.Ok(p:Publish, _) => true } val consumer = ConnectionOwner.createChildActor(conn, Consumer.props(listener = Some(probe.ref), autoack = false), timeout = 5000 millis) @@ -125,9 +137,9 @@ class ConsumerSpec extends ChannelSpec { probe.expectMsg(1 second, ChannelOwner.Connected) consumer ! AddQueue(queue1) - val Amqp.Ok(AddQueue(_), Some(consumerTag1: String)) = receiveOne(1 second) + val Amqp.Ok(AddQueue(_), Some(consumerTag1: String)) = probe.receiveOne(1 second) consumer ! AddQueue(queue2) - val Amqp.Ok(AddQueue(_), Some(consumerTag2: String)) = receiveOne(1 second) + val Amqp.Ok(AddQueue(_), Some(consumerTag2: String)) = probe.receiveOne(1 second) producer ! Publish("", queue1.name, "test1".getBytes("UTF-8")) val delivery1: Delivery = probe.expectMsgClass(classOf[Delivery]) @@ -138,16 +150,17 @@ class ConsumerSpec extends ChannelSpec { assert(delivery2.consumerTag === consumerTag2) consumer ! CancelConsumer(consumerTag1) - val Amqp.Ok(CancelConsumer(_), _) = receiveOne(1 second) + val Amqp.Ok(CancelConsumer(_), _) = probe.receiveOne(1 second) producer ! Publish("", queue1.name, "test1".getBytes("UTF-8")) probe.expectNoMsg() } "send consumer cancellation notifications" in { - ignoreMsg { + val probe = TestProbe() + implicit val sender = probe.ref + probe.ignoreMsg { case Amqp.Ok(p:Publish, _) => true } - val probe = TestProbe() val queue = randomQueue val consumer = ConnectionOwner.createChildActor(conn, Consumer.props(listener = Some(probe.ref), autoack = false), timeout = 5000 millis) val producer = ConnectionOwner.createChildActor(conn, ChannelOwner.props()) @@ -157,21 +170,22 @@ class ConsumerSpec extends ChannelSpec { probe.expectMsg(1 second, ChannelOwner.Connected) consumer ! AddQueue(queue) - val Amqp.Ok(AddQueue(_), Some(consumerTag: String)) = receiveOne(1 second) + val Amqp.Ok(AddQueue(_), Some(consumerTag: String)) = probe.receiveOne(1 second) producer ! Publish("", queue.name, "test".getBytes("UTF-8")) val delivery: Delivery = probe.expectMsgClass(classOf[Delivery]) assert(delivery.consumerTag === consumerTag) producer ! DeleteQueue(queue.name) - val Ok(DeleteQueue(_, _, _), result) = receiveOne(1 second) + val Ok(DeleteQueue(_, _, _), result) = probe.receiveOne(1 second) probe.expectMsg(1 second, ConsumerCancelled(consumerTag)) } "create exclusive consumers" in { - ignoreMsg { + val probe = TestProbe() + implicit val sender = probe.ref + probe.ignoreMsg { case Amqp.Ok(p:Publish, _) => true } - val probe = TestProbe() val queue = randomQueue val consumer = ConnectionOwner.createChildActor(conn, Consumer.props( listener = Some(probe.ref), autoack = false, init = Seq.empty[Request], channelParams = None, @@ -185,7 +199,7 @@ class ConsumerSpec extends ChannelSpec { probe.expectMsg(1 second, ChannelOwner.Connected) consumer ! AddQueue(queue) - val Amqp.Ok(AddQueue(_), Some(consumerTag: String)) = receiveOne(1 second) + val Amqp.Ok(AddQueue(_), Some(consumerTag: String)) = probe.receiveOne(1 second) producer ! Publish("", queue.name, "test".getBytes("UTF-8")) val delivery: Delivery = probe.expectMsgClass(classOf[Delivery]) @@ -200,7 +214,7 @@ class ConsumerSpec extends ChannelSpec { // you cannot have more than 1 exclusive consumer on the same queue consumer1 ! AddQueue(queue) - val Amqp.Error(_, reason) = receiveOne(1 second) + val Amqp.Error(_, reason) = probe.receiveOne(1 second) } } } From ff9870d1e61acb999c706bdfcf210a2e17940968 Mon Sep 17 00:00:00 2001 From: Ilya Brin Date: Wed, 31 Jan 2018 14:56:44 -0800 Subject: [PATCH 6/7] add a test --- .../com/github/sstone/amqp/ConsumerSpec.scala | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala b/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala index 51ebfbb..c71a35d 100644 --- a/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala +++ b/src/test/scala/com/github/sstone/amqp/ConsumerSpec.scala @@ -1,13 +1,16 @@ package com.github.sstone.amqp import akka.actor.ActorSystem +import akka.pattern.{ask, gracefulStop} import akka.testkit.TestProbe +import akka.util.Timeout import com.github.sstone.amqp.Amqp._ - -import concurrent.duration._ import org.junit.runner.RunWith import org.scalatest.junit.JUnitRunner +import scala.concurrent.Await +import scala.concurrent.duration._ + @RunWith(classOf[JUnitRunner]) class ConsumerSpec extends ChannelSpecNoTestKit { override implicit val system: ActorSystem = ActorSystem("ConsumerSpec") @@ -216,5 +219,15 @@ class ConsumerSpec extends ChannelSpecNoTestKit { consumer1 ! AddQueue(queue) val Amqp.Error(_, reason) = probe.receiveOne(1 second) } + "save sender for requests while disconnected" in { + val declareExchange = DeclareExchange(ExchangeParameters(name = "amq.direct", passive = true, exchangeType = "")) + val conn = system.actorOf(ConnectionOwner.props(connFactory, 1 second)) + try { + val channelOwner = ConnectionOwner.createChildActor(conn, ChannelOwner.props()) + val Ok(`declareExchange`, Some(_)) = Await.result(channelOwner.ask(Record(declareExchange))(Timeout(1 second)), 1 second) + } finally { + Await.result(gracefulStop(conn, 5 seconds), 6 seconds) + } + } } } From ae267e8e9474e2c4d80f33c99430184a709146a0 Mon Sep 17 00:00:00 2001 From: Ilya Brin Date: Wed, 31 Jan 2018 14:57:47 -0800 Subject: [PATCH 7/7] newline at end --- src/test/scala/com/github/sstone/amqp/ChannelSpec.scala | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/src/test/scala/com/github/sstone/amqp/ChannelSpec.scala b/src/test/scala/com/github/sstone/amqp/ChannelSpec.scala index 3c2016f..14afee4 100644 --- a/src/test/scala/com/github/sstone/amqp/ChannelSpec.scala +++ b/src/test/scala/com/github/sstone/amqp/ChannelSpec.scala @@ -45,6 +45,4 @@ trait ChannelSpecNoTestKit extends WordSpecLike with ShouldMatchers with BeforeA } } -class ChannelSpec extends TestKit(ActorSystem("TestSystem")) with ChannelSpecNoTestKit with ImplicitSender { - -} \ No newline at end of file +class ChannelSpec extends TestKit(ActorSystem("TestSystem")) with ChannelSpecNoTestKit with ImplicitSender