fix review comments
This commit is contained in:
@@ -72,6 +72,7 @@ open class ConfiguredUpstreams(
|
|||||||
override fun run(args: ApplicationArguments) {
|
override fun run(args: ApplicationArguments) {
|
||||||
log.debug("Starting upstreams")
|
log.debug("Starting upstreams")
|
||||||
val defaultOptions = buildDefaultOptions(config)
|
val defaultOptions = buildDefaultOptions(config)
|
||||||
|
val observerScheduler = Schedulers.newParallel("status-observer", 3)
|
||||||
config.upstreams.forEach { up ->
|
config.upstreams.forEach { up ->
|
||||||
if (!up.isEnabled) {
|
if (!up.isEnabled) {
|
||||||
log.debug("Upstream ${up.id} is disabled")
|
log.debug("Upstream ${up.id} is disabled")
|
||||||
@@ -106,23 +107,18 @@ open class ConfiguredUpstreams(
|
|||||||
options
|
options
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
else -> {
|
|
||||||
log.error("Chain is unsupported: ${up.chain}")
|
|
||||||
return@forEach
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
upstream?.let {
|
upstream?.let {
|
||||||
Flux.concat(Mono.just(UpstreamChangeEvent.ChangeType.ADDED), upstream.observeStatus())
|
Flux.concat(Mono.just(UpstreamChangeEvent.ChangeType.ADDED), upstream.observeStatus())
|
||||||
.distinctUntilChanged()
|
.distinctUntilChanged()
|
||||||
.subscribeOn(Schedulers.boundedElastic())
|
.subscribeOn(observerScheduler)
|
||||||
.subscribe { status ->
|
.subscribe { status ->
|
||||||
when (status) {
|
when (status) {
|
||||||
UpstreamAvailability.UNAVAILABLE -> UpstreamChangeEvent.ChangeType.REMOVED
|
UpstreamAvailability.UNAVAILABLE -> UpstreamChangeEvent.ChangeType.REMOVED
|
||||||
else -> UpstreamChangeEvent.ChangeType.REVALIDATED
|
else -> UpstreamChangeEvent.ChangeType.REVALIDATED
|
||||||
}.let { eventType ->
|
}.let { eventType ->
|
||||||
if (eventType == UpstreamChangeEvent.ChangeType.REMOVED) {
|
if (eventType == UpstreamChangeEvent.ChangeType.REMOVED) {
|
||||||
log.warn("Remove upstream ${upstream.getId()} due to $it")
|
log.warn("Remove upstream ${it::class.java.simpleName}:${upstream.getId()} due to $status")
|
||||||
}
|
}
|
||||||
eventPublisher.publishEvent(UpstreamChangeEvent(chain, upstream, eventType))
|
eventPublisher.publishEvent(UpstreamChangeEvent(chain, upstream, eventType))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -136,10 +136,11 @@ abstract class Multistream(
|
|||||||
|
|
||||||
fun removeUpstream(id: String): Boolean =
|
fun removeUpstream(id: String): Boolean =
|
||||||
upstreams.removeIf { up ->
|
upstreams.removeIf { up ->
|
||||||
up.takeIf { up.getId() == id }
|
(up.getId() == id).also {
|
||||||
?.also { removed[id] = up }
|
if (it) {
|
||||||
?.let { true }
|
removed[id] = up
|
||||||
?: false
|
}
|
||||||
|
}
|
||||||
}.also {
|
}.also {
|
||||||
if (it) {
|
if (it) {
|
||||||
onUpstreamsUpdated()
|
onUpstreamsUpdated()
|
||||||
|
|||||||
Reference in New Issue
Block a user