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"