Compare commits

..

14 Commits

Author SHA1 Message Date
rob
bcc71c6ca0 build-image.sh: hard gates - submodule init (https rewrite), foundation republish to mavenLocal, clean build, in-image resource check (ss1 shipped without compatible-clients.yaml and crash-looped the canary) 2026-07-30 07:49:13 +00:00
rob
7f2cc0e739 rebase onto v0.79.11 (bearer-token auth + registry bump; published 2026-07-27) 2026-07-30 07:14:22 +00:00
rob
89b331bdf3 track upstream base tag for the release watcher 2026-07-30 07:14:22 +00:00
rob
0cbe8ba2e8 build-image.sh: wire-correct version builds (tag-move so getVersion() exact-matches; dRPC edges track provider version strings) 2026-07-30 07:14:22 +00:00
rob
a287e74bd7 fix: SIGHUP reload supports upstream and method removal
Reload no longer crashes on Chain.UNSPECIFIED, removes gRPC upstreams
by prefixed id with lifecycle cleanup, applies methods.disabled changes
synchronously, pushes status via existing SubscribeChainStatus streams,
and optionally closes client RPCs so edges reconnect. NativeCall returns
-32601 (CODE_METHOD_NOT_EXIST) for disabled/unknown methods.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-30 07:14:22 +00: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
KirillPamPam
f2c1bd0628 Update deps (#877) 2026-06-30 16:59:58 +04:00
a10zn8
2b79458bf2 Add node_getChainTips fallback for Aztec v5 (#873) 2026-06-25 14:07:55 +03:00
Vadim Filin
634291bc98 Fix hl (#869)
* fix(hyperliquid): derive native-tx routing labels from the ?hl= URL flag

The include_hl_native_tx/exclude_hl_native_tx detector classified a node by scanning the last 300 blocks for a system (native) topup tx from a per-chain address. That cannot work on testnet:

- the configured testnet address 0x6ed35e7d6de4b45f4efb8a91eff31afa49362569 was a regular bot (non-zero gasPrice, not filtered by hl-compliant mode, present on both node types);
- real testnet system txs (from 0x2222...) are far too sparse and bursty (median gap ~1200 blocks, max ~9000 = ~2.5h at ~1s/block) for any practical window;
- eth_getLogs is identical between compliant and non-compliant modes, so there is no cheap wide-range signal either.

Our hl-node upstreams already encode the mode in the URL (?hl=false serves native txs, ?hl=true is compliant). Read that flag directly:
- GenericUpstream captures the configured RPC/WS URL and exposes getRpcConnectionUrl();
- detectHlNativeTx emits include/exclude_hl_native_tx straight from ?hl= when present (cheap, exact, drift-free), for both mainnet and testnet;
- it falls back to the recent-blocks scan only when there is no ?hl= flag, and only on mainnet (testnet without the flag is too sparse to classify);
- the bogus HL_NATIVE_TX_FROM_TESTNET constant is removed.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>

* Fix hl

---------

Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-05 13:41:41 +02:00
Vadim Filin
2ddd9c8889 Update emerald-grpc and foundation submodule references (#866)
* add Humanity chain (grpcId 1167)

Bump emerald-grpc and foundation public submodules to include
CHAIN_HUMANITY__MAINNET. Chain.kt is generated at build time.

* add Humanity testnet (grpcId 10202)

bump emerald-grpc and foundation public submodules

---------

Co-authored-by: Claude <noreply@anthropic.com>
2026-06-02 14:47:19 +02:00
msizov
6dd2ee566f add kite mainnet, robinhood mainnet (#865) 2026-05-29 16:19:47 +05:00
KirillPamPam
998e1ac933 newPedningTransactions validator (#864) 2026-05-26 17:13:31 +04:00
a10zn8
f93f350a29 Add method label to upstream.rpc.conn metric (#863) 2026-05-26 12:11:13 +03:00
59 changed files with 1499 additions and 77 deletions

1
.upstream-base Normal file
View File

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

150
FINDINGS.md Normal file
View File

@@ -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.

61
build-image.sh Executable file
View File

@@ -0,0 +1,61 @@
#!/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
# Submodules are REQUIRED at build time (SSH urls -> rewrite to https; controller
# has no github key). The `public` submodule feeds foundation's resources
# (compatible-clients.yaml etc.) - foundation is consumed as a mavenLocal
# artifact, so it must be REPUBLISHED whenever the submodule moves, and the
# main build cleaned so jib can't reuse stale layers. Skipping any of this
# shipped a boot-crashing image on 2026-07-30 (ss1: Spring died on missing
# compatible-clients.yaml) - hence the hard gate below.
git config url."https://github.com/".insteadOf "git@github.com:" 2>/dev/null || true
git submodule sync -q && git submodule update --init --recursive
test -f foundation/src/main/resources/public/compatible-clients.yaml \
|| { echo "FATAL: public submodule missing compatible-clients.yaml"; exit 1; }
export JAVA_HOME=${JAVA_HOME:-/usr/lib/jvm/java-21-openjdk-amd64}
./gradlew -p foundation publishToMavenLocal -q
./gradlew clean -q
docker build -q -t drpc-dshackle . >/dev/null
./gradlew jibDockerBuild -Djib.to.image="stakesquid/dshackle:${VER}-${SUFFIX}" -x test
echo "--- resource gate (foundation jar inside the image):"
docker run --rm --entrypoint sh "stakesquid/dshackle:${VER}-${SUFFIX}" -c \
'ls /app/libs/foundation-*.jar >/dev/null' || { echo FATAL: no foundation jar; exit 1; }
python3 - "$VER" "$SUFFIX" <<'PY'
import subprocess, sys, zipfile, io
ver, suf = sys.argv[1], sys.argv[2]
jar = subprocess.run(["docker", "run", "--rm", "--entrypoint", "sh",
f"stakesquid/dshackle:{ver}-{suf}", "-c",
"cat /app/libs/foundation-*.jar"], capture_output=True).stdout
names = zipfile.ZipFile(io.BytesIO(jar)).namelist()
assert "public/compatible-clients.yaml" in names, "FATAL: compatible-clients.yaml not packaged"
print("resource gate OK: public/compatible-clients.yaml present in foundation jar")
PY
echo "--- wire version check (must equal stock ${VER}):"
docker run --rm --entrypoint cat "stakesquid/dshackle:${VER}-${SUFFIX}" \
/app/resources/version.properties

View File

@@ -142,6 +142,7 @@ open class CodeGen(private val config: ChainsConfig) {
"kadena" -> "BlockchainType.KADENA"
"avm" -> "BlockchainType.AVM"
"app" -> "BlockchainType.ETHEREUM"
"aptos" -> "BlockchainType.ETHEREUM"
else -> throw IllegalArgumentException("unknown blockchain type $type")
}
}

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

@@ -20,6 +20,7 @@ class ChainOptions {
val disableBoundValidation: Boolean = false,
val valdateErigonBug: Boolean,
val disableLogIndexValidation: Boolean = false,
val disablePendingTxValidation: Boolean = false,
)
data class DefaultOptions(
@@ -43,7 +44,8 @@ class ChainOptions {
var disableLivenessSubscriptionValidation: Boolean? = null,
var disableBoundValidation: Boolean? = null,
var validateErigonBug: Boolean? = null,
var disableLogIndexValidation: Boolean? = null
var disableLogIndexValidation: Boolean? = null,
var disablePendingTxValidation: Boolean? = null
) {
companion object {
@JvmStatic
@@ -76,6 +78,7 @@ class ChainOptions {
copy.disableBoundValidation = overwrites.disableBoundValidation ?: this.disableBoundValidation
copy.validateErigonBug = overwrites.validateErigonBug ?: this.validateErigonBug
copy.disableLogIndexValidation = overwrites.disableLogIndexValidation ?: this.disableLogIndexValidation
copy.disablePendingTxValidation = overwrites.disablePendingTxValidation ?: this.disablePendingTxValidation
return copy
}
@@ -97,6 +100,7 @@ class ChainOptions {
this.disableBoundValidation ?: false,
this.validateErigonBug ?: true,
this.disableLogIndexValidation ?: false,
this.disablePendingTxValidation ?: false,
)
}
}

View File

@@ -61,6 +61,9 @@ class ChainOptionsReader : YamlConfigReader<ChainOptions.PartialOptions>() {
getValueAsBool(values, "disable-log-index-validation")?.let {
options.disableLogIndexValidation = it
}
getValueAsBool(values, "disable-pending-tx-validation")?.let {
options.disablePendingTxValidation = it
}
return options
}
}

View File

@@ -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")

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,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

View File

@@ -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
}
}

View File

@@ -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,
)

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,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<Pair<String, Chain>>()
val removed = mutableSetOf<Pair<String, Chain>>()
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<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>,
@@ -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)
}

View File

@@ -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<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(
@@ -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<Chain>,
upstreamsToRemove: Set<Pair<String, Chain>>,
grpcIdsToStop: Set<String>,
): Set<Chain> {
val usedChains = mutableSetOf<Chain>()
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<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

@@ -32,11 +32,12 @@ class ChainEventMapper {
}
fun mapCapabilities(capabilities: Collection<Capability>): BlockchainOuterClass.ChainEvent {
val caps = capabilities.map {
val caps = capabilities.filter { it != Capability.WS_PENDING_TX }.map {
when (it) {
Capability.RPC -> BlockchainOuterClass.Capabilities.CAP_CALLS
Capability.BALANCE -> BlockchainOuterClass.Capabilities.CAP_BALANCE
Capability.WS_HEAD -> BlockchainOuterClass.Capabilities.CAP_WS_HEAD
else -> null
}
}

View File

@@ -63,11 +63,12 @@ class Describe(
}
capabilities.addAll(chainUpstreams.getCapabilities())
chainDescription.addAllCapabilities(
capabilities.map {
capabilities.filter { it != Capability.WS_PENDING_TX }.map {
when (it) {
Capability.RPC -> BlockchainOuterClass.Capabilities.CAP_CALLS
Capability.BALANCE -> BlockchainOuterClass.Capabilities.CAP_BALANCE
Capability.WS_HEAD -> BlockchainOuterClass.Capabilities.CAP_WS_HEAD
else -> null
}
},
)

View File

@@ -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<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.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)

View File

@@ -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<UpstreamsConfig.GrpcConnection>,
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

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)
@@ -34,11 +35,14 @@ class BasicHttpFactory(
Tag.of("chain", chain.chainCode),
)
val metrics = RequestMetrics(
Timer.builder("upstream.rpc.conn")
.description("Request time through a HTTP JSON RPC connection")
.tags(metricsTags)
.publishPercentileHistogram()
.register(Metrics.globalRegistry),
{ method ->
Timer.builder("upstream.rpc.conn")
.description("Request time through a HTTP JSON RPC connection")
.tags(metricsTags)
.tag("method", method ?: "unknown")
.publishPercentileHistogram()
.register(Metrics.globalRegistry)
},
Counter.builder("upstream.rpc.fail")
.description("Number of failures of HTTP JSON RPC requests")
.tags(metricsTags)
@@ -47,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

@@ -4,4 +4,5 @@ enum class Capability {
RPC,
BALANCE,
WS_HEAD,
WS_PENDING_TX,
}

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) ->
@@ -101,7 +109,7 @@ abstract class HttpReader(
open fun onStop() {
if (metrics != null) {
Metrics.globalRegistry.remove(metrics.timer)
metrics.registeredTimers().forEach { Metrics.globalRegistry.remove(it) }
Metrics.globalRegistry.remove(metrics.fails)
}
}

View File

@@ -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) {

View File

@@ -17,9 +17,20 @@ package io.emeraldpay.dshackle.upstream
import io.micrometer.core.instrument.Counter
import io.micrometer.core.instrument.Timer
import java.util.concurrent.ConcurrentHashMap
class RequestMetrics(
val timer: Timer,
private val timerFactory: (String?) -> Timer,
val fails: Counter,
val nettyMetricsEnabled: Boolean,
)
) {
private val timerCache = ConcurrentHashMap<String, Timer>()
constructor(timer: Timer, fails: Counter, nettyMetricsEnabled: Boolean) :
this({ _ -> timer }, fails, nettyMetricsEnabled)
fun timer(method: String? = null): Timer =
timerCache.computeIfAbsent(method ?: "") { timerFactory(method) }
fun registeredTimers(): Collection<Timer> = timerCache.values
}

View File

@@ -1,6 +1,7 @@
package io.emeraldpay.dshackle.upstream.aztec
import com.fasterxml.jackson.databind.JsonNode
import com.github.benmanes.caffeine.cache.Caffeine
import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.config.ChainsConfig.ChainConfig
@@ -8,6 +9,7 @@ import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.foundation.ChainOptions.Options
import io.emeraldpay.dshackle.reader.ChainReader
import io.emeraldpay.dshackle.upstream.ChainException
import io.emeraldpay.dshackle.upstream.ChainRequest
import io.emeraldpay.dshackle.upstream.GenericSingleCallValidator
import io.emeraldpay.dshackle.upstream.SingleValidator
@@ -21,16 +23,36 @@ import io.emeraldpay.dshackle.upstream.rpcclient.ListParams
import org.slf4j.LoggerFactory
import reactor.core.publisher.Mono
import java.math.BigInteger
import java.time.Duration
import java.time.Instant
object AztecChainSpecific : AbstractPollChainSpecific() {
private val log = LoggerFactory.getLogger(AztecChainSpecific::class.java)
// node_getL2Tips reshaped between Aztec versions:
private const val METHOD_NOT_FOUND = -32601
// Aztec v5 (v5.0.0-rc.1) renamed the tips RPC: node_getL2Tips -> node_getChainTips.
// Older nodes (incl. current mainnet) expose only the legacy name; v5+ nodes expose
// only the new one. We probe the legacy method first (most upstreams are still on it)
// and fall back to the new one on "method not found", remembering the working method
// per upstream so we stop probing the dead one on every poll.
private val LEGACY_TIPS_REQUEST = ChainRequest("node_getL2Tips", ListParams())
private val CHAIN_TIPS_REQUEST = ChainRequest("node_getChainTips", ListParams())
// Bounded so per-upstream entries can't accumulate without limit (e.g. across config
// reloads): a hard size cap plus idle expiry evict stale ids, and an evicted entry just
// costs one re-probe. Only upstreams that actually fall back take a slot — legacy-only
// ones keep using the default and never populate it.
private val workingTipsRequest = Caffeine.newBuilder()
.maximumSize(1024)
.expireAfterAccess(Duration.ofHours(1))
.build<String, ChainRequest>()
// The tips response reshaped between Aztec versions:
// v3 (and earlier): {proposed: {number, hash}, proven: {number, hash}, checkpointed: {number, hash}}
// v4: proven/finalized/checkpointed each became {block: {number, hash}, checkpoint: {number, hash}}
// proposed stayed flat. We always look at the v4 nested path first and fall back
// to the flat v3 path so an upstream on either version is parsed correctly.
// v4/v5: proven/finalized/checkpointed each became {block: {number, hash}, checkpoint: {number, hash}}
// proposed stayed flat across all versions. We look at the flat path first and fall back
// to the nested path so an upstream on any version is parsed correctly.
private val PROPOSED_NUMBER = arrayOf("proposed.number", "proposed.block.number")
private val PROPOSED_HASH = arrayOf("proposed.hash", "proposed.block.hash")
@@ -122,8 +144,43 @@ object AztecChainSpecific : AbstractPollChainSpecific() {
return AztecLowerBoundService(chain, upstream)
}
override fun latestBlockRequest(): ChainRequest =
ChainRequest("node_getL2Tips", ListParams())
override fun latestBlockRequest(): ChainRequest = LEGACY_TIPS_REQUEST
// Try the per-upstream remembered method (legacy by default); on "method not found"
// fall back to the other one and remember whichever succeeds, so subsequent polls go
// straight to the working method. Any other error propagates as before.
override fun getLatestBlock(api: ChainReader, upstreamId: String): Mono<BlockContainer> {
val preferred = workingTipsRequest.getIfPresent(upstreamId) ?: LEGACY_TIPS_REQUEST
val fallback = if (preferred === LEGACY_TIPS_REQUEST) CHAIN_TIPS_REQUEST else LEGACY_TIPS_REQUEST
return fetchTips(api, upstreamId, preferred)
.onErrorResume { err ->
if (isMethodNotFound(err)) {
log.info(
"Aztec upstream {} does not support {}, falling back to {}",
upstreamId,
preferred.method,
fallback.method,
)
fetchTips(api, upstreamId, fallback)
.doOnNext { workingTipsRequest.put(upstreamId, fallback) }
} else {
Mono.error(err)
}
}
}
private fun fetchTips(api: ChainReader, upstreamId: String, request: ChainRequest): Mono<BlockContainer> {
return api.read(request).flatMap {
parseBlock(it.getResult(), upstreamId, api)
}
}
private fun isMethodNotFound(err: Throwable): Boolean {
if (err is ChainException && err.error.code == METHOD_NOT_FOUND) {
return true
}
return err.message?.contains("method not found", ignoreCase = true) ?: false
}
override fun upstreamSettingsDetector(
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

@@ -16,18 +16,24 @@ class DefaultAztecMethods : CallMethods {
"node_getBlockNumber",
"node_getProvenBlockNumber",
"node_getL2Tips",
"node_getChainTips",
"node_getBlock",
"node_getBlocks",
"node_getBlockData",
"node_getBlockHeader",
"node_getBlockByArchive",
"node_getBlockByHash",
"node_getBlockHeaderByArchive",
// Checkpoints
// Checkpoints / consensus
"node_getCheckpointNumber",
"node_getCheckpoint",
"node_getCheckpointedBlockNumber",
"node_getCheckpointedBlocks",
"node_getCheckpoints",
"node_getCheckpointsData",
"node_getCheckpointAttestationsForSlot",
"node_getProposalsForSlot",
// Transactions
"node_sendTx",
@@ -45,6 +51,11 @@ class DefaultAztecMethods : CallMethods {
"node_getWorldStateSyncStatus",
"node_findLeavesIndexes",
// Sync status
"node_getSyncedL1Timestamp",
"node_getSyncedL2EpochNumber",
"node_getSyncedL2SlotNumber",
// Sibling paths
"node_getNullifierSiblingPath",
"node_getNoteHashSiblingPath",
@@ -65,12 +76,14 @@ class DefaultAztecMethods : CallMethods {
"node_getL1ToL2MessageCheckpoint",
"node_isL1ToL2MessageSynced",
"node_getL2ToL1Messages",
"node_getL2ToL1MembershipWitness",
// Logs
"node_getPrivateLogs",
"node_getPrivateLogsByTags",
"node_getPublicLogs",
"node_getPublicLogsByTagsFromContract",
"node_getPublicLogsByTags",
"node_getContractClassLogs",
"node_getLogsByTags",
@@ -86,12 +99,14 @@ class DefaultAztecMethods : CallMethods {
"node_getVersion",
"node_getChainId",
"node_getL1ContractAddresses",
"node_getL1Constants",
"node_getProtocolContractAddresses",
"node_getEncodedEnr",
// Fees
"node_getCurrentBaseFees",
"node_getCurrentMinFees",
"node_getPredictedMinFees",
"node_getMaxPriorityFees",
// Validators

View File

@@ -11,6 +11,7 @@ import io.emeraldpay.dshackle.upstream.ethereum.domain.Address
import io.emeraldpay.dshackle.upstream.ethereum.hex.Hex32
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectLogs
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectNewHeads
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.NoPendingTxes
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource
import io.emeraldpay.dshackle.upstream.rpcclient.ListParams
import org.slf4j.LoggerFactory
@@ -72,7 +73,7 @@ open class EthereumEgressSubscription(
} else {
listOf()
}
return if (pendingTxesSource != null) {
return if (pendingTxesSource != null && pendingTxesSource !is NoPendingTxes && upstream.getCapabilities().contains(Capability.WS_PENDING_TX)) {
subs.plus(listOf(METHOD_PENDING_TXES, METHOD_DRPC_PENDING_TXES))
} else {
subs

View File

@@ -8,6 +8,7 @@ import io.emeraldpay.dshackle.upstream.ChainRequest
import io.emeraldpay.dshackle.upstream.ChainResponse
import io.emeraldpay.dshackle.upstream.NodeTypeRequest
import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.generic.GenericUpstream
import io.emeraldpay.dshackle.upstream.rpcclient.ListParams
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
@@ -15,7 +16,6 @@ import java.util.concurrent.atomic.AtomicInteger
const val ZERO_ADDRESS = "0x0000000000000000000000000000000000000000"
const val HL_NATIVE_TX_FROM_MAINNET = "0x2222222222222222222222222222222222222222"
const val HL_NATIVE_TX_FROM_TESTNET = "0x6ed35e7d6de4b45f4efb8a91eff31afa49362569"
class EthereumUpstreamSettingsDetector(
private val _upstream: Upstream,
@@ -179,15 +179,17 @@ class EthereumUpstreamSettingsDetector(
Some clients on hyperliquid don't include system topup transactions, set either one of labels
*/
private fun detectHlNativeTx(): Flux<Pair<String, String>> {
// Only run HL native tx detection on Hyperliquid chains
if (chain != Chain.HYPERLIQUID__MAINNET && chain != Chain.HYPERLIQUID__TESTNET) {
return Flux.empty()
}
val hlNativeTxFrom = when (chain) {
Chain.HYPERLIQUID__MAINNET -> HL_NATIVE_TX_FROM_MAINNET
Chain.HYPERLIQUID__TESTNET -> HL_NATIVE_TX_FROM_TESTNET
else -> return Flux.empty()
// prefer the explicit ?hl= flag from the upstream URL
hlNativeTxLabelsFromUrl()?.let { return Flux.fromIterable(it) }
// no ?hl= flag: block-scan fallback (reliable only on mainnet)
if (chain != Chain.HYPERLIQUID__MAINNET) {
return Flux.empty()
}
val hlNativeTxFrom = HL_NATIVE_TX_FROM_MAINNET
if (detectCounter.get() % 5 != 1) {
return Flux.empty() // reduce frequency of detection
}
@@ -255,6 +257,19 @@ class EthereumUpstreamSettingsDetector(
}
}
// maps the upstream URL's ?hl= flag to routing labels, or null if absent
private fun hlNativeTxLabelsFromUrl(): List<Pair<String, String>>? {
val url = (upstream as? GenericUpstream)?.getRpcConnectionUrl()?.toString() ?: return null
// ?hl=false => serves native txs (include); ?hl=true => hl-node compliant (exclude)
return when {
url.contains(Regex("[?&]hl=false\\b")) ->
listOf("include_hl_native_tx" to "true", "exclude_hl_native_tx" to "false")
url.contains(Regex("[?&]hl=true\\b")) ->
listOf("exclude_hl_native_tx" to "true", "include_hl_native_tx" to "false")
else -> null
}
}
private fun detectArchiveNode(notArchived: Boolean): Mono<Pair<String, String>> {
if (notArchived) {
return Mono.empty()

View File

@@ -0,0 +1,66 @@
package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.reader.ChainReader
import io.emeraldpay.dshackle.upstream.ChainRequest
import io.emeraldpay.dshackle.upstream.ChainResponse
import io.emeraldpay.dshackle.upstream.rpcclient.ListParams
import org.slf4j.LoggerFactory
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
import java.time.Duration
const val BASE_TX_LIMIT = 1000L
interface PendingTransactionValidator {
fun pendingTxExists(): Flux<Boolean>
}
class NoopPendingTransactionValidator : PendingTransactionValidator {
override fun pendingTxExists(): Flux<Boolean> {
return Flux.just(true)
}
}
class PendingTransactionValidatorImpl(
private val upstreamId: String,
private val directReader: ChainReader,
private val interval: Duration,
private val txLimit: Long,
) : PendingTransactionValidator {
private val log = LoggerFactory.getLogger(this::class.java)
override fun pendingTxExists(): Flux<Boolean> {
return Flux.interval(
Duration.ofSeconds(15),
interval,
)
.flatMap {
directReader.read(ChainRequest("txpool_content", ListParams()))
.flatMap(ChainResponse::requireResult)
.map {
val node = Global.objectMapper.readTree(it)
val pendingTxsNode = node.get("pending")
val queuedTxsNode = node.get("queued")
val pendingTxsCount = if (pendingTxsNode != null) {
pendingTxsNode.fieldNames().asSequence().toList().size
} else {
0
}
val queuedTxsCount = if (queuedTxsNode != null) {
queuedTxsNode.fieldNames().asSequence().toList().size
} else {
0
}
((pendingTxsCount + queuedTxsCount) >= txLimit)
}
.timeout(Duration.ofSeconds(30))
.onErrorResume {
log.error("unable to read txs from txpool_content of upstream {}", upstreamId, it)
Mono.just(false)
}
}
}
}

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)
}
@@ -404,7 +410,7 @@ open class WsConnectionImpl(
return Mono.from(onResponse.asMono()).or(failOnDisconnect)
.doOnSubscribe { sendRpc(request) }
.take(Defaults.timeout)
.doOnNext { requestMetrics?.timer?.record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS) }
.doOnNext { requestMetrics?.timer()?.record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS) }
.doOnError { requestMetrics?.fails?.increment() }
.map { it.copyWithId(ChainResponse.Id.from(originalId)) }
.switchIfEmpty(
@@ -424,7 +430,7 @@ open class WsConnectionImpl(
it.close()
Metrics.globalRegistry.remove(it)
}
requestMetrics?.timer?.let {
requestMetrics?.registeredTimers()?.forEach {
it.close()
Metrics.globalRegistry.remove(it)
}

View File

@@ -43,6 +43,7 @@ import reactor.core.Disposable
import reactor.core.publisher.Flux
import reactor.core.publisher.Sinks
import reactor.core.scheduler.Schedulers
import java.net.URI
import java.time.Duration
import java.util.concurrent.Executors
import java.util.concurrent.atomic.AtomicBoolean
@@ -101,6 +102,8 @@ open class GenericUpstream(
) {
rpcMethodsDetector = upstreamRpcMethodsDetectorBuilder(this, config)
detectRpcMethods(config, buildMethods)
rpcConnectionUrl = (config.connection as? UpstreamsConfig.RpcConnection)
?.let { it.rpc?.url ?: it.ws?.url }
}
private val validator: UpstreamValidator? = validatorBuilder(chain, this, getOptions(), chainConfig, versionRules)
@@ -110,6 +113,8 @@ open class GenericUpstream(
private val settingsDetectorSubscription = AtomicReference<Disposable?>()
private val hasLiveSubscriptionHead: AtomicBoolean = AtomicBoolean(getOptions().disableLivenessSubscriptionValidation)
private val hasPendingTxs = AtomicBoolean(getOptions().disablePendingTxValidation)
protected val connector: GenericConnector = connectorFactory.create(this, chain)
.also { upConnector ->
Gauge.builder("upstream_head", upConnector.getHead()) {
@@ -120,9 +125,13 @@ open class GenericUpstream(
.register(Metrics.globalRegistry)
}
private val livenessSubscription = AtomicReference<Disposable?>()
private val pendingTxSubscription = AtomicReference<Disposable?>()
private val settingsDetector = upstreamSettingsDetectorBuilder(chain, this)
private var rpcMethodsDetector: UpstreamRpcMethodsDetector? = null
// configured RPC/WS URL (carries query flags like ?hl=)
private var rpcConnectionUrl: URI? = null
private val lowerBoundService = lowerBoundServiceBuilder(chain, this)
private val started = AtomicBoolean(false)
@@ -150,11 +159,14 @@ open class GenericUpstream(
// outdated, looks like applicable only for bitcoin and our ws_head trick
override fun getCapabilities(): Set<Capability> {
return if (hasLiveSubscriptionHead.get()) {
setOf(Capability.RPC, Capability.BALANCE, Capability.WS_HEAD)
} else {
setOf(Capability.RPC, Capability.BALANCE)
val caps = mutableSetOf(Capability.RPC, Capability.BALANCE)
if (hasLiveSubscriptionHead.get()) {
caps.add(Capability.WS_HEAD)
}
if (hasPendingTxs.get()) {
caps.add(Capability.WS_PENDING_TX)
}
return caps
}
override fun isGrpc(): Boolean {
@@ -178,6 +190,8 @@ open class GenericUpstream(
)
}
fun getRpcConnectionUrl(): URI? = rpcConnectionUrl
@Suppress("UNCHECKED_CAST")
override fun <T : Upstream> cast(selfType: Class<T>): T {
if (!selfType.isAssignableFrom(this.javaClass)) {
@@ -358,6 +372,14 @@ open class GenericUpstream(
),
)
}
if (!getOptions().disablePendingTxValidation) {
pendingTxSubscription.set(
connector.pendingTxEvents().subscribe {
hasPendingTxs.set(it)
sendUpstreamStateEvent(UPDATED)
},
)
}
detectSettings()
if (!getOptions().disableBoundValidation) {
@@ -380,6 +402,7 @@ open class GenericUpstream(
lowerBlockDetectorSubscription.getAndSet(null)?.dispose()
finalizationDetectorSubscription.getAndSet(null)?.dispose()
settingsDetectorSubscription.getAndSet(null)?.dispose()
pendingTxSubscription.getAndSet(null)?.dispose()
connector.getHead().stop()
}

View File

@@ -15,4 +15,6 @@ interface GenericConnector : Lifecycle {
fun getIngressReader(): ChainReader
fun getIngressSubscription(): IngressSubscription
fun pendingTxEvents(): Flux<Boolean>
}

View File

@@ -13,11 +13,15 @@ import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.MergedHead
import io.emeraldpay.dshackle.upstream.NoIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.AlwaysHeadLivenessValidator
import io.emeraldpay.dshackle.upstream.ethereum.BASE_TX_LIMIT
import io.emeraldpay.dshackle.upstream.ethereum.GenericWsHead
import io.emeraldpay.dshackle.upstream.ethereum.HeadLivenessState
import io.emeraldpay.dshackle.upstream.ethereum.HeadLivenessValidator
import io.emeraldpay.dshackle.upstream.ethereum.HeadLivenessValidatorImpl
import io.emeraldpay.dshackle.upstream.ethereum.NoHeadLivenessValidator
import io.emeraldpay.dshackle.upstream.ethereum.NoopPendingTransactionValidator
import io.emeraldpay.dshackle.upstream.ethereum.PendingTransactionValidator
import io.emeraldpay.dshackle.upstream.ethereum.PendingTransactionValidatorImpl
import io.emeraldpay.dshackle.upstream.ethereum.WsConnectionPool
import io.emeraldpay.dshackle.upstream.ethereum.WsConnectionPoolFactory
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions
@@ -58,6 +62,7 @@ class GenericRpcConnector(
private val head: Head
private val liveness: HeadLivenessValidator
private val jsonRpcWsClient: JsonRpcWsClient?
private val pendingTxValidator: PendingTransactionValidator
companion object {
private val log = LoggerFactory.getLogger(GenericRpcConnector::class.java)
@@ -136,6 +141,16 @@ class GenericRpcConnector(
)
}
}
pendingTxValidator = if (chain == Chain.BASE__MAINNET) {
PendingTransactionValidatorImpl(
upstream.getId(),
getIngressReader(),
Duration.ofMinutes(5),
BASE_TX_LIMIT,
)
} else {
NoopPendingTransactionValidator()
}
liveness = if (connectorType != RPC_ONLY && isSpecialChain(chain)) {
AlwaysHeadLivenessValidator()
@@ -186,6 +201,10 @@ class GenericRpcConnector(
return ingressSubscription ?: NoIngressSubscription()
}
override fun pendingTxEvents(): Flux<Boolean> {
return pendingTxValidator.pendingTxExists()
}
override fun getHead(): Head {
return head
}

View File

@@ -6,9 +6,13 @@ import io.emeraldpay.dshackle.upstream.BlockValidator
import io.emeraldpay.dshackle.upstream.DefaultUpstream
import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.BASE_TX_LIMIT
import io.emeraldpay.dshackle.upstream.ethereum.GenericWsHead
import io.emeraldpay.dshackle.upstream.ethereum.HeadLivenessState
import io.emeraldpay.dshackle.upstream.ethereum.HeadLivenessValidatorImpl
import io.emeraldpay.dshackle.upstream.ethereum.NoopPendingTransactionValidator
import io.emeraldpay.dshackle.upstream.ethereum.PendingTransactionValidator
import io.emeraldpay.dshackle.upstream.ethereum.PendingTransactionValidatorImpl
import io.emeraldpay.dshackle.upstream.ethereum.WsConnectionPool
import io.emeraldpay.dshackle.upstream.ethereum.WsConnectionPoolFactory
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptionsImpl
@@ -36,6 +40,8 @@ class GenericWsConnector(
private val head: GenericWsHead
private val subscriptions: IngressSubscription
private val liveness: HeadLivenessValidatorImpl
private val pendingTxValidator: PendingTransactionValidator
init {
pool = wsFactory.create(upstream)
reader = JsonRpcWsClient(pool)
@@ -54,6 +60,17 @@ class GenericWsConnector(
)
liveness = HeadLivenessValidatorImpl(head, expectedBlockTime, headLivenessScheduler, upstream.getId())
subscriptions = chainSpecific.makeIngressSubscription(chain, wsSubscriptions)
pendingTxValidator = if (chain == Chain.BASE__MAINNET) {
PendingTransactionValidatorImpl(
upstream.getId(),
getIngressReader(),
Duration.ofMinutes(10),
BASE_TX_LIMIT,
)
} else {
NoopPendingTransactionValidator()
}
}
override fun headLivenessEvents(): Flux<HeadLivenessState> {
@@ -81,6 +98,10 @@ class GenericWsConnector(
return subscriptions
}
override fun pendingTxEvents(): Flux<Boolean> {
return pendingTxValidator.pendingTxExists()
}
override fun getHead(): Head {
return head
}

View File

@@ -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<Chain, DefaultUpstream>()
private val lock = ReentrantLock()
private val statusSubscriptions = mutableMapOf<Chain, Disposable>()
fun start(): Flux<UpstreamChangeEvent> {
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<Chain, Disposable>()
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<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

@@ -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<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
@@ -59,7 +60,7 @@ class RestHttpReader(
.flatMap(this::execute)
.doOnNext {
if (startTime.isStarted) {
metrics?.timer?.record(startTime.nanoTime, TimeUnit.NANOSECONDS)
metrics?.timer(key.method)?.record(startTime.nanoTime, TimeUnit.NANOSECONDS)
}
}
.handle { it, sink ->

View File

@@ -114,7 +114,7 @@ class JsonRpcGrpcClient(
}
.doOnNext {
if (timer.isStarted) {
metrics?.timer?.record(timer.getTime(TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS)
metrics?.timer()?.record(timer.getTime(TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS)
}
}
}

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()
@@ -92,7 +93,7 @@ class JsonRpcHttpReader(
.flatMap(this@JsonRpcHttpReader::execute)
.doOnNext {
if (startTime.isStarted) {
metrics?.timer?.record(startTime.nanoTime, TimeUnit.NANOSECONDS)
metrics?.timer(key.method)?.record(startTime.nanoTime, TimeUnit.NANOSECONDS)
}
}
.transform(asJsonRpcResponse(key))

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")
@@ -636,7 +668,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def options = partialOptions.buildOptions()
then:
options == new ChainOptions.Options(
false, false, 30, Duration.ofSeconds(60), null, true, 1, true, true, true, true, 1_000_000, false, false, true, false
false, false, 30, Duration.ofSeconds(60), null, true, 1, true, true, true, true, 1_000_000, false, false, true, false, false
)
}
}

View File

@@ -14,11 +14,13 @@ class GenericConnectorMock implements GenericConnector {
Reader<ChainRequest, ChainResponse> api
Head head
Flux<HeadLivenessState> liveness
Flux<Boolean> pendingTxs
GenericConnectorMock(Reader<ChainRequest, ChainResponse> api, Head head) {
this.api = api
this.head = head
this.liveness = Flux.just(HeadLivenessState.NON_CONSECUTIVE)
this.pendingTxs = Flux.just(false)
}
@Override
@@ -51,4 +53,9 @@ class GenericConnectorMock implements GenericConnector {
IngressSubscription getIngressSubscription() {
return NoEthereumIngressSubscription.DEFAULT
}
@Override
Flux<Boolean> pendingTxEvents() {
return pendingTxs
}
}

View File

@@ -17,6 +17,8 @@ package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.test.TestingCommons
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.AggregatedPendingTxes
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.NoPendingTxes
import io.emeraldpay.dshackle.upstream.generic.GenericMultistream
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource
import io.emeraldpay.dshackle.upstream.ethereum.domain.Address
@@ -216,6 +218,7 @@ class EthereumEgressSubscriptionSpec extends Specification {
when:
def up3 = TestingCommons.upstream("test")
up3.getConnectorMock().setLiveness(Flux.just(HeadLivenessState.OK))
up3.getConnectorMock().setPendingTxs(Flux.just(true))
up3.stop()
up3.start()
def ethereumSubscribe3 = new EthereumEgressSubscription(TestingCommons.multistream(up3) as GenericMultistream, Schedulers.boundedElastic(), Stub(PendingTxesSource))
@@ -230,5 +233,23 @@ class EthereumEgressSubscriptionSpec extends Specification {
then:
ethereumSubscribe4.getAvailableTopics().toSet() == [EthereumEgressSubscription.METHOD_NEW_HEADS].toSet()
when:
def up5 = TestingCommons.upstream("test")
up5.getConnectorMock().setLiveness(Flux.just(HeadLivenessState.OK))
up5.stop()
up5.start()
def ethereumSubscribe5 = new EthereumEgressSubscription(TestingCommons.multistream(up5) as GenericMultistream, Schedulers.boundedElastic(), Stub(PendingTxesSource))
then:
ethereumSubscribe5.getAvailableTopics().toSet() == [EthereumEgressSubscription.METHOD_LOGS, EthereumEgressSubscription.METHOD_NEW_HEADS].toSet()
when:
def up6 = TestingCommons.upstream("test")
up6.getConnectorMock().setLiveness(Flux.just(HeadLivenessState.OK))
up6.getConnectorMock().setPendingTxs(Flux.just(true))
up6.stop()
up6.start()
def ethereumSubscribe6 = new EthereumEgressSubscription(TestingCommons.multistream(up6) as GenericMultistream, Schedulers.boundedElastic(), new NoPendingTxes())
then:
ethereumSubscribe6.getAvailableTopics().toSet() == [EthereumEgressSubscription.METHOD_LOGS, EthereumEgressSubscription.METHOD_NEW_HEADS].toSet()
}
}

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,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<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() {
@@ -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<ReloadConfigUpstreamService>()
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<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,146 @@
package io.emeraldpay.dshackle.upstream.aztec
import io.emeraldpay.dshackle.reader.ChainReader
import io.emeraldpay.dshackle.upstream.ChainCallError
import io.emeraldpay.dshackle.upstream.ChainCallUpstreamException
import io.emeraldpay.dshackle.upstream.ChainRequest
import io.emeraldpay.dshackle.upstream.ChainResponse
import org.assertj.core.api.Assertions
import org.junit.jupiter.api.Test
import reactor.core.publisher.Mono
import java.util.concurrent.atomic.AtomicInteger
// v5 (v5.0.0-rc.1) node_getChainTips response: proposed stays flat, the rest nested.
private val chainTipsResponse = """
{
"proposed": {"number": 12345, "hash": "0xaaaa"},
"checkpointed": {"block": {"number": 12340, "hash": "0xbbbb"}, "checkpoint": {"number": 100, "hash": "0x1111"}},
"proven": {"block": {"number": 12330, "hash": "0xcccc"}, "checkpoint": {"number": 99, "hash": "0x2222"}},
"finalized": {"block": {"number": 12320, "hash": "0xdddd"}, "checkpoint": {"number": 98, "hash": "0x3333"}}
}
""".trimIndent()
// legacy node_getL2Tips response (v3-era flat proposed)
private val l2TipsResponse = """
{
"proposed": {"number": 999, "hash": "0x0999"},
"proven": {"number": 990, "hash": "0x0990"},
"checkpointed": {"number": 980, "hash": "0x0980"}
}
""".trimIndent()
// defensive: some versions nest proposed under .block
private val nestedProposedResponse = """
{
"proposed": {"block": {"number": 777, "hash": "0x0777"}}
}
""".trimIndent()
private fun methodNotFound(method: String) =
Mono.error<ChainResponse>(
ChainCallUpstreamException(
ChainResponse.NumberId(1),
ChainCallError(-32601, "Method not found: $method"),
),
)
class AztecChainSpecificTest {
@Test
fun latestBlockRequestUsesL2Tips() {
Assertions.assertThat(AztecChainSpecific.latestBlockRequest().method)
.isEqualTo("node_getL2Tips")
}
@Test
fun parseBlockReadsFlatProposed() {
val result = AztecChainSpecific.parseBlock(
chainTipsResponse.toByteArray(),
"up-flat",
noopReader(),
).block()!!
Assertions.assertThat(result.height).isEqualTo(12345L)
Assertions.assertThat(result.hash.toHex()).contains("aaaa")
}
@Test
fun parseBlockReadsNestedProposed() {
val result = AztecChainSpecific.parseBlock(
nestedProposedResponse.toByteArray(),
"up-nested",
noopReader(),
).block()!!
Assertions.assertThat(result.height).isEqualTo(777L)
Assertions.assertThat(result.hash.toHex()).contains("0777")
}
@Test
fun getLatestBlockUsesL2TipsWhenAvailable() {
val calls = mutableListOf<String>()
val reader = object : ChainReader {
override fun read(key: ChainRequest): Mono<ChainResponse> {
calls += key.method
return Mono.just(ChainResponse(l2TipsResponse.toByteArray(), null))
}
}
val result = AztecChainSpecific.getLatestBlock(reader, "up-legacy").block()!!
Assertions.assertThat(result.height).isEqualTo(999L)
Assertions.assertThat(calls).containsExactly("node_getL2Tips")
}
@Test
fun getLatestBlockFallsBackToChainTipsAndCaches() {
val calls = mutableListOf<String>()
val reader = object : ChainReader {
override fun read(key: ChainRequest): Mono<ChainResponse> {
calls += key.method
return when (key.method) {
"node_getL2Tips" -> methodNotFound("node_getL2Tips")
"node_getChainTips" -> Mono.just(ChainResponse(chainTipsResponse.toByteArray(), null))
else -> Mono.error(IllegalStateException("unexpected ${key.method}"))
}
}
}
// first poll: probes legacy, falls back to v5
val first = AztecChainSpecific.getLatestBlock(reader, "up-v5").block()!!
Assertions.assertThat(first.height).isEqualTo(12345L)
Assertions.assertThat(calls).containsExactly("node_getL2Tips", "node_getChainTips")
// second poll: must hit the cached working method directly, no dead probe
calls.clear()
val second = AztecChainSpecific.getLatestBlock(reader, "up-v5").block()!!
Assertions.assertThat(second.height).isEqualTo(12345L)
Assertions.assertThat(calls).containsExactly("node_getChainTips")
}
@Test
fun getLatestBlockDoesNotFallBackOnOtherErrors() {
val attempts = AtomicInteger(0)
val reader = object : ChainReader {
override fun read(key: ChainRequest): Mono<ChainResponse> {
attempts.incrementAndGet()
return Mono.error(
ChainCallUpstreamException(
ChainResponse.NumberId(1),
ChainCallError(-32000, "internal error"),
),
)
}
}
val thrown = runCatching { AztecChainSpecific.getLatestBlock(reader, "up-err").block() }
Assertions.assertThat(thrown.isFailure).isTrue()
// only the primary method is attempted; no fallback probe on a non-method-not-found error
Assertions.assertThat(attempts.get()).isEqualTo(1)
}
private fun noopReader() = object : ChainReader {
override fun read(key: ChainRequest): Mono<ChainResponse> =
Mono.error(IllegalStateException("not expected"))
}
}

View File

@@ -0,0 +1,211 @@
package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.reader.ChainReader
import io.emeraldpay.dshackle.upstream.ChainCallError
import io.emeraldpay.dshackle.upstream.ChainRequest
import io.emeraldpay.dshackle.upstream.ChainResponse
import io.emeraldpay.dshackle.upstream.rpcclient.ListParams
import org.junit.jupiter.api.Test
import org.mockito.kotlin.doReturn
import org.mockito.kotlin.mock
import reactor.core.publisher.Mono
import reactor.test.StepVerifier
import java.time.Duration
class PendingTransactionValidatorTest {
@Test
fun `noop validator always emits true`() {
val validator = NoopPendingTransactionValidator()
StepVerifier.create(validator.pendingTxExists())
.expectNext(true)
.expectComplete()
.verify(Duration.ofSeconds(3))
}
@Test
fun `emits true when pending plus queued exceeds limit`() {
val reader = mockReader(
response = txpoolContent(pendingAddresses = 6, queuedAddresses = 0),
)
val validator = PendingTransactionValidatorImpl(
upstreamId = "test-upstream",
directReader = reader,
interval = Duration.ofSeconds(30),
txLimit = 5,
)
StepVerifier.withVirtualTime { validator.pendingTxExists() }
.expectSubscription()
.expectNoEvent(Duration.ofSeconds(15))
.expectNext(true)
.thenCancel()
.verify(Duration.ofSeconds(3))
}
@Test
fun `emits true when pending plus queued is at limit`() {
val reader = mockReader(
response = txpoolContent(pendingAddresses = 3, queuedAddresses = 2),
)
val validator = PendingTransactionValidatorImpl(
upstreamId = "test-upstream",
directReader = reader,
interval = Duration.ofSeconds(30),
txLimit = 5,
)
StepVerifier.withVirtualTime { validator.pendingTxExists() }
.expectSubscription()
.expectNoEvent(Duration.ofSeconds(15))
.expectNext(true)
.thenCancel()
.verify(Duration.ofSeconds(3))
}
@Test
fun `counts both pending and queued buckets`() {
// 3 pending + 4 queued = 7 > limit of 5
val reader = mockReader(
response = txpoolContent(pendingAddresses = 3, queuedAddresses = 4),
)
val validator = PendingTransactionValidatorImpl(
upstreamId = "test-upstream",
directReader = reader,
interval = Duration.ofSeconds(30),
txLimit = 5,
)
StepVerifier.withVirtualTime { validator.pendingTxExists() }
.expectSubscription()
.expectNoEvent(Duration.ofSeconds(15))
.expectNext(true)
.thenCancel()
.verify(Duration.ofSeconds(3))
}
@Test
fun `handles missing pending field`() {
val reader = mockReader(response = """{"queued": {"0xaa": {}, "0xbb": {}}}""")
val validator = PendingTransactionValidatorImpl(
upstreamId = "test-upstream",
directReader = reader,
interval = Duration.ofSeconds(30),
txLimit = 5,
)
// Only 2 queued (no pending node) -> below limit -> false
StepVerifier.withVirtualTime { validator.pendingTxExists() }
.expectSubscription()
.expectNoEvent(Duration.ofSeconds(15))
.expectNext(false)
.thenCancel()
.verify(Duration.ofSeconds(3))
}
@Test
fun `handles missing queued field`() {
val reader = mockReader(response = """{"pending": {"0xaa": {}, "0xbb": {}, "0xcc": {}}}""")
val validator = PendingTransactionValidatorImpl(
upstreamId = "test-upstream",
directReader = reader,
interval = Duration.ofSeconds(30),
txLimit = 2,
)
// 3 pending (no queued node) > limit 2 -> true
StepVerifier.withVirtualTime { validator.pendingTxExists() }
.expectSubscription()
.expectNoEvent(Duration.ofSeconds(15))
.expectNext(true)
.thenCancel()
.verify(Duration.ofSeconds(3))
}
@Test
fun `emits false on rpc error`() {
val reader = mock<ChainReader> {
on { read(ChainRequest("txpool_content", ListParams())) } doReturn
Mono.just(ChainResponse(null, ChainCallError(-32000, "method not supported")))
}
val validator = PendingTransactionValidatorImpl(
upstreamId = "test-upstream",
directReader = reader,
interval = Duration.ofSeconds(30),
txLimit = 5,
)
StepVerifier.withVirtualTime { validator.pendingTxExists() }
.expectSubscription()
.expectNoEvent(Duration.ofSeconds(15))
.expectNext(false)
.thenCancel()
.verify(Duration.ofSeconds(3))
}
@Test
fun `continues polling after an error response`() {
val errorResponse = Mono.just(ChainResponse(null, ChainCallError(-32000, "method not supported")))
val okResponse = Mono.just(
ChainResponse(
txpoolContent(pendingAddresses = 10, queuedAddresses = 0).toByteArray(),
null,
),
)
val reader = mock<ChainReader> {
on { read(ChainRequest("txpool_content", ListParams())) } doReturn errorResponse doReturn okResponse
}
val validator = PendingTransactionValidatorImpl(
upstreamId = "test-upstream",
directReader = reader,
interval = Duration.ofSeconds(30),
txLimit = 5,
)
StepVerifier.withVirtualTime { validator.pendingTxExists() }
.expectSubscription()
.expectNoEvent(Duration.ofSeconds(15))
.expectNext(false) // first poll: error -> false
.expectNoEvent(Duration.ofSeconds(30))
.expectNext(true) // second poll: success
.thenCancel()
.verify(Duration.ofSeconds(3))
}
@Test
fun `polls at configured interval`() {
val reader = mockReader(
response = txpoolContent(pendingAddresses = 10, queuedAddresses = 0),
)
val validator = PendingTransactionValidatorImpl(
upstreamId = "test-upstream",
directReader = reader,
interval = Duration.ofSeconds(30),
txLimit = 5,
)
StepVerifier.withVirtualTime { validator.pendingTxExists() }
.expectSubscription()
.expectNoEvent(Duration.ofSeconds(15))
.expectNext(true)
.expectNoEvent(Duration.ofSeconds(30))
.expectNext(true)
.expectNoEvent(Duration.ofSeconds(30))
.expectNext(true)
.thenCancel()
.verify(Duration.ofSeconds(3))
}
private fun mockReader(response: String): ChainReader =
mock<ChainReader> {
on { read(ChainRequest("txpool_content", ListParams())) } doReturn
Mono.just(ChainResponse(response.toByteArray(), null))
}
private fun txpoolContent(pendingAddresses: Int, queuedAddresses: Int): String {
val pending = (0 until pendingAddresses).joinToString(",") { """"0xp$it": {}""" }
val queued = (0 until queuedAddresses).joinToString(",") { """"0xq$it": {}""" }
return """{"pending": {$pending}, "queued": {$queued}}"""
}
}

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

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

View File

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