Compare commits

..

3 Commits

Author SHA1 Message Date
Artem Rootman
302036c2c7 chore: bump emerald-grpc submodule (adds Mova chain refs) (#883)
Some checks failed
Tests / unit-test (push) Has been cancelled
Bumps emerald-grpc 5191731 -> 111ff26, i.e. p2p-org/emerald-grpc#165, which
adds CHAIN_MOVA__MAINNET = 1170 and CHAIN_MOVA__TESTNET = 10205 to the ChainRef
enum.

Describe.kt:38 and SubscribeChainStatus.kt:31 filter chains with
'Common.ChainRef.forNumber(it.id) != null'. mova's grpcId is 1170, which had no
enum entry, so the chain was dropped from the gRPC surface even though the
mova-nodecore upstream is healthy in all three public-multiregion regions.
Completes #882, which added the chain to the embedded chains.yaml.

Submodule range contains exactly that one commit (+2 lines in
proto/common.proto).
2026-07-27 18:01:40 +03:00
Artem Rootman
66d9e8078a chore: bump public submodule to 2120281 (adds Mova chain) (#882)
Bumps foundation/src/main/resources/public 0a070aac -> 2120281, pulling in
drpcorg/public#237, #238, #240, #241:

- Add Mova chain (MOVA_MAINNET 0xf1cc / grpcId 1170, MOVA_TESTNET 0x2853 /
  grpcId 10205) — the deployed dshackle image did not know the chain, so the
  mova-nodecore upstream in dshackle-public-multiregion was silently skipped
- aptos: fork-choice: height
- jovay: disable-log-index-validation
- tron: support-safe-block-tag: false (new optional schema key, ignored by
  ChainsConfigReader which reads keys explicitly)

Diff is additive only (+28 lines): no chain-id/grpcId changes to existing chains.
2026-07-27 13:16:54 +00:00
a10zn8
481c8f1721 Add bearer token authorization for upstreams (#878) 2026-07-03 17:54:11 +03:00
38 changed files with 224 additions and 624 deletions

View File

@@ -1 +0,0 @@
v0.79.10

View File

@@ -1,150 +0,0 @@
# 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.

View File

@@ -1,32 +0,0 @@
#!/bin/bash
# build-image.sh — build the StakeSquid dshackle fork image with a WIRE-CORRECT version.
#
# dRPC edges record each provider's dshackle version (Global.version, sent in every
# Describe response and SubscribeChainStatus BuildInfo) and deprioritize providers
# whose string doesn't match the expected release. gradle's getVersion() only yields
# the clean release string on an EXACT tag match; one patch commit past the tag would
# report "<next>-SNAPSHOT" to every edge. So: force-move the base release tag onto the
# patch HEAD locally before building (never push that tag), which makes our build
# report exactly what stock drpcorg/dshackle:<version> reports (verified against a
# live gateway: version.app=0.79.10).
#
# Usage: ./build-image.sh [base-tag] [suffix]
# base-tag: upstream release our branch is based on (default v0.79.10)
# suffix: local image tag suffix (default ss1) -> stakesquid/dshackle:<ver>-<suffix>
#
# Per upstream release: rebase reload-removal onto the new tag, bump the default here,
# rerun. Distribute with: docker save stakesquid/dshackle:<ver>-<suffix> | ssh <gw> docker load
set -euo pipefail
cd "$(dirname "$0")"
BASE_TAG=${1:-v0.79.10}
SUFFIX=${2:-ss1}
VER=${BASE_TAG#v}
git tag -f "$BASE_TAG" HEAD >/dev/null # local only - makes getVersion() exact-match
docker build -q -t drpc-dshackle . >/dev/null
export JAVA_HOME=${JAVA_HOME:-/usr/lib/jvm/java-21-openjdk-amd64}
./gradlew jibDockerBuild -Djib.to.image="stakesquid/dshackle:${VER}-${SUFFIX}" -x test
echo "--- wire version check (must equal stock ${VER}):"
docker run --rm --entrypoint cat "stakesquid/dshackle:${VER}-${SUFFIX}" \
/app/resources/version.properties

View File

@@ -849,6 +849,18 @@ rpc:
password: "${ETH_PASSWORD}"
----
| `rpc.bearer-auth` + `rpc.bearer-auth.token`
a| HTTP Bearer token authorization (`Authorization: Bearer <token>` header), if required by the remote server. +
Cannot be used together with `basic-auth`.
Value can also reference env variables, for example:
[source,yaml]
----
rpc:
url: "https://ethereum.com:8545"
bearer-auth:
token: "${ETH_TOKEN}"
----
| `ws.url`
| WebSocket URL to connect to.
Optional, but optimizes performance if it's available.
@@ -859,6 +871,9 @@ Optional, but optimizes performance if it's available.
| `ws.basic-auth` + ...
| WebSocket Basic Auth configuration, if required by the remote server
| `ws.bearer-auth` + `ws.bearer-auth.token`
| WebSocket Bearer token authorization, if required by the remote server
| `ws.frameSize`
| WebSocket frame size limit.
Ex `1kb`, `1024` (same as `1kb), `2mb`, etc.

View File

@@ -19,7 +19,6 @@ 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
@@ -47,7 +46,6 @@ 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
@@ -91,7 +89,6 @@ 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")

View File

@@ -33,6 +33,10 @@ class AuthConfig {
val password: String,
) : ClientAuth()
class ClientBearerAuth(
val token: String,
) : ClientAuth()
class ClientTlsAuth(
var ca: String? = null,
var certificate: String? = null,

View File

@@ -40,6 +40,18 @@ class AuthConfigReader : YamlConfigReader<AuthConfig>() {
}
}
fun readClientBearerAuth(node: MappingNode?): AuthConfig.ClientBearerAuth? {
return getMapping(node, "bearer-auth")?.let { authNode ->
val token = getValueAsString(authNode, "token")
if (token != null) {
AuthConfig.ClientBearerAuth(token)
} else {
log.warn("Bearer auth is not fully configured, token is required")
null
}
}
}
fun readClientTls(node: MappingNode?): AuthConfig.ClientTlsAuth? {
return getMapping(node, "tls")?.let { authNode ->
val auth = AuthConfig.ClientTlsAuth()

View File

@@ -44,7 +44,6 @@ 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

View File

@@ -83,18 +83,6 @@ 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
}
}

View File

@@ -1,9 +0,0 @@
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,
)

View File

@@ -136,12 +136,14 @@ data class UpstreamsConfig(
constructor(url: URI) : this(url, DEFAULT_MAX_CONNECTIONS, DEFAULT_QUEUE_SIZE)
var basicAuth: AuthConfig.ClientBasicAuth? = null
var bearerAuth: AuthConfig.ClientBearerAuth? = null
var tls: AuthConfig.ClientTlsAuth? = null
}
data class WsEndpoint(val url: URI) {
var origin: URI? = null
var basicAuth: AuthConfig.ClientBasicAuth? = null
var bearerAuth: AuthConfig.ClientBearerAuth? = null
var frameSize: Int? = null
var msgSize: Int? = null
var connections: Int = 1

View File

@@ -147,6 +147,7 @@ class UpstreamsConfigReader(
getValueAsString(node, "url")?.let { url ->
val http = UpstreamsConfig.HttpEndpoint(URI(url), DEFAULT_MAX_CONNECTIONS, DEFAULT_QUEUE_SIZE)
http.basicAuth = authConfigReader.readClientBasicAuth(node)
http.bearerAuth = readBearerAuth(node, http.basicAuth, url)
http.tls = authConfigReader.readClientTls(node)
connection.esplora = http
}
@@ -173,6 +174,19 @@ class UpstreamsConfigReader(
return connection
}
private fun readBearerAuth(
node: MappingNode?,
basicAuth: AuthConfig.ClientBasicAuth?,
url: String,
): AuthConfig.ClientBearerAuth? {
val bearerAuth = authConfigReader.readClientBearerAuth(node)
if (bearerAuth != null && basicAuth != null) {
log.warn("Both basic-auth and bearer-auth are configured for $url, basic-auth is used")
return null
}
return bearerAuth
}
private fun readRpcConfig(connConfigNode: MappingNode): UpstreamsConfig.HttpEndpoint? {
return getMapping(connConfigNode, "rpc")?.let { node ->
val maxConnections = getValueAsInt(node, "max-connections") ?: DEFAULT_MAX_CONNECTIONS
@@ -181,6 +195,7 @@ class UpstreamsConfigReader(
getValueAsString(node, "url")?.let { url ->
val http = UpstreamsConfig.HttpEndpoint(URI(url), maxConnections, queueSize)
http.basicAuth = authConfigReader.readClientBasicAuth(node)
http.bearerAuth = readBearerAuth(node, http.basicAuth, url)
http.tls = authConfigReader.readClientTls(node)
http
}
@@ -224,6 +239,7 @@ class UpstreamsConfigReader(
ws.origin = URI(origin)
}
ws.basicAuth = authConfigReader.readClientBasicAuth(node)
ws.bearerAuth = readBearerAuth(node, ws.basicAuth, url)
getValueAsBytes(node, "frameSize")?.let {
if (it < 65_535) {

View File

@@ -2,11 +2,8 @@ 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
@@ -19,12 +16,7 @@ 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()
@@ -43,29 +35,13 @@ 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)
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")
}
reloadConfigUpstreamService.reloadUpstreams(chainsToReload, upstreamsToRemove, upstreamsToAdd, newUpstreamsConfig)
return true
}
@@ -83,8 +59,8 @@ class UpstreamConfigReloadConfigProcessor(
}
val reloaded = mutableSetOf<Pair<String, Chain>>()
val removed = mutableSetOf<Pair<String, Chain>>()
val currentUpstreamsMap = upstreamKeyMap(currentUpstreams)
val newUpstreamsMap = upstreamKeyMap(newUpstreams)
val currentUpstreamsMap = currentUpstreams.associateBy { it.id!! to chainById(it.chain) }
val newUpstreamsMap = newUpstreams.associateBy { it.id!! to chainById(it.chain) }
currentUpstreamsMap.forEach {
val newUpstream = newUpstreamsMap[it.key]
@@ -103,19 +79,6 @@ class UpstreamConfigReloadConfigProcessor(
return UpstreamAnalyzeData(added, removed.plus(reloaded))
}
private fun upstreamKeyMap(
upstreams: List<UpstreamsConfig.Upstream<*>>,
): Map<Pair<String, Chain>, 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<ChainOptions.DefaultOptions>,
newDefaultOptions: List<ChainOptions.DefaultOptions>,
@@ -134,21 +97,13 @@ class UpstreamConfigReloadConfigProcessor(
currentOptions.forEach {
val newChainOption = newOptions[it.key]
if (newChainOption == null) {
val chain = chainById(it.key)
if (chain != Chain.UNSPECIFIED) {
removed.add(chain)
}
removed.add(chainById(it.key))
} else if (newChainOption != it.value) {
val chain = chainById(it.key)
if (chain != Chain.UNSPECIFIED) {
chainsToReload.add(chain)
}
chainsToReload.add(chainById(it.key))
}
}
val added = newOptions.minus(currentOptions.keys).mapNotNull {
chainById(it.key).takeIf { chain -> chain != Chain.UNSPECIFIED }
}
val added = newOptions.minus(currentOptions.keys).map { chainById(it.key) }
return chainsToReload.plus(added).plus(removed)
}

View File

@@ -7,38 +7,21 @@ 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<Chain>,
)
@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<Chain>,
upstreamsToRemove: Set<Pair<String, Chain>>,
upstreamsToAdd: Set<Pair<String, Chain>>,
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(
@@ -47,61 +30,38 @@ open class ReloadConfigUpstreamService(
),
)
val usedChains = removeUpstreams(safeChainsToReload, safeUpstreamsToRemove, grpcIdsToStop)
val usedChains = removeUpstreams(chainsToReload, upstreamsToRemove)
addUpstreams(newUpstreamsConfig, safeChainsToReload, upstreamsToAdd.map { it.first }.toSet())
addUpstreams(newUpstreamsConfig, chainsToReload, upstreamsToAdd.map { it.first }.toSet())
usedChains.filterValidChains().forEach { chain ->
if (newUpstreamsCount[chain] == null || newUpstreamsCount[chain] == 0L) {
usedChains.forEach { chain ->
if (newUpstreamsCount[chain] == null) {
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<Chain>,
upstreamsToRemove: Set<Pair<String, Chain>>,
grpcIdsToStop: Set<String>,
): Set<Chain> {
val usedChains = mutableSetOf<Chain>()
grpcIdsToStop.forEach { id ->
grpcUpstreamsRegistry.stop(id).forEach { event ->
usedChains.add(event.chain)
}
}
chainsToReload.forEach { chain ->
usedChains.add(chain)
val ms = multistreamHolder.getUpstream(chain)
chainsToReload.forEach {
usedChains.add(it)
val ms = multistreamHolder.getUpstream(it)
ms.getAll()
.toList()
.forEach { up ->
ms.processUpstreamsEventsSync(
UpstreamChangeEvent(chain, up, UpstreamChangeEvent.ChangeType.REMOVED),
)
ms.processUpstreamsEvents(UpstreamChangeEvent(it, up, UpstreamChangeEvent.ChangeType.REMOVED))
}
}
upstreamsToRemove.forEach { pair ->
if (grpcUpstreamsRegistry.isRegistered(pair.first)) {
return@forEach
}
usedChains.add(pair.second)
val ms = multistreamHolder.getUpstream(pair.second)
findUpstreamsToRemove(ms.getAll(), pair.first)
.forEach { up ->
ms.processUpstreamsEventsSync(
UpstreamChangeEvent(pair.second, up, UpstreamChangeEvent.ChangeType.REMOVED),
)
ms.getAll()
.find { pair.first == it.getId() }
?.let { up ->
ms.processUpstreamsEvents(UpstreamChangeEvent(pair.second, up, UpstreamChangeEvent.ChangeType.REMOVED))
}
}
@@ -121,22 +81,4 @@ open class ReloadConfigUpstreamService(
)
configuredUpstreams.processUpstreams(configToReload)
}
companion object {
internal fun findUpstreamsToRemove(upstreams: List<Upstream>, configId: String): List<Upstream> {
val prefix = "${configId}_"
return upstreams.filter { up ->
up.getId() == configId || up.getId().startsWith(prefix)
}
}
private fun Set<Chain>.filterValidChains(): Set<Chain> =
filter { it != Chain.UNSPECIFIED }.toSet()
private fun Set<Pair<String, Chain>>.filterValidPairs(): Set<Pair<String, Chain>> =
filter { it.second != Chain.UNSPECIFIED }.toSet()
private fun Collection<Chain>.filterValidChains(): List<Chain> =
filter { it != Chain.UNSPECIFIED }
}
}

View File

@@ -1,55 +0,0 @@
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<ServerCall<*, *>>()
override fun <ReqT : Any, RespT : Any> interceptCall(
call: ServerCall<ReqT, RespT>,
headers: Metadata,
next: ServerCallHandler<ReqT, RespT>,
): ServerCall.Listener<ReqT> {
activeCalls.add(call)
val trackedCall = object : ForwardingServerCall.SimpleForwardingServerCall<ReqT, RespT>(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)
}
}
}
}

View File

@@ -193,7 +193,7 @@ open class NativeCall(
if (it.isError()) {
it.error?.let { error ->
result.setErrorMessage(error.message)
.setItemErrorCode(error.itemErrorCode())
.setItemErrorCode(error.id)
error.data?.let { data ->
result.setErrorData(data)
@@ -750,8 +750,6 @@ open class NativeCall(
val errorAsIs: ByteArray? = null,
) {
fun itemErrorCode(): Int = upstreamError?.code ?: id
companion object {
private val log = LoggerFactory.getLogger(CallError::class.java)

View File

@@ -24,7 +24,6 @@ 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
@@ -36,7 +35,6 @@ 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)
@@ -76,11 +74,11 @@ open class ConfiguredUpstreams(
)
}
} else {
val grpcUpstreams = upstreamFactory.createGrpcUpstream(
upstreamFactory.createGrpcUpstream(
up as UpstreamsConfig.Upstream<UpstreamsConfig.GrpcConnection>,
chainsConfig,
)
val subscription = grpcUpstreams.start()
.start()
.doOnNext {
log.info("Chain ${it.chain} ${it.type} through gRPC at ${up.connection?.host}:${up.connection?.port}. With caps: ${it.upstream.getCapabilities()}")
}
@@ -88,7 +86,6 @@ open class ConfiguredUpstreams(
multistreamHolder.getUpstream(it.chain)
.processUpstreamsEvents(it)
}
grpcUpstreamsRegistry.register(up.id!!, grpcUpstreams, subscription)
}
}
}
@@ -98,10 +95,6 @@ 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

View File

@@ -52,7 +52,7 @@ class BitcoinUpstreamCreator(
fileResolver.resolve(ca).readBytes()
}
}
EsploraClient(endpoint.url, endpoint.basicAuth, tls)
EsploraClient(endpoint.url, endpoint.basicAuth, tls, endpoint.bearerAuth)
}
val extractBlock = ExtractBlock()

View File

@@ -84,6 +84,7 @@ open class GenericConnectorFactoryCreator(
monitoringCfg.nettyMetricsConfig.enabled,
httpScheduler,
customHeaders,
conn.bearerAuth,
)
}
}
@@ -106,6 +107,7 @@ open class GenericConnectorFactoryCreator(
).apply {
config = endpoint
basicAuth = endpoint.basicAuth
bearerAuth = endpoint.bearerAuth
this.customHeaders = customHeaders
}
val wsApi = WsConnectionPoolFactory(

View File

@@ -21,6 +21,7 @@ class BasicHttpFactory(
private val nettyMetricsEnabled: Boolean,
private val httpScheduler: Scheduler,
private val customHeaders: Map<String, String> = emptyMap(),
private val bearerAuth: AuthConfig.ClientBearerAuth? = null,
) : HttpFactory {
private val log = LoggerFactory.getLogger(this::class.java)
@@ -50,8 +51,8 @@ class BasicHttpFactory(
)
if (chain.type.apiType == ApiType.REST) {
return RestHttpReader(url, maxConnections, queueSize, metrics, httpScheduler, chain, basicAuth, tls, customHeaders)
return RestHttpReader(url, maxConnections, queueSize, metrics, httpScheduler, chain, basicAuth, tls, customHeaders, bearerAuth)
}
return JsonRpcHttpReader(url, maxConnections, queueSize, metrics, httpScheduler, basicAuth, tls, customHeaders)
return JsonRpcHttpReader(url, maxConnections, queueSize, metrics, httpScheduler, basicAuth, tls, customHeaders, bearerAuth)
}
}

View File

@@ -28,6 +28,7 @@ abstract class HttpReader(
basicAuth: AuthConfig.ClientBasicAuth? = null,
tlsCAAuth: ByteArray? = null,
customHeaders: Map<String, String> = emptyMap(),
bearerAuth: AuthConfig.ClientBearerAuth? = null,
) : ChainReader {
constructor() : this("", 1500, 1000, null)
@@ -66,6 +67,13 @@ abstract class HttpReader(
build = build.headers(headers)
}
if (basicAuth == null) {
bearerAuth?.let { auth ->
val headers = Consumer { h: HttpHeaders -> h.add(HttpHeaderNames.AUTHORIZATION, "Bearer ${auth.token}") }
build = build.headers(headers)
}
}
if (customHeaders.isNotEmpty()) {
val headers = Consumer { h: HttpHeaders ->
customHeaders.forEach { (key, value) ->

View File

@@ -444,14 +444,6 @@ 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) {

View File

@@ -34,10 +34,11 @@ import java.security.cert.X509Certificate
import java.util.Base64
import java.util.function.Consumer
class EsploraClient(
class EsploraClient @JvmOverloads constructor(
private val url: URI,
basicAuth: AuthConfig.ClientBasicAuth? = null,
tlsCAAuth: ByteArray? = null,
bearerAuth: AuthConfig.ClientBearerAuth? = null,
) {
companion object {
@@ -62,6 +63,13 @@ class EsploraClient(
build = build.headers(headers)
}
if (basicAuth == null) {
bearerAuth?.let { auth ->
val headers = Consumer { h: HttpHeaders -> h.add(HttpHeaderNames.AUTHORIZATION, "Bearer ${auth.token}") }
build = build.headers(headers)
}
}
tlsCAAuth?.let { auth ->
val cf = CertificateFactory.getInstance("X.509")
val cert = cf.generateCertificate(ByteArrayInputStream(auth)) as X509Certificate

View File

@@ -21,6 +21,7 @@ open class WsConnectionFactory(
) {
var basicAuth: AuthConfig.ClientBasicAuth? = null
var bearerAuth: AuthConfig.ClientBearerAuth? = null
var config: UpstreamsConfig.WsEndpoint? = null
var customHeaders: Map<String, String> = emptyMap()
@@ -47,7 +48,7 @@ open class WsConnectionFactory(
}
open fun createWsConnection(connIndex: Int = 0): WsConnection =
WsConnectionImpl(uri, origin, basicAuth, metrics(connIndex), scheduler, eventsScheduler, customHeaders).also { ws ->
WsConnectionImpl(uri, origin, basicAuth, metrics(connIndex), scheduler, eventsScheduler, customHeaders, bearerAuth).also { ws ->
config?.frameSize?.let {
ws.frameSize = it
}

View File

@@ -67,6 +67,7 @@ open class WsConnectionImpl(
private val scheduler: Scheduler,
private val eventsScheduler: Scheduler,
private val customHeaders: Map<String, String> = emptyMap(),
private val bearerAuth: AuthConfig.ClientBearerAuth? = null,
) : AutoCloseable, WsConnection, Cloneable {
companion object {
@@ -227,6 +228,11 @@ open class WsConnectionImpl(
val base64password = Base64.getEncoder().encodeToString(tmp.toByteArray())
headers.add(HttpHeaderNames.AUTHORIZATION, "Basic $base64password")
}
if (basicAuth == null) {
bearerAuth?.let { auth ->
headers.add(HttpHeaderNames.AUTHORIZATION, "Bearer ${auth.token}")
}
}
customHeaders.forEach { (key, value) ->
headers.add(key, value)
}

View File

@@ -43,7 +43,6 @@ 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
@@ -65,7 +64,6 @@ 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
@@ -96,10 +94,8 @@ class GrpcUpstreams(
var timeout = Defaults.timeout
lateinit var client: ReactorBlockchainGrpc.ReactorBlockchainStub
private var channel: ManagedChannel? = null
private val known = HashMap<Chain, DefaultUpstream>()
private val lock = ReentrantLock()
private val statusSubscriptions = mutableMapOf<Chain, Disposable>()
fun start(): Flux<UpstreamChangeEvent> {
val chanelBuilder = NettyChannelBuilder.forAddress(host, port)
@@ -126,9 +122,8 @@ class GrpcUpstreams(
chanelBuilder.usePlaintext()
}
val builtChannel = chanelBuilder.build()
this.channel = builtChannel
var client = ReactorBlockchainGrpc.newReactorStub(builtChannel)
val channel = chanelBuilder.build()
var client = ReactorBlockchainGrpc.newReactorStub(channel)
if (compression) {
client = client.withCompression(Codec.Gzip().messageEncoding)
}
@@ -137,7 +132,7 @@ class GrpcUpstreams(
val grpcUpstreamsAuth =
if (tokenAuth != null && authorizationConfig.enabled) {
GrpcUpstreamsAuth(
ReactorAuthGrpc.newReactorStub(builtChannel),
ReactorAuthGrpc.newReactorStub(channel),
authorizationConfig,
grpcAuthContext,
tokenAuth.publicKeyPath!!,
@@ -146,6 +141,8 @@ class GrpcUpstreams(
null
}
val statusSubscriptions = mutableMapOf<Chain, Disposable>()
return Flux.interval(Duration.ZERO, Duration.ofSeconds(20))
.flatMap {
authAndDescribe(grpcUpstreamsAuth)
@@ -359,31 +356,4 @@ 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<UpstreamChangeEvent> {
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()
}
}
}
}

View File

@@ -1,55 +0,0 @@
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<String, Entry>()
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<UpstreamChangeEvent> {
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<String>): 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()
}
}

View File

@@ -36,7 +36,8 @@ class RestHttpReader(
basicAuth: AuthConfig.ClientBasicAuth? = null,
tlsCAAuth: ByteArray? = null,
customHeaders: Map<String, String> = emptyMap(),
) : HttpReader(target, maxConnections, queueSize, metrics, basicAuth, tlsCAAuth, customHeaders) {
bearerAuth: AuthConfig.ClientBearerAuth? = null,
) : HttpReader(target, maxConnections, queueSize, metrics, basicAuth, tlsCAAuth, customHeaders, bearerAuth) {
private val parser = ResponseRpcParser()
private val requestParser = RestRequestParser

View File

@@ -36,7 +36,7 @@ import java.util.function.Function
/**
* JSON RPC client
*/
class JsonRpcHttpReader(
class JsonRpcHttpReader @JvmOverloads constructor(
target: String,
maxConnections: Int,
queueSize: Int,
@@ -45,7 +45,8 @@ class JsonRpcHttpReader(
basicAuth: AuthConfig.ClientBasicAuth? = null,
tlsCAAuth: ByteArray? = null,
customHeaders: Map<String, String> = emptyMap(),
) : HttpReader(target, maxConnections, queueSize, metrics, basicAuth, tlsCAAuth, customHeaders) {
bearerAuth: AuthConfig.ClientBearerAuth? = null,
) : HttpReader(target, maxConnections, queueSize, metrics, basicAuth, tlsCAAuth, customHeaders, bearerAuth) {
private val parser = ResponseRpcParser()
private val streamParser = JsonRpcStreamParser()

View File

@@ -35,6 +35,29 @@ class AuthConfigReaderSpec extends Specification {
act.password == "258fe4149c199ad8f2811a68f20154fc"
}
def "Read bearer-auth for client"() {
setup:
def yaml =
"bearer-auth:\n" +
" token: 9c199ad8f281f20154fc258fe41a6814"
when:
def act = reader.readClientBearerAuth(reader.readNode(yaml))
then:
act != null
act.token == "9c199ad8f281f20154fc258fe41a6814"
}
def "Read bearer-auth without token as null"() {
setup:
def yaml =
"bearer-auth:\n" +
" something: else"
when:
def act = reader.readClientBearerAuth(reader.readNode(yaml))
then:
act == null
}
def "Read tls for client"() {
setup:
def yaml =

View File

@@ -75,6 +75,38 @@ class UpstreamsConfigReaderSpec extends Specification {
}
}
def "Parse config with bearer auth"() {
setup:
def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-bearer-auth.yaml")
when:
def act = reader.readInternal(config)
then:
act != null
act.upstreams.size() == 2
with(act.upstreams.get(0)) {
id == "local"
connection instanceof UpstreamsConfig.RpcConnection
with((UpstreamsConfig.RpcConnection) connection) {
rpc.basicAuth == null
rpc.bearerAuth != null
rpc.bearerAuth.token == "9c199ad8f281f20154fc258fe41a6814"
ws != null
ws.basicAuth == null
ws.bearerAuth != null
ws.bearerAuth.token == "258fe4149c199ad8f2811a68f20154fc"
}
}
with(act.upstreams.get(1)) {
id == "both-auth"
connection instanceof UpstreamsConfig.RpcConnection
with((UpstreamsConfig.RpcConnection) connection) {
// when both are configured basic-auth wins and bearer-auth is dropped
rpc.basicAuth != null
rpc.bearerAuth == null
}
}
}
def "Parse websocket-only config"() {
setup:
def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-ws-only.yaml")

View File

@@ -16,6 +16,7 @@
package io.emeraldpay.dshackle.upstream.rpcclient
import io.emeraldpay.dshackle.config.AuthConfig
import io.emeraldpay.dshackle.test.TestingCommons
import io.emeraldpay.dshackle.upstream.ChainException
import io.emeraldpay.dshackle.upstream.ChainRequest
@@ -69,6 +70,28 @@ class JsonRpcHttpReaderSpec extends Specification {
new String(act.result) == '"0x98de45"'
}
def "Make a request with bearer auth"() {
setup:
def bearerAuth = new AuthConfig.ClientBearerAuth("test-token-123")
JsonRpcHttpReader client = new JsonRpcHttpReader("localhost:${mockServer.port}", 50, 50, metrics, Schedulers.boundedElastic(), null, null, [:], bearerAuth)
def resp = '{' +
' "jsonrpc": "2.0",' +
' "result": "0x98de45",' +
' "error": null,' +
' "id": 15' +
'}'
mockServer.when(
HttpRequest.request().withHeader("authorization", "Bearer test-token-123")
).respond(
HttpResponse.response(resp)
)
when:
def act = client.read(new ChainRequest("test", new ListParams())).block(Duration.ofSeconds(5))
then:
act.error == null
new String(act.result) == '"0x98de45"'
}
def "Produces RPC Exception on error status code"() {
setup:
def client = new JsonRpcHttpReader("localhost:${mockServer.port}", 50, 50, metrics, Schedulers.boundedElastic(), null, null, [:])

View File

@@ -10,7 +10,6 @@ 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
@@ -19,7 +18,6 @@ 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
@@ -48,8 +46,6 @@ class ReloadConfigTest {
private val config = mock<Config>()
private val reloadConfigService = ReloadConfigService(config, fileResolver, mainConfig)
private val configuredUpstreams = mock<ConfiguredUpstreams>()
private val grpcUpstreamsRegistry = mock<GrpcUpstreamsRegistry>()
private val grpcClientRefreshService = mock<GrpcClientRefreshService>()
@BeforeEach
fun setupTests() {
@@ -92,14 +88,8 @@ 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()
@@ -109,8 +99,8 @@ class ReloadConfigTest {
reloadConfig.handle(Signal("HUP"))
verify(msEth).processUpstreamsEventsSync(UpstreamChangeEvent(ETHEREUM__MAINNET, up1, UpstreamChangeEvent.ChangeType.REMOVED))
verify(msPoly).processUpstreamsEventsSync(UpstreamChangeEvent(POLYGON__MAINNET, up3, UpstreamChangeEvent.ChangeType.REMOVED))
verify(msEth).processUpstreamsEvents(UpstreamChangeEvent(ETHEREUM__MAINNET, up1, UpstreamChangeEvent.ChangeType.REMOVED))
verify(msPoly).processUpstreamsEvents(UpstreamChangeEvent(POLYGON__MAINNET, up3, UpstreamChangeEvent.ChangeType.REMOVED))
verify(configuredUpstreams).processUpstreams(
UpstreamsConfig(
newConfig.defaultOptions,
@@ -150,14 +140,8 @@ 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)!!
@@ -188,12 +172,7 @@ class ReloadConfigTest {
val reloadConfigUpstreamService = mock<ReloadConfigUpstreamService>()
val upstreamCfgReloadProcessor = UpstreamConfigReloadConfigProcessor(
reloadConfigService,
reloadConfigUpstreamService,
mainConfig,
grpcClientRefreshService,
)
val upstreamCfgReloadProcessor = UpstreamConfigReloadConfigProcessor(reloadConfigService, reloadConfigUpstreamService)
val reloadConfig = ReloadConfigSetup(listOf(upstreamCfgReloadProcessor))
whenever(config.getConfigPath()).thenReturn(initialConfigFile)
@@ -205,78 +184,6 @@ 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<ReloadConfigUpstreamService>()
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<Multistream> {
on { getAll() } doReturn emptyList()
}
val currentMultistreamHolder = mock<CurrentMultistreamHolder> {
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(

View File

@@ -0,0 +1,27 @@
version: v1
upstreams:
- id: local
chain: ethereum
connection:
ethereum:
rpc:
url: "http://localhost:8545"
bearer-auth:
token: 9c199ad8f281f20154fc258fe41a6814
ws:
url: "ws://localhost:8546"
origin: "http://localhost"
bearer-auth:
token: 258fe4149c199ad8f2811a68f20154fc
- id: both-auth
chain: ethereum
connection:
ethereum:
rpc:
url: "http://localhost:8545"
basic-auth:
username: 4fc258fe41a68149c199ad8f281f2015
password: 1a68f20154fc258fe4149c199ad8f281
bearer-auth:
token: 9c199ad8f281f20154fc258fe41a6814

View File

@@ -1,12 +0,0 @@
cluster:
upstreams:
- id: local1
chain: ethereum
methods:
disabled:
- name: eth_call
connection:
ethereum-pos:
execution:
rpc:
url: "http://localhost"

View File

@@ -1,9 +0,0 @@
cluster:
upstreams:
- id: local1
chain: ethereum
connection:
ethereum-pos:
execution:
rpc:
url: "http://localhost"