{-# LANGUAGE BangPatterns #-}
{-# LANGUAGE DuplicateRecordFields #-}
{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE LambdaCase #-}
{-# LANGUAGE NamedFieldPuns #-}
{-# LANGUAGE RecordWildCards #-}
{-# LANGUAGE ScopedTypeVariables #-}

module Ouroboros.Consensus.NodeKernel
  ( -- * Node kernel
    MempoolCapacityBytesOverride (..)
  , NodeKernel (..)
  , NodeKernelArgs (..)
  , TraceForgeEvent (..)
  , getImmTipSlot
  , getMempoolReader
  , getMempoolWriter
  , getPeersFromCurrentLedger
  , getPeersFromCurrentLedgerAfterSlot
  , initNodeKernel
  , toConsensusMode
  ) where

import Cardano.Base.FeatureFlags (CardanoFeatureFlag)
import Cardano.Network.ConsensusMode (ConsensusMode (..))
import Cardano.Network.LedgerStateJudgement (LedgerStateJudgement (..))
import Cardano.Network.NodeToNode
  ( ConnectionId
  , MiniProtocolParameters (..)
  )
import Cardano.Network.PeerSelection.Bootstrap (UseBootstrapPeers)
import Cardano.Network.PeerSelection.LocalRootPeers
  ( OutboundConnectionsState (..)
  )
import qualified Control.Concurrent.Class.MonadSTM as LazySTM
import qualified Control.Concurrent.Class.MonadSTM.Strict as StrictSTM
import Control.DeepSeq (force)
import Control.Monad
import qualified Control.Monad.Class.MonadTimer.SI as SI
import Control.ResourceRegistry
import Control.Tracer
import Data.Bifunctor (second)
import Data.Data (Typeable)
import Data.Either (partitionEithers)
import Data.Foldable (traverse_)
import Data.Function (on)
import Data.Functor ((<&>))
import Data.Hashable (Hashable)
import Data.List.NonEmpty (NonEmpty)
import Data.Set (Set)
import qualified Data.Text as Text
import Data.Void (Void)
import Ouroboros.Consensus.Block hiding (blockMatchesHeader)
import Ouroboros.Consensus.BlockchainTime
import Ouroboros.Consensus.Config
import Ouroboros.Consensus.Genesis.Governor (gddWatcher)
import Ouroboros.Consensus.HeaderValidation
import Ouroboros.Consensus.Ledger.Abstract
import Ouroboros.Consensus.Ledger.Extended
import Ouroboros.Consensus.Ledger.SupportsMempool
import Ouroboros.Consensus.Ledger.SupportsPeerSelection
import Ouroboros.Consensus.Mempool
import qualified Ouroboros.Consensus.MiniProtocol.BlockFetch.ClientInterface as BlockFetchClientInterface
import Ouroboros.Consensus.MiniProtocol.ChainSync.Client
  ( ChainSyncClientHandle (..)
  , ChainSyncClientHandleCollection (..)
  , ChainSyncState (..)
  , newChainSyncClientHandleCollection
  )
import Ouroboros.Consensus.MiniProtocol.ChainSync.Client.HistoricityCheck
  ( HistoricityCheck
  )
import Ouroboros.Consensus.MiniProtocol.ChainSync.Client.InFutureCheck
  ( SomeHeaderInFutureCheck
  )
import Ouroboros.Consensus.Node.GSM (GsmNodeKernelArgs (..))
import qualified Ouroboros.Consensus.Node.GSM as GSM
import Ouroboros.Consensus.Node.Genesis
  ( GenesisNodeKernelArgs (..)
  , LoEAndGDDConfig (..)
  , LoEAndGDDNodeKernelArgs (..)
  , setGetLoEFragment
  )
import Ouroboros.Consensus.Node.Run
import Ouroboros.Consensus.Node.Tracers
import Ouroboros.Consensus.NodeKernel.Forge (forge)
import Ouroboros.Consensus.Protocol.Abstract
import Ouroboros.Consensus.Storage.ChainDB.API
  ( ChainDB
  )
import qualified Ouroboros.Consensus.Storage.ChainDB.API as ChainDB
import Ouroboros.Consensus.Storage.ChainDB.Init (InitChainDB)
import qualified Ouroboros.Consensus.Storage.ChainDB.Init as InitChainDB
import Ouroboros.Consensus.Util.AnchoredFragment
  ( preferAnchoredCandidate
  )
import Ouroboros.Consensus.Util.EarlyExit
import Ouroboros.Consensus.Util.IOLike
import Ouroboros.Consensus.Util.LeakyBucket
  ( atomicallyWithMonotonicTime
  )
import Ouroboros.Consensus.Util.Orphans ()
import Ouroboros.Consensus.Util.STM
import qualified Ouroboros.Network.AnchoredFragment as AF
import Ouroboros.Network.Block (castTip, tipFromHeader)
import Ouroboros.Network.BlockFetch
import Ouroboros.Network.BlockFetch.ClientState
  ( mapTraceFetchClientState
  )
import Ouroboros.Network.BlockFetch.Decision.Trace
  ( TraceDecisionEvent (..)
  )
import Ouroboros.Network.PeerSelection.Governor.Types
  ( PublicPeerSelectionState
  )
import Ouroboros.Network.PeerSharing
  ( PeerSharingAPI
  , PeerSharingRegistry
  , newPeerSharingAPI
  , newPeerSharingRegistry
  , ps_POLICY_PEER_SHARE_MAX_PEERS
  , ps_POLICY_PEER_SHARE_STICKY_TIME
  )
import Ouroboros.Network.TxSubmission.Inbound.V1
  ( TxSubmissionInitDelay
  , TxSubmissionMempoolWriter
  )
import qualified Ouroboros.Network.TxSubmission.Inbound.V1 as Inbound
import Ouroboros.Network.TxSubmission.Inbound.V2 (TxDecisionPolicy)
import Ouroboros.Network.TxSubmission.Inbound.V2.Registry
  ( PeerTxRegistry
  , SharedTxStateVar
  , TxSubmissionCountersVar
  , newPeerTxRegistry
  , newSharedTxStateVar
  , newTxSubmissionCountersVar
  , txCountersThreadV2
  )
import Ouroboros.Network.TxSubmission.Inbound.V2.Types
  ( emptySharedTxState
  )
import Ouroboros.Network.TxSubmission.Mempool.Reader
  ( TxSubmissionMempoolReader
  )
import qualified Ouroboros.Network.TxSubmission.Mempool.Reader as MempoolReader
import System.Random (StdGen)

{-------------------------------------------------------------------------------
  Relay node
-------------------------------------------------------------------------------}

-- | Interface against running relay node
data NodeKernel m addrNTN addrNTC blk = NodeKernel
  { forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> ChainDB m blk
getChainDB :: ChainDB m blk
  -- ^ The 'ChainDB' of the node
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> Mempool m blk
getMempool :: Mempool m blk
  -- ^ The node's mempool
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> TopLevelConfig blk
getTopLevelConfig :: TopLevelConfig blk
  -- ^ The node's top-level static configuration
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk
-> FetchClientRegistry
     (ConnectionId addrNTN) (HeaderWithTime blk) blk m
getFetchClientRegistry :: FetchClientRegistry (ConnectionId addrNTN) (HeaderWithTime blk) blk m
  -- ^ The fetch client registry, used for the block fetch clients.
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk
-> KeepAliveRegistry (ConnectionId addrNTN) m
getKeepAliveRegistry :: KeepAliveRegistry (ConnectionId addrNTN) m
  -- ^ The keep alive registry, used for the block fetch clients.
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> STM m FetchMode
getFetchMode :: STM m FetchMode
  -- ^ The fetch mode, used by diffusion.
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> STM m GsmState
getGsmState :: STM m GSM.GsmState
  -- ^ The GSM state, used by diffusion. A ledger judgement can be derived
  -- from it with 'GSM.gsmStateToLedgerJudgement'.
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk
-> ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
getChainSyncHandles :: ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
  -- ^ The kill handle and exposed state for each ChainSync client.
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> PeerSharingRegistry addrNTN m
getPeerSharingRegistry :: PeerSharingRegistry addrNTN m
  -- ^ Read the current peer sharing registry, used for interacting with
  -- the PeerSharing protocol
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk
-> Tracers m (ConnectionId addrNTN) addrNTC blk
getTracers :: Tracers m (ConnectionId addrNTN) addrNTC blk
  -- ^ The node's tracers
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> [MkBlockForging m blk] -> m ()
setBlockForging :: [MkBlockForging m blk] -> m ()
  -- ^ Set block forging
  --
  -- When set with the empty list '[]' block forging will be disabled.
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> PeerSharingAPI addrNTN StdGen m
getPeerSharingAPI :: PeerSharingAPI addrNTN StdGen m
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk
-> StrictTVar m OutboundConnectionsState
getOutboundConnectionsState ::
      StrictTVar m OutboundConnectionsState
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> DiffusionPipeliningSupport
getDiffusionPipeliningSupport ::
      DiffusionPipeliningSupport
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> BlockchainTime m
getBlockchainTime :: BlockchainTime m
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk
-> PeerTxRegistry m (ConnectionId addrNTN)
getPeerTxRegistry :: PeerTxRegistry m (ConnectionId addrNTN)
  -- ^ Per-peer tx-submission state: in-flight tracking and counters.
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk
-> SharedTxStateVar m (ConnectionId addrNTN) (GenTxId blk)
getSharedTxStateVar :: SharedTxStateVar m (ConnectionId addrNTN) (GenTxId blk)
  -- ^ Shared state of all `TxSubmission` clients.
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> TxSubmissionCountersVar m
getTxCountersVar :: TxSubmissionCountersVar m
  -- ^ Accumulator for tx-submission counters of disconnected peers.
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> TxDecisionPolicy
getTxDecisionPolicy :: TxDecisionPolicy
  -- ^ Tx-submission decision policy.
  }

-- | Arguments required when initializing a node
data NodeKernelArgs m addrNTN addrNTC blk = NodeKernelArgs
  { forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> Tracers m (ConnectionId addrNTN) addrNTC blk
tracers :: Tracers m (ConnectionId addrNTN) addrNTC blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> ResourceRegistry m
registry :: ResourceRegistry m
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> TopLevelConfig blk
cfg :: TopLevelConfig blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> Set CardanoFeatureFlag
featureFlags :: Set CardanoFeatureFlag
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> BlockchainTime m
btime :: BlockchainTime m
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> SystemTime m
systemTime :: SystemTime m
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> ChainDB m blk
chainDB :: ChainDB m blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> StorageConfig blk -> InitChainDB m blk -> m ()
initChainDB :: StorageConfig blk -> InitChainDB m blk -> m ()
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> SomeHeaderInFutureCheck m blk
chainSyncFutureCheck :: SomeHeaderInFutureCheck m blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> m GsmState -> HistoricityCheck m blk
chainSyncHistoricityCheck ::
      m GSM.GsmState ->
      HistoricityCheck m blk
  -- ^ See 'HistoricityCheck' for details.
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> Header blk -> SizeInBytes
blockFetchSize :: Header blk -> SizeInBytes
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> MempoolCapacityBytesOverride
mempoolCapacityOverride :: MempoolCapacityBytesOverride
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> Maybe MempoolTimeoutConfig
mempoolTimeoutConfig :: Maybe MempoolTimeoutConfig
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> MiniProtocolParameters
miniProtocolParameters :: MiniProtocolParameters
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> BlockFetchConfiguration
blockFetchConfiguration :: BlockFetchConfiguration
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> StdGen
keepAliveRng :: StdGen
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> GsmNodeKernelArgs m blk
gsmArgs :: GsmNodeKernelArgs m blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> STM m UseBootstrapPeers
getUseBootstrapPeers :: STM m UseBootstrapPeers
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> StdGen
peerSharingRng :: StdGen
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> TxSubmissionInitDelay
txSubmissionInitDelay :: TxSubmissionInitDelay
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> StrictTVar m (PublicPeerSelectionState addrNTN)
publicPeerSelectionStateVar ::
      StrictSTM.StrictTVar m (PublicPeerSelectionState addrNTN)
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> GenesisNodeKernelArgs m blk
genesisArgs :: GenesisNodeKernelArgs m blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> DiffusionPipeliningSupport
getDiffusionPipeliningSupport :: DiffusionPipeliningSupport
  }

initNodeKernel ::
  forall m addrNTN addrNTC blk.
  ( IOLike m
  , SI.MonadTimer m
  , RunNode blk
  , Ord addrNTN
  , Hashable addrNTN
  , Typeable addrNTN
  ) =>
  NodeKernelArgs m addrNTN addrNTC blk ->
  m (NodeKernel m addrNTN addrNTC blk)
initNodeKernel :: forall (m :: * -> *) addrNTN addrNTC blk.
(IOLike m, MonadTimer m, RunNode blk, Ord addrNTN,
 Hashable addrNTN, Typeable addrNTN) =>
NodeKernelArgs m addrNTN addrNTC blk
-> m (NodeKernel m addrNTN addrNTC blk)
initNodeKernel
  args :: NodeKernelArgs m addrNTN addrNTC blk
args@NodeKernelArgs
    { ResourceRegistry m
registry :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> ResourceRegistry m
registry :: ResourceRegistry m
registry
    , TopLevelConfig blk
cfg :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> TopLevelConfig blk
cfg :: TopLevelConfig blk
cfg
    , Tracers m (ConnectionId addrNTN) addrNTC blk
tracers :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> Tracers m (ConnectionId addrNTN) addrNTC blk
tracers :: Tracers m (ConnectionId addrNTN) addrNTC blk
tracers
    , ChainDB m blk
chainDB :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> ChainDB m blk
chainDB :: ChainDB m blk
chainDB
    , StorageConfig blk -> InitChainDB m blk -> m ()
initChainDB :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> StorageConfig blk -> InitChainDB m blk -> m ()
initChainDB :: StorageConfig blk -> InitChainDB m blk -> m ()
initChainDB
    , BlockFetchConfiguration
blockFetchConfiguration :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> BlockFetchConfiguration
blockFetchConfiguration :: BlockFetchConfiguration
blockFetchConfiguration
    , BlockchainTime m
btime :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> BlockchainTime m
btime :: BlockchainTime m
btime
    , GsmNodeKernelArgs m blk
gsmArgs :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> GsmNodeKernelArgs m blk
gsmArgs :: GsmNodeKernelArgs m blk
gsmArgs
    , StdGen
peerSharingRng :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> StdGen
peerSharingRng :: StdGen
peerSharingRng
    , StrictTVar m (PublicPeerSelectionState addrNTN)
publicPeerSelectionStateVar :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> StrictTVar m (PublicPeerSelectionState addrNTN)
publicPeerSelectionStateVar :: StrictTVar m (PublicPeerSelectionState addrNTN)
publicPeerSelectionStateVar
    , GenesisNodeKernelArgs m blk
genesisArgs :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> GenesisNodeKernelArgs m blk
genesisArgs :: GenesisNodeKernelArgs m blk
genesisArgs
    , DiffusionPipeliningSupport
getDiffusionPipeliningSupport :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> DiffusionPipeliningSupport
getDiffusionPipeliningSupport :: DiffusionPipeliningSupport
getDiffusionPipeliningSupport
    , MiniProtocolParameters
miniProtocolParameters :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> MiniProtocolParameters
miniProtocolParameters :: MiniProtocolParameters
miniProtocolParameters
    } = do
    -- using a lazy 'TVar', 'BlockForging' does not have a 'NoThunks' instance.
    blockForgingVar :: LazySTM.TMVar m [MkBlockForging m blk] <- [MkBlockForging m blk] -> m (TMVar m [MkBlockForging m blk])
forall a. a -> m (TMVar m a)
forall (m :: * -> *) a. MonadSTM m => a -> m (TMVar m a)
LazySTM.newTMVarIO []
    initChainDB (configStorage cfg) (InitChainDB.fromFull chainDB)

    st <- initInternalState args
    let IS
          { blockFetchInterface
          , fetchClientRegistry
          , keepAliveRegistry
          , mempool
          , peerSharingRegistry
          , varChainSyncHandles
          , varGsmState
          } = st

    varOutboundConnectionsState <- newTVarIO UntrustedState

    do
      let GsmNodeKernelArgs{..} = gsmArgs
          gsmTracerArgs =
            ( Tip (HeaderWithTime blk) -> Tip blk
forall {k1} {k2} (a :: k1) (b :: k2).
(HeaderHash a ~ HeaderHash b) =>
Tip a -> Tip b
castTip (Tip (HeaderWithTime blk) -> Tip blk)
-> ((AnchoredSeq
       (WithOrigin SlotNo)
       (Anchor (HeaderWithTime blk))
       (HeaderWithTime blk),
     b)
    -> Tip (HeaderWithTime blk))
-> (AnchoredSeq
      (WithOrigin SlotNo)
      (Anchor (HeaderWithTime blk))
      (HeaderWithTime blk),
    b)
-> Tip blk
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (Anchor (HeaderWithTime blk) -> Tip (HeaderWithTime blk))
-> (HeaderWithTime blk -> Tip (HeaderWithTime blk))
-> Either (Anchor (HeaderWithTime blk)) (HeaderWithTime blk)
-> Tip (HeaderWithTime blk)
forall a c b. (a -> c) -> (b -> c) -> Either a b -> c
either Anchor (HeaderWithTime blk) -> Tip (HeaderWithTime blk)
forall a b. (HeaderHash a ~ HeaderHash b) => Anchor a -> Tip b
AF.anchorToTip HeaderWithTime blk -> Tip (HeaderWithTime blk)
forall a. HasHeader a => a -> Tip a
tipFromHeader (Either (Anchor (HeaderWithTime blk)) (HeaderWithTime blk)
 -> Tip (HeaderWithTime blk))
-> ((AnchoredSeq
       (WithOrigin SlotNo)
       (Anchor (HeaderWithTime blk))
       (HeaderWithTime blk),
     b)
    -> Either (Anchor (HeaderWithTime blk)) (HeaderWithTime blk))
-> (AnchoredSeq
      (WithOrigin SlotNo)
      (Anchor (HeaderWithTime blk))
      (HeaderWithTime blk),
    b)
-> Tip (HeaderWithTime blk)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. AnchoredSeq
  (WithOrigin SlotNo)
  (Anchor (HeaderWithTime blk))
  (HeaderWithTime blk)
-> Either (Anchor (HeaderWithTime blk)) (HeaderWithTime blk)
forall v a b. Anchorable v a b => AnchoredSeq v a b -> Either a b
AF.head (AnchoredSeq
   (WithOrigin SlotNo)
   (Anchor (HeaderWithTime blk))
   (HeaderWithTime blk)
 -> Either (Anchor (HeaderWithTime blk)) (HeaderWithTime blk))
-> ((AnchoredSeq
       (WithOrigin SlotNo)
       (Anchor (HeaderWithTime blk))
       (HeaderWithTime blk),
     b)
    -> AnchoredSeq
         (WithOrigin SlotNo)
         (Anchor (HeaderWithTime blk))
         (HeaderWithTime blk))
-> (AnchoredSeq
      (WithOrigin SlotNo)
      (Anchor (HeaderWithTime blk))
      (HeaderWithTime blk),
    b)
-> Either (Anchor (HeaderWithTime blk)) (HeaderWithTime blk)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (AnchoredSeq
   (WithOrigin SlotNo)
   (Anchor (HeaderWithTime blk))
   (HeaderWithTime blk),
 b)
-> AnchoredSeq
     (WithOrigin SlotNo)
     (Anchor (HeaderWithTime blk))
     (HeaderWithTime blk)
forall a b. (a, b) -> a
fst
            , Tracers m (ConnectionId addrNTN) addrNTC blk
-> Tracer m (TraceGsmEvent (Tip blk))
forall remotePeer localPeer blk (f :: * -> *).
Tracers' remotePeer localPeer blk f -> f (TraceGsmEvent (Tip blk))
gsmTracer Tracers m (ConnectionId addrNTN) addrNTC blk
tracers
            )

      let gsm =
            ((AnchoredSeq
    (WithOrigin SlotNo)
    (Anchor (HeaderWithTime blk))
    (HeaderWithTime blk),
  LedgerState blk EmptyMK)
 -> Tip blk,
 Tracer m (TraceGsmEvent (Tip blk)))
-> GsmView
     m
     (ConnectionId addrNTN)
     (AnchoredSeq
        (WithOrigin SlotNo)
        (Anchor (HeaderWithTime blk))
        (HeaderWithTime blk),
      LedgerState blk EmptyMK)
     (ChainSyncState blk)
-> GsmEntryPoints m
forall (m :: * -> *) upstreamPeer selection tracedSelection
       candidate.
(MonadDelay m, MonadTimer m) =>
(selection -> tracedSelection,
 Tracer m (TraceGsmEvent tracedSelection))
-> GsmView m upstreamPeer selection candidate -> GsmEntryPoints m
GSM.realGsmEntryPoints
              ((AnchoredSeq
    (WithOrigin SlotNo)
    (Anchor (HeaderWithTime blk))
    (HeaderWithTime blk),
  LedgerState blk EmptyMK)
 -> Tip blk,
 Tracer m (TraceGsmEvent (Tip blk)))
forall {b}.
((AnchoredSeq
    (WithOrigin SlotNo)
    (Anchor (HeaderWithTime blk))
    (HeaderWithTime blk),
  b)
 -> Tip blk,
 Tracer m (TraceGsmEvent (Tip blk)))
gsmTracerArgs
              GSM.GsmView
                { antiThunderingHerd :: Maybe StdGen
GSM.antiThunderingHerd = StdGen -> Maybe StdGen
forall a. a -> Maybe a
Just StdGen
gsmAntiThunderingHerd
                , getCandidateOverSelection :: STM
  m
  ((AnchoredSeq
      (WithOrigin SlotNo)
      (Anchor (HeaderWithTime blk))
      (HeaderWithTime blk),
    LedgerState blk EmptyMK)
   -> ChainSyncState blk -> CandidateVersusSelection)
GSM.getCandidateOverSelection = do
                    weights <- ChainDB m blk -> STM m (WithFingerprint (PerasWeightSnapshot blk))
forall (m :: * -> *) blk.
ChainDB m blk -> STM m (WithFingerprint (PerasWeightSnapshot blk))
ChainDB.getPerasWeightSnapshot ChainDB m blk
chainDB
                    pure $ \(AnchoredSeq
  (WithOrigin SlotNo)
  (Anchor (HeaderWithTime blk))
  (HeaderWithTime blk)
headers, LedgerState blk EmptyMK
_lst) ChainSyncState blk
state ->
                      case AnchoredSeq
  (WithOrigin SlotNo)
  (Anchor (HeaderWithTime blk))
  (HeaderWithTime blk)
-> AnchoredSeq
     (WithOrigin SlotNo)
     (Anchor (HeaderWithTime blk))
     (HeaderWithTime blk)
-> Maybe (Point (HeaderWithTime blk))
forall block1 block2.
(HasHeader block1, HasHeader block2,
 HeaderHash block1 ~ HeaderHash block2) =>
AnchoredFragment block1
-> AnchoredFragment block2 -> Maybe (Point block1)
AF.intersectionPoint AnchoredSeq
  (WithOrigin SlotNo)
  (Anchor (HeaderWithTime blk))
  (HeaderWithTime blk)
headers (ChainSyncState blk
-> AnchoredSeq
     (WithOrigin SlotNo)
     (Anchor (HeaderWithTime blk))
     (HeaderWithTime blk)
forall blk.
ChainSyncState blk -> AnchoredFragment (HeaderWithTime blk)
csCandidate ChainSyncState blk
state) of
                        Maybe (Point (HeaderWithTime blk))
Nothing -> CandidateVersusSelection
GSM.CandidateDoesNotIntersect
                        Just{} ->
                          Bool -> CandidateVersusSelection
GSM.WhetherCandidateIsBetter (Bool -> CandidateVersusSelection)
-> Bool -> CandidateVersusSelection
forall a b. (a -> b) -> a -> b
$ -- precondition requires intersection
                            ShouldSwitch
  (Either
     (WithEmptyFragmentReasonForSwitch
        (WeightedSelectView (BlockProtocol blk)))
     (SelectViewReasonForSwitch (BlockProtocol blk)))
-> Bool
forall reason. ShouldSwitch reason -> Bool
shouldSwitch
                              ( BlockConfig blk
-> PerasWeightSnapshot blk
-> AnchoredSeq
     (WithOrigin SlotNo)
     (Anchor (HeaderWithTime blk))
     (HeaderWithTime blk)
-> AnchoredSeq
     (WithOrigin SlotNo)
     (Anchor (HeaderWithTime blk))
     (HeaderWithTime blk)
-> ShouldSwitch (ReasonForSwitch' blk)
forall blk (h :: * -> *) (h' :: * -> *).
(BlockSupportsProtocol blk, HasCallStack, GetHeader1 h,
 GetHeader1 h', HeaderHash (h blk) ~ HeaderHash blk,
 HeaderHash (h blk) ~ HeaderHash (h' blk), HasHeader (h blk),
 HasHeader (h' blk)) =>
BlockConfig blk
-> PerasWeightSnapshot blk
-> AnchoredFragment (h blk)
-> AnchoredFragment (h' blk)
-> ShouldSwitch (ReasonForSwitch' blk)
preferAnchoredCandidate
                                  (TopLevelConfig blk -> BlockConfig blk
forall blk. TopLevelConfig blk -> BlockConfig blk
configBlock TopLevelConfig blk
cfg)
                                  (WithFingerprint (PerasWeightSnapshot blk)
-> PerasWeightSnapshot blk
forall a. WithFingerprint a -> a
forgetFingerprint WithFingerprint (PerasWeightSnapshot blk)
weights)
                                  AnchoredSeq
  (WithOrigin SlotNo)
  (Anchor (HeaderWithTime blk))
  (HeaderWithTime blk)
headers
                                  (ChainSyncState blk
-> AnchoredSeq
     (WithOrigin SlotNo)
     (Anchor (HeaderWithTime blk))
     (HeaderWithTime blk)
forall blk.
ChainSyncState blk -> AnchoredFragment (HeaderWithTime blk)
csCandidate ChainSyncState blk
state)
                              )
                , peerIsIdle :: ChainSyncState blk -> Bool
GSM.peerIsIdle = ChainSyncState blk -> Bool
forall blk. ChainSyncState blk -> Bool
csIdling
                , durationUntilTooOld :: Maybe
  ((AnchoredSeq
      (WithOrigin SlotNo)
      (Anchor (HeaderWithTime blk))
      (HeaderWithTime blk),
    LedgerState blk EmptyMK)
   -> m DurationFromNow)
GSM.durationUntilTooOld =
                    Maybe (WrapDurationUntilTooOld m blk)
gsmDurationUntilTooOld
                      Maybe (WrapDurationUntilTooOld m blk)
-> (WrapDurationUntilTooOld m blk
    -> (AnchoredSeq
          (WithOrigin SlotNo)
          (Anchor (HeaderWithTime blk))
          (HeaderWithTime blk),
        LedgerState blk EmptyMK)
    -> m DurationFromNow)
-> Maybe
     ((AnchoredSeq
         (WithOrigin SlotNo)
         (Anchor (HeaderWithTime blk))
         (HeaderWithTime blk),
       LedgerState blk EmptyMK)
      -> m DurationFromNow)
forall (f :: * -> *) a b. Functor f => f a -> (a -> b) -> f b
<&> \WrapDurationUntilTooOld m blk
wd (AnchoredSeq
  (WithOrigin SlotNo)
  (Anchor (HeaderWithTime blk))
  (HeaderWithTime blk)
_headers, LedgerState blk EmptyMK
lst) ->
                        WrapDurationUntilTooOld m blk
-> WithOrigin SlotNo -> m DurationFromNow
forall (m :: * -> *) blk.
WrapDurationUntilTooOld m blk
-> WithOrigin SlotNo -> m DurationFromNow
GSM.getDurationUntilTooOld WrapDurationUntilTooOld m blk
wd (LedgerState blk EmptyMK -> WithOrigin SlotNo
forall (l :: LedgerStateKind) (mk :: MapKind).
GetTip l =>
l mk -> WithOrigin SlotNo
getTipSlot LedgerState blk EmptyMK
lst)
                , equivalent :: (AnchoredSeq
   (WithOrigin SlotNo)
   (Anchor (HeaderWithTime blk))
   (HeaderWithTime blk),
 LedgerState blk EmptyMK)
-> (AnchoredSeq
      (WithOrigin SlotNo)
      (Anchor (HeaderWithTime blk))
      (HeaderWithTime blk),
    LedgerState blk EmptyMK)
-> Bool
GSM.equivalent = Point (HeaderWithTime blk) -> Point (HeaderWithTime blk) -> Bool
forall a. Eq a => a -> a -> Bool
(==) (Point (HeaderWithTime blk) -> Point (HeaderWithTime blk) -> Bool)
-> ((AnchoredSeq
       (WithOrigin SlotNo)
       (Anchor (HeaderWithTime blk))
       (HeaderWithTime blk),
     LedgerState blk EmptyMK)
    -> Point (HeaderWithTime blk))
-> (AnchoredSeq
      (WithOrigin SlotNo)
      (Anchor (HeaderWithTime blk))
      (HeaderWithTime blk),
    LedgerState blk EmptyMK)
-> (AnchoredSeq
      (WithOrigin SlotNo)
      (Anchor (HeaderWithTime blk))
      (HeaderWithTime blk),
    LedgerState blk EmptyMK)
-> Bool
forall b c a. (b -> b -> c) -> (a -> b) -> a -> a -> c
`on` (AnchoredSeq
  (WithOrigin SlotNo)
  (Anchor (HeaderWithTime blk))
  (HeaderWithTime blk)
-> Point (HeaderWithTime blk)
forall block.
HasHeader block =>
AnchoredFragment block -> Point block
AF.headPoint (AnchoredSeq
   (WithOrigin SlotNo)
   (Anchor (HeaderWithTime blk))
   (HeaderWithTime blk)
 -> Point (HeaderWithTime blk))
-> ((AnchoredSeq
       (WithOrigin SlotNo)
       (Anchor (HeaderWithTime blk))
       (HeaderWithTime blk),
     LedgerState blk EmptyMK)
    -> AnchoredSeq
         (WithOrigin SlotNo)
         (Anchor (HeaderWithTime blk))
         (HeaderWithTime blk))
-> (AnchoredSeq
      (WithOrigin SlotNo)
      (Anchor (HeaderWithTime blk))
      (HeaderWithTime blk),
    LedgerState blk EmptyMK)
-> Point (HeaderWithTime blk)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (AnchoredSeq
   (WithOrigin SlotNo)
   (Anchor (HeaderWithTime blk))
   (HeaderWithTime blk),
 LedgerState blk EmptyMK)
-> AnchoredSeq
     (WithOrigin SlotNo)
     (Anchor (HeaderWithTime blk))
     (HeaderWithTime blk)
forall a b. (a, b) -> a
fst)
                , getChainSyncStates :: STM
  m (Map (ConnectionId addrNTN) (StrictTVar m (ChainSyncState blk)))
GSM.getChainSyncStates = (ChainSyncClientHandle m blk -> StrictTVar m (ChainSyncState blk))
-> Map (ConnectionId addrNTN) (ChainSyncClientHandle m blk)
-> Map (ConnectionId addrNTN) (StrictTVar m (ChainSyncState blk))
forall a b.
(a -> b)
-> Map (ConnectionId addrNTN) a -> Map (ConnectionId addrNTN) b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap ChainSyncClientHandle m blk -> StrictTVar m (ChainSyncState blk)
forall (m :: * -> *) blk.
ChainSyncClientHandle m blk -> StrictTVar m (ChainSyncState blk)
cschState (Map (ConnectionId addrNTN) (ChainSyncClientHandle m blk)
 -> Map (ConnectionId addrNTN) (StrictTVar m (ChainSyncState blk)))
-> STM m (Map (ConnectionId addrNTN) (ChainSyncClientHandle m blk))
-> STM
     m (Map (ConnectionId addrNTN) (StrictTVar m (ChainSyncState blk)))
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
-> STM m (Map (ConnectionId addrNTN) (ChainSyncClientHandle m blk))
forall peer (m :: * -> *) blk.
ChainSyncClientHandleCollection peer m blk
-> STM m (Map peer (ChainSyncClientHandle m blk))
cschcMap ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
varChainSyncHandles
                , getCurrentSelection :: STM
  m
  (AnchoredSeq
     (WithOrigin SlotNo)
     (Anchor (HeaderWithTime blk))
     (HeaderWithTime blk),
   LedgerState blk EmptyMK)
GSM.getCurrentSelection = do
                    headers <- ChainDB m blk
-> STM
     m
     (AnchoredSeq
        (WithOrigin SlotNo)
        (Anchor (HeaderWithTime blk))
        (HeaderWithTime blk))
forall (m :: * -> *) blk.
ChainDB m blk -> STM m (AnchoredFragment (HeaderWithTime blk))
ChainDB.getCurrentChainWithTime ChainDB m blk
chainDB
                    extLedgerState <- ChainDB.getCurrentLedger chainDB
                    return (headers, ledgerState extLedgerState)
                , minCaughtUpDuration :: NominalDiffTime
GSM.minCaughtUpDuration = NominalDiffTime
gsmMinCaughtUpDuration
                , setCaughtUpPersistentMark :: Bool -> m ()
GSM.setCaughtUpPersistentMark = \Bool
upd ->
                    (if Bool
upd then MarkerFileView m -> m ()
forall (m :: * -> *). MarkerFileView m -> m ()
GSM.touchMarkerFile else MarkerFileView m -> m ()
forall (m :: * -> *). MarkerFileView m -> m ()
GSM.removeMarkerFile)
                      MarkerFileView m
gsmMarkerFileView
                , writeGsmState :: GsmState -> m ()
GSM.writeGsmState = \GsmState
gsmState ->
                    (Time -> STM m ()) -> m ()
forall (m :: * -> *) b.
(MonadMonotonicTime m, MonadSTM m) =>
(Time -> STM m b) -> m b
atomicallyWithMonotonicTime ((Time -> STM m ()) -> m ()) -> (Time -> STM m ()) -> m ()
forall a b. (a -> b) -> a -> b
$ \Time
time -> do
                      StrictTVar m GsmState -> GsmState -> STM m ()
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
StrictTVar m a -> a -> STM m ()
writeTVar StrictTVar m GsmState
varGsmState GsmState
gsmState
                      handles <- ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
-> STM m (Map (ConnectionId addrNTN) (ChainSyncClientHandle m blk))
forall peer (m :: * -> *) blk.
ChainSyncClientHandleCollection peer m blk
-> STM m (Map peer (ChainSyncClientHandle m blk))
cschcMap ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
varChainSyncHandles
                      traverse_ (($ time) . ($ gsmState) . cschOnGsmStateChanged) handles
                , isHaaSatisfied :: STM m Bool
GSM.isHaaSatisfied = do
                    StrictTVar m OutboundConnectionsState
-> STM m OutboundConnectionsState
forall (m :: * -> *) a. MonadSTM m => StrictTVar m a -> STM m a
readTVar StrictTVar m OutboundConnectionsState
varOutboundConnectionsState STM m OutboundConnectionsState
-> (OutboundConnectionsState -> Bool) -> STM m Bool
forall (f :: * -> *) a b. Functor f => f a -> (a -> b) -> f b
<&> \case
                      -- See the upstream Haddocks for the exact conditions under
                      -- which the diffusion layer is in this state.
                      OutboundConnectionsState
TrustedStateWithExternalPeers -> Bool
True
                      OutboundConnectionsState
UntrustedState -> Bool
False
                }
      judgment <- GSM.gsmStateToLedgerJudgement <$> readTVarIO varGsmState
      void $ forkLinkedThread registry "NodeKernel.GSM" $ case judgment of
        LedgerStateJudgement
TooOld -> GsmEntryPoints m -> forall neverTerminates. m neverTerminates
forall (m :: * -> *).
GsmEntryPoints m -> forall neverTerminates. m neverTerminates
GSM.enterPreSyncing GsmEntryPoints m
gsm
        LedgerStateJudgement
YoungEnough -> GsmEntryPoints m -> forall neverTerminates. m neverTerminates
forall (m :: * -> *).
GsmEntryPoints m -> forall neverTerminates. m neverTerminates
GSM.enterCaughtUp GsmEntryPoints m
gsm

    peerSharingAPI <-
      newPeerSharingAPI
        publicPeerSelectionStateVar
        peerSharingRng
        ps_POLICY_PEER_SHARE_STICKY_TIME
        ps_POLICY_PEER_SHARE_MAX_PEERS

    sharedTxStateVar <- newSharedTxStateVar emptySharedTxState
    peerTxRegistry <- newPeerTxRegistry
    txCountersVar <- newTxSubmissionCountersVar mempty

    case gnkaLoEAndGDDArgs genesisArgs of
      LoEAndGDDConfig (LoEAndGDDNodeKernelArgs m blk)
LoEAndGDDDisabled -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
      LoEAndGDDEnabled LoEAndGDDNodeKernelArgs m blk
lgArgs -> do
        varLoEFragment <- AnchoredSeq
  (WithOrigin SlotNo)
  (Anchor (HeaderWithTime blk))
  (HeaderWithTime blk)
-> m (StrictTVar
        m
        (AnchoredSeq
           (WithOrigin SlotNo)
           (Anchor (HeaderWithTime blk))
           (HeaderWithTime blk)))
forall (m :: * -> *) a.
(HasCallStack, MonadSTM m, NoThunks a) =>
a -> m (StrictTVar m a)
newTVarIO (AnchoredSeq
   (WithOrigin SlotNo)
   (Anchor (HeaderWithTime blk))
   (HeaderWithTime blk)
 -> m (StrictTVar
         m
         (AnchoredSeq
            (WithOrigin SlotNo)
            (Anchor (HeaderWithTime blk))
            (HeaderWithTime blk))))
-> AnchoredSeq
     (WithOrigin SlotNo)
     (Anchor (HeaderWithTime blk))
     (HeaderWithTime blk)
-> m (StrictTVar
        m
        (AnchoredSeq
           (WithOrigin SlotNo)
           (Anchor (HeaderWithTime blk))
           (HeaderWithTime blk)))
forall a b. (a -> b) -> a -> b
$ Anchor (HeaderWithTime blk)
-> AnchoredSeq
     (WithOrigin SlotNo)
     (Anchor (HeaderWithTime blk))
     (HeaderWithTime blk)
forall v a b. Anchorable v a b => a -> AnchoredSeq v a b
AF.Empty Anchor (HeaderWithTime blk)
forall block. Anchor block
AF.AnchorGenesis
        setGetLoEFragment
          (readTVar varGsmState)
          (readTVar varLoEFragment)
          (lgnkaLoEFragmentTVar lgArgs)

        void $
          forkLinkedWatcher registry "NodeKernel.GDD" $
            gddWatcher
              cfg
              (gddTracer tracers)
              chainDB
              (lgnkaGDDRateLimit lgArgs)
              (readTVar varGsmState)
              (cschcMap varChainSyncHandles)
              varLoEFragment

    void $
      forkLinkedThread registry "NodeKernel.blockForging" $
        blockForgingController st (LazySTM.takeTMVar blockForgingVar)

    -- Run the block fetch logic in the background. This will call
    -- 'addFetchedBlock' whenever a new block is downloaded.
    void $
      forkLinkedThread registry "NodeKernel.blockFetchLogic" $
        blockFetchLogic
          (contramap castTraceFetchDecision $ blockFetchDecisionTracer tracers)
          (contramap (fmap castTraceFetchClientState) $ blockFetchClientTracer tracers)
          blockFetchInterface
          fetchClientRegistry
          keepAliveRegistry
          blockFetchConfiguration

    void $
      forkLinkedThread registry "NodeKernel.txCountersThreadV2" $
        txCountersThreadV2
          (txDecisionPolicy miniProtocolParameters)
          (txCountersTracer tracers)
          (txLogicTracer tracers)
          txCountersVar
          sharedTxStateVar
          peerTxRegistry

    return
      NodeKernel
        { getChainDB = chainDB
        , getMempool = mempool
        , getTopLevelConfig = cfg
        , getFetchClientRegistry = fetchClientRegistry
        , getKeepAliveRegistry = keepAliveRegistry
        , getFetchMode = readFetchMode blockFetchInterface
        , getGsmState = readTVar varGsmState
        , getChainSyncHandles = varChainSyncHandles
        , getPeerSharingRegistry = peerSharingRegistry
        , getTracers = tracers
        , setBlockForging = \[MkBlockForging m blk]
a -> STM m () -> m ()
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m () -> m ())
-> ([MkBlockForging m blk] -> STM m ())
-> [MkBlockForging m blk]
-> m ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. TMVar m [MkBlockForging m blk]
-> [MkBlockForging m blk] -> STM m ()
forall a. TMVar m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TMVar m a -> a -> STM m ()
LazySTM.putTMVar TMVar m [MkBlockForging m blk]
blockForgingVar ([MkBlockForging m blk] -> m ()) -> [MkBlockForging m blk] -> m ()
forall a b. (a -> b) -> a -> b
$! [MkBlockForging m blk]
a
        , getPeerSharingAPI = peerSharingAPI
        , getOutboundConnectionsState =
            varOutboundConnectionsState
        , getDiffusionPipeliningSupport
        , getBlockchainTime = btime
        , getPeerTxRegistry = peerTxRegistry
        , getSharedTxStateVar = sharedTxStateVar
        , getTxCountersVar = txCountersVar
        , getTxDecisionPolicy = txDecisionPolicy miniProtocolParameters
        }
   where
    blockForgingController ::
      InternalState m remotePeer localPeer blk ->
      STM m [MkBlockForging m blk] ->
      m Void
    blockForgingController :: forall remotePeer localPeer.
InternalState m remotePeer localPeer blk
-> STM m [MkBlockForging m blk] -> m Void
blockForgingController InternalState m remotePeer localPeer blk
st STM m [MkBlockForging m blk]
getBlockForging = [Thread m Void] -> m Void
go []
     where
      go :: [Thread m Void] -> m Void
      go :: [Thread m Void] -> m Void
go ![Thread m Void]
forgingThreads = do
        blockForging <- STM m [MkBlockForging m blk] -> m [MkBlockForging m blk]
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically STM m [MkBlockForging m blk]
getBlockForging
        traverse_ cancelThread forgingThreads
        blockForging' <- traverse (forkBlockForging st) blockForging
        go blockForging'

castTraceFetchDecision ::
  forall remotePeer blk.
  TraceDecisionEvent remotePeer (HeaderWithTime blk) -> TraceDecisionEvent remotePeer (Header blk)
castTraceFetchDecision :: forall remotePeer blk.
TraceDecisionEvent remotePeer (HeaderWithTime blk)
-> TraceDecisionEvent remotePeer (Header blk)
castTraceFetchDecision = \case
  PeersFetch [TraceLabelPeer
   remotePeer (FetchDecision [Point (HeaderWithTime blk)])]
xs -> [TraceLabelPeer remotePeer (FetchDecision [Point (Header blk)])]
-> TraceDecisionEvent remotePeer (Header blk)
forall peer header.
[TraceLabelPeer peer (FetchDecision [Point header])]
-> TraceDecisionEvent peer header
PeersFetch ((TraceLabelPeer
   remotePeer (FetchDecision [Point (HeaderWithTime blk)])
 -> TraceLabelPeer remotePeer (FetchDecision [Point (Header blk)]))
-> [TraceLabelPeer
      remotePeer (FetchDecision [Point (HeaderWithTime blk)])]
-> [TraceLabelPeer remotePeer (FetchDecision [Point (Header blk)])]
forall a b. (a -> b) -> [a] -> [b]
map ((FetchDecision [Point (HeaderWithTime blk)]
 -> FetchDecision [Point (Header blk)])
-> TraceLabelPeer
     remotePeer (FetchDecision [Point (HeaderWithTime blk)])
-> TraceLabelPeer remotePeer (FetchDecision [Point (Header blk)])
forall a b.
(a -> b)
-> TraceLabelPeer remotePeer a -> TraceLabelPeer remotePeer b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap (([Point (HeaderWithTime blk)] -> [Point (Header blk)])
-> FetchDecision [Point (HeaderWithTime blk)]
-> FetchDecision [Point (Header blk)]
forall b c a. (b -> c) -> Either a b -> Either a c
forall (p :: MapKind) b c a.
Bifunctor p =>
(b -> c) -> p a b -> p a c
second ((Point (HeaderWithTime blk) -> Point (Header blk))
-> [Point (HeaderWithTime blk)] -> [Point (Header blk)]
forall a b. (a -> b) -> [a] -> [b]
map Point (HeaderWithTime blk) -> Point (Header blk)
forall {k1} {k2} (b :: k1) (b' :: k2).
Coercible (HeaderHash b) (HeaderHash b') =>
Point b -> Point b'
castPoint))) [TraceLabelPeer
   remotePeer (FetchDecision [Point (HeaderWithTime blk)])]
xs)
  PeerStarvedUs remotePeer
peer -> remotePeer -> TraceDecisionEvent remotePeer (Header blk)
forall peer header. peer -> TraceDecisionEvent peer header
PeerStarvedUs remotePeer
peer

castTraceFetchClientState ::
  forall blk.
  HasHeader (Header blk) =>
  TraceFetchClientState (HeaderWithTime blk) -> TraceFetchClientState (Header blk)
castTraceFetchClientState :: forall blk.
HasHeader (Header blk) =>
TraceFetchClientState (HeaderWithTime blk)
-> TraceFetchClientState (Header blk)
castTraceFetchClientState = (HeaderWithTime blk -> Header blk)
-> TraceFetchClientState (HeaderWithTime blk)
-> TraceFetchClientState (Header blk)
forall h1 h2.
(HeaderHash h1 ~ HeaderHash h2, HasHeader h2) =>
(h1 -> h2) -> TraceFetchClientState h1 -> TraceFetchClientState h2
mapTraceFetchClientState HeaderWithTime blk -> Header blk
forall blk. HeaderWithTime blk -> Header blk
hwtHeader

{-------------------------------------------------------------------------------
  Internal node components
-------------------------------------------------------------------------------}

data InternalState m addrNTN addrNTC blk = IS
  { forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> Tracers m (ConnectionId addrNTN) addrNTC blk
tracers :: Tracers m (ConnectionId addrNTN) addrNTC blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> TopLevelConfig blk
cfg :: TopLevelConfig blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> ResourceRegistry m
registry :: ResourceRegistry m
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> BlockchainTime m
btime :: BlockchainTime m
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> ChainDB m blk
chainDB :: ChainDB m blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> BlockFetchConsensusInterface
     (ConnectionId addrNTN) (HeaderWithTime blk) blk m
blockFetchInterface ::
      BlockFetchConsensusInterface (ConnectionId addrNTN) (HeaderWithTime blk) blk m
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> FetchClientRegistry
     (ConnectionId addrNTN) (HeaderWithTime blk) blk m
fetchClientRegistry :: FetchClientRegistry (ConnectionId addrNTN) (HeaderWithTime blk) blk m
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> KeepAliveRegistry (ConnectionId addrNTN) m
keepAliveRegistry :: KeepAliveRegistry (ConnectionId addrNTN) m
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
varChainSyncHandles :: ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> StrictTVar m GsmState
varGsmState :: StrictTVar m GSM.GsmState
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> Mempool m blk
mempool :: Mempool m blk
  , forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> PeerSharingRegistry addrNTN m
peerSharingRegistry :: PeerSharingRegistry addrNTN m
  }

initInternalState ::
  forall m addrNTN addrNTC blk.
  ( IOLike m
  , SI.MonadTimer m
  , Ord addrNTN
  , Typeable addrNTN
  , RunNode blk
  ) =>
  NodeKernelArgs m addrNTN addrNTC blk ->
  m (InternalState m addrNTN addrNTC blk)
initInternalState :: forall (m :: * -> *) addrNTN addrNTC blk.
(IOLike m, MonadTimer m, Ord addrNTN, Typeable addrNTN,
 RunNode blk) =>
NodeKernelArgs m addrNTN addrNTC blk
-> m (InternalState m addrNTN addrNTC blk)
initInternalState
  NodeKernelArgs
    { Tracers m (ConnectionId addrNTN) addrNTC blk
tracers :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> Tracers m (ConnectionId addrNTN) addrNTC blk
tracers :: Tracers m (ConnectionId addrNTN) addrNTC blk
tracers
    , ChainDB m blk
chainDB :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> ChainDB m blk
chainDB :: ChainDB m blk
chainDB
    , ResourceRegistry m
registry :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> ResourceRegistry m
registry :: ResourceRegistry m
registry
    , TopLevelConfig blk
cfg :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> TopLevelConfig blk
cfg :: TopLevelConfig blk
cfg
    , Header blk -> SizeInBytes
blockFetchSize :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> Header blk -> SizeInBytes
blockFetchSize :: Header blk -> SizeInBytes
blockFetchSize
    , BlockchainTime m
btime :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> BlockchainTime m
btime :: BlockchainTime m
btime
    , MempoolCapacityBytesOverride
mempoolCapacityOverride :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk
-> MempoolCapacityBytesOverride
mempoolCapacityOverride :: MempoolCapacityBytesOverride
mempoolCapacityOverride
    , Maybe MempoolTimeoutConfig
mempoolTimeoutConfig :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> Maybe MempoolTimeoutConfig
mempoolTimeoutConfig :: Maybe MempoolTimeoutConfig
mempoolTimeoutConfig
    , GsmNodeKernelArgs m blk
gsmArgs :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> GsmNodeKernelArgs m blk
gsmArgs :: GsmNodeKernelArgs m blk
gsmArgs
    , STM m UseBootstrapPeers
getUseBootstrapPeers :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> STM m UseBootstrapPeers
getUseBootstrapPeers :: STM m UseBootstrapPeers
getUseBootstrapPeers
    , DiffusionPipeliningSupport
getDiffusionPipeliningSupport :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> DiffusionPipeliningSupport
getDiffusionPipeliningSupport :: DiffusionPipeliningSupport
getDiffusionPipeliningSupport
    , GenesisNodeKernelArgs m blk
genesisArgs :: forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernelArgs m addrNTN addrNTC blk -> GenesisNodeKernelArgs m blk
genesisArgs :: GenesisNodeKernelArgs m blk
genesisArgs
    } = do
    varGsmState <- do
      let GsmNodeKernelArgs{Maybe (WrapDurationUntilTooOld m blk)
StdGen
NominalDiffTime
MarkerFileView m
gsmMinCaughtUpDuration :: forall (m :: * -> *) blk.
GsmNodeKernelArgs m blk -> NominalDiffTime
gsmMarkerFileView :: forall (m :: * -> *) blk.
GsmNodeKernelArgs m blk -> MarkerFileView m
gsmDurationUntilTooOld :: forall (m :: * -> *) blk.
GsmNodeKernelArgs m blk -> Maybe (WrapDurationUntilTooOld m blk)
gsmAntiThunderingHerd :: forall (m :: * -> *) blk. GsmNodeKernelArgs m blk -> StdGen
gsmAntiThunderingHerd :: StdGen
gsmDurationUntilTooOld :: Maybe (WrapDurationUntilTooOld m blk)
gsmMarkerFileView :: MarkerFileView m
gsmMinCaughtUpDuration :: NominalDiffTime
..} = GsmNodeKernelArgs m blk
gsmArgs
      gsmState <-
        m (LedgerState blk EmptyMK)
-> Maybe (WrapDurationUntilTooOld m blk)
-> MarkerFileView m
-> m GsmState
forall blk (m :: * -> *).
(GetTip (LedgerState blk), Monad m) =>
m (LedgerState blk EmptyMK)
-> Maybe (WrapDurationUntilTooOld m blk)
-> MarkerFileView m
-> m GsmState
GSM.initializationGsmState
          (STM m (LedgerState blk EmptyMK) -> m (LedgerState blk EmptyMK)
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m (LedgerState blk EmptyMK) -> m (LedgerState blk EmptyMK))
-> STM m (LedgerState blk EmptyMK) -> m (LedgerState blk EmptyMK)
forall a b. (a -> b) -> a -> b
$ ExtLedgerState blk EmptyMK -> LedgerState blk EmptyMK
forall blk (mk :: MapKind).
ExtLedgerState blk mk -> LedgerState blk mk
ledgerState (ExtLedgerState blk EmptyMK -> LedgerState blk EmptyMK)
-> STM m (ExtLedgerState blk EmptyMK)
-> STM m (LedgerState blk EmptyMK)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> ChainDB m blk -> STM m (ExtLedgerState blk EmptyMK)
forall (m :: * -> *) blk.
ChainDB m blk -> STM m (ExtLedgerState blk EmptyMK)
ChainDB.getCurrentLedger ChainDB m blk
chainDB)
          Maybe (WrapDurationUntilTooOld m blk)
gsmDurationUntilTooOld
          MarkerFileView m
gsmMarkerFileView
      newTVarIO gsmState

    varChainSyncHandles <- atomically newChainSyncClientHandleCollection
    mempool <-
      openMempool
        registry
        (chainDBLedgerInterface chainDB)
        (configLedger cfg)
        mempoolCapacityOverride
        mempoolTimeoutConfig
        (mempoolTracer tracers)

    fetchClientRegistry <- newFetchClientRegistry
    keepAliveRegistry <- newKeepAliveRegistry

    let readFetchMode =
          ConsensusMode
-> BlockchainTime m
-> STM m (AnchoredFragment (Header blk))
-> STM m UseBootstrapPeers
-> STM m LedgerStateJudgement
-> STM m FetchMode
forall (m :: * -> *) blk.
(MonadSTM m, HasHeader blk) =>
ConsensusMode
-> BlockchainTime m
-> STM m (AnchoredFragment blk)
-> STM m UseBootstrapPeers
-> STM m LedgerStateJudgement
-> STM m FetchMode
BlockFetchClientInterface.readFetchModeDefault
            (LoEAndGDDConfig (LoEAndGDDNodeKernelArgs m blk) -> ConsensusMode
forall a. LoEAndGDDConfig a -> ConsensusMode
toConsensusMode (LoEAndGDDConfig (LoEAndGDDNodeKernelArgs m blk) -> ConsensusMode)
-> LoEAndGDDConfig (LoEAndGDDNodeKernelArgs m blk) -> ConsensusMode
forall a b. (a -> b) -> a -> b
$ GenesisNodeKernelArgs m blk
-> LoEAndGDDConfig (LoEAndGDDNodeKernelArgs m blk)
forall (m :: * -> *) blk.
GenesisNodeKernelArgs m blk
-> LoEAndGDDConfig (LoEAndGDDNodeKernelArgs m blk)
gnkaLoEAndGDDArgs GenesisNodeKernelArgs m blk
genesisArgs)
            BlockchainTime m
btime
            (ChainDB m blk -> STM m (AnchoredFragment (Header blk))
forall (m :: * -> *) blk.
ChainDB m blk -> STM m (AnchoredFragment (Header blk))
ChainDB.getCurrentChain ChainDB m blk
chainDB)
            STM m UseBootstrapPeers
getUseBootstrapPeers
            (GsmState -> LedgerStateJudgement
GSM.gsmStateToLedgerJudgement (GsmState -> LedgerStateJudgement)
-> STM m GsmState -> STM m LedgerStateJudgement
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> StrictTVar m GsmState -> STM m GsmState
forall (m :: * -> *) a. MonadSTM m => StrictTVar m a -> STM m a
readTVar StrictTVar m GsmState
varGsmState)
        blockFetchInterface ::
          BlockFetchConsensusInterface (ConnectionId addrNTN) (HeaderWithTime blk) blk m
        blockFetchInterface =
          Tracer m (TraceEventDbf (ConnectionId addrNTN))
-> BlockConfig blk
-> ChainDbView m blk
-> ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
-> (Header blk -> SizeInBytes)
-> STM m FetchMode
-> DiffusionPipeliningSupport
-> BlockFetchConsensusInterface
     (ConnectionId addrNTN) (HeaderWithTime blk) blk m
forall (m :: * -> *) peer blk.
(IOLike m, BlockSupportsDiffusionPipelining blk, Ord peer,
 LedgerSupportsProtocol blk, ConfigSupportsNode blk) =>
Tracer m (TraceEventDbf peer)
-> BlockConfig blk
-> ChainDbView m blk
-> ChainSyncClientHandleCollection peer m blk
-> (Header blk -> SizeInBytes)
-> STM m FetchMode
-> DiffusionPipeliningSupport
-> BlockFetchConsensusInterface peer (HeaderWithTime blk) blk m
BlockFetchClientInterface.mkBlockFetchConsensusInterface
            (Tracers m (ConnectionId addrNTN) addrNTC blk
-> Tracer m (TraceEventDbf (ConnectionId addrNTN))
forall remotePeer localPeer blk (f :: * -> *).
Tracers' remotePeer localPeer blk f -> f (TraceEventDbf remotePeer)
dbfTracer Tracers m (ConnectionId addrNTN) addrNTC blk
tracers)
            (TopLevelConfig blk -> BlockConfig blk
forall blk. TopLevelConfig blk -> BlockConfig blk
configBlock TopLevelConfig blk
cfg)
            (ChainDB m blk -> ChainDbView m blk
forall (m :: * -> *) blk. ChainDB m blk -> ChainDbView m blk
BlockFetchClientInterface.defaultChainDbView ChainDB m blk
chainDB)
            ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
varChainSyncHandles
            Header blk -> SizeInBytes
blockFetchSize
            STM m FetchMode
readFetchMode
            DiffusionPipeliningSupport
getDiffusionPipeliningSupport

    peerSharingRegistry <- newPeerSharingRegistry

    return IS{..}

toConsensusMode :: forall a. LoEAndGDDConfig a -> ConsensusMode
toConsensusMode :: forall a. LoEAndGDDConfig a -> ConsensusMode
toConsensusMode = \case
  LoEAndGDDConfig a
LoEAndGDDDisabled -> ConsensusMode
PraosMode
  LoEAndGDDEnabled a
_ -> ConsensusMode
GenesisMode

forkBlockForging ::
  forall m addrNTN addrNTC blk.
  (IOLike m, RunNode blk) =>
  InternalState m addrNTN addrNTC blk ->
  MkBlockForging m blk ->
  m (Thread m Void)
forkBlockForging :: forall (m :: * -> *) addrNTN addrNTC blk.
(IOLike m, RunNode blk) =>
InternalState m addrNTN addrNTC blk
-> MkBlockForging m blk -> m (Thread m Void)
forkBlockForging IS{StrictTVar m GsmState
TopLevelConfig blk
Mempool m blk
BlockchainTime m
ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
ChainDB m blk
ResourceRegistry m
BlockFetchConsensusInterface
  (ConnectionId addrNTN) (HeaderWithTime blk) blk m
KeepAliveRegistry (ConnectionId addrNTN) m
FetchClientRegistry
  (ConnectionId addrNTN) (HeaderWithTime blk) blk m
PeerSharingRegistry addrNTN m
Tracers m (ConnectionId addrNTN) addrNTC blk
blockFetchInterface :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> BlockFetchConsensusInterface
     (ConnectionId addrNTN) (HeaderWithTime blk) blk m
fetchClientRegistry :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> FetchClientRegistry
     (ConnectionId addrNTN) (HeaderWithTime blk) blk m
keepAliveRegistry :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> KeepAliveRegistry (ConnectionId addrNTN) m
mempool :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> Mempool m blk
peerSharingRegistry :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> PeerSharingRegistry addrNTN m
varChainSyncHandles :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
varGsmState :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> StrictTVar m GsmState
tracers :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk
-> Tracers m (ConnectionId addrNTN) addrNTC blk
cfg :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> TopLevelConfig blk
registry :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> ResourceRegistry m
btime :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> BlockchainTime m
chainDB :: forall (m :: * -> *) addrNTN addrNTC blk.
InternalState m addrNTN addrNTC blk -> ChainDB m blk
tracers :: Tracers m (ConnectionId addrNTN) addrNTC blk
cfg :: TopLevelConfig blk
registry :: ResourceRegistry m
btime :: BlockchainTime m
chainDB :: ChainDB m blk
blockFetchInterface :: BlockFetchConsensusInterface
  (ConnectionId addrNTN) (HeaderWithTime blk) blk m
fetchClientRegistry :: FetchClientRegistry
  (ConnectionId addrNTN) (HeaderWithTime blk) blk m
keepAliveRegistry :: KeepAliveRegistry (ConnectionId addrNTN) m
varChainSyncHandles :: ChainSyncClientHandleCollection (ConnectionId addrNTN) m blk
varGsmState :: StrictTVar m GsmState
mempool :: Mempool m blk
peerSharingRegistry :: PeerSharingRegistry addrNTN m
..} (MkBlockForging m (BlockForging m blk)
blockForgingM) =
  ResourceRegistry m
-> String
-> m (BlockForging m blk)
-> (BlockForging m blk -> m ())
-> (BlockForging m blk -> Watcher m SlotNo SlotNo)
-> m (Thread m Void)
forall (m :: * -> *) r a fp.
(IOLike m, Eq fp, HasCallStack) =>
ResourceRegistry m
-> String
-> m r
-> (r -> m ())
-> (r -> Watcher m a fp)
-> m (Thread m Void)
forkLinkedWatcherAllocate
    ResourceRegistry m
registry
    String
label
    m (BlockForging m blk)
blockForgingMLabel
    BlockForging m blk -> m ()
forall (m :: * -> *) blk. BlockForging m blk -> m ()
finalize
    ( \BlockForging m blk
bf -> BlockchainTime m -> (SlotNo -> m ()) -> Watcher m SlotNo SlotNo
forall (m :: * -> *).
IOLike m =>
BlockchainTime m -> (SlotNo -> m ()) -> Watcher m SlotNo SlotNo
knownSlotWatcher BlockchainTime m
btime ((SlotNo -> m ()) -> Watcher m SlotNo SlotNo)
-> (SlotNo -> m ()) -> Watcher m SlotNo SlotNo
forall a b. (a -> b) -> a -> b
$
        \SlotNo
currentSlot ->
          WithEarlyExit m () -> m ()
forall (m :: * -> *). Functor m => WithEarlyExit m () -> m ()
withEarlyExit_ (WithEarlyExit m () -> m ()) -> WithEarlyExit m () -> m ()
forall a b. (a -> b) -> a -> b
$
            Tracer m (TraceLabelCreds (TraceForgeEvent blk))
-> Tracer m (TraceLabelCreds (ForgeStateInfo blk))
-> TopLevelConfig blk
-> ChainDB m blk
-> Mempool m blk
-> BlockForging m blk
-> SlotNo
-> WithEarlyExit m ()
forall (m :: * -> *) blk.
(IOLike m, RunNode blk) =>
Tracer m (TraceLabelCreds (TraceForgeEvent blk))
-> Tracer m (TraceLabelCreds (ForgeStateInfo blk))
-> TopLevelConfig blk
-> ChainDB m blk
-> Mempool m blk
-> BlockForging m blk
-> SlotNo
-> WithEarlyExit m ()
forge
              (Tracers m (ConnectionId addrNTN) addrNTC blk
-> Tracer m (TraceLabelCreds (TraceForgeEvent blk))
forall remotePeer localPeer blk (f :: * -> *).
Tracers' remotePeer localPeer blk f
-> f (TraceLabelCreds (TraceForgeEvent blk))
forgeTracer Tracers m (ConnectionId addrNTN) addrNTC blk
tracers)
              (Tracers m (ConnectionId addrNTN) addrNTC blk
-> Tracer m (TraceLabelCreds (ForgeStateInfo blk))
forall remotePeer localPeer blk (f :: * -> *).
Tracers' remotePeer localPeer blk f
-> f (TraceLabelCreds (ForgeStateInfo blk))
forgeStateInfoTracer Tracers m (ConnectionId addrNTN) addrNTC blk
tracers)
              TopLevelConfig blk
cfg
              ChainDB m blk
chainDB
              Mempool m blk
mempool
              BlockForging m blk
bf
              SlotNo
currentSlot
    )
 where
  label :: String
  label :: String
label = String
"NodeKernel.blockForging"

  blockForgingMLabel :: m (BlockForging m blk)
blockForgingMLabel = do
    bf <- m (BlockForging m blk)
blockForgingM
    labelThisThread $ Text.unpack $ forgeLabel bf
    pure bf

{-------------------------------------------------------------------------------
  TxSubmission integration
-------------------------------------------------------------------------------}

getMempoolReader ::
  forall m blk.
  ( LedgerSupportsMempool blk
  , IOLike m
  , HasTxId (GenTx blk)
  ) =>
  Mempool m blk ->
  TxSubmissionMempoolReader (GenTxId blk) (Validated (GenTx blk)) TicketNo m
getMempoolReader :: forall (m :: * -> *) blk.
(LedgerSupportsMempool blk, IOLike m, HasTxId (GenTx blk)) =>
Mempool m blk
-> TxSubmissionMempoolReader
     (GenTxId blk) (Validated (GenTx blk)) TicketNo m
getMempoolReader Mempool m blk
mempool =
  MempoolReader.TxSubmissionMempoolReader
    { mempoolZeroIdx :: TicketNo
mempoolZeroIdx = TicketNo
zeroTicketNo
    , mempoolGetSnapshot :: STM
  m (MempoolSnapshot (GenTxId blk) (Validated (GenTx blk)) TicketNo)
mempoolGetSnapshot = MempoolSnapshot blk
-> MempoolSnapshot (GenTxId blk) (Validated (GenTx blk)) TicketNo
convertSnapshot (MempoolSnapshot blk
 -> MempoolSnapshot (GenTxId blk) (Validated (GenTx blk)) TicketNo)
-> STM m (MempoolSnapshot blk)
-> STM
     m (MempoolSnapshot (GenTxId blk) (Validated (GenTx blk)) TicketNo)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Mempool m blk -> STM m (MempoolSnapshot blk)
forall (m :: * -> *) blk.
Mempool m blk -> STM m (MempoolSnapshot blk)
getSnapshot Mempool m blk
mempool
    }
 where
  convertSnapshot ::
    MempoolSnapshot blk ->
    MempoolReader.MempoolSnapshot (GenTxId blk) (Validated (GenTx blk)) TicketNo
  convertSnapshot :: MempoolSnapshot blk
-> MempoolSnapshot (GenTxId blk) (Validated (GenTx blk)) TicketNo
convertSnapshot
    MempoolSnapshot
      { TicketNo -> [(Validated (GenTx blk), TicketNo, TxMeasure blk)]
snapshotTxsAfter :: TicketNo -> [(Validated (GenTx blk), TicketNo, TxMeasure blk)]
snapshotTxsAfter :: forall blk.
MempoolSnapshot blk
-> TicketNo -> [(Validated (GenTx blk), TicketNo, TxMeasure blk)]
snapshotTxsAfter
      , TicketNo -> Maybe (Validated (GenTx blk))
snapshotLookupTx :: TicketNo -> Maybe (Validated (GenTx blk))
snapshotLookupTx :: forall blk.
MempoolSnapshot blk -> TicketNo -> Maybe (Validated (GenTx blk))
snapshotLookupTx
      , GenTxId blk -> Bool
snapshotHasTx :: GenTxId blk -> Bool
snapshotHasTx :: forall blk. MempoolSnapshot blk -> GenTxId blk -> Bool
snapshotHasTx
      } =
      MempoolReader.MempoolSnapshot
        { mempoolTxIdsAfter :: TicketNo -> [(GenTxId blk, TicketNo, SizeInBytes)]
mempoolTxIdsAfter = \TicketNo
idx ->
            [ ( GenTx blk -> GenTxId blk
forall tx. HasTxId tx => tx -> TxId tx
txId (Validated (GenTx blk) -> GenTx blk
forall blk.
LedgerSupportsMempool blk =>
Validated (GenTx blk) -> GenTx blk
txForgetValidated Validated (GenTx blk)
tx)
              , TicketNo
idx'
              , GenTx blk -> SizeInBytes
forall blk. TxLimits blk => GenTx blk -> SizeInBytes
txWireSize (GenTx blk -> SizeInBytes) -> GenTx blk -> SizeInBytes
forall a b. (a -> b) -> a -> b
$ Validated (GenTx blk) -> GenTx blk
forall blk.
LedgerSupportsMempool blk =>
Validated (GenTx blk) -> GenTx blk
txForgetValidated Validated (GenTx blk)
tx
              )
            | (Validated (GenTx blk)
tx, TicketNo
idx', TxMeasure blk
_msr) <- TicketNo -> [(Validated (GenTx blk), TicketNo, TxMeasure blk)]
snapshotTxsAfter TicketNo
idx
            ]
        , mempoolLookupTx :: TicketNo -> Maybe (Validated (GenTx blk))
mempoolLookupTx = TicketNo -> Maybe (Validated (GenTx blk))
snapshotLookupTx
        , mempoolHasTx :: GenTxId blk -> Bool
mempoolHasTx = GenTxId blk -> Bool
snapshotHasTx
        }

getMempoolWriter ::
  forall blk m.
  ( LedgerSupportsMempool blk
  , IOLike m
  , HasTxId (GenTx blk)
  ) =>
  Mempool m blk ->
  TxSubmissionMempoolWriter (GenTxId blk) (GenTx blk) TicketNo m ()
getMempoolWriter :: forall blk (m :: * -> *).
(LedgerSupportsMempool blk, IOLike m, HasTxId (GenTx blk)) =>
Mempool m blk
-> TxSubmissionMempoolWriter
     (GenTxId blk) (GenTx blk) TicketNo m ()
getMempoolWriter Mempool m blk
mempool =
  Inbound.TxSubmissionMempoolWriter
    { txId :: GenTx blk -> TxId (GenTx blk)
Inbound.txId = GenTx blk -> TxId (GenTx blk)
forall tx. HasTxId tx => tx -> TxId tx
txId
    , mempoolAddTxs :: [GenTx blk] -> m ([TxId (GenTx blk)], [(TxId (GenTx blk), ())])
mempoolAddTxs = \[GenTx blk]
txs ->
        [MempoolAddTxResult blk]
-> ([TxId (GenTx blk)], [(TxId (GenTx blk), ())])
partitionTxIds ([MempoolAddTxResult blk]
 -> ([TxId (GenTx blk)], [(TxId (GenTx blk), ())]))
-> m [MempoolAddTxResult blk]
-> m ([TxId (GenTx blk)], [(TxId (GenTx blk), ())])
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Mempool m blk -> [GenTx blk] -> m [MempoolAddTxResult blk]
forall (m :: * -> *) blk (t :: * -> *).
(MonadSTM m, Traversable t) =>
Mempool m blk -> t (GenTx blk) -> m (t (MempoolAddTxResult blk))
addTxs Mempool m blk
mempool [GenTx blk]
txs
    }
 where
  partitionTxIds ::
    [MempoolAddTxResult blk] -> ([TxId (GenTx blk)], [(TxId (GenTx blk), ())])
  partitionTxIds :: [MempoolAddTxResult blk]
-> ([TxId (GenTx blk)], [(TxId (GenTx blk), ())])
partitionTxIds = [Either (TxId (GenTx blk)) (TxId (GenTx blk), ())]
-> ([TxId (GenTx blk)], [(TxId (GenTx blk), ())])
forall a b. [Either a b] -> ([a], [b])
partitionEithers ([Either (TxId (GenTx blk)) (TxId (GenTx blk), ())]
 -> ([TxId (GenTx blk)], [(TxId (GenTx blk), ())]))
-> ([MempoolAddTxResult blk]
    -> [Either (TxId (GenTx blk)) (TxId (GenTx blk), ())])
-> [MempoolAddTxResult blk]
-> ([TxId (GenTx blk)], [(TxId (GenTx blk), ())])
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (MempoolAddTxResult blk
 -> Either (TxId (GenTx blk)) (TxId (GenTx blk), ()))
-> [MempoolAddTxResult blk]
-> [Either (TxId (GenTx blk)) (TxId (GenTx blk), ())]
forall a b. (a -> b) -> [a] -> [b]
map MempoolAddTxResult blk
-> Either (TxId (GenTx blk)) (TxId (GenTx blk), ())
getTxId

  getTxId :: MempoolAddTxResult blk -> Either (TxId (GenTx blk)) (TxId (GenTx blk), ())
  getTxId :: MempoolAddTxResult blk
-> Either (TxId (GenTx blk)) (TxId (GenTx blk), ())
getTxId = \case
    MempoolTxAdded Validated (GenTx blk)
tx LedgerTables blk DiffMK
_ -> TxId (GenTx blk)
-> Either (TxId (GenTx blk)) (TxId (GenTx blk), ())
forall a b. a -> Either a b
Left (TxId (GenTx blk)
 -> Either (TxId (GenTx blk)) (TxId (GenTx blk), ()))
-> TxId (GenTx blk)
-> Either (TxId (GenTx blk)) (TxId (GenTx blk), ())
forall a b. (a -> b) -> a -> b
$ GenTx blk -> TxId (GenTx blk)
forall tx. HasTxId tx => tx -> TxId tx
txId (Validated (GenTx blk) -> GenTx blk
forall blk.
LedgerSupportsMempool blk =>
Validated (GenTx blk) -> GenTx blk
txForgetValidated Validated (GenTx blk)
tx)
    MempoolTxRejected GenTx blk
tx ApplyTxErr blk
_reason -> (TxId (GenTx blk), ())
-> Either (TxId (GenTx blk)) (TxId (GenTx blk), ())
forall a b. b -> Either a b
Right ((TxId (GenTx blk), ())
 -> Either (TxId (GenTx blk)) (TxId (GenTx blk), ()))
-> (TxId (GenTx blk), ())
-> Either (TxId (GenTx blk)) (TxId (GenTx blk), ())
forall a b. (a -> b) -> a -> b
$ (GenTx blk -> TxId (GenTx blk)
forall tx. HasTxId tx => tx -> TxId tx
txId GenTx blk
tx, ())

{-------------------------------------------------------------------------------
  PeerSelection integration
-------------------------------------------------------------------------------}

-- | Retrieve the peers registered in the current chain/ledger state by
-- descending stake.
--
-- For example, for Shelley, this will return the stake pool relays ordered by
-- descending stake.
--
-- Only returns a 'Just' when the given predicate returns 'True'. This predicate
-- can for example check whether the slot of the ledger state is older or newer
-- than some slot number.
--
-- We don't use the ledger state at the tip of the chain, but the ledger state
-- @k@ blocks back, i.e., at the tip of the immutable chain, because any stake
-- pools registered in that ledger state are guaranteed to be stable. This
-- justifies merging the future and current stake pools.
getPeersFromCurrentLedger ::
  (IOLike m, LedgerSupportsPeerSelection blk) =>
  NodeKernel m addrNTN addrNTC blk ->
  (LedgerState blk EmptyMK -> Bool) ->
  STM m (Maybe [(PoolStake, NonEmpty LedgerRelayAccessPoint)])
getPeersFromCurrentLedger :: forall (m :: * -> *) blk addrNTN addrNTC.
(IOLike m, LedgerSupportsPeerSelection blk) =>
NodeKernel m addrNTN addrNTC blk
-> (LedgerState blk EmptyMK -> Bool)
-> STM m (Maybe [(PoolStake, NonEmpty LedgerRelayAccessPoint)])
getPeersFromCurrentLedger NodeKernel m addrNTN addrNTC blk
kernel LedgerState blk EmptyMK -> Bool
p = do
  immutableLedger <-
    ExtLedgerState blk EmptyMK -> LedgerState blk EmptyMK
forall blk (mk :: MapKind).
ExtLedgerState blk mk -> LedgerState blk mk
ledgerState (ExtLedgerState blk EmptyMK -> LedgerState blk EmptyMK)
-> STM m (ExtLedgerState blk EmptyMK)
-> STM m (LedgerState blk EmptyMK)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> ChainDB m blk -> STM m (ExtLedgerState blk EmptyMK)
forall (m :: * -> *) blk.
ChainDB m blk -> STM m (ExtLedgerState blk EmptyMK)
ChainDB.getImmutableLedger (NodeKernel m addrNTN addrNTC blk -> ChainDB m blk
forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> ChainDB m blk
getChainDB NodeKernel m addrNTN addrNTC blk
kernel)
  return $ do
    guard (p immutableLedger)
    return $
      map (second (fmap stakePoolRelayAccessPoint)) $
        force $
          getPeers immutableLedger

-- | Like 'getPeersFromCurrentLedger' but with a \"after slot number X\"
-- condition.
getPeersFromCurrentLedgerAfterSlot ::
  forall m blk addrNTN addrNTC.
  ( IOLike m
  , LedgerSupportsPeerSelection blk
  , UpdateLedger blk
  ) =>
  NodeKernel m addrNTN addrNTC blk ->
  SlotNo ->
  STM m (Maybe [(PoolStake, NonEmpty LedgerRelayAccessPoint)])
getPeersFromCurrentLedgerAfterSlot :: forall (m :: * -> *) blk addrNTN addrNTC.
(IOLike m, LedgerSupportsPeerSelection blk, UpdateLedger blk) =>
NodeKernel m addrNTN addrNTC blk
-> SlotNo
-> STM m (Maybe [(PoolStake, NonEmpty LedgerRelayAccessPoint)])
getPeersFromCurrentLedgerAfterSlot NodeKernel m addrNTN addrNTC blk
kernel SlotNo
slotNo =
  NodeKernel m addrNTN addrNTC blk
-> (LedgerState blk EmptyMK -> Bool)
-> STM m (Maybe [(PoolStake, NonEmpty LedgerRelayAccessPoint)])
forall (m :: * -> *) blk addrNTN addrNTC.
(IOLike m, LedgerSupportsPeerSelection blk) =>
NodeKernel m addrNTN addrNTC blk
-> (LedgerState blk EmptyMK -> Bool)
-> STM m (Maybe [(PoolStake, NonEmpty LedgerRelayAccessPoint)])
getPeersFromCurrentLedger NodeKernel m addrNTN addrNTC blk
kernel LedgerState blk EmptyMK -> Bool
forall (mk :: MapKind). LedgerState blk mk -> Bool
afterSlotNo
 where
  afterSlotNo :: LedgerState blk mk -> Bool
  afterSlotNo :: forall (mk :: MapKind). LedgerState blk mk -> Bool
afterSlotNo LedgerState blk mk
st =
    case LedgerState blk mk -> WithOrigin SlotNo
forall blk (mk :: MapKind).
UpdateLedger blk =>
LedgerState blk mk -> WithOrigin SlotNo
ledgerTipSlot LedgerState blk mk
st of
      WithOrigin SlotNo
Origin -> Bool
False
      NotOrigin SlotNo
tip -> SlotNo
tip SlotNo -> SlotNo -> Bool
forall a. Ord a => a -> a -> Bool
> SlotNo
slotNo

-- | Retrieve the slot of the immutable tip
getImmTipSlot ::
  ( IOLike m
  , UpdateLedger blk
  ) =>
  NodeKernel m addrNTN addrNTC blk ->
  STM m (WithOrigin SlotNo)
getImmTipSlot :: forall (m :: * -> *) blk addrNTN addrNTC.
(IOLike m, UpdateLedger blk) =>
NodeKernel m addrNTN addrNTC blk -> STM m (WithOrigin SlotNo)
getImmTipSlot NodeKernel m addrNTN addrNTC blk
kernel =
  ExtLedgerState blk EmptyMK -> WithOrigin SlotNo
forall (l :: LedgerStateKind) (mk :: MapKind).
GetTip l =>
l mk -> WithOrigin SlotNo
getTipSlot
    (ExtLedgerState blk EmptyMK -> WithOrigin SlotNo)
-> STM m (ExtLedgerState blk EmptyMK) -> STM m (WithOrigin SlotNo)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> ChainDB m blk -> STM m (ExtLedgerState blk EmptyMK)
forall (m :: * -> *) blk.
ChainDB m blk -> STM m (ExtLedgerState blk EmptyMK)
ChainDB.getImmutableLedger (NodeKernel m addrNTN addrNTC blk -> ChainDB m blk
forall (m :: * -> *) addrNTN addrNTC blk.
NodeKernel m addrNTN addrNTC blk -> ChainDB m blk
getChainDB NodeKernel m addrNTN addrNTC blk
kernel)