Wait for head for the full response (#523)
This commit is contained in:
@@ -12,6 +12,7 @@ import org.springframework.stereotype.Service
|
|||||||
import reactor.core.publisher.Flux
|
import reactor.core.publisher.Flux
|
||||||
import reactor.core.publisher.Mono
|
import reactor.core.publisher.Mono
|
||||||
import reactor.kotlin.core.publisher.switchIfEmpty
|
import reactor.kotlin.core.publisher.switchIfEmpty
|
||||||
|
import java.time.Duration
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
class SubscribeChainStatus(
|
class SubscribeChainStatus(
|
||||||
@@ -79,8 +80,13 @@ class SubscribeChainStatus(
|
|||||||
.map { toFullResponse(it!!, ms) }
|
.map { toFullResponse(it!!, ms) }
|
||||||
.switchIfEmpty {
|
.switchIfEmpty {
|
||||||
// in case if there is still no head we mush wait until we get it
|
// in case if there is still no head we mush wait until we get it
|
||||||
ms.getHead()
|
// also we have to use 2 approaches due to the head's flux can be stopped
|
||||||
.getFlux()
|
Flux.concat(
|
||||||
|
ms.getHead()
|
||||||
|
.getFlux(),
|
||||||
|
Flux.interval(Duration.ofSeconds(3))
|
||||||
|
.mapNotNull { ms.getHead().getCurrent() },
|
||||||
|
)
|
||||||
.next()
|
.next()
|
||||||
.map { toFullResponse(it!!, ms) }
|
.map { toFullResponse(it!!, ms) }
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -63,7 +63,7 @@ class SubscribeChainStatusTest {
|
|||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `first full event if there is already an ms head`() {
|
fun `first full event if there is already a ms head`() {
|
||||||
val head = mock<Head> {
|
val head = mock<Head> {
|
||||||
on { getCurrent() } doReturn head(550)
|
on { getCurrent() } doReturn head(550)
|
||||||
on { getFlux() } doReturn Flux.empty()
|
on { getFlux() } doReturn Flux.empty()
|
||||||
@@ -106,6 +106,29 @@ class SubscribeChainStatusTest {
|
|||||||
.verify(Duration.ofSeconds(1))
|
.verify(Duration.ofSeconds(1))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `first full event with awaiting a head from polling`() {
|
||||||
|
val head = mock<Head> {
|
||||||
|
on { getCurrent() } doReturn null doReturn head(550)
|
||||||
|
on { getFlux() } doReturn Flux.empty()
|
||||||
|
}
|
||||||
|
val ms = spy<TestMultistream> {
|
||||||
|
on { getHead() } doReturn head
|
||||||
|
on { stateEvents() } doReturn Flux.empty()
|
||||||
|
}
|
||||||
|
val msHolder = mock<MultistreamHolder> {
|
||||||
|
on { all() } doReturn listOf(ms)
|
||||||
|
}
|
||||||
|
val subscribeChainStatus = SubscribeChainStatus(msHolder, chainEventMapper)
|
||||||
|
|
||||||
|
StepVerifier.withVirtualTime { subscribeChainStatus.chainStatuses() }
|
||||||
|
.expectSubscription()
|
||||||
|
.expectNoEvent(Duration.ofSeconds(3))
|
||||||
|
.expectNext(response(true))
|
||||||
|
.thenCancel()
|
||||||
|
.verify(Duration.ofSeconds(1))
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `first full event with awaiting a head from a head stream and then state events`() {
|
fun `first full event with awaiting a head from a head stream and then state events`() {
|
||||||
val head = mock<Head> {
|
val head = mock<Head> {
|
||||||
|
|||||||
Reference in New Issue
Block a user