diff --git a/app/mqttClient/MessageProcessorActor.scala b/app/mqttClient/MessageProcessorActor.scala new file mode 100644 index 0000000..de2e979 --- /dev/null +++ b/app/mqttClient/MessageProcessorActor.scala @@ -0,0 +1,46 @@ +package mqttClient + +import com.google.inject.Provides +import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient +import com.hivemq.client.mqtt.mqtt5.message.publish.Mqtt5Publish +import jakarta.inject.{Inject, Singleton} +import org.apache.pekko.actor.typed.{ActorRef, Behavior} +import org.apache.pekko.actor.typed.scaladsl.Behaviors +import play.api.Logging +import play.api.libs.concurrent.ActorModule + +import java.nio.charset.StandardCharsets + +object MessageProcessorActor extends ActorModule with Logging { + sealed trait Command + case object Stop extends Command + private case class MqttMessageReceived(publish: Mqtt5Publish) extends Command + + override type Message = Command + + @Singleton + class Reg @Inject()(private val mqttMessageProcessorActor: ActorRef[Command]) + + @Provides + def create(client: Mqtt5AsyncClient): Behavior[Command] = Behaviors.setup { context => + client.connect().whenComplete { (_action, _throwable) => + client.subscribeWith() + .topicFilter("test/+") + .callback{ (publish) => + context.self ! MqttMessageReceived(publish) + } + .send() + } + + Behaviors.receiveMessage[Command] { + case MqttMessageReceived(publish) => + val topic = publish.getTopic + val payload = String(publish.getPayloadAsBytes, StandardCharsets.UTF_8) + logger.info(s"Received message $topic $payload") + Behaviors.same + case Stop => Behaviors.stopped{ () => + client.disconnect() + } + } + } +} diff --git a/app/mqttClient/Module.scala b/app/mqttClient/Module.scala new file mode 100644 index 0000000..5bab635 --- /dev/null +++ b/app/mqttClient/Module.scala @@ -0,0 +1,36 @@ +package mqttClient + +import com.google.inject.{AbstractModule, Provides} +import com.hivemq.client.mqtt.mqtt5.{Mqtt5AsyncClient, Mqtt5Client} +import jakarta.inject.Inject +import play.api.libs.concurrent.PekkoGuiceSupport +import play.api.{Configuration, Logging} + +import java.nio.charset.StandardCharsets + +class Module extends AbstractModule with PekkoGuiceSupport with Logging { + override def configure(): Unit = { + bindTypedActor(MessageProcessorActor, "mqtt-message-processor-actor") + bind(classOf[MessageProcessorActor.Reg]).asEagerSingleton() + } + + @Provides + def mqtt5Client(config: Configuration): Mqtt5AsyncClient = { + val mqttConfig = config.get[Configuration]("mqtt") + + var builder = Mqtt5Client.builder() + .identifier(mqttConfig.get[String]("clientId")) + .serverHost(mqttConfig.get[String]("host")) + .serverPort(mqttConfig.get("port")) + + mqttConfig.getOptional[String]("username").foreach{username => + val password = mqttConfig.get[String]("password").getBytes(StandardCharsets.US_ASCII) + builder = builder.simpleAuth() + .username(username) + .password(password) + .applySimpleAuth() + } + + builder.buildAsync() + } +} diff --git a/build.sbt b/build.sbt index 7ce3f6b..143801d 100644 --- a/build.sbt +++ b/build.sbt @@ -21,6 +21,9 @@ libraryDependencies += "org.playframework" %% "play-slick-evolutions" % "6.2.0" // Source: https://mvnrepository.com/artifact/de.mkammerer/argon2-jvm libraryDependencies += "de.mkammerer" % "argon2-jvm" % "2.12" +// Source: https://mvnrepository.com/artifact/com.hivemq/hivemq-mqtt-client +libraryDependencies += "com.hivemq" % "hivemq-mqtt-client" % "1.3.17" + // Adds additional packages into Twirl //TwirlKeys.templateImports += "me.artemis.controllers._" diff --git a/conf/application.conf b/conf/application.conf index 90c5950..32acc96 100644 --- a/conf/application.conf +++ b/conf/application.conf @@ -2,6 +2,8 @@ play.http.filters=controllers.filters.Filters +play.modules.enabled += "mqttClient.Module" + slick.dbs.default = { profile = "slick.jdbc.PostgresProfile$" db = { @@ -26,5 +28,20 @@ pekko { } throughput = 1 } + + mqtt-dispatcher { + type = Dispatcher + executor = "fork-join-executor" + fork-join-executor { + } + } } +} + +mqtt { + host = "localhost" + port = 1883 + clientId = "neon-vertigo-api" + username = "artemis" + password = "artemis" } \ No newline at end of file diff --git a/conf/logback.xml b/conf/logback.xml index ab6c2b1..61ea438 100644 --- a/conf/logback.xml +++ b/conf/logback.xml @@ -41,6 +41,7 @@ +