Skip to content
Merged
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
30 changes: 30 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -210,6 +210,36 @@ object Example extends IOApp {

Both backends implement the same `AmqpClient[F]` trait, so switching between them should only require changing the client creation and imports.

# Migrating to 4.0.3 and above — connection recovery fix (fs2-rabbit backend)

Versions **4.0.0-M1 through 4.0.2** contain a connection-recovery bug that affects **`Fs2RabbitAmqpClient` only**. The `JavaBackendAmqpClient` is not affected.

## What went wrong

The fs2-rabbit backend ran each consumer as an fs2 `Stream` inside a `.background` fiber, but silently ignored the fiber outcome. When a network blip caused the underlying AMQP channel to close, the stream terminated with an error that was swallowed — no restart was ever attempted. After the Java AMQP client automatically recovered the connection, messages were delivered to the channel's internal buffer but nobody was draining it, so the consumer appeared permanently stuck.

The v3 Java backend (and `JavaBackendAmqpClient` in v4) use a `DefaultConsumer` callback on an `AutorecoveringChannel`, which the AMQP Java client automatically re-registers after recovery. The fs2-rabbit backend had no equivalent mechanism.

A secondary issue in the same code path: handler exceptions propagated through the stream and also killed it, rather than being resolved via `exceptionalAction`.

## Fix in 4.0.3

- The consumer loop now wraps each run in a retry: on any failure it logs a warning, waits for `networkRecoveryInterval` (default 3 seconds, matching the Java AMQP client's recovery interval), then creates a **fresh channel and consumer** before resuming. A fresh channel is required because the auto-recovered channel's dead internal fs2-rabbit queue would otherwise accumulate un-acked messages indefinitely.
- Handler exceptions are now caught and resolved via `exceptionalAction` rather than terminating the stream.

## Action required

If you are using **`Fs2RabbitAmqpClient`** on any 4.x release prior to 4.0.3, upgrade to **4.0.3**:

```scala
// For fs2-rabbit Backend
libraryDependencies += "com.itv" %% "bucky-backend-fs2-rabbit" % "4.0.3"
```

No code changes are required — the fix is entirely internal to `registerConsumer`.

If you are using **`JavaBackendAmqpClient`** you are not affected and no action is needed.

# Releasing a new version

1. Merge your change into master
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import com.itv.bucky.{
publish
}
import com.rabbitmq.client.LongString
import com.typesafe.scalalogging.StrictLogging
import dev.profunktor.fs2rabbit.arguments.SafeArg
import dev.profunktor.fs2rabbit.config.Fs2RabbitConfig
import dev.profunktor.fs2rabbit.config.declaration._
Expand All @@ -48,21 +49,24 @@ import dev.profunktor.fs2rabbit.model.AmqpFieldValue.{
TimestampVal
}
import dev.profunktor.fs2rabbit.model.{AMQPChannel, HeaderKey, Headers, PublishingFlag, ShortString}
import fs2.Stream
import scodec.bits.ByteVector

import java.util.{Date, UUID}
import scala.concurrent.duration.FiniteDuration
import scala.concurrent.duration._
import scala.jdk.CollectionConverters._
import scala.language.higherKinds
import Fs2RabbitAmqpClient._
import cats.effect.kernel.Temporal

class Fs2RabbitAmqpClient[F[_]: Async: Temporal](
config: AmqpClientConfig,
client: RabbitClient[F],
connection: model.AMQPConnection,
publishChannel: model.AMQPChannel,
amqpClientConnectionManager: AmqpClientConnectionManager[F]
) extends AmqpClient[F] {
) extends AmqpClient[F]
with StrictLogging {

override def declare(declarations: decl.Declaration*): F[Unit] = declare(declarations.toList)

Expand Down Expand Up @@ -173,6 +177,18 @@ class Fs2RabbitAmqpClient[F[_]: Async: Temporal](
_ <- if (ended) Async[F].unit else Temporal[F].sleep(sleep) *> repeatUntil(eval)(pred)(sleep)
} yield ()

/** Creates a fresh channel and registers a consumer on it. Extracted as a
* protected method so tests can override it with a controlled stream without
* needing a real RabbitMQ connection.
*/
protected def acquireConsumerStream(
queueName: bucky.QueueName
): Resource[F, (model.AckResult => F[Unit], Stream[F, model.AmqpEnvelope[consume.Delivery]])] =
client.createChannel(connection).evalMap { implicit channel =>
implicit val decoder: EnvelopeDecoder[F, consume.Delivery] = deliveryDecoder(queueName)
client.createAckerConsumer[consume.Delivery](model.QueueName(queueName.value))
}

override def registerConsumer(
queueName: bucky.QueueName,
handler: Handler[F, consume.Delivery],
Expand All @@ -181,37 +197,62 @@ class Fs2RabbitAmqpClient[F[_]: Async: Temporal](
shutdownTimeout: FiniteDuration,
shutdownRetry: FiniteDuration
): Resource[F, Unit] =
client.createChannel(connection).flatMap { implicit channel =>
implicit val decoder: EnvelopeDecoder[F, consume.Delivery] = deliveryDecoder(queueName)
Resource.eval(Ref.of[F, Set[UUID]](Set.empty)).flatMap { consumptionIds =>
Resource.eval(client.createAckerConsumer[consume.Delivery](model.QueueName(queueName.value))).flatMap { case (acker, consumer) =>
Resource.eval(Ref.of[F, Set[UUID]](Set.empty)).flatMap { consumptionIds =>

// Create a fresh channel and consumer for each run. This ensures that after
// an auto-recovery the old (dead) fs2-rabbit stream is discarded and a new
// one is registered on the recovered channel rather than relying on the
// channel's internal queue which nobody is draining.
def runConsumer: F[Unit] =
acquireConsumerStream(queueName).use { case (acker, consumer) =>
consumer
.evalMap(delivery =>
for {
uuid <- Async[F].delay(UUID.randomUUID())
_ <- consumptionIds.update(set => set + uuid)
res <- handler(delivery.payload).attempt
tag = delivery.deliveryTag
_ <- consumptionIds.update(set => set - uuid)
result <- Async[F].fromEither(res)
} yield (result, tag)
)
.evalMap { delivery =>
val tag = delivery.deliveryTag

Async[F]
.bracket {
Async[F].delay(UUID.randomUUID()).flatTap(uuid => consumptionIds.update(_ + uuid))
} { uuid =>
handler(delivery.payload).attempt.flatMap {
case Right(action) => Async[F].pure((action, tag))
case Left(e) =>
Async[F].delay(logger.error(s"Handler exception for queue ${queueName.value}: ${e.getMessage}", e)) *>
Async[F].pure((exceptionalAction, tag))
}
} { uuid =>
consumptionIds.update(_ - uuid)
}
}
.evalMap {
case (consume.Ack, tag) => acker(model.AckResult.Ack(tag))
case (consume.DeadLetter, tag) => acker(model.AckResult.NAck(tag))
case (consume.RequeueImmediately, tag) => acker(model.AckResult.Reject(tag))
}
.compile
.drain
.background
.flatMap { _ =>
Resource.onFinalize(
repeatUntil(consumptionIds.get)(_.isEmpty)(shutdownRetry).timeout(shutdownTimeout)
)
}
.map(_ => ())
}
}

// The delay between retry attempts. Aligned with the Java AMQP client's
// automatic-recovery interval so we wait long enough for the connection
// to be re-established before trying to open a new channel.
val recoveryDelay: FiniteDuration = config.networkRecoveryInterval.getOrElse(3.seconds)

// Retry the consumer indefinitely on failure. This mirrors the
// AutorecoveringChannel behaviour of the v3 Java backend: when a network
// blip closes the channel, the consumer is re-registered automatically
// after the connection is recovered.
def consumerWithRecovery: F[Unit] =
runConsumer.handleErrorWith { error =>
Async[F].delay(logger.warn(s"Consumer for queue ${queueName.value} failed, will retry after $recoveryDelay: ${error.getMessage}")) *>
Temporal[F].sleep(recoveryDelay) *>
consumerWithRecovery
}

consumerWithRecovery.background.flatMap { _ =>
Resource.onFinalize(
repeatUntil(consumptionIds.get)(_.isEmpty)(shutdownRetry).timeout(shutdownTimeout)
)
}.map(_ => ())
}

override def isConnectionOpen: F[Boolean] = Async[F].pure(connection.value.isOpen)
Expand Down Expand Up @@ -246,7 +287,7 @@ object Fs2RabbitAmqpClient {
amqpChannel = publishChannel
)
)
} yield new Fs2RabbitAmqpClient(client, connection, publishChannel, amqpClientConnectionManager)
} yield new Fs2RabbitAmqpClient(config, client, connection, publishChannel, amqpClientConnectionManager)
}

implicit def deliveryEncoder[F[_]: Async]: MessageEncoder[F, PublishCommand] =
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,167 @@
package fs2rabbit

import cats.effect.{Deferred, IO, Ref, Resource}
import cats.effect.testing.scalatest.AsyncIOSpec
import cats.implicits._
import com.itv.bucky.{AmqpClientConfig, Envelope, ExchangeName, Payload, QueueName, RoutingKey, consume, publish}
import com.itv.bucky.backend.fs2rabbit.{AmqpClientConnectionManager, Fs2RabbitAmqpClient}
import dev.profunktor.fs2rabbit.interpreter.RabbitClient
import dev.profunktor.fs2rabbit.model
import dev.profunktor.fs2rabbit.model.{AmqpEnvelope, AmqpProperties, DeliveryTag}
import fs2.Stream
import org.scalatest.matchers.should.Matchers
import org.scalatest.wordspec.AsyncWordSpec

import java.io.IOException
import scala.concurrent.duration._

/** Regression tests for the connection-recovery bug introduced in v4.
*
* v3 used `DefaultConsumer` on an `AutorecoveringChannel`, which the Java AMQP
* client transparently re-registers after a network blip. v4's fs2-rabbit
* backend ran the consumer as an fs2 Stream in a `.background` fiber and ignored
* the fiber outcome. When the stream terminated due to a connection drop the
* error was silently swallowed and no restart was attempted.
*
* These tests exercise the actual `Fs2RabbitAmqpClient.registerConsumer` method
* by subclassing it and overriding `acquireConsumerStream` with controlled
* in-memory streams — no real RabbitMQ connection is required.
*/
class ConsumerConnectionRecoverySpec extends AsyncWordSpec with AsyncIOSpec with Matchers {

/** Build a minimal `AmqpEnvelope[consume.Delivery]` suitable for use in tests. */
private def makeEnvelope(tag: Long): model.AmqpEnvelope[consume.Delivery] =
AmqpEnvelope(
DeliveryTag(tag),
consume.Delivery(
Payload("test".getBytes),
consume.ConsumerTag("ctag"),
Envelope(tag, redeliver = false, ExchangeName("ex"), RoutingKey("rk")),
publish.MessageProperties.minimalBasic
),
AmqpProperties.empty,
model.ExchangeName("ex"),
model.RoutingKey("rk"),
redelivered = false
)

/** Config with a short recovery interval so tests don't wait 3 seconds. */
private val testConfig: AmqpClientConfig =
AmqpClientConfig("localhost", 5672, "guest", "guest", networkRecoveryInterval = Some(50.millis))

/** Create a test double for `Fs2RabbitAmqpClient` that serves controlled
* `(acker, stream)` pairs from `streamsRef` in order. No real AMQP
* connection is needed: the overridden `acquireConsumerStream` never calls
* `client` or `connection`.
*/
private def makeTestClient(
streamsRef: Ref[IO, List[(model.AckResult => IO[Unit], Stream[IO, model.AmqpEnvelope[consume.Delivery]])]]
): Fs2RabbitAmqpClient[IO] =
new Fs2RabbitAmqpClient[IO](
testConfig,
null.asInstanceOf[RabbitClient[IO]],
null.asInstanceOf[model.AMQPConnection],
null.asInstanceOf[model.AMQPChannel],
null.asInstanceOf[AmqpClientConnectionManager[IO]]
) {
override protected def acquireConsumerStream(queueName: QueueName) =
Resource.eval(streamsRef.modify {
case head :: tail => (tail, head)
case Nil => sys.error("No more test streams available")
})
}

"Fs2RabbitAmqpClient.registerConsumer" when {

"the consumer stream fails (simulated connection drop)" should {

"retry and resume processing messages after recovery" in {
// stream1 emits one message then fails, simulating a dropped connection
val stream1 = Stream.emit(makeEnvelope(1L)) ++ Stream.raiseError[IO](new IOException("Connection reset by peer"))
val acker1 = (_: model.AckResult) => IO.unit
// stream2 emits one more message then completes normally
val stream2 = Stream.emit(makeEnvelope(2L))
val acker2 = (_: model.AckResult) => IO.unit

for {
processedTags <- IO.ref(List.empty[Long])
done <- IO.deferred[Unit]
streamsRef <- IO.ref(List((acker1, stream1), (acker2, stream2)))
client = makeTestClient(streamsRef)
handler = (delivery: consume.Delivery) =>
processedTags.update(_ :+ delivery.envelope.deliveryTag) *>
processedTags.get.flatMap(tags => if (tags.size >= 2) done.complete(()).void else IO.unit) *>
IO.pure(consume.Ack)
_ <- client
.registerConsumer(QueueName("test"), handler, consume.DeadLetter, 10, 500.millis, 100.millis)
.use(_ => done.get.timeout(5.seconds))
tags <- processedTags.get
} yield {
tags should contain(1L)
tags should contain(2L)
}
}
}

"a handler throws an exception" should {

"use exceptionalAction rather than killing the consumer stream" in {
for {
ackerResults <- IO.ref(List.empty[model.AckResult])
done <- IO.deferred[Unit]
acker = (result: model.AckResult) =>
ackerResults.update(_ :+ result) *>
ackerResults.get.flatMap(rs => if (rs.size >= 2) done.complete(()).void else IO.unit)
stream = Stream.emits(List(makeEnvelope(1L), makeEnvelope(2L))).covary[IO]
streamsRef <- IO.ref(List((acker, stream)))
client = makeTestClient(streamsRef)
// Handler throws for the first message; succeeds (Ack) for the second
handler = (delivery: consume.Delivery) =>
if (delivery.envelope.deliveryTag == 1L)
IO.raiseError[consume.ConsumeAction](new RuntimeException("handler explosion"))
else
IO.pure(consume.Ack)
_ <- client
.registerConsumer(QueueName("test"), handler, consume.DeadLetter, 10, 500.millis, 100.millis)
.use(_ => done.get.timeout(5.seconds))
results <- ackerResults.get
} yield {
// tag 1: handler threw → exceptionalAction (DeadLetter) → NAck
results should contain(model.AckResult.NAck(DeliveryTag(1L)))
// tag 2: handler returned Ack → Ack
results should contain(model.AckResult.Ack(DeliveryTag(2L)))
}
}
}

"messages are processed normally" should {

"route Ack / DeadLetter / RequeueImmediately to the correct AckResult" in {
for {
ackerResults <- IO.ref(List.empty[model.AckResult])
done <- IO.deferred[Unit]
acker = (result: model.AckResult) =>
ackerResults.update(_ :+ result) *>
ackerResults.get.flatMap(rs => if (rs.size >= 3) done.complete(()).void else IO.unit)
stream = Stream.emits(List(makeEnvelope(1L), makeEnvelope(2L), makeEnvelope(3L))).covary[IO]
streamsRef <- IO.ref(List((acker, stream)))
client = makeTestClient(streamsRef)
handler = (delivery: consume.Delivery) =>
delivery.envelope.deliveryTag match {
case 1L => IO.pure(consume.Ack)
case 2L => IO.pure(consume.DeadLetter)
case _ => IO.pure(consume.RequeueImmediately)
}
_ <- client
.registerConsumer(QueueName("test"), handler, consume.DeadLetter, 10, 500.millis, 100.millis)
.use(_ => done.get.timeout(5.seconds))
results <- ackerResults.get
} yield {
results should contain(model.AckResult.Ack(DeliveryTag(1L)))
results should contain(model.AckResult.NAck(DeliveryTag(2L)))
results should contain(model.AckResult.Reject(DeliveryTag(3L)))
}
}
}
}
}