This is an automated email from the ASF dual-hosted git repository.

fanningpj pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-pekko-samples.git


The following commit(s) were added to refs/heads/main by this push:
     new d66906e  Created WorkResultConsumerActor for handling published events 
(#67)
d66906e is described below

commit d66906ecbd6271ef7b30c5341cae49ee7547b3ee
Author: Mario Renau - Alstom 
<[email protected]>
AuthorDate: Tue Aug 15 22:54:57 2023 +0200

    Created WorkResultConsumerActor for handling published events (#67)
    
    * Created WorkResultConsumerActor for handling published events
    
    * reorder imports
    
    ---------
    
    Co-authored-by: mario <[email protected]>
    Co-authored-by: PJ Fanning <[email protected]>
---
 .../src/main/scala/worker/Main.scala               | 14 +++++----
 .../src/main/scala/worker/WorkManager.scala        |  3 ++
 .../scala/worker/WorkResultConsumerActor.scala     | 33 ++++++++++++++++++++++
 3 files changed, 45 insertions(+), 5 deletions(-)

diff --git 
a/pekko-sample-distributed-workers-scala/src/main/scala/worker/Main.scala 
b/pekko-sample-distributed-workers-scala/src/main/scala/worker/Main.scala
index 42a9558..03002fd 100644
--- a/pekko-sample-distributed-workers-scala/src/main/scala/worker/Main.scala
+++ b/pekko-sample-distributed-workers-scala/src/main/scala/worker/Main.scala
@@ -2,14 +2,12 @@ package worker
 
 import java.io.File
 import java.util.concurrent.CountDownLatch
-
 import org.apache.pekko.actor.typed.ActorSystem
-import org.apache.pekko.actor.typed.scaladsl.Behaviors
-import org.apache.pekko.cluster.typed.Cluster
+import org.apache.pekko.actor.typed.eventstream.EventStream
+import org.apache.pekko.actor.typed.scaladsl.{ ActorContext, Behaviors }
+import org.apache.pekko.cluster.typed.{ Cluster, SelfUp, Subscribe }
 import org.apache.pekko.persistence.cassandra.testkit.CassandraLauncher
 import com.typesafe.config.{ Config, ConfigFactory }
-import org.apache.pekko.cluster.typed.SelfUp
-import org.apache.pekko.cluster.typed.Subscribe
 
 object Main {
 
@@ -39,6 +37,11 @@ object Main {
     }
   }
 
+  private def createWorkResultsConsumer(ctx: ActorContext[SelfUp]): Unit = {
+    val eventListener = ctx.spawn(WorkResultConsumerActor(), 
"work-consumer-actor")
+    ctx.system.eventStream ! EventStream.Subscribe(eventListener)
+  }
+
   def startClusterInSameJvm(): Unit = {
     startCassandraDatabase()
     // two backend nodes
@@ -57,6 +60,7 @@ object Main {
       Behaviors.setup[SelfUp](ctx => {
         val cluster = Cluster(ctx.system)
         cluster.subscriptions ! Subscribe(ctx.self, classOf[SelfUp])
+        createWorkResultsConsumer(ctx)
         Behaviors.receiveMessage {
           case SelfUp(_) =>
             ctx.log.info("Node is up")
diff --git 
a/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkManager.scala
 
b/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkManager.scala
index a2a07f0..11a2611 100644
--- 
a/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkManager.scala
+++ 
b/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkManager.scala
@@ -119,6 +119,9 @@ object WorkManager {
           // Any in progress work from the previous incarnation is retried
           ctx.self ! ResetWorkInProgress
       }
+        // Publish events to the system event stream as PublishedEvent after 
they have been persisted
+        .withEventPublishing(enabled = true)
+
     }
 
 }
diff --git 
a/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkResultConsumerActor.scala
 
b/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkResultConsumerActor.scala
new file mode 100644
index 0000000..dbaf878
--- /dev/null
+++ 
b/pekko-sample-distributed-workers-scala/src/main/scala/worker/WorkResultConsumerActor.scala
@@ -0,0 +1,33 @@
+package worker
+
+import org.apache.pekko.actor.typed.Behavior
+import org.apache.pekko.actor.typed.scaladsl.{ ActorContext, Behaviors }
+import org.apache.pekko.persistence.typed.PublishedEvent
+import worker.WorkState.{ WorkAccepted, WorkCompleted, WorkInProgressReset, 
WorkStarted }
+
+object WorkResultConsumerActor {
+  def apply(): Behavior[PublishedEvent] =
+    Behaviors.setup { context =>
+      context.log.info("WorkResultConsumerActor started")
+
+      Behaviors.receiveMessage { message =>
+        handleReceivedEvent(context, message)
+        Behaviors.same
+      }
+    }
+
+  private def handleReceivedEvent(context: ActorContext[PublishedEvent], 
event: PublishedEvent): Unit = {
+    val actualEvent = event.event
+    event.event match {
+      case WorkInProgressReset =>
+        context.log.info(s"Received published event 
[${actualEvent.getClass.getCanonicalName}]: {}", event)
+      case WorkCompleted(workId) =>
+        context.log.info(s"Received published event 
[${actualEvent.getClass.getCanonicalName}]: workId {}", workId)
+      case WorkStarted(workId) =>
+        context.log.info(s"Received published event 
[${actualEvent.getClass.getCanonicalName}]: workId {}", workId)
+      case WorkAccepted(workId) =>
+        context.log.info(s"Received published event 
[${actualEvent.getClass.getCanonicalName}]: workId {}", workId)
+      case _ => context.log.warn("Message not supported")
+    }
+  }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to