diff --git a/README.md b/README.md index 67b809dc..279b8239 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/backendFs2Rabbit/src/main/scala/com/itv/bucky/backend/fs2rabbit/Fs2RabbitAmqpClient.scala b/backendFs2Rabbit/src/main/scala/com/itv/bucky/backend/fs2rabbit/Fs2RabbitAmqpClient.scala index b6321e12..2042ab75 100644 --- a/backendFs2Rabbit/src/main/scala/com/itv/bucky/backend/fs2rabbit/Fs2RabbitAmqpClient.scala +++ b/backendFs2Rabbit/src/main/scala/com/itv/bucky/backend/fs2rabbit/Fs2RabbitAmqpClient.scala @@ -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._ @@ -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) @@ -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], @@ -181,21 +197,32 @@ 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)) @@ -203,15 +230,29 @@ class Fs2RabbitAmqpClient[F[_]: Async: Temporal]( } .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) @@ -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] = diff --git a/backendFs2Rabbit/src/test/scala/fs2rabbit/ConsumerConnectionRecoverySpec.scala b/backendFs2Rabbit/src/test/scala/fs2rabbit/ConsumerConnectionRecoverySpec.scala new file mode 100644 index 00000000..44a418ec --- /dev/null +++ b/backendFs2Rabbit/src/test/scala/fs2rabbit/ConsumerConnectionRecoverySpec.scala @@ -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))) + } + } + } + } +}