packages feed

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 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"