Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions src/main/scala/units/ConsensusClient.scala
Original file line number Diff line number Diff line change
Expand Up @@ -98,14 +98,14 @@ object ConsensusClient {
context.time,
context.wallet,
context.settings.blockchainSettings.functionalitySettings.unitsRegistryAddressParsed.explicitGet(),
blockObserver.loadBlock,
blockObserver.requestBlockFromPeers,
context.broadcastTransaction,
eluScheduler,
globalScheduler
)

private val blocksStreamCancelable: CancelableFuture[Unit] =
blockObserver.getBlockStream.foreach { case (ch, block) => elu.executionBlockReceived(block, ch) }(using globalScheduler)
blockObserver.blockStream.foreach { case (ch, block) => elu.executionBlockReceived(block, ch) }(using globalScheduler)

override def close(): Unit = {
blocksStreamCancelable.cancel()
Expand Down
45 changes: 16 additions & 29 deletions src/main/scala/units/ELUpdater.scala
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,8 @@ import com.wavesplatform.wallet.Wallet
import io.netty.channel.Channel
import io.netty.channel.group.DefaultChannelGroup
import kamon.Kamon
import monix.execution.Scheduler
import monix.execution.cancelables.SerialCancelable
import monix.execution.{CancelableFuture, Scheduler}
import play.api.libs.json.*
import units.ELUpdater.State.*
import units.ELUpdater.State.ChainStatus.{FollowingChain, Mining, WaitForNewChain}
Expand All @@ -41,6 +41,7 @@ import units.util.HexBytesConverter.toHexNoPrefix

import java.math.BigInteger
import scala.annotation.tailrec
import scala.concurrent.Future
import scala.concurrent.duration.*
import scala.util.*

Expand All @@ -52,7 +53,7 @@ class ELUpdater(
time: Time,
wallet: Wallet,
registryAddress: Option[Address],
requestBlockFromPeers: BlockHash => CancelableFuture[BlockWithChannel],
requestBlockFromPeers: BlockHash => Future[BlockWithChannel],
broadcastTx: Transaction => TracedResult[ValidationError, Boolean],
scheduler: Scheduler,
globalScheduler: Scheduler
Expand All @@ -74,7 +75,7 @@ class ELUpdater(
val now = time.correctedTime() / 1000
if (block.timestamp - now <= MaxTimeDrift) {
state match {
case WaitingForSyncHead(target, _) if block.hash == target.hash =>
case WaitingForSyncHead(target) if block.hash == target.hash =>
val syncStarted = for {
_ <- engineApiClient.newPayload(block.payload)
fcuStatus <- confirmBlock(target, target)
Expand Down Expand Up @@ -502,17 +503,15 @@ class ELUpdater(
else if (chainContractClient.getAllActualMiners.isEmpty) logger.debug("Waiting for at least one miner to join")
else {
val finalizedBlock = chainContractClient.getFinalizedBlock
logger.debug(s"Finalized block is ${finalizedBlock.hash}")
engineApiClient.getBlockByHash(finalizedBlock.hash) match {
case Left(error) => logger.error(s"Could not load finalized block", error)
case Left(error) => logger.error(s"Could not load finalized block $finalizedBlock", error)
case Right(Some(finalizedEcBlock)) =>
logger.trace(s"Finalized block ${finalizedBlock.hash} is at height ${finalizedEcBlock.height}")
(for {
newEpochInfo <- chainContractClient.calculateEpochInfo(blockchain)
mainChainInfo <- chainContractClient.getMainChainInfo.toRight("Can't get main chain info")
lastEcBlock <- engineApiClient.getLastExecutionBlock().leftMap(_.message)
} yield {
logger.trace(s"Following main chain ${mainChainInfo.id}")
logger.trace(s"Finalized block ${finalizedBlock.hash} is at height ${finalizedEcBlock.height}, following main chain ${mainChainInfo.id}")
val fullValidationStatus = FullValidationStatus(
lastValidatedBlock = finalizedBlock,
lastElWithdrawalIndex = None
Expand All @@ -534,7 +533,8 @@ class ELUpdater(
)
case Right(None) =>
logger.trace(s"Finalized block ${finalizedBlock.hash} is not in EC, requesting from peers")
setState("updateStartingState", WaitingForSyncHead(finalizedBlock, requestAndProcessBlock(finalizedBlock.hash)))
requestAndProcessBlock(finalizedBlock.hash)
setState("updateStartingState", WaitingForSyncHead(finalizedBlock))
}
}
}
Expand All @@ -544,7 +544,7 @@ class ELUpdater(
_ <- Either.raiseWhen(chainContractClient.isStopped)(s"Chain $contractAddress is stopped")
_ <- registryAddress.toLeft(()).leftFlatMap { addr =>
Either.raiseUnless(blockchain.accountData(addr, registryKey(contractAddress)).contains(BooleanDataEntry(registryKey(contractAddress), true)))(
s"Chain ${contractAddress} is not enabled in the registry $addr"
s"Chain $contractAddress is not enabled in the registry $addr"
)
}
} yield ()) match {
Expand Down Expand Up @@ -767,7 +767,8 @@ class ELUpdater(
requestMainChainBlock()
case Right(None) =>
logger.trace(s"Finalized block ${finalizedBlock.hash} is not in EC, requesting from peers")
setState("updateWorkingState, finalized", WaitingForSyncHead(finalizedBlock, requestAndProcessBlock(finalizedBlock.hash)))
requestAndProcessBlock(finalizedBlock.hash)
setState("updateWorkingState, finalized", WaitingForSyncHead(finalizedBlock))
}
}

Expand Down Expand Up @@ -837,12 +838,11 @@ class ELUpdater(
}
}

private def requestAndProcessBlock(hash: BlockHash): CancelableFuture[(Channel, NetworkL2Block)] = {
requestBlockFromPeers(hash).andThen {
private def requestAndProcessBlock(hash: BlockHash): Unit =
requestBlockFromPeers(hash).onComplete {
case Success((ch, block)) => executionBlockReceived(block, ch)
case Failure(exception) => logger.error(s"Error requesting block $hash from peers", exception)
}(using globalScheduler)
}

private def updateToFollowChain(
prevState: Working[ChainStatus],
Expand Down Expand Up @@ -1157,23 +1157,11 @@ class ELUpdater(
confirmBlockAndFollowChain(networkBlock.toEcBlock, prevState, nodeChainInfo, returnToMainChainInfo)
}

private def findBlockChild(parent: BlockHash, lastBlockHash: BlockHash): Either[String, ContractBlock] = {
@tailrec
def loop(b: BlockHash): Option[ContractBlock] = chainContractClient.getBlock(b) match {
case None => None
case Some(cb) =>
if (cb.parentHash == parent) Some(cb)
else loop(cb.parentHash)
}

loop(lastBlockHash).toRight(s"Could not find child of $parent")
}

@tailrec
private def maybeRequestNextBlock(prevState: Working[FollowingChain], finalizedBlock: ContractBlock): Working[FollowingChain] = {
if (prevState.lastEcBlock.height < prevState.chainStatus.nodeChainInfo.lastBlock.height) {
logger.debug(s"EC chain is not synced, trying to find next block to request")
findBlockChild(prevState.lastEcBlock.hash, prevState.chainStatus.nodeChainInfo.lastBlock.hash) match {
chainContractClient.findBlockChild(prevState.lastEcBlock.hash, prevState.chainStatus.nodeChainInfo.lastBlock) match {
case Left(error) =>
logger.error(s"Could not find child of ${prevState.lastEcBlock.hash} on contract: $error")
prevState
Expand Down Expand Up @@ -1807,9 +1795,8 @@ object ELUpdater {
}
}

case class WaitingForSyncHead(target: ContractBlock, task: CancelableFuture[BlockWithChannel]) extends State {
override def toString: String = s"WaitingForSyncHead($target)"
}
case class WaitingForSyncHead(target: ContractBlock) extends State

case class SyncingToFinalizedBlock(target: BlockHash) extends State
}

Expand Down
53 changes: 33 additions & 20 deletions src/main/scala/units/client/contract/ChainContractClient.scala
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ trait ChainContractClient {
blockMeta <- getBlock(hash)
} yield blockMeta

def getFirstBlockMeta(chainId: Long): Option[ContractBlock] =
private def getFirstBlockMeta(chainId: Long): Option[ContractBlock] =
for {
hash <- getFirstBlockHash(chainId)
blockMeta <- getBlock(hash)
Expand All @@ -55,7 +55,7 @@ trait ChainContractClient {
}
}

def getElRewardAddress(miner: Address): Option[EthAddress] = getElRewardAddress(ByteStr(miner.bytes))
private def getElRewardAddress(miner: Address): Option[EthAddress] = getElRewardAddress(ByteStr(miner.bytes))
private def getElRewardAddress(minerAddress: ByteStr): Option[EthAddress] =
extractData(s"miner_${minerAddress}_RewardAddress")
.orElse(extractData(s"miner${minerAddress}RewardAddress"))
Expand Down Expand Up @@ -111,13 +111,13 @@ trait ChainContractClient {
}
}

def getLastChainId: Long =
private def getLastChainId: Long =
getLongData("lastChainId").getOrElse(DefaultMainChainId)

def getFirstValidAltChainId: Long =
private def getFirstValidAltChainId: Long =
getLongData("firstValidAltChainId").getOrElse(DefaultMainChainId)

def getMainChainIdOpt: Option[Long] =
private def getMainChainIdOpt: Option[Long] =
getLongData(MainChainIdKey)

def getMainChainId: Long =
Expand Down Expand Up @@ -153,7 +153,7 @@ trait ChainContractClient {
else fail(s"Can't get chain $chainId info, one of fields is empty, first block: $firstBlock, last block: $lastBlock")
}

def getFinalizedBlockHash: BlockHash =
private def getFinalizedBlockHash: BlockHash =
getStringData("finalizedBlock")
.map(hash => BlockHash(s"0x$hash"))
.getOrElse(throw new IllegalStateException("Can't get finalized block hash: not found at contract"))
Expand Down Expand Up @@ -288,7 +288,7 @@ trait ChainContractClient {

// Asset transfer, before strict transfers activation
// {destElAddressHex with 0x}_{fromClAddressHex with 0x}_{amount}_{assetRegistryIndex}
case Array(EthAddress(destElAddress), EthAddress(fromAddress), rawAmount, AssetIndex(assetIndex)) => {
case Array(EthAddress(destElAddress), EthAddress(fromAddress), rawAmount, AssetIndex(assetIndex)) =>
val asset = getRegisteredAsset(assetIndex)
val assetData = getRegisteredAssetData(asset)

Expand All @@ -299,15 +299,14 @@ trait ChainContractClient {
to = destElAddress,
amount =
try WAmount(rawAmount).scale(assetData.exponent)
catch { case e: ArithmeticException => fail(s"Expected an integer amount of a native transfer, got: ${rawAmount}", e) },
catch { case e: ArithmeticException => fail(s"Expected an integer amount of a native transfer, got: $rawAmount", e) },
tokenAddress = assetData.erc20Address,
asset
)
}

// Asset transfer, after strict transfers activation
// {epoch}_{destElAddressHex with 0x}_{fromClAddressHex with 0x}_{amount}_{assetRegistryIndex}
case Array(Epoch(epoch), EthAddress(destElAddress), EthAddress(fromAddress), rawAmount, AssetIndex(assetIndex)) => {
case Array(Epoch(epoch), EthAddress(destElAddress), EthAddress(fromAddress), rawAmount, AssetIndex(assetIndex)) =>
val asset = getRegisteredAsset(assetIndex)
val assetData = getRegisteredAssetData(asset)

Expand All @@ -318,17 +317,16 @@ trait ChainContractClient {
to = destElAddress,
amount =
try WAmount(rawAmount).scale(assetData.exponent)
catch { case e: ArithmeticException => fail(s"Expected an integer amount of a native transfer, got: ${rawAmount}", e) },
catch { case e: ArithmeticException => fail(s"Expected an integer amount of a native transfer, got: $rawAmount", e) },
tokenAddress = assetData.erc20Address,
asset
)
}

case _ => fail(s"Expected one of ContractTransfer variants in a transfer key '$key', got: $raw")
}
}

def getRegisteredAssetData(asset: Asset): Registry.RegisteredAsset = {
private def getRegisteredAssetData(asset: Asset): Registry.RegisteredAsset = {
val key = s"assetRegistry_${Registry.stringifyAsset(asset)}"
val raw = getStringData(key).getOrElse(fail(s"Can't find a registered asset $asset at $key"))
val parts = raw.split(Sep)
Expand All @@ -351,7 +349,7 @@ trait ChainContractClient {
.map(getRegisteredAssetData)
.toList

def getPrevEpochLastBlockHash(thisEpoch: Int): Either[String, Option[BlockHash]] = {
private def getPrevEpochLastBlockHash(thisEpoch: Int): Either[String, Option[BlockHash]] = {
@tailrec
def loop(curEpochNumber: Int): Either[String, Option[BlockHash]] = {
if (curEpochNumber <= 0) {
Expand Down Expand Up @@ -397,21 +395,36 @@ trait ChainContractClient {
}
}

def findBlockChild(lastExecutionBlockHash: BlockHash, lastBlock: ContractBlock): Either[String, ContractBlock] = {
@tailrec
def loop(b: BlockHash): Option[ContractBlock] = this.getBlock(b) match {
case None => None
case Some(cb) =>
if (cb.parentHash == lastExecutionBlockHash) Some(cb)
else loop(cb.parentHash)
}

this
.getBlock(lastExecutionBlockHash)
.toRight(s"Last EC block $lastExecutionBlockHash not found on contract")
.flatMap(_ => loop(lastBlock.hash).toRight(s"Could not find child of $lastExecutionBlockHash"))
}

private def getLastBlockHash(chainId: Long): Option[BlockHash] = getChainMeta(chainId).map(_._2)

protected def getFirstBlockHash(chainId: Long): Option[BlockHash] =
getBlockHash(s"chain${chainId}FirstBlock")

protected def getBinaryData(key: String): Option[ByteStr] =
private def getBinaryData(key: String): Option[ByteStr] =
extractBinaryValue(key, extractData(key))

protected def getStringData(key: String): Option[String] =
extractStringValue(key, extractData(key))

protected def getLongData(key: String): Option[Long] =
private def getLongData(key: String): Option[Long] =
extractLongValue(key, extractData(key))

protected def getBooleanData(key: String): Option[Boolean] =
private def getBooleanData(key: String): Option[Boolean] =
extractBooleanValue(key, extractData(key))

private def extractLongValue(context: String, extractedDataEntry: Option[DataEntry[?]]): Option[Long] =
Expand All @@ -437,13 +450,13 @@ trait ChainContractClient {
}

object ChainContractClient {
val MinMinerBalance: Long = 20000_00000000L
val DefaultMainChainId = 0L
private val MinMinerBalance: Long = 20000_00000000L
val DefaultMainChainId = 0L

private val AllMinersKey = "allMiners"
private val MainChainIdKey = "mainChainId"
private val BlockHashBytesSize = 32
val Sep = ","
private val Sep = ","

private class InconsistentContractData(message: String, cause: Throwable = null)
extends IllegalStateException(s"Probably, you have to upgrade your client. $message", cause)
Expand Down
12 changes: 5 additions & 7 deletions src/main/scala/units/network/BlocksObserver.scala
Original file line number Diff line number Diff line change
@@ -1,15 +1,13 @@
package units.network

import units.network.BlocksObserverImpl.BlockWithChannel
import com.wavesplatform.network.ChannelObservable
import monix.eval.Task
import monix.execution.CancelableFuture
import units.network.BlocksObserverImpl.BlockWithChannel
import units.{BlockHash, NetworkL2Block}

trait BlocksObserver {
def getBlockStream: ChannelObservable[NetworkL2Block]
import scala.concurrent.Future

def requestBlock(req: BlockHash): Task[BlockWithChannel]
trait BlocksObserver {
def blockStream: ChannelObservable[NetworkL2Block]

def loadBlock(req: BlockHash): CancelableFuture[BlockWithChannel]
def requestBlockFromPeers(req: BlockHash): Future[BlockWithChannel]
}
Loading