haskoin-node 1.3.1 → 1.4.0
raw patch · 7 files changed
+277/−253 lines, 7 filesdep +monad-loggerdep +unliftiodep −loggingdep −temporaryPVP ok
version bump matches the API change (PVP)
Dependencies added: monad-logger, unliftio
Dependencies removed: logging, temporary
API changes (from Hackage documentation)
- Haskoin.Node: chainBlockMain :: Chain -> BlockHash -> IO Bool
+ Haskoin.Node: chainBlockMain :: MonadIO m => Chain -> BlockHash -> m Bool
- Haskoin.Node: chainGetAncestor :: Chain -> BlockHeight -> BlockNode -> IO (Maybe BlockNode)
+ Haskoin.Node: chainGetAncestor :: MonadIO m => Chain -> BlockHeight -> BlockNode -> m (Maybe BlockNode)
- Haskoin.Node: chainGetBest :: Chain -> IO BlockNode
+ Haskoin.Node: chainGetBest :: MonadIO m => Chain -> m BlockNode
- Haskoin.Node: chainGetBlock :: Chain -> BlockHash -> IO (Maybe BlockNode)
+ Haskoin.Node: chainGetBlock :: MonadIO m => Chain -> BlockHash -> m (Maybe BlockNode)
- Haskoin.Node: chainGetParents :: Chain -> BlockHeight -> BlockNode -> IO [BlockNode]
+ Haskoin.Node: chainGetParents :: MonadIO m => Chain -> BlockHeight -> BlockNode -> m [BlockNode]
- Haskoin.Node: chainGetSplitBlock :: Chain -> BlockNode -> BlockNode -> IO BlockNode
+ Haskoin.Node: chainGetSplitBlock :: MonadIO m => Chain -> BlockNode -> BlockNode -> m BlockNode
- Haskoin.Node: chainHeaders :: Chain -> Peer -> [BlockHeader] -> IO ()
+ Haskoin.Node: chainHeaders :: MonadIO m => Chain -> Peer -> [BlockHeader] -> m ()
- Haskoin.Node: chainIsSynced :: Chain -> IO Bool
+ Haskoin.Node: chainIsSynced :: MonadIO m => Chain -> m Bool
- Haskoin.Node: chainPeerConnected :: Chain -> Peer -> IO ()
+ Haskoin.Node: chainPeerConnected :: MonadIO m => Chain -> Peer -> m ()
- Haskoin.Node: chainPeerDisconnected :: Chain -> Peer -> IO ()
+ Haskoin.Node: chainPeerDisconnected :: MonadIO m => Chain -> Peer -> m ()
- Haskoin.Node: getBlocks :: Network -> Int -> Peer -> [BlockHash] -> IO (Maybe [Block])
+ Haskoin.Node: getBlocks :: MonadUnliftIO m => Network -> Int -> Peer -> [BlockHash] -> m (Maybe [Block])
- Haskoin.Node: getBusy :: Peer -> IO Bool
+ Haskoin.Node: getBusy :: MonadIO m => Peer -> m Bool
- Haskoin.Node: getData :: Int -> Peer -> GetData -> IO (Maybe [Either Tx Block])
+ Haskoin.Node: getData :: MonadUnliftIO m => Int -> Peer -> GetData -> m (Maybe [Either Tx Block])
- Haskoin.Node: getOnlinePeer :: PeerMgr -> Peer -> IO (Maybe OnlinePeer)
+ Haskoin.Node: getOnlinePeer :: MonadIO m => PeerMgr -> Peer -> m (Maybe OnlinePeer)
- Haskoin.Node: getPeers :: PeerMgr -> IO [OnlinePeer]
+ Haskoin.Node: getPeers :: MonadIO m => PeerMgr -> m [OnlinePeer]
- Haskoin.Node: getTxs :: Network -> Int -> Peer -> [TxHash] -> IO (Maybe [Tx])
+ Haskoin.Node: getTxs :: MonadUnliftIO m => Network -> Int -> Peer -> [TxHash] -> m (Maybe [Tx])
- Haskoin.Node: killPeer :: Peer -> IO ()
+ Haskoin.Node: killPeer :: MonadIO m => Peer -> m ()
- Haskoin.Node: peer :: PeerConfig -> TVar Bool -> Inbox PeerMessage -> IO ()
+ Haskoin.Node: peer :: (MonadLoggerIO m, MonadUnliftIO m) => PeerConfig -> TVar Bool -> Inbox PeerMessage -> m ()
- Haskoin.Node: peerMgrAddrs :: PeerMgr -> Peer -> [NetworkAddress] -> IO ()
+ Haskoin.Node: peerMgrAddrs :: MonadIO m => PeerMgr -> Peer -> [NetworkAddress] -> m ()
- Haskoin.Node: peerMgrBest :: PeerMgr -> BlockHeight -> IO ()
+ Haskoin.Node: peerMgrBest :: MonadIO m => PeerMgr -> BlockHeight -> m ()
- Haskoin.Node: peerMgrPing :: PeerMgr -> Peer -> Word64 -> IO ()
+ Haskoin.Node: peerMgrPing :: MonadIO m => PeerMgr -> Peer -> Word64 -> m ()
- Haskoin.Node: peerMgrPong :: PeerMgr -> Peer -> Word64 -> IO ()
+ Haskoin.Node: peerMgrPong :: MonadIO m => PeerMgr -> Peer -> Word64 -> m ()
- Haskoin.Node: peerMgrVerAck :: PeerMgr -> Peer -> IO ()
+ Haskoin.Node: peerMgrVerAck :: MonadIO m => PeerMgr -> Peer -> m ()
- Haskoin.Node: peerMgrVersion :: PeerMgr -> Peer -> Version -> IO ()
+ Haskoin.Node: peerMgrVersion :: MonadIO m => PeerMgr -> Peer -> Version -> m ()
- Haskoin.Node: pingPeer :: Int -> Peer -> IO Bool
+ Haskoin.Node: pingPeer :: MonadUnliftIO m => Int -> Peer -> m Bool
- Haskoin.Node: sendMessage :: Message -> Peer -> IO ()
+ Haskoin.Node: sendMessage :: MonadIO m => Message -> Peer -> m ()
- Haskoin.Node: setBusy :: Peer -> IO Bool
+ Haskoin.Node: setBusy :: MonadIO m => Peer -> m Bool
- Haskoin.Node: setFree :: Peer -> IO ()
+ Haskoin.Node: setFree :: MonadIO m => Peer -> m ()
- Haskoin.Node: ticklePeer :: PeerMgr -> Peer -> IO ()
+ Haskoin.Node: ticklePeer :: MonadLoggerIO m => PeerMgr -> Peer -> m ()
- Haskoin.Node: withChain :: ChainConfig -> (Chain -> IO a) -> IO a
+ Haskoin.Node: withChain :: (MonadLoggerIO m, MonadUnliftIO m) => ChainConfig -> (Chain -> m a) -> m a
- Haskoin.Node: withNode :: NodeConfig -> (Node -> IO a) -> IO a
+ Haskoin.Node: withNode :: (MonadUnliftIO m, MonadLoggerIO m) => NodeConfig -> (Node -> m a) -> m a
- Haskoin.Node: withPeerMgr :: PeerMgrConfig -> (PeerMgr -> IO a) -> IO a
+ Haskoin.Node: withPeerMgr :: (MonadLoggerIO m, MonadUnliftIO m) => PeerMgrConfig -> (PeerMgr -> m a) -> m a
- Haskoin.Node: wrapPeer :: PeerConfig -> TVar Bool -> Mailbox PeerMessage -> IO Peer
+ Haskoin.Node: wrapPeer :: PeerConfig -> TVar Bool -> Mailbox PeerMessage -> Peer
Files
- CHANGELOG.md +6/−0
- haskoin-node.cabal +5/−4
- src/Haskoin/Node.hs +13/−12
- src/Haskoin/Node/Chain.hs +74/−65
- src/Haskoin/Node/Peer.hs +72/−67
- src/Haskoin/Node/PeerMgr.hs +92/−93
- test/Haskoin/NodeSpec.hs +15/−12
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.4.0] - 2026-08-13++### Changed++- Add monad transformers and monad-logger again because logging is too slow.+ ## [1.3.1] - 2026-08-12 ### Added
haskoin-node.cabal view
@@ -5,7 +5,7 @@ -- see: https://github.com/sol/hpack name: haskoin-node-version: 1.3.1+version: 1.4.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@@ -45,7 +45,7 @@ , data-default , hashable , haskoin-core >=1.3.0- , logging+ , monad-logger , mtl , network , nqe >=0.6.3@@ -57,6 +57,7 @@ , text , time , transformers+ , unliftio , unordered-containers default-language: Haskell2010 @@ -83,7 +84,7 @@ , haskoin-core >=1.3.0 , haskoin-node , hspec- , logging+ , monad-logger , mtl , network , nqe >=0.6.3@@ -93,10 +94,10 @@ , 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
@@ -20,12 +20,9 @@ ) 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 import Control.Monad.Trans (lift) import Data.Conduit.Network import Data.String.Conversions (cs)@@ -38,6 +35,7 @@ import NQE import Network.Socket import Text.Read (readMaybe)+import UnliftIO -- | General node configuration. data NodeConfig = NodeConfig@@ -77,16 +75,16 @@ withConnection :: SockAddr -> WithConnection withConnection na f = fromSockAddr na >>= \case- Nothing -> errorSL "Node" ("Peer address invalid: " <> cs (show na))+ Nothing -> error $ "Invalid address " ++ show na Just cset -> runTCPClient cset $ \ad -> f (Conduits (appSource ad) (appSink ad)) -fromSockAddr :: SockAddr -> IO (Maybe ClientSettings)+fromSockAddr :: (MonadUnliftIO m) => SockAddr -> m (Maybe ClientSettings) fromSockAddr sa = go `catch` e where go = do- (maybe_host, maybe_port) <- getNameInfo flags True True sa+ (maybe_host, maybe_port) <- liftIO $ getNameInfo flags True True sa return $ clientSettings <$> (readMaybe =<< maybe_port)@@ -95,7 +93,8 @@ e :: (Monad m) => SomeException -> m (Maybe a) e _ = return Nothing -chainEvents :: PeerMgr -> Inbox ChainEvent -> Publisher NodeEvent -> IO ()+chainEvents ::+ (MonadIO m) => PeerMgr -> Inbox ChainEvent -> Publisher NodeEvent -> m () chainEvents mgr input output = forever $ do event <- receive input case event of@@ -104,7 +103,8 @@ publish (ChainEvent event) output peerEvents ::- Chain -> PeerMgr -> Inbox PeerEvent -> Publisher NodeEvent -> IO ()+ (MonadLoggerIO m) =>+ Chain -> PeerMgr -> Inbox PeerEvent -> Publisher NodeEvent -> m () peerEvents ch mgr input output = forever $ do event <- receive input case event of@@ -131,7 +131,8 @@ publish (PeerEvent event) output -- | Launch node process in the foreground.-withNode :: NodeConfig -> (Node -> IO a) -> IO a+withNode ::+ (MonadUnliftIO m, MonadLoggerIO m) => NodeConfig -> (Node -> m a) -> m a withNode NodeConfig {..} action = flip runContT return $ do peerPub <- ContT withPublisher peerSub <- ContT (withSubscription peerPub)@@ -141,6 +142,6 @@ let chainCfg = ChainConfig {pub = chainPub, ..} chain <- ContT (withChain chainCfg) peerMgr <- ContT $ withPeerMgr peerMgrCfg- lift . link =<< ContT (withAsync $ chainEvents peerMgr chainSub pub)- lift . link =<< ContT (withAsync $ peerEvents chain peerMgr peerSub pub)+ link =<< ContT (withAsync $ chainEvents peerMgr chainSub pub)+ link =<< ContT (withAsync $ peerEvents chain peerMgr peerSub pub) lift $ action Node {..}
src/Haskoin/Node/Chain.hs view
@@ -36,12 +36,8 @@ ) 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.IO.Class (MonadIO, liftIO)+import Control.Monad.Logger import Control.Monad.Trans.Reader import Data.ByteString qualified as B import Data.Function (on)@@ -60,6 +56,8 @@ import Haskoin.Node.PeerMgr (myVersion) import NQE import System.Random (randomRIO)+import UnliftIO+import UnliftIO.Concurrent (threadDelay) -- | Configuration for chain syncing process. data ChainConfig = ChainConfig@@ -158,7 +156,7 @@ runChainT :: ChainT m a -> Chain -> m a runChainT = runReaderT -instance MonadIO m => BlockHeaders (ReaderT ChainConfig m) where+instance (MonadIO m) => BlockHeaders (ReaderT ChainConfig m) where addBlockHeader bn = ReaderT $ \ChainConfig {db, cf} -> liftIO $ do case cf of Nothing -> insert db (BlockHeaderKey h) bn@@ -183,17 +181,17 @@ Nothing -> insertOp (BlockHeaderKey (h bn)) bn Just cf' -> insertOpCF cf' (BlockHeaderKey (h bn)) bn -instance MonadIO m => BlockHeaders (ChainT m) where+instance (MonadIO m) => BlockHeaders (ChainT m) where addBlockHeader bn = withReaderT (.config) (addBlockHeader bn) getBlockHeader bh = withReaderT (.config) (getBlockHeader bh) getBestBlockHeader = withReaderT (.config) getBestBlockHeader setBestBlockHeader bn = withReaderT (.config) (setBestBlockHeader bn) addBlockHeaders bns = withReaderT (.config) (addBlockHeaders bns) -withChain :: ChainConfig -> (Chain -> IO a) -> IO a+withChain ::+ (MonadLoggerIO m, MonadUnliftIO m) => ChainConfig -> (Chain -> m a) -> m a withChain cfg action = do (inbox, mailbox) <- newMailbox- debugS "Chain" "Starting chain actor" st <- newTVarIO ChainState@@ -212,31 +210,37 @@ runReaderT getBestBlockHeader cfg >>= chainEvent cfg.pub . ChainBestBlock forever $ do- debugS "Chain" "Awaiting event..."+ $(logDebugS) "Chain" "Awaiting event..." msg <- receive inbox chainConfigMessage ch msg -chainEvent :: Publisher ChainEvent -> ChainEvent -> IO ()+chainEvent :: (MonadLoggerIO m) => Publisher ChainEvent -> ChainEvent -> m () chainEvent pub e = do case e of ChainBestBlock b ->- logS "Chain" ("Best block header at height: " <> cs (show b.height))+ $(logInfoS)+ "Chain"+ ("Best block header at height " <> cs (show b.height)) ChainSynced b ->- logS "Chain" ("Headers in sync at height: " <> cs (show b.height))+ $(logInfoS)+ "Chain"+ ("Headers in sync at height " <> cs (show b.height)) publish e pub -processHeaders :: Chain -> Peer -> [BlockHeader] -> IO ()+processHeaders :: (MonadLoggerIO m) => Chain -> Peer -> [BlockHeader] -> m () processHeaders ch p hs = do let len = length hs- debugS+ $(logDebugS) "Chain"- ("Processing " <> cs (show len) <> " headers from peer: " <> p.label)+ ("Processing " <> cs (show len) <> " headers from peer " <> p.label) let net = ch.config.net- now <- getCurrentTime+ now <- liftIO getCurrentTime pbest <- runReaderT getBestBlockHeader ch.config importHeaders ch now hs >>= \case Nothing -> do- warnS "Chain" ("Could not connect headers from peer: " <> p.label)+ $(logWarnS)+ "Chain"+ ("Could not connect headers from peer " <> p.label) killPeer p Just done -> do setLastReceived ch.state@@ -251,7 +255,7 @@ syncNotif ch else syncPeer ch p -syncNewPeer :: Chain -> IO ()+syncNewPeer :: (MonadLoggerIO m) => Chain -> m () syncNewPeer ch = getSyncingPeer ch.state >>= \case Just _ -> return ()@@ -259,10 +263,10 @@ nextPeer ch.state >>= \case Nothing -> return () Just p -> do- debugS "Chain" ("Syncing against peer: " <> p.label)+ $(logDebugS) "Chain" ("Syncing against peer " <> p.label) syncPeer ch p -syncNotif :: Chain -> IO ()+syncNotif :: (MonadLoggerIO m) => Chain -> m () syncNotif ch = notifySynced ch >>= \case False -> return ()@@ -270,9 +274,9 @@ runReaderT getBestBlockHeader ch.config >>= chainEvent ch.config.pub . ChainSynced -syncPeer :: Chain -> Peer -> IO ()+syncPeer :: (MonadLoggerIO m) => Chain -> Peer -> m () syncPeer ch p = do- t <- getCurrentTime+ t <- liftIO getCurrentTime m <- chainSyncingPeer ch.state >>= \case Just ChainSync {peer = s, best = m}@@ -280,16 +284,16 @@ | otherwise -> return Nothing Nothing -> syncing_new t forM_ m $ \g -> do- debugS+ $(logDebugS) "Chain"- ("Requesting headers from peer: " <> p.label)+ ("Requesting headers from peer " <> p.label) MGetHeaders g `sendMessage` p where syncing_new t = setSyncingPeer ch.state p >>= \case False -> return Nothing True -> do- debugS "Chain" ("Locked peer: " <> p.label)+ $(logDebugS) "Chain" ("Locked peer " <> p.label) h <- runReaderT getBestBlockHeader ch.config Just <$> syncHeaders ch t h p syncing_me t m = do@@ -298,32 +302,32 @@ Just h -> return h Just <$> syncHeaders ch t h p -chainConfigMessage :: Chain -> ChainMessage -> IO ()+chainConfigMessage :: (MonadLoggerIO m) => Chain -> ChainMessage -> m () chainConfigMessage ch (ChainHeaders p hs) = processHeaders ch p hs chainConfigMessage ch (ChainPeerConnected p) = do- debugS "Chain" ("Peer connected: " <> p.label)+ $(logDebugS) "Chain" ("Connected peer " <> p.label) addPeer ch.state p syncNewPeer ch chainConfigMessage ch (ChainPeerDisconnected p) = do- debugS "Chain" ("Peer disconnected: " <> p.label)+ $(logDebugS) "Chain" ("Disconnected peer " <> p.label) finishPeer ch.state p syncNewPeer ch chainConfigMessage ch ChainPing = do- debugS "Chain" "Internal clock event"+ $(logDebugS) "Chain" "Internal clock event" let to = ch.config.timeout- now <- getCurrentTime+ now <- liftIO getCurrentTime chainSyncingPeer ch.state >>= \case Just ChainSync {peer = p, timestamp = t} | now `diffUTCTime` t > to -> do- warnS+ $(logWarnS) "Chain"- ("Syncing peer timed out: " <> p.label)+ ("Syncing peer " <> p.label <> " timed out") killPeer p | otherwise -> return () Nothing -> syncNewPeer ch -withSyncLoop :: Mailbox ChainMessage -> IO a -> IO a+withSyncLoop :: (MonadUnliftIO m) => Mailbox ChainMessage -> m a -> m a withSyncLoop mbox mf = withAsync go $ \a -> link a >> mf@@ -339,8 +343,8 @@ -- | Initialize header database. If version is different from current, the -- database is purged of conflicting elements first.-initChainDB :: ChainConfig -> IO ()-initChainDB cfg@ChainConfig {db, cf, net} = do+initChainDB :: (MonadIO m) => ChainConfig -> m ()+initChainDB cfg@ChainConfig {db, cf, net} = liftIO $ do ver <- retrieveCommon db cf ChainDataVersionKey when (ver /= Just dataVersion) $ purgeChainDB cfg >>= writeBatch db case cf of@@ -353,8 +357,8 @@ -- | Purge database of elements having keys that may conflict with those used in -- this module.-purgeChainDB :: ChainConfig -> IO [R.BatchOp]-purgeChainDB ChainConfig {db, cf} = do+purgeChainDB :: (MonadIO m) => ChainConfig -> m [R.BatchOp]+purgeChainDB ChainConfig {db, cf} = liftIO $ do with_iter $ \it -> do R.iterSeek it (B.singleton 0x90) recurse_delete it@@ -376,7 +380,8 @@ -- | 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 :: Chain -> UTCTime -> [BlockHeader] -> IO (Maybe Bool)+importHeaders ::+ (MonadIO m) => Chain -> UTCTime -> [BlockHeader] -> m (Maybe Bool) importHeaders ch now hs = connect >>= \case Left _ -> return Nothing@@ -400,10 +405,10 @@ -- 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 :: Chain -> IO Bool+notifySynced :: (MonadIO m) => Chain -> m Bool notifySynced ch = do bb <- runReaderT getBestBlockHeader ch.config- df <- (`diffUTCTime` block_time bb) <$> getCurrentTime+ df <- (`diffUTCTime` block_time bb) <$> liftIO getCurrentTime atomically $ do s <- readTVar ch.state if@@ -419,7 +424,7 @@ posixSecondsToUTCTime . fromIntegral . (.header.timestamp) -- | Get next peer to sync against from the queue.-nextPeer :: TVar ChainState -> IO (Maybe Peer)+nextPeer :: (MonadLoggerIO m) => TVar ChainState -> m (Maybe Peer) nextPeer st = fmap (.peers) (readTVarIO st) >>= go where go [] = return Nothing@@ -430,7 +435,8 @@ -- | Set a syncing peer and generate a 'GetHeaders' data structure with a block -- locator to send to that peer for syncing.-syncHeaders :: Chain -> UTCTime -> BlockNode -> Peer -> IO GetHeaders+syncHeaders ::+ (MonadIO m) => Chain -> UTCTime -> BlockNode -> Peer -> m GetHeaders syncHeaders ch now bb p = do atomically $ modifyTVar ch.state $ \s ->@@ -455,40 +461,40 @@ z = "0000000000000000000000000000000000000000000000000000000000000000" -- | Set the time of last received data to now if a syncing peer is active.-setLastReceived :: TVar ChainState -> IO ()+setLastReceived :: (MonadIO m) => TVar ChainState -> m () setLastReceived st = do- now <- getCurrentTime+ now <- liftIO getCurrentTime let f ChainSync {..} = ChainSync {timestamp = now, ..} atomically (modifyTVar st (\s -> s {syncing = f <$> s.syncing})) -- | Add a new peer to the queue of peers to sync against.-addPeer :: TVar ChainState -> Peer -> IO ()+addPeer :: (MonadIO m) => TVar ChainState -> Peer -> m () addPeer st p = do atomically (modifyTVar st (\s -> s {peers = nub (p : s.peers)})) -- | Get syncing peer if there is one.-getSyncingPeer :: TVar ChainState -> IO (Maybe Peer)+getSyncingPeer :: (MonadIO m) => TVar ChainState -> m (Maybe Peer) getSyncingPeer st = readTVarIO st >>= \case ChainState {syncing = Just ChainSync {peer}} -> return (Just peer) _ -> return Nothing -setSyncingPeer :: TVar ChainState -> Peer -> IO Bool+setSyncingPeer :: (MonadLoggerIO m) => TVar ChainState -> Peer -> m Bool setSyncingPeer st p = setBusy p >>= \case False -> do- debugS+ $(logDebugS) "Chain" ("Could not lock peer: " <> p.label) return False True -> do- debugS "Chain" $+ $(logDebugS) "Chain" $ ("Locked peer: " <> p.label) set_it return True where set_it = do- now <- getCurrentTime+ now <- liftIO getCurrentTime atomically $ modifyTVar st $ \s -> s { syncing =@@ -502,15 +508,15 @@ -- | Remove a peer from the queue of peers to sync and unset the syncing peer if -- it is set to the provided peer.-finishPeer :: TVar ChainState -> Peer -> IO ()+finishPeer :: (MonadLoggerIO m) => TVar ChainState -> Peer -> m () finishPeer st p = remove_peer >>= \case False ->- debugS+ $(logDebugS) "Chain" ("Removed peer from queue: " <> p.label) True -> do- debugS+ $(logDebugS) "Chain" ("Releasing syncing peer: " <> p.label) setFree p@@ -533,23 +539,25 @@ x {peers = delete p x.peers} -- | Return syncing peer data.-chainSyncingPeer :: TVar ChainState -> IO (Maybe ChainSync)+chainSyncingPeer :: (MonadIO m) => TVar ChainState -> m (Maybe ChainSync) chainSyncingPeer st = (.syncing) <$> readTVarIO st -- | Get a block header from the block chain.-chainGetBlock :: Chain -> BlockHash -> IO (Maybe BlockNode)+chainGetBlock :: (MonadIO m) => Chain -> BlockHash -> m (Maybe BlockNode) chainGetBlock ch bh = runReaderT (getBlockHeader bh) ch.config -- | Get best block header from chain process.-chainGetBest :: Chain -> IO BlockNode+chainGetBest :: (MonadIO m) => Chain -> m BlockNode chainGetBest ch = runReaderT getBestBlockHeader ch.config -- | Get ancestor of 'BlockNode' at 'BlockHeight' from chain process.-chainGetAncestor :: Chain -> BlockHeight -> BlockNode -> IO (Maybe BlockNode)+chainGetAncestor ::+ (MonadIO m) => Chain -> BlockHeight -> BlockNode -> m (Maybe BlockNode) chainGetAncestor ch h bn = runReaderT (getAncestor h bn) ch.config -- | Get parents of 'BlockNode' starting at 'BlockHeight' from chain process.-chainGetParents :: Chain -> BlockHeight -> BlockNode -> IO [BlockNode]+chainGetParents ::+ (MonadIO m) => Chain -> BlockHeight -> BlockNode -> m [BlockNode] chainGetParents ch height top = go [] top where@@ -562,21 +570,22 @@ Just p -> go (p : acc) p -- | Get last common block from chain process.-chainGetSplitBlock :: Chain -> BlockNode -> BlockNode -> IO BlockNode+chainGetSplitBlock ::+ (MonadIO m) => Chain -> BlockNode -> BlockNode -> m BlockNode chainGetSplitBlock ch l r = runReaderT (splitPoint l r) ch.config -- | Notify chain that a new peer is connected.-chainPeerConnected :: Chain -> Peer -> IO ()+chainPeerConnected :: (MonadIO m) => Chain -> Peer -> m () chainPeerConnected ch p = ChainPeerConnected p `send` ch.mailbox -- | Notify chain that a peer has disconnected.-chainPeerDisconnected :: Chain -> Peer -> IO ()+chainPeerDisconnected :: (MonadIO m) => Chain -> Peer -> m () chainPeerDisconnected ch p = ChainPeerDisconnected p `send` ch.mailbox -- | Is given 'BlockHash' in the main chain?-chainBlockMain :: Chain -> BlockHash -> IO Bool+chainBlockMain :: (MonadIO m) => Chain -> BlockHash -> m Bool chainBlockMain ch bh = chainGetBest ch >>= \bb -> chainGetBlock ch bh >>= \case@@ -586,11 +595,11 @@ (== bm) <$> chainGetAncestor ch bn.height bb -- | Is chain in sync with network?-chainIsSynced :: Chain -> IO Bool+chainIsSynced :: (MonadIO m) => Chain -> m Bool chainIsSynced ch = (.beenInSync) <$> readTVarIO (ch.state) -- | Peer sends a bunch of headers to the chain process.-chainHeaders :: Chain -> Peer -> [BlockHeader] -> IO ()+chainHeaders :: (MonadIO m) => Chain -> Peer -> [BlockHeader] -> m () chainHeaders ch p hs = ChainHeaders p hs `send` ch.mailbox
src/Haskoin/Node/Peer.hs view
@@ -34,11 +34,9 @@ where import Conduit-import Control.Concurrent.Async-import Control.Concurrent.STM import Control.Exception-import Control.Logging import Control.Monad+import Control.Monad.Logger import Control.Monad.Trans.Maybe import Data.Bool (bool) import Data.ByteString (ByteString)@@ -52,8 +50,8 @@ import Data.Word (Word32) import Haskoin import NQE-import System.Random (randomIO)-import System.Timeout+import System.Random+import UnliftIO data Conduits = Conduits { inboundConduit :: ConduitT () ByteString IO (),@@ -115,37 +113,37 @@ PeerConfig -> TVar Bool -> Mailbox PeerMessage ->- IO Peer+ Peer wrapPeer cfg busy mbox =- return- Peer- { mailbox = mbox,- pub = cfg.pub,- label = cfg.label,- busy = busy- }+ Peer+ { mailbox = mbox,+ pub = cfg.pub,+ label = cfg.label,+ busy = busy+ } -- | Run peer process in current thread. peer ::+ (MonadLoggerIO m, MonadUnliftIO m) => PeerConfig -> TVar Bool -> Inbox PeerMessage ->- IO ()-peer cfg@PeerConfig {..} busy inbox = do- p <- wrapPeer cfg busy (inboxToMailbox inbox)- connect $ peer_session p+ m ()+peer cfg@PeerConfig {label, net, connect, pub} busy inbox = do+ let p = wrapPeer cfg busy (inboxToMailbox inbox)+ withRunInIO $ \run -> connect (run . peer_session p) where- go = forever $ do- lift $ debugS "Peer" $ label <> " awaiting event..."+ go = do+ $(logDebugS) "Peer" $ label <> " awaiting event..." msg <- receive inbox- dispatchMessage cfg msg+ dispatchMessage cfg msg >>= bool (return ()) go peer_session p ad = do- let ins = ad.inboundConduit- ons = ad.outboundConduit+ let ins = transPipe liftIO ad.inboundConduit+ ons = transPipe liftIO ad.outboundConduit src = runConduit $ ins- .| inPeerConduit net cfg label+ .| inPeerConduit cfg .| mapM_C (send_msg p) snk = outPeerConduit net .| ons withAsync src $ \as -> do@@ -155,80 +153,86 @@ -- | Internal function to dispatch peer messages. dispatchMessage ::+ (MonadLoggerIO m) => PeerConfig -> PeerMessage ->- ConduitT i Message IO ()+ ConduitT i Message m Bool dispatchMessage PeerConfig {label} (SendMessage msg) = do- lift $ debugS "Peer" $ label <> " sending: " <> cs (show msg)+ $(logDebugS) "Peer" (label <> " sending: " <> cs (show msg)) yield msg+ return True dispatchMessage PeerConfig {label} KillPeer = do- lift $ errorSL "Peer" $ label <> " disconnecting"+ $(logInfoS) "Peer" (label <> " disconnecting")+ return False -- | Internal conduit to parse messages coming from peer.-inPeerConduit ::- Network ->- PeerConfig ->- Text ->- ConduitT ByteString Message IO ()-inPeerConduit net PeerConfig {label} a =- forever $ do- lift $ debugS "Peer" $ label <> " awaiting network message..."- x <- takeCE 24 .| foldC- when (B.null x) $ do- lift $ errorSL "Peer" $ label <> " empty header"- case decode x of- Left e -> do- lift $ errorSL "Peer" $ label <> " error decoding header"- Right (MessageHeader _ cmd len _) -> do- 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- lift $- errorSL "Peer" $- label- <> " could not decode payload for cmd: "- <> cs (show cmd)- Right msg -> do- lift $ debugS "Peer" $ label <> " forwarding: " <> cs (show msg)- yield msg+inPeerConduit :: (MonadLoggerIO m) => PeerConfig -> ConduitT ByteString Message m ()+inPeerConduit pc@PeerConfig {label, net} = do+ $(logDebugS) "Peer" (label <> ": awaiting message...")+ x <- takeCE 24 .| foldC+ when (B.null x) $ do+ $(logWarnS) "Peer" (label <> " sent empty header")+ case decode x of+ Left e -> do+ $(logWarnS) "Peer" (label <> " sent invalid header")+ Right (MessageHeader _ cmd len _)+ | len > 32 * 2 ^ (20 :: Int) ->+ $(logWarnS) "Peer" $+ label+ <> " wants to send "+ <> cs (show len)+ <> " bytes (too large) for cmd "+ <> cs (show cmd)+ | otherwise -> do+ $(logDebugS) "Peer" (label <> " sent cmd " <> cs (show cmd))+ y <- takeCE (fromIntegral len) .| foldC+ case runGet (getMessage net) $ x `B.append` y of+ Left e ->+ $(logErrorS)+ "Peer"+ (label <> ": sent invalid payload for cmd " <> cs (show cmd))+ Right msg -> do+ $(logDebugS)+ "Peer"+ (label <> " sent full message for cmd " <> cs (show cmd))+ yield msg+ inPeerConduit pc -- | Outgoing peer conduit to serialize and send messages. outPeerConduit :: (Monad m) => Network -> ConduitT Message ByteString m () outPeerConduit net = awaitForever $ yield . runPut . putMessage net -- | Kill a peer with the provided exception.-killPeer :: Peer -> IO ()+killPeer :: (MonadIO m) => Peer -> m () killPeer p = KillPeer `send` p.mailbox -- | Send a network message to peer.-sendMessage :: Message -> Peer -> IO ()+sendMessage :: (MonadIO m) => Message -> Peer -> m () sendMessage msg p = SendMessage msg `send` p.mailbox -getBusy :: Peer -> IO Bool+getBusy :: (MonadIO m) => Peer -> m Bool getBusy p = readTVarIO p.busy -setBusy :: Peer -> IO Bool+setBusy :: (MonadIO m) => Peer -> m Bool setBusy p = atomically $ do b <- readTVar p.busy unless b $ writeTVar p.busy True return $ not b -setFree :: Peer -> IO ()+setFree :: (MonadIO m) => Peer -> m () setFree p = atomically $ writeTVar p.busy False -- | Request full blocks from peer. Will return 'Nothing' if the list of blocks -- returned by the peer is incomplete, comes out of order, or a timeout is -- reached. getBlocks ::+ (MonadUnliftIO m) => Network -> Int -> Peer -> [BlockHash] ->- IO (Maybe [Block])+ m (Maybe [Block]) getBlocks net time p bhs = runMaybeT $ mapM f =<< MaybeT (getData time p (GetData ivs)) where@@ -243,11 +247,12 @@ -- transactions returned by the peer is incomplete, comes out of order, or a -- timeout is reached. getTxs ::+ (MonadUnliftIO m) => Network -> Int -> Peer -> [TxHash] ->- IO (Maybe [Tx])+ m (Maybe [Tx]) getTxs net time p ths = runMaybeT $ mapM f =<< MaybeT (getData time p (GetData ivs)) where@@ -262,10 +267,10 @@ -- single inventory fails to be retrieved, if they come out of order, or if -- timeout is reached. getData ::- Int -> Peer -> GetData -> IO (Maybe [Either Tx Block])+ (MonadUnliftIO m) => Int -> Peer -> GetData -> m (Maybe [Either Tx Block]) getData seconds p gd@(GetData ivs) = withSubscription p.pub $ \inb -> do- r <- liftIO randomIO+ r <- randomIO MGetData gd `sendMessage` p MPing (Ping r) `sendMessage` p fmap join@@ -303,17 +308,17 @@ -- | Ping a peer and await response. Return 'False' if response not received -- before timeout.-pingPeer :: Int -> Peer -> IO Bool+pingPeer :: (MonadUnliftIO m) => Int -> Peer -> m Bool pingPeer time p = fmap isJust . withSubscription p.pub $ \sub -> do- r <- liftIO randomIO+ r <- randomIO MPing (Ping r) `sendMessage` p receiveMatchS time sub $ \case PeerMessage p' (MPong (Pong r')) | p == p' && r == r' -> Just () _ -> Nothing -filterReceive :: Peer -> Inbox PeerEvent -> IO Message+filterReceive :: (MonadIO m) => Peer -> Inbox PeerEvent -> m Message filterReceive p inb = receive inb >>= \case PeerMessage p' msg | p == p' -> return msg
src/Haskoin/Node/PeerMgr.hs view
@@ -36,13 +36,8 @@ 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-import Control.Monad.Except+import Control.Monad.Logger import Data.Bits ((.&.)) import Data.Function (on) import Data.List (dropWhileEnd, elemIndex, find, nub, sort)@@ -58,6 +53,8 @@ import NQE import Network.Socket import System.Random (randomIO, randomRIO)+import UnliftIO+import UnliftIO.Concurrent data PeerMgrConfig = PeerMgrConfig { maxPeers :: !Int,@@ -117,9 +114,10 @@ f OnlinePeer {pings = pings} = fromMaybe 60 (median pings) withPeerMgr ::+ (MonadLoggerIO m, MonadUnliftIO m) => PeerMgrConfig ->- (PeerMgr -> IO a) ->- IO a+ (PeerMgr -> m a) ->+ m a withPeerMgr cfg action = do ibx <- newInbox withSupervisor (Notify (death ibx)) $ \sup -> do@@ -141,207 +139,211 @@ go mgr ibx = withAsync (peerManager mgr ibx) $ \a -> withConnectLoop mgr (link a >> action mgr) -peerManager :: PeerMgr -> Inbox PeerMgrMessage -> IO ()+peerManager ::+ (MonadLoggerIO m, MonadUnliftIO m) => PeerMgr -> Inbox PeerMgrMessage -> m () peerManager mgr ibx = do- debugS "PeerMgr" "Getting best block"+ $(logDebugS) "PeerMgr" "Getting best block..." putBestBlock mgr <=< receiveMatch ibx $ \case ManagerBest b -> Just b _ -> Nothing- debugS "PeerMgr" "Starting peer manager actor" forever $ do- debugS "PeerMgr" "Awaiting event..."+ $(logDebugS) "PeerMgr" "Awaiting event..." dispatch mgr =<< receive ibx -putBestBlock :: PeerMgr -> BlockHeight -> IO ()+putBestBlock :: (MonadIO m) => PeerMgr -> BlockHeight -> m () putBestBlock mgr bb = atomically $ writeTVar mgr.best bb -getBestBlock :: PeerMgr -> IO BlockHeight+getBestBlock :: (MonadIO m) => PeerMgr -> m BlockHeight getBestBlock mgr = readTVarIO mgr.best getNetwork :: PeerMgr -> Network getNetwork mgr = mgr.config.net -loadPeers :: PeerMgr -> IO ()+loadPeers :: (MonadIO m) => PeerMgr -> m () loadPeers mgr = do loadStaticPeers mgr loadNetSeeds mgr -loadStaticPeers :: PeerMgr -> IO ()+loadStaticPeers :: (MonadIO m) => PeerMgr -> m () loadStaticPeers mgr = mapM_ (newPeer mgr) . concat- =<< mapM (toSockAddr mgr.config.net) mgr.config.peers+ =<< mapM (liftIO . toSockAddr mgr.config.net) mgr.config.peers -loadNetSeeds :: PeerMgr -> IO ()+loadNetSeeds :: (MonadIO m) => PeerMgr -> m () loadNetSeeds mgr = when mgr.config.discover $ do- ss <- concat <$> mapM (toSockAddr mgr.config.net) mgr.config.net.seeds+ ss <-+ concat+ <$> mapM (liftIO . toSockAddr mgr.config.net) mgr.config.net.seeds mapM_ (newPeer mgr) ss -logConnectedPeers :: PeerMgr -> IO ()+logConnectedPeers :: (MonadLoggerIO m) => PeerMgr -> m () logConnectedPeers mgr = do let m = mgr.config.maxPeers l <- length <$> getConnectedPeers mgr- logS "PeerMgr" $ "Peers connected: " <> cs (show l) <> "/" <> cs (show m)+ $(logInfoS)+ "PeerMgr"+ ("Peers connected: " <> cs (show l) <> "/" <> cs (show m)) -getOnlinePeers :: PeerMgr -> IO [OnlinePeer]+getOnlinePeers :: (MonadIO m) => PeerMgr -> m [OnlinePeer] getOnlinePeers mgr = readTVarIO mgr.peers -getConnectedPeers :: PeerMgr -> IO [OnlinePeer]+getConnectedPeers :: (MonadIO m) => PeerMgr -> m [OnlinePeer] getConnectedPeers mgr = filter (.online) <$> getOnlinePeers mgr -managerEvent :: PeerMgr -> PeerEvent -> IO ()+managerEvent :: (MonadIO m) => PeerMgr -> PeerEvent -> m () managerEvent mgr e = publish e mgr.config.pub -dispatch :: PeerMgr -> PeerMgrMessage -> IO ()+dispatch ::+ (MonadLoggerIO m, MonadUnliftIO m) => PeerMgr -> PeerMgrMessage -> m () dispatch mgr (PeerVersion p v) = do- debugS+ $(logDebugS) "PeerMgr"- ("Received peer " <> p.label <> " version: " <> cs (show v))+ ("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+ $(logDebugS) "PeerMgr"- ("Sending version ack to peer: " <> p.label)+ ("Sending version ack to peer " <> p.label) MVerAck `sendMessage` p Nothing -> do- warnS+ $(logWarnS) "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- debugS "PeerMgr" ("Received version ack from peer: " <> p.label)+ $(logDebugS) "PeerMgr" ("Received version ack from peer " <> p.label) when o.online (announcePeer mgr p) Nothing -> do- warnS+ $(logWarnS) "PeerMgr"- ("Received verack from unknown peer: " <> p.label)+ ("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+ $(logDebugS) "PeerMgr" ("Received " <> cs (show len) <> " addresses from peer " <> p.label) forM_ sas (newPeer mgr)- | otherwise = debugS "PeerMgr" ("Ignoring addresses from peer " <> p.label)+ | otherwise =+ $(logDebugS) "PeerMgr" ("Ignoring addresses from peer " <> p.label) dispatch mgr (PeerPong p n) = do- debugS+ $(logDebugS) "PeerMgr"- ("Received pong " <> cs (show n) <> " from: " <> p.label)- now <- getCurrentTime+ ("Received pong " <> cs (show n) <> " from peer " <> p.label)+ now <- liftIO getCurrentTime atomically (gotPong mgr.peers n now p) dispatch _mgr (PeerPing p n) = do- debugS+ $(logDebugS) "PeerMgr"- ("Responding to ping " <> cs (show n) <> " from: " <> p.label)+ ("Responding to ping " <> cs (show n) <> " from peer " <> p.label) MPong (Pong n) `sendMessage` p dispatch mgr (ManagerBest h) = do- debugS "PeerMgr" ("Setting best block to " <> cs (show h))+ $(logDebugS) "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)+ $(logDebugS) "PeerMgr" ("Housekeeping for peer " <> p.label) checkPeer mgr p -ticklePeer :: PeerMgr -> Peer -> IO ()+ticklePeer :: (MonadLoggerIO m) => PeerMgr -> Peer -> m () ticklePeer m p = do- debugS "PeerMgr" ("Tickle peer " <> p.label)- t <- getCurrentTime+ $(logDebugS) "PeerMgr" ("Tickling peer " <> p.label)+ t <- liftIO getCurrentTime atomically (modifyPeer m.peers p (\o -> o {tickled = t})) -checkPeer :: PeerMgr -> Peer -> IO ()+checkPeer :: (MonadLoggerIO m) => PeerMgr -> Peer -> m () checkPeer mgr p = atomically (mgr.peers `findPeer` p) >>= \case Just o -> do- now <- getCurrentTime+ now <- liftIO getCurrentTime let expired = mgr.config.maxPeerLife `addUTCTime` o.connected let deadline = mgr.config.timeout `addUTCTime` o.tickled if | now > expired -> do- warnS "PeerMgr" ("Killing old peer: " <> p.label)+ $(logWarnS) "PeerMgr" ("Peer " <> p.label <> " is too old") killPeer p | now > deadline -> do- warnS "PeerMgr" ("Peer timeout: " <> p.label)+ $(logWarnS) "PeerMgr" ("Peer " <> p.label <> " timed out") killPeer p | isNothing o.ping -> sendPing mgr p | otherwise -> return () _ -> return () -sendPing :: PeerMgr -> Peer -> IO ()+sendPing :: (MonadLoggerIO m) => PeerMgr -> Peer -> m () sendPing mgr p = do atomically (mgr.peers `findPeer` p) >>= \case Nothing ->- warnS "PeerMgr" ("Will not ping unknown peer: " <> p.label)+ $(logWarnS) "PeerMgr" ("Will not ping unknown peer " <> p.label) Just o | o.online -> do n <- randomIO- now <- getCurrentTime+ now <- liftIO getCurrentTime atomically (setPeerPing mgr.peers n now p)- debugS+ $(logDebugS) "PeerMgr"- ("Sending ping " <> cs (show n) <> " to: " <> p.label)+ ("Sending ping " <> cs (show n) <> " to peer " <> p.label) MPing (Ping n) `sendMessage` p | otherwise ->- debugS+ $(logDebugS) "PeerMgr"- ("Will not ping offline peer: " <> p.label)+ ("Will not ping offline peer " <> p.label) -processPeerOffline :: PeerMgr -> Child -> IO ()+processPeerOffline :: (MonadLoggerIO m) => PeerMgr -> Child -> m () processPeerOffline mgr a = do atomically (findPeerAsync mgr.peers a) >>= \case- Nothing -> warnS "PeerMgr" "Disconnected unknown peer"+ Nothing -> $(logWarnS) "PeerMgr" "Disconnected unknown peer" Just o -> do if o.online then do- warnS "PeerMgr" ("Disconnected peer: " <> o.mailbox.label)+ $(logWarnS)+ "PeerMgr"+ ("Disconnected peer " <> o.mailbox.label) managerEvent mgr (PeerDisconnected o.mailbox) else- warnS "PeerMgr" ("Could not connect to peer: " <> o.mailbox.label)+ $(logWarnS)+ "PeerMgr"+ ("Could not connect to peer " <> o.mailbox.label) atomically (removePeer mgr.peers o.mailbox) logConnectedPeers mgr -announcePeer :: PeerMgr -> Peer -> IO ()+announcePeer :: (MonadLoggerIO m) => PeerMgr -> Peer -> m () announcePeer mgr p = do atomically (findPeer mgr.peers p) >>= \case Just OnlinePeer {online = True} -> do- logS- "PeerMgr"- ("Connected to peer " <> p.label)+ $(logInfoS) "PeerMgr" ("Connected to peer " <> p.label) managerEvent mgr (PeerConnected p) logConnectedPeers mgr Just OnlinePeer {online = False} -> return () Nothing ->- warnS- "PeerMgr"- ("Not announcing disconnected peer: " <> p.label)+ $(logWarnS) "PeerMgr" ("Not announcing disconnected peer " <> p.label) -getNewPeer :: PeerMgr -> IO (Maybe SockAddr)+getNewPeer :: (MonadIO m) => PeerMgr -> m (Maybe SockAddr) getNewPeer mgr = do loadPeers mgr atomically . stateTVar mgr.addresses $ first (listToMaybe . Set.elems) . Set.splitAt 1 -connectPeer :: PeerMgr -> SockAddr -> IO ()+connectPeer :: (MonadUnliftIO m, MonadLoggerIO m) => PeerMgr -> SockAddr -> m () connectPeer mgr sa = do atomically (findPeerAddress mgr.peers sa) >>= \case Just _ ->- warnS- "PeerMgr"- ("Attempted to connect to peer twice: " <> cs (show sa))+ $(logWarnS) "PeerMgr" ("Duplicate connection to peer " <> cs (show sa)) Nothing -> do- logS "PeerMgr" ("Connecting to " <> cs (show sa))+ $(logWarnS) "PeerMgr" ("Connecting to peer " <> cs (show sa)) nonce <- randomIO bb <- getBestBlock mgr- now <- getCurrentTime+ now <- liftIO getCurrentTime let rmt = NetworkAddress (srv mgr.config.net) (sockToHostAddress sa) unix = floor (utcTimeToPOSIXSeconds now) ver = buildVersion mgr.config.net nonce bb mgr.config.address rmt unix@@ -355,8 +357,9 @@ connect = mgr.config.connect sa } busy <- newTVarIO False- p <- wrapPeer pc busy mailbox- a <- mgr.supervisor `addChild` launch pc busy inbox p+ let p = wrapPeer pc busy mailbox+ a <- withRunInIO $ \io ->+ mgr.supervisor `addChild` io (launch pc busy inbox p) MVersion ver `sendMessage` p atomically $ insertPeer@@ -381,11 +384,7 @@ launch pc busy inbox p = withPeerLoop mgr p (\a -> link a >> peer pc busy inbox) -withPeerLoop ::- PeerMgr ->- Peer ->- (Async a -> IO a) ->- IO a+withPeerLoop :: (MonadUnliftIO m) => PeerMgr -> Peer -> (Async a -> m a) -> m a withPeerLoop mgr p = withAsync . forever $ do let timeout = mgr.config.timeout@@ -394,7 +393,7 @@ threadDelay r managerCheck mgr p -withConnectLoop :: PeerMgr -> IO a -> IO a+withConnectLoop :: (MonadUnliftIO m) => PeerMgr -> m a -> m a withConnectLoop mgr act = withAsync go $ \a -> link a >> act where@@ -405,7 +404,7 @@ (getNewPeer mgr >>= mapM_ (managerConnect mgr)) threadDelay =<< randomRIO (10 ^ 5, 5 * 10 ^ 6) -newPeer :: PeerMgr -> SockAddr -> IO ()+newPeer :: (MonadIO m) => PeerMgr -> SockAddr -> m () newPeer mgr sa = atomically $ findPeerAddress mgr.peers sa >>= \case@@ -473,34 +472,34 @@ findPeerAddress :: TVar [OnlinePeer] -> SockAddr -> STM (Maybe OnlinePeer) findPeerAddress b a = find ((== a) . (.address)) <$> readTVar b -getPeers :: PeerMgr -> IO [OnlinePeer]+getPeers :: (MonadIO m) => PeerMgr -> m [OnlinePeer] getPeers = getConnectedPeers -getOnlinePeer :: PeerMgr -> Peer -> IO (Maybe OnlinePeer)+getOnlinePeer :: (MonadIO m) => PeerMgr -> Peer -> m (Maybe OnlinePeer) getOnlinePeer mgr p = atomically (mgr.peers `findPeer` p) -managerCheck :: PeerMgr -> Peer -> IO ()+managerCheck :: (MonadIO m) => PeerMgr -> Peer -> m () managerCheck mgr p = CheckPeer p `send` mgr.mailbox -managerConnect :: PeerMgr -> SockAddr -> IO ()+managerConnect :: (MonadIO m) => PeerMgr -> SockAddr -> m () managerConnect mgr sa = Connect sa `send` mgr.mailbox -peerMgrBest :: PeerMgr -> BlockHeight -> IO ()+peerMgrBest :: (MonadIO m) => PeerMgr -> BlockHeight -> m () peerMgrBest mgr bh = ManagerBest bh `send` mgr.mailbox -peerMgrVerAck :: PeerMgr -> Peer -> IO ()+peerMgrVerAck :: (MonadIO m) => PeerMgr -> Peer -> m () peerMgrVerAck mgr p = PeerVerAck p `send` mgr.mailbox -peerMgrVersion :: PeerMgr -> Peer -> Version -> IO ()+peerMgrVersion :: (MonadIO m) => PeerMgr -> Peer -> Version -> m () peerMgrVersion mgr p ver = PeerVersion p ver `send` mgr.mailbox -peerMgrPing :: PeerMgr -> Peer -> Word64 -> IO ()+peerMgrPing :: (MonadIO m) => PeerMgr -> Peer -> Word64 -> m () peerMgrPing mgr p nonce = PeerPing p nonce `send` mgr.mailbox -peerMgrPong :: PeerMgr -> Peer -> Word64 -> IO ()+peerMgrPong :: (MonadIO m) => PeerMgr -> Peer -> Word64 -> m () peerMgrPong mgr p nonce = PeerPong p nonce `send` mgr.mailbox -peerMgrAddrs :: PeerMgr -> Peer -> [NetworkAddress] -> IO ()+peerMgrAddrs :: (MonadIO m) => PeerMgr -> Peer -> [NetworkAddress] -> m () peerMgrAddrs mgr p addrs = PeerAddrs p addrs `send` mgr.mailbox toHostService :: String -> (Maybe String, Maybe String)
test/Haskoin/NodeSpec.hs view
@@ -11,10 +11,9 @@ module Haskoin.NodeSpec (spec) where import Conduit-import Control.Concurrent.Async-import Control.Logging import Control.Monad (forM_, forever, replicateM) import Control.Monad.Cont+import Control.Monad.Logger import Control.Monad.Trans (lift) import Data.ByteString (ByteString) import Data.ByteString qualified as B@@ -22,7 +21,7 @@ import Data.Default (def) import Data.Either (fromRight) import Data.List (find)-import Data.Maybe (isJust, mapMaybe)+import Data.Maybe (isJust, listToMaybe, mapMaybe) import Data.Serialize (decode, get, runGet, runPut) import Data.Time.Clock.POSIX (getPOSIXTime) import Database.RocksDB qualified as R@@ -30,10 +29,10 @@ import Haskoin.Node import NQE import Network.Socket (AddrInfo (addrAddress), SockAddr (..))-import System.IO.Temp import System.Random (randomIO) import Test.Hspec import Test.Hspec.QuickCheck+import UnliftIO data TestNode = TestNode { testMgr :: PeerMgr,@@ -106,8 +105,8 @@ if b then SockAddrInet p w1 else SockAddrInet6 p 0 (w1, w2, w3, w4) 0- s <- head <$> toSockAddr net (show a)- s `shouldBe` a+ s <- toSockAddr net (show a)+ s `shouldStartWith` [a] it "reads some specific addresses" $ do toHostService "localhost" `shouldBe` (Just "localhost", Nothing) toHostService "::1" `shouldBe` (Just "::1", Nothing)@@ -184,13 +183,17 @@ PeerEvent (PeerConnected p) -> Just p _ -> Nothing -withTestNode :: Network -> String -> (TestNode -> IO a) -> IO a-withTestNode net str f = withStderrLogging $ flip runContT return $ do- lift $ setLogLevel LevelError+withLifted :: (MonadUnliftIO m) => ((a -> IO b) -> IO b) -> (a -> m b) -> m b+withLifted w m = withRunInIO $ \io -> w (io . m)++withTestNode ::+ (MonadUnliftIO m) =>+ Network -> String -> (TestNode -> m a) -> m a+withTestNode net str f = runNoLoggingT $ flip runContT return $ do w <- ContT $ withSystemTempDirectory ("haskoin-node-test-" <> str <> "-") pub <- ContT withPublisher sub <- ContT $ withSubscription pub- db <- ContT $ R.withDBCF w cfg cols+ db <- ContT $ withLifted (R.withDBCF w cfg cols) let ad = NetworkAddress nodeNetwork@@ -203,7 +206,7 @@ NodeConfig { maxPeers = 20, db = db,- cf = Just (head (R.columnFamilies db)),+ cf = listToMaybe (R.columnFamilies db), peers = ["[::1]:17486"], discover = False, address = na,@@ -214,7 +217,7 @@ connect = dummyPeerConnect net ad } Node mgr ch <- ContT $ withNode cfg'- lift $+ lift . lift $ f TestNode { testMgr = mgr,