Sitelet https://github.com/profunktor/fs2-rabbit/commit/dcd8644f56d4404fb71e1421c0f99394796d5c6c
Skip to content

Commit dcd8644

Browse files
simpadjosimpadjo
andauthored
Terminate the stream when the server closes the connection (#737)
Co-authored-by: simpadjo <simpadjo@gmai.com>
1 parent 6e43d99 commit dcd8644

5 files changed

Lines changed: 37 additions & 1 deletion

File tree

‎core/src/main/scala/dev/profunktor/fs2rabbit/algebra/Consume.scala‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@ import cats.effect.std.Dispatcher
2121
import cats.syntax.flatMap._
2222
import cats.syntax.functor._
2323
import cats.{Applicative, Functor}
24-
import com.rabbitmq.client.{AMQP, Consumer, DefaultConsumer, Envelope}
24+
import com.rabbitmq.client.{AMQP, Consumer, DefaultConsumer, Envelope, ShutdownSignalException}
2525
import dev.profunktor.fs2rabbit.arguments.{Arguments, _}
2626
import dev.profunktor.fs2rabbit.model._
2727

@@ -106,6 +106,11 @@ object Consume {
106106
}
107107
}
108108
}
109+
110+
override def handleShutdownSignal(consumerTag: String, sig: ShutdownSignalException): Unit =
111+
if (!sig.isInitiatedByApplication) {
112+
internals.queue.foreach(q => dispatcher.unsafeRunAndForget(q.offer(Left(sig))))
113+
}
109114
}
110115
}
111116

‎docker-compose.yml‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,3 +5,5 @@ RabbitMQ:
55
- "5672:5672"
66
environment:
77
- DEBUG=false
8+
volumes:
9+
- ./rabbit-test-config/:/etc/rabbitmq/

‎rabbit-test-config/advanced.config‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
[
2+
{rabbit, [
3+
{channel_tick_interval, 500},
4+
{loopback_users, []}
5+
]}
6+
].

‎rabbit-test-config/rabbitmq.conf‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
consumer_timeout = 1000

‎tests/src/test/scala/dev/profunktor/fs2rabbit/interpreter/Fs2RabbitSpec.scala‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -711,6 +711,28 @@ trait Fs2RabbitSpec { self: BaseSpec =>
711711
}
712712
}
713713

714+
it should "shutdown the stream when the server closes the channel" in withRabbit { interpreter =>
715+
import interpreter._
716+
val msg = "will-not-be-acked"
717+
createConnectionChannel.use { implicit channel =>
718+
for {
719+
qxrk <- randomQueueData
720+
(q, x, rk) = qxrk
721+
_ <- declareExchange(x, ExchangeType.Topic)
722+
_ <- declareQueue(DeclarationQueueConfig.default(q))
723+
_ <- bindQueue(q, x, rk, QueueBindingArgs(Map.empty))
724+
publisher <- createPublisher[String](x, rk)
725+
_ <- publisher(msg)
726+
stream <- createAckerConsumer(q).map(_._2)
727+
results <- stream.attempt.compile.toList.timeoutAndForget(Duration(10, "s"))
728+
} yield {
729+
results.size shouldEqual 2
730+
results.head.map(_.payload) shouldEqual Right(msg)
731+
results.last.isLeft shouldEqual true
732+
}
733+
}
734+
}
735+
714736
it should "preserve order of published messages" in withRabbit { interpreter =>
715737
import dev.profunktor.fs2rabbit.effects.{EnvelopeDecoder, MessageEncoder}
716738
import interpreter._

0 commit comments

Comments
 (0)