eth_subscribe: fix preparing filter (#635)
This commit is contained in:
@@ -37,6 +37,7 @@ open class ConnectLogs(
|
|||||||
}
|
}
|
||||||
|
|
||||||
constructor(upstream: Multistream, scheduler: Scheduler) : this(upstream, ConnectBlockUpdates(upstream, scheduler))
|
constructor(upstream: Multistream, scheduler: Scheduler) : this(upstream, ConnectBlockUpdates(upstream, scheduler))
|
||||||
|
|
||||||
private val produceLogs = ProduceLogs(upstream)
|
private val produceLogs = ProduceLogs(upstream)
|
||||||
|
|
||||||
fun start(matcher: Selector.Matcher): Flux<LogMessage> {
|
fun start(matcher: Selector.Matcher): Flux<LogMessage> {
|
||||||
@@ -65,12 +66,15 @@ open class ConnectLogs(
|
|||||||
logs.filter {
|
logs.filter {
|
||||||
val goodAddress =
|
val goodAddress =
|
||||||
sortedAddresses.isEmpty() || sortedAddresses.binarySearch(it.address, ADDR_COMPARATOR) >= 0
|
sortedAddresses.isEmpty() || sortedAddresses.binarySearch(it.address, ADDR_COMPARATOR) >= 0
|
||||||
val goodTopic = sortedTopics.isEmpty() || (
|
val goodTopic = when {
|
||||||
it.topics.isNotEmpty() && sortedTopics.binarySearch(
|
sortedTopics.isEmpty() -> true
|
||||||
it.topics[0],
|
it.topics.size < sortedTopics.size -> false
|
||||||
TOPIC_COMPARATOR,
|
else -> sortedTopics.indices.all { index ->
|
||||||
) >= 0
|
it.topics[index].let { logTopic ->
|
||||||
)
|
sortedTopics.binarySearch(logTopic, TOPIC_COMPARATOR) >= 0
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
goodAddress && goodTopic
|
goodAddress && goodTopic
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -88,6 +88,22 @@ class ConnectLogsSpec extends Specification {
|
|||||||
"upstream"
|
"upstream"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
def log5 = new LogMessage(
|
||||||
|
Address.from("0x63bc4a36c66c64acb3d695298d492e8c1d909d3f"),
|
||||||
|
BlockHash.empty(),
|
||||||
|
100L,
|
||||||
|
HexData.empty(),
|
||||||
|
1L,
|
||||||
|
[
|
||||||
|
Hex32.from("0x952ba7f163c4a11628f55a4df523b3efddf252ad1be2c89b69c2b068fc378daa"),
|
||||||
|
Hex32.from("0x00000000000000000000000088e6a0c2ddd26feeb64f039a2c41296fcb3f5640")
|
||||||
|
],
|
||||||
|
TransactionId.from("0xb5e554178a94fd993111f2ae64cb708cb0899d7b5182024e70d5c468164a8bec"),
|
||||||
|
1L,
|
||||||
|
false,
|
||||||
|
"upstream"
|
||||||
|
)
|
||||||
|
|
||||||
def "Filter is empty"() {
|
def "Filter is empty"() {
|
||||||
setup:
|
setup:
|
||||||
def connectLogs = new ConnectLogs(TestingCommons.emptyMultistream(), Schedulers.boundedElastic())
|
def connectLogs = new ConnectLogs(TestingCommons.emptyMultistream(), Schedulers.boundedElastic())
|
||||||
@@ -152,4 +168,26 @@ class ConnectLogsSpec extends Specification {
|
|||||||
act.size() == 1
|
act.size() == 1
|
||||||
act[0] == log3
|
act[0] == log3
|
||||||
}
|
}
|
||||||
|
|
||||||
|
def "Filter by address and two topics"() {
|
||||||
|
setup:
|
||||||
|
def connectLogs = new ConnectLogs(TestingCommons.emptyMultistream(), Schedulers.boundedElastic())
|
||||||
|
when:
|
||||||
|
def input = Flux.fromIterable([
|
||||||
|
log1, log2, log3, log4, log5
|
||||||
|
])
|
||||||
|
|
||||||
|
def act = input.transform(connectLogs.filtered(
|
||||||
|
[Address.from("0x63bc4a36c66c64acb3d695298d492e8c1d909d3f")],
|
||||||
|
[
|
||||||
|
Hex32.from("0x952ba7f163c4a11628f55a4df523b3efddf252ad1be2c89b69c2b068fc378daa"),
|
||||||
|
Hex32.from("0x00000000000000000000000088e6a0c2ddd26feeb64f039a2c41296fcb3f5640"),
|
||||||
|
]
|
||||||
|
))
|
||||||
|
.collectList().block()
|
||||||
|
|
||||||
|
then:
|
||||||
|
act.size() == 1
|
||||||
|
act[0] == log5
|
||||||
|
}
|
||||||
}
|
}
|
||||||
Reference in New Issue
Block a user