packages feed

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