haskoin-node 1.2.0 → 1.3.0
raw patch · 7 files changed
+612/−1257 lines, 7 filesdep +asyncdep +loggingdep +stmdep −monad-loggerdep −resourcetdep −unliftiodep ~rocksdb-haskell-jpruppdep ~rocksdb-queryPVP ok
version bump matches the API change (PVP)
Dependencies added: async, logging, stm, temporary
Dependencies removed: monad-logger, resourcet, unliftio
Dependency ranges changed: rocksdb-haskell-jprupp, rocksdb-query
API changes (from Hackage documentation)
- Haskoin.Node: chainBlockMain :: MonadIO m => BlockHash -> Chain -> m Bool
+ Haskoin.Node: chainBlockMain :: Chain -> BlockHash -> IO Bool
- Haskoin.Node: chainGetAncestor :: MonadIO m => BlockHeight -> BlockNode -> Chain -> m (Maybe BlockNode)
+ Haskoin.Node: chainGetAncestor :: Chain -> BlockHeight -> BlockNode -> IO (Maybe BlockNode)
- Haskoin.Node: chainGetBest :: MonadIO m => Chain -> m BlockNode
+ Haskoin.Node: chainGetBest :: Chain -> IO BlockNode
- Haskoin.Node: chainGetBlock :: MonadIO m => BlockHash -> Chain -> m (Maybe BlockNode)
+ Haskoin.Node: chainGetBlock :: Chain -> BlockHash -> IO (Maybe BlockNode)
- Haskoin.Node: chainGetParents :: MonadIO m => BlockHeight -> BlockNode -> Chain -> m [BlockNode]
+ Haskoin.Node: chainGetParents :: Chain -> BlockHeight -> BlockNode -> IO [BlockNode]
- Haskoin.Node: chainGetSplitBlock :: MonadIO m => BlockNode -> BlockNode -> Chain -> m BlockNode
+ Haskoin.Node: chainGetSplitBlock :: Chain -> BlockNode -> BlockNode -> IO BlockNode
- Haskoin.Node: chainHeaders :: MonadIO m => Peer -> [BlockHeader] -> Chain -> m ()
+ Haskoin.Node: chainHeaders :: Chain -> Peer -> [BlockHeader] -> IO ()
- Haskoin.Node: chainIsSynced :: MonadIO m => Chain -> m Bool
+ Haskoin.Node: chainIsSynced :: Chain -> IO Bool
- Haskoin.Node: chainPeerConnected :: MonadIO m => Peer -> Chain -> m ()
+ Haskoin.Node: chainPeerConnected :: Chain -> Peer -> IO ()
- Haskoin.Node: chainPeerDisconnected :: MonadIO m => Peer -> Chain -> m ()
+ Haskoin.Node: chainPeerDisconnected :: Chain -> Peer -> IO ()
- Haskoin.Node: getBlocks :: MonadUnliftIO m => Network -> Int -> Peer -> [BlockHash] -> m (Maybe [Block])
+ Haskoin.Node: getBlocks :: Network -> Int -> Peer -> [BlockHash] -> IO (Maybe [Block])
- Haskoin.Node: getBusy :: MonadIO m => Peer -> m Bool
+ Haskoin.Node: getBusy :: Peer -> IO Bool
- Haskoin.Node: getData :: MonadUnliftIO m => Int -> Peer -> GetData -> m (Maybe [Either Tx Block])
+ Haskoin.Node: getData :: Int -> Peer -> GetData -> IO (Maybe [Either Tx Block])
- Haskoin.Node: getOnlinePeer :: MonadIO m => Peer -> PeerMgr -> m (Maybe OnlinePeer)
+ Haskoin.Node: getOnlinePeer :: PeerMgr -> Peer -> IO (Maybe OnlinePeer)
- Haskoin.Node: getPeers :: MonadIO m => PeerMgr -> m [OnlinePeer]
+ Haskoin.Node: getPeers :: PeerMgr -> IO [OnlinePeer]
- Haskoin.Node: getTxs :: MonadUnliftIO m => Network -> Int -> Peer -> [TxHash] -> m (Maybe [Tx])
+ Haskoin.Node: getTxs :: Network -> Int -> Peer -> [TxHash] -> IO (Maybe [Tx])
- Haskoin.Node: killPeer :: MonadIO m => PeerException -> Peer -> m ()
+ Haskoin.Node: killPeer :: Peer -> IO ()
- Haskoin.Node: peer :: (MonadUnliftIO m, MonadLoggerIO m) => PeerConfig -> TVar Bool -> Inbox PeerMessage -> m ()
+ Haskoin.Node: peer :: PeerConfig -> TVar Bool -> Inbox PeerMessage -> IO ()
- Haskoin.Node: peerMgrAddrs :: MonadIO m => Peer -> [NetworkAddress] -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrAddrs :: PeerMgr -> Peer -> [NetworkAddress] -> IO ()
- Haskoin.Node: peerMgrBest :: MonadIO m => BlockHeight -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrBest :: PeerMgr -> BlockHeight -> IO ()
- Haskoin.Node: peerMgrPing :: MonadIO m => Peer -> Word64 -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrPing :: PeerMgr -> Peer -> Word64 -> IO ()
- Haskoin.Node: peerMgrPong :: MonadIO m => Peer -> Word64 -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrPong :: PeerMgr -> Peer -> Word64 -> IO ()
- Haskoin.Node: peerMgrVerAck :: MonadIO m => Peer -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrVerAck :: PeerMgr -> Peer -> IO ()
- Haskoin.Node: peerMgrVersion :: MonadIO m => Peer -> Version -> PeerMgr -> m ()
+ Haskoin.Node: peerMgrVersion :: PeerMgr -> Peer -> Version -> IO ()
- Haskoin.Node: pingPeer :: MonadUnliftIO m => Int -> Peer -> m Bool
+ Haskoin.Node: pingPeer :: Int -> Peer -> IO Bool
- Haskoin.Node: sendMessage :: MonadIO m => Message -> Peer -> m ()
+ Haskoin.Node: sendMessage :: Message -> Peer -> IO ()
- Haskoin.Node: setBusy :: MonadIO m => Peer -> m Bool
+ Haskoin.Node: setBusy :: Peer -> IO Bool
- Haskoin.Node: setFree :: MonadIO m => Peer -> m ()
+ Haskoin.Node: setFree :: Peer -> IO ()
- Haskoin.Node: ticklePeer :: MonadLoggerIO m => PeerMgr -> Peer -> m ()
+ Haskoin.Node: ticklePeer :: PeerMgr -> Peer -> IO ()
- Haskoin.Node: toSockAddr :: MonadUnliftIO m => Network -> String -> m [SockAddr]
+ Haskoin.Node: toSockAddr :: Network -> String -> IO [SockAddr]
- Haskoin.Node: withChain :: (MonadUnliftIO m, MonadLoggerIO m) => ChainConfig -> (Chain -> m a) -> m a
+ Haskoin.Node: withChain :: ChainConfig -> (Chain -> IO a) -> IO a
- Haskoin.Node: withNode :: (MonadLoggerIO m, MonadUnliftIO m) => NodeConfig -> (Node -> m a) -> m a
+ Haskoin.Node: withNode :: NodeConfig -> (Node -> IO a) -> IO a
- Haskoin.Node: withPeerMgr :: (MonadUnliftIO m, MonadLoggerIO m) => PeerMgrConfig -> (PeerMgr -> m a) -> m a
+ Haskoin.Node: withPeerMgr :: PeerMgrConfig -> (PeerMgr -> IO a) -> IO a
- Haskoin.Node: wrapPeer :: MonadIO m => PeerConfig -> TVar Bool -> Mailbox PeerMessage -> m Peer
+ Haskoin.Node: wrapPeer :: PeerConfig -> TVar Bool -> Mailbox PeerMessage -> IO Peer
Files
- CHANGELOG.md +6/−0
- haskoin-node.cabal +13/−12
- src/Haskoin/Node.hs +20/−67
- src/Haskoin/Node/Chain.hs +249/−435
- src/Haskoin/Node/Peer.hs +43/−128
- src/Haskoin/Node/PeerMgr.hs +266/−546
- test/Haskoin/NodeSpec.hs +15/−69
CHANGELOG.md view
@@ -4,6 +4,12 @@ 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). +## [1.3.0] - 2026-08-12++### Changed++- Simplify modules to make compatible with upstream changes.+ ## [1.2.0] - 2026-08-12 ### Changed
haskoin-node.cabal view
@@ -5,7 +5,7 @@ -- see: https://github.com/sol/hpack name: haskoin-node-version: 1.2.0+version: 1.3.0 synopsis: P2P library for Bitcoin and Bitcoin Cash description: Please see the README on GitHub at <https://github.com/jprupp/haskoin-node#readme> category: Bitcoin, Finance, Network@@ -35,7 +35,8 @@ hs-source-dirs: src build-depends:- base >=4.9 && <5+ async+ , base >=4.9 && <5 , bytestring , cereal , conduit@@ -44,19 +45,18 @@ , data-default , hashable , haskoin-core >=1.3.0- , monad-logger+ , logging , mtl , network , nqe >=0.6.3 , random- , resourcet- , rocksdb-haskell-jprupp >=2.2.0- , rocksdb-query >=0.5.0+ , rocksdb-haskell-jprupp >=2.3.0+ , rocksdb-query >=0.6.0+ , stm , string-conversions , text , time , transformers- , unliftio , unordered-containers default-language: Haskell2010 @@ -70,6 +70,7 @@ test build-depends: HUnit+ , async , base >=4.9 && <5 , base64 , bytestring@@ -82,20 +83,20 @@ , haskoin-core >=1.3.0 , haskoin-node , hspec- , monad-logger+ , logging , mtl , network , nqe >=0.6.3 , random- , resourcet- , rocksdb-haskell-jprupp >=2.2.0- , rocksdb-query >=0.5.0+ , rocksdb-haskell-jprupp >=2.3.0+ , rocksdb-query >=0.6.0 , safe+ , stm , string-conversions+ , temporary , text , time , transformers- , unliftio , unordered-containers default-language: Haskell2010 build-tool-depends: hspec-discover:hspec-discover
src/Haskoin/Node.hs view
@@ -4,6 +4,7 @@ {-# LANGUAGE LambdaCase #-} {-# LANGUAGE MultiParamTypeClasses #-} {-# LANGUAGE OverloadedRecordDot #-}+{-# LANGUAGE OverloadedStrings #-} {-# LANGUAGE RecordWildCards #-} {-# LANGUAGE NoFieldSelectors #-} @@ -19,56 +20,24 @@ ) where +import Control.Concurrent+import Control.Concurrent.Async+import Control.Exception+import Control.Logging import Control.Monad (forever) import Control.Monad.Cont (ContT (..), MonadCont (callCC), cont, runCont, runContT)-import Control.Monad.Logger (MonadLoggerIO) import Control.Monad.Trans (lift) 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@@ -108,17 +77,16 @@ withConnection :: SockAddr -> WithConnection withConnection na f = fromSockAddr na >>= \case- Nothing -> throwIO PeerAddressInvalid+ Nothing -> errorSL "Node" ("Peer address invalid: " <> cs (show na)) Just cset -> runTCPClient cset $ \ad -> f (Conduits (appSource ad) (appSink ad)) -fromSockAddr ::- (MonadUnliftIO m) => SockAddr -> m (Maybe ClientSettings)+fromSockAddr :: SockAddr -> IO (Maybe ClientSettings) fromSockAddr sa = go `catch` e where go = do- (maybe_host, maybe_port) <- liftIO (getNameInfo flags True True sa)+ (maybe_host, maybe_port) <- getNameInfo flags True True sa return $ clientSettings <$> (readMaybe =<< maybe_port)@@ -127,58 +95,43 @@ e :: (Monad m) => SomeException -> m (Maybe a) e _ = return Nothing -chainEvents ::- (MonadUnliftIO m, MonadLoggerIO m) =>- PeerMgr ->- Inbox ChainEvent ->- Publisher NodeEvent ->- m ()+chainEvents :: PeerMgr -> Inbox ChainEvent -> Publisher NodeEvent -> IO () chainEvents mgr input output = forever $ do event <- receive input case event of- ChainBestBlock bb ->- peerMgrBest bb.height mgr+ ChainBestBlock bb -> peerMgrBest mgr bb.height _ -> return () publish (ChainEvent event) output peerEvents ::- (MonadUnliftIO m, MonadLoggerIO m) =>- Chain ->- PeerMgr ->- Inbox PeerEvent ->- Publisher NodeEvent ->- m ()+ Chain -> PeerMgr -> Inbox PeerEvent -> Publisher NodeEvent -> IO () peerEvents ch mgr input output = forever $ do event <- receive input case event of PeerConnected p ->- chainPeerConnected p ch+ chainPeerConnected ch p PeerDisconnected p ->- chainPeerDisconnected p ch+ chainPeerDisconnected ch p PeerMessage p msg -> do case msg of MVersion v ->- peerMgrVersion p v mgr+ peerMgrVersion mgr p v MVerAck ->- peerMgrVerAck p mgr+ peerMgrVerAck mgr p MPing (Ping n) ->- peerMgrPing p n mgr+ peerMgrPing mgr p n MPong (Pong n) ->- peerMgrPong p n mgr+ peerMgrPong mgr p n MAddr (Addr ns) ->- peerMgrAddrs p (map snd ns) mgr+ peerMgrAddrs mgr p (map snd ns) MHeaders (Headers hs) ->- chainHeaders p (map fst hs) ch+ chainHeaders ch p (map fst hs) _ -> return () ticklePeer mgr p publish (PeerEvent event) output -- | Launch node process in the foreground.-withNode ::- (MonadLoggerIO m, MonadUnliftIO m) =>- NodeConfig ->- (Node -> m a) ->- m a+withNode :: NodeConfig -> (Node -> IO a) -> IO a withNode NodeConfig {..} action = flip runContT return $ do peerPub <- ContT withPublisher peerSub <- ContT (withSubscription peerPub)
src/Haskoin/Node/Chain.hs view
@@ -3,8 +3,11 @@ {-# LANGUAGE ExistentialQuantification #-} {-# LANGUAGE FlexibleContexts #-} {-# LANGUAGE FlexibleInstances #-}+{-# LANGUAGE ImportQualifiedPost #-} {-# LANGUAGE LambdaCase #-} {-# LANGUAGE MultiParamTypeClasses #-}+{-# LANGUAGE MultiWayIf #-}+{-# LANGUAGE NamedFieldPuns #-} {-# LANGUAGE OverloadedRecordDot #-} {-# LANGUAGE OverloadedStrings #-} {-# LANGUAGE RecordWildCards #-}@@ -31,109 +34,30 @@ ) where +import Control.Concurrent+import Control.Concurrent.Async+import Control.Concurrent.STM+import Control.Logging 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 Control.Monad.Trans.Reader+import Data.ByteString qualified 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 qualified 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 hiding (Key) 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@@ -162,20 +86,16 @@ 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+data Chain = Chain+ { config :: !ChainConfig,+ mailbox :: !(Mailbox ChainMessage), state :: !(TVar ChainState) } +instance Eq Chain where+ (==) = (==) `on` (.mailbox)+ -- | Database key for version. data ChainDataVersionKey = ChainDataVersionKey deriving (Eq, Ord, Show)@@ -230,58 +150,40 @@ return BestBlockKey put BestBlockKey = putWord8 0x91 -instance (MonadIO m) => BlockHeaders (ReaderT ChainConfig m) where- addBlockHeader bn = do- db <- asks (.db)- asks (.cf) >>= \case+type ChainM = ReaderT ChainConfig IO++chainM :: ChainConfig -> ChainM a -> IO a+chainM cfg m = runReaderT m cfg++instance BlockHeaders ChainM where+ addBlockHeader bn = ReaderT $ \ChainConfig {db, cf} -> do+ case cf of Nothing -> insert db (BlockHeaderKey h) bn- Just cf -> insertCF db cf (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+ getBlockHeader bh = ReaderT $ \ChainConfig {db, cf} -> do+ retrieveCommon db cf (BlockHeaderKey bh)+ getBestBlockHeader = ReaderT $ \ChainConfig {db, cf} -> do+ retrieveCommon db cf BestBlockKey >>= \case Nothing -> error "Could not get best block from database" Just b -> return b- setBestBlockHeader bn = do- db <- asks (.db)- asks (.cf) >>= \case+ setBestBlockHeader bn = ReaderT $ \ChainConfig {db, cf} -> do+ case cf of 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)+ Just cf' -> insertCF db cf' BestBlockKey bn+ addBlockHeaders bns = ReaderT $ \ChainConfig {db, cf} -> do+ writeBatch db (map (f cf) 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+ f cf bn = case cf of+ Nothing -> insertOp (BlockHeaderKey (h bn)) bn+ Just cf' -> insertOpCF cf' (BlockHeaderKey (h bn)) bn -withChain ::- (MonadUnliftIO m, MonadLoggerIO m) =>- ChainConfig ->- (Chain -> m a) ->- m a+withChain :: ChainConfig -> (Chain -> IO a) -> IO a withChain cfg action = do (inbox, mailbox) <- newMailbox- $(logDebugS) "Chain" "Starting chain actor"+ debugS "Chain" "Starting chain actor" st <- newTVarIO ChainState@@ -289,161 +191,137 @@ 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 ->+ let ch = Chain {config = cfg, mailbox = mailbox, state = st}+ initChainDB cfg+ withAsync (main_loop ch 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+ main_loop ch inbox =+ withSyncLoop ch.mailbox (run ch inbox)+ run ch inbox = do+ chainM cfg getBestBlockHeader+ >>= chainEvent cfg.pub . ChainBestBlock forever $ do- $(logDebugS) "Chain" "Awaiting event..."+ debugS "Chain" "Awaiting event..." msg <- receive inbox- chainMessage msg+ chainMessage ch msg -chainEvent :: (MonadChain m) => ChainEvent -> m ()-chainEvent e = do- pub <- asks (.config.pub)+chainEvent :: Publisher ChainEvent -> ChainEvent -> IO ()+chainEvent pub e = do case e of ChainBestBlock b ->- $(logInfoS) "Chain" $- "Best block header at height: "- <> cs (show b.height)+ logS "Chain" ("Best block header at height: " <> cs (show b.height)) ChainSynced b ->- $(logInfoS) "Chain" $- "Headers in sync at height: "- <> cs (show b.height)+ logS "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+processHeaders :: Chain -> Peer -> [BlockHeader] -> IO ()+processHeaders ch p hs = do+ let len = length hs+ debugS+ "Chain"+ ("Processing " <> cs (show len) <> " headers from peer: " <> p.label)+ let net = ch.config.net+ now <- getCurrentTime+ pbest <- chainM ch.config getBestBlockHeader+ importHeaders ch now hs >>= \case+ Nothing -> do+ warnS "Chain" ("Could not connect headers from peer: " <> p.label)+ killPeer p+ Just done -> do+ setLastReceived ch.state+ best <- chainM ch.config getBestBlockHeader when (pbest.header /= best.header) $- chainEvent (ChainBestBlock best)+ chainEvent ch.config.pub (ChainBestBlock best) if done then do MSendHeaders `sendMessage` p- finishPeer p- syncNewPeer- syncNotif- else syncPeer p+ finishPeer ch.state p+ syncNewPeer ch+ syncNotif ch+ else syncPeer ch p -syncNewPeer :: (MonadChain m) => m ()-syncNewPeer =- getSyncingPeer >>= \case+syncNewPeer :: Chain -> IO ()+syncNewPeer ch =+ getSyncingPeer ch.state >>= \case Just _ -> return () Nothing ->- nextPeer >>= \case+ nextPeer ch.state >>= \case Nothing -> return () Just p -> do- $(logDebugS) "Chain" $- "Syncing against peer: " <> p.label- syncPeer p+ debugS "Chain" ("Syncing against peer: " <> p.label)+ syncPeer ch p -syncNotif :: (MonadChain m) => m ()-syncNotif =- notifySynced >>= \case+syncNotif :: Chain -> IO ()+syncNotif ch =+ notifySynced ch >>= \case False -> return () True ->- withBlockHeaders getBestBlockHeader- >>= chainEvent . ChainSynced+ chainM ch.config getBestBlockHeader+ >>= chainEvent ch.config.pub . ChainSynced -syncPeer :: (MonadChain m) => Peer -> m ()-syncPeer p = do- t <- liftIO getCurrentTime+syncPeer :: Chain -> Peer -> IO ()+syncPeer ch p = do+ t <- getCurrentTime m <-- chainSyncingPeer >>= \case- Just- ChainSync- { peer = s,- best = m- }- | p == s -> syncing_me t m- | otherwise -> return Nothing+ chainSyncingPeer ch.state >>= \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+ debugS+ "Chain"+ ("Requesting headers from peer: " <> p.label) MGetHeaders g `sendMessage` p where syncing_new t =- setSyncingPeer p >>= \case+ setSyncingPeer ch.state p >>= \case False -> return Nothing True -> do- $(logDebugS) "Chain" $- "Locked peer: " <> p.label- h <- withBlockHeaders getBestBlockHeader- Just <$> syncHeaders t h p+ debugS "Chain" ("Locked peer: " <> p.label)+ h <- chainM ch.config getBestBlockHeader+ Just <$> syncHeaders ch t h p syncing_me t m = do h <- case m of- Nothing -> withBlockHeaders getBestBlockHeader+ Nothing -> chainM ch.config getBestBlockHeader Just h -> return h- Just <$> syncHeaders t h p+ Just <$> syncHeaders ch 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+chainMessage :: Chain -> ChainMessage -> IO ()+chainMessage ch (ChainHeaders p hs) =+ processHeaders ch p hs+chainMessage ch (ChainPeerConnected p) = do+ debugS "Chain" ("Peer connected: " <> p.label)+ addPeer ch.state p+ syncNewPeer ch+chainMessage ch (ChainPeerDisconnected p) = do+ debugS "Chain" ("Peer disconnected: " <> p.label)+ finishPeer ch.state p+ syncNewPeer ch+chainMessage ch ChainPing = do+ debugS "Chain" "Internal clock event"+ let to = ch.config.timeout+ now <- getCurrentTime+ chainSyncingPeer ch.state >>= \case Just ChainSync {peer = p, timestamp = t} | now `diffUTCTime` t > to -> do- $(logErrorS) "Chain" $- "Syncing peer timed out: " <> p.label- PeerTimeout `killPeer` p+ warnS+ "Chain"+ ("Syncing peer timed out: " <> p.label)+ killPeer p | otherwise -> return ()- Nothing -> syncNewPeer+ Nothing -> syncNewPeer ch -withSyncLoop ::- (MonadUnliftIO m, MonadLoggerIO m) =>- Chain ->- m a ->- m a-withSyncLoop ch f =+withSyncLoop :: Mailbox ChainMessage -> IO a -> IO a+withSyncLoop mbox mf = withAsync go $ \a ->- link a >> f+ link a >> mf where go = forever $ do- delay <-- liftIO $- randomRIO- ( 2 * 1000 * 1000,- 20 * 1000 * 1000- )+ delay <- randomRIO (2 * 10 ^ 6, 2 * 10 ^ 7) threadDelay delay- ChainPing `send` ch.mailbox+ ChainPing `send` mbox -- | Version of the database. dataVersion :: Word32@@ -451,73 +329,59 @@ -- | 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+initChainDB :: ChainConfig -> IO ()+initChainDB cfg@ChainConfig {db, cf, net} = do+ ver <- retrieveCommon db cf ChainDataVersionKey+ when (ver /= Just dataVersion) $ purgeChainDB cfg >>= writeBatch db+ case cf 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)+ Just cf' -> insertCF db cf' ChainDataVersionKey dataVersion+ retrieveCommon db cf BestBlockKey >>= \b ->+ when (isNothing (b :: Maybe BlockNode)) . chainM cfg $ 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+purgeChainDB :: ChainConfig -> IO [R.BatchOp]+purgeChainDB ChainConfig {db, cf} = do+ with_iter $ \it -> do+ R.iterSeek it (B.singleton 0x90)+ recurse_delete it where- f db Nothing = liftIO . R.withIter db- f db (Just cf) = liftIO . R.withIterCF db cf- recurse_delete it db mcf =- liftIO (R.iterKey it) >>= \case+ with_iter = case cf of+ Nothing -> R.withIter db+ Just cf' -> R.withIterCF db cf'+ recurse_delete it =+ R.iterKey it >>= \case Just k | B.head k == 0x90 || B.head k == 0x91 -> do- case mcf of- Nothing -> liftIO $ R.delete db k- Just cf -> liftIO $ R.deleteCF db cf k- liftIO $ R.iterNext it- (R.Del k :) <$> recurse_delete it db mcf+ case cf of+ Nothing -> R.delete db k+ Just cf' -> R.deleteCF db cf' k+ R.iterNext it+ (R.Del k :) <$> recurse_delete it _ -> 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+importHeaders :: Chain -> UTCTime -> [BlockHeader] -> IO (Maybe Bool)+importHeaders ch now hs =+ connect >>= \case+ Left _ -> return Nothing+ Right _+ | null hs -> return (Just False)+ | otherwise -> do+ bb <- get_last+ atomically . modifyTVar ch.state $ \s ->+ s {syncing = set_best bb <$> s.syncing}+ return (Just (length hs == 2000)) where+ set_best bb ChainSync {..} = ChainSync {best = bb, ..} timestamp = floor (utcTimeToPOSIXSeconds now)- connect = withBlockHeaders $ connectBlocks net timestamp hs- get_last = withBlockHeaders . getBlockHeader . headerHash $ last hs+ connect = chainM ch.config (connectBlocks ch.config.net timestamp hs)+ get_last = chainM ch.config $ 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@@ -526,49 +390,40 @@ -- 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 ()+notifySynced :: Chain -> IO Bool+notifySynced ch = do+ bb <- chainM ch.config getBestBlockHeader+ df <- (`diffUTCTime` block_time bb) <$> getCurrentTime+ atomically $ do+ s <- readTVar ch.state+ if+ | df > 7200 -> return False+ | isJust s.syncing -> return False+ | not (null s.peers) -> return False+ | s.beenInSync -> return False+ | otherwise -> do+ writeTVar ch.state s {beenInSync = True}+ return True 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+nextPeer :: TVar ChainState -> IO (Maybe Peer)+nextPeer st = fmap (.peers) (readTVarIO st) >>= go where go [] = return Nothing go (p : ps) =- setSyncingPeer p >>= \case+ setSyncingPeer st 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)+syncHeaders :: Chain -> UTCTime -> BlockNode -> Peer -> IO GetHeaders+syncHeaders ch now bb p = do atomically $- modifyTVar st $ \s ->+ modifyTVar ch.state $ \s -> s { syncing = Just@@ -579,7 +434,7 @@ }, peers = delete p s.peers }- loc <- withBlockHeaders $ blockLocator bb+ loc <- chainM ch.config (blockLocator bb) return GetHeaders { version = myVersion,@@ -590,43 +445,41 @@ 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)+setLastReceived :: TVar ChainState -> IO ()+setLastReceived st = do+ now <- getCurrentTime let f ChainSync {..} = ChainSync {timestamp = now, ..}- atomically . modifyTVar st $ \s ->- s {syncing = f <$> s.syncing}+ 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)}+addPeer :: TVar ChainState -> Peer -> IO ()+addPeer st p = do+ 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))+getSyncingPeer :: TVar ChainState -> IO (Maybe Peer)+getSyncingPeer st =+ readTVarIO st >>= \case+ ChainState {syncing = Just ChainSync {peer}} -> return (Just peer)+ _ -> return Nothing -setSyncingPeer :: (MonadChain m) => Peer -> m Bool-setSyncingPeer p =+setSyncingPeer :: TVar ChainState -> Peer -> IO Bool+setSyncingPeer st p = setBusy p >>= \case False -> do- $(logDebugS) "Chain" $- "Could not lock peer: " <> p.label+ debugS+ "Chain"+ ("Could not lock peer: " <> p.label) return False True -> do- $(logDebugS) "Chain" $- "Locked peer: " <> p.label+ debugS "Chain" $+ ("Locked peer: " <> p.label) set_it return True where set_it = do- now <- liftIO getCurrentTime- box <- asks (.state)- atomically $ modifyTVar box $ \s ->+ now <- getCurrentTime+ atomically $ modifyTVar st $ \s -> s { syncing = Just@@ -639,134 +492,95 @@ -- | 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+finishPeer :: TVar ChainState -> Peer -> IO ()+finishPeer st p =+ remove_peer >>= \case False ->- $(logDebugS) "Chain" $- "Removed peer from queue: " <> p.label+ debugS+ "Chain"+ ("Removed peer from queue: " <> p.label) True -> do- $(logDebugS) "Chain" $- "Releasing syncing peer: " <> p.label+ debugS+ "Chain"+ ("Releasing syncing peer: " <> p.label) setFree p where- remove_peer st =+ remove_peer = atomically $ readTVar st >>= \s -> case s.syncing of Just ChainSync {peer = p'} | p == p' -> do- unset_syncing st+ unset_syncing return True _ -> do- remove_from_queue st+ remove_from_queue return False- unset_syncing st =+ unset_syncing = modifyTVar st $ \x -> x {syncing = Nothing}- remove_from_queue st =+ remove_from_queue = modifyTVar st $ \x -> x {peers = delete p x.peers} -- | Return syncing peer data.-chainSyncingPeer :: (MonadChain m) => m (Maybe ChainSync)-chainSyncingPeer =- (.syncing) <$> (readTVarIO =<< asks (.state))+chainSyncingPeer :: TVar ChainState -> IO (Maybe ChainSync)+chainSyncingPeer st = (.syncing) <$> readTVarIO st --- | 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 a block header from the block chain.+chainGetBlock :: Chain -> BlockHash -> IO (Maybe BlockNode)+chainGetBlock ch bh = chainM ch.config (getBlockHeader bh) -- | Get best block header from chain process.-chainGetBest :: (MonadIO m) => Chain -> m BlockNode-chainGetBest ch =- runReaderT getBestBlockHeader ch.reader.config+chainGetBest :: Chain -> IO BlockNode+chainGetBest ch = chainM ch.config getBestBlockHeader -- | 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+chainGetAncestor :: Chain -> BlockHeight -> BlockNode -> IO (Maybe BlockNode)+chainGetAncestor ch h bn = chainM ch.config (getAncestor h bn) -- | Get parents of 'BlockNode' starting at 'BlockHeight' from chain process.-chainGetParents ::- (MonadIO m) =>- BlockHeight ->- BlockNode ->- Chain ->- m [BlockNode]-chainGetParents height top ch =+chainGetParents :: Chain -> BlockHeight -> BlockNode -> IO [BlockNode]+chainGetParents ch height top = go [] top where go acc b | height >= b.height = return acc | otherwise = do- m <- chainGetBlock b.header.prev ch+ m <- chainGetBlock ch b.header.prev 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+chainGetSplitBlock :: Chain -> BlockNode -> BlockNode -> IO BlockNode+chainGetSplitBlock ch l r = chainM ch.config (splitPoint l r) -- | Notify chain that a new peer is connected.-chainPeerConnected ::- (MonadIO m) =>- Peer ->- Chain ->- m ()-chainPeerConnected p ch =+chainPeerConnected :: Chain -> Peer -> IO ()+chainPeerConnected ch p = ChainPeerConnected p `send` ch.mailbox -- | Notify chain that a peer has disconnected.-chainPeerDisconnected ::- (MonadIO m) =>- Peer ->- Chain ->- m ()-chainPeerDisconnected p ch =+chainPeerDisconnected :: Chain -> Peer -> IO ()+chainPeerDisconnected ch p = ChainPeerDisconnected p `send` ch.mailbox -- | Is given 'BlockHash' in the main chain?-chainBlockMain ::- (MonadIO m) =>- BlockHash ->- Chain ->- m Bool-chainBlockMain bh ch =+chainBlockMain :: Chain -> BlockHash -> IO Bool+chainBlockMain ch bh = chainGetBest ch >>= \bb ->- chainGetBlock bh ch >>= \case+ chainGetBlock ch bh >>= \case Nothing -> return False bm@(Just bn) ->- (== bm) <$> chainGetAncestor bn.height bb ch+ (== bm) <$> chainGetAncestor ch bn.height bb -- | Is chain in sync with network?-chainIsSynced :: (MonadIO m) => Chain -> m Bool+chainIsSynced :: Chain -> IO Bool chainIsSynced ch =- (.beenInSync) <$> readTVarIO (ch.reader.state)+ (.beenInSync) <$> readTVarIO (ch.state) -- | Peer sends a bunch of headers to the chain process.-chainHeaders ::- (MonadIO m) =>- Peer ->- [BlockHeader] ->- Chain ->- m ()-chainHeaders p hs ch =+chainHeaders :: Chain -> Peer -> [BlockHeader] -> IO ()+chainHeaders ch p hs = ChainHeaders p hs `send` ch.mailbox
src/Haskoin/Node/Peer.hs view
@@ -34,20 +34,12 @@ where 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 Control.Concurrent.Async+import Control.Concurrent.STM+import Control.Exception+import Control.Logging+import Control.Monad+import Control.Monad.Trans.Maybe import Data.Bool (bool) import Data.ByteString (ByteString) import Data.ByteString qualified as B@@ -59,55 +51,9 @@ 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 System.Timeout data Conduits = Conduits { inboundConduit :: ConduitT () ByteString IO (),@@ -146,26 +92,6 @@ | 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"--instance Exception PeerException- -- | Mailbox for a peer. data Peer = Peer { mailbox :: !(Mailbox PeerMessage),@@ -182,15 +108,14 @@ -- | Incoming messages that a peer accepts. data PeerMessage- = KillPeer !PeerException+ = KillPeer | SendMessage !Message wrapPeer ::- (MonadIO m) => PeerConfig -> TVar Bool -> Mailbox PeerMessage ->- m Peer+ IO Peer wrapPeer cfg busy mbox = return Peer@@ -202,23 +127,21 @@ -- | Run peer process in current thread. peer ::- (MonadUnliftIO m, MonadLoggerIO m) => PeerConfig -> TVar Bool -> Inbox PeerMessage ->- m ()+ IO () peer cfg@PeerConfig {..} busy inbox = do p <- wrapPeer cfg busy (inboxToMailbox inbox)- withRunInIO $ \restore -> do- connect (restore . peer_session p)+ connect $ peer_session p where go = forever $ do- $(logDebugS) "Peer" $ label <> " awaiting event..."+ lift $ debugS "Peer" $ label <> " awaiting event..." msg <- receive inbox dispatchMessage cfg msg peer_session p ad = do- let ins = transPipe liftIO ad.inboundConduit- ons = transPipe liftIO ad.outboundConduit+ let ins = ad.inboundConduit+ ons = ad.outboundConduit src = runConduit $ ins@@ -232,50 +155,44 @@ -- | Internal function to dispatch peer messages. dispatchMessage ::- (MonadLoggerIO m) => PeerConfig -> PeerMessage ->- ConduitT i Message m ()+ ConduitT i Message IO () dispatchMessage PeerConfig {label} (SendMessage msg) = do- $(logDebugS) "Peer" $ label <> " sending: " <> cs (show msg)+ lift $ debugS "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+dispatchMessage PeerConfig {label} KillPeer = do+ lift $ errorSL "Peer" $ label <> " disconnecting" -- | Internal conduit to parse messages coming from peer. inPeerConduit ::- (MonadLoggerIO m) => Network -> PeerConfig -> Text ->- ConduitT ByteString Message m ()+ ConduitT ByteString Message IO () inPeerConduit net PeerConfig {label} a = forever $ do- $(logDebugS) "Peer" $ label <> " awaiting network message..."+ lift $ debugS "Peer" $ label <> " awaiting network message..." x <- takeCE 24 .| foldC when (B.null x) $ do- $(logErrorS) "Peer" $ label <> " empty header"- throwIO EmptyHeader+ lift $ errorSL "Peer" $ label <> " empty header" case decode x of Left e -> do- $(logErrorS) "Peer" $ label <> " error decoding header"- throwIO DecodeHeaderError+ lift $ errorSL "Peer" $ label <> " error decoding header" 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+ lift $ debugS "Peer" $ label <> " received: " <> cs (show cmd)+ when (len > 32 * 2 ^ (20 :: Int)) . lift . errorSL "Peer" $+ label <> " payload too large: " <> cs (show 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)+ lift $+ errorSL "Peer" $+ label+ <> " could not decode payload for cmd: "+ <> cs (show cmd) Right msg -> do- $(logDebugS) "Peer" $ label <> " forwarding: " <> cs (show msg)+ lift $ debugS "Peer" $ label <> " forwarding: " <> cs (show msg) yield msg -- | Outgoing peer conduit to serialize and send messages.@@ -283,36 +200,35 @@ 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` p.mailbox+killPeer :: Peer -> IO ()+killPeer p = KillPeer `send` p.mailbox -- | Send a network message to peer.-sendMessage :: (MonadIO m) => Message -> Peer -> m ()+sendMessage :: Message -> Peer -> IO () sendMessage msg p = SendMessage msg `send` p.mailbox -getBusy :: (MonadIO m) => Peer -> m Bool+getBusy :: Peer -> IO Bool getBusy p = readTVarIO p.busy -setBusy :: (MonadIO m) => Peer -> m Bool+setBusy :: Peer -> IO Bool setBusy p = atomically $ do b <- readTVar p.busy unless b $ writeTVar p.busy True return $ not b -setFree :: (MonadIO m) => Peer -> m ()+setFree :: Peer -> IO () 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])+ IO (Maybe [Block]) getBlocks net time p bhs = runMaybeT $ mapM f =<< MaybeT (getData time p (GetData ivs)) where@@ -327,12 +243,11 @@ -- 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])+ IO (Maybe [Tx]) getTxs net time p ths = runMaybeT $ mapM f =<< MaybeT (getData time p (GetData ivs)) where@@ -347,7 +262,7 @@ -- 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])+ Int -> Peer -> GetData -> IO (Maybe [Either Tx Block]) getData seconds p gd@(GetData ivs) = withSubscription p.pub $ \inb -> do r <- liftIO randomIO@@ -361,7 +276,7 @@ get_thing _inb _r acc [] = return $ reverse acc get_thing inb r acc hss@(InvVector t h : hs) =- filterReceive p inb >>= \case+ lift (filterReceive p inb) >>= \case MTx tx | is_tx t && (txHash tx).get == h -> get_thing inb r (Left tx : acc) hs@@ -388,7 +303,7 @@ -- | Ping a peer and await response. Return 'False' if response not received -- before timeout.-pingPeer :: (MonadUnliftIO m) => Int -> Peer -> m Bool+pingPeer :: Int -> Peer -> IO Bool pingPeer time p = fmap isJust . withSubscription p.pub $ \sub -> do r <- liftIO randomIO@@ -398,7 +313,7 @@ | p == p' && r == r' -> Just () _ -> Nothing -filterReceive :: (MonadIO m) => Peer -> Inbox PeerEvent -> m Message+filterReceive :: Peer -> Inbox PeerEvent -> IO Message filterReceive p inb = receive inb >>= \case PeerMessage p' msg | p == p' -> return msg
src/Haskoin/Node/PeerMgr.hs view
@@ -2,6 +2,7 @@ {-# LANGUAGE DuplicateRecordFields #-} {-# LANGUAGE FlexibleContexts #-} {-# LANGUAGE FlexibleInstances #-}+{-# LANGUAGE ImportQualifiedPost #-} {-# LANGUAGE LambdaCase #-} {-# LANGUAGE MultiParamTypeClasses #-} {-# LANGUAGE MultiWayIf #-}@@ -35,118 +36,29 @@ import Control.Applicative ((<|>)) import Control.Arrow+import Control.Concurrent+import Control.Concurrent.Async+import Control.Concurrent.STM+import Control.Exception+import Control.Logging 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, isNothing)+import Data.Maybe import Data.Set (Set)-import qualified Data.Set as Set+import Data.Set qualified 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],@@ -171,7 +83,7 @@ data PeerMgrMessage = Connect !SockAddr | CheckPeer !Peer- | PeerDied !Child !(Maybe SomeException)+ | PeerDied !Child | ManagerBest !BlockHeight | PeerVerAck !Peer | PeerVersion !Peer !Version@@ -205,364 +117,250 @@ f OnlinePeer {pings = pings} = fromMaybe 60 (median pings) withPeerMgr ::- (MonadUnliftIO m, MonadLoggerIO m) => PeerMgrConfig ->- (PeerMgr -> m a) ->- m a+ (PeerMgr -> IO a) ->+ IO a withPeerMgr cfg action = do- inbox <- newInbox- let mgr = inboxToMailbox inbox- withSupervisor (Notify (death mgr)) $ \sup -> do+ ibx <- newInbox+ withSupervisor (Notify (death ibx)) $ \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- }+ let mgr =+ PeerMgr+ { config = cfg,+ supervisor = sup,+ mailbox = inboxToMailbox ibx,+ best = bb,+ addresses = kp,+ peers = ob+ }+ go mgr ibx where- death mgr (a, ex) = PeerDied a ex `sendSTM` mgr- go inbox =- withAsync (peerManager inbox) $ \a ->- withConnectLoop $- link a >> ReaderT action+ death ibx (a, _e) = PeerDied a `sendSTM` ibx+ go mgr ibx = withAsync (peerManager mgr ibx) $ \a ->+ withConnectLoop mgr (link a >> action mgr) -peerManager ::- ( MonadUnliftIO m,- MonadManager m,- MonadLoggerIO m- ) =>- Inbox PeerMgrMessage ->- m ()-peerManager inb = do- $(logDebugS) "PeerMgr" "Awaiting best block"- putBestBlock <=< receiveMatch inb $ \case+peerManager :: PeerMgr -> Inbox PeerMgrMessage -> IO ()+peerManager mgr ibx = do+ debugS "PeerMgr" "Getting best block"+ putBestBlock mgr <=< receiveMatch ibx $ \case ManagerBest b -> Just b _ -> Nothing- $(logDebugS) "PeerMgr" "Starting peer manager actor"+ debugS "PeerMgr" "Starting peer manager actor" forever $ do- $(logDebugS) "PeerMgr" "Awaiting event..."- dispatch =<< receive inb+ debugS "PeerMgr" "Awaiting event..."+ dispatch mgr =<< receive ibx -putBestBlock :: (MonadManager m) => BlockHeight -> m ()-putBestBlock bb = do- b <- asks (.best)- atomically $ writeTVar b bb+putBestBlock :: PeerMgr -> BlockHeight -> IO ()+putBestBlock mgr bb = atomically $ writeTVar mgr.best bb -getBestBlock :: (MonadManager m) => m BlockHeight-getBestBlock =- asks (.best) >>= readTVarIO+getBestBlock :: PeerMgr -> IO BlockHeight+getBestBlock mgr = readTVarIO mgr.best -getNetwork :: (MonadManager m) => m Network-getNetwork =- asks (.config.net)+getNetwork :: PeerMgr -> Network+getNetwork mgr = mgr.config.net -loadPeers :: (MonadUnliftIO m, MonadManager m) => m ()-loadPeers = do- loadStaticPeers- loadNetSeeds+loadPeers :: PeerMgr -> IO ()+loadPeers mgr = do+ loadStaticPeers mgr+ loadNetSeeds mgr -loadStaticPeers :: (MonadUnliftIO m, MonadManager m) => m ()-loadStaticPeers = do- net <- asks (.config.net)- xs <- asks (.config.peers)- mapM_ newPeer . concat =<< mapM (toSockAddr net) xs+loadStaticPeers :: PeerMgr -> IO ()+loadStaticPeers mgr =+ mapM_ (newPeer mgr) . concat+ =<< mapM (toSockAddr mgr.config.net) mgr.config.peers -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+loadNetSeeds :: PeerMgr -> IO ()+loadNetSeeds mgr =+ when mgr.config.discover $ do+ ss <- concat <$> mapM (toSockAddr mgr.config.net) mgr.config.net.seeds+ mapM_ (newPeer mgr) 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)+logConnectedPeers :: PeerMgr -> IO ()+logConnectedPeers mgr = do+ let m = mgr.config.maxPeers+ l <- length <$> getConnectedPeers mgr+ logS "PeerMgr" $ "Peers connected: " <> cs (show l) <> "/" <> cs (show m) -getOnlinePeers :: (MonadManager m) => m [OnlinePeer]-getOnlinePeers =- asks (.peers) >>= readTVarIO+getOnlinePeers :: PeerMgr -> IO [OnlinePeer]+getOnlinePeers mgr = readTVarIO mgr.peers -getConnectedPeers :: (MonadManager m) => m [OnlinePeer]-getConnectedPeers =- filter (.online) <$> getOnlinePeers+getConnectedPeers :: PeerMgr -> IO [OnlinePeer]+getConnectedPeers mgr = filter (.online) <$> getOnlinePeers mgr -managerEvent :: (MonadManager m) => PeerEvent -> m ()-managerEvent e =- publish e =<< asks (.config.pub)+managerEvent :: PeerMgr -> PeerEvent -> IO ()+managerEvent mgr e = publish e mgr.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+dispatch :: PeerMgr -> PeerMgrMessage -> IO ()+dispatch mgr (PeerVersion p v) = do+ debugS+ "PeerMgr"+ ("Received peer " <> p.label <> " version: " <> cs (show v))+ let b = mgr.peers+ atomically (setPeerVersion b p v) >>= \case+ Just o -> do+ when o.online (announcePeer mgr p)+ debugS+ "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+ Nothing -> do+ warnS+ "PeerMgr"+ ("Version rejected for peer " <> p.label <> ": " <> cs (show v))+ killPeer p+dispatch mgr (PeerVerAck p) = do+ atomically (setPeerVerAck mgr.peers p) >>= \case Just o -> do- $(logDebugS) "PeerMgr" $- "Received version ack from peer: "- <> p.label- when o.online $- announcePeer p+ debugS "PeerMgr" ("Received version ack from peer: " <> p.label)+ when o.online (announcePeer mgr 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+ warnS+ "PeerMgr"+ ("Received verack from unknown peer: " <> p.label)+ killPeer p+dispatch mgr (PeerAddrs p nas)+ | mgr.config.discover = do+ let sas = map (hostToSockAddr . (.address)) nas+ len = length sas+ debugS+ "PeerMgr"+ ("Received " <> cs (show len) <> " addresses from peer " <> p.label)+ forM_ sas (newPeer mgr)+ | otherwise = debugS "PeerMgr" ("Ignoring addresses from peer " <> p.label)+dispatch mgr (PeerPong p n) = do+ debugS+ "PeerMgr"+ ("Received pong " <> cs (show n) <> " from: " <> p.label)+ now <- getCurrentTime+ atomically (gotPong mgr.peers n now p)+dispatch _mgr (PeerPing p n) = do+ debugS+ "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) "PeerMgr" $- "Housekeeping for peer " <> p.label- checkPeer p+dispatch mgr (ManagerBest h) = do+ debugS "PeerMgr" ("Setting best block to " <> cs (show h))+ putBestBlock mgr h+dispatch mgr (Connect sa) = do+ connectPeer mgr sa+dispatch mgr (PeerDied a) = do+ processPeerOffline mgr a+dispatch mgr (CheckPeer p) = do+ debugS "PeerMgr" ("Housekeeping for peer: " <> p.label)+ checkPeer mgr p -ticklePeer :: (MonadLoggerIO m) => PeerMgr -> Peer -> m ()+ticklePeer :: PeerMgr -> Peer -> IO () ticklePeer m p = do- $(logDebugS) "PeerMgr" $ "Tickle peer " <> p.label- t <- liftIO getCurrentTime- atomically $ modifyPeer m.peers p $ \o -> o {tickled = t}+ debugS "PeerMgr" ("Tickle peer " <> p.label)+ t <- getCurrentTime+ atomically (modifyPeer m.peers p (\o -> o {tickled = t})) -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 ()+checkPeer :: PeerMgr -> Peer -> IO ()+checkPeer mgr p =+ atomically (mgr.peers `findPeer` p) >>= \case Just o -> do- now <- liftIO getCurrentTime- maxLife <- asks (.config.maxPeerLife)- let expired = maxLife `addUTCTime` o.connected- timeout <- asks (.config.timeout)- let deadline = timeout `addUTCTime` o.tickled+ now <- getCurrentTime+ let expired = mgr.config.maxPeerLife `addUTCTime` o.connected+ let deadline = mgr.config.timeout `addUTCTime` o.tickled if | now > expired -> do- $(logErrorS) "PeerMgr" $- "Disconnecting old peer "- <> p.label- <> " online since "- <> cs (show o.connected)- killPeer PeerTooOld p+ warnS "PeerMgr" ("Killing old peer: " <> p.label)+ killPeer p | now > deadline -> do- $(logWarnS) "PeerMgr" $ "Peer timeout: " <> p.label- killPeer PeerTimeout p- | isNothing o.ping ->- sendPing p+ warnS "PeerMgr" ("Peer timeout: " <> p.label)+ killPeer p+ | isNothing o.ping -> sendPing mgr p | otherwise -> return ()+ _ -> return () -sendPing :: (MonadManager m, MonadLoggerIO m) => Peer -> m ()-sendPing p = do- b <- asks (.peers)- atomically (findPeer b p) >>= \case+sendPing :: PeerMgr -> Peer -> IO ()+sendPing mgr p = do+ atomically (mgr.peers `findPeer` p) >>= \case Nothing ->- $(logWarnS) "PeerMgr" $- "Will not ping unknown peer: " <> p.label+ warnS "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) "PeerMgr" $- "Sending ping "- <> cs (show n)- <> " to: "- <> p.label+ n <- randomIO+ now <- getCurrentTime+ atomically (setPeerPing mgr.peers n now p)+ debugS+ "PeerMgr"+ ("Sending ping " <> cs (show n) <> " to: " <> p.label) MPing (Ping n) `sendMessage` p- | otherwise -> return ()+ | otherwise ->+ debugS+ "PeerMgr"+ ("Will not ping offline peer: " <> p.label) -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+processPeerOffline :: PeerMgr -> Child -> IO ()+processPeerOffline mgr a = do+ atomically (findPeerAsync mgr.peers a) >>= \case+ Nothing -> warnS "PeerMgr" "Disconnected unknown peer" 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)+ warnS "PeerMgr" ("Disconnected peer: " <> o.mailbox.label)+ managerEvent mgr (PeerDisconnected o.mailbox)+ else+ warnS "PeerMgr" ("Could not connect to peer: " <> o.mailbox.label)+ atomically (removePeer mgr.peers o.mailbox)+ logConnectedPeers mgr -announcePeer :: (MonadManager m, MonadLoggerIO m) => Peer -> m ()-announcePeer p = do- b <- asks (.peers)- atomically (findPeer b p) >>= \case+announcePeer :: PeerMgr -> Peer -> IO ()+announcePeer mgr p = do+ atomically (findPeer mgr.peers p) >>= \case Just OnlinePeer {online = True} -> do- $(logInfoS) "PeerMgr" $- "Connected to peer " <> p.label- managerEvent $ PeerConnected p- logConnectedPeers+ logS+ "PeerMgr"+ ("Connected to peer " <> p.label)+ managerEvent mgr (PeerConnected p)+ logConnectedPeers mgr Just OnlinePeer {online = False} -> return () Nothing ->- $(logErrorS) "PeerMgr" $- "Not announcing disconnected peer: "- <> p.label+ warnS+ "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+getNewPeer :: PeerMgr -> IO (Maybe SockAddr)+getNewPeer mgr = do+ loadPeers mgr+ atomically . stateTVar mgr.addresses $+ first (listToMaybe . Set.elems) . Set.splitAt 1 -connectPeer ::- ( MonadUnliftIO m,- MonadManager m,- MonadLoggerIO m- ) =>- SockAddr ->- m ()-connectPeer sa = do- os <- asks (.peers)- atomically (findPeerAddress os sa) >>= \case+connectPeer :: PeerMgr -> SockAddr -> IO ()+connectPeer mgr sa = do+ atomically (findPeerAddress mgr.peers sa) >>= \case Just _ ->- $(logErrorS) "PeerMgr" $- "Attempted to connect to peer twice: " <> cs (show sa)+ warnS+ "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)+ logS "PeerMgr" ("Connecting to " <> cs (show sa))+ nonce <- randomIO+ bb <- getBestBlock mgr+ now <- getCurrentTime+ let rmt = NetworkAddress (srv mgr.config.net) (sockToHostAddress sa) unix = floor (utcTimeToPOSIXSeconds now)- ver = buildVersion net nonce bb ad rmt unix+ ver = buildVersion mgr.config.net nonce bb mgr.config.address rmt unix text = cs (show sa) (inbox, mailbox) <- newMailbox let pc = PeerConfig- { pub = pub,- net = net,+ { pub = mgr.config.pub,+ net = mgr.config.net, label = text,- connect = conn sa+ connect = mgr.config.connect sa } busy <- newTVarIO False p <- wrapPeer pc busy mailbox- a <- withRunInIO $ \io ->- sup `addChild` io (launch pc busy inbox p)+ a <- mgr.supervisor `addChild` launch pc busy inbox p MVersion ver `sendMessage` p- b <- asks (.peers) atomically $ insertPeer- b+ mgr.peers OnlinePeer { address = sa, verack = False,@@ -581,104 +379,75 @@ | net.segWit = 8 | otherwise = 0 launch pc busy inbox p =- ask >>= \mgr ->- withPeerLoop p mgr $ \a ->- link a >> peer pc busy inbox+ withPeerLoop mgr p (\a -> link a >> peer pc busy inbox) withPeerLoop ::- (MonadUnliftIO m, MonadLogger m) =>- Peer -> PeerMgr ->- (Async a -> m a) ->- m a-withPeerLoop p mgr =+ Peer ->+ (Async a -> IO a) ->+ IO a+withPeerLoop mgr p = withAsync . forever $ do let timeout = mgr.config.timeout ms = floor (timeout * 1000 * 1000)- r <- liftIO $ randomRIO (ms `div` 4, ms `div` 2)+ r <- randomRIO (ms `div` 4, ms `div` 2) threadDelay r- managerCheck p mgr+ managerCheck mgr p -withConnectLoop ::- (MonadUnliftIO m, MonadManager m) =>- m a ->- m a-withConnectLoop act =- withAsync go $ \a ->- link a >> act+withConnectLoop :: PeerMgr -> IO a -> IO a+withConnectLoop mgr 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+ l <- length <$> getOnlinePeers mgr+ when+ (l < mgr.config.maxPeers)+ (getNewPeer mgr >>= mapM_ (managerConnect mgr))+ threadDelay =<< randomRIO (10 ^ 5, 5 * 10 ^ 6) -newPeer :: (MonadIO m, MonadManager m) => SockAddr -> m ()-newPeer sa = do- b <- asks (.addresses)- o <- asks (.peers)+newPeer :: PeerMgr -> SockAddr -> IO ()+newPeer mgr sa = atomically $- findPeerAddress o sa >>= \case+ findPeerAddress mgr.peers sa >>= \case Just _ -> return ()- Nothing -> modifyTVar b $ Set.insert sa+ Nothing -> modifyTVar mgr.addresses $ 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- }+gotPong b nonce now p =+ findPeer b p >>= \case+ Just o@OnlinePeer {ping = Just (time, nonce')} | nonce' == nonce -> do+ let d = now `diffUTCTime` time+ o' = o {ping = Nothing, pings = sort (take 11 (d : o.pings))}+ insertPeer b o'+ _ -> return () 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+setPeerVersion :: TVar [OnlinePeer] -> Peer -> Version -> STM (Maybe OnlinePeer)+setPeerVersion b p v+ | v.services .&. nodeNetwork == 0 = return Nothing+ | otherwise =+ readTVar b >>= \ops ->+ if any (\o -> v.nonce == o.nonce) ops+ then return Nothing -- peer is myself+ else+ findPeer b p >>= \case+ Nothing -> return Nothing -- peer not found+ Just o -> do+ let n = o {version = Just v, online = o.verack}+ insertPeer b n+ return (Just 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+setPeerVerAck b p =+ findPeer b p >>= \case+ Just o -> do+ let o' = o {verack = True, online = isJust o.version}+ insertPeer b o'+ return (Just o')+ Nothing -> return Nothing findPeer :: TVar [OnlinePeer] -> Peer -> STM (Maybe OnlinePeer) findPeer b p =@@ -689,99 +458,50 @@ insertPeer b o = modifyTVar b $ \x -> sort . nub $ o : x -modifyPeer ::- TVar [OnlinePeer] ->- Peer ->- (OnlinePeer -> OnlinePeer) ->- STM ()+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))+removePeer b p = modifyTVar b (filter ((/= p) . (.mailbox))) -findPeerAsync ::- TVar [OnlinePeer] ->- Async () ->- STM (Maybe OnlinePeer)-findPeerAsync b a =- find ((== a) . (.async))- <$> readTVar b+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+findPeerAddress :: TVar [OnlinePeer] -> SockAddr -> STM (Maybe OnlinePeer)+findPeerAddress b a = find ((== a) . (.address)) <$> readTVar b -getPeers :: (MonadIO m) => PeerMgr -> m [OnlinePeer]-getPeers = runReaderT getConnectedPeers+getPeers :: PeerMgr -> IO [OnlinePeer]+getPeers = getConnectedPeers -getOnlinePeer ::- (MonadIO m) =>- Peer ->- PeerMgr ->- m (Maybe OnlinePeer)-getOnlinePeer p =- runReaderT $ asks (.peers) >>= atomically . (`findPeer` p)+getOnlinePeer :: PeerMgr -> Peer -> IO (Maybe OnlinePeer)+getOnlinePeer mgr p = atomically (mgr.peers `findPeer` p) -managerCheck :: (MonadIO m) => Peer -> PeerMgr -> m ()-managerCheck p mgr =- CheckPeer p `send` mgr.mailbox+managerCheck :: PeerMgr -> Peer -> IO ()+managerCheck mgr p = CheckPeer p `send` mgr.mailbox -managerConnect :: (MonadIO m) => SockAddr -> PeerMgr -> m ()-managerConnect sa mgr =- Connect sa `send` mgr.mailbox+managerConnect :: PeerMgr -> SockAddr -> IO ()+managerConnect mgr sa = Connect sa `send` mgr.mailbox -peerMgrBest :: (MonadIO m) => BlockHeight -> PeerMgr -> m ()-peerMgrBest bh mgr =- ManagerBest bh `send` mgr.mailbox+peerMgrBest :: PeerMgr -> BlockHeight -> IO ()+peerMgrBest mgr bh = ManagerBest bh `send` mgr.mailbox -peerMgrVerAck :: (MonadIO m) => Peer -> PeerMgr -> m ()-peerMgrVerAck p mgr =- PeerVerAck p `send` mgr.mailbox+peerMgrVerAck :: PeerMgr -> Peer -> IO ()+peerMgrVerAck mgr p = PeerVerAck p `send` mgr.mailbox -peerMgrVersion ::- (MonadIO m) =>- Peer ->- Version ->- PeerMgr ->- m ()-peerMgrVersion p ver mgr =- PeerVersion p ver `send` mgr.mailbox+peerMgrVersion :: PeerMgr -> Peer -> Version -> IO ()+peerMgrVersion mgr p ver = PeerVersion p ver `send` mgr.mailbox -peerMgrPing ::- (MonadIO m) =>- Peer ->- Word64 ->- PeerMgr ->- m ()-peerMgrPing p nonce mgr =- PeerPing p nonce `send` mgr.mailbox+peerMgrPing :: PeerMgr -> Peer -> Word64 -> IO ()+peerMgrPing mgr p nonce = PeerPing p nonce `send` mgr.mailbox -peerMgrPong ::- (MonadIO m) =>- Peer ->- Word64 ->- PeerMgr ->- m ()-peerMgrPong p nonce mgr =- PeerPong p nonce `send` mgr.mailbox+peerMgrPong :: PeerMgr -> Peer -> Word64 -> IO ()+peerMgrPong mgr p nonce = PeerPong p nonce `send` mgr.mailbox -peerMgrAddrs ::- (MonadIO m) =>- Peer ->- [NetworkAddress] ->- PeerMgr ->- m ()-peerMgrAddrs p addrs mgr =- PeerAddrs p addrs `send` mgr.mailbox+peerMgrAddrs :: PeerMgr -> Peer -> [NetworkAddress] -> IO ()+peerMgrAddrs mgr p addrs = PeerAddrs p addrs `send` mgr.mailbox toHostService :: String -> (Maybe String, Maybe String) toHostService str =@@ -801,17 +521,17 @@ (x : xs) | x == '[' -> do i <- elemIndex ']' xs- return $ second tail $ splitAt i xs+ return $ second (drop 1) $ splitAt i xs | x == ':' -> do return (str, "") _ -> Nothing in (host, srv) -toSockAddr :: (MonadUnliftIO m) => Network -> String -> m [SockAddr]+toSockAddr :: Network -> String -> IO [SockAddr] toSockAddr net str = go `catch` e where- go = fmap (map addrAddress) $ liftIO $ getAddrInfo Nothing host srv+ go = fmap (map addrAddress) $ getAddrInfo Nothing host srv (host, srv) = second (<|> Just (show net.defaultPort)) $ toHostService str
test/Haskoin/NodeSpec.hs view
@@ -1,5 +1,6 @@ {-# LANGUAGE DuplicateRecordFields #-} {-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE ImportQualifiedPost #-} {-# LANGUAGE LambdaCase #-} {-# LANGUAGE MultiParamTypeClasses #-} {-# LANGUAGE OverloadedRecordDot #-}@@ -10,21 +11,13 @@ module Haskoin.NodeSpec (spec) where import Conduit- ( awaitForever,- concatMapC,- foldC,- mapMC,- runConduit,- takeCE,- yield,- (.|),- )+import Control.Concurrent.Async+import Control.Logging 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 qualified as B import Data.ByteString.Base64 (decodeBase64Lenient) import Data.Default (def) import Data.Either (fromRight)@@ -32,58 +25,15 @@ 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 Database.RocksDB qualified 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.IO.Temp import System.Random (randomIO) import Test.Hspec import Test.Hspec.QuickCheck-import UnliftIO- ( MonadIO,- MonadUnliftIO,- liftIO,- throwString,- withAsync,- withSystemTempDirectory,- ) data TestNode = TestNode { testMgr :: PeerMgr,@@ -108,7 +58,7 @@ go :: Inbox ByteString -> Mailbox ByteString -> IO () go r s = do nonce <- randomIO- now <- round <$> liftIO getPOSIXTime+ now <- round <$> getPOSIXTime let rmt = NetworkAddress 0 (sockToHostAddress sa) ver = buildVersion net nonce 0 ad rmt now runPut (putMessage net (MVersion ver)) `send` s@@ -173,7 +123,7 @@ withTestNode net "connect-one-peer" $ \TestNode {..} -> do p <- waitForPeer nodeEvents Just OnlinePeer {version = Just Version {version = ver}} <-- getOnlinePeer p testMgr+ getOnlinePeer testMgr p ver `shouldSatisfy` (>= 70002) it "downloads some blocks" $ withTestNode net "get-blocks" $ \TestNode {..} -> do@@ -206,8 +156,8 @@ bb <- chainGetBest testChain bb.height `shouldSatisfy` (== 15) an <-- maybe (throwString "No ancestor found") return- =<< chainGetAncestor 10 bn testChain+ maybe (error "No ancestor found") return+ =<< chainGetAncestor testChain 10 bn headerHash bn.header `shouldBe` bh headerHash an.header `shouldBe` ah it "downloads some block parents" $@@ -223,7 +173,7 @@ ChainEvent (ChainBestBlock bn) -> Just bn _ -> Nothing bn.height `shouldBe` 15- ps <- chainGetParents 12 bn testChain+ ps <- chainGetParents testChain 12 bn length ps `shouldBe` 3 forM_ (zip ps hs) $ \(p, h) -> headerHash p.header `shouldBe` h@@ -234,13 +184,9 @@ PeerEvent (PeerConnected p) -> Just p _ -> Nothing -withTestNode ::- (MonadUnliftIO m) =>- Network ->- String ->- (TestNode -> m a) ->- m a-withTestNode net str f = runNoLoggingT $ flip runContT return $ do+withTestNode :: Network -> String -> (TestNode -> IO a) -> IO a+withTestNode net str f = withStderrLogging $ flip runContT return $ do+ lift $ setLogLevel LevelError w <- ContT $ withSystemTempDirectory ("haskoin-node-test-" <> str <> "-") pub <- ContT withPublisher sub <- ContT $ withSubscription pub@@ -268,7 +214,7 @@ connect = dummyPeerConnect net ad } Node mgr ch <- ContT $ withNode cfg'- lift . lift $+ lift $ f TestNode { testMgr = mgr,