diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsSubscriptionsImpl.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsSubscriptionsImpl.kt index 6b7be8c2..089aa291 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsSubscriptionsImpl.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsSubscriptionsImpl.kt @@ -19,6 +19,8 @@ import io.emeraldpay.dshackle.upstream.ChainException import io.emeraldpay.dshackle.upstream.ChainRequest import io.emeraldpay.dshackle.upstream.ChainResponse import io.emeraldpay.dshackle.upstream.rpcclient.ListParams +import io.emeraldpay.dshackle.upstream.rpcclient.ObjectParams +import io.emeraldpay.dshackle.upstream.rpcclient.RippleCommandParams import org.slf4j.LoggerFactory import reactor.core.publisher.Flux import reactor.core.publisher.Mono @@ -47,8 +49,10 @@ class WsSubscriptionsImpl( log.warn("Failed to establish subscription: ${it.error?.message}") Mono.error(ChainException(it.id, it.error!!)) } else { - val id = if (it.getResultAsRawString() == "{}") { - request.id.toString() // in case empty result - match by request id + val rawResult = it.getResultAsRawString() + val id = if (rawResult.startsWith("{") || rawResult == "{}") { + // Ripple returns object result; use stream type as subscription ID + extractRippleStreamType(request) ?: request.id.toString() } else { it.getResultAsProcessedString() } @@ -60,6 +64,26 @@ class WsSubscriptionsImpl( return WsSubscriptions.SubscribeData(message, conn.connectionId(), subscriptionId) } + /** + * Extract expected event type from Ripple subscribe request. + * For `subscribe` with `streams: ["ledger"]`, returns "ledgerClosed" as this is the event type + * that will be received in subscription messages. + */ + private fun extractRippleStreamType(request: ChainRequest): String? { + if (request.method == "subscribe") { + val params = request.params + val streams: List<*>? = when (params) { + is RippleCommandParams -> params.params["streams"] as? List<*> + is ObjectParams -> params.obj["streams"] as? List<*> + else -> null + } + if (streams?.contains("ledger") == true) { + return "ledgerClosed" + } + } + return null + } + override fun unsubscribe(request: ChainRequest): Mono { if (request.params is ListParams && (request.params.list.isEmpty() || request.params.list.contains("")) ) { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ripple/RippleChainSpecific.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ripple/RippleChainSpecific.kt index 4a77d2ce..3ce071bb 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ripple/RippleChainSpecific.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ripple/RippleChainSpecific.kt @@ -12,70 +12,83 @@ import io.emeraldpay.dshackle.foundation.ChainOptions.Options import io.emeraldpay.dshackle.reader.ChainReader import io.emeraldpay.dshackle.upstream.BasicUpstreamSettingsDetector import io.emeraldpay.dshackle.upstream.ChainRequest +import io.emeraldpay.dshackle.upstream.EgressSubscription import io.emeraldpay.dshackle.upstream.GenericSingleCallValidator +import io.emeraldpay.dshackle.upstream.IngressSubscription +import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.NodeTypeRequest import io.emeraldpay.dshackle.upstream.SingleValidator import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamSettingsDetector import io.emeraldpay.dshackle.upstream.ValidateUpstreamSettingsResult +import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions import io.emeraldpay.dshackle.upstream.generic.AbstractPollChainSpecific +import io.emeraldpay.dshackle.upstream.generic.GenericEgressSubscription +import io.emeraldpay.dshackle.upstream.generic.GenericIngressSubscription import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundService import io.emeraldpay.dshackle.upstream.rpcclient.ListParams -import org.slf4j.LoggerFactory +import io.emeraldpay.dshackle.upstream.rpcclient.RippleCommandParams import reactor.core.publisher.Flux import reactor.core.publisher.Mono +import reactor.core.scheduler.Scheduler import java.math.BigInteger import java.time.Instant object RippleChainSpecific : AbstractPollChainSpecific() { - private val log = LoggerFactory.getLogger(RippleChainSpecific::class.java) - override fun parseBlock(data: ByteArray, upstreamId: String, api: ChainReader): Mono { + // Parse ledger_closed response: { "ledger_hash": "...", "ledger_index": 123 } val jsonNode = Global.objectMapper.readTree(data) - - val block: RippleClosedLedger = if (jsonNode.has("ledger")) { - Global.objectMapper.treeToValue(jsonNode.get("ledger"), RippleClosedLedger::class.java) - } else { - val result = Global.objectMapper.readValue(data, RippleBlock::class.java) - result.closed.ledger - } - - var height: Long = 0 - try { - height = block.ledgerIndex.toLong() - } catch (e: NumberFormatException) { - log.error("Invalid ledgerIndex ${block.ledgerIndex}, upstreamId:$upstreamId") - } + val ledgerHash = jsonNode.get("ledger_hash").asText() + val ledgerIndex = jsonNode.get("ledger_index").asLong() return Mono.just( BlockContainer( - height = height, - hash = BlockId.from(block.ledgerHash), + height = ledgerIndex, + hash = BlockId.from(ledgerHash), difficulty = BigInteger.ZERO, timestamp = Instant.EPOCH, full = false, json = data, - parsed = block, + parsed = jsonNode, transactions = emptyList(), upstreamId = upstreamId, - parentHash = BlockId.from(block.parentHash), + parentHash = null, ), ) } override fun getFromHeader(data: ByteArray, upstreamId: String, api: ChainReader): Mono { - throw NotImplementedError() + // Parse Ripple ledger stream event: { "type": "ledgerClosed", "ledger_hash": "...", "ledger_index": 123, ... } + val event = Global.objectMapper.readValue(data, RippleLedgerStreamEvent::class.java) + + // Ripple epoch starts at 2000-01-01 00:00:00 UTC (946684800 seconds after Unix epoch) + val timestamp = event.ledgerTime?.let { + Instant.ofEpochSecond(it + 946684800L) + } ?: Instant.EPOCH + + return Mono.just( + BlockContainer( + height = event.ledgerIndex, + hash = BlockId.from(event.ledgerHash), + difficulty = BigInteger.ZERO, + timestamp = timestamp, + full = false, + json = data, + parsed = event, + transactions = emptyList(), + upstreamId = upstreamId, + parentHash = null, // Not available in stream event + ), + ) } - override fun listenNewHeadsRequest(): ChainRequest { - throw NotImplementedError() - } + override fun listenNewHeadsRequest(): ChainRequest = + ChainRequest("subscribe", RippleCommandParams("streams" to listOf("ledger"))) - override fun unsubscribeNewHeadsRequest(subId: Any): ChainRequest { - throw NotImplementedError() - } + override fun unsubscribeNewHeadsRequest(subId: Any): ChainRequest = + ChainRequest("unsubscribe", RippleCommandParams("streams" to listOf("ledger"))) override fun upstreamValidators( chain: Chain, @@ -131,11 +144,19 @@ object RippleChainSpecific : AbstractPollChainSpecific() { } override fun latestBlockRequest(): ChainRequest = - ChainRequest("ledger", ListParams()) + ChainRequest("ledger_closed", ListParams()) override fun upstreamSettingsDetector(chain: Chain, upstream: Upstream): UpstreamSettingsDetector? { return RippleUpstreamSettingsDetector(upstream) } + + override fun makeIngressSubscription(chain: Chain, ws: WsSubscriptions): IngressSubscription { + return GenericIngressSubscription(chain, ws, listOf("ledger")) + } + + override fun subscriptionBuilder(headScheduler: Scheduler): (Multistream) -> EgressSubscription { + return { ms -> GenericEgressSubscription(ms, headScheduler) } + } } class RippleUpstreamSettingsDetector(val upstream: Upstream) : BasicUpstreamSettingsDetector(upstream) { @@ -189,47 +210,6 @@ data class RippleInfo( @param:JsonProperty("build_version") var buildVersion: String?, ) -@JsonIgnoreProperties(ignoreUnknown = true) -data class RippleBlock( - @param:JsonProperty("closed") var closed: RippleClosed, - @param:JsonProperty("open") var open: RippleOpen?, - @param:JsonProperty("status") var status: String?, -) - -@JsonIgnoreProperties(ignoreUnknown = true) -data class RippleClosed( - @param:JsonProperty("ledger") var ledger: RippleClosedLedger, -) - -@JsonIgnoreProperties(ignoreUnknown = true) -data class RippleOpen( - @param:JsonProperty("ledger") var ledger: RippleOpenLedger, -) - -@JsonIgnoreProperties(ignoreUnknown = true) -data class RippleClosedLedger( - @param:JsonProperty("account_hash") var accountHash: String, - @param:JsonProperty("close_flags") var closeFlags: Short, - @param:JsonProperty("close_time") var closeTime: Long, - @param:JsonProperty("close_time_human") var closeTimeHuman: String, - @param:JsonProperty("close_time_iso") var closeTimeIso: String, - @param:JsonProperty("close_time_resolution") var closeTimeResolution: Short, - @param:JsonProperty("closed") var closed: Boolean, - @param:JsonProperty("ledger_hash") var ledgerHash: String, - @param:JsonProperty("ledger_index") var ledgerIndex: String, - @param:JsonProperty("parent_close_time") var parentCloseTime: Long, - @param:JsonProperty("parent_hash") var parentHash: String, - @param:JsonProperty("total_coins") var totalCoins: String, - @param:JsonProperty("transaction_hash") var transactionHash: String, -) - -@JsonIgnoreProperties(ignoreUnknown = true) -data class RippleOpenLedger( - @param:JsonProperty("closed") var closed: Boolean, - @param:JsonProperty("ledger_index") var ledgerIndex: String, - @param:JsonProperty("parent_hash") var parentHash: String, -) - @JsonIgnoreProperties(ignoreUnknown = true) data class RippleState( @param:JsonProperty("state") var state: RippleServerState, @@ -301,3 +281,21 @@ data class ValidatedLedger( @param:JsonProperty("reserve_inc") val reserveInc: Long, @param:JsonProperty("seq") val seq: Long, ) + +/** + * Ripple ledger stream subscription event. + * Received when subscribed to the "ledger" stream via WebSocket. + */ +@JsonIgnoreProperties(ignoreUnknown = true) +data class RippleLedgerStreamEvent( + @param:JsonProperty("type") val type: String, // "ledgerClosed" + @param:JsonProperty("ledger_hash") val ledgerHash: String, + @param:JsonProperty("ledger_index") val ledgerIndex: Long, + @param:JsonProperty("ledger_time") val ledgerTime: Long? = null, + @param:JsonProperty("txn_count") val txnCount: Int? = null, + @param:JsonProperty("validated_ledgers") val validatedLedgers: String? = null, + @param:JsonProperty("reserve_base") val reserveBase: Long? = null, + @param:JsonProperty("reserve_inc") val reserveInc: Long? = null, + @param:JsonProperty("fee_base") val feeBase: Int? = null, + @param:JsonProperty("fee_ref") val feeRef: Int? = null, +) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/CallParams.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/CallParams.kt index bd031edd..e9764fc8 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/CallParams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/CallParams.kt @@ -39,6 +39,23 @@ data class ObjectParams(val obj: Map) : JsonRpcParams() { } } +/** + * Ripple native WebSocket command format. + * Unlike JSON-RPC, Ripple uses "command" field and flat structure. + */ +data class RippleCommandParams(val params: Map) : CallParams { + constructor(vararg pairs: Pair) : this(mapOf(*pairs)) + + override fun toJson(id: Int, method: String): ByteArray { + val json = mutableMapOf( + "id" to id, + "command" to method, + ) + json.putAll(params) + return Global.objectMapper.writeValueAsBytes(json) + } +} + data class RestParams( val headers: List>, val queryParams: List>, diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParser.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParser.kt index b99c65c2..3f5a7dcd 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParser.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParser.kt @@ -55,6 +55,18 @@ class ResponseWSParser : ResponseParser() { val method = parser.valueAsString return state.copy(subMethod = method) } + // Handle Ripple subscription format: { "type": "ledgerClosed", "ledger_hash": "...", ... } + if ("type" == field) { + parser.nextToken() + val type = parser.valueAsString + // "ledgerClosed" is a Ripple subscription notification + if (type == "ledgerClosed") { + // Use type as subscription identifier, entire JSON as result + return state.copy(subId = type, result = json) + } + // "response" is a normal RPC response, let it pass through + return state + } if ("params" == field) { // example: // newHeads diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParserSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParserSpec.groovy index 460e5e02..a0f91a6f 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParserSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParserSpec.groovy @@ -103,4 +103,47 @@ class ResponseWSParserSpec extends Specification { act.error == null new String(act.value) == "null" } + + def "Parse Ripple ledgerClosed subscription event"() { + setup: + def msg = '''{ + "type": "ledgerClosed", + "fee_base": 10, + "fee_ref": 10, + "ledger_hash": "17ACB57A0F73B5160713E81FE72B2AC9F6064541004E272BD09F257D57C30C02", + "ledger_index": 6643099, + "ledger_time": 780804221, + "reserve_base": 10000000, + "reserve_inc": 2000000, + "txn_count": 5, + "validated_ledgers": "6643000-6643099" + }''' + when: + def act = parser.parse(msg.bytes) + then: + act.type == ResponseWSParser.Type.SUBSCRIPTION + act.id.asString() == "ledgerClosed" + act.error == null + with(new String(act.value)) { + it.contains("\"ledger_hash\"") + it.contains("17ACB57A0F73B5160713E81FE72B2AC9F6064541004E272BD09F257D57C30C02") + it.contains("\"ledger_index\": 6643099") + } + } + + def "Parse Ripple RPC response with type field"() { + setup: + def msg = '''{ + "id": 1, + "status": "success", + "type": "response", + "result": {} + }''' + when: + def act = parser.parse(msg.bytes) + then: + act.type == ResponseWSParser.Type.RPC + act.id.asNumber() == 1L + act.error == null + } } diff --git a/src/test/kotlin/io/emeraldpay/dshackle/upstream/ripple/RippleChainSpecificTest.kt b/src/test/kotlin/io/emeraldpay/dshackle/upstream/ripple/RippleChainSpecificTest.kt new file mode 100644 index 00000000..343aa84c --- /dev/null +++ b/src/test/kotlin/io/emeraldpay/dshackle/upstream/ripple/RippleChainSpecificTest.kt @@ -0,0 +1,168 @@ +package io.emeraldpay.dshackle.upstream.ripple + +import io.emeraldpay.dshackle.Chain +import io.emeraldpay.dshackle.data.BlockId +import io.emeraldpay.dshackle.reader.ChainReader +import io.emeraldpay.dshackle.upstream.ChainRequest +import io.emeraldpay.dshackle.upstream.ChainResponse +import io.emeraldpay.dshackle.upstream.UpstreamAvailability +import io.emeraldpay.dshackle.upstream.ValidateUpstreamSettingsResult +import org.assertj.core.api.Assertions.assertThat +import org.junit.jupiter.api.Test +import reactor.core.publisher.Mono +import java.time.Instant + +// ledger_closed response format +val ledgerClosedResponse = """ +{ + "ledger_hash": "17ACB57A0F73B5160713E81FE72B2AC9F6064541004E272BD09F257D57C30C02", + "ledger_index": 6643099 +} +""".trimIndent() + +// ledgerClosed WebSocket subscription event +val ledgerClosedEvent = """ +{ + "type": "ledgerClosed", + "fee_base": 10, + "fee_ref": 10, + "ledger_hash": "17ACB57A0F73B5160713E81FE72B2AC9F6064541004E272BD09F257D57C30C02", + "ledger_index": 6643099, + "ledger_time": 780804221, + "reserve_base": 10000000, + "reserve_inc": 2000000, + "txn_count": 5, + "validated_ledgers": "6643000-6643099" +} +""".trimIndent() + +// server_state response for validation +val serverStateOk = """ +{ + "state": { + "build_version": "1.12.0", + "complete_ledgers": "32570-6643099", + "io_latency_ms": 1, + "jq_trans_overflow": "0", + "last_close": { + "converge_time": 2000, + "proposers": 34 + }, + "load_base": 256, + "load_factor": 256, + "load_factor_fee_escalation": 256, + "load_factor_fee_queue": 256, + "load_factor_fee_reference": 256, + "load_factor_server": 256, + "network_id": 0, + "peers": 21, + "ports": [], + "server_state": "full", + "time": "2024-Jan-01 00:00:00", + "uptime": 123456, + "validated_ledger": { + "base_fee": 10, + "close_time": 780804221, + "hash": "17ACB57A0F73B5160713E81FE72B2AC9F6064541004E272BD09F257D57C30C02", + "reserve_base": 10000000, + "reserve_inc": 2000000, + "seq": 6643099 + }, + "validation_quorum": 28 + } +} +""".trimIndent() + +val serverStateSyncing = """ +{ + "state": { + "build_version": "1.12.0", + "complete_ledgers": "32570-6643099", + "io_latency_ms": 1, + "jq_trans_overflow": "0", + "last_close": { + "converge_time": 2000, + "proposers": 34 + }, + "load_base": 256, + "load_factor": 256, + "load_factor_fee_escalation": 256, + "load_factor_fee_queue": 256, + "load_factor_fee_reference": 256, + "load_factor_server": 256, + "network_id": 0, + "peers": 21, + "ports": [], + "server_state": "connected", + "time": "2024-Jan-01 00:00:00", + "uptime": 123456, + "validated_ledger": { + "base_fee": 10, + "close_time": 780804221, + "hash": "17ACB57A0F73B5160713E81FE72B2AC9F6064541004E272BD09F257D57C30C02", + "reserve_base": 10000000, + "reserve_inc": 2000000, + "seq": 6643099 + }, + "validation_quorum": 28 + } +} +""".trimIndent() + +class RippleChainSpecificTest { + + private val dummyReader = object : ChainReader { + override fun read(key: ChainRequest): Mono = Mono.empty() + } + + @Test + fun `parseBlock with ledger_closed format`() { + val result = RippleChainSpecific.parseBlock( + ledgerClosedResponse.toByteArray(), + "test-upstream", + dummyReader, + ).block()!! + + assertThat(result.height).isEqualTo(6643099) + assertThat(result.hash) + .isEqualTo(BlockId.from("17ACB57A0F73B5160713E81FE72B2AC9F6064541004E272BD09F257D57C30C02")) + assertThat(result.upstreamId).isEqualTo("test-upstream") + assertThat(result.parentHash).isNull() + } + + @Test + fun `getFromHeader with ledgerClosed WS event`() { + val result = RippleChainSpecific.getFromHeader( + ledgerClosedEvent.toByteArray(), + "test-upstream", + dummyReader, + ).block()!! + + assertThat(result.height).isEqualTo(6643099) + assertThat(result.hash) + .isEqualTo(BlockId.from("17ACB57A0F73B5160713E81FE72B2AC9F6064541004E272BD09F257D57C30C02")) + assertThat(result.upstreamId).isEqualTo("test-upstream") + assertThat(result.parentHash).isNull() + // Ripple epoch (2000-01-01) + 780804221 seconds = expected timestamp + assertThat(result.timestamp).isEqualTo(Instant.ofEpochSecond(780804221 + 946684800L)) + } + + @Test + fun `validate returns OK for full server state`() { + val result = RippleChainSpecific.validate(serverStateOk.toByteArray()) + assertThat(result).isEqualTo(UpstreamAvailability.OK) + } + + @Test + fun `validate returns SYNCING for connected server state`() { + val result = RippleChainSpecific.validate(serverStateSyncing.toByteArray()) + assertThat(result).isEqualTo(UpstreamAvailability.SYNCING) + } + + @Test + fun `validateSettings returns VALID for matching network`() { + val chain = Chain.RIPPLE__MAINNET + val result = RippleChainSpecific.validateSettings(serverStateOk.toByteArray(), chain) + assertThat(result).isEqualTo(ValidateUpstreamSettingsResult.UPSTREAM_VALID) + } +}