1717
1818package org .apache .pekko .persistence .journal
1919
20- import scala .collection .immutable
20+ import org .apache .pekko .persistence .journal .AsyncWriteJournalResponseOrderSpec ._
21+
22+ import scala .collection .{ immutable , mutable }
2123import scala .concurrent .{ ExecutionContext , Future , Promise }
2224import scala .util .Try
23-
2425import org .apache .pekko .persistence .{ AtomicWrite , JournalProtocol , PersistenceSpec , PersistentRepr }
2526import org .apache .pekko .testkit .ImplicitSender
2627
@@ -34,90 +35,103 @@ class AsyncWriteJournalResponseOrderSpec
3435 PersistenceSpec .config(
3536 plugin = " " , // we will provide explicit plugin IDs later
3637 test = classOf [AsyncWriteJournalResponseOrderSpec ].getSimpleName,
38+ // using the default system dispatcher to make sure write-response-global-order works with it
39+ // see: https://github.com/apache/pekko/pull/2434
3740 extraConfig = Some (
3841 s """
39- |pekko.persistence.journal.reverse-plugin {
42+ | ${ ControlledWriteCompletionPlugin . BaseId } {
4043 | with-global-order {
41- | class = " ${classOf [AsyncWriteJournalResponseOrderSpec . ReversePlugin ].getName}"
42- |
44+ | class = " ${classOf [ControlledWriteCompletionPlugin ].getName}"
45+ | plugin-dispatcher = "pekko.actor.default-dispatcher"
4346 | write-response-global-order = on
4447 | }
4548 | no-global-order {
46- | class = " ${classOf [AsyncWriteJournalResponseOrderSpec . ReversePlugin ].getName}"
47- |
49+ | class = " ${classOf [ControlledWriteCompletionPlugin ].getName}"
50+ | plugin-dispatcher = "pekko.actor.default-dispatcher"
4851 | write-response-global-order = off
4952 | }
5053 |}
5154 | """ .stripMargin
5255 ))) with ImplicitSender {
5356
54- import AsyncWriteJournalResponseOrderSpec ._
55-
5657 " AsyncWriteJournal" must {
5758 " return write responses in request order if global response order is enabled" in {
5859 val pluginRef =
59- extension.journalFor(journalPluginId = " pekko.persistence.journal.reverse-plugin.with-global-order" )
60-
61- pluginRef ! mkWriteMessages(1 )
62- pluginRef ! mkWriteMessages(2 )
63- pluginRef ! mkWriteMessages(3 )
64-
65- pluginRef ! CompleteWriteOps
66-
67- getMessageNumsFromResponses(receiveN(6 )) shouldEqual Vector (1 , 2 , 3 )
60+ extension.journalFor(journalPluginId = s " ${ControlledWriteCompletionPlugin .BaseId }.with-global-order " )
61+
62+ // request writes for persistence Ids 1..9
63+ 1 .to(9 ).foreach { persistenceId =>
64+ pluginRef ! mkWriteMessages(persistenceId = persistenceId)
65+ }
66+ // complete writes for all but the first persistence ID in reverse order
67+ pluginRef ! CompleteWriteOps (persistenceIdsInOrder = 2 .to(9 ).toVector.reverse)
68+ // AsyncWriteJournal should hold the responses yet to preserve the response order
69+ expectNoMessage()
70+ // complete write for the first persistence ID
71+ pluginRef ! CompleteWriteOps (persistenceIdsInOrder = Vector (1 ))
72+ // now we should receive all responses in the order in which they have been requested
73+ getPersistenceIdsFromResponses(receiveN(18 )) shouldEqual 1 .to(9 ).toVector
6874 }
6975
70- " return write responses in completion order if global response order is disabled" in {
76+ " return write responses as soon as the operation is complete if global response order is disabled" in {
7177 val pluginRef =
72- extension.journalFor(journalPluginId = " pekko.persistence.journal.reverse-plugin.no-global-order" )
73-
74- pluginRef ! mkWriteMessages(1 )
75- pluginRef ! mkWriteMessages(2 )
76- pluginRef ! mkWriteMessages(3 )
77-
78- pluginRef ! CompleteWriteOps
79-
80- getMessageNumsFromResponses(receiveN(6 )) shouldEqual Vector (3 , 2 , 1 )
78+ extension.journalFor(journalPluginId = s " ${ControlledWriteCompletionPlugin .BaseId }.no-global-order " )
79+
80+ // request writes for persistence Ids 1..9
81+ 1 .to(9 ).foreach { persistenceId =>
82+ pluginRef ! mkWriteMessages(persistenceId = persistenceId)
83+ }
84+ // complete writes for all but the first persistence ID in reverse order
85+ pluginRef ! CompleteWriteOps (persistenceIdsInOrder = 2 .to(9 ).toVector.reverse)
86+ // AsyncWriteJournal sends out write responses for 2..9 right away without waiting
87+ getPersistenceIdsFromResponses(receiveN(16 )).toSet shouldEqual 2 .to(9 ).toSet
88+ // complete write for the first persistence ID
89+ pluginRef ! CompleteWriteOps (persistenceIdsInOrder = Vector (1 ))
90+ // and now we finally receive the response for persistence ID 1
91+ getPersistenceIdsFromResponses(receiveN(2 )) shouldEqual Vector (1 )
8192 }
8293 }
8394
84- private def mkWriteMessages (num : Int ): JournalProtocol .WriteMessages = JournalProtocol .WriteMessages (
95+ private def mkWriteMessages (persistenceId : Int ): JournalProtocol .WriteMessages = JournalProtocol .WriteMessages (
8596 messages = Vector (AtomicWrite (PersistentRepr (
86- payload = num ,
97+ payload = " " ,
8798 sequenceNr = 0L ,
88- persistenceId = num .toString
99+ persistenceId = persistenceId .toString
89100 ))),
90101 persistentActor = self,
91102 actorInstanceId = 1
92103 )
93104
94- private def getMessageNumsFromResponses (responses : Seq [AnyRef ]): Vector [Int ] = responses.collect {
105+ private def getPersistenceIdsFromResponses (responses : Seq [AnyRef ]): Vector [Int ] = responses.collect {
95106 case successResponse : JournalProtocol .WriteMessageSuccess =>
96- successResponse.persistent.payload. asInstanceOf [ Int ]
107+ successResponse.persistent.persistenceId.toInt
97108 }.toVector
98109}
99110
100111private object AsyncWriteJournalResponseOrderSpec {
101- case object CompleteWriteOps
112+ final case class CompleteWriteOps ( persistenceIdsInOrder : Vector [ Int ])
102113
103114 /**
104- * Accumulates asyncWriteMessages requests and completes them in reverse receive order on [[CompleteWriteOps ]] command
115+ * Accumulates asyncWriteMessages requests (one for each persistence ID is expected)
116+ * and completes them in requested order on [[CompleteWriteOps ]] command.
105117 */
106- class ReversePlugin extends AsyncWriteJournal {
118+ final class ControlledWriteCompletionPlugin extends AsyncWriteJournal {
107119
108120 private implicit val ec : ExecutionContext = context.dispatcher
109121
110- private var pendingOps : Vector [ Promise [Unit ]] = Vector .empty
122+ private val pendingOps : mutable. HashMap [ Int , Promise [Unit ]] = mutable. HashMap .empty
111123
112124 override def receivePluginInternal : Receive = {
113- case CompleteWriteOps =>
114- pendingOps.reverse.foreach(_.success(()))
115- pendingOps = Vector .empty
125+ case cmd : CompleteWriteOps =>
126+ cmd.persistenceIdsInOrder.foreach { persistenceId =>
127+ pendingOps(persistenceId).success(())
128+ pendingOps -= persistenceId
129+ }
116130 }
117131
118132 override def asyncWriteMessages (messages : immutable.Seq [AtomicWrite ]): Future [immutable.Seq [Try [Unit ]]] = {
119133 val responsePromise = Promise [Unit ]()
120- pendingOps = pendingOps :+ responsePromise
134+ pendingOps.put(messages.head.persistenceId.toInt, responsePromise)
121135 responsePromise.future.map(_ => Vector .empty)
122136 }
123137
@@ -128,4 +142,8 @@ private object AsyncWriteJournalResponseOrderSpec {
128142
129143 override def asyncReadHighestSequenceNr (persistenceId : String , fromSequenceNr : Long ): Future [Long ] = ???
130144 }
145+
146+ object ControlledWriteCompletionPlugin {
147+ val BaseId : String = " pekko.persistence.journal.controlled-write-completion-plugin"
148+ }
131149}
0 commit comments