From a520565d9e2358f0bab4ad75e368407b6774b543 Mon Sep 17 00:00:00 2001 From: huntc Date: Fri, 30 Aug 2019 10:07:07 +1000 Subject: [PATCH 1/2] Document and improve PubAck handling This improves the doc around how to back-pressure the publication of Publish commands to avoid buffer overflow in QoS 1+ cases. Implements the ask pattern for this purpose. --- .../mqtt/streaming/javadsl/MqttSession.scala | 22 +++++++++++++ .../mqtt/streaming/scaladsl/MqttSession.scala | 31 +++++++++++++++++++ .../test/java/docs/javadsl/MqttFlowTest.java | 20 +++++++----- .../scala/docs/scaladsl/MqttFlowSpec.scala | 9 ++++-- 4 files changed, 71 insertions(+), 11 deletions(-) diff --git a/mqtt-streaming/src/main/scala/akka/stream/alpakka/mqtt/streaming/javadsl/MqttSession.scala b/mqtt-streaming/src/main/scala/akka/stream/alpakka/mqtt/streaming/javadsl/MqttSession.scala index b2409c44f3..f8be4c3cb0 100644 --- a/mqtt-streaming/src/main/scala/akka/stream/alpakka/mqtt/streaming/javadsl/MqttSession.scala +++ b/mqtt-streaming/src/main/scala/akka/stream/alpakka/mqtt/streaming/javadsl/MqttSession.scala @@ -5,6 +5,8 @@ package akka.stream.alpakka.mqtt.streaming package javadsl +import java.util.concurrent.CompletionStage + import akka.NotUsed import akka.actor.ActorSystem import akka.stream.Materializer @@ -16,6 +18,8 @@ import akka.stream.alpakka.mqtt.streaming.scaladsl.{ } import akka.stream.javadsl.Source +import scala.compat.java8.FutureConverters._ + /** * Represents MQTT session state for both clients or servers. Session * state can survive across connections i.e. their lifetime is @@ -32,6 +36,18 @@ abstract class MqttSession { */ def tell[A](cp: Command[A]): Unit + /** + * Ask the session to perform a command regardless of the state it is + * in. This is important for sending Publish messages in particular, + * as a connection may not have been established with a session. + * @param cp The command to perform + * @tparam A The type of any carry for the command. + * @return A future indicating when the command has completed. Completion + * is defined as when it has been acknowledged by the recipient + * endpoint. + */ + def ask[A](cp: Command[A]): CompletionStage[A] + /** * Shutdown the session gracefully */ @@ -47,6 +63,9 @@ abstract class MqttClientSession extends MqttSession { override def tell[A](cp: Command[A]): Unit = underlying ! cp + override def ask[A](cp: Command[A]): CompletionStage[A] = + (underlying ? cp).toJava + override def shutdown(): Unit = underlying.shutdown() } @@ -91,6 +110,9 @@ abstract class MqttServerSession extends MqttSession { override def tell[A](cp: Command[A]): Unit = underlying ! cp + override def ask[A](cp: Command[A]): CompletionStage[A] = + (underlying ? cp).toJava + override def shutdown(): Unit = underlying.shutdown() } diff --git a/mqtt-streaming/src/main/scala/akka/stream/alpakka/mqtt/streaming/scaladsl/MqttSession.scala b/mqtt-streaming/src/main/scala/akka/stream/alpakka/mqtt/streaming/scaladsl/MqttSession.scala index 03e35841e6..9df9631e33 100644 --- a/mqtt-streaming/src/main/scala/akka/stream/alpakka/mqtt/streaming/scaladsl/MqttSession.scala +++ b/mqtt-streaming/src/main/scala/akka/stream/alpakka/mqtt/streaming/scaladsl/MqttSession.scala @@ -55,6 +55,31 @@ abstract class MqttSession { */ def ![A](cp: Command[A]): Unit + /** + * Ask the session to perform a command regardless of the state it is + * in. This is important for sending Publish messages in particular, + * as a connection may not have been established with a session. + * @param cp The command to perform + * @tparam A The type of any carry for the command. + * @return A future indicating when the command has completed. Completion + * is defined as when it has been acknowledged by the recipient + * endpoint. + */ + final def ask[A](cp: Command[A]): Future[A] = + this ? cp + + /** + * Ask the session to perform a command regardless of the state it is + * in. This is important for sending Publish messages in particular, + * as a connection may not have been established with a session. + * @param cp The command to perform + * @tparam A The type of any carry for the command. + * @return A future indicating when the command has completed. Completion + * is defined as when it has been acknowledged by the recipient + * endpoint. + */ + def ?[A](cp: Command[A]): Future[A] + /** * Shutdown the session gracefully */ @@ -171,6 +196,9 @@ final class ActorMqttClientSession(settings: MqttSessionSettings)(implicit mat: case c: Command[A] => throw new IllegalStateException(c + " is not a client command that can be sent directly") } + override def ?[A](cp: Command[A]): Future[A] = + ??? + override def shutdown(): Unit = { system.stop(clientConnector.toClassic) system.stop(consumerPacketRouter.toClassic) @@ -513,6 +541,9 @@ final class ActorMqttServerSession(settings: MqttSessionSettings)(implicit mat: case c: Command[A] => throw new IllegalStateException(c + " is not a server command that can be sent directly") } + override def ?[A](cp: Command[A]): Future[A] = + ??? + override def shutdown(): Unit = { system.stop(serverConnector.toClassic) system.stop(consumerPacketRouter.toClassic) diff --git a/mqtt-streaming/src/test/java/docs/javadsl/MqttFlowTest.java b/mqtt-streaming/src/test/java/docs/javadsl/MqttFlowTest.java index eb5746bc18..74ff56491b 100644 --- a/mqtt-streaming/src/test/java/docs/javadsl/MqttFlowTest.java +++ b/mqtt-streaming/src/test/java/docs/javadsl/MqttFlowTest.java @@ -129,16 +129,20 @@ public Publish apply(DecodeErrorOrEvent x, boolean isCheck) { SourceQueueWithComplete> commands = run.first(); commands.offer(new Command<>(new Connect(clientId, ConnectFlags.CleanSession()))); commands.offer(new Command<>(new Subscribe(topic))); - session.tell( - new Command<>( - new Publish( - ControlPacketFlags.RETAIN() | ControlPacketFlags.QoSAtLeastOnceDelivery(), - topic, - ByteString.fromString("ohi")))); + CompletionStage publishDone = + session.ask( + new Command<>( + new Publish( + ControlPacketFlags.RETAIN() | ControlPacketFlags.QoSAtLeastOnceDelivery(), + topic, + ByteString.fromString("ohi")), + Done.getInstance())); // #run-streaming-flow - CompletionStage event = run.second(); - Publish publishEvent = event.toCompletableFuture().get(TIMEOUT_SECONDS, TimeUnit.SECONDS); + publishDone.toCompletableFuture().get(TIMEOUT_SECONDS, TimeUnit.SECONDS); + + CompletionStage events = run.second(); + Publish publishEvent = events.toCompletableFuture().get(TIMEOUT_SECONDS, TimeUnit.SECONDS); assertEquals(publishEvent.topicName(), topic); assertEquals(publishEvent.payload(), ByteString.fromString("ohi")); diff --git a/mqtt-streaming/src/test/scala/docs/scaladsl/MqttFlowSpec.scala b/mqtt-streaming/src/test/scala/docs/scaladsl/MqttFlowSpec.scala index 8cdb0ff45e..4ff7bfd900 100644 --- a/mqtt-streaming/src/test/scala/docs/scaladsl/MqttFlowSpec.scala +++ b/mqtt-streaming/src/test/scala/docs/scaladsl/MqttFlowSpec.scala @@ -70,11 +70,14 @@ trait MqttFlowSpec extends WordSpecLike with Matchers with BeforeAndAfterAll wit commands.offer(Command(Connect(clientId, ConnectFlags.CleanSession))) commands.offer(Command(Subscribe(topic))) - session ! Command( - Publish(ControlPacketFlags.RETAIN | ControlPacketFlags.QoSAtLeastOnceDelivery, topic, ByteString("ohi")) - ) + val publishDone = session ? Command( + Publish(ControlPacketFlags.RETAIN | ControlPacketFlags.QoSAtLeastOnceDelivery, topic, ByteString("ohi")), + Done + ) + //#run-streaming-flow + publishDone.futureValue shouldBe Done events.futureValue match { case Publish(_, `topic`, _, bytes) => bytes shouldBe ByteString("ohi") case e => fail("Unexpected event: " + e) From 0fe6ec446183acd1399c3d16fdd0147fa2ca48b6 Mon Sep 17 00:00:00 2001 From: huntc Date: Wed, 20 Nov 2019 10:00:43 +1100 Subject: [PATCH 2/2] New APIs, breaking backward compatibility --- .../src/main/mima-filters/2.0.0-M1.backwards.excludes | 4 ++++ 1 file changed, 4 insertions(+) create mode 100644 mqtt-streaming/src/main/mima-filters/2.0.0-M1.backwards.excludes diff --git a/mqtt-streaming/src/main/mima-filters/2.0.0-M1.backwards.excludes b/mqtt-streaming/src/main/mima-filters/2.0.0-M1.backwards.excludes new file mode 100644 index 0000000000..669cde88ff --- /dev/null +++ b/mqtt-streaming/src/main/mima-filters/2.0.0-M1.backwards.excludes @@ -0,0 +1,4 @@ +# PR #1908 +# https://github.com/akka/alpakka/pull/1908 +ProblemFilters.exclude[ReversedMissingMethodProblem]("akka.stream.alpakka.mqtt.streaming.javadsl.MqttSession.ask") +ProblemFilters.exclude[ReversedMissingMethodProblem]("akka.stream.alpakka.mqtt.streaming.scaladsl.MqttSession.?") \ No newline at end of file