From a287e74bd76926d1739465c75f0b7271d2fa7863 Mon Sep 17 00:00:00 2001 From: rob Date: Thu, 30 Jul 2026 05:59:19 +0000 Subject: [PATCH] fix: SIGHUP reload supports upstream and method removal Reload no longer crashes on Chain.UNSPECIFIED, removes gRPC upstreams by prefixed id with lifecycle cleanup, applies methods.disabled changes synchronously, pushes status via existing SubscribeChainStatus streams, and optionally closes client RPCs so edges reconnect. NativeCall returns -32601 (CODE_METHOD_NOT_EXIST) for disabled/unknown methods. Co-authored-by: Cursor --- FINDINGS.md | 150 ++++++++++++++++++ .../io/emeraldpay/dshackle/GrpcServer.kt | 3 + .../emeraldpay/dshackle/config/MainConfig.kt | 1 + .../dshackle/config/MainConfigReader.kt | 12 ++ .../dshackle/config/ReloadConfig.kt | 9 ++ .../config/reload/ReloadConfigProcessor.kt | 57 ++++++- .../reload/ReloadConfigUpstreamService.kt | 84 ++++++++-- .../dshackle/rpc/GrpcClientRefreshService.kt | 55 +++++++ .../io/emeraldpay/dshackle/rpc/NativeCall.kt | 4 +- .../dshackle/startup/ConfiguredUpstreams.kt | 11 +- .../dshackle/upstream/Multistream.kt | 8 + .../dshackle/upstream/grpc/GrpcUpstreams.kt | 40 ++++- .../upstream/grpc/GrpcUpstreamsRegistry.kt | 55 +++++++ .../config/reload/ReloadConfigTest.kt | 103 +++++++++++- .../upstreams-methods-reload-changed.yaml | 12 ++ .../upstreams-methods-reload-initial.yaml | 9 ++ 16 files changed, 581 insertions(+), 32 deletions(-) create mode 100644 FINDINGS.md create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/config/ReloadConfig.kt create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/rpc/GrpcClientRefreshService.kt create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreamsRegistry.kt create mode 100644 src/test/resources/configs/upstreams-methods-reload-changed.yaml create mode 100644 src/test/resources/configs/upstreams-methods-reload-initial.yaml diff --git a/FINDINGS.md b/FINDINGS.md new file mode 100644 index 00000000..17eafb6d --- /dev/null +++ b/FINDINGS.md @@ -0,0 +1,150 @@ +# dshackle v0.79.10 — SIGHUP reload removal investigation + +**Base tag:** `v0.79.10` (`f2c1bd06`, "Update deps (#877)") +**Branch:** `reload-removal` + +--- + +## (a) Hot-reload codepath and the `Key UNSPECIFIED missing` failure + +### SIGHUP entry + +| Step | Location | +|------|----------| +| OS signal registration | `ReloadConfigSetup.kt:21-22` — `Signal.handle(signalHup, this)` | +| Handler / lock | `ReloadConfigSetup.kt:25-58` — catches exceptions, logs `Config is not reloaded, cause - ${e.message}` | +| Upstream processor | `UpstreamConfigReloadConfigProcessor` in `ReloadConfigProcessor.kt:20-46` | + +### Diff logic + +1. **Full-config equality short-circuit** — `ReloadConfigProcessor.kt:24-26` compares deserialized YAML to `mainConfig.initialConfig` (`ReloadConfigService.kt:27`, `MainConfig.kt:28-38`). Any change (including `methods.disabled`) marks the upstream entry as modified. + +2. **Per-upstream diff** — `ReloadConfigProcessor.kt:53-79`: + - Missing in new config → `removed` + - Present but not equal (data class) → `reloaded` (treated as remove **and** re-add) + - New keys → `added` + +3. **Default-options diff** — `ReloadConfigProcessor.kt:82-108`: changed/removed chain names in `defaultOptions` populate `chainsToReload`, which triggers **bulk removal of all upstreams on that chain** before re-adding (`ReloadConfigUpstreamService.kt:50-56` pre-patch). + +4. **Removal execution (v0.79.10 bugs)** — `ReloadConfigUpstreamService.kt:44-68` (original): + - Calls `multistreamHolder.getUpstream(chain)` for every chain in `chainsToReload` and every `upstreamsToRemove` pair. + - `CurrentMultistreamHolder.getUpstream()` uses `chainMapping.getValue(chain)` (`CurrentMultistreamHolder.kt:33-34`). + - Multistreams are created only for `Chain.entries.filterNot { UNSPECIFIED }` (`MultistreamsConfig.kt:30-31`). + - **`Chain.UNSPECIFIED` is not in the map** → Kotlin throws `NoSuchElementException: Key UNSPECIFIED missing`. + + **When `UNSPECIFIED` enters the reload set:** + - `Global.chainById(unknownChainName)` returns `UNSPECIFIED` (`Global.kt:56-62`). + - `analyzeDefaultOptions` / `analyzeUpstreams` used `chainById` without filtering (`ReloadConfigProcessor.kt:62-63, 100, 106` pre-patch). + - `buildDefaultOptions` in `ConfiguredUpstreams.kt:97-105` (pre-patch) could store options under `UNSPECIFIED` when a default-options chain name was unknown. + +5. **Additional v0.79.10 removal gaps (fixed in patch):** + - **Async removal:** `Multistream.processUpstreamsEvents()` emits to a Reactor sink processed on a scheduler (`Multistream.kt:153-157, 441-444`). Reload called `addUpstreams()` immediately after firing REMOVED events, so add/remove could race. + - **gRPC upstream IDs:** Config `id: foo` registers per-chain upstreams as `foo_ethereum`, `foo_polygon`, etc. (`GenericGrpcUpstream.kt:69`). Removal looked up exact config id (`ReloadConfigUpstreamService.kt:61-64` pre-patch) → **no match, upstream never removed**. + - **gRPC client lifecycle:** `ConfiguredUpstreams` started `GrpcUpstreams.start()` with no registry/stop on config removal (`ConfiguredUpstreams.kt:77-88` pre-patch) — outbound describe loop kept running after YAML entry deleted. + +### `methods.disabled` on reload + +- Parsed into `ManagedCallMethods` at upstream creation (`UpstreamCreator.kt:97-112`). +- A `methods.disabled` change makes the upstream data class unequal → `reloaded` set → remove + re-add path (`ReloadConfigProcessor.kt:69-70`). +- **Does apply on reload** (not startup-only), provided the YAML change is detected. +- Aggregated method list on the multistream is recalculated in `MultistreamState.updateMethods()` (`MultistreamState.kt:81-97`) via `onUpstreamsUpdated()` after REMOVED/ADDED events. + +--- + +## (b) How method availability is announced to gRPC clients + +| Mechanism | RPC | Method list? | Live updates? | +|-----------|-----|--------------|---------------| +| **Describe** | `BlockchainRpc.describe` → `Describe.kt:33-79` | Yes — `DescribeChain.supportedMethods` from `multistream.getMethods().getSupportedMethods()` (`Describe.kt:42-46`) | **No** — unary snapshot at call time | +| **SubscribeChainStatus** | `BlockchainRpc.subscribeChainStatus` → `SubscribeChainStatus.kt:27-47` | Yes — `SupportedMethodsEvent` in `ChainEvent` (`SubscribeChainStatus.kt:71`, `ChainEventMapper.kt:124-125`) | **Yes** — `multistream.stateEvents()` pushes `MethodsEvent` when methods change (`MultistreamStateHandler.kt:11-12`, `MultistreamState.kt:74-76`) | +| **SubscribeStatus** | `SubscribeStatus.kt:34-48` | **No** — availability + quorum only | Yes for status | +| **NativeCall** | — | Enforces at call time via `upstream.getMethods().isAvailable(method)` (`NativeCall.kt:313-327`) | N/A | + +**Implication for dRPC edge gateways** + +- Edges on an active **`SubscribeChainStatus`** stream receive updated method lists without reconnect (comment at `SubscribeChainStatus.kt:28-29` explicitly mentions hot reload). +- Edges that only call **`Describe` at connect** will **not** see removals until reconnect. +- **`Describe` is not streamed.** + +--- + +## (c) NativeCall typed error for disabled / unknown methods + +Proto field: `NativeCallReplyItem.item_error_code` (int), set in `NativeCall.buildResponse()` (`NativeCall.kt:193-196`). + +| Condition | Behaviour (v0.79.10) | +|-----------|------------------------| +| Method not in aggregated list | `InvalidCallContext` with `ChainCallError(RpcResponseError.CODE_METHOD_NOT_EXIST, …)` (`NativeCall.kt:314-327`) | +| Method not callable on remote path | `RpcException(CODE_METHOD_NOT_EXIST, …)` (`NativeCall.kt:445-446`) | +| JSON-RPC standard code | `RpcResponseError.CODE_METHOD_NOT_EXIST = -32601` (`RpcResponseError.java:28`) — "method does not exist / is not available" | + +**Bug (pre-patch):** `InvalidCallContext` stored the **request item id** in `CallError.id`, and `buildResponse` emitted that as `item_error_code` instead of `-32601`. Upstream failures via `ChainException` already used the RPC code as `CallError.id` (`NativeCall.kt:768`). + +**Patch:** `CallError.itemErrorCode()` returns `upstreamError?.code ?: id` (`NativeCall.kt`). + +Distinct codes exist: `-32601` (method missing) vs `-32603` (internal) vs `-32000..` (upstream) per `RpcResponseError.java`. + +--- + +## (d) Graceful client disconnect mechanism (v0.79.10) + +| Mechanism | Present? | Location | +|-----------|----------|----------| +| Server `@PreDestroy shutdown()` | Yes — full server stop only | `GrpcServer.kt:124-129` | +| `maxConnectionIdle(3600s)` | Yes — idle timeout, not reload-triggered | `GrpcServer.kt:81`, `Defaults.kt:33` | +| Per-connection GOAWAY / drain API | **No** | — | +| HTTP/2 GOAWAY on demand | **No** public API in grpc-java server used here | — | + +**Conclusion:** v0.79.10 had **no** reload-triggered client refresh. Edges holding stale `Describe` data need either live `SubscribeChainStatus` or a forced reconnect. + +### Patch approach (`reload-removal` branch) + +1. **`GrpcClientRefreshService`** — server interceptor tracking active `ServerCall`s; `disconnectClients()` closes them with `Status.UNAVAILABLE` (`GrpcClientRefreshService.kt`). +2. **Config flag** — `reload.disconnect-clients-on-removal` (default `true`) in `ReloadConfig.kt` / `MainConfigReader.kt`. +3. **Triggered after** upstream/method removal in `UpstreamConfigReloadConfigProcessor.reload()` (`ReloadConfigProcessor.kt`). + +This closes active RPCs (including long-lived `SubscribeChainStatus`), prompting client reconnect. It is not a raw Netty GOAWAY on the transport, but achieves the same operational outcome for grpc-java clients. + +### Recommendation + +| Edge behaviour | Approach | +|----------------|----------| +| Uses **SubscribeChainStatus** | Live `MethodsEvent` propagation is sufficient for method list; optional disconnect still helps for `Describe`-cached state | +| Uses **Describe-at-connect only** | **Requires reconnect** — use `disconnectClientsOnRemoval` (default on) or accept stale method lists until natural reconnect | +| Calls removed methods | **NativeCall `-32601`** — immediate typed failure even without reconnect | + +**GOAWAY vs error-code:** Use **both**: error-code for in-flight calls to removed methods; disconnect (or SubscribeChainStatus stream) for updating the edge's cached capability set. Error-code alone leaves the edge believing the method still exists. + +--- + +## Patch summary (`reload-removal`) + +1. Filter `Chain.UNSPECIFIED` from reload diffs and default-options processing. +2. Synchronous removal via `Multistream.processUpstreamsEventsSync()`. +3. Match gRPC upstreams by `{configId}_` prefix; `GrpcUpstreamsRegistry` stops outbound clients on removal. +4. Propagate removals through existing `SubscribeChainStatus` / `MultistreamState` path. +5. Optional client disconnect on removal (configurable, default on). +6. Fix `item_error_code` for disabled methods (`-32601`). +7. Tests extended in `ReloadConfigTest.kt`. + +--- + +## Build / test verification + +```text +# Prerequisites (this workspace) +# - foundation published to mavenLocal (./foundation/gradlew publishToMavenLocal) +# - foundation/src/main/resources/public/chains.yaml (drpcorg/public submodule or raw fetch) +# - emerald-grpc/ proto submodule cloned + +export JAVA_HOME=/usr/lib/jvm/java-21-openjdk-amd64 +./gradlew test --tests "io.emeraldpay.dshackle.config.reload.*" +# BUILD SUCCESSFUL — 8 tests (ReloadConfigTest + HealthReloadConfigProcessorTest) + +./gradlew build +# compiles + ktlint pass; full suite: 1635 tests, 2 failed (pre-existing on v0.79.10): +# - GenericWsHeadSpec (WS disconnect liveness) +# - IntegrationTest (Spring context / environment) +``` + +Reload-specific tests all pass. Full `./gradlew build` reports two failures that also occur on unmodified `v0.79.10` in this environment. diff --git a/src/main/kotlin/io/emeraldpay/dshackle/GrpcServer.kt b/src/main/kotlin/io/emeraldpay/dshackle/GrpcServer.kt index e144620b..5ba5ffd8 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/GrpcServer.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/GrpcServer.kt @@ -19,6 +19,7 @@ package io.emeraldpay.dshackle import io.emeraldpay.dshackle.auth.AuthInterceptor import io.emeraldpay.dshackle.config.MainConfig import io.emeraldpay.dshackle.monitoring.accesslog.AccessHandlerGrpc +import io.emeraldpay.dshackle.rpc.GrpcClientRefreshService import io.grpc.Codec import io.grpc.Server import io.grpc.ServerCall @@ -46,6 +47,7 @@ open class GrpcServer( private val accessHandler: AccessHandlerGrpc, private val grpcServerBraveInterceptor: ServerInterceptor, private val authInterceptor: AuthInterceptor, + private val grpcClientRefreshService: GrpcClientRefreshService, ) { @Value("\${spring.application.max-metadata-size}") private var maxMetadataSize: Int = Defaults.maxMetadataSize @@ -89,6 +91,7 @@ open class GrpcServer( } serverBuilder.intercept(grpcServerBraveInterceptor) + serverBuilder.intercept(grpcClientRefreshService) if (mainConfig.authorization.enabled && mainConfig.authorization.hasServerConfig()) { serverBuilder.intercept(authInterceptor) log.info("Token authorization is turned on") diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/MainConfig.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/MainConfig.kt index 0d579aa3..25709bc3 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/config/MainConfig.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/MainConfig.kt @@ -44,6 +44,7 @@ class MainConfig { var compression: CompressionConfig = CompressionConfig.default() var chains: ChainsConfig = ChainsConfig.default() var authorization: AuthorizationConfig = AuthorizationConfig.default() + var reload: ReloadConfig = ReloadConfig() var initialConfig: UpstreamsConfig? = null private set diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/MainConfigReader.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/MainConfigReader.kt index 5664eaf6..2832bf41 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/config/MainConfigReader.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/MainConfigReader.kt @@ -83,6 +83,18 @@ class MainConfigReader( authorizationConfigReader.read(input).let { config.authorization = it } + readReloadConfig(input)?.let { + config.reload = it + } return config } + + private fun readReloadConfig(input: MappingNode?): ReloadConfig? { + val reloadNode = getMapping(input, "reload") ?: return null + val reload = ReloadConfig() + getValueAsBool(reloadNode, "disconnect-clients-on-removal")?.let { + reload.disconnectClientsOnRemoval = it + } + return reload + } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/ReloadConfig.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/ReloadConfig.kt new file mode 100644 index 00000000..3f75aab8 --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/ReloadConfig.kt @@ -0,0 +1,9 @@ +package io.emeraldpay.dshackle.config + +data class ReloadConfig( + /** + * When true (default), gracefully close active gRPC client RPCs after upstream or + * method removal so clients reconnect and re-fetch Describe / SubscribeChainStatus. + */ + var disconnectClientsOnRemoval: Boolean = true, +) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigProcessor.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigProcessor.kt index c886d447..48d08603 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigProcessor.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigProcessor.kt @@ -2,8 +2,11 @@ package io.emeraldpay.dshackle.config.reload import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Global.Companion.chainById +import io.emeraldpay.dshackle.config.MainConfig import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.foundation.ChainOptions +import io.emeraldpay.dshackle.rpc.GrpcClientRefreshService +import org.slf4j.LoggerFactory import org.springframework.stereotype.Component import java.util.stream.Collectors @@ -16,7 +19,12 @@ interface ReloadConfigProcessor { class UpstreamConfigReloadConfigProcessor( private val reloadConfigService: ReloadConfigService, private val reloadConfigUpstreamService: ReloadConfigUpstreamService, + private val mainConfig: MainConfig, + private val grpcClientRefreshService: GrpcClientRefreshService, ) : ReloadConfigProcessor { + + private val log = LoggerFactory.getLogger(UpstreamConfigReloadConfigProcessor::class.java) + override fun reload(): Boolean { val newUpstreamsConfig = reloadConfigService.readUpstreamsConfig() val currentUpstreamsConfig = reloadConfigService.currentUpstreamsConfig() @@ -35,13 +43,29 @@ class UpstreamConfigReloadConfigProcessor( ) val upstreamsToRemove = upstreamsAnalyzeData.removed + .filter { it.second != Chain.UNSPECIFIED } .filterNot { chainsToReload.contains(it.second) } .toSet() val upstreamsToAdd = upstreamsAnalyzeData.added + .filter { it.second != Chain.UNSPECIFIED } + .toSet() reloadConfigService.updateUpstreamsConfig(newUpstreamsConfig) - reloadConfigUpstreamService.reloadUpstreams(chainsToReload, upstreamsToRemove, upstreamsToAdd, newUpstreamsConfig) + val reloadResult = reloadConfigUpstreamService.reloadUpstreams( + chainsToReload, + upstreamsToRemove, + upstreamsToAdd, + newUpstreamsConfig, + ) + + if (reloadResult.hadRemovals && mainConfig.reload.disconnectClientsOnRemoval) { + log.info( + "Disconnecting gRPC clients after upstream/method removal on chains {}", + reloadResult.affectedChains.map { it.chainCode }, + ) + grpcClientRefreshService.disconnectClients("upstream or method removed via config reload") + } return true } @@ -59,8 +83,8 @@ class UpstreamConfigReloadConfigProcessor( } val reloaded = mutableSetOf>() val removed = mutableSetOf>() - val currentUpstreamsMap = currentUpstreams.associateBy { it.id!! to chainById(it.chain) } - val newUpstreamsMap = newUpstreams.associateBy { it.id!! to chainById(it.chain) } + val currentUpstreamsMap = upstreamKeyMap(currentUpstreams) + val newUpstreamsMap = upstreamKeyMap(newUpstreams) currentUpstreamsMap.forEach { val newUpstream = newUpstreamsMap[it.key] @@ -79,6 +103,19 @@ class UpstreamConfigReloadConfigProcessor( return UpstreamAnalyzeData(added, removed.plus(reloaded)) } + private fun upstreamKeyMap( + upstreams: List>, + ): Map, UpstreamsConfig.Upstream<*>> { + return upstreams.mapNotNull { up -> + val chain = chainById(up.chain) + if (chain == Chain.UNSPECIFIED) { + null + } else { + (up.id!! to chain) to up + } + }.toMap() + } + private fun analyzeDefaultOptions( currentDefaultOptions: List, newDefaultOptions: List, @@ -97,13 +134,21 @@ class UpstreamConfigReloadConfigProcessor( currentOptions.forEach { val newChainOption = newOptions[it.key] if (newChainOption == null) { - removed.add(chainById(it.key)) + val chain = chainById(it.key) + if (chain != Chain.UNSPECIFIED) { + removed.add(chain) + } } else if (newChainOption != it.value) { - chainsToReload.add(chainById(it.key)) + val chain = chainById(it.key) + if (chain != Chain.UNSPECIFIED) { + chainsToReload.add(chain) + } } } - val added = newOptions.minus(currentOptions.keys).map { chainById(it.key) } + val added = newOptions.minus(currentOptions.keys).mapNotNull { + chainById(it.key).takeIf { chain -> chain != Chain.UNSPECIFIED } + } return chainsToReload.plus(added).plus(removed) } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigUpstreamService.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigUpstreamService.kt index 81b3f3b7..a0bd47b7 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigUpstreamService.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigUpstreamService.kt @@ -7,21 +7,38 @@ import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.startup.ConfiguredUpstreams import io.emeraldpay.dshackle.startup.UpstreamChangeEvent import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder +import io.emeraldpay.dshackle.upstream.Upstream +import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstreamsRegistry +import org.slf4j.LoggerFactory import org.springframework.stereotype.Component import java.util.stream.Collectors +data class UpstreamReloadResult( + val hadRemovals: Boolean, + val affectedChains: Set, +) + @Component open class ReloadConfigUpstreamService( private val multistreamHolder: CurrentMultistreamHolder, private val configuredUpstreams: ConfiguredUpstreams, + private val grpcUpstreamsRegistry: GrpcUpstreamsRegistry, ) { + private val log = LoggerFactory.getLogger(ReloadConfigUpstreamService::class.java) + fun reloadUpstreams( chainsToReload: Set, upstreamsToRemove: Set>, upstreamsToAdd: Set>, newUpstreamsConfig: UpstreamsConfig, - ) { + ): UpstreamReloadResult { + val safeChainsToReload = chainsToReload.filterValidChains() + val safeUpstreamsToRemove = upstreamsToRemove.filterValidPairs() + val grpcIdsToStop = safeUpstreamsToRemove.map { it.first } + .filter { grpcUpstreamsRegistry.isRegistered(it) } + .toSet() + val newUpstreamsCount = newUpstreamsConfig.upstreams.stream() .collect( Collectors.groupingBy( @@ -30,38 +47,61 @@ open class ReloadConfigUpstreamService( ), ) - val usedChains = removeUpstreams(chainsToReload, upstreamsToRemove) + val usedChains = removeUpstreams(safeChainsToReload, safeUpstreamsToRemove, grpcIdsToStop) - addUpstreams(newUpstreamsConfig, chainsToReload, upstreamsToAdd.map { it.first }.toSet()) + addUpstreams(newUpstreamsConfig, safeChainsToReload, upstreamsToAdd.map { it.first }.toSet()) - usedChains.forEach { chain -> - if (newUpstreamsCount[chain] == null) { + usedChains.filterValidChains().forEach { chain -> + if (newUpstreamsCount[chain] == null || newUpstreamsCount[chain] == 0L) { multistreamHolder.getUpstream(chain).stop() } } + + val hadRemovals = safeUpstreamsToRemove.isNotEmpty() || + safeChainsToReload.isNotEmpty() || + grpcIdsToStop.isNotEmpty() + return UpstreamReloadResult( + hadRemovals = hadRemovals, + affectedChains = usedChains.filterValidChains().toSet(), + ) } private fun removeUpstreams( chainsToReload: Set, upstreamsToRemove: Set>, + grpcIdsToStop: Set, ): Set { val usedChains = mutableSetOf() - chainsToReload.forEach { - usedChains.add(it) - val ms = multistreamHolder.getUpstream(it) + grpcIdsToStop.forEach { id -> + grpcUpstreamsRegistry.stop(id).forEach { event -> + usedChains.add(event.chain) + } + } + + chainsToReload.forEach { chain -> + usedChains.add(chain) + val ms = multistreamHolder.getUpstream(chain) ms.getAll() + .toList() .forEach { up -> - ms.processUpstreamsEvents(UpstreamChangeEvent(it, up, UpstreamChangeEvent.ChangeType.REMOVED)) + ms.processUpstreamsEventsSync( + UpstreamChangeEvent(chain, up, UpstreamChangeEvent.ChangeType.REMOVED), + ) } } + upstreamsToRemove.forEach { pair -> + if (grpcUpstreamsRegistry.isRegistered(pair.first)) { + return@forEach + } usedChains.add(pair.second) val ms = multistreamHolder.getUpstream(pair.second) - ms.getAll() - .find { pair.first == it.getId() } - ?.let { up -> - ms.processUpstreamsEvents(UpstreamChangeEvent(pair.second, up, UpstreamChangeEvent.ChangeType.REMOVED)) + findUpstreamsToRemove(ms.getAll(), pair.first) + .forEach { up -> + ms.processUpstreamsEventsSync( + UpstreamChangeEvent(pair.second, up, UpstreamChangeEvent.ChangeType.REMOVED), + ) } } @@ -81,4 +121,22 @@ open class ReloadConfigUpstreamService( ) configuredUpstreams.processUpstreams(configToReload) } + + companion object { + internal fun findUpstreamsToRemove(upstreams: List, configId: String): List { + val prefix = "${configId}_" + return upstreams.filter { up -> + up.getId() == configId || up.getId().startsWith(prefix) + } + } + + private fun Set.filterValidChains(): Set = + filter { it != Chain.UNSPECIFIED }.toSet() + + private fun Set>.filterValidPairs(): Set> = + filter { it.second != Chain.UNSPECIFIED }.toSet() + + private fun Collection.filterValidChains(): List = + filter { it != Chain.UNSPECIFIED } + } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/GrpcClientRefreshService.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/GrpcClientRefreshService.kt new file mode 100644 index 00000000..0a809ad4 --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/GrpcClientRefreshService.kt @@ -0,0 +1,55 @@ +package io.emeraldpay.dshackle.rpc + +import io.grpc.ForwardingServerCall +import io.grpc.Metadata +import io.grpc.ServerCall +import io.grpc.ServerCallHandler +import io.grpc.ServerInterceptor +import io.grpc.Status +import org.slf4j.LoggerFactory +import org.springframework.stereotype.Component +import java.util.concurrent.ConcurrentHashMap + +/** + * Tracks active inbound gRPC server calls and can close them on config reload so clients + * reconnect and pick up updated Describe / SubscribeChainStatus announcements. + */ +@Component +class GrpcClientRefreshService : ServerInterceptor { + + private val log = LoggerFactory.getLogger(GrpcClientRefreshService::class.java) + + private val activeCalls = ConcurrentHashMap.newKeySet>() + + override fun interceptCall( + call: ServerCall, + headers: Metadata, + next: ServerCallHandler, + ): ServerCall.Listener { + activeCalls.add(call) + val trackedCall = object : ForwardingServerCall.SimpleForwardingServerCall(call) { + override fun close(status: Status, trailers: Metadata) { + activeCalls.remove(call) + super.close(status, trailers) + } + } + return next.startCall(trackedCall, headers) + } + + fun disconnectClients(reason: String = "config reload") { + val status = Status.UNAVAILABLE.withDescription(reason) + val snapshot = activeCalls.toList() + if (snapshot.isEmpty()) { + log.info("No active gRPC client calls to close for reload propagation") + return + } + log.info("Closing {} active gRPC client call(s) after reload: {}", snapshot.size, reason) + snapshot.forEach { call -> + try { + call.close(status, Metadata()) + } catch (e: Exception) { + log.debug("Failed to close gRPC client call during reload", e) + } + } + } +} diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt index 32169a30..50b8d6f6 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt @@ -193,7 +193,7 @@ open class NativeCall( if (it.isError()) { it.error?.let { error -> result.setErrorMessage(error.message) - .setItemErrorCode(error.id) + .setItemErrorCode(error.itemErrorCode()) error.data?.let { data -> result.setErrorData(data) @@ -750,6 +750,8 @@ open class NativeCall( val errorAsIs: ByteArray? = null, ) { + fun itemErrorCode(): Int = upstreamError?.code ?: id + companion object { private val log = LoggerFactory.getLogger(CallError::class.java) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt index 11a2fff5..27b7edec 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt @@ -24,6 +24,7 @@ import io.emeraldpay.dshackle.foundation.ChainOptions import io.emeraldpay.dshackle.startup.configure.UpstreamCreationData import io.emeraldpay.dshackle.startup.configure.UpstreamFactory import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder +import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstreamsRegistry import org.slf4j.LoggerFactory import org.springframework.boot.ApplicationArguments import org.springframework.boot.ApplicationRunner @@ -35,6 +36,7 @@ open class ConfiguredUpstreams( private val config: UpstreamsConfig, private val multistreamHolder: CurrentMultistreamHolder, private val chainsConfig: ChainsConfig, + private val grpcUpstreamsRegistry: GrpcUpstreamsRegistry, ) : ApplicationRunner { private val log = LoggerFactory.getLogger(ConfiguredUpstreams::class.java) @@ -74,11 +76,11 @@ open class ConfiguredUpstreams( ) } } else { - upstreamFactory.createGrpcUpstream( + val grpcUpstreams = upstreamFactory.createGrpcUpstream( up as UpstreamsConfig.Upstream, chainsConfig, ) - .start() + val subscription = grpcUpstreams.start() .doOnNext { log.info("Chain ${it.chain} ${it.type} through gRPC at ${up.connection?.host}:${up.connection?.port}. With caps: ${it.upstream.getCapabilities()}") } @@ -86,6 +88,7 @@ open class ConfiguredUpstreams( multistreamHolder.getUpstream(it.chain) .processUpstreamsEvents(it) } + grpcUpstreamsRegistry.register(up.id!!, grpcUpstreams, subscription) } } } @@ -95,6 +98,10 @@ open class ConfiguredUpstreams( config.defaultOptions.forEach { defaultsConfig -> defaultsConfig.chains?.forEach { chainName -> Global.chainById(chainName).let { chain -> + if (chain == Chain.UNSPECIFIED) { + log.warn("Skipping unknown chain in default options: $chainName") + return@forEach + } defaultsConfig.options?.let { options -> if (!defaultOptions.containsKey(chain)) { defaultOptions[chain] = options diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index 7c7bb332..530ced3a 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -444,6 +444,14 @@ abstract class Multistream( ) { _, res -> res == Sinks.EmitResult.FAIL_NON_SERIALIZED } } + /** + * Applies an upstream change immediately on the calling thread. Used during config reload + * so removals complete before new upstreams are registered. + */ + fun processUpstreamsEventsSync(event: UpstreamChangeEvent) { + onUpstreamChange(event) + } + private fun onUpstreamChange(event: UpstreamChangeEvent) { val chain = event.chain if (this.chain == chain) { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt index e0819849..439df6da 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt @@ -43,6 +43,7 @@ import io.emeraldpay.dshackle.upstream.grpc.auth.GrpcUpstreamsAuth import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient import io.grpc.ClientInterceptor import io.grpc.Codec +import io.grpc.ManagedChannel import io.grpc.Status import io.grpc.StatusRuntimeException import io.grpc.netty.NettyChannelBuilder @@ -64,6 +65,7 @@ import reactor.core.scheduler.Scheduler import java.io.IOException import java.time.Duration import java.util.concurrent.Executor +import java.util.concurrent.TimeUnit import java.util.concurrent.locks.ReentrantLock import kotlin.concurrent.withLock @@ -94,8 +96,10 @@ class GrpcUpstreams( var timeout = Defaults.timeout lateinit var client: ReactorBlockchainGrpc.ReactorBlockchainStub + private var channel: ManagedChannel? = null private val known = HashMap() private val lock = ReentrantLock() + private val statusSubscriptions = mutableMapOf() fun start(): Flux { val chanelBuilder = NettyChannelBuilder.forAddress(host, port) @@ -122,8 +126,9 @@ class GrpcUpstreams( chanelBuilder.usePlaintext() } - val channel = chanelBuilder.build() - var client = ReactorBlockchainGrpc.newReactorStub(channel) + val builtChannel = chanelBuilder.build() + this.channel = builtChannel + var client = ReactorBlockchainGrpc.newReactorStub(builtChannel) if (compression) { client = client.withCompression(Codec.Gzip().messageEncoding) } @@ -132,7 +137,7 @@ class GrpcUpstreams( val grpcUpstreamsAuth = if (tokenAuth != null && authorizationConfig.enabled) { GrpcUpstreamsAuth( - ReactorAuthGrpc.newReactorStub(channel), + ReactorAuthGrpc.newReactorStub(builtChannel), authorizationConfig, grpcAuthContext, tokenAuth.publicKeyPath!!, @@ -141,8 +146,6 @@ class GrpcUpstreams( null } - val statusSubscriptions = mutableMapOf() - return Flux.interval(Duration.ZERO, Duration.ofSeconds(20)) .flatMap { authAndDescribe(grpcUpstreamsAuth) @@ -356,4 +359,31 @@ class GrpcUpstreams( } private fun describe() = this.client.describe(DescribeRequest.newBuilder().build()) + + /** + * Stops the outbound gRPC client, disposes status subscriptions, and returns REMOVED events + * for all chains previously served through this connection. + */ + fun stopAndRemoveAll(): List { + lock.withLock { + statusSubscriptions.values.forEach { it.dispose() } + statusSubscriptions.clear() + channel?.let { ch -> + ch.shutdown() + try { + ch.awaitTermination(5, TimeUnit.SECONDS) + } catch (_: InterruptedException) { + Thread.currentThread().interrupt() + ch.shutdownNow() + } + } + channel = null + return known.map { (chain, upstream) -> + upstream.stop() + UpstreamChangeEvent(chain, upstream, UpstreamChangeEvent.ChangeType.REMOVED) + }.also { + known.clear() + } + } + } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreamsRegistry.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreamsRegistry.kt new file mode 100644 index 00000000..ac019926 --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreamsRegistry.kt @@ -0,0 +1,55 @@ +package io.emeraldpay.dshackle.upstream.grpc + +import io.emeraldpay.dshackle.Chain +import io.emeraldpay.dshackle.startup.UpstreamChangeEvent +import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder +import org.slf4j.LoggerFactory +import org.springframework.stereotype.Component +import reactor.core.Disposable +import java.util.concurrent.ConcurrentHashMap + +@Component +class GrpcUpstreamsRegistry( + private val multistreamHolder: CurrentMultistreamHolder, +) { + private val log = LoggerFactory.getLogger(GrpcUpstreamsRegistry::class.java) + + private data class Entry( + val upstreams: GrpcUpstreams, + val subscription: Disposable, + ) + + private val active = ConcurrentHashMap() + + fun register(id: String, upstreams: GrpcUpstreams, subscription: Disposable) { + active[id]?.let { previous -> + log.warn("Replacing existing gRPC upstream registration for id=$id") + dispose(previous) + } + active[id] = Entry(upstreams, subscription) + } + + fun stop(id: String): Collection { + val entry = active.remove(id) ?: return emptyList() + dispose(entry) + return entry.upstreams.stopAndRemoveAll() + .also { events -> + events.forEach { event -> + if (event.chain != Chain.UNSPECIFIED) { + multistreamHolder.getUpstream(event.chain) + .processUpstreamsEventsSync(event) + } + } + } + } + + fun stopAll(ids: Collection): Boolean { + return ids.map { stop(it) }.any { it.isNotEmpty() } || ids.any { active.containsKey(it) } + } + + fun isRegistered(id: String): Boolean = active.containsKey(id) + + private fun dispose(entry: Entry) { + entry.subscription.dispose() + } +} diff --git a/src/test/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigTest.kt b/src/test/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigTest.kt index 4c82b2d1..71777afa 100644 --- a/src/test/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigTest.kt +++ b/src/test/kotlin/io/emeraldpay/dshackle/config/reload/ReloadConfigTest.kt @@ -10,6 +10,7 @@ import io.emeraldpay.dshackle.config.MainConfig import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfigReader import io.emeraldpay.dshackle.foundation.ChainOptionsReader +import io.emeraldpay.dshackle.rpc.GrpcClientRefreshService import io.emeraldpay.dshackle.startup.ConfiguredUpstreams import io.emeraldpay.dshackle.startup.UpstreamChangeEvent import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder @@ -18,6 +19,7 @@ import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.generic.ChainSpecificRegistry import io.emeraldpay.dshackle.upstream.generic.GenericMultistream import io.emeraldpay.dshackle.upstream.generic.GenericUpstream +import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstreamsRegistry import org.junit.jupiter.api.Assertions.assertEquals import org.junit.jupiter.api.Assertions.assertFalse import org.junit.jupiter.api.Assertions.assertTrue @@ -46,6 +48,8 @@ class ReloadConfigTest { private val config = mock() private val reloadConfigService = ReloadConfigService(config, fileResolver, mainConfig) private val configuredUpstreams = mock() + private val grpcUpstreamsRegistry = mock() + private val grpcClientRefreshService = mock() @BeforeEach fun setupTests() { @@ -88,8 +92,14 @@ class ReloadConfigTest { val reloadConfigUpstreamService = ReloadConfigUpstreamService( currentMultistreamHolder, configuredUpstreams, + grpcUpstreamsRegistry, + ) + val upstreamCfgReloadProcessor = UpstreamConfigReloadConfigProcessor( + reloadConfigService, + reloadConfigUpstreamService, + mainConfig, + grpcClientRefreshService, ) - val upstreamCfgReloadProcessor = UpstreamConfigReloadConfigProcessor(reloadConfigService, reloadConfigUpstreamService) val reloadConfig = ReloadConfigSetup(listOf(upstreamCfgReloadProcessor)) val initialConfigIs = ResourceUtils.getFile("classpath:configs/upstreams-initial.yaml").inputStream() @@ -99,8 +109,8 @@ class ReloadConfigTest { reloadConfig.handle(Signal("HUP")) - verify(msEth).processUpstreamsEvents(UpstreamChangeEvent(ETHEREUM__MAINNET, up1, UpstreamChangeEvent.ChangeType.REMOVED)) - verify(msPoly).processUpstreamsEvents(UpstreamChangeEvent(POLYGON__MAINNET, up3, UpstreamChangeEvent.ChangeType.REMOVED)) + verify(msEth).processUpstreamsEventsSync(UpstreamChangeEvent(ETHEREUM__MAINNET, up1, UpstreamChangeEvent.ChangeType.REMOVED)) + verify(msPoly).processUpstreamsEventsSync(UpstreamChangeEvent(POLYGON__MAINNET, up3, UpstreamChangeEvent.ChangeType.REMOVED)) verify(configuredUpstreams).processUpstreams( UpstreamsConfig( newConfig.defaultOptions, @@ -140,8 +150,14 @@ class ReloadConfigTest { val reloadConfigUpstreamService = ReloadConfigUpstreamService( currentMultistreamHolder, configuredUpstreams, + grpcUpstreamsRegistry, + ) + val upstreamCfgReloadProcessor = UpstreamConfigReloadConfigProcessor( + reloadConfigService, + reloadConfigUpstreamService, + mainConfig, + grpcClientRefreshService, ) - val upstreamCfgReloadProcessor = UpstreamConfigReloadConfigProcessor(reloadConfigService, reloadConfigUpstreamService) val reloadConfig = ReloadConfigSetup(listOf(upstreamCfgReloadProcessor)) val initialConfigIs = ResourceUtils.getFile("classpath:configs/upstreams-initial.yaml").inputStream() val initialConfig = upstreamsConfigReader.read(initialConfigIs)!! @@ -172,7 +188,12 @@ class ReloadConfigTest { val reloadConfigUpstreamService = mock() - val upstreamCfgReloadProcessor = UpstreamConfigReloadConfigProcessor(reloadConfigService, reloadConfigUpstreamService) + val upstreamCfgReloadProcessor = UpstreamConfigReloadConfigProcessor( + reloadConfigService, + reloadConfigUpstreamService, + mainConfig, + grpcClientRefreshService, + ) val reloadConfig = ReloadConfigSetup(listOf(upstreamCfgReloadProcessor)) whenever(config.getConfigPath()).thenReturn(initialConfigFile) @@ -184,6 +205,78 @@ class ReloadConfigTest { assertEquals(initialConfig, mainConfig.upstreams) } + @Test + fun `reload upstream removal disconnects grpc clients by default`() { + val initialConfigFile = ResourceUtils.getFile("classpath:configs/upstreams-initial.yaml") + val newConfigFile = ResourceUtils.getFile("classpath:configs/upstreams-changed-upstreams-removed.yaml") + whenever(config.getConfigPath()).thenReturn(newConfigFile) + + val reloadConfigUpstreamService = mock() + whenever( + reloadConfigUpstreamService.reloadUpstreams(any(), any(), any(), any()), + ).thenReturn(UpstreamReloadResult(hadRemovals = true, affectedChains = setOf(ETHEREUM__MAINNET))) + + val upstreamCfgReloadProcessor = UpstreamConfigReloadConfigProcessor( + reloadConfigService, + reloadConfigUpstreamService, + mainConfig, + grpcClientRefreshService, + ) + val reloadConfig = ReloadConfigSetup(listOf(upstreamCfgReloadProcessor)) + + val initialConfig = upstreamsConfigReader.read(initialConfigFile.inputStream())!! + mainConfig.upstreams = initialConfig + + reloadConfig.handle(Signal("HUP")) + + verify(grpcClientRefreshService).disconnectClients("upstream or method removed via config reload") + } + + @Test + fun `findUpstreamsToRemove matches grpc upstream ids`() { + val ethUpstream = upstream("provider_ethereum") + val polyUpstream = upstream("other_polygon") + val found = ReloadConfigUpstreamService.findUpstreamsToRemove( + listOf(ethUpstream, polyUpstream), + "provider", + ) + assertEquals(listOf(ethUpstream), found) + } + + @Test + fun `reload methods disabled change triggers upstream reload`() { + val initialConfigFile = ResourceUtils.getFile("classpath:configs/upstreams-methods-reload-initial.yaml") + val changedConfigFile = ResourceUtils.getFile("classpath:configs/upstreams-methods-reload-changed.yaml") + val initialConfig = upstreamsConfigReader.read(initialConfigFile.inputStream())!! + mainConfig.upstreams = initialConfig + + val msEth = mock { + on { getAll() } doReturn emptyList() + } + val currentMultistreamHolder = mock { + on { getUpstream(ETHEREUM__MAINNET) } doReturn msEth + } + val reloadConfigUpstreamService = ReloadConfigUpstreamService( + currentMultistreamHolder, + configuredUpstreams, + grpcUpstreamsRegistry, + ) + val upstreamCfgReloadProcessor = UpstreamConfigReloadConfigProcessor( + reloadConfigService, + reloadConfigUpstreamService, + mainConfig, + grpcClientRefreshService, + ) + val reloadConfig = ReloadConfigSetup(listOf(upstreamCfgReloadProcessor)) + + whenever(config.getConfigPath()).thenReturn(changedConfigFile) + + reloadConfig.handle(Signal("HUP")) + + verify(configuredUpstreams).processUpstreams(any()) + verify(grpcClientRefreshService).disconnectClients("upstream or method removed via config reload") + } + private fun multistream(chain: Chain): Multistream { val cs = ChainSpecificRegistry.resolve(chain) return GenericMultistream( diff --git a/src/test/resources/configs/upstreams-methods-reload-changed.yaml b/src/test/resources/configs/upstreams-methods-reload-changed.yaml new file mode 100644 index 00000000..b6fd0df0 --- /dev/null +++ b/src/test/resources/configs/upstreams-methods-reload-changed.yaml @@ -0,0 +1,12 @@ +cluster: + upstreams: + - id: local1 + chain: ethereum + methods: + disabled: + - name: eth_call + connection: + ethereum-pos: + execution: + rpc: + url: "http://localhost" diff --git a/src/test/resources/configs/upstreams-methods-reload-initial.yaml b/src/test/resources/configs/upstreams-methods-reload-initial.yaml new file mode 100644 index 00000000..1967cc3d --- /dev/null +++ b/src/test/resources/configs/upstreams-methods-reload-initial.yaml @@ -0,0 +1,9 @@ +cluster: + upstreams: + - id: local1 + chain: ethereum + connection: + ethereum-pos: + execution: + rpc: + url: "http://localhost"