haskoin-node 0.19.0 → 1.0.0
raw patch · 8 files changed
+2624/−2180 lines, 8 filesdep ~haskoin-corePVP ok
version bump matches the API change (PVP)
Dependency ranges changed: haskoin-core
API changes (from Hackage documentation)
- Haskoin.Node: PeerManagerConfig :: !Int -> ![String] -> !Bool -> !NetworkAddress -> !Network -> !Publisher PeerEvent -> !NominalDiffTime -> !NominalDiffTime -> !SockAddr -> WithConnection -> !Publisher (Peer, Message) -> PeerManagerConfig
- Haskoin.Node: [chainConfColumnFamily] :: ChainConfig -> !Maybe ColumnFamily
- Haskoin.Node: [chainConfDB] :: ChainConfig -> !DB
- Haskoin.Node: [chainConfEvents] :: ChainConfig -> !Publisher ChainEvent
- Haskoin.Node: [chainConfNetwork] :: ChainConfig -> !Network
- Haskoin.Node: [chainConfTimeout] :: ChainConfig -> !NominalDiffTime
- Haskoin.Node: [inboundConduit] :: Conduits -> ConduitT () ByteString IO ()
- Haskoin.Node: [nodeChain] :: Node -> !Chain
- Haskoin.Node: [nodeConfColumnFamily] :: NodeConfig -> !Maybe ColumnFamily
- Haskoin.Node: [nodeConfConnect] :: NodeConfig -> !SockAddr -> WithConnection
- Haskoin.Node: [nodeConfDB] :: NodeConfig -> !DB
- Haskoin.Node: [nodeConfDiscover] :: NodeConfig -> !Bool
- Haskoin.Node: [nodeConfEvents] :: NodeConfig -> !Publisher NodeEvent
- Haskoin.Node: [nodeConfMaxPeers] :: NodeConfig -> !Int
- Haskoin.Node: [nodeConfNetAddr] :: NodeConfig -> !NetworkAddress
- Haskoin.Node: [nodeConfNet] :: NodeConfig -> !Network
- Haskoin.Node: [nodeConfPeerMaxLife] :: NodeConfig -> !NominalDiffTime
- Haskoin.Node: [nodeConfPeers] :: NodeConfig -> ![String]
- Haskoin.Node: [nodeConfTimeout] :: NodeConfig -> !NominalDiffTime
- Haskoin.Node: [nodeManager] :: Node -> !PeerManager
- Haskoin.Node: [onlinePeerAddress] :: OnlinePeer -> !SockAddr
- Haskoin.Node: [onlinePeerAsync] :: OnlinePeer -> !Async ()
- Haskoin.Node: [onlinePeerConnectTime] :: OnlinePeer -> !UTCTime
- Haskoin.Node: [onlinePeerConnected] :: OnlinePeer -> !Bool
- Haskoin.Node: [onlinePeerDisconnect] :: OnlinePeer -> !UTCTime
- Haskoin.Node: [onlinePeerMailbox] :: OnlinePeer -> !Peer
- Haskoin.Node: [onlinePeerNonce] :: OnlinePeer -> !Word64
- Haskoin.Node: [onlinePeerPing] :: OnlinePeer -> !Maybe (UTCTime, Word64)
- Haskoin.Node: [onlinePeerPings] :: OnlinePeer -> ![NominalDiffTime]
- Haskoin.Node: [onlinePeerTickled] :: OnlinePeer -> !UTCTime
- Haskoin.Node: [onlinePeerVerAck] :: OnlinePeer -> !Bool
- Haskoin.Node: [onlinePeerVersion] :: OnlinePeer -> !Maybe Version
- Haskoin.Node: [outboundConduit] :: Conduits -> ConduitT ByteString Void IO ()
- Haskoin.Node: [peerConfConnect] :: PeerConfig -> !WithConnection
- Haskoin.Node: [peerConfNetwork] :: PeerConfig -> !Network
- Haskoin.Node: [peerConfPub] :: PeerConfig -> !Publisher (Peer, Message)
- Haskoin.Node: [peerConfText] :: PeerConfig -> !Text
- Haskoin.Node: [peerManagerConnect] :: PeerManagerConfig -> !SockAddr -> WithConnection
- Haskoin.Node: [peerManagerDiscover] :: PeerManagerConfig -> !Bool
- Haskoin.Node: [peerManagerEvents] :: PeerManagerConfig -> !Publisher PeerEvent
- Haskoin.Node: [peerManagerMaxLife] :: PeerManagerConfig -> !NominalDiffTime
- Haskoin.Node: [peerManagerMaxPeers] :: PeerManagerConfig -> !Int
- Haskoin.Node: [peerManagerNetAddr] :: PeerManagerConfig -> !NetworkAddress
- Haskoin.Node: [peerManagerNetwork] :: PeerManagerConfig -> !Network
- Haskoin.Node: [peerManagerPeers] :: PeerManagerConfig -> ![String]
- Haskoin.Node: [peerManagerPub] :: PeerManagerConfig -> !Publisher (Peer, Message)
- Haskoin.Node: [peerManagerTimeout] :: PeerManagerConfig -> !NominalDiffTime
- Haskoin.Node: data PeerManager
- Haskoin.Node: data PeerManagerConfig
- Haskoin.Node: managerAddrs :: MonadIO m => Peer -> [NetworkAddress] -> PeerManager -> m ()
- Haskoin.Node: managerBest :: MonadIO m => BlockHeight -> PeerManager -> m ()
- Haskoin.Node: managerPing :: MonadIO m => Peer -> Word64 -> PeerManager -> m ()
- Haskoin.Node: managerPong :: MonadIO m => Peer -> Word64 -> PeerManager -> m ()
- Haskoin.Node: managerTickle :: MonadIO m => Peer -> PeerManager -> m ()
- Haskoin.Node: managerVerAck :: MonadIO m => Peer -> PeerManager -> m ()
- Haskoin.Node: managerVersion :: MonadIO m => Peer -> Version -> PeerManager -> m ()
- Haskoin.Node: peerPublisher :: Peer -> Publisher (Peer, Message)
- Haskoin.Node: peerText :: Peer -> Text
- Haskoin.Node: withPeerManager :: (MonadUnliftIO m, MonadLoggerIO m) => PeerManagerConfig -> (PeerManager -> m a) -> m a
+ Haskoin.Node: Peer :: !Mailbox PeerMessage -> !Publisher PeerEvent -> !Text -> !TVar Bool -> Peer
+ Haskoin.Node: PeerMgrConfig :: !Int -> ![String] -> !Bool -> !NetworkAddress -> !Network -> !Publisher PeerEvent -> !NominalDiffTime -> !NominalDiffTime -> !SockAddr -> WithConnection -> PeerMgrConfig
+ Haskoin.Node: [$sel:address:NodeConfig] :: NodeConfig -> !NetworkAddress
+ Haskoin.Node: [$sel:address:OnlinePeer] :: OnlinePeer -> !SockAddr
+ Haskoin.Node: [$sel:address:PeerMgrConfig] :: PeerMgrConfig -> !NetworkAddress
+ Haskoin.Node: [$sel:async:OnlinePeer] :: OnlinePeer -> !Async ()
+ Haskoin.Node: [$sel:busy:Peer] :: Peer -> !TVar Bool
+ Haskoin.Node: [$sel:cf:ChainConfig] :: ChainConfig -> !Maybe ColumnFamily
+ Haskoin.Node: [$sel:cf:NodeConfig] :: NodeConfig -> !Maybe ColumnFamily
+ Haskoin.Node: [$sel:chain:Node] :: Node -> !Chain
+ Haskoin.Node: [$sel:connect:NodeConfig] :: NodeConfig -> !SockAddr -> WithConnection
+ Haskoin.Node: [$sel:connect:PeerConfig] :: PeerConfig -> !WithConnection
+ Haskoin.Node: [$sel:connect:PeerMgrConfig] :: PeerMgrConfig -> !SockAddr -> WithConnection
+ Haskoin.Node: [$sel:connected:OnlinePeer] :: OnlinePeer -> !UTCTime
+ Haskoin.Node: [$sel:db:ChainConfig] :: ChainConfig -> !DB
+ Haskoin.Node: [$sel:db:NodeConfig] :: NodeConfig -> !DB
+ Haskoin.Node: [$sel:discover:NodeConfig] :: NodeConfig -> !Bool
+ Haskoin.Node: [$sel:discover:PeerMgrConfig] :: PeerMgrConfig -> !Bool
+ Haskoin.Node: [$sel:inboundConduit:Conduits] :: Conduits -> ConduitT () ByteString IO ()
+ Haskoin.Node: [$sel:label:PeerConfig] :: PeerConfig -> !Text
+ Haskoin.Node: [$sel:label:Peer] :: Peer -> !Text
+ Haskoin.Node: [$sel:mailbox:OnlinePeer] :: OnlinePeer -> !Peer
+ Haskoin.Node: [$sel:mailbox:Peer] :: Peer -> !Mailbox PeerMessage
+ Haskoin.Node: [$sel:maxPeerLife:NodeConfig] :: NodeConfig -> !NominalDiffTime
+ Haskoin.Node: [$sel:maxPeerLife:PeerMgrConfig] :: PeerMgrConfig -> !NominalDiffTime
+ Haskoin.Node: [$sel:maxPeers:NodeConfig] :: NodeConfig -> !Int
+ Haskoin.Node: [$sel:maxPeers:PeerMgrConfig] :: PeerMgrConfig -> !Int
+ Haskoin.Node: [$sel:net:ChainConfig] :: ChainConfig -> !Network
+ Haskoin.Node: [$sel:net:NodeConfig] :: NodeConfig -> !Network
+ Haskoin.Node: [$sel:net:PeerConfig] :: PeerConfig -> !Network
+ Haskoin.Node: [$sel:net:PeerMgrConfig] :: PeerMgrConfig -> !Network
+ Haskoin.Node: [$sel:nonce:OnlinePeer] :: OnlinePeer -> !Word64
+ Haskoin.Node: [$sel:online:OnlinePeer] :: OnlinePeer -> !Bool
+ Haskoin.Node: [$sel:outboundConduit:Conduits] :: Conduits -> ConduitT ByteString Void IO ()
+ Haskoin.Node: [$sel:peerMgr:Node] :: Node -> !PeerMgr
+ Haskoin.Node: [$sel:peers:NodeConfig] :: NodeConfig -> ![String]
+ Haskoin.Node: [$sel:peers:PeerMgrConfig] :: PeerMgrConfig -> ![String]
+ Haskoin.Node: [$sel:ping:OnlinePeer] :: OnlinePeer -> !Maybe (UTCTime, Word64)
+ Haskoin.Node: [$sel:pings:OnlinePeer] :: OnlinePeer -> ![NominalDiffTime]
+ Haskoin.Node: [$sel:pub:ChainConfig] :: ChainConfig -> !Publisher ChainEvent
+ Haskoin.Node: [$sel:pub:NodeConfig] :: NodeConfig -> !Publisher NodeEvent
+ Haskoin.Node: [$sel:pub:PeerConfig] :: PeerConfig -> !Publisher PeerEvent
+ Haskoin.Node: [$sel:pub:PeerMgrConfig] :: PeerMgrConfig -> !Publisher PeerEvent
+ Haskoin.Node: [$sel:pub:Peer] :: Peer -> !Publisher PeerEvent
+ Haskoin.Node: [$sel:tickled:OnlinePeer] :: OnlinePeer -> !UTCTime
+ Haskoin.Node: [$sel:timeout:ChainConfig] :: ChainConfig -> !NominalDiffTime
+ Haskoin.Node: [$sel:timeout:NodeConfig] :: NodeConfig -> !NominalDiffTime
+ Haskoin.Node: [$sel:timeout:PeerMgrConfig] :: PeerMgrConfig -> !NominalDiffTime
+ Haskoin.Node: [$sel:verack:OnlinePeer] :: OnlinePeer -> !Bool
+ Haskoin.Node: [$sel:version:OnlinePeer] :: OnlinePeer -> !Maybe Version
+ Haskoin.Node: data PeerMgr
+ Haskoin.Node: data PeerMgrConfig
+ Haskoin.Node: peerMgrAddrs :: MonadIO m => Peer -> [NetworkAddress] -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrBest :: MonadIO m => BlockHeight -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrPing :: MonadIO m => Peer -> Word64 -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrPong :: MonadIO m => Peer -> Word64 -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrTickle :: MonadIO m => Peer -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrVerAck :: MonadIO m => Peer -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrVersion :: MonadIO m => Peer -> Version -> PeerMgr -> m ()
+ Haskoin.Node: withPeerMgr :: (MonadUnliftIO m, MonadLoggerIO m) => PeerMgrConfig -> (PeerMgr -> m a) -> m a
- Haskoin.Node: Node :: !PeerManager -> !Chain -> Node
+ Haskoin.Node: Node :: !PeerMgr -> !Chain -> Node
- Haskoin.Node: OnlinePeer :: !SockAddr -> !Bool -> !Bool -> !Maybe Version -> !Async () -> !Peer -> !Word64 -> !Maybe (UTCTime, Word64) -> ![NominalDiffTime] -> !UTCTime -> !UTCTime -> !UTCTime -> OnlinePeer
+ Haskoin.Node: OnlinePeer :: !SockAddr -> !Bool -> !Bool -> !Maybe Version -> !Async () -> !Peer -> !Word64 -> !Maybe (UTCTime, Word64) -> ![NominalDiffTime] -> !UTCTime -> !UTCTime -> OnlinePeer
- Haskoin.Node: PeerConfig :: !Publisher (Peer, Message) -> !Network -> !Text -> !WithConnection -> PeerConfig
+ Haskoin.Node: PeerConfig :: !Publisher PeerEvent -> !Network -> !Text -> !WithConnection -> PeerConfig
- Haskoin.Node: PeerMessage :: !Peer -> !Message -> NodeEvent
+ Haskoin.Node: PeerMessage :: !Peer -> !Message -> PeerEvent
- Haskoin.Node: getOnlinePeer :: MonadIO m => Peer -> PeerManager -> m (Maybe OnlinePeer)
+ Haskoin.Node: getOnlinePeer :: MonadIO m => Peer -> PeerMgr -> m (Maybe OnlinePeer)
- Haskoin.Node: getPeers :: MonadIO m => PeerManager -> m [OnlinePeer]
+ Haskoin.Node: getPeers :: MonadIO m => PeerMgr -> m [OnlinePeer]
Files
- CHANGELOG.md +186/−56
- haskoin-node.cabal +6/−6
- src/Haskoin/Node.hs +167/−165
- src/Haskoin/Node/Chain.hs +772/−683
- src/Haskoin/Node/Manager.hs +0/−767
- src/Haskoin/Node/Peer.hs +326/−249
- src/Haskoin/Node/PeerMgr.hs +867/−0
- test/Haskoin/NodeSpec.hs +300/−254
CHANGELOG.md view
@@ -4,243 +4,358 @@ The format is based on [Keep a Changelog](http://keepachangelog.com/en/1.0.0/) and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0.html). -## 0.18.1+## [1.0.0] - 2023-07-28++### Changed++- Make compatible with latest haskoin-core.+- Use DuplicateRecordFields and OverloadedRecordDot.+- Simplify pub/sub queues.+- Multiple refactoring passes.++## [0.18.1] - 2022-07-27+ ### Fixed+ - Set default port for peers where it is unset. -## 0.18.0+## [0.18.0] - 2022-07-27+ ### Added+ - Support setting up and connecting to IPv6 peers. -## 0.17.14+## [0.17.14] - 2021-08-14+ ### Fixed+ - Reduce verbosity on incoming header decode test. - Show appropriate error message upon receiving empty headers. -## 0.17.13+## [0.17.13] - 2021-08-14+ ### Added+ - Display more details about invalid incoming headers. -## 0.17.12+## [0.17.12] - 2021-05-17+ ### Fixed+ - Do not connect to more than the maximum number of peers. -## 0.17.11+## [0.17.11] - 2021-05-17+ ### Added+ - Display message command that disconnects a peer. -## 0.17.10+## [0.17.10] - 2021-05-17+ ### Fixed+ - Correct disconnect timeout algorithm bug. -## 0.17.9+## [0.17.9] - 2021-05-17+ ### Fixed+ - Add randomised timeouts to avoid disconnecting all peers. -## 0.17.2+## [0.17.2] - 2021-03-09+ ### Fixed+ - Do not start chain actor until database initialized. -## 0.17.1+## [0.17.1] - 2021-01-08+ ### Changed+ - Depend on haskoin-core-0.17.3. -## 0.17.0+## [0.17.0] - 2020-10-21+ ### Added+ - Support for Bitcoin Cash November 2020 hard fork. - BlockHeaders instance for ReaderT Chain. -## 0.16.0+## [0.16.0] - 2020-07-23+ ### Changed+ - Add support for column families. -## 0.15.0+## [0.15.0] - 2020-07-20+ ### Changed+ - Use new Haskell bindings for RocksDB. -## 0.14.1+## [0.14.1] - 2020-06-19+ ### Fixed+ - Correct flawed peer locking logic in Chain actor. -## 0.14.0+## [0.14.0] - 2020-06-18+ ### Changed+ - Massively refactor everything in a non-backwards-compatible way. - Use MIT license. - Bump haskoin-core. - Bump secp256k1-haskell. ### Fixed+ - Fix getting stuck on a single peer. -## 0.13.0+## [0.13.0] - 2020-05-08+ ### Changed+ - Depend on Haskoin Store 0.13.3. - Better code organisation. -## 0.12.0+## [0.12.0] - 2020-05-06+ ### Changed+ - Add a test suite that simulates network instead of connecting to real one. -## 0.11.3+## [0.11.3] - 2020-05-03+ ### Changed+ - Revert including multiline decoding error text in logs. -## 0.11.2+## [0.11.2] - 2020-05-03+ ### Changed+ - Include header decoding error text in logs. -## 0.11.1+## [0.11.1] - 2020-05-03+ ### Changed+ - Improve logging. -## 0.11.0+## [0.11.0] - 2020-05-03+ ### Changed+ - Set peer too old time -## 0.10.1+## [0.10.1] - 2020-05-03+ ### Changed+ - Disconnect old peers after 48 hours instead of 30 minutes. -## 0.10.0+## [0.10.0] - 2020-05-03+ ### Changed+ - Move modules out of Network.Haskoin namespace. - Add better and more logging. - Change Manager module name and related values to PeerManager. -## 0.9.21+## [0.9.21] - 2020-04-07+ ### Removed+ - Remove unnecessary logging. -## 0.9.20+## [0.9.20] - 2020-04-07+ ### Changed+ - Better log messages. - Less verbose debug logging. -## 0.9.19+## [0.9.19] - 2020-04-07+ ### Changed - Better log messages. -## 0.9.18+## [0.9.18] - 2020-04-07+ ### Added+ - More aggressive peer discovery. -## 0.9.17+## [0.9.17] - 2020-04-07+ ### Added+ - Peers are disconnected automatically after awhile. -## 0.9.16+## [0.9.16] - 2020-02-08+ ### Added+ - Lower bound versions for some dependencies. -## 0.9.15+## [0.9.15] - 2020-01-15+ ### Changed+ - Update to support new `NetworkAddress` data structure from `haskoin-core`. -## 0.9.14+## [0.9.14] - 2019-12-10+ ### Removed+ - No longer support storing peers in db as performance tradeoff doesn't justify it. -## 0.9.13+## [0.9.13] - 2019-10-08+ ### Changed+ - Really store peers in db. -## 0.9.12+## [0.9.12] - 2019-10-08+ ### Changed+ - Demote some logging to debug level. -## 0.9.11+## [0.9.11] - 2019-10-02+ ### Added+ - Add `-O2` optimisations to GHC. -## 0.9.10+## [0.9.10] - 2019-04-19+ ### Added+ - Increase debugging information where application freezes. -## 0.9.9+## [0.9.9] - 2019-04-12+ ### Added+ - Increase debugging information. -## 0.9.8+## [0.9.8] - 2019-04-12+ ### Changed+ - Increase version of haskoin-core to 0.9.0. - Fix some tests. -## 0.9.7+## [0.9.7] - 2019-04-12+ ### Added+ - More debugging. ### Changed+ - Be defensive against duplicate peers. - Increase interval between housekeeping pings. - Replace peers in database atomically. -## 0.9.6+## [0.9.6] - 2019-04-01+ ### Changed+ - Randomize known peers instead of keeping scores. - Simplify peer management code to avoid freezes. - Merge logic for chain and manager. -## 0.9.5+## [0.9.5] - 2018-11-14+ ### Changed+ - Do not record new peers in database when peer discovery is disabled. -## 0.9.4+## [0.9.4] - 2018-11-01+ ### Changed+ - Don't spam best block events. -## 0.9.3+## [0.9.3] - 2018-10-22+ ### Changed+ - Correct display of milliseconds in log. - Correct bug when receiving headers from unknown peer. - Simplify chain syncing code. -## 0.9.2+## [0.9.2] - 2018-10-18+ ### Changed+ - Peer dies immediately when receiving a bad message. -## 0.9.1+## [0.9.1] - 2018-10-18+ ### Changed+ - Keep track of last synced header from a peer to avoid endless loops on large reorgs. -## 0.9.0+## [0.9.0] - 2018-10-17+ ### Changed+ - Use an STM listener instead of a publisher for the node API. -## 0.8.2+## [0.8.2] - 2018-10-17+ ### Added+ - Expose `ChainMessage` and `ManagerMessage` types from `Haskoin.Node` module. -## 0.8.1+## [0.8.1] - 2018-10-11+ ### Changed+ - Corrected documentation for `killPeer` function. - Leave time out of logic code. -## 0.8.0+## [0.8.0] - 2018-10-09+ ### Changed+ - Peers are now killed directly instead of through peer manager. ### Removed+ - Chain no longer needs peer manager. -## 0.7.2+## [0.7.2] - 2018-10-09+ ### Added+ - Compatibility with base 4.12. ### Changed+ - Update base to 4.9. -## 0.7.1+## [0.7.1] - 2018-10-09+ ### Added+ - Allow to easily obtain a peer's publisher. -## 0.7.0+## [0.7.0] - 2018-10-09+ ### Added+ - Versioning for chain and peer database. - Automatic purging of chain and peer database when version changes. - Add extra timers. - Add publishers to every peer. ### Changed+ - Full reimplementation of node API. - Simplify peer selection and management. - Merge manager and peer events.@@ -248,6 +363,7 @@ - Separate logic from actors for peer manager and chain. ### Removed+ - Remove irrelevant fields from peer information. - Remove unreliable peer block head tracking. - Remove dependency on deprecated binary conduits.@@ -255,33 +371,45 @@ - Remove unreliable peer request tracking code. - Remove separate manager events. -## 0.6.1+## [0.6.1] - 2018-09-14+ ### Changed+ - Fix bug where peer height did not update in certain cases. -## 0.6.0+## [0.6.0] - 2018-09-14+ ### Added+ - Documentation everywhere. ### Changed+ - Make compatible with NQE 0.5. - Use supervisor only in peer manager. - API quality of life changes. - Exposed module is now only `Haskoin.Node`. ### Removed+ - No more direct access to internals. -## 0.5.2+## [0.5.2] - 2018-09-10+ ### Changed+ - Improve dependency definitions. -## 0.5.1+## [0.5.1] - 2018-09-10+ ### Changed+ - Dependency `sec256k1` changes to `secp256k1-haskell`. -## 0.5.0+## [0.5.0] - 2018-09-09+ ### Added+ - New `CHANGELOG.md` file. - Use `nqe` for concurrency. - Peer discovery.@@ -289,11 +417,13 @@ - Support for Merkle blocks. ### Changed+ - Split out of former `haskoin` repository. - Use hpack and `package.yaml`. - Old `haskoin-node` package now renamed to `old-haskoin-node` and deprecated. ### Removed+ - Removed Old Haskoin Node package completely. - Removed Stylish Haskell configuration file. - Remvoed `haskoin-core` and `haskoin-wallet` packages from this repository.
haskoin-node.cabal view
@@ -1,13 +1,13 @@ cabal-version: 1.12 --- This file has been generated from package.yaml by hpack version 0.35.1.+-- This file has been generated from package.yaml by hpack version 0.35.2. -- -- see: https://github.com/sol/hpack ----- hash: a756b44df814f3c6d44c22345389b0b87b74ae7d921e18f027d27d41c1ad38ed+-- hash: 2f4be0a55ae78de70b6694a331bfa9338f4a9ac58590c07e9f7a044fe58fd9b9 name: haskoin-node-version: 0.19.0+version: 1.0.0 synopsis: P2P library for Bitcoin and Bitcoin Cash description: Please see the README on GitHub at <https://github.com/haskoin/haskoin-node#readme> category: Bitcoin, Finance, Network@@ -31,8 +31,8 @@ Haskoin.Node other-modules: Haskoin.Node.Chain- Haskoin.Node.Manager Haskoin.Node.Peer+ Haskoin.Node.PeerMgr Paths_haskoin_node hs-source-dirs: src@@ -45,7 +45,7 @@ , containers , data-default , hashable- , haskoin-core >=0.22.0+ , haskoin-core >=1.0.0 , monad-logger , mtl , network@@ -81,7 +81,7 @@ , containers , data-default , hashable- , haskoin-core >=0.22.0+ , haskoin-core >=1.0.0 , haskoin-node , hspec , monad-logger
src/Haskoin/Node.hs view
@@ -1,190 +1,192 @@-{-# LANGUAGE FlexibleContexts #-}-{-# LANGUAGE GADTs #-}-{-# LANGUAGE LambdaCase #-}+{-# LANGUAGE DuplicateRecordFields #-}+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE LambdaCase #-} {-# LANGUAGE MultiParamTypeClasses #-}+{-# LANGUAGE OverloadedRecordDot #-}+{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE NoFieldSelectors #-}+ module Haskoin.Node- ( module Haskoin.Node.Peer- , module Haskoin.Node.Manager- , module Haskoin.Node.Chain- , NodeConfig (..)- , NodeEvent (..)- , Node (..)- , withNode- , withConnection- ) where+ ( module Haskoin.Node.Peer,+ module Haskoin.Node.PeerMgr,+ module Haskoin.Node.Chain,+ NodeConfig (..),+ NodeEvent (..),+ Node (..),+ withNode,+ withConnection,+ )+where -import Control.Monad (forever)-import Control.Monad.Logger (MonadLoggerIO)-import Data.Conduit.Network (appSink, appSource, clientSettings,- runTCPClient, ClientSettings)-import Data.String.Conversions (cs)-import Data.Time.Clock (NominalDiffTime)-import Database.RocksDB (ColumnFamily, DB)-import Haskoin (Addr (..), BlockNode (..),- Headers (..), Message (..), Network,- NetworkAddress, Ping (..), Pong (..))-import Haskoin.Node.Chain-import Haskoin.Node.Manager-import Haskoin.Node.Peer-import Network.Socket (NameInfoFlag (..), SockAddr,- getNameInfo)-import NQE (Inbox, Publisher, publish, receive,- withPublisher, withSubscription)-import Text.Read (readMaybe)-import UnliftIO (MonadUnliftIO, SomeException, catch,- liftIO, link, throwIO, withAsync)+import Control.Monad (forever)+import Control.Monad.Cont (ContT (..), MonadCont (callCC), cont, lift, runCont, runContT)+import Control.Monad.Logger (MonadLoggerIO)+import Data.Conduit.Network+ ( ClientSettings,+ appSink,+ appSource,+ clientSettings,+ runTCPClient,+ )+import Data.String.Conversions (cs)+import Data.Time.Clock (NominalDiffTime)+import Database.RocksDB (ColumnFamily, DB)+import Haskoin+ ( Addr (..),+ BlockNode (..),+ Headers (..),+ Message (..),+ Network,+ NetworkAddress,+ Ping (..),+ Pong (..),+ )+import Haskoin.Node.Chain+import Haskoin.Node.Peer+import Haskoin.Node.PeerMgr+import NQE+ ( Inbox,+ Publisher,+ publish,+ receive,+ withPublisher,+ withSubscription,+ )+import Network.Socket+ ( NameInfoFlag (..),+ SockAddr,+ getNameInfo,+ )+import Text.Read (readMaybe)+import UnliftIO+ ( MonadUnliftIO,+ SomeException,+ catch,+ liftIO,+ link,+ throwIO,+ withAsync,+ ) -- | General node configuration. data NodeConfig = NodeConfig- { nodeConfMaxPeers :: !Int- -- ^ maximum number of connected peers allowed- , nodeConfDB :: !DB- -- ^ database handler- , nodeConfColumnFamily :: !(Maybe ColumnFamily)- -- ^ database column family- , nodeConfPeers :: ![String]- -- ^ static list of peers to connect to- , nodeConfDiscover :: !Bool- -- ^ activate peer discovery- , nodeConfNetAddr :: !NetworkAddress- -- ^ network address for the local host- , nodeConfNet :: !Network- -- ^ network constants- , nodeConfEvents :: !(Publisher NodeEvent)- -- ^ node events are sent to this publisher- , nodeConfTimeout :: !NominalDiffTime- -- ^ timeout in seconds- , nodeConfPeerMaxLife :: !NominalDiffTime- -- ^ peer disconnect after seconds- , nodeConfConnect :: !(SockAddr -> WithConnection)- }+ { -- | maximum number of connected peers allowed+ maxPeers :: !Int,+ -- | database handler+ db :: !DB,+ -- | database column family+ cf :: !(Maybe ColumnFamily),+ -- | static list of peers to connect to+ peers :: ![String],+ -- | activate peer discovery+ discover :: !Bool,+ -- | network address for the local host+ address :: !NetworkAddress,+ -- | network constants+ net :: !Network,+ -- | node events are sent to this publisher+ pub :: !(Publisher NodeEvent),+ -- | timeout in seconds+ timeout :: !NominalDiffTime,+ -- | peer disconnect after seconds+ maxPeerLife :: !NominalDiffTime,+ connect :: !(SockAddr -> WithConnection)+ } -data Node = Node { nodeManager :: !PeerManager- , nodeChain :: !Chain- }+data Node = Node+ { peerMgr :: !PeerMgr,+ chain :: !Chain+ } data NodeEvent- = ChainEvent !ChainEvent- | PeerEvent !PeerEvent- | PeerMessage !Peer !Message- deriving Eq+ = ChainEvent !ChainEvent+ | PeerEvent !PeerEvent+ deriving (Eq) withConnection :: SockAddr -> WithConnection withConnection na f =- fromSockAddr na >>= \case- Nothing -> throwIO PeerAddressInvalid- Just cset ->- runTCPClient cset $ \ad ->- f (Conduits (appSource ad) (appSink ad))+ fromSockAddr na >>= \case+ Nothing -> throwIO PeerAddressInvalid+ Just cset ->+ runTCPClient cset $ \ad ->+ f (Conduits (appSource ad) (appSink ad)) fromSockAddr ::- (MonadUnliftIO m) => SockAddr -> m (Maybe ClientSettings)+ (MonadUnliftIO m) => SockAddr -> m (Maybe ClientSettings) fromSockAddr sa = go `catch` e where go = do- (maybe_host, maybe_port) <- liftIO (getNameInfo flags True True sa)- return $- clientSettings - <$> (readMaybe =<< maybe_port) + (maybe_host, maybe_port) <- liftIO (getNameInfo flags True True sa)+ return $+ clientSettings+ <$> (readMaybe =<< maybe_port) <*> (cs <$> maybe_host) flags = [NI_NUMERICHOST, NI_NUMERICSERV]- e :: Monad m => SomeException -> m (Maybe a)+ e :: (Monad m) => SomeException -> m (Maybe a) e _ = return Nothing -chainForwarder :: MonadLoggerIO m- => PeerManager- -> Publisher NodeEvent- -> Inbox ChainEvent- -> m ()-chainForwarder mgr pub inbox =- forever $ receive inbox >>= \event -> do- case event of- ChainBestBlock bb ->- managerBest (nodeHeight bb) mgr- _ -> return ()- publish (ChainEvent event) pub--managerForwarder :: MonadLoggerIO m- => Chain- -> Publisher NodeEvent- -> Inbox PeerEvent- -> m ()-managerForwarder ch pub inbox =- forever $ receive inbox >>= \event -> do- case event of- PeerConnected p ->- chainPeerConnected p ch- PeerDisconnected p ->- chainPeerDisconnected p ch- publish (PeerEvent event) pub+chainEvents ::+ (MonadUnliftIO m, MonadLoggerIO m) =>+ PeerMgr ->+ Inbox ChainEvent ->+ Publisher NodeEvent ->+ m ()+chainEvents mgr input output = forever $ do+ event <- receive input+ case event of+ ChainBestBlock bb ->+ peerMgrBest bb.height mgr+ _ -> return ()+ publish (ChainEvent event) output -peerForwarder :: MonadLoggerIO m- => Chain- -> PeerManager- -> Publisher NodeEvent- -> Inbox (Peer, Message)- -> m ()-peerForwarder ch mgr pub inbox =- forever $ receive inbox >>= \(p, msg) -> do- case msg of- MVersion v ->- managerVersion p v mgr- MVerAck ->- managerVerAck p mgr- MPing (Ping n) ->- managerPing p n mgr- MPong (Pong n) ->- managerPong p n mgr- MAddr (Addr ns) ->- managerAddrs p (map snd ns) mgr- MHeaders (Headers hs) ->- chainHeaders p (map fst hs) ch- _ -> return ()- managerTickle p mgr- publish (PeerMessage p msg) pub+peerEvents ::+ (MonadUnliftIO m, MonadLoggerIO m) =>+ Chain ->+ PeerMgr ->+ Inbox PeerEvent ->+ Publisher NodeEvent ->+ m ()+peerEvents ch mgr input output = forever $ do+ event <- receive input+ case event of+ PeerConnected p ->+ chainPeerConnected p ch+ PeerDisconnected p ->+ chainPeerDisconnected p ch+ PeerMessage p msg -> do+ case msg of+ MVersion v ->+ peerMgrVersion p v mgr+ MVerAck ->+ peerMgrVerAck p mgr+ MPing (Ping n) ->+ peerMgrPing p n mgr+ MPong (Pong n) ->+ peerMgrPong p n mgr+ MAddr (Addr ns) ->+ peerMgrAddrs p (map snd ns) mgr+ MHeaders (Headers hs) ->+ chainHeaders p (map fst hs) ch+ _ -> return ()+ peerMgrTickle p mgr+ publish (PeerEvent event) output -- | Launch node process in the foreground. withNode ::- ( MonadLoggerIO m- , MonadUnliftIO m- )- => NodeConfig- -> (Node -> m a)- -> m a-withNode cfg action =- withPublisher $ \peer_pub ->- withPublisher $ \mgr_pub ->- withPublisher $ \ch_pub ->- withSubscription peer_pub $ \peer_inbox ->- withSubscription mgr_pub $ \mgr_inbox ->- withSubscription ch_pub $ \ch_inbox ->- withPeerManager (mgr_config mgr_pub peer_pub) $ \mgr ->- withChain (chain_config ch_pub) $ \ch ->- withAsync (peerForwarder ch mgr pub peer_inbox) $ \a ->- withAsync (managerForwarder ch pub mgr_inbox) $ \b ->- withAsync (chainForwarder mgr pub ch_inbox) $ \c ->- link a >> link b >> link c >>- action Node { nodeManager = mgr, nodeChain = ch }- where- pub = nodeConfEvents cfg- chain_config ch_pub =- ChainConfig- { chainConfDB = nodeConfDB cfg- , chainConfColumnFamily = nodeConfColumnFamily cfg- , chainConfNetwork = nodeConfNet cfg- , chainConfEvents = ch_pub- , chainConfTimeout = nodeConfTimeout cfg- }- mgr_config mgr_pub peer_pub =- PeerManagerConfig- { peerManagerMaxPeers = nodeConfMaxPeers cfg- , peerManagerPeers = nodeConfPeers cfg- , peerManagerDiscover = nodeConfDiscover cfg- , peerManagerNetAddr = nodeConfNetAddr cfg- , peerManagerNetwork = nodeConfNet cfg- , peerManagerEvents = mgr_pub- , peerManagerMaxLife = nodeConfPeerMaxLife cfg- , peerManagerTimeout = nodeConfTimeout cfg- , peerManagerConnect = nodeConfConnect cfg- , peerManagerPub = peer_pub- }+ (MonadLoggerIO m, MonadUnliftIO m) =>+ NodeConfig ->+ (Node -> m a) ->+ m a+withNode NodeConfig {..} action = flip runContT return $ do+ peerPub <- ContT withPublisher+ peerSub <- ContT (withSubscription peerPub)+ chainPub <- ContT withPublisher+ chainSub <- ContT (withSubscription chainPub)+ let peerMgrCfg = PeerMgrConfig {pub = peerPub, ..}+ let chainCfg = ChainConfig {pub = chainPub, ..}+ chain <- ContT (withChain chainCfg)+ peerMgr <- ContT $ withPeerMgr peerMgrCfg+ lift . link =<< ContT (withAsync $ chainEvents peerMgr chainSub pub)+ lift . link =<< ContT (withAsync $ peerEvents chain peerMgr peerSub pub)+ lift $ action Node {..}
src/Haskoin/Node/Chain.hs view
@@ -1,683 +1,772 @@-{-# LANGUAGE ConstraintKinds #-}-{-# LANGUAGE ExistentialQuantification #-}-{-# LANGUAGE FlexibleContexts #-}-{-# LANGUAGE FlexibleInstances #-}-{-# LANGUAGE LambdaCase #-}-{-# LANGUAGE MultiParamTypeClasses #-}-{-# LANGUAGE OverloadedStrings #-}-{-# LANGUAGE TemplateHaskell #-}-{-# LANGUAGE UndecidableInstances #-}-{-# OPTIONS_GHC -fno-warn-orphans #-}-module Haskoin.Node.Chain- ( ChainConfig (..)- , ChainEvent (..)- , Chain- , withChain- , chainGetBlock- , chainGetBest- , chainGetAncestor- , chainGetParents- , chainGetSplitBlock- , chainPeerConnected- , chainPeerDisconnected- , chainIsSynced- , chainBlockMain- , chainHeaders- ) where--import Control.Monad (forM_, forever, guard, when)-import Control.Monad.Except (runExceptT, throwError)-import Control.Monad.Logger (MonadLoggerIO, logDebugS, logErrorS,- logInfoS)-import Control.Monad.Reader (MonadReader, ReaderT (..), asks,- runReaderT)-import Control.Monad.Trans (lift)-import Control.Monad.Trans.Maybe (MaybeT (..), runMaybeT)-import qualified Data.ByteString as B-import Data.Function (on)-import Data.List (delete, nub)-import Data.Maybe (isJust, isNothing)-import Data.Serialize (Serialize, get, getWord8, put,- putWord8)-import Data.String.Conversions (cs)-import Data.Time.Clock (NominalDiffTime, UTCTime,- diffUTCTime, getCurrentTime)-import Data.Time.Clock.POSIX (posixSecondsToUTCTime,- utcTimeToPOSIXSeconds)-import Data.Word (Word32)-import Database.RocksDB (ColumnFamily, DB)-import qualified Database.RocksDB as R-import Database.RocksDB.Query (Key, KeyValue, insert, insertCF,- insertOp, insertOpCF,- retrieveCommon, writeBatch)-import Haskoin (BlockHash, BlockHeader (..),- BlockHeaders (..), BlockHeight,- BlockNode (..), GetHeaders (..),- Message (..), Network, blockLocator,- connectBlocks, genesisNode,- getAncestor, headerHash, splitPoint)-import Haskoin.Node.Manager (myVersion)-import Haskoin.Node.Peer-import NQE (Mailbox, Publisher, newMailbox,- publish, receive, send)-import System.Random (randomRIO)-import UnliftIO (MonadIO, MonadUnliftIO, TVar,- atomically, liftIO, link,- modifyTVar, newTVarIO, readTVar,- readTVarIO, withAsync, writeTVar)-import UnliftIO.Concurrent (threadDelay)---- | Mailbox for chain header syncing process.-data Chain = Chain { chainMailbox :: !(Mailbox ChainMessage)- , chainReader :: !ChainReader- }--instance Eq Chain where- (==) = (==) `on` chainMailbox---- | Configuration for chain syncing process.-data ChainConfig =- ChainConfig- { chainConfDB :: !DB- -- ^ database handle- , chainConfColumnFamily :: !(Maybe ColumnFamily)- -- ^ column family- , chainConfNetwork :: !Network- -- ^ network constants- , chainConfEvents :: !(Publisher ChainEvent)- -- ^ send header chain events here- , chainConfTimeout :: !NominalDiffTime- -- ^ timeout in seconds- }--data ChainMessage- = ChainHeaders !Peer ![BlockHeader]- | ChainPeerConnected !Peer- | ChainPeerDisconnected !Peer- | ChainPing---- | Events originating from chain syncing process.-data ChainEvent- = ChainBestBlock !BlockNode- -- ^ chain has new best block- | ChainSynced !BlockNode- -- ^ chain is in sync with the network- deriving (Eq, Show)--type MonadChain m =- ( MonadLoggerIO m- , MonadUnliftIO m- , MonadReader ChainReader m )---- | Reader for header synchronization code.-data ChainReader = ChainReader- { myConfig :: !ChainConfig- -- ^ placeholder for upstream data- , chainState :: !(TVar ChainState)- -- ^ mutable state for header synchronization- }---- | Database key for version.-data ChainDataVersionKey = ChainDataVersionKey- deriving (Eq, Ord, Show)--instance Key ChainDataVersionKey-instance KeyValue ChainDataVersionKey Word32--instance Serialize ChainDataVersionKey where- get = do- guard . (== 0x92) =<< getWord8- return ChainDataVersionKey- put ChainDataVersionKey = putWord8 0x92--data ChainSync = ChainSync- { chainSyncPeer :: !Peer- , chainTimestamp :: !UTCTime- , chainHighest :: !(Maybe BlockNode)- }---- | Mutable state for the header chain process.-data ChainState = ChainState- { chainSyncing :: !(Maybe ChainSync)- -- ^ peer to sync against and time of last received message- , newPeers :: ![Peer]- -- ^ queue of peers to sync against- , mySynced :: !Bool- -- ^ has the header chain ever been considered synced?- }---- | Key for block header in database.-newtype BlockHeaderKey = BlockHeaderKey BlockHash deriving (Eq, Show)--instance Serialize BlockHeaderKey where- get = do- guard . (== 0x90) =<< getWord8- BlockHeaderKey <$> get- put (BlockHeaderKey bh) = do- putWord8 0x90- put bh---- | Key for best block in database.-data BestBlockKey = BestBlockKey deriving (Eq, Show)--instance KeyValue BlockHeaderKey BlockNode-instance KeyValue BestBlockKey BlockNode--instance Serialize BestBlockKey where- get = do- guard . (== 0x91) =<< getWord8- return BestBlockKey- put BestBlockKey = putWord8 0x91--instance MonadIO m => BlockHeaders (ReaderT ChainConfig m) where- addBlockHeader bn = do- db <- asks chainConfDB- asks chainConfColumnFamily >>= \case- Nothing -> insert db (BlockHeaderKey h) bn- Just cf -> insertCF db cf (BlockHeaderKey h) bn- where- h = headerHash (nodeHeader bn)- getBlockHeader bh = do- db <- asks chainConfDB- mcf <- asks chainConfColumnFamily- retrieveCommon db mcf (BlockHeaderKey bh)- getBestBlockHeader = do- db <- asks chainConfDB- mcf <- asks chainConfColumnFamily- retrieveCommon db mcf BestBlockKey >>= \case- Nothing -> error "Could not get best block from database"- Just b -> return b- setBestBlockHeader bn = do- db <- asks chainConfDB- asks chainConfColumnFamily >>= \case- Nothing -> insert db BestBlockKey bn- Just cf -> insertCF db cf BestBlockKey bn- addBlockHeaders bns = do- db <- asks chainConfDB- mcf <- asks chainConfColumnFamily- writeBatch db (map (f mcf) bns)- where- h bn = headerHash (nodeHeader bn)- f Nothing bn = insertOp (BlockHeaderKey (h bn)) bn- f (Just cf) bn = insertOpCF cf (BlockHeaderKey (h bn)) bn--instance MonadIO m => BlockHeaders (ReaderT Chain m) where- getBlockHeader bh = ReaderT $ chainGetBlock bh- getBestBlockHeader = ReaderT chainGetBest- addBlockHeader _ = undefined- setBestBlockHeader _ = undefined- addBlockHeaders _ = undefined--withBlockHeaders :: MonadChain m => ReaderT ChainConfig m a -> m a-withBlockHeaders f = do- cfg <- asks myConfig- runReaderT f cfg--withChain ::- (MonadUnliftIO m, MonadLoggerIO m)- => ChainConfig- -> (Chain -> m a)- -> m a-withChain cfg action = do- (inbox, mailbox) <- newMailbox- $(logDebugS) "Chain" "Starting chain actor"- st <- newTVarIO ChainState { chainSyncing = Nothing- , mySynced = False- , newPeers = []- }- let rd = ChainReader { myConfig = cfg- , chainState = st- }- ch = Chain { chainReader = rd- , chainMailbox = mailbox- }- runReaderT initChainDB rd- withAsync (main_loop ch rd inbox) $ \a ->- link a >> action ch- where- main_loop ch rd inbox = withSyncLoop ch $- runReaderT (run inbox) rd- run inbox = do- withBlockHeaders getBestBlockHeader >>=- chainEvent . ChainBestBlock- forever $ do- msg <- receive inbox- chainMessage msg--chainEvent :: MonadChain m => ChainEvent -> m ()-chainEvent e = do- pub <- asks (chainConfEvents . myConfig)- case e of- ChainBestBlock b ->- $(logInfoS) "Chain" $- "Best block header at height: "- <> cs (show (nodeHeight b))- ChainSynced b ->- $(logInfoS) "Chain" $- "Headers in sync at height: "- <> cs (show (nodeHeight b))- publish e pub--processHeaders :: MonadChain m => Peer -> [BlockHeader] -> m ()-processHeaders p hs = do- $(logDebugS) "Chain" $- "Processing " <> cs (show (length hs))- <> " headers from peer: " <> peerText p- net <- asks (chainConfNetwork . myConfig)- now <- liftIO getCurrentTime- pbest <- withBlockHeaders getBestBlockHeader- importHeaders net now hs >>= \case- Left e -> do- $(logErrorS) "Chain" $- "Could not connect headers from peer: "- <> peerText p- e `killPeer` p- Right done -> do- setLastReceived- best <- withBlockHeaders getBestBlockHeader- when (nodeHeader pbest /= nodeHeader best) $- chainEvent (ChainBestBlock best)- if done- then do- MSendHeaders `sendMessage` p- finishPeer p- syncNewPeer- syncNotif- else syncPeer p--syncNewPeer :: MonadChain m => m ()-syncNewPeer = getSyncingPeer >>= \case- Just _ -> return ()- Nothing -> nextPeer >>= \case- Nothing -> return ()- Just p -> do- $(logDebugS) "Chain" $- "Syncing against peer: " <> peerText p- syncPeer p--syncNotif :: MonadChain m => m ()-syncNotif =- notifySynced >>= \case- False -> return ()- True -> withBlockHeaders getBestBlockHeader >>=- chainEvent . ChainSynced--syncPeer :: MonadChain m => Peer -> m ()-syncPeer p = do- t <- liftIO getCurrentTime- m <- chainSyncingPeer >>= \case- Just ChainSync { chainSyncPeer = s- , chainHighest = m- }- | p == s -> syncing_me t m- | otherwise -> return Nothing- Nothing -> syncing_new t- forM_ m $ \g -> do- $(logDebugS) "Chain" $- "Requesting headers from peer: "- <> peerText p- MGetHeaders g `sendMessage` p- where- syncing_new t =- setSyncingPeer p >>= \case- False -> return Nothing- True -> do- $(logDebugS) "Chain" $- "Locked peer: " <> peerText p- h <- withBlockHeaders getBestBlockHeader- Just <$> syncHeaders t h p- syncing_me t m = do- h <- case m of- Nothing -> withBlockHeaders getBestBlockHeader- Just h -> return h- Just <$> syncHeaders t h p--chainMessage :: MonadChain m => ChainMessage -> m ()--chainMessage (ChainHeaders p hs) =- processHeaders p hs--chainMessage (ChainPeerConnected p) = do- $(logDebugS) "Chain" $ "Peer connected: " <> peerText p- addPeer p- syncNewPeer--chainMessage (ChainPeerDisconnected p) = do- $(logDebugS) "Chain" $ "Peer disconnected: " <> peerText p- finishPeer p- syncNewPeer--chainMessage ChainPing = do- to <- asks (chainConfTimeout . myConfig)- now <- liftIO getCurrentTime- chainSyncingPeer >>= \case- Just ChainSync {chainSyncPeer = p, chainTimestamp = t}- | now `diffUTCTime` t > to -> do- $(logErrorS) "Chain" $- "Syncing peer timed out: " <> peerText p- PeerTimeout `killPeer` p- | otherwise -> return ()- Nothing -> syncNewPeer--withSyncLoop :: (MonadUnliftIO m, MonadLoggerIO m)- => Chain -> m a -> m a-withSyncLoop ch f =- withAsync go $ \a ->- link a >> f- where- go = forever $ do- delay <- liftIO $- randomRIO ( 2 * 1000 * 1000- , 20 * 1000 * 1000 )- threadDelay delay- ChainPing `send` chainMailbox ch---- | Version of the database.-dataVersion :: Word32-dataVersion = 1---- | Initialize header database. If version is different from current, the--- database is purged of conflicting elements first.-initChainDB :: MonadChain m => m ()-initChainDB = do- db <- asks (chainConfDB . myConfig)- mcf <- asks (chainConfColumnFamily . myConfig)- net <- asks (chainConfNetwork . myConfig)- ver <- retrieveCommon db mcf ChainDataVersionKey- when (ver /= Just dataVersion) $ purgeChainDB >>= writeBatch db- case mcf of- Nothing -> insert db ChainDataVersionKey dataVersion- Just cf -> insertCF db cf ChainDataVersionKey dataVersion- retrieveCommon db mcf BestBlockKey >>= \b ->- when (isNothing (b :: Maybe BlockNode)) $- withBlockHeaders $ do- addBlockHeader (genesisNode net)- setBestBlockHeader (genesisNode net)---- | Purge database of elements having keys that may conflict with those used in--- this module.-purgeChainDB :: MonadChain m => m [R.BatchOp]-purgeChainDB = do- db <- asks (chainConfDB . myConfig)- mcf <- asks (chainConfColumnFamily . myConfig)- f db mcf $ \it -> do- R.iterSeek it $ B.singleton 0x90- recurse_delete it db mcf- where- f db Nothing = R.withIter db- f db (Just cf) = R.withIterCF db cf- recurse_delete it db mcf =- R.iterKey it >>= \case- Just k- | B.head k == 0x90 || B.head k == 0x91 -> do- case mcf of- Nothing -> R.delete db k- Just cf -> R.deleteCF db cf k- R.iterNext it- (R.Del k :) <$> recurse_delete it db mcf- _ -> return []---- | Import a bunch of continuous headers. Returns 'True' if the number of--- headers is 2000, which means that there are possibly more headers to sync--- from whatever peer delivered these.-importHeaders :: MonadChain m- => Network- -> UTCTime- -> [BlockHeader]- -> m (Either PeerException Bool)-importHeaders net now hs =- runExceptT $- lift connect >>= \case- Right _ -> do- case hs of- [] -> return ()- _ -> do- bb <- lift get_last- box <- asks chainState- atomically . modifyTVar box $ \s ->- s { chainSyncing =- (\x -> x {chainHighest = bb})- <$> chainSyncing s- }- case length hs of- 2000 -> return False- _ -> return True- Left _ -> throwError PeerSentBadHeaders- where- timestamp = floor (utcTimeToPOSIXSeconds now)- connect = withBlockHeaders $ connectBlocks net timestamp hs- get_last = withBlockHeaders . getBlockHeader . headerHash $ last hs---- | Check if best block header is in sync with the rest of the block chain by--- comparing the best block with the current time, verifying that there are no--- peers in the queue to be synced, and no peer is being synced at the moment.--- This function will only return 'True' once. It should be used to decide--- whether to notify other processes that the header chain has been synced. The--- state of the chain will be flipped to synced when this function returns--- 'True'.-notifySynced :: MonadChain m => m Bool-notifySynced =- fmap isJust $- runMaybeT $ do- bb <- lift $ withBlockHeaders getBestBlockHeader- now <- liftIO getCurrentTime- guard $ now `diffUTCTime` block_time bb > 7200- st <- asks chainState- MaybeT . atomically . runMaybeT $ do- s <- lift $ readTVar st- guard . isNothing $ chainSyncing s- guard . null $ newPeers s- guard . not $ mySynced s- lift $ writeTVar st s {mySynced = True}- return ()- where- block_time =- posixSecondsToUTCTime . fromIntegral . blockTimestamp . nodeHeader---- | Get next peer to sync against from the queue.-nextPeer :: MonadChain m => m (Maybe Peer)-nextPeer = do- ps <- newPeers <$> (readTVarIO =<< asks chainState)- go ps- where- go [] = return Nothing- go (p:ps) =- setSyncingPeer p >>= \case- True -> return (Just p)- False -> go ps---- | Set a syncing peer and generate a 'GetHeaders' data structure with a block--- locator to send to that peer for syncing.-syncHeaders ::- MonadChain m- => UTCTime- -> BlockNode- -> Peer- -> m GetHeaders-syncHeaders now bb p = do- st <- asks chainState- atomically $- modifyTVar st $ \s ->- s { chainSyncing =- Just- ChainSync- { chainSyncPeer = p- , chainTimestamp = now- , chainHighest = Nothing- }- , newPeers = delete p (newPeers s)- }- loc <- withBlockHeaders $ blockLocator bb- return- GetHeaders- { getHeadersVersion = myVersion- , getHeadersBL = loc- , getHeadersHashStop = z- }- where- z = "0000000000000000000000000000000000000000000000000000000000000000"---- | Set the time of last received data to now if a syncing peer is active.-setLastReceived :: MonadChain m => m ()-setLastReceived = do- now <- liftIO getCurrentTime- st <- asks chainState- let f p = p { chainTimestamp = now }- atomically . modifyTVar st $ \s ->- s { chainSyncing = f <$> chainSyncing s }---- | Add a new peer to the queue of peers to sync against.-addPeer :: MonadChain m => Peer -> m ()-addPeer p = do- st <- asks chainState- atomically . modifyTVar st $ \s -> s {newPeers = nub (p : newPeers s)}---- | Get syncing peer if there is one.-getSyncingPeer :: MonadChain m => m (Maybe Peer)-getSyncingPeer =- fmap chainSyncPeer . chainSyncing <$> (readTVarIO =<< asks chainState)--setSyncingPeer :: MonadChain m => Peer -> m Bool-setSyncingPeer p =- setBusy p >>= \case- False -> do- $(logDebugS) "Chain" $- "Could not lock peer: " <> peerText p- return False- True -> do- $(logDebugS) "Chain" $- "Locked peer: " <> peerText p- set_it- return True- where- set_it = do- now <- liftIO getCurrentTime- box <- asks chainState- atomically $ modifyTVar box $ \s ->- s { chainSyncing =- Just ChainSync { chainSyncPeer = p- , chainTimestamp = now- , chainHighest = Nothing- }- }----- | Remove a peer from the queue of peers to sync and unset the syncing peer if--- it is set to the provided peer.-finishPeer :: MonadChain m => Peer -> m ()-finishPeer p =- asks chainState >>= remove_peer >>= \case- False ->- $(logDebugS) "Chain" $- "Removed peer from queue: " <> peerText p- True -> do- $(logDebugS) "Chain" $- "Releasing syncing peer: " <> peerText p- setFree p- where- remove_peer st = atomically $- readTVar st >>= \s -> case chainSyncing s of- Just ChainSync { chainSyncPeer = p' }- | p == p' -> do- unset_syncing st- return True- _ -> do- remove_from_queue st- return False- unset_syncing st =- modifyTVar st $ \x ->- x { chainSyncing = Nothing }- remove_from_queue st =- modifyTVar st $ \x ->- x { newPeers = delete p (newPeers x) }---- | Return syncing peer data.-chainSyncingPeer :: MonadChain m => m (Maybe ChainSync)-chainSyncingPeer =- chainSyncing <$> (readTVarIO =<< asks chainState)---- | Get a block header from 'Chain' process.-chainGetBlock :: MonadIO m- => BlockHash -> Chain -> m (Maybe BlockNode)-chainGetBlock bh ch =- runReaderT (getBlockHeader bh) (myConfig (chainReader ch))---- | Get best block header from chain process.-chainGetBest :: MonadIO m => Chain -> m BlockNode-chainGetBest ch =- runReaderT getBestBlockHeader (myConfig (chainReader ch))---- | Get ancestor of 'BlockNode' at 'BlockHeight' from chain process.-chainGetAncestor :: MonadIO m- => BlockHeight- -> BlockNode- -> Chain- -> m (Maybe BlockNode)-chainGetAncestor h bn ch =- runReaderT (getAncestor h bn) (myConfig (chainReader ch))---- | Get parents of 'BlockNode' starting at 'BlockHeight' from chain process.-chainGetParents :: MonadIO m- => BlockHeight- -> BlockNode- -> Chain- -> m [BlockNode]-chainGetParents height top ch =- go [] top- where- go acc b- | height >= nodeHeight b = return acc- | otherwise = do- m <- chainGetBlock (prevBlock $ nodeHeader b) ch- case m of- Nothing -> return acc- Just p -> go (p : acc) p---- | Get last common block from chain process.-chainGetSplitBlock :: MonadIO m- => BlockNode- -> BlockNode- -> Chain- -> m BlockNode-chainGetSplitBlock l r ch =- runReaderT (splitPoint l r) (myConfig (chainReader ch))---- | Notify chain that a new peer is connected.-chainPeerConnected :: MonadIO m- => Peer- -> Chain- -> m ()-chainPeerConnected p ch =- ChainPeerConnected p `send` chainMailbox ch---- | Notify chain that a peer has disconnected.-chainPeerDisconnected :: MonadIO m- => Peer- -> Chain- -> m ()-chainPeerDisconnected p ch =- ChainPeerDisconnected p `send` chainMailbox ch---- | Is given 'BlockHash' in the main chain?-chainBlockMain :: MonadIO m- => BlockHash- -> Chain- -> m Bool-chainBlockMain bh ch =- chainGetBest ch >>= \bb ->- chainGetBlock bh ch >>= \case- Nothing ->- return False- bm@(Just bn) ->- (== bm) <$> chainGetAncestor (nodeHeight bn) bb ch---- | Is chain in sync with network?-chainIsSynced :: MonadIO m => Chain -> m Bool-chainIsSynced ch =- mySynced <$> readTVarIO (chainState (chainReader ch))---- | Peer sends a bunch of headers to the chain process.-chainHeaders :: MonadIO m- => Peer -> [BlockHeader] -> Chain -> m ()-chainHeaders p hs ch =- ChainHeaders p hs `send` chainMailbox ch+{-# LANGUAGE ConstraintKinds #-}+{-# LANGUAGE DuplicateRecordFields #-}+{-# LANGUAGE ExistentialQuantification #-}+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE FlexibleInstances #-}+{-# LANGUAGE LambdaCase #-}+{-# LANGUAGE MultiParamTypeClasses #-}+{-# LANGUAGE OverloadedRecordDot #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE UndecidableInstances #-}+{-# LANGUAGE NoFieldSelectors #-}+{-# OPTIONS_GHC -fno-warn-orphans #-}++module Haskoin.Node.Chain+ ( ChainConfig (..),+ ChainEvent (..),+ Chain,+ withChain,+ chainGetBlock,+ chainGetBest,+ chainGetAncestor,+ chainGetParents,+ chainGetSplitBlock,+ chainPeerConnected,+ chainPeerDisconnected,+ chainIsSynced,+ chainBlockMain,+ chainHeaders,+ )+where++import Control.Monad (forM_, forever, guard, when)+import Control.Monad.Except (runExceptT, throwError)+import Control.Monad.Logger+ ( MonadLoggerIO,+ logDebugS,+ logErrorS,+ logInfoS,+ )+import Control.Monad.Reader+ ( MonadReader,+ ReaderT (..),+ asks,+ runReaderT,+ )+import Control.Monad.Trans (lift)+import Control.Monad.Trans.Maybe (MaybeT (..), runMaybeT)+import qualified Data.ByteString as B+import Data.Function (on)+import Data.List (delete, nub)+import Data.Maybe (isJust, isNothing)+import Data.Serialize+ ( Serialize,+ get,+ getWord8,+ put,+ putWord8,+ )+import Data.String.Conversions (cs)+import Data.Time.Clock+ ( NominalDiffTime,+ UTCTime,+ diffUTCTime,+ getCurrentTime,+ )+import Data.Time.Clock.POSIX+ ( posixSecondsToUTCTime,+ utcTimeToPOSIXSeconds,+ )+import Data.Word (Word32)+import Database.RocksDB (ColumnFamily, DB)+import qualified Database.RocksDB as R+import Database.RocksDB.Query+ ( Key,+ KeyValue,+ insert,+ insertCF,+ insertOp,+ insertOpCF,+ retrieveCommon,+ writeBatch,+ )+import Haskoin+ ( BlockHash,+ BlockHeader (..),+ BlockHeaders (..),+ BlockHeight,+ BlockNode (..),+ GetHeaders (..),+ Message (..),+ Network,+ blockLocator,+ connectBlocks,+ genesisNode,+ getAncestor,+ headerHash,+ splitPoint,+ )+import Haskoin.Node.Peer+import Haskoin.Node.PeerMgr (myVersion)+import NQE+ ( Mailbox,+ Publisher,+ newMailbox,+ publish,+ receive,+ send,+ )+import System.Random (randomRIO)+import UnliftIO+ ( MonadIO,+ MonadUnliftIO,+ TVar,+ atomically,+ liftIO,+ link,+ modifyTVar,+ newTVarIO,+ readTVar,+ readTVarIO,+ withAsync,+ writeTVar,+ )+import UnliftIO.Concurrent (threadDelay)++-- | Mailbox for chain header syncing process.+data Chain = Chain+ { mailbox :: !(Mailbox ChainMessage),+ reader :: !ChainReader+ }++instance Eq Chain where+ (==) = (==) `on` (.mailbox)++-- | Configuration for chain syncing process.+data ChainConfig = ChainConfig+ { -- | database handle+ db :: !DB,+ -- | column family+ cf :: !(Maybe ColumnFamily),+ -- | network constants+ net :: !Network,+ -- | send header chain events here+ pub :: !(Publisher ChainEvent),+ -- | timeout in seconds+ timeout :: !NominalDiffTime+ }++data ChainMessage+ = ChainHeaders !Peer ![BlockHeader]+ | ChainPeerConnected !Peer+ | ChainPeerDisconnected !Peer+ | ChainPing++-- | Events originating from chain syncing process.+data ChainEvent+ = -- | chain has new best block+ ChainBestBlock !BlockNode+ | -- | chain is in sync with the network+ ChainSynced !BlockNode+ deriving (Eq, Show)++type MonadChain m =+ ( MonadLoggerIO m,+ MonadUnliftIO m,+ MonadReader ChainReader m+ )++-- | State and configuration.+data ChainReader = ChainReader+ { -- | placeholder for upstream data+ config :: !ChainConfig,+ -- | mutable state for header synchronization+ state :: !(TVar ChainState)+ }++-- | Database key for version.+data ChainDataVersionKey = ChainDataVersionKey+ deriving (Eq, Ord, Show)++instance Key ChainDataVersionKey++instance KeyValue ChainDataVersionKey Word32++instance Serialize ChainDataVersionKey where+ get = do+ guard . (== 0x92) =<< getWord8+ return ChainDataVersionKey+ put ChainDataVersionKey = putWord8 0x92++data ChainSync = ChainSync+ { peer :: !Peer,+ timestamp :: !UTCTime,+ best :: !(Maybe BlockNode)+ }++-- | Mutable state for the header chain process.+data ChainState = ChainState+ { -- | peer to sync against and time of last received message+ syncing :: !(Maybe ChainSync),+ -- | queue of peers to sync against+ peers :: ![Peer],+ -- | has the header chain ever been considered synced?+ beenInSync :: !Bool+ }++-- | Key for block header in database.+newtype BlockHeaderKey = BlockHeaderKey BlockHash deriving (Eq, Show)++instance Serialize BlockHeaderKey where+ get = do+ guard . (== 0x90) =<< getWord8+ BlockHeaderKey <$> get+ put (BlockHeaderKey bh) = do+ putWord8 0x90+ put bh++-- | Key for best block in database.+data BestBlockKey = BestBlockKey deriving (Eq, Show)++instance KeyValue BlockHeaderKey BlockNode++instance KeyValue BestBlockKey BlockNode++instance Serialize BestBlockKey where+ get = do+ guard . (== 0x91) =<< getWord8+ return BestBlockKey+ put BestBlockKey = putWord8 0x91++instance (MonadIO m) => BlockHeaders (ReaderT ChainConfig m) where+ addBlockHeader bn = do+ db <- asks (.db)+ asks (.cf) >>= \case+ Nothing -> insert db (BlockHeaderKey h) bn+ Just cf -> insertCF db cf (BlockHeaderKey h) bn+ where+ h = headerHash bn.header+ getBlockHeader bh = do+ db <- asks (.db)+ mcf <- asks (.cf)+ retrieveCommon db mcf (BlockHeaderKey bh)+ getBestBlockHeader = do+ db <- asks (.db)+ mcf <- asks (.cf)+ retrieveCommon db mcf BestBlockKey >>= \case+ Nothing -> error "Could not get best block from database"+ Just b -> return b+ setBestBlockHeader bn = do+ db <- asks (.db)+ asks (.cf) >>= \case+ Nothing -> insert db BestBlockKey bn+ Just cf -> insertCF db cf BestBlockKey bn+ addBlockHeaders bns = do+ db <- asks (.db)+ mcf <- asks (.cf)+ writeBatch db (map (f mcf) bns)+ where+ h bn = headerHash bn.header+ f Nothing bn = insertOp (BlockHeaderKey (h bn)) bn+ f (Just cf) bn = insertOpCF cf (BlockHeaderKey (h bn)) bn++instance (MonadIO m) => BlockHeaders (ReaderT Chain m) where+ getBlockHeader bh = ReaderT $ chainGetBlock bh+ getBestBlockHeader = ReaderT chainGetBest+ addBlockHeader _ = undefined+ setBestBlockHeader _ = undefined+ addBlockHeaders _ = undefined++withBlockHeaders :: (MonadChain m) => ReaderT ChainConfig m a -> m a+withBlockHeaders f = do+ cfg <- asks (.config)+ runReaderT f cfg++withChain ::+ (MonadUnliftIO m, MonadLoggerIO m) =>+ ChainConfig ->+ (Chain -> m a) ->+ m a+withChain cfg action = do+ (inbox, mailbox) <- newMailbox+ $(logDebugS) "Chain" "Starting chain actor"+ st <-+ newTVarIO+ ChainState+ { syncing = Nothing,+ beenInSync = False,+ peers = []+ }+ let rd = ChainReader {config = cfg, state = st}+ ch = Chain {reader = rd, mailbox = mailbox}+ runReaderT initChainDB rd+ withAsync (main_loop ch rd inbox) $ \a ->+ link a >> action ch+ where+ main_loop ch rd inbox =+ withSyncLoop ch $+ runReaderT (run inbox) rd+ run inbox = do+ withBlockHeaders getBestBlockHeader+ >>= chainEvent . ChainBestBlock+ forever $ do+ $(logDebugS) "Chain" "Awaiting event..."+ msg <- receive inbox+ chainMessage msg++chainEvent :: (MonadChain m) => ChainEvent -> m ()+chainEvent e = do+ pub <- asks (.config.pub)+ case e of+ ChainBestBlock b ->+ $(logInfoS) "Chain" $+ "Best block header at height: "+ <> cs (show b.height)+ ChainSynced b ->+ $(logInfoS) "Chain" $+ "Headers in sync at height: "+ <> cs (show b.height)+ publish e pub++processHeaders :: (MonadChain m) => Peer -> [BlockHeader] -> m ()+processHeaders p hs = do+ $(logDebugS) "Chain" $+ "Processing "+ <> cs (show (length hs))+ <> " headers from peer: "+ <> p.label+ net <- asks (.config.net)+ now <- liftIO getCurrentTime+ pbest <- withBlockHeaders getBestBlockHeader+ importHeaders net now hs >>= \case+ Left e -> do+ $(logErrorS) "Chain" $+ "Could not connect headers from peer: "+ <> p.label+ e `killPeer` p+ Right done -> do+ setLastReceived+ best <- withBlockHeaders getBestBlockHeader+ when (pbest.header /= best.header) $+ chainEvent (ChainBestBlock best)+ if done+ then do+ MSendHeaders `sendMessage` p+ finishPeer p+ syncNewPeer+ syncNotif+ else syncPeer p++syncNewPeer :: (MonadChain m) => m ()+syncNewPeer =+ getSyncingPeer >>= \case+ Just _ -> return ()+ Nothing ->+ nextPeer >>= \case+ Nothing -> return ()+ Just p -> do+ $(logDebugS) "Chain" $+ "Syncing against peer: " <> p.label+ syncPeer p++syncNotif :: (MonadChain m) => m ()+syncNotif =+ notifySynced >>= \case+ False -> return ()+ True ->+ withBlockHeaders getBestBlockHeader+ >>= chainEvent . ChainSynced++syncPeer :: (MonadChain m) => Peer -> m ()+syncPeer p = do+ t <- liftIO getCurrentTime+ m <-+ chainSyncingPeer >>= \case+ Just+ ChainSync+ { peer = s,+ best = m+ }+ | p == s -> syncing_me t m+ | otherwise -> return Nothing+ Nothing -> syncing_new t+ forM_ m $ \g -> do+ $(logDebugS) "Chain" $+ "Requesting headers from peer: "+ <> p.label+ MGetHeaders g `sendMessage` p+ where+ syncing_new t =+ setSyncingPeer p >>= \case+ False -> return Nothing+ True -> do+ $(logDebugS) "Chain" $+ "Locked peer: " <> p.label+ h <- withBlockHeaders getBestBlockHeader+ Just <$> syncHeaders t h p+ syncing_me t m = do+ h <- case m of+ Nothing -> withBlockHeaders getBestBlockHeader+ Just h -> return h+ Just <$> syncHeaders t h p++chainMessage :: (MonadChain m) => ChainMessage -> m ()+chainMessage (ChainHeaders p hs) =+ processHeaders p hs+chainMessage (ChainPeerConnected p) = do+ $(logDebugS) "Chain" $ "Peer connected: " <> p.label+ addPeer p+ syncNewPeer+chainMessage (ChainPeerDisconnected p) = do+ $(logDebugS) "Chain" $ "Peer disconnected: " <> p.label+ finishPeer p+ syncNewPeer+chainMessage ChainPing = do+ $(logDebugS) "Chain" "Internal clock event"+ to <- asks (.config.timeout)+ now <- liftIO getCurrentTime+ chainSyncingPeer >>= \case+ Just ChainSync {peer = p, timestamp = t}+ | now `diffUTCTime` t > to -> do+ $(logErrorS) "Chain" $+ "Syncing peer timed out: " <> p.label+ PeerTimeout `killPeer` p+ | otherwise -> return ()+ Nothing -> syncNewPeer++withSyncLoop ::+ (MonadUnliftIO m, MonadLoggerIO m) =>+ Chain ->+ m a ->+ m a+withSyncLoop ch f =+ withAsync go $ \a ->+ link a >> f+ where+ go = forever $ do+ delay <-+ liftIO $+ randomRIO+ ( 2 * 1000 * 1000,+ 20 * 1000 * 1000+ )+ threadDelay delay+ ChainPing `send` ch.mailbox++-- | Version of the database.+dataVersion :: Word32+dataVersion = 1++-- | Initialize header database. If version is different from current, the+-- database is purged of conflicting elements first.+initChainDB :: (MonadChain m) => m ()+initChainDB = do+ db <- asks (.config.db)+ mcf <- asks (.config.cf)+ net <- asks (.config.net)+ ver <- retrieveCommon db mcf ChainDataVersionKey+ when (ver /= Just dataVersion) $ purgeChainDB >>= writeBatch db+ case mcf of+ Nothing -> insert db ChainDataVersionKey dataVersion+ Just cf -> insertCF db cf ChainDataVersionKey dataVersion+ retrieveCommon db mcf BestBlockKey >>= \b ->+ when (isNothing (b :: Maybe BlockNode)) $+ withBlockHeaders $ do+ addBlockHeader (genesisNode net)+ setBestBlockHeader (genesisNode net)++-- | Purge database of elements having keys that may conflict with those used in+-- this module.+purgeChainDB :: (MonadChain m) => m [R.BatchOp]+purgeChainDB = do+ db <- asks (.config.db)+ mcf <- asks (.config.cf)+ f db mcf $ \it -> do+ R.iterSeek it $ B.singleton 0x90+ recurse_delete it db mcf+ where+ f db Nothing = R.withIter db+ f db (Just cf) = R.withIterCF db cf+ recurse_delete it db mcf =+ R.iterKey it >>= \case+ Just k+ | B.head k == 0x90 || B.head k == 0x91 -> do+ case mcf of+ Nothing -> R.delete db k+ Just cf -> R.deleteCF db cf k+ R.iterNext it+ (R.Del k :) <$> recurse_delete it db mcf+ _ -> return []++-- | Import a bunch of continuous headers. Returns 'True' if the number of+-- headers is 2000, which means that there are possibly more headers to sync+-- from whatever peer delivered these.+importHeaders ::+ (MonadChain m) =>+ Network ->+ UTCTime ->+ [BlockHeader] ->+ m (Either PeerException Bool)+importHeaders net now hs =+ runExceptT $+ lift connect >>= \case+ Right _ -> do+ case hs of+ [] -> return ()+ _ -> do+ bb <- lift get_last+ box <- asks (.state)+ atomically . modifyTVar box $ \s ->+ s {syncing = (\x -> x {best = bb}) <$> s.syncing}+ case length hs of+ 2000 -> return False+ _ -> return True+ Left _ -> throwError PeerSentBadHeaders+ where+ timestamp = floor (utcTimeToPOSIXSeconds now)+ connect = withBlockHeaders $ connectBlocks net timestamp hs+ get_last = withBlockHeaders . getBlockHeader . headerHash $ last hs++-- | Check if best block header is in sync with the rest of the block chain by+-- comparing the best block with the current time, verifying that there are no+-- peers in the queue to be synced, and no peer is being synced at the moment.+-- This function will only return 'True' once. It should be used to decide+-- whether to notify other processes that the header chain has been synced. The+-- state of the chain will be flipped to synced when this function returns+-- 'True'.+notifySynced :: (MonadChain m) => m Bool+notifySynced =+ fmap isJust $+ runMaybeT $ do+ bb <- lift $ withBlockHeaders getBestBlockHeader+ now <- liftIO getCurrentTime+ guard $ now `diffUTCTime` block_time bb > 7200+ st <- asks (.state)+ MaybeT . atomically . runMaybeT $ do+ s <- lift $ readTVar st+ guard $ isNothing s.syncing+ guard $ null s.peers+ guard $ not s.beenInSync+ lift $ writeTVar st s {beenInSync = True}+ return ()+ where+ block_time =+ posixSecondsToUTCTime . fromIntegral . (.header.timestamp)++-- | Get next peer to sync against from the queue.+nextPeer :: (MonadChain m) => m (Maybe Peer)+nextPeer = do+ ps <- (.peers) <$> (readTVarIO =<< asks (.state))+ go ps+ where+ go [] = return Nothing+ go (p : ps) =+ setSyncingPeer p >>= \case+ True -> return (Just p)+ False -> go ps++-- | Set a syncing peer and generate a 'GetHeaders' data structure with a block+-- locator to send to that peer for syncing.+syncHeaders ::+ (MonadChain m) =>+ UTCTime ->+ BlockNode ->+ Peer ->+ m GetHeaders+syncHeaders now bb p = do+ st <- asks (.state)+ atomically $+ modifyTVar st $ \s ->+ s+ { syncing =+ Just+ ChainSync+ { peer = p,+ timestamp = now,+ best = Nothing+ },+ peers = delete p s.peers+ }+ loc <- withBlockHeaders $ blockLocator bb+ return+ GetHeaders+ { version = myVersion,+ locator = loc,+ stop = z+ }+ where+ z = "0000000000000000000000000000000000000000000000000000000000000000"++-- | Set the time of last received data to now if a syncing peer is active.+setLastReceived :: (MonadChain m) => m ()+setLastReceived = do+ now <- liftIO getCurrentTime+ st <- asks (.state)+ let f ChainSync {..} = ChainSync {timestamp = now, ..}+ atomically . modifyTVar st $ \s ->+ s {syncing = f <$> s.syncing}++-- | Add a new peer to the queue of peers to sync against.+addPeer :: (MonadChain m) => Peer -> m ()+addPeer p = do+ st <- asks (.state)+ atomically . modifyTVar st $ \s -> s {peers = nub (p : s.peers)}++-- | Get syncing peer if there is one.+getSyncingPeer :: (MonadChain m) => m (Maybe Peer)+getSyncingPeer =+ fmap (.peer) . (.syncing)+ <$> (readTVarIO =<< asks (.state))++setSyncingPeer :: (MonadChain m) => Peer -> m Bool+setSyncingPeer p =+ setBusy p >>= \case+ False -> do+ $(logDebugS) "Chain" $+ "Could not lock peer: " <> p.label+ return False+ True -> do+ $(logDebugS) "Chain" $+ "Locked peer: " <> p.label+ set_it+ return True+ where+ set_it = do+ now <- liftIO getCurrentTime+ box <- asks (.state)+ atomically $ modifyTVar box $ \s ->+ s+ { syncing =+ Just+ ChainSync+ { peer = p,+ timestamp = now,+ best = Nothing+ }+ }++-- | Remove a peer from the queue of peers to sync and unset the syncing peer if+-- it is set to the provided peer.+finishPeer :: (MonadChain m) => Peer -> m ()+finishPeer p =+ asks (.state) >>= remove_peer >>= \case+ False ->+ $(logDebugS) "Chain" $+ "Removed peer from queue: " <> p.label+ True -> do+ $(logDebugS) "Chain" $+ "Releasing syncing peer: " <> p.label+ setFree p+ where+ remove_peer st =+ atomically $+ readTVar st >>= \s -> case s.syncing of+ Just ChainSync {peer = p'}+ | p == p' -> do+ unset_syncing st+ return True+ _ -> do+ remove_from_queue st+ return False+ unset_syncing st =+ modifyTVar st $ \x ->+ x {syncing = Nothing}+ remove_from_queue st =+ modifyTVar st $ \x ->+ x {peers = delete p x.peers}++-- | Return syncing peer data.+chainSyncingPeer :: (MonadChain m) => m (Maybe ChainSync)+chainSyncingPeer =+ (.syncing) <$> (readTVarIO =<< asks (.state))++-- | Get a block header from 'Chain' process.+chainGetBlock ::+ (MonadIO m) =>+ BlockHash ->+ Chain ->+ m (Maybe BlockNode)+chainGetBlock bh ch =+ runReaderT (getBlockHeader bh) (ch.reader.config)++-- | Get best block header from chain process.+chainGetBest :: (MonadIO m) => Chain -> m BlockNode+chainGetBest ch =+ runReaderT getBestBlockHeader ch.reader.config++-- | Get ancestor of 'BlockNode' at 'BlockHeight' from chain process.+chainGetAncestor ::+ (MonadIO m) =>+ BlockHeight ->+ BlockNode ->+ Chain ->+ m (Maybe BlockNode)+chainGetAncestor h bn ch =+ runReaderT (getAncestor h bn) ch.reader.config++-- | Get parents of 'BlockNode' starting at 'BlockHeight' from chain process.+chainGetParents ::+ (MonadIO m) =>+ BlockHeight ->+ BlockNode ->+ Chain ->+ m [BlockNode]+chainGetParents height top ch =+ go [] top+ where+ go acc b+ | height >= b.height = return acc+ | otherwise = do+ m <- chainGetBlock b.header.prev ch+ case m of+ Nothing -> return acc+ Just p -> go (p : acc) p++-- | Get last common block from chain process.+chainGetSplitBlock ::+ (MonadIO m) =>+ BlockNode ->+ BlockNode ->+ Chain ->+ m BlockNode+chainGetSplitBlock l r ch =+ runReaderT (splitPoint l r) ch.reader.config++-- | Notify chain that a new peer is connected.+chainPeerConnected ::+ (MonadIO m) =>+ Peer ->+ Chain ->+ m ()+chainPeerConnected p ch =+ ChainPeerConnected p `send` ch.mailbox++-- | Notify chain that a peer has disconnected.+chainPeerDisconnected ::+ (MonadIO m) =>+ Peer ->+ Chain ->+ m ()+chainPeerDisconnected p ch =+ ChainPeerDisconnected p `send` ch.mailbox++-- | Is given 'BlockHash' in the main chain?+chainBlockMain ::+ (MonadIO m) =>+ BlockHash ->+ Chain ->+ m Bool+chainBlockMain bh ch =+ chainGetBest ch >>= \bb ->+ chainGetBlock bh ch >>= \case+ Nothing ->+ return False+ bm@(Just bn) ->+ (== bm) <$> chainGetAncestor bn.height bb ch++-- | Is chain in sync with network?+chainIsSynced :: (MonadIO m) => Chain -> m Bool+chainIsSynced ch =+ (.beenInSync) <$> readTVarIO (ch.reader.state)++-- | Peer sends a bunch of headers to the chain process.+chainHeaders ::+ (MonadIO m) =>+ Peer ->+ [BlockHeader] ->+ Chain ->+ m ()+chainHeaders p hs ch =+ ChainHeaders p hs `send` ch.mailbox
− src/Haskoin/Node/Manager.hs
@@ -1,767 +0,0 @@-{-# LANGUAGE ConstraintKinds #-}-{-# LANGUAGE FlexibleContexts #-}-{-# LANGUAGE FlexibleInstances #-}-{-# LANGUAGE LambdaCase #-}-{-# LANGUAGE MultiParamTypeClasses #-}-{-# LANGUAGE OverloadedStrings #-}-{-# LANGUAGE TemplateHaskell #-}-{-# LANGUAGE TupleSections #-}-module Haskoin.Node.Manager- ( PeerManagerConfig (..)- , PeerEvent (..)- , OnlinePeer (..)- , PeerManager- , withPeerManager- , managerBest- , managerVersion- , managerPing- , managerPong- , managerAddrs- , managerVerAck- , managerTickle- , getPeers- , getOnlinePeer- , buildVersion- , myVersion- , toSockAddr- , toHostService- ) where--import Control.Arrow-import Control.Monad (forM_, forever, guard, void, when,- (<=<))-import Control.Monad.Except (ExceptT (..), runExceptT,- throwError)-import Control.Monad.Logger (MonadLogger, MonadLoggerIO,- logDebugS, logErrorS, logInfoS,- logWarnS)-import Control.Monad.Reader (MonadReader, ReaderT (ReaderT), ask,- asks, runReaderT)-import Control.Monad.Trans (lift)-import Control.Monad.Trans.Maybe (MaybeT (..), runMaybeT)-import Data.Bits ((.&.))-import Data.Function (on)-import Data.List (find, nub, sort, dropWhileEnd, elemIndex)-import Data.Maybe (fromMaybe, isJust)-import Data.Set (Set)-import qualified Data.Set as Set-import Data.String.Conversions (cs)-import Data.Time.Clock (NominalDiffTime, UTCTime,- addUTCTime, diffUTCTime,- getCurrentTime)-import Data.Time.Clock.POSIX (utcTimeToPOSIXSeconds)-import Data.Word (Word32, Word64)-import Haskoin (BlockHeight, Message (..),- Network (..), NetworkAddress (..),- Ping (..), Pong (..),- VarString (..), Version (..),- hostToSockAddr, nodeNetwork,- sockToHostAddress)-import Haskoin.Node.Peer-import NQE (Child, Inbox, Mailbox, Publisher,- Strategy (..), Supervisor, addChild,- inboxToMailbox, newInbox,- newMailbox, publish, receive,- receiveMatch, send, sendSTM,- withSupervisor)-import Network.Socket (AddrInfo (..), AddrInfoFlag (..),- Family (..), SockAddr (..),- SocketType (..), defaultHints,- getAddrInfo)-import System.Random (randomIO, randomRIO)-import UnliftIO (Async, MonadIO, MonadUnliftIO, STM,- SomeException, TVar, atomically,- catch, liftIO, link, modifyTVar,- newTVarIO, readTVar, readTVarIO,- withAsync, withRunInIO, writeTVar)-import UnliftIO.Concurrent (threadDelay)-import Control.Applicative ((<|>))--type MonadManager m = (MonadIO m, MonadReader PeerManager m)--data PeerEvent- = PeerConnected !Peer- | PeerDisconnected !Peer- deriving Eq--data PeerManagerConfig =- PeerManagerConfig- { peerManagerMaxPeers :: !Int- , peerManagerPeers :: ![String]- , peerManagerDiscover :: !Bool- , peerManagerNetAddr :: !NetworkAddress- , peerManagerNetwork :: !Network- , peerManagerEvents :: !(Publisher PeerEvent)- , peerManagerTimeout :: !NominalDiffTime- , peerManagerMaxLife :: !NominalDiffTime- , peerManagerConnect :: !(SockAddr -> WithConnection)- , peerManagerPub :: !(Publisher (Peer, Message))- }--data PeerManager =- PeerManager- { myConfig :: !PeerManagerConfig- , mySupervisor :: !Supervisor- , myMailbox :: !(Mailbox ManagerMessage)- , myBestBlock :: !(TVar BlockHeight)- , knownPeers :: !(TVar (Set SockAddr))- , onlinePeers :: !(TVar [OnlinePeer])- }--data ManagerMessage- = Connect !SockAddr- | CheckPeer !Peer- | PeerDied !Child !(Maybe SomeException)- | ManagerBest !BlockHeight- | PeerVerAck !Peer- | PeerVersion !Peer !Version- | PeerPing !Peer !Word64- | PeerPong !Peer !Word64- | PeerAddrs !Peer ![NetworkAddress]- | PeerTickle !Peer---- | Data structure representing an online peer.-data OnlinePeer =- OnlinePeer- { onlinePeerAddress :: !SockAddr- , onlinePeerVerAck :: !Bool- , onlinePeerConnected :: !Bool- , onlinePeerVersion :: !(Maybe Version)- , onlinePeerAsync :: !(Async ())- , onlinePeerMailbox :: !Peer- , onlinePeerNonce :: !Word64- , onlinePeerPing :: !(Maybe (UTCTime, Word64))- , onlinePeerPings :: ![NominalDiffTime]- , onlinePeerConnectTime :: !UTCTime- , onlinePeerTickled :: !UTCTime- , onlinePeerDisconnect :: !UTCTime- }--instance Eq OnlinePeer where- (==) = (==) `on` f- where- f OnlinePeer {onlinePeerMailbox = p} = p--instance Ord OnlinePeer where- compare = compare `on` f- where- f OnlinePeer {onlinePeerPings = pings} = fromMaybe 60 (median pings)--withPeerManager :: (MonadUnliftIO m, MonadLoggerIO m)- => PeerManagerConfig- -> (PeerManager -> m a)- -> m a-withPeerManager cfg action = do- inbox <- newInbox- let mgr = inboxToMailbox inbox- withSupervisor (Notify (death mgr)) $ \sup -> do- bb <- newTVarIO 0- kp <- newTVarIO Set.empty- ob <- newTVarIO []- let rd = PeerManager { myConfig = cfg- , mySupervisor = sup- , myMailbox = mgr- , myBestBlock = bb- , knownPeers = kp- , onlinePeers = ob- }- go inbox `runReaderT` rd- where- death mgr (a, ex) = PeerDied a ex `sendSTM` mgr- go inbox =- withAsync (peerManager inbox) $ \a ->- withConnectLoop $- link a >> ReaderT action--peerManager :: ( MonadUnliftIO m- , MonadManager m- , MonadLoggerIO m )- => Inbox ManagerMessage- -> m ()-peerManager inb = do- $(logDebugS) "PeerManager" "Awaiting best block"- putBestBlock <=< receiveMatch inb $ \case- ManagerBest b -> Just b- _ -> Nothing- $(logDebugS) "PeerManager" "Starting peer manager actor"- forever loop- where- loop = do- m <- receive inb- managerMessage m--putBestBlock :: MonadManager m => BlockHeight -> m ()-putBestBlock bb = do- b <- asks myBestBlock- atomically $ writeTVar b bb--getBestBlock :: MonadManager m => m BlockHeight-getBestBlock =- asks myBestBlock >>= readTVarIO--getNetwork :: MonadManager m => m Network-getNetwork =- asks (peerManagerNetwork . myConfig)--loadPeers :: (MonadUnliftIO m, MonadManager m) => m ()-loadPeers = do- loadStaticPeers- loadNetSeeds--loadStaticPeers :: (MonadUnliftIO m, MonadManager m) => m ()-loadStaticPeers = do- net <- asks (peerManagerNetwork . myConfig)- xs <- asks (peerManagerPeers . myConfig)- mapM_ newPeer . concat =<< mapM (toSockAddr net) xs--loadNetSeeds :: (MonadUnliftIO m, MonadManager m) => m ()-loadNetSeeds =- asks (peerManagerDiscover . myConfig) >>= \discover ->- when discover $ do- net <- getNetwork- ss <- concat <$> mapM (toSockAddr net) (getSeeds net)- mapM_ newPeer ss--logConnectedPeers :: (MonadManager m, MonadLoggerIO m) => m ()-logConnectedPeers = do- m <- asks (peerManagerMaxPeers . myConfig)- l <- length <$> getConnectedPeers- $(logInfoS) "PeerManager" $- "Peers connected: " <> cs (show l) <> "/" <> cs (show m)--getOnlinePeers :: MonadManager m => m [OnlinePeer]-getOnlinePeers =- asks onlinePeers >>= readTVarIO--getConnectedPeers :: MonadManager m => m [OnlinePeer]-getConnectedPeers =- filter onlinePeerConnected <$> getOnlinePeers--managerEvent :: MonadManager m => PeerEvent -> m ()-managerEvent e =- publish e =<< asks (peerManagerEvents . myConfig)--managerMessage :: ( MonadUnliftIO m- , MonadManager m- , MonadLoggerIO m )- => ManagerMessage- -> m ()--managerMessage (PeerVersion p v) = do- b <- asks onlinePeers- e <- runExceptT $ do- o <- ExceptT . atomically $ setPeerVersion b p v- when (onlinePeerConnected o) $ announcePeer p- case e of- Right () -> do- $(logDebugS) "PeerManager" $- "Sending version ack to peer: " <> peerText p- MVerAck `sendMessage` p- Left x -> do- $(logErrorS) "PeerManager" $- "Version rejected for peer "- <> peerText p <> ": " <> cs (show x)- killPeer x p--managerMessage (PeerVerAck p) = do- b <- asks onlinePeers- atomically (setPeerVerAck b p) >>= \case- Just o -> do- $(logDebugS) "PeerManager" $- "Received version ack from peer: "- <> peerText p- when (onlinePeerConnected o) $- announcePeer p- Nothing -> do- $(logErrorS) "PeerManager" $- "Received verack from unknown peer: "- <> peerText p- killPeer UnknownPeer p--managerMessage (PeerAddrs p nas) = do- discover <- asks (peerManagerDiscover . myConfig)- when discover $ do- let sas = map (hostToSockAddr . naAddress) nas- forM_ (zip [(1 :: Int) ..] sas) $ \(i, a) -> do- $(logDebugS) "PeerManager" $- "Got peer address "- <> cs (show i) <> "/" <> cs (show (length sas))- <> ": " <> cs (show a)- <> " from peer " <> peerText p- newPeer a--managerMessage (PeerPong p n) = do- b <- asks onlinePeers- $(logDebugS) "PeerManager" $- "Received pong "- <> cs (show n)- <> " from: " <> peerText p- now <- liftIO getCurrentTime- atomically (gotPong b n now p)--managerMessage (PeerPing p n) = do- $(logDebugS) "PeerManager" $- "Responding to ping "- <> cs (show n)- <> " from: " <> peerText p- MPong (Pong n) `sendMessage` p--managerMessage (ManagerBest h) =- putBestBlock h--managerMessage (Connect sa) =- connectPeer sa--managerMessage (PeerDied a e) =- processPeerOffline a e--managerMessage (CheckPeer p) =- checkPeer p--managerMessage (PeerTickle p) = do- b <- asks onlinePeers- now <- liftIO getCurrentTime- atomically $- modifyPeer b p $ \o ->- o { onlinePeerTickled = now }--checkPeer :: (MonadManager m, MonadLoggerIO m) => Peer -> m ()-checkPeer p =- getBusy p >>= \case- True -> return ()- False -> do- to <- asks (peerManagerTimeout . myConfig)- b <- asks onlinePeers- atomically (findPeer b p) >>= \case- Nothing -> return ()- Just o -> do- now <- liftIO getCurrentTime- check_conn now o- when (check_tickle now to o) (check_ping o)- where- check_tickle now to o =- now `diffUTCTime` onlinePeerTickled o > to- check_conn now o =- when (now `diffUTCTime` onlinePeerDisconnect o > 0) $- killPeer PeerTooOld p- check_ping o =- case onlinePeerPing o of- Nothing ->- sendPing p- Just _ -> do- $(logWarnS) "PeerManager" $- "Peer ping timeout: " <> peerText p- killPeer PeerTimeout p--sendPing :: (MonadManager m, MonadLoggerIO m) => Peer -> m ()-sendPing p = do- b <- asks onlinePeers- atomically (findPeer b p) >>= \case- Nothing ->- $(logWarnS) "PeerManager" $- "Will not ping unknown peer: " <> peerText p- Just o- | onlinePeerConnected o -> do- n <- liftIO randomIO- now <- liftIO getCurrentTime- atomically (setPeerPing b n now p)- $(logDebugS)" PeerManager" $- "Sending ping " <> cs (show n)- <> " to: " <> peerText p- MPing (Ping n) `sendMessage` p- | otherwise -> return ()--processPeerOffline :: (MonadManager m, MonadLoggerIO m)- => Child -> Maybe SomeException -> m ()-processPeerOffline a e = do- b <- asks onlinePeers- atomically (findPeerAsync b a) >>= \case- Nothing -> log_unknown e- Just o -> do- let p = onlinePeerMailbox o- if onlinePeerConnected o- then do- log_disconnected p e- managerEvent $ PeerDisconnected p- else log_not_connect p e- atomically $ removePeer b p- logConnectedPeers- where- log_unknown Nothing =- $(logErrorS) "PeerManager"- "Disconnected unknown peer"- log_unknown (Just x) =- $(logErrorS) "PeerManager" $- "Unknown peer died: " <> cs (show x)- log_disconnected p Nothing =- $(logWarnS) "PeerManager" $- "Disconnected peer: " <> peerText p- log_disconnected p (Just x) =- $(logErrorS) "PeerManager" $- "Peer " <> peerText p <> " died: " <> cs (show x)- log_not_connect p Nothing =- $(logWarnS) "PeerManager" $- "Could not connect to peer " <> peerText p- log_not_connect p (Just x) =- $(logErrorS) "PeerManager" $- "Could not connect to peer "- <> peerText p <> ": " <> cs (show x)--announcePeer :: (MonadManager m, MonadLoggerIO m) => Peer -> m ()-announcePeer p = do- b <- asks onlinePeers- atomically (findPeer b p) >>= \case- Just OnlinePeer {onlinePeerConnected = True} -> do- $(logInfoS) "PeerManager" $- "Connected to peer " <> peerText p- managerEvent $ PeerConnected p- logConnectedPeers- Just OnlinePeer {onlinePeerConnected = False} ->- return ()- Nothing ->- $(logErrorS) "PeerManager" $- "Not announcing disconnected peer: "- <> peerText p--getNewPeer :: (MonadUnliftIO m, MonadManager m) => m (Maybe SockAddr)-getNewPeer =- runMaybeT $ lift loadPeers >> go- where- go = do- b <- asks knownPeers- ks <- readTVarIO b- guard . not $ Set.null ks- let xs = Set.toList ks- a <- liftIO $ randomRIO (0, length xs - 1)- let p = xs !! a- o <- asks onlinePeers- m <- atomically $ do- modifyTVar b $ Set.delete p- findPeerAddress o p- maybe (return p) (const go) m---connectPeer :: ( MonadUnliftIO m- , MonadManager m- , MonadLoggerIO m- )- => SockAddr- -> m ()-connectPeer sa = do- os <- asks onlinePeers- atomically (findPeerAddress os sa) >>= \case- Just _ ->- $(logErrorS) "PeerManager" $- "Attempted to connect to peer twice: " <> cs (show sa)- Nothing -> do- $(logInfoS) "PeerManager" $ "Connecting to " <> cs (show sa)- PeerManagerConfig { peerManagerNetAddr = ad- , peerManagerNetwork = net- } <- asks myConfig- sup <- asks mySupervisor- conn <- asks (peerManagerConnect . myConfig)- pub <- asks (peerManagerPub . myConfig)- nonce <- liftIO randomIO- bb <- getBestBlock- now <- liftIO getCurrentTime- let rmt = NetworkAddress (srv net) (sockToHostAddress sa)- unix = floor (utcTimeToPOSIXSeconds now)- ver = buildVersion net nonce bb ad rmt unix- text = cs (show sa)- (inbox, mailbox) <- newMailbox- let pc = PeerConfig { peerConfPub = pub- , peerConfNetwork = net- , peerConfText = text- , peerConfConnect = conn sa- }- busy <- newTVarIO False- p <- wrapPeer pc busy mailbox- a <- withRunInIO $ \io ->- sup `addChild` io (launch pc busy inbox p)- MVersion ver `sendMessage` p- b <- asks onlinePeers- max_life <- asks (peerManagerMaxLife . myConfig)- rand <- liftIO $ toRational <$> randomRIO (0.75 :: Double, 1.00)- let life = max_life * fromRational rand- let dc = life `addUTCTime` now- _ <- atomically $ newOnlinePeer b sa nonce p a now dc- return ()- where- srv net- | getSegWit net = 8- | otherwise = 0- launch pc busy inbox p =- ask >>= \mgr ->- withPeerLoop sa p mgr $ \a ->- link a >> peer pc busy inbox--withPeerLoop ::- (MonadUnliftIO m, MonadLogger m)- => SockAddr- -> Peer- -> PeerManager- -> (Async a -> m a)- -> m a-withPeerLoop _ p mgr =- withAsync . forever $ do- let x = peerManagerTimeout (myConfig mgr)- y = floor (x * 1000000)- r <- liftIO $ randomRIO (y * 3 `div` 4, y)- threadDelay r- managerCheck p mgr--withConnectLoop :: (MonadUnliftIO m, MonadManager m)- => m a- -> m a-withConnectLoop act =- withAsync go $ \a ->- link a >> act- where- go = forever $ do- l <- length <$> getOnlinePeers- x <- asks (peerManagerMaxPeers . myConfig)- when (l < x) $- getNewPeer >>= mapM_ (\sa -> ask >>= managerConnect sa)- delay <- liftIO $- randomRIO ( 100 * 1000- , 10 * 500 * 1000 )- threadDelay delay--newPeer :: (MonadIO m, MonadManager m) => SockAddr -> m ()-newPeer sa = do- b <- asks knownPeers- o <- asks onlinePeers- atomically $- findPeerAddress o sa >>= \case- Just _ -> return ()- Nothing -> modifyTVar b $ Set.insert sa--gotPong :: TVar [OnlinePeer] -> Word64 -> UTCTime -> Peer -> STM ()-gotPong b nonce now p = void . runMaybeT $ do- o <- MaybeT $ findPeer b p- (time, old_nonce) <- MaybeT . return $ onlinePeerPing o- guard $ nonce == old_nonce- let diff = now `diffUTCTime` time- lift $- insertPeer- b- o { onlinePeerPing = Nothing- , onlinePeerPings = sort $ take 11 $ diff : onlinePeerPings o- }--setPeerPing :: TVar [OnlinePeer] -> Word64 -> UTCTime -> Peer -> STM ()-setPeerPing b nonce now p =- modifyPeer b p $ \o -> o {onlinePeerPing = Just (now, nonce)}--setPeerVersion ::- TVar [OnlinePeer]- -> Peer- -> Version- -> STM (Either PeerException OnlinePeer)-setPeerVersion b p v = runExceptT $ do- when (services v .&. nodeNetwork == 0) $- throwError NotNetworkPeer- ops <- lift $ readTVar b- when (any ((verNonce v ==) . onlinePeerNonce) ops) $- throwError PeerIsMyself- lift (findPeer b p) >>= \case- Nothing -> throwError UnknownPeer- Just o -> do- let n = o { onlinePeerVersion = Just v- , onlinePeerConnected = onlinePeerVerAck o }- lift $ insertPeer b n- return n--setPeerVerAck :: TVar [OnlinePeer] -> Peer -> STM (Maybe OnlinePeer)-setPeerVerAck b p = runMaybeT $ do- o <- MaybeT $ findPeer b p- let n = o { onlinePeerVerAck = True- , onlinePeerConnected = isJust (onlinePeerVersion o) }- lift $ insertPeer b n- return n--newOnlinePeer ::- TVar [OnlinePeer]- -> SockAddr- -> Word64- -> Peer- -> Async ()- -> UTCTime- -> UTCTime- -> STM OnlinePeer-newOnlinePeer box addr nonce p peer_async connect_time dc = do- let op = OnlinePeer- { onlinePeerAddress = addr- , onlinePeerVerAck = False- , onlinePeerConnected = False- , onlinePeerVersion = Nothing- , onlinePeerAsync = peer_async- , onlinePeerMailbox = p- , onlinePeerNonce = nonce- , onlinePeerPings = []- , onlinePeerPing = Nothing- , onlinePeerConnectTime = connect_time- , onlinePeerTickled = connect_time- , onlinePeerDisconnect = dc- }- insertPeer box op- return op--findPeer :: TVar [OnlinePeer] -> Peer -> STM (Maybe OnlinePeer)-findPeer b p =- find ((== p) . onlinePeerMailbox)- <$> readTVar b--insertPeer :: TVar [OnlinePeer] -> OnlinePeer -> STM ()-insertPeer b o =- modifyTVar b $ \x -> sort . nub $ o : x--modifyPeer :: TVar [OnlinePeer]- -> Peer- -> (OnlinePeer -> OnlinePeer)- -> STM ()-modifyPeer b p f =- findPeer b p >>= \case- Nothing -> return ()- Just o -> insertPeer b $ f o--removePeer :: TVar [OnlinePeer] -> Peer -> STM ()-removePeer b p =- modifyTVar b $- filter ((/= p) . onlinePeerMailbox)--findPeerAsync :: TVar [OnlinePeer]- -> Async ()- -> STM (Maybe OnlinePeer)-findPeerAsync b a =- find ((== a) . onlinePeerAsync)- <$> readTVar b--findPeerAddress :: TVar [OnlinePeer]- -> SockAddr- -> STM (Maybe OnlinePeer)-findPeerAddress b a =- find ((== a) . onlinePeerAddress)- <$> readTVar b--getPeers :: MonadIO m => PeerManager -> m [OnlinePeer]-getPeers = runReaderT getConnectedPeers--getOnlinePeer :: MonadIO m- => Peer- -> PeerManager- -> m (Maybe OnlinePeer)-getOnlinePeer p =- runReaderT $ asks onlinePeers >>= atomically . (`findPeer` p)--managerCheck :: MonadIO m => Peer -> PeerManager -> m ()-managerCheck p mgr =- CheckPeer p `send` myMailbox mgr--managerConnect :: MonadIO m => SockAddr -> PeerManager -> m ()-managerConnect sa mgr =- Connect sa `send` myMailbox mgr--managerBest :: MonadIO m => BlockHeight -> PeerManager -> m ()-managerBest bh mgr =- ManagerBest bh `send` myMailbox mgr--managerVerAck :: MonadIO m => Peer -> PeerManager -> m ()-managerVerAck p mgr =- PeerVerAck p `send` myMailbox mgr--managerVersion :: MonadIO m- => Peer -> Version -> PeerManager -> m ()-managerVersion p ver mgr =- PeerVersion p ver `send` myMailbox mgr--managerPing :: MonadIO m- => Peer -> Word64 -> PeerManager -> m ()-managerPing p nonce mgr =- PeerPing p nonce `send` myMailbox mgr--managerPong :: MonadIO m- => Peer -> Word64 -> PeerManager -> m ()-managerPong p nonce mgr =- PeerPong p nonce `send` myMailbox mgr--managerAddrs :: MonadIO m- => Peer -> [NetworkAddress] -> PeerManager -> m ()-managerAddrs p addrs mgr =- PeerAddrs p addrs `send` myMailbox mgr--managerTickle :: MonadIO m- => Peer -> PeerManager -> m ()-managerTickle p mgr =- PeerTickle p `send` myMailbox mgr--toHostService :: String -> (Maybe String, Maybe String)-toHostService str =- let host = case m6 of- Just (x, _) -> Just x- Nothing -> case takeWhile (/= ':') str of- [] -> Nothing- xs -> Just xs- srv = case m6 of- Just (_, y) -> s y- Nothing -> s str- s xs =- case dropWhile (/= ':') xs of- [] -> Nothing- _ : ys -> Just ys- m6 = case str of- (x : xs)- | x == '[' -> do- i <- elemIndex ']' xs- return $ second tail $ splitAt i xs- | x == ':' -> do- return (str, "")- _ -> Nothing- in (host, srv)--toSockAddr :: MonadUnliftIO m => Network -> String -> m [SockAddr]-toSockAddr net str = - go `catch` e- where- go = fmap (map addrAddress) $ liftIO $ getAddrInfo Nothing host srv- (host, srv) = - second (<|> Just (show (getDefaultPort net))) $- toHostService str- e :: Monad m => SomeException -> m [SockAddr]- e _ = return []--median :: (Ord a, Fractional a) => [a] -> Maybe a-median ls- | null ls =- Nothing- | even (length ls) =- Just . (/ 2) . sum . take 2 $- drop (length ls `div` 2 - 1) ls'- | otherwise =- Just (ls' !! (length ls `div` 2))- where- ls' = sort ls--buildVersion- :: Network- -> Word64- -> BlockHeight- -> NetworkAddress- -> NetworkAddress- -> Word64- -> Version-buildVersion net nonce height loc rmt time =- Version- { version = myVersion- , services = naServices loc- , timestamp = time- , addrRecv = rmt- , addrSend = loc- , verNonce = nonce- , userAgent = VarString (getHaskoinUserAgent net)- , startHeight = height- , relay = True- }--myVersion :: Word32-myVersion = 70012
src/Haskoin/Node/Peer.hs view
@@ -1,328 +1,405 @@-{-# LANGUAGE ConstraintKinds #-}-{-# LANGUAGE FlexibleContexts #-}-{-# LANGUAGE LambdaCase #-}+{-# LANGUAGE ConstraintKinds #-}+{-# LANGUAGE DuplicateRecordFields #-}+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE ImportQualifiedPost #-}+{-# LANGUAGE LambdaCase #-} {-# LANGUAGE MultiParamTypeClasses #-}-{-# LANGUAGE OverloadedStrings #-}-{-# LANGUAGE RecordWildCards #-}-{-# LANGUAGE TypeFamilies #-}+{-# LANGUAGE NamedFieldPuns #-}+{-# LANGUAGE OverloadedRecordDot #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE TypeFamilies #-}+{-# LANGUAGE NoFieldSelectors #-}+ module Haskoin.Node.Peer- ( PeerConfig(..)- , Conduits(..)- , PeerException(..)- , WithConnection- , Peer- , peer- , wrapPeer- , peerPublisher- , peerText- , sendMessage- , killPeer- , getBlocks- , getTxs- , getData- , pingPeer- , getBusy- , setBusy- , setFree- ) where+ ( PeerConfig (..),+ PeerEvent (..),+ Conduits (..),+ PeerException (..),+ WithConnection,+ Peer (..),+ peer,+ wrapPeer,+ sendMessage,+ killPeer,+ getBlocks,+ getTxs,+ getData,+ pingPeer,+ getBusy,+ setBusy,+ setFree,+ )+where -import Conduit (ConduitT, Void, awaitForever, foldC,- mapM_C, runConduit, takeCE,- transPipe, yield, (.|))-import Control.Monad (forever, join, when)-import Control.Monad.Logger (MonadLoggerIO, logErrorS)-import Control.Monad.Trans.Maybe (MaybeT (MaybeT), runMaybeT)-import Data.ByteString (ByteString)-import qualified Data.ByteString as B-import Data.Function (on)-import Data.List (union)-import Data.Maybe (isJust)-import Data.Serialize (decode, runGet, runPut)-import Data.String.Conversions (cs)-import Data.Text (Text)-import Data.Word (Word32)-import Haskoin (Block (..), BlockHash (..),- GetData (..), InvType (..),- InvVector (..), Message (..),- MessageCommand (..),- MessageHeader (..), Network (..),- NotFound (..), Ping (..), Pong (..),- Tx, TxHash (..), commandToString,- encodeHex, getMessage, headerHash,- putMessage, txHash)-import NQE (Inbox, Mailbox, Publisher,- inboxToMailbox, publish, receive,- receiveMatchS, send,- withSubscription)-import System.Random (randomIO)-import UnliftIO (Exception, MonadIO, MonadUnliftIO,- TVar, atomically, liftIO, link,- readTVar, readTVarIO, throwIO,- timeout, withAsync, withRunInIO,- writeTVar)+import Conduit+ ( ConduitT,+ Void,+ awaitForever,+ foldC,+ mapM_C,+ runConduit,+ takeCE,+ transPipe,+ yield,+ (.|),+ )+import Control.Monad (forever, join, unless, when)+import Control.Monad.Logger (MonadLoggerIO, logDebugS, logErrorS, logInfoS)+import Control.Monad.Trans.Maybe (MaybeT (MaybeT), runMaybeT)+import Data.Bool (bool)+import Data.ByteString (ByteString)+import Data.ByteString qualified as B+import Data.Function (on)+import Data.List (union)+import Data.Maybe (isJust)+import Data.Serialize (decode, runGet, runPut)+import Data.String.Conversions (cs)+import Data.Text (Text)+import Data.Word (Word32)+import Haskoin+ ( Block (..),+ BlockHash (..),+ GetData (..),+ InvType (..),+ InvVector (..),+ Message (..),+ MessageCommand (..),+ MessageHeader (..),+ Network (..),+ NotFound (..),+ Ping (..),+ Pong (..),+ Tx,+ TxHash (..),+ commandToString,+ encodeHex,+ getMessage,+ headerHash,+ putMessage,+ txHash,+ )+import NQE+ ( Inbox,+ Mailbox,+ Publisher,+ inboxToMailbox,+ publish,+ receive,+ receiveMatchS,+ send,+ withSubscription,+ )+import System.Random (randomIO)+import UnliftIO+ ( Exception,+ MonadIO,+ MonadUnliftIO,+ TVar,+ atomically,+ liftIO,+ link,+ readTVar,+ readTVarIO,+ throwIO,+ timeout,+ withAsync,+ withRunInIO,+ writeTVar,+ ) -data Conduits =- Conduits- { inboundConduit :: ConduitT () ByteString IO ()- , outboundConduit :: ConduitT ByteString Void IO ()- }+data Conduits = Conduits+ { inboundConduit :: ConduitT () ByteString IO (),+ outboundConduit :: ConduitT ByteString Void IO ()+ } type WithConnection = (Conduits -> IO ()) -> IO () data PeerConfig = PeerConfig- { peerConfPub :: !(Publisher (Peer, Message))- , peerConfNetwork :: !Network- , peerConfText :: !Text- , peerConfConnect :: !WithConnection- }+ { pub :: !(Publisher PeerEvent),+ net :: !Network,+ label :: !Text,+ connect :: !WithConnection+ } +data PeerEvent+ = PeerConnected !Peer+ | PeerDisconnected !Peer+ | PeerMessage !Peer !Message+ deriving (Eq)+ data PeerException- = PeerMisbehaving !String- | DuplicateVersion- | DecodeHeaderError- | CannotDecodePayload !MessageCommand- | PeerIsMyself- | PayloadTooLarge !Word32- | PeerAddressInvalid- | PeerSentBadHeaders- | NotNetworkPeer- | PeerNoSegWit- | PeerTimeout- | UnknownPeer- | PeerTooOld- | EmptyHeader- deriving Eq+ = PeerMisbehaving !String+ | DuplicateVersion+ | DecodeHeaderError+ | CannotDecodePayload !MessageCommand+ | PeerIsMyself+ | PayloadTooLarge !Word32+ | PeerAddressInvalid+ | PeerSentBadHeaders+ | NotNetworkPeer+ | PeerNoSegWit+ | PeerTimeout+ | UnknownPeer+ | PeerTooOld+ | EmptyHeader+ deriving (Eq) instance Show PeerException where- show (PeerMisbehaving s) = "Peer misbehaving: " <> s- show DuplicateVersion = "Duplicate version"- show DecodeHeaderError = "Error decoding header"- show (CannotDecodePayload c) = "Cannot decode payload: " <>- cs (commandToString c)- show PeerIsMyself = "Peer is myself"- show (PayloadTooLarge s) = "Payload too large: " <> show s- show PeerAddressInvalid = "Peer address invalid"- show PeerSentBadHeaders = "Peer sent bad headers"- show NotNetworkPeer = "Not network peer"- show PeerNoSegWit = "Segwit not supported by peer"- show PeerTimeout = "Peer timed out"- show UnknownPeer = "Unknown peer"- show PeerTooOld = "Peer too old"- show EmptyHeader = "Empty header"+ show (PeerMisbehaving s) = "Peer misbehaving: " <> s+ show DuplicateVersion = "Duplicate version"+ show DecodeHeaderError = "Error decoding header"+ show (CannotDecodePayload c) =+ "Cannot decode payload: "+ <> cs (commandToString c)+ show PeerIsMyself = "Peer is myself"+ show (PayloadTooLarge s) = "Payload too large: " <> show s+ show PeerAddressInvalid = "Peer address invalid"+ show PeerSentBadHeaders = "Peer sent bad headers"+ show NotNetworkPeer = "Not network peer"+ show PeerNoSegWit = "Segwit not supported by peer"+ show PeerTimeout = "Peer timed out"+ show UnknownPeer = "Unknown peer"+ show PeerTooOld = "Peer too old"+ show EmptyHeader = "Empty header" instance Exception PeerException -- | Mailbox for a peer.-data Peer = Peer { peerMailbox :: !(Mailbox PeerMessage)- , peerPublisher :: !(Publisher (Peer, Message))- , peerText :: !Text- , peerBusy :: !(TVar Bool)- }+data Peer = Peer+ { mailbox :: !(Mailbox PeerMessage),+ pub :: !(Publisher PeerEvent),+ label :: !Text,+ busy :: !(TVar Bool)+ } instance Eq Peer where- (==) = (==) `on` peerMailbox+ (==) = (==) `on` (.mailbox) instance Show Peer where- show = cs . peerText+ show = cs . (.label) -- | Incoming messages that a peer accepts. data PeerMessage- = KillPeer !PeerException- | SendMessage !Message+ = KillPeer !PeerException+ | SendMessage !Message -wrapPeer :: MonadIO m- => PeerConfig- -> TVar Bool- -> Mailbox PeerMessage- -> m Peer+wrapPeer ::+ (MonadIO m) =>+ PeerConfig ->+ TVar Bool ->+ Mailbox PeerMessage ->+ m Peer wrapPeer cfg busy mbox =- return Peer { peerMailbox = mbox- , peerPublisher = peerConfPub cfg- , peerText = peerConfText cfg- , peerBusy = busy- }+ return+ Peer+ { mailbox = mbox,+ pub = cfg.pub,+ label = cfg.label,+ busy = busy+ } -- | Run peer process in current thread.-peer :: (MonadUnliftIO m, MonadLoggerIO m)- => PeerConfig- -> TVar Bool- -> Inbox PeerMessage- -> m ()-peer cfg@PeerConfig{..} busy inbox = do- p <- wrapPeer cfg busy (inboxToMailbox inbox)- withRunInIO $ \restore -> do- peerConfConnect (peer_session p)+peer ::+ (MonadUnliftIO m, MonadLoggerIO m) =>+ PeerConfig ->+ TVar Bool ->+ Inbox PeerMessage ->+ m ()+peer cfg@PeerConfig {..} busy inbox = do+ p <- wrapPeer cfg busy (inboxToMailbox inbox)+ withRunInIO $ \restore -> do+ connect (restore . peer_session p) where- go = forever $ receive inbox >>= dispatchMessage cfg+ go = forever $ do+ $(logDebugS) "Peer" $ label <> " awaiting event..."+ msg <- receive inbox+ dispatchMessage cfg msg peer_session p ad = do- let ins = transPipe liftIO (inboundConduit ad)- ons = transPipe liftIO (outboundConduit ad)- src = runConduit $- ins- .| inPeerConduit peerConfNetwork peerConfText+ let ins = transPipe liftIO ad.inboundConduit+ ons = transPipe liftIO ad.outboundConduit+ src =+ runConduit $+ ins+ .| inPeerConduit net cfg label .| mapM_C (send_msg p)- snk = outPeerConduit peerConfNetwork .| ons- withAsync src $ \as -> do- link as- runConduit (go .| snk)- send_msg p msg = publish (p, msg) peerConfPub+ snk = outPeerConduit net .| ons+ withAsync src $ \as -> do+ link as+ runConduit (go .| snk)+ send_msg p msg = publish (PeerMessage p msg) pub -- | Internal function to dispatch peer messages.-dispatchMessage :: MonadIO m- => PeerConfig- -> PeerMessage- -> ConduitT i Message m ()-dispatchMessage _ (SendMessage msg) = yield msg-dispatchMessage _ (KillPeer e) = throwIO e+dispatchMessage ::+ (MonadLoggerIO m) =>+ PeerConfig ->+ PeerMessage ->+ ConduitT i Message m ()+dispatchMessage PeerConfig {label} (SendMessage msg) = do+ $(logDebugS) "Peer" $ label <> " sending: " <> cs (show msg)+ yield msg+dispatchMessage PeerConfig {label} (KillPeer e) = do+ $(logInfoS) "Peer" $ label <> " killing with error: " <> cs (show e)+ throwIO e -- | Internal conduit to parse messages coming from peer.-inPeerConduit :: MonadIO m- => Network- -> Text- -> ConduitT ByteString Message m ()-inPeerConduit net a =- forever $ do- x <- takeCE 24 .| foldC- when (B.null x) $ do- throwIO EmptyHeader- case decode x of- Left e -> do- throwIO DecodeHeaderError- Right (MessageHeader _ cmd len _) -> do- when (len > 32 * 2 ^ (20 :: Int)) $ do- throwIO $ PayloadTooLarge len- y <- takeCE (fromIntegral len) .| foldC- case runGet (getMessage net) $ x `B.append` y of- Left e -> do- throwIO (CannotDecodePayload cmd)- Right msg -> yield msg+inPeerConduit ::+ (MonadLoggerIO m) =>+ Network ->+ PeerConfig ->+ Text ->+ ConduitT ByteString Message m ()+inPeerConduit net PeerConfig {label} a =+ forever $ do+ $(logDebugS) "Peer" $ label <> " awaiting network message..."+ x <- takeCE 24 .| foldC+ when (B.null x) $ do+ $(logErrorS) "Peer" $ label <> " empty header"+ throwIO EmptyHeader+ case decode x of+ Left e -> do+ $(logErrorS) "Peer" $ label <> " error decoding header"+ throwIO DecodeHeaderError+ Right (MessageHeader _ cmd len _) -> do+ $(logDebugS) "Peer" $ label <> " received: " <> cs (show cmd)+ when (len > 32 * 2 ^ (20 :: Int)) $ do+ $(logErrorS) "Peer" $ label <> " payload too large: " <> cs (show len)+ throwIO $ PayloadTooLarge len+ y <- takeCE (fromIntegral len) .| foldC+ case runGet (getMessage net) $ x `B.append` y of+ Left e -> do+ $(logErrorS) "Peer" $+ label+ <> " could not decode payload for cmd: "+ <> cs (show cmd)+ throwIO (CannotDecodePayload cmd)+ Right msg -> do+ $(logDebugS) "Peer" $ label <> " forwarding: " <> cs (show msg)+ yield msg -- | Outgoing peer conduit to serialize and send messages.-outPeerConduit :: Monad m => Network -> ConduitT Message ByteString m ()+outPeerConduit :: (Monad m) => Network -> ConduitT Message ByteString m () outPeerConduit net = awaitForever $ yield . runPut . putMessage net -- | Kill a peer with the provided exception.-killPeer :: MonadIO m => PeerException -> Peer -> m ()-killPeer e p = KillPeer e `send` peerMailbox p+killPeer :: (MonadIO m) => PeerException -> Peer -> m ()+killPeer e p = KillPeer e `send` p.mailbox -- | Send a network message to peer.-sendMessage :: MonadIO m => Message -> Peer -> m ()-sendMessage msg p = SendMessage msg `send` peerMailbox p+sendMessage :: (MonadIO m) => Message -> Peer -> m ()+sendMessage msg p = SendMessage msg `send` p.mailbox -getBusy :: MonadIO m => Peer -> m Bool-getBusy p = readTVarIO (peerBusy p)+getBusy :: (MonadIO m) => Peer -> m Bool+getBusy p = readTVarIO p.busy -setBusy :: MonadIO m => Peer -> m Bool+setBusy :: (MonadIO m) => Peer -> m Bool setBusy p =- atomically $- readTVar (peerBusy p) >>= \case- True -> return False- False -> writeTVar (peerBusy p) True >>- return True+ atomically $ do+ b <- readTVar p.busy+ unless b $ writeTVar p.busy True+ return $ not b -setFree :: MonadIO m => Peer -> m ()-setFree p = atomically $ writeTVar (peerBusy p) False+setFree :: (MonadIO m) => Peer -> m ()+setFree p = atomically $ writeTVar p.busy False -- | Request full blocks from peer. Will return 'Nothing' if the list of blocks -- returned by the peer is incomplete, comes out of order, or a timeout is -- reached.-getBlocks :: MonadUnliftIO m- => Network- -> Int- -> Peer- -> [BlockHash]- -> m (Maybe [Block])+getBlocks ::+ (MonadUnliftIO m) =>+ Network ->+ Int ->+ Peer ->+ [BlockHash] ->+ m (Maybe [Block]) getBlocks net time p bhs =- runMaybeT $ mapM f =<< MaybeT (getData time p (GetData ivs))+ runMaybeT $ mapM f =<< MaybeT (getData time p (GetData ivs)) where f (Right b) = return b- f (Left _) = MaybeT $ return Nothing+ f (Left _) = MaybeT $ return Nothing c- | getSegWit net = InvWitnessBlock- | otherwise = InvBlock- ivs = map (InvVector c . getBlockHash) bhs+ | net.segWit = InvWitnessBlock+ | otherwise = InvBlock+ ivs = map (InvVector c . (.get)) bhs -- | Request transactions from peer. Will return 'Nothing' if the list of -- transactions returned by the peer is incomplete, comes out of order, or a -- timeout is reached.-getTxs :: MonadUnliftIO m- => Network- -> Int- -> Peer- -> [TxHash]- -> m (Maybe [Tx])+getTxs ::+ (MonadUnliftIO m) =>+ Network ->+ Int ->+ Peer ->+ [TxHash] ->+ m (Maybe [Tx]) getTxs net time p ths =- runMaybeT $ mapM f =<< MaybeT (getData time p (GetData ivs))+ runMaybeT $ mapM f =<< MaybeT (getData time p (GetData ivs)) where f (Right _) = MaybeT $ return Nothing- f (Left t) = return t+ f (Left t) = return t c- | getSegWit net = InvWitnessTx- | otherwise = InvTx- ivs = map (InvVector c . getTxHash) ths+ | net.segWit = InvWitnessTx+ | otherwise = InvTx+ ivs = map (InvVector c . (.get)) ths -- | Request transactions and/or blocks from peer. Return 'Nothing' if any -- single inventory fails to be retrieved, if they come out of order, or if -- timeout is reached.-getData :: MonadUnliftIO m- => Int- -> Peer- -> GetData- -> m (Maybe [Either Tx Block])+getData ::+ (MonadUnliftIO m) => Int -> Peer -> GetData -> m (Maybe [Either Tx Block]) getData seconds p gd@(GetData ivs) =- withSubscription (peerPublisher p) $ \inb -> do- MGetData gd `sendMessage` p+ withSubscription p.pub $ \inb -> do r <- liftIO randomIO+ MGetData gd `sendMessage` p MPing (Ping r) `sendMessage` p- fmap join . timeout (seconds * 1000 * 1000) .- runMaybeT $ get_thing inb r [] ivs+ fmap join+ . timeout (seconds * 1000 * 1000)+ . runMaybeT+ $ get_thing inb r [] ivs where get_thing _inb _r acc [] =- return $ reverse acc+ return $ reverse acc get_thing inb r acc hss@(InvVector t h : hs) =- filterReceive p inb >>= \case- MTx tx- | is_tx t && getTxHash (txHash tx) == h ->- get_thing inb r (Left tx : acc) hs- MBlock b@(Block bh _)- | is_block t && getBlockHash (headerHash bh) == h ->- get_thing inb r (Right b : acc) hs- MNotFound (NotFound nvs)- | not (null (nvs `union` hs)) ->- MaybeT $ return Nothing- MPong (Pong r')- | r == r' ->- MaybeT $ return Nothing- _- | null acc ->- get_thing inb r acc hss- | otherwise ->- MaybeT $ return Nothing+ filterReceive p inb >>= \case+ MTx tx+ | is_tx t && (txHash tx).get == h ->+ get_thing inb r (Left tx : acc) hs+ MBlock b@(Block bh _)+ | is_block t && (headerHash bh).get == h ->+ get_thing inb r (Right b : acc) hs+ MNotFound (NotFound nvs)+ | not (null (nvs `union` hs)) ->+ MaybeT $ return Nothing+ MPong (Pong r')+ | r == r' ->+ MaybeT $ return Nothing+ _+ | null acc ->+ get_thing inb r acc hss+ | otherwise ->+ MaybeT $ return Nothing is_tx InvWitnessTx = True- is_tx InvTx = True- is_tx _ = False+ is_tx InvTx = True+ is_tx _ = False is_block InvWitnessBlock = True- is_block InvBlock = True- is_block _ = False+ is_block InvBlock = True+ is_block _ = False -- | Ping a peer and await response. Return 'False' if response not received -- before timeout.-pingPeer :: MonadUnliftIO m => Int -> Peer -> m Bool+pingPeer :: (MonadUnliftIO m) => Int -> Peer -> m Bool pingPeer time p =- fmap isJust . withSubscription (peerPublisher p) $ \sub -> do- r <- liftIO randomIO- MPing (Ping r) `sendMessage` p- receiveMatchS time sub $ \case- (p', MPong (Pong r'))- | p == p' && r == r' -> Just ()- _ -> Nothing---- | Peer string for logging-peerLog :: Text -> Text-peerLog = mappend "Peer|"+ fmap isJust . withSubscription p.pub $ \sub -> do+ r <- liftIO randomIO+ MPing (Ping r) `sendMessage` p+ receiveMatchS time sub $ \case+ PeerMessage p' (MPong (Pong r'))+ | p == p' && r == r' -> Just ()+ _ -> Nothing -filterReceive :: MonadIO m => Peer -> Inbox (Peer, Message) -> m Message+filterReceive :: (MonadIO m) => Peer -> Inbox PeerEvent -> m Message filterReceive p inb =- receive inb >>= \case- (p', msg) | p == p' -> return msg- _ -> filterReceive p inb+ receive inb >>= \case+ PeerMessage p' msg | p == p' -> return msg+ _ -> filterReceive p inb
+ src/Haskoin/Node/PeerMgr.hs view
@@ -0,0 +1,867 @@+{-# LANGUAGE ConstraintKinds #-}+{-# LANGUAGE DuplicateRecordFields #-}+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE FlexibleInstances #-}+{-# LANGUAGE LambdaCase #-}+{-# LANGUAGE MultiParamTypeClasses #-}+{-# LANGUAGE OverloadedRecordDot #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE TupleSections #-}+{-# LANGUAGE NoFieldSelectors #-}++module Haskoin.Node.PeerMgr+ ( PeerMgrConfig (..),+ PeerEvent (..),+ OnlinePeer (..),+ PeerMgr,+ withPeerMgr,+ peerMgrBest,+ peerMgrVersion,+ peerMgrPing,+ peerMgrPong,+ peerMgrAddrs,+ peerMgrVerAck,+ peerMgrTickle,+ getPeers,+ getOnlinePeer,+ buildVersion,+ myVersion,+ toSockAddr,+ toHostService,+ )+where++import Control.Applicative ((<|>))+import Control.Arrow+import Control.Monad+ ( forM_,+ forever,+ guard,+ unless,+ void,+ when,+ (<=<),+ )+import Control.Monad.Except+ ( ExceptT (..),+ runExceptT,+ throwError,+ )+import Control.Monad.Logger+ ( MonadLogger,+ MonadLoggerIO,+ logDebugS,+ logErrorS,+ logInfoS,+ logWarnS,+ )+import Control.Monad.Reader+ ( MonadReader,+ ReaderT (ReaderT),+ ask,+ asks,+ runReaderT,+ )+import Control.Monad.Trans (lift)+import Control.Monad.Trans.Maybe (MaybeT (..), runMaybeT)+import Data.Bits ((.&.))+import Data.Function (on)+import Data.List (dropWhileEnd, elemIndex, find, nub, sort)+import Data.Maybe (fromMaybe, isJust)+import Data.Set (Set)+import qualified Data.Set as Set+import Data.String.Conversions (cs)+import Data.Time.Clock+ ( NominalDiffTime,+ UTCTime,+ addUTCTime,+ diffUTCTime,+ getCurrentTime,+ )+import Data.Time.Clock.POSIX (utcTimeToPOSIXSeconds)+import Data.Word (Word32, Word64)+import Haskoin+ ( BlockHeight,+ Message (..),+ Network (..),+ NetworkAddress (..),+ Ping (..),+ Pong (..),+ VarString (..),+ Version (..),+ hostToSockAddr,+ nodeNetwork,+ sockToHostAddress,+ )+import Haskoin.Node.Peer+import NQE+ ( Child,+ Inbox,+ Mailbox,+ Publisher,+ Strategy (..),+ Supervisor,+ addChild,+ inboxToMailbox,+ newInbox,+ newMailbox,+ publish,+ receive,+ receiveMatch,+ send,+ sendSTM,+ withSupervisor,+ )+import Network.Socket+ ( AddrInfo (..),+ AddrInfoFlag (..),+ Family (..),+ SockAddr (..),+ SocketType (..),+ defaultHints,+ getAddrInfo,+ )+import System.Random (randomIO, randomRIO)+import UnliftIO+ ( Async,+ MonadIO,+ MonadUnliftIO,+ STM,+ SomeException,+ TVar,+ atomically,+ catch,+ liftIO,+ link,+ modifyTVar,+ newTVarIO,+ readTVar,+ readTVarIO,+ withAsync,+ withRunInIO,+ writeTVar,+ )+import UnliftIO.Concurrent (threadDelay)++type MonadManager m = (MonadIO m, MonadReader PeerMgr m)++data PeerMgrConfig = PeerMgrConfig+ { maxPeers :: !Int,+ peers :: ![String],+ discover :: !Bool,+ address :: !NetworkAddress,+ net :: !Network,+ pub :: !(Publisher PeerEvent),+ timeout :: !NominalDiffTime,+ maxPeerLife :: !NominalDiffTime,+ connect :: !(SockAddr -> WithConnection)+ }++data PeerMgr = PeerMgr+ { config :: !PeerMgrConfig,+ supervisor :: !Supervisor,+ mailbox :: !(Mailbox PeerMgrMessage),+ best :: !(TVar BlockHeight),+ addresses :: !(TVar (Set SockAddr)),+ peers :: !(TVar [OnlinePeer])+ }++data PeerMgrMessage+ = Connect !SockAddr+ | CheckPeer !Peer+ | PeerDied !Child !(Maybe SomeException)+ | ManagerBest !BlockHeight+ | PeerVerAck !Peer+ | PeerVersion !Peer !Version+ | PeerPing !Peer !Word64+ | PeerPong !Peer !Word64+ | PeerAddrs !Peer ![NetworkAddress]+ | PeerTickle !Peer++-- | Data structure representing an online peer.+data OnlinePeer = OnlinePeer+ { address :: !SockAddr,+ verack :: !Bool,+ online :: !Bool,+ version :: !(Maybe Version),+ async :: !(Async ()),+ mailbox :: !Peer,+ nonce :: !Word64,+ ping :: !(Maybe (UTCTime, Word64)),+ pings :: ![NominalDiffTime],+ connected :: !UTCTime,+ tickled :: !UTCTime+ }++instance Eq OnlinePeer where+ (==) = (==) `on` f+ where+ f OnlinePeer {mailbox = p} = p++instance Ord OnlinePeer where+ compare = compare `on` f+ where+ f OnlinePeer {pings = pings} = fromMaybe 60 (median pings)++withPeerMgr ::+ (MonadUnliftIO m, MonadLoggerIO m) =>+ PeerMgrConfig ->+ (PeerMgr -> m a) ->+ m a+withPeerMgr cfg action = do+ inbox <- newInbox+ let mgr = inboxToMailbox inbox+ withSupervisor (Notify (death mgr)) $ \sup -> do+ bb <- newTVarIO 0+ kp <- newTVarIO Set.empty+ ob <- newTVarIO []+ runReaderT+ (go inbox)+ PeerMgr+ { config = cfg,+ supervisor = sup,+ mailbox = mgr,+ best = bb,+ addresses = kp,+ peers = ob+ }+ where+ death mgr (a, ex) = PeerDied a ex `sendSTM` mgr+ go inbox =+ withAsync (peerManager inbox) $ \a ->+ withConnectLoop $+ link a >> ReaderT action++peerManager ::+ ( MonadUnliftIO m,+ MonadManager m,+ MonadLoggerIO m+ ) =>+ Inbox PeerMgrMessage ->+ m ()+peerManager inb = do+ $(logDebugS) "PeerMgr" "Awaiting best block"+ putBestBlock <=< receiveMatch inb $ \case+ ManagerBest b -> Just b+ _ -> Nothing+ $(logDebugS) "PeerMgr" "Starting peer manager actor"+ forever $ do+ $(logDebugS) "PeerMgr" "Awaiting event..."+ dispatch =<< receive inb++putBestBlock :: (MonadManager m) => BlockHeight -> m ()+putBestBlock bb = do+ b <- asks (.best)+ atomically $ writeTVar b bb++getBestBlock :: (MonadManager m) => m BlockHeight+getBestBlock =+ asks (.best) >>= readTVarIO++getNetwork :: (MonadManager m) => m Network+getNetwork =+ asks (.config.net)++loadPeers :: (MonadUnliftIO m, MonadManager m) => m ()+loadPeers = do+ loadStaticPeers+ loadNetSeeds++loadStaticPeers :: (MonadUnliftIO m, MonadManager m) => m ()+loadStaticPeers = do+ net <- asks (.config.net)+ xs <- asks (.config.peers)+ mapM_ newPeer . concat =<< mapM (toSockAddr net) xs++loadNetSeeds :: (MonadUnliftIO m, MonadManager m) => m ()+loadNetSeeds =+ asks (.config.discover) >>= \discover ->+ when discover $ do+ net <- getNetwork+ ss <- concat <$> mapM (toSockAddr net) net.seeds+ mapM_ newPeer ss++logConnectedPeers :: (MonadManager m, MonadLoggerIO m) => m ()+logConnectedPeers = do+ m <- asks (.config.maxPeers)+ l <- length <$> getConnectedPeers+ $(logInfoS) "PeerMgr" $+ "Peers connected: " <> cs (show l) <> "/" <> cs (show m)++getOnlinePeers :: (MonadManager m) => m [OnlinePeer]+getOnlinePeers =+ asks (.peers) >>= readTVarIO++getConnectedPeers :: (MonadManager m) => m [OnlinePeer]+getConnectedPeers =+ filter (.online) <$> getOnlinePeers++managerEvent :: (MonadManager m) => PeerEvent -> m ()+managerEvent e =+ publish e =<< asks (.config.pub)++dispatch ::+ ( MonadUnliftIO m,+ MonadManager m,+ MonadLoggerIO m+ ) =>+ PeerMgrMessage ->+ m ()+dispatch (PeerVersion p v) = do+ $(logDebugS) "PeerMgr" $+ "Received peer " <> p.label <> " version: " <> cs (show v)+ b <- asks (.peers)+ e <- runExceptT $ do+ o <- ExceptT . atomically $ setPeerVersion b p v+ when o.online $ announcePeer p+ case e of+ Right () -> do+ $(logDebugS) "PeerMgr" $+ "Sending version ack to peer: " <> p.label+ MVerAck `sendMessage` p+ Left x -> do+ $(logErrorS) "PeerMgr" $+ "Version rejected for peer "+ <> p.label+ <> ": "+ <> cs (show x)+ killPeer x p+dispatch (PeerVerAck p) = do+ b <- asks (.peers)+ atomically (setPeerVerAck b p) >>= \case+ Just o -> do+ $(logDebugS) "PeerMgr" $+ "Received version ack from peer: "+ <> p.label+ when o.online $+ announcePeer p+ Nothing -> do+ $(logErrorS) "PeerMgr" $+ "Received verack from unknown peer: "+ <> p.label+ killPeer UnknownPeer p+dispatch (PeerAddrs p nas) = do+ $(logDebugS) "PeerMgr" $+ "Received addresses from peer " <> p.label+ discover <- asks (.config.discover)+ when discover $ do+ let sas = map (hostToSockAddr . (.address)) nas+ forM_ (zip [(1 :: Int) ..] sas) $ \(i, a) -> do+ $(logDebugS) "PeerMgr" $+ "Got peer address "+ <> cs (show i)+ <> "/"+ <> cs (show (length sas))+ <> ": "+ <> cs (show a)+ <> " from peer "+ <> p.label+ newPeer a+dispatch (PeerPong p n) = do+ b <- asks (.peers)+ $(logDebugS) "PeerMgr" $+ "Received pong "+ <> cs (show n)+ <> " from: "+ <> p.label+ now <- liftIO getCurrentTime+ atomically (gotPong b n now p)+dispatch (PeerPing p n) = do+ $(logDebugS) "PeerMgr" $+ "Responding to ping "+ <> cs (show n)+ <> " from: "+ <> p.label+ MPong (Pong n) `sendMessage` p+dispatch (ManagerBest h) = do+ $(logDebugS) "PeerMgr" $+ "Setting best block to " <> cs (show h)+ putBestBlock h+dispatch (Connect sa) = do+ connectPeer sa+dispatch (PeerDied a e) = do+ processPeerOffline a e+dispatch (CheckPeer p) = do+ $(logDebugS) "PeerManager" $+ "Housekeeping for peer " <> p.label+ checkPeer p+dispatch (PeerTickle p) = do+ $(logDebugS) "PeerMgr" $+ "Tickled peer " <> p.label+ b <- asks (.peers)+ now <- liftIO getCurrentTime+ atomically $+ modifyPeer b p $ \o ->+ o {tickled = now}++checkPeer :: (MonadManager m, MonadLoggerIO m) => Peer -> m ()+checkPeer p = do+ busy <- getBusy p+ b <- asks (.peers)+ mp <- asks (.peers) >>= atomically . flip findPeer p+ case mp of+ Nothing -> return ()+ Just o ->+ liftIO getCurrentTime >>= \now -> do+ maxLife <- asks (.config.maxPeerLife)+ let disconnect = maxLife `addUTCTime` o.connected+ when (now > disconnect) $ do+ $(logErrorS) "PeerMgr" $+ "Disconnecting old peer "+ <> p.label+ <> " online since "+ <> cs (show o.connected)+ killPeer PeerTooOld p+ timeout <- asks (.config.timeout)+ let pingTime = timeout `addUTCTime` o.tickled+ when (now > pingTime) $+ case o.ping of+ Nothing ->+ sendPing p+ Just _ -> do+ $(logWarnS) "PeerMgr" $+ "Peer ping timeout: " <> p.label+ killPeer PeerTimeout p++sendPing :: (MonadManager m, MonadLoggerIO m) => Peer -> m ()+sendPing p = do+ b <- asks (.peers)+ atomically (findPeer b p) >>= \case+ Nothing ->+ $(logWarnS) "PeerMgr" $+ "Will not ping unknown peer: " <> p.label+ Just o+ | o.online -> do+ n <- liftIO randomIO+ now <- liftIO getCurrentTime+ atomically (setPeerPing b n now p)+ $(logDebugS) " PeerManager" $+ "Sending ping "+ <> cs (show n)+ <> " to: "+ <> p.label+ MPing (Ping n) `sendMessage` p+ | otherwise -> return ()++processPeerOffline ::+ (MonadManager m, MonadLoggerIO m) =>+ Child ->+ Maybe SomeException ->+ m ()+processPeerOffline a e = do+ b <- asks (.peers)+ atomically (findPeerAsync b a) >>= \case+ Nothing -> log_unknown e+ Just o -> do+ let p = o.mailbox+ if o.online+ then do+ log_disconnected p e+ managerEvent $ PeerDisconnected p+ else log_not_connect p e+ atomically $ removePeer b p+ logConnectedPeers+ where+ log_unknown Nothing =+ $(logErrorS)+ "PeerMgr"+ "Disconnected unknown peer"+ log_unknown (Just x) =+ $(logErrorS) "PeerMgr" $+ "Unknown peer died: " <> cs (show x)+ log_disconnected p Nothing =+ $(logWarnS) "PeerMgr" $+ "Disconnected peer: " <> p.label+ log_disconnected p (Just x) =+ $(logErrorS) "PeerMgr" $+ "Peer " <> p.label <> " died: " <> cs (show x)+ log_not_connect p Nothing =+ $(logWarnS) "PeerMgr" $+ "Could not connect to peer " <> p.label+ log_not_connect p (Just x) =+ $(logErrorS) "PeerMgr" $+ "Could not connect to peer "+ <> p.label+ <> ": "+ <> cs (show x)++announcePeer :: (MonadManager m, MonadLoggerIO m) => Peer -> m ()+announcePeer p = do+ b <- asks (.peers)+ atomically (findPeer b p) >>= \case+ Just OnlinePeer {online = True} -> do+ $(logInfoS) "PeerMgr" $+ "Connected to peer " <> p.label+ managerEvent $ PeerConnected p+ logConnectedPeers+ Just OnlinePeer {online = False} ->+ return ()+ Nothing ->+ $(logErrorS) "PeerMgr" $+ "Not announcing disconnected peer: "+ <> p.label++getNewPeer :: (MonadUnliftIO m, MonadManager m) => m (Maybe SockAddr)+getNewPeer =+ runMaybeT $ lift loadPeers >> go+ where+ go = do+ b <- asks (.addresses)+ ks <- readTVarIO b+ guard . not $ Set.null ks+ let xs = Set.toList ks+ a <- liftIO $ randomRIO (0, length xs - 1)+ let p = xs !! a+ o <- asks (.peers)+ m <- atomically $ do+ modifyTVar b $ Set.delete p+ findPeerAddress o p+ maybe (return p) (const go) m++connectPeer ::+ ( MonadUnliftIO m,+ MonadManager m,+ MonadLoggerIO m+ ) =>+ SockAddr ->+ m ()+connectPeer sa = do+ os <- asks (.peers)+ atomically (findPeerAddress os sa) >>= \case+ Just _ ->+ $(logErrorS) "PeerMgr" $+ "Attempted to connect to peer twice: " <> cs (show sa)+ Nothing -> do+ $(logInfoS) "PeerMgr" $ "Connecting to " <> cs (show sa)+ PeerMgrConfig+ { address = ad,+ net = net+ } <-+ asks (.config)+ sup <- asks (.supervisor)+ conn <- asks (.config.connect)+ pub <- asks (.config.pub)+ nonce <- liftIO randomIO+ bb <- getBestBlock+ now <- liftIO getCurrentTime+ let rmt = NetworkAddress (srv net) (sockToHostAddress sa)+ unix = floor (utcTimeToPOSIXSeconds now)+ ver = buildVersion net nonce bb ad rmt unix+ text = cs (show sa)+ (inbox, mailbox) <- newMailbox+ let pc =+ PeerConfig+ { pub = pub,+ net = net,+ label = text,+ connect = conn sa+ }+ busy <- newTVarIO False+ p <- wrapPeer pc busy mailbox+ a <- withRunInIO $ \io ->+ sup `addChild` io (launch pc busy inbox p)+ MVersion ver `sendMessage` p+ b <- asks (.peers)+ atomically $+ insertPeer+ b+ OnlinePeer+ { address = sa,+ verack = False,+ online = False,+ version = Nothing,+ async = a,+ mailbox = p,+ nonce = nonce,+ pings = [],+ ping = Nothing,+ connected = now,+ tickled = now+ }+ where+ srv net+ | net.segWit = 8+ | otherwise = 0+ launch pc busy inbox p =+ ask >>= \mgr ->+ withPeerLoop sa p mgr $ \a ->+ link a >> peer pc busy inbox++withPeerLoop ::+ (MonadUnliftIO m, MonadLogger m) =>+ SockAddr ->+ Peer ->+ PeerMgr ->+ (Async a -> m a) ->+ m a+withPeerLoop _ p mgr =+ withAsync . forever $ do+ let x = mgr.config.timeout+ y = floor (x * 1000000)+ r <- liftIO $ randomRIO (y * 3 `div` 4, y)+ threadDelay r+ managerCheck p mgr++withConnectLoop ::+ (MonadUnliftIO m, MonadManager m) =>+ m a ->+ m a+withConnectLoop act =+ withAsync go $ \a ->+ link a >> act+ where+ go = forever $ do+ l <- length <$> getOnlinePeers+ x <- asks (.config.maxPeers)+ when (l < x) $+ getNewPeer >>= mapM_ (\sa -> ask >>= managerConnect sa)+ delay <-+ liftIO $+ randomRIO+ ( 100 * 1000,+ 10 * 500 * 1000+ )+ threadDelay delay++newPeer :: (MonadIO m, MonadManager m) => SockAddr -> m ()+newPeer sa = do+ b <- asks (.addresses)+ o <- asks (.peers)+ atomically $+ findPeerAddress o sa >>= \case+ Just _ -> return ()+ Nothing -> modifyTVar b $ Set.insert sa++gotPong :: TVar [OnlinePeer] -> Word64 -> UTCTime -> Peer -> STM ()+gotPong b nonce now p = void . runMaybeT $ do+ o <- MaybeT (findPeer b p)+ (time, old_nonce) <- MaybeT (return o.ping)+ guard $ nonce == old_nonce+ let diff = now `diffUTCTime` time+ lift $+ insertPeer+ b+ o+ { ping = Nothing,+ pings = sort $ take 11 $ diff : o.pings+ }++setPeerPing :: TVar [OnlinePeer] -> Word64 -> UTCTime -> Peer -> STM ()+setPeerPing b nonce now p =+ modifyPeer b p $ \o -> o {ping = Just (now, nonce)}++setPeerVersion ::+ TVar [OnlinePeer] ->+ Peer ->+ Version ->+ STM (Either PeerException OnlinePeer)+setPeerVersion b p v = runExceptT $ do+ when (v.services .&. nodeNetwork == 0) $+ throwError NotNetworkPeer+ ops <- lift $ readTVar b+ when (any ((v.nonce ==) . (.nonce)) ops) $+ throwError PeerIsMyself+ lift (findPeer b p) >>= \case+ Nothing -> throwError UnknownPeer+ Just o -> do+ let n =+ o+ { version = Just v,+ online = o.verack+ }+ lift $ insertPeer b n+ return n++setPeerVerAck :: TVar [OnlinePeer] -> Peer -> STM (Maybe OnlinePeer)+setPeerVerAck b p = runMaybeT $ do+ o <- MaybeT $ findPeer b p+ let n =+ o+ { verack = True,+ online = isJust o.version+ }+ lift $ insertPeer b n+ return n++findPeer :: TVar [OnlinePeer] -> Peer -> STM (Maybe OnlinePeer)+findPeer b p =+ find ((== p) . (.mailbox))+ <$> readTVar b++insertPeer :: TVar [OnlinePeer] -> OnlinePeer -> STM ()+insertPeer b o =+ modifyTVar b $ \x -> sort . nub $ o : x++modifyPeer ::+ TVar [OnlinePeer] ->+ Peer ->+ (OnlinePeer -> OnlinePeer) ->+ STM ()+modifyPeer b p f =+ findPeer b p >>= \case+ Nothing -> return ()+ Just o -> insertPeer b $ f o++removePeer :: TVar [OnlinePeer] -> Peer -> STM ()+removePeer b p =+ modifyTVar b $+ filter ((/= p) . (.mailbox))++findPeerAsync ::+ TVar [OnlinePeer] ->+ Async () ->+ STM (Maybe OnlinePeer)+findPeerAsync b a =+ find ((== a) . (.async))+ <$> readTVar b++findPeerAddress ::+ TVar [OnlinePeer] ->+ SockAddr ->+ STM (Maybe OnlinePeer)+findPeerAddress b a =+ find ((== a) . (.address))+ <$> readTVar b++getPeers :: (MonadIO m) => PeerMgr -> m [OnlinePeer]+getPeers = runReaderT getConnectedPeers++getOnlinePeer ::+ (MonadIO m) =>+ Peer ->+ PeerMgr ->+ m (Maybe OnlinePeer)+getOnlinePeer p =+ runReaderT $ asks (.peers) >>= atomically . (`findPeer` p)++managerCheck :: (MonadIO m) => Peer -> PeerMgr -> m ()+managerCheck p mgr =+ CheckPeer p `send` mgr.mailbox++managerConnect :: (MonadIO m) => SockAddr -> PeerMgr -> m ()+managerConnect sa mgr =+ Connect sa `send` mgr.mailbox++peerMgrBest :: (MonadIO m) => BlockHeight -> PeerMgr -> m ()+peerMgrBest bh mgr =+ ManagerBest bh `send` mgr.mailbox++peerMgrVerAck :: (MonadIO m) => Peer -> PeerMgr -> m ()+peerMgrVerAck p mgr =+ PeerVerAck p `send` mgr.mailbox++peerMgrVersion ::+ (MonadIO m) =>+ Peer ->+ Version ->+ PeerMgr ->+ m ()+peerMgrVersion p ver mgr =+ PeerVersion p ver `send` mgr.mailbox++peerMgrPing ::+ (MonadIO m) =>+ Peer ->+ Word64 ->+ PeerMgr ->+ m ()+peerMgrPing p nonce mgr =+ PeerPing p nonce `send` mgr.mailbox++peerMgrPong ::+ (MonadIO m) =>+ Peer ->+ Word64 ->+ PeerMgr ->+ m ()+peerMgrPong p nonce mgr =+ PeerPong p nonce `send` mgr.mailbox++peerMgrAddrs ::+ (MonadIO m) =>+ Peer ->+ [NetworkAddress] ->+ PeerMgr ->+ m ()+peerMgrAddrs p addrs mgr =+ PeerAddrs p addrs `send` mgr.mailbox++peerMgrTickle ::+ (MonadIO m) =>+ Peer ->+ PeerMgr ->+ m ()+peerMgrTickle p mgr =+ PeerTickle p `send` mgr.mailbox++toHostService :: String -> (Maybe String, Maybe String)+toHostService str =+ let host = case m6 of+ Just (x, _) -> Just x+ Nothing -> case takeWhile (/= ':') str of+ [] -> Nothing+ xs -> Just xs+ srv = case m6 of+ Just (_, y) -> s y+ Nothing -> s str+ s xs =+ case dropWhile (/= ':') xs of+ [] -> Nothing+ _ : ys -> Just ys+ m6 = case str of+ (x : xs)+ | x == '[' -> do+ i <- elemIndex ']' xs+ return $ second tail $ splitAt i xs+ | x == ':' -> do+ return (str, "")+ _ -> Nothing+ in (host, srv)++toSockAddr :: (MonadUnliftIO m) => Network -> String -> m [SockAddr]+toSockAddr net str =+ go `catch` e+ where+ go = fmap (map addrAddress) $ liftIO $ getAddrInfo Nothing host srv+ (host, srv) =+ second (<|> Just (show net.defaultPort)) $+ toHostService str+ e :: (Monad m) => SomeException -> m [SockAddr]+ e _ = return []++median :: (Ord a, Fractional a) => [a] -> Maybe a+median ls+ | null ls =+ Nothing+ | even (length ls) =+ Just . (/ 2) . sum . take 2 $+ drop (length ls `div` 2 - 1) ls'+ | otherwise =+ Just (ls' !! (length ls `div` 2))+ where+ ls' = sort ls++buildVersion ::+ Network ->+ Word64 ->+ BlockHeight ->+ NetworkAddress ->+ NetworkAddress ->+ Word64 ->+ Version+buildVersion net nonce height loc rmt time =+ Version+ { version = myVersion,+ services = loc.services,+ timestamp = time,+ addrRecv = rmt,+ addrSend = loc,+ nonce = nonce,+ userAgent = VarString net.userAgent,+ startHeight = height,+ relay = True+ }++myVersion :: Word32+myVersion = 70012
test/Haskoin/NodeSpec.hs view
@@ -1,294 +1,340 @@-{-# LANGUAGE FlexibleContexts #-}-{-# LANGUAGE LambdaCase #-}+{-# LANGUAGE DuplicateRecordFields #-}+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE LambdaCase #-} {-# LANGUAGE MultiParamTypeClasses #-}-{-# LANGUAGE OverloadedStrings #-}-{-# LANGUAGE RecordWildCards #-}-module Haskoin.NodeSpec- ( spec- ) where+{-# LANGUAGE OverloadedRecordDot #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE NoFieldSelectors #-} -import Conduit (awaitForever, concatMapC, foldC, mapMC,- runConduit, takeCE, yield, (.|))-import Control.Monad (forM_, forever, replicateM)-import Control.Monad.Logger (runNoLoggingT)-import Control.Monad.Trans (lift)-import Data.ByteString (ByteString)-import qualified Data.ByteString as B-import Data.ByteString.Base64 (decodeBase64Lenient)-import Data.Default (def)-import Data.Either (fromRight)-import Data.List (find)-import Data.Maybe (isJust, mapMaybe)-import Data.Serialize (decode, get, runGet, runPut)-import Data.Time.Clock.POSIX (getPOSIXTime)-import qualified Database.RocksDB as R-import Haskoin (Block (..), BlockHash (..),- BlockHeader (..), BlockNode (..),- GetData (..), GetHeaders (..),- Headers (..), InvType (..),- InvVector (..), Message (..),- MessageHeader (..), Network (..),- NetworkAddress (..), Ping (..),- Pong (..), VarInt (..), Version (..),- bchRegTest, buildMerkleRoot,- getMessage, headerHash, nodeNetwork,- putMessage, sockToHostAddress, txHash)-import Haskoin.Node-import Network.Socket (SockAddr (..), AddrInfo (addrAddress))-import NQE (Inbox, Mailbox, inboxToMailbox,- newInbox, receive, receiveMatch, send,- withPublisher, withSubscription)-import System.Random (randomIO)-import Test.Hspec-import Test.Hspec.QuickCheck-import UnliftIO (MonadIO, MonadUnliftIO, liftIO,- throwString, withAsync,- withSystemTempDirectory)+module Haskoin.NodeSpec (spec) where +import Conduit+ ( awaitForever,+ concatMapC,+ foldC,+ mapMC,+ runConduit,+ takeCE,+ yield,+ (.|),+ )+import Control.Monad (forM_, forever, replicateM)+import Control.Monad.Cont+import Control.Monad.Logger (runNoLoggingT)+import Control.Monad.Trans (lift)+import Data.ByteString (ByteString)+import qualified Data.ByteString as B+import Data.ByteString.Base64 (decodeBase64Lenient)+import Data.Default (def)+import Data.Either (fromRight)+import Data.List (find)+import Data.Maybe (isJust, mapMaybe)+import Data.Serialize (decode, get, runGet, runPut)+import Data.Time.Clock.POSIX (getPOSIXTime)+import qualified Database.RocksDB as R+import Haskoin+ ( Block (..),+ BlockHash (..),+ BlockHeader (..),+ BlockNode (..),+ GetData (..),+ GetHeaders (..),+ Headers (..),+ InvType (..),+ InvVector (..),+ Message (..),+ MessageHeader (..),+ Network (..),+ NetworkAddress (..),+ Ping (..),+ Pong (..),+ VarInt (..),+ Version (..),+ bchRegTest,+ buildMerkleRoot,+ getMessage,+ headerHash,+ nodeNetwork,+ putMessage,+ sockToHostAddress,+ txHash,+ )+import Haskoin.Node+import NQE+ ( Inbox,+ Mailbox,+ inboxToMailbox,+ newInbox,+ receive,+ receiveMatch,+ send,+ withPublisher,+ withSubscription,+ )+import Network.Socket (AddrInfo (addrAddress), SockAddr (..))+import System.Random (randomIO)+import Test.Hspec+import Test.Hspec.QuickCheck+import UnliftIO+ ( MonadIO,+ MonadUnliftIO,+ liftIO,+ throwString,+ withAsync,+ withSystemTempDirectory,+ )+ data TestNode = TestNode- { testMgr :: PeerManager- , testChain :: Chain- , nodeEvents :: Inbox NodeEvent- }+ { testMgr :: PeerMgr,+ testChain :: Chain,+ nodeEvents :: Inbox NodeEvent+ } -dummyPeerConnect :: Network- -> NetworkAddress- -> SockAddr- -> WithConnection+dummyPeerConnect ::+ Network ->+ NetworkAddress ->+ SockAddr ->+ WithConnection dummyPeerConnect net ad sa f = do- r <- newInbox- s <- newInbox- let s' = inboxToMailbox s- withAsync (go r s') $ \_ -> do- let o = awaitForever (`send` r)- i = forever (receive s >>= yield)- f (Conduits i o) :: IO ()+ r <- newInbox+ s <- newInbox+ let s' = inboxToMailbox s+ withAsync (go r s') $ \_ -> do+ let o = awaitForever (`send` r)+ i = forever (receive s >>= yield)+ f (Conduits i o) :: IO () where go :: Inbox ByteString -> Mailbox ByteString -> IO () go r s = do- nonce <- randomIO- now <- round <$> liftIO getPOSIXTime- let rmt = NetworkAddress 0 (sockToHostAddress sa)- ver = buildVersion net nonce 0 ad rmt now- runPut (putMessage net (MVersion ver)) `send` s- runConduit $- forever (receive r >>= yield) .| inc .| concatMapC mockPeerReact .|- outc .|- awaitForever (`send` s)+ nonce <- randomIO+ now <- round <$> liftIO getPOSIXTime+ let rmt = NetworkAddress 0 (sockToHostAddress sa)+ ver = buildVersion net nonce 0 ad rmt now+ runPut (putMessage net (MVersion ver)) `send` s+ runConduit $+ forever (receive r >>= yield)+ .| inc+ .| concatMapC mockPeerReact+ .| outc+ .| awaitForever (`send` s) outc = mapMC $ \msg -> return $ runPut (putMessage net msg) inc =- forever $ do- x <- takeCE 24 .| foldC- case decode x of- Left _ -> error "Dummy peer not decode message header"- Right (MessageHeader _ _ len _) -> do- y <- takeCE (fromIntegral len) .| foldC- case runGet (getMessage net) $ x `B.append` y of- Left e ->- error $- "Dummy peer could not decode payload: " <> show e- Right msg -> yield msg+ forever $ do+ x <- takeCE 24 .| foldC+ y <- case decode x of+ Left _ -> error "Dummy peer not decode message header"+ Right (MessageHeader _ _ len _) ->+ takeCE (fromIntegral len) .| foldC+ case runGet (getMessage net) $ x `B.append` y of+ Left e ->+ error $+ "Dummy peer could not decode payload: " <> show e+ Right msg -> yield msg mockPeerReact :: Message -> [Message] mockPeerReact (MPing (Ping n)) = [MPong (Pong n)] mockPeerReact (MVersion _) = [MVerAck] mockPeerReact (MGetHeaders (GetHeaders _ _hs _)) = [MHeaders (Headers hs')] where- f b = (blockHeader b, VarInt (fromIntegral (length (blockTxns b))))+ f b = (b.header, (VarInt . fromIntegral . length) b.txs) hs' = map f allBlocks mockPeerReact (MGetData (GetData ivs)) = mapMaybe f ivs where f (InvVector InvBlock h) = MBlock <$> find (l h) allBlocks- f _ = Nothing- l h b = headerHash (blockHeader b) == BlockHash h+ f _ = Nothing+ l h b = headerHash b.header == BlockHash h mockPeerReact _ = [] - spec :: Spec spec = do- let net = bchRegTest- describe "reads address/port combinations" $ do- prop "reads arbitrary addresses" $ \(e, w1, w2, w3, w4, b) -> do- let p = toEnum (e `mod` 65536)- a = if b - then SockAddrInet p w1- else SockAddrInet6 p 0 (w1, w2, w3, w4) 0- s <- head <$> toSockAddr net (show a)- s `shouldBe` a- it "reads some specific addresses" $ do- toHostService "localhost" `shouldBe` (Just "localhost", Nothing)- toHostService "::1" `shouldBe` (Just "::1", Nothing)- toHostService "localhost:8080" `shouldBe` (Just "localhost", Just "8080")- toHostService "example.com" `shouldBe` (Just "example.com", Nothing)- toHostService "api.example.com:443" `shouldBe` (Just "api.example.com", Just "443")- toHostService "api.example.com:http" `shouldBe` (Just "api.example.com", Just "http")- toHostService "[::1]" `shouldBe` (Just "::1", Nothing)- toHostService "[::1]:8080" `shouldBe` (Just "::1", Just "8080")- toHostService "[2002::dead:beef]:ssh" `shouldBe` (Just "2002::dead:beef", Just "ssh")- describe "peer manager on test network" $ do- it "connects to a peer" $- withTestNode net "connect-one-peer" $ \TestNode {..} -> do- p <- waitForPeer nodeEvents- Just OnlinePeer {onlinePeerVersion = Just Version {version = ver}} <-- getOnlinePeer p testMgr- ver `shouldSatisfy` (>= 70002)- it "downloads some blocks" $- withTestNode net "get-blocks" $ \TestNode {..} -> do- let h1 =- "3094ed3592a06f3d8e099eed2d9c1192329944f5df4a48acb29e08f12cfbb660"- h2 =- "0c89955fc5c9f98ecc71954f167b938138c90c6a094c4737f2e901669d26763f"- p <- waitForPeer nodeEvents- pbs <- getBlocks net 10 p [h1, h2]- pbs `shouldSatisfy` isJust- let Just [b1, b2] = pbs- headerHash (blockHeader b1) `shouldBe` h1- headerHash (blockHeader b2) `shouldBe` h2- let testMerkle b =- merkleRoot (blockHeader b) `shouldBe`- buildMerkleRoot (map txHash (blockTxns b))- testMerkle b1- testMerkle b2- describe "chain on test network" $ do- it "syncs some headers" $- withTestNode net "connect-sync" $ \TestNode {..} -> do- let bh =- "3bfa0c6da615fc45aa44ddea6854ac19d16f3ca167e0e21ac2cc262a49c9b002"- ah =- "7dc835a78a55fa76f9184dc4f6663a73e418c7afec789c5ae25e432fd7fc8467"- bn <-- receiveMatch nodeEvents $ \case- ChainEvent (ChainBestBlock bn)- | nodeHeight bn > 0 -> Just bn- _ -> Nothing- bb <- chainGetBest testChain- nodeHeight bb `shouldSatisfy` (== 15)- an <-- maybe (throwString "No ancestor found") return =<<- chainGetAncestor 10 bn testChain- headerHash (nodeHeader bn) `shouldBe` bh- headerHash (nodeHeader an) `shouldBe` ah- it "downloads some block parents" $- withTestNode net "parents" $ \TestNode {..} -> do- let hs =- [ "52e886df7b166d961ac2d3d2d561d806325d51a609dc0a5d9d5fcb65d47906d7"- , "2537a081b9e2b24d217fac2886f387758cb3aa4e4956b3be7ed229bafbb71b0f"- , "7c72f306215a296f9714320a497b1f2cb5f9b99f162d7e04333c243fac9a54d8"- ]- [_, bn] <-- replicateM 2 $- receiveMatch nodeEvents $ \case- ChainEvent (ChainBestBlock bn) -> Just bn- _ -> Nothing- nodeHeight bn `shouldBe` 15- ps <- chainGetParents 12 bn testChain- length ps `shouldBe` 3- forM_ (zip ps hs) $ \(p, h) ->- headerHash (nodeHeader p) `shouldBe` h+ let net = bchRegTest+ describe "reads address/port combinations" $ do+ prop "reads arbitrary addresses" $ \(e, w1, w2, w3, w4, b) -> do+ let p = toEnum (e `mod` 65536)+ a =+ if b+ then SockAddrInet p w1+ else SockAddrInet6 p 0 (w1, w2, w3, w4) 0+ s <- head <$> toSockAddr net (show a)+ s `shouldBe` a+ it "reads some specific addresses" $ do+ toHostService "localhost" `shouldBe` (Just "localhost", Nothing)+ toHostService "::1" `shouldBe` (Just "::1", Nothing)+ toHostService "localhost:8080" `shouldBe` (Just "localhost", Just "8080")+ toHostService "example.com" `shouldBe` (Just "example.com", Nothing)+ toHostService "api.example.com:443" `shouldBe` (Just "api.example.com", Just "443")+ toHostService "api.example.com:http" `shouldBe` (Just "api.example.com", Just "http")+ toHostService "[::1]" `shouldBe` (Just "::1", Nothing)+ toHostService "[::1]:8080" `shouldBe` (Just "::1", Just "8080")+ toHostService "[2002::dead:beef]:ssh" `shouldBe` (Just "2002::dead:beef", Just "ssh")+ describe "peer manager on test network" $ do+ it "connects to a peer" $+ withTestNode net "connect-one-peer" $ \TestNode {..} -> do+ p <- waitForPeer nodeEvents+ Just OnlinePeer {version = Just Version {version = ver}} <-+ getOnlinePeer p testMgr+ ver `shouldSatisfy` (>= 70002)+ it "downloads some blocks" $+ withTestNode net "get-blocks" $ \TestNode {..} -> do+ let h1 =+ "3094ed3592a06f3d8e099eed2d9c1192329944f5df4a48acb29e08f12cfbb660"+ h2 =+ "0c89955fc5c9f98ecc71954f167b938138c90c6a094c4737f2e901669d26763f"+ p <- waitForPeer nodeEvents+ pbs <- getBlocks net 10 p [h1, h2]+ pbs `shouldSatisfy` isJust+ let Just [b1, b2] = pbs+ headerHash b1.header `shouldBe` h1+ headerHash b2.header `shouldBe` h2+ let ths b = map txHash b.txs+ testMerkle b = b.header.merkle `shouldBe` buildMerkleRoot (ths b)+ testMerkle b1+ testMerkle b2+ describe "chain on test network" $ do+ it "syncs some headers" $+ withTestNode net "connect-sync" $ \TestNode {..} -> do+ let bh =+ "3bfa0c6da615fc45aa44ddea6854ac19d16f3ca167e0e21ac2cc262a49c9b002"+ ah =+ "7dc835a78a55fa76f9184dc4f6663a73e418c7afec789c5ae25e432fd7fc8467"+ bn <-+ receiveMatch nodeEvents $ \case+ ChainEvent (ChainBestBlock bn)+ | bn.height > 0 -> Just bn+ _ -> Nothing+ bb <- chainGetBest testChain+ bb.height `shouldSatisfy` (== 15)+ an <-+ maybe (throwString "No ancestor found") return+ =<< chainGetAncestor 10 bn testChain+ headerHash bn.header `shouldBe` bh+ headerHash an.header `shouldBe` ah+ it "downloads some block parents" $+ withTestNode net "parents" $ \TestNode {..} -> do+ let hs =+ [ "52e886df7b166d961ac2d3d2d561d806325d51a609dc0a5d9d5fcb65d47906d7",+ "2537a081b9e2b24d217fac2886f387758cb3aa4e4956b3be7ed229bafbb71b0f",+ "7c72f306215a296f9714320a497b1f2cb5f9b99f162d7e04333c243fac9a54d8"+ ]+ [_, bn] <-+ replicateM 2 $+ receiveMatch nodeEvents $ \case+ ChainEvent (ChainBestBlock bn) -> Just bn+ _ -> Nothing+ bn.height `shouldBe` 15+ ps <- chainGetParents 12 bn testChain+ length ps `shouldBe` 3+ forM_ (zip ps hs) $ \(p, h) ->+ headerHash p.header `shouldBe` h -waitForPeer :: MonadIO m => Inbox NodeEvent -> m Peer+waitForPeer :: (MonadIO m) => Inbox NodeEvent -> m Peer waitForPeer inbox =- receiveMatch inbox $ \case- PeerEvent (PeerConnected p) -> Just p- _ -> Nothing+ receiveMatch inbox $ \case+ PeerEvent (PeerConnected p) -> Just p+ _ -> Nothing withTestNode ::- MonadUnliftIO m- => Network- -> String- -> (TestNode -> m a)- -> m a-withTestNode net str f =- runNoLoggingT $- withSystemTempDirectory ("haskoin-node-test-" <> str <> "-") $ \w ->- withPublisher $ \pub ->- R.withDBCF w cfg cols $ \db -> do- let ad = NetworkAddress- nodeNetwork- (sockToHostAddress (SockAddrInet 0 0))- na = NetworkAddress- 0- (sockToHostAddress (SockAddrInet 0 0))- cfg' = NodeConfig- { nodeConfMaxPeers = 20- , nodeConfDB = db- , nodeConfColumnFamily = Just (head (R.columnFamilies db))- , nodeConfPeers = ["[::1]:17486"]- , nodeConfDiscover = False- , nodeConfNetAddr = na- , nodeConfNet = net- , nodeConfEvents = pub- , nodeConfTimeout = 120- , nodeConfPeerMaxLife = 48 * 3600- , nodeConfConnect = dummyPeerConnect net ad- }- withNode cfg' $ \(Node mgr ch) ->- withSubscription pub $ \sub ->- lift $- f TestNode { testMgr = mgr- , testChain = ch- , nodeEvents = sub- }+ (MonadUnliftIO m) =>+ Network ->+ String ->+ (TestNode -> m a) ->+ m a+withTestNode net str f = runNoLoggingT $ flip runContT return $ do+ w <- ContT $ withSystemTempDirectory ("haskoin-node-test-" <> str <> "-")+ pub <- ContT withPublisher+ sub <- ContT $ withSubscription pub+ db <- ContT $ R.withDBCF w cfg cols+ let ad =+ NetworkAddress+ nodeNetwork+ (sockToHostAddress (SockAddrInet 0 0))+ na =+ NetworkAddress+ 0+ (sockToHostAddress (SockAddrInet 0 0))+ cfg' =+ NodeConfig+ { maxPeers = 20,+ db = db,+ cf = Just (head (R.columnFamilies db)),+ peers = ["[::1]:17486"],+ discover = False,+ address = na,+ net = net,+ pub = pub,+ timeout = 120,+ maxPeerLife = 48 * 3600,+ connect = dummyPeerConnect net ad+ }+ Node mgr ch <- ContT $ withNode cfg'+ lift . lift $+ f+ TestNode+ { testMgr = mgr,+ testChain = ch,+ nodeEvents = sub+ } where- cfg = def{R.createIfMissing = True, R.errorIfExists = True}+ cfg = def {R.createIfMissing = True, R.errorIfExists = True} cols = [("node", def)] allBlocks :: [Block] allBlocks =- fromRight (error "Could not decode blocks") $+ fromRight (error "Could not decode blocks") $ runGet f (decodeBase64Lenient allBlocksBase64) where f = mapM (const get) [(1 :: Int) .. 15] allBlocksBase64 :: ByteString allBlocksBase64 =- "AAAAIAYibkYRGgtZyq8SYEPrW78ow086XjMqH8eytzzxiJEPakRJalmWTFwdvzNuH8fHLZEjn+4N\- \FNMANdB7ez2M4a3TFbNe//9/IAMAAAABAgAAAAEAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA\- \AAAAAP////8MUQEBCC9FQjMyLjAv/////wEA8gUqAQAAACMhAwTspkCjMezKs47BPpafou1jjsHf\- \1OHjgkqxnwEYkK9zrAAAAAAAAAAge0RDjOrqVayGUoQsbNTJcTXUM+psaHpmuiFy6hwo2T8yn0CL\- \7WDJw9hxl1kf5c4JySq3WJF8OPsoguzF7mXH3tQVs17//38gAAAAAAECAAAAAQAAAAAAAAAAAAAA\- \AAAAAAAAAAAAAAAAAAAAAAAAAAAA/////wxSAQEIL0VCMzIuMC//////AQDyBSoBAAAAIyEDBOym\- \QKMx7MqzjsE+lp+i7WOOwd/U4eOCSrGfARiQr3OsAAAAAAAAACCKlhzDaFkrsmO2FhmeQS9ONS8D\- \QsU4H97yNxVhyIXYJuG3a9cyQpdeETjCQ6JybgkwI0OOfa4eYazf7WWI5UAk1BWzXv//fyAEAAAA\- \AQIAAAABAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAD/////DFMBAQgvRUIzMi4wL///\- \//8BAPIFKgEAAAAjIQME7KZAozHsyrOOwT6Wn6LtY47B39Th44JKsZ8BGJCvc6wAAAAAAAAAIP/S\- \XiIJZqvUyBY90z72dv6+/GG50R3vc3UAK8AHP89wChmkVP6nefjOt+sNyhbKk9zia47F08oTNtC0\- \OG1zyuXVFbNe//9/IAEAAAABAgAAAAEAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAP//\- \//8MVAEBCC9FQjMyLjAv/////wEA8gUqAQAAACMhAwTspkCjMezKs47BPpafou1jjsHf1OHjgkqx\- \nwEYkK9zrAAAAAAAAAAgeQtE1s3YV/uS2jUouo3S9DJAVf5OGk+Nyx+No1mPH24b5JCkr/tSP0E/\- \NYVkVcE0ZHxbO/fu5wOd+8VolvPQYtUVs17//38gAAAAAAECAAAAAQAAAAAAAAAAAAAAAAAAAAAA\- \AAAAAAAAAAAAAAAAAAAA/////wxVAQEIL0VCMzIuMC//////AQDyBSoBAAAAIyEDBOymQKMx7Mqz\- \jsE+lp+i7WOOwd/U4eOCSrGfARiQr3OsAAAAAAAAACBgtvss8QiesqxISt/1RJkykhGcLe2eCY49\- \b6CSNe2UMOVYGZ++uRCKvaJ2+jo7akr7XsdXCYSAmuw6DwSO8lvF1RWzXv//fyAAAAAAAQIAAAAB\- \AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAD/////DFYBAQgvRUIzMi4wL/////8BAPIF\- \KgEAAAAjIQME7KZAozHsyrOOwT6Wn6LtY47B39Th44JKsZ8BGJCvc6wAAAAAAAAAID92Jp1mAeny\- \N0dMCWoMyTiBk3sWT5VxzI75ycVflYkMCnXLFhuwrMdBbZmXJinAJBUpN7BV0XvlM2PRmb7HQebV\- \FbNe//9/IAEAAAABAgAAAAEAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAP////8MVwEB\- \CC9FQjMyLjAv/////wEA8gUqAQAAACMhAwTspkCjMezKs47BPpafou1jjsHf1OHjgkqxnwEYkK9z\- \rAAAAAAAAAAgxEgEkhjf5p+ql8dETmdSCdCdk+vB26+V2SGLEuE1+kA1acGCdQoQBqec8P/knItJ\- \M213OIrDX6U5IB6fgIas7dYVs17//38gAQAAAAECAAAAAQAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA\- \AAAAAAAAAAAA/////wxYAQEIL0VCMzIuMC//////AQDyBSoBAAAAIyEDBOymQKMx7MqzjsE+lp+i\- \7WOOwd/U4eOCSrGfARiQr3OsAAAAAAAAACDku4EB5X7htWpHg+aMzzW1AABttpNQTew7K3Aj2fh/\- \OuOCPhJApmcXq5o42tkksFSuhYvcfqaSHCuuFgjo6ohz1hWzXv//fyAAAAAAAQIAAAABAAAAAAAA\- \AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAD/////DFkBAQgvRUIzMi4wL/////8BAPIFKgEAAAAj\- \IQME7KZAozHsyrOOwT6Wn6LtY47B39Th44JKsZ8BGJCvc6wAAAAAAAAAIKWpAhOWbkEN9vWf1uCu\- \eXtVOZIE9V1OE87iC+H9atBRtY4LPgaWUSVMNh9SeZK1NViIFMklbjsfqYiC4eA/VuLWFbNe//9/\- \IAAAAAABAgAAAAEAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAP////8MWgEBCC9FQjMy\- \LjAv/////wEA8gUqAQAAACMhAwTspkCjMezKs47BPpafou1jjsHf1OHjgkqxnwEYkK9zrAAAAAAA\- \AAAgZ4T81y9DXuJanHjsr8cY5HM6ZvbETRj5dvpViqc1yH0oN9OOruaO5mjdITJwweVCzjSQ5Wsl\- \vSOKaKvEX5j9l9YVs17//38gAAAAAAECAAAAAQAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA\- \AAAA/////wxbAQEIL0VCMzIuMC//////AQDyBSoBAAAAIyEDBOymQKMx7MqzjsE+lp+i7WOOwd/U\- \4eOCSrGfARiQr3OsAAAAAAAAACCV3J2A3qneSJ7Q/RuF8OPd8O1izIXvKElR/xg/+InGNEafu0Ul\- \3VYJR93zbAQuns9hUfAhA8MTBPk8bbDabDfo1hWzXv//fyAAAAAAAQIAAAABAAAAAAAAAAAAAAAA\- \AAAAAAAAAAAAAAAAAAAAAAAAAAD/////DFwBAQgvRUIzMi4wL/////8BAPIFKgEAAAAjIQME7KZA\- \ozHsyrOOwT6Wn6LtY47B39Th44JKsZ8BGJCvc6wAAAAAAAAAINcGedRly1+dXQrcCaZRXTIG2GHV\- \0tPCGpZtFnvfhuhSx8d3Azdv/MXRJgsb56qqmD5gsXiWUdi7ia7wsBZVylvWFbNe//9/IAEAAAAB\- \AgAAAAEAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAP////8MXQEBCC9FQjMyLjAv////\- \/wEA8gUqAQAAACMhAwTspkCjMezKs47BPpafou1jjsHf1OHjgkqxnwEYkK9zrAAAAAAAAAAgDxu3\- \+7op0n6+s1ZJTqqzjHWH84YorH8hTbLiuYGgNyWIkhaj0zR7Vc+fSRm4UYUaPsefRhq3fUt8glyS\- \D8P/5tcVs17//38gAwAAAAECAAAAAQAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA////\- \/wxeAQEIL0VCMzIuMC//////AQDyBSoBAAAAIyEDBOymQKMx7MqzjsE+lp+i7WOOwd/U4eOCSrGf\- \ARiQr3OsAAAAAAAAACDYVJqsPyQ8MwR+LRafufm1LB97SQoyFJdvKVohBvNyfD4/FxT2i0rlYQcS\- \TQAvTnehousK2P8T9c0qx4Yj72lT1xWzXv//fyAAAAAAAQIAAAABAAAAAAAAAAAAAAAAAAAAAAAA\- \AAAAAAAAAAAAAAAAAAD/////DF8BAQgvRUIzMi4wL/////8BAPIFKgEAAAAjIQME7KZAozHsyrOO\- \wT6Wn6LtY47B39Th44JKsZ8BGJCvc6wAAAAA"+ "AAAAIAYibkYRGgtZyq8SYEPrW78ow086XjMqH8eytzzxiJEPakRJalmWTFwdvzNuH8fHLZEjn+4N\+ \FNMANdB7ez2M4a3TFbNe//9/IAMAAAABAgAAAAEAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA\+ \AAAAAP////8MUQEBCC9FQjMyLjAv/////wEA8gUqAQAAACMhAwTspkCjMezKs47BPpafou1jjsHf\+ \1OHjgkqxnwEYkK9zrAAAAAAAAAAge0RDjOrqVayGUoQsbNTJcTXUM+psaHpmuiFy6hwo2T8yn0CL\+ \7WDJw9hxl1kf5c4JySq3WJF8OPsoguzF7mXH3tQVs17//38gAAAAAAECAAAAAQAAAAAAAAAAAAAA\+ \AAAAAAAAAAAAAAAAAAAAAAAAAAAA/////wxSAQEIL0VCMzIuMC//////AQDyBSoBAAAAIyEDBOym\+ \QKMx7MqzjsE+lp+i7WOOwd/U4eOCSrGfARiQr3OsAAAAAAAAACCKlhzDaFkrsmO2FhmeQS9ONS8D\+ \QsU4H97yNxVhyIXYJuG3a9cyQpdeETjCQ6JybgkwI0OOfa4eYazf7WWI5UAk1BWzXv//fyAEAAAA\+ \AQIAAAABAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAD/////DFMBAQgvRUIzMi4wL///\+ \//8BAPIFKgEAAAAjIQME7KZAozHsyrOOwT6Wn6LtY47B39Th44JKsZ8BGJCvc6wAAAAAAAAAIP/S\+ \XiIJZqvUyBY90z72dv6+/GG50R3vc3UAK8AHP89wChmkVP6nefjOt+sNyhbKk9zia47F08oTNtC0\+ \OG1zyuXVFbNe//9/IAEAAAABAgAAAAEAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAP//\+ \//8MVAEBCC9FQjMyLjAv/////wEA8gUqAQAAACMhAwTspkCjMezKs47BPpafou1jjsHf1OHjgkqx\+ \nwEYkK9zrAAAAAAAAAAgeQtE1s3YV/uS2jUouo3S9DJAVf5OGk+Nyx+No1mPH24b5JCkr/tSP0E/\+ \NYVkVcE0ZHxbO/fu5wOd+8VolvPQYtUVs17//38gAAAAAAECAAAAAQAAAAAAAAAAAAAAAAAAAAAA\+ \AAAAAAAAAAAAAAAAAAAA/////wxVAQEIL0VCMzIuMC//////AQDyBSoBAAAAIyEDBOymQKMx7Mqz\+ \jsE+lp+i7WOOwd/U4eOCSrGfARiQr3OsAAAAAAAAACBgtvss8QiesqxISt/1RJkykhGcLe2eCY49\+ \b6CSNe2UMOVYGZ++uRCKvaJ2+jo7akr7XsdXCYSAmuw6DwSO8lvF1RWzXv//fyAAAAAAAQIAAAAB\+ \AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAD/////DFYBAQgvRUIzMi4wL/////8BAPIF\+ \KgEAAAAjIQME7KZAozHsyrOOwT6Wn6LtY47B39Th44JKsZ8BGJCvc6wAAAAAAAAAID92Jp1mAeny\+ \N0dMCWoMyTiBk3sWT5VxzI75ycVflYkMCnXLFhuwrMdBbZmXJinAJBUpN7BV0XvlM2PRmb7HQebV\+ \FbNe//9/IAEAAAABAgAAAAEAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAP////8MVwEB\+ \CC9FQjMyLjAv/////wEA8gUqAQAAACMhAwTspkCjMezKs47BPpafou1jjsHf1OHjgkqxnwEYkK9z\+ \rAAAAAAAAAAgxEgEkhjf5p+ql8dETmdSCdCdk+vB26+V2SGLEuE1+kA1acGCdQoQBqec8P/knItJ\+ \M213OIrDX6U5IB6fgIas7dYVs17//38gAQAAAAECAAAAAQAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA\+ \AAAAAAAAAAAA/////wxYAQEIL0VCMzIuMC//////AQDyBSoBAAAAIyEDBOymQKMx7MqzjsE+lp+i\+ \7WOOwd/U4eOCSrGfARiQr3OsAAAAAAAAACDku4EB5X7htWpHg+aMzzW1AABttpNQTew7K3Aj2fh/\+ \OuOCPhJApmcXq5o42tkksFSuhYvcfqaSHCuuFgjo6ohz1hWzXv//fyAAAAAAAQIAAAABAAAAAAAA\+ \AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAD/////DFkBAQgvRUIzMi4wL/////8BAPIFKgEAAAAj\+ \IQME7KZAozHsyrOOwT6Wn6LtY47B39Th44JKsZ8BGJCvc6wAAAAAAAAAIKWpAhOWbkEN9vWf1uCu\+ \eXtVOZIE9V1OE87iC+H9atBRtY4LPgaWUSVMNh9SeZK1NViIFMklbjsfqYiC4eA/VuLWFbNe//9/\+ \IAAAAAABAgAAAAEAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAP////8MWgEBCC9FQjMy\+ \LjAv/////wEA8gUqAQAAACMhAwTspkCjMezKs47BPpafou1jjsHf1OHjgkqxnwEYkK9zrAAAAAAA\+ \AAAgZ4T81y9DXuJanHjsr8cY5HM6ZvbETRj5dvpViqc1yH0oN9OOruaO5mjdITJwweVCzjSQ5Wsl\+ \vSOKaKvEX5j9l9YVs17//38gAAAAAAECAAAAAQAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA\+ \AAAA/////wxbAQEIL0VCMzIuMC//////AQDyBSoBAAAAIyEDBOymQKMx7MqzjsE+lp+i7WOOwd/U\+ \4eOCSrGfARiQr3OsAAAAAAAAACCV3J2A3qneSJ7Q/RuF8OPd8O1izIXvKElR/xg/+InGNEafu0Ul\+ \3VYJR93zbAQuns9hUfAhA8MTBPk8bbDabDfo1hWzXv//fyAAAAAAAQIAAAABAAAAAAAAAAAAAAAA\+ \AAAAAAAAAAAAAAAAAAAAAAAAAAD/////DFwBAQgvRUIzMi4wL/////8BAPIFKgEAAAAjIQME7KZA\+ \ozHsyrOOwT6Wn6LtY47B39Th44JKsZ8BGJCvc6wAAAAAAAAAINcGedRly1+dXQrcCaZRXTIG2GHV\+ \0tPCGpZtFnvfhuhSx8d3Azdv/MXRJgsb56qqmD5gsXiWUdi7ia7wsBZVylvWFbNe//9/IAEAAAAB\+ \AgAAAAEAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAP////8MXQEBCC9FQjMyLjAv////\+ \/wEA8gUqAQAAACMhAwTspkCjMezKs47BPpafou1jjsHf1OHjgkqxnwEYkK9zrAAAAAAAAAAgDxu3\+ \+7op0n6+s1ZJTqqzjHWH84YorH8hTbLiuYGgNyWIkhaj0zR7Vc+fSRm4UYUaPsefRhq3fUt8glyS\+ \D8P/5tcVs17//38gAwAAAAECAAAAAQAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA////\+ \/wxeAQEIL0VCMzIuMC//////AQDyBSoBAAAAIyEDBOymQKMx7MqzjsE+lp+i7WOOwd/U4eOCSrGf\+ \ARiQr3OsAAAAAAAAACDYVJqsPyQ8MwR+LRafufm1LB97SQoyFJdvKVohBvNyfD4/FxT2i0rlYQcS\+ \TQAvTnehousK2P8T9c0qx4Yj72lT1xWzXv//fyAAAAAAAQIAAAABAAAAAAAAAAAAAAAAAAAAAAAA\+ \AAAAAAAAAAAAAAAAAAD/////DF8BAQgvRUIzMi4wL/////8BAPIFKgEAAAAjIQME7KZAozHsyrOO\+ \wT6Wn6LtY47B39Th44JKsZ8BGJCvc6wAAAAA"