haskoin-store 0.34.7 → 0.35.0
raw patch · 10 files changed
+327/−290 lines, 10 filesdep +rocksdb-haskell-jpruppdep −rocksdb-haskelldep ~haskoin-storedep ~haskoin-store-datadep ~rocksdb-query
Dependencies added: rocksdb-haskell-jprupp
Dependencies removed: rocksdb-haskell
Dependency ranges changed: haskoin-store, haskoin-store-data, rocksdb-query
Files
- app/Main.hs +15/−0
- haskoin-store.cabal +10/−10
- src/Haskoin/Store/BlockStore.hs +21/−12
- src/Haskoin/Store/Cache.hs +163/−159
- src/Haskoin/Store/Database/Reader.hs +57/−77
- src/Haskoin/Store/Database/Writer.hs +8/−8
- src/Haskoin/Store/Logic.hs +21/−2
- src/Haskoin/Store/Manager.hs +6/−6
- src/Haskoin/Store/Web.hs +24/−15
- test/Haskoin/StoreSpec.hs +2/−1
app/Main.hs view
@@ -65,6 +65,7 @@ , configPeerTimeout :: !Int , configPeerMaxLife :: !Int , configMaxPeers :: !Int+ , configMaxDiff :: !Int } instance Default Config where@@ -88,6 +89,7 @@ , configPeerTimeout = defPeerTimeout , configPeerMaxLife = defPeerMaxLife , configMaxPeers = defMaxPeers+ , configMaxDiff = defMaxDiff } defEnv :: MonadIO m => String -> a -> (String -> Maybe a) -> m a@@ -207,6 +209,11 @@ defEnv "PEER_MAX_LIFE" (48 * 3600) readMaybe {-# NOINLINE defPeerMaxLife #-} +defMaxDiff :: Int+defMaxDiff = unsafePerformIO $+ defEnv "MAX_DIFF" 2 readMaybe+{-# NOINLINE defMaxDiff #-}+ netNames :: String netNames = intercalate "|" (map getNetworkName allNets) @@ -390,6 +397,12 @@ flag (configWipeMempool def) True $ long "wipe-mempool" <> help "Wipe mempool at start"+ configMaxDiff <-+ option auto $+ metavar "INT"+ <> long "max-diff"+ <> help "Maximum difference between headers and blocks"+ <> value (configMaxDiff def) pure Config { configWebLimits = WebLimits {..}@@ -458,6 +471,7 @@ , configPeerTimeout = peertimeout , configPeerMaxLife = peerlife , configMaxPeers = maxpeers+ , configMaxDiff = maxdiff } = runStderrLoggingT . filterLogger l $ do $(logInfoS) "Main" $@@ -495,6 +509,7 @@ , webTimeouts = tos , webMaxPending = pend , webVersion = version+ , webMaxDiff = maxdiff } where l _ lvl
haskoin-store.cabal view
@@ -4,10 +4,10 @@ -- -- see: https://github.com/sol/hpack ----- hash: 419bd834bc9fa07c6d1097df06a7b19c9b3b4c10f0b6a2f3e95d0849206567fd+-- hash: f1615ec7d8469f9d8c159a8a4c0a62904bfdd0ea8143e0d6c38f4cc0f7da7ad9 name: haskoin-store-version: 0.34.7+version: 0.35.0 synopsis: Storage and index for Bitcoin and Bitcoin Cash description: Please see the README on GitHub at <https://github.com/haskoin/haskoin-store#readme> category: Bitcoin, Finance, Network@@ -53,7 +53,7 @@ , hashable >=1.3.0.0 , haskoin-core >=0.13.6 , haskoin-node >=0.14.1- , haskoin-store-data ==0.34.7+ , haskoin-store-data ==0.35.0 , hedis >=0.12.13 , http-types >=0.12.3 , monad-logger >=0.3.32@@ -61,8 +61,8 @@ , network >=3.1.1.1 , nqe >=0.6.1 , random >=1.1- , rocksdb-haskell >=1.0.1- , rocksdb-query >=0.3.1+ , rocksdb-haskell-jprupp >=2.0.0+ , rocksdb-query >=0.4.0 , scotty >=0.11.5 , string-conversions >=0.4.0.1 , text >=1.2.4.0@@ -96,7 +96,7 @@ , haskoin-core >=0.13.6 , haskoin-node >=0.14.1 , haskoin-store- , haskoin-store-data ==0.34.7+ , haskoin-store-data ==0.35.0 , monad-logger >=0.3.32 , mtl >=2.2.2 , nqe >=0.6.1@@ -135,8 +135,8 @@ , hashable >=1.3.0.0 , haskoin-core >=0.13.6 , haskoin-node >=0.14.1- , haskoin-store ==0.34.7- , haskoin-store-data ==0.34.7+ , haskoin-store ==0.35.0+ , haskoin-store-data ==0.35.0 , hedis >=0.12.13 , hspec >=2.7.1 , http-types >=0.12.3@@ -145,8 +145,8 @@ , network >=3.1.1.1 , nqe >=0.6.1 , random >=1.1- , rocksdb-haskell >=1.0.1- , rocksdb-query >=0.3.1+ , rocksdb-haskell-jprupp >=2.0.0+ , rocksdb-query >=0.4.0 , scotty >=0.11.5 , string-conversions >=0.4.0.1 , text >=1.2.4.0
src/Haskoin/Store/BlockStore.hs view
@@ -84,8 +84,9 @@ initBest, newMempoolTx, revertBlock) import NQE (Listen, Mailbox, Publisher,- newMailbox, publish, query,- receive, send, sendSTM)+ inboxToMailbox, newInbox,+ publish, query, receive, send,+ sendSTM) import System.Random (randomRIO) import UnliftIO (Exception, MonadIO, MonadUnliftIO, STM, TVar, async,@@ -136,6 +137,7 @@ , myPeer :: !(TVar (Maybe Syncing)) , myTxs :: !(TVar (HashMap TxHash PendingTx)) , requested :: !(TVar (HashSet TxHash))+ , mempooled :: !(TVar Bool) } -- | Configuration for a block store.@@ -213,21 +215,23 @@ pb <- newTVarIO Nothing ts <- newTVarIO HashMap.empty rq <- newTVarIO HashSet.empty- (inbox, mbox) <- newMailbox- let r = BlockStore { myMailbox = mbox+ inbox <- newInbox+ mem <- newTVarIO False+ let r = BlockStore { myMailbox = inboxToMailbox inbox , myConfig = cfg , myPeer = pb , myTxs = ts , requested = rq+ , mempooled = mem }- withAsync (runReaderT (go inbox mbox) r) $ \a -> do+ withAsync (runReaderT (go inbox) r) $ \a -> do link a action r where- go inbox mbox = do+ go inbox = do ini wipe- run inbox mbox+ run inbox del txs = do $(logInfoS) "BlockStore" $ "Deleting " <> cs (show (length txs)) <> " transactions"@@ -253,7 +257,7 @@ "Could not initialize: " <> cs (show e) throwIO e Right () -> return ()- run inbox mbox = withAsync (pingMe mbox)+ run inbox = withAsync (pingMe (inboxToMailbox inbox)) $ const $ forever $ receive inbox >>=@@ -271,11 +275,16 @@ then clearSyncingState >> return True else return False -mempool :: MonadLoggerIO m => Peer -> m ()+mempool :: MonadLoggerIO m => Peer -> BlockT m () mempool p = do- $(logDebugS) "BlockStore" $- "Requesting mempool from peer: " <> peerText p- MMempool `sendMessage` p+ mem <- readTVarIO =<< asks mempooled+ if mem+ then return ()+ else do+ $(logDebugS) "BlockStore" $+ "Requesting mempool from peer: " <> peerText p+ MMempool `sendMessage` p+ atomically . (`writeTVar` True) =<< asks mempooled processBlock :: MonadLoggerIO m => Peer -> Block -> BlockT m () processBlock peer block = void . runMaybeT $ do
src/Haskoin/Store/Cache.hs view
@@ -26,7 +26,8 @@ ) where import Control.DeepSeq (NFData)-import Control.Monad (forM, forM_, forever, unless, void)+import Control.Monad (forM, forM_, forever, unless, void,+ when) import Control.Monad.Logger (MonadLoggerIO, logDebugS, logErrorS, logInfoS, logWarnS) import Control.Monad.Reader (ReaderT (..), ask, asks)@@ -43,7 +44,7 @@ import qualified Data.IntMap.Strict as I import Data.List (sort) import qualified Data.Map.Strict as Map-import Data.Maybe (catMaybes, mapMaybe)+import Data.Maybe (catMaybes, isNothing, mapMaybe) import Data.Serialize (Serialize, decode, encode) import Data.String.Conversions (cs) import Data.Text (Text)@@ -425,12 +426,6 @@ mempoolSetKey :: ByteString mempoolSetKey = "mempool" -lockKey :: ByteString-lockKey = "xpub"--importLockKey :: ByteString-importLockKey = "import"- addrPfx :: ByteString addrPfx = "a" @@ -460,26 +455,26 @@ x <- receive inbox cacheWriterReact x -lockIt :: MonadLoggerIO m => ByteString -> CacheX m (Maybe Word32)-lockIt k = do+lockIt :: MonadLoggerIO m => CacheX m (Maybe Word32)+lockIt = do rnd <- liftIO randomIO go rnd >>= \case Right Redis.Ok -> do $(logDebugS) "Cache" $- "Acquired " <> cs k <> " lock with value " <> cs (show rnd)+ "Acquired lock with value " <> cs (show rnd) return (Just rnd) Right Redis.Pong -> do- $(logErrorS) "Cache" $- "Unexpected pong when acquiring " <> cs k <> " lock"+ $(logErrorS) "Cache"+ "Unexpected pong when acquiring lock" return Nothing Right (Redis.Status s) -> do $(logErrorS) "Cache" $- "Unexpected status acquiring " <> cs k <> " lock: " <> cs s+ "Unexpected status acquiring lock: " <> cs s return Nothing Left (Redis.Bulk Nothing) -> return Nothing Left e -> do- $(logErrorS) "Cache" $- "Error when trying to acquire " <> cs k <> " lock"+ $(logErrorS) "Cache"+ "Error when trying to acquire lock" throwIO (RedisError e) where go rnd = do@@ -491,55 +486,60 @@ , Redis.setMilliseconds = Nothing , Redis.setCondition = Just Redis.Nx }- Redis.setOpts k (cs (show rnd)) opts+ Redis.setOpts "lock" (cs (show rnd)) opts -unlockIt :: MonadLoggerIO m => ByteString -> Maybe Word32 -> CacheX m ()-unlockIt _ Nothing = return ()-unlockIt k (Just i) =- runRedis (Redis.get k) >>= \case+unlockIt :: MonadLoggerIO m => Maybe Word32 -> CacheX m ()+unlockIt Nothing = return ()+unlockIt (Just i) =+ runRedis (Redis.get "lock") >>= \case Nothing -> $(logErrorS) "Cache" $- "Not releasing " <> cs k <> " lock with value " <> cs (show i) <>+ "Not releasing lock with value " <> cs (show i) <> ": not locked" Just bs -> if read (cs bs) == i then do- void $ runRedis (Redis.del [k])+ void $ runRedis (Redis.del ["lock"]) $(logDebugS) "Cache" $- "Released " <> cs k <> " lock with value " <>+ "Released lock with value " <> cs (show i) else $(logErrorS) "Cache" $- "Could not release " <> cs k <> " lock: value is not " <>+ "Could not release lock: value is not " <> cs (show i) withLock :: (MonadLoggerIO m, MonadUnliftIO m)- => ByteString- -> CacheX m a+ => CacheX m a -> CacheX m (Maybe a)-withLock k f =- bracket (lockIt k) (unlockIt k) $ \case+withLock f =+ bracket lockIt unlockIt $ \case Just _ -> Just <$> f Nothing -> return Nothing withLockWait :: (MonadLoggerIO m, MonadUnliftIO m)- => ByteString- -> CacheX m a+ => Integer -> CacheX m a-withLockWait k f = do- $(logDebugS) "Cache" $ "Acquiring " <> cs k <> " lock…"- go+ -> CacheX m (Maybe a)+withLockWait n f = do+ $(logDebugS) "Cache" "Acquiring lock"+ go n where- go =- withLock k f >>= \case- Just x -> return x- Nothing -> do- r <- liftIO $ randomRIO (500, 10000)- threadDelay r- go+ go n'+ | n' <= 0 = do+ $(logDebugS) "Cache" "Failed to acquire lock"+ return Nothing+ | otherwise = withLock f >>= \case+ Just x -> return $ Just x+ Nothing -> do+ $(logDebugS) "Cache" $+ "Retrying acquisition of lock " <>+ cs (show n') <> " times"+ r <- liftIO $ randomRIO (5000, 10000)+ threadDelay r+ go (n' - 1) pruneDB :: (MonadUnliftIO m, MonadLoggerIO m, StoreReadBase m) => CacheX m Integer@@ -556,14 +556,17 @@ flush n = case min 1000 (n `div` 64) of 0 -> return 0- x ->- withLockWait lockKey $ do+ x -> do+ m <- withLockWait 20 $ do ks <- fmap (map fst) . runRedis $ getFromSortedSet maxKey Nothing 0 (fromIntegral x) $(logDebugS) "Cache" $- "Pruning " <> cs (show (length ks)) <> " old xpubs…"+ "Pruning " <> cs (show (length ks)) <> " old xpubs" delXPubKeys ks+ case m of+ Nothing -> return 0+ Just y -> return y touchKeys :: MonadLoggerIO m => [XPubSpec] -> CacheX m () touchKeys xpubs = do@@ -580,7 +583,7 @@ => CacheWriterMessage -> CacheX m () cacheWriterReact CacheNewBlock = newBlockC cacheWriterReact (CacheNewTx th) = newTxC th-cacheWriterReact CachePing = syncMempoolC+cacheWriterReact CachePing = newBlockC >> syncMempoolC lenNotNull :: [XPubBal] -> Int lenNotNull bals = length $ filter (not . nullBalance . xPubBal) bals@@ -596,17 +599,16 @@ Just _ -> do bals <- lift $ xPubBals xpub x <- asks cacheMin- xpubtxt <- xpubText xpub+ t <- xpubText xpub let n = lenNotNull bals if x <= n- then do- go bals- return (bals, True)+ then go bals >>= \case+ Just _ -> return (bals, True)+ Nothing -> return (bals, False) else do $(logDebugS) "Cache" $ "Not caching xpub with " <> cs (show n) <>- " used addresses: " <>- xpubtxt+ " used addresses: " <> t return (bals, False) where op XPubUnspent {xPubUnspent = u} = (unspentPoint u, unspentBlock u)@@ -627,9 +629,9 @@ "Caching xpub with " <> cs (show (length xtxs)) <> " txs: " <> xpubtxt now <- systemSeconds <$> liftIO getSystemTime- withLockWait lockKey $ do+ withLockWait 20 $ do $(logDebugS) "Cache" $- "Running Redis pipeline to cache xpub: " <> xpubtxt <> "…"+ "Running Redis pipeline to cache xpub: " <> xpubtxt <> "" runRedis $ do b <- redisTouchKeys now [xpub] c <- redisAddXPubBalances xpub bals@@ -640,9 +642,8 @@ newTxC :: (MonadUnliftIO m, MonadLoggerIO m, StoreReadExtra m) => TxHash -> CacheX m ()-newTxC th =- withLockWait importLockKey $- isAnythingCached >>= \case+newTxC th = do+ m <- withLockWait 20 $ isAnythingCached >>= \case False -> $(logDebugS) "Cache" $ "Not importing tx " <> txHashToHex th <> " because cache is empty"@@ -660,67 +661,67 @@ Just _ -> $(logDebugS) "Cache" $ "Tx already in cache: " <> txHashToHex th+ when (isNothing m) $+ $(logErrorS) "Cache" $+ "Could not get lock to cache tx: " <> txHashToHex th newBlockC :: (MonadUnliftIO m, MonadLoggerIO m, StoreReadExtra m) => CacheX m ()-newBlockC = withLockWait importLockKey f+newBlockC =+ withLockWait 20 f >>= \m ->+ when (isNothing m) $+ $(logErrorS) "Cache" "Could not get lock to add add block to cache" where- f =- isAnythingCached >>= \case- False ->- $(logDebugS) "Cache" "Not syncing blocks because cache is empty"- True ->- cachePrime >>= \case- Nothing -> return ()- Just cachehead ->- lift getBestBlock >>= \case- Nothing -> return ()- Just newhead -> go newhead cachehead+ f = isAnythingCached >>= \case+ False ->+ $(logDebugS) "Cache" "Not syncing blocks because cache is empty"+ True -> cachePrime >>= \case+ Nothing -> return ()+ Just cachehead -> lift getBestBlock >>= \case+ Nothing -> return ()+ Just newhead -> go newhead cachehead go newhead cachehead | cachehead == newhead = do $(logDebugS) "Cache" $ "Blocks in sync: " <> blockHashToHex cachehead syncMempoolC- | otherwise =- asks cacheChain >>= \ch ->- chainBlockMain newhead ch >>= \case- False ->+ | otherwise = do+ ch <- asks cacheChain+ chainBlockMain newhead ch >>= \case+ False ->+ $(logErrorS) "Cache" $+ "New head not in main chain: " <> blockHashToHex newhead+ True -> chainGetBlock cachehead ch >>= \case+ Nothing -> $(logErrorS) "Cache" $- "New head not in main chain: " <> blockHashToHex newhead- True ->- chainGetBlock cachehead ch >>= \case- Nothing ->- $(logErrorS) "Cache" $- "Cache head block node not found: " <>+ "Cache head block node not found: " <>+ blockHashToHex cachehead+ Just cacheheadnode -> chainBlockMain cachehead ch >>= \case+ False -> do+ $(logDebugS) "Cache" $+ "Reverting cache head not in main chain: " <> blockHashToHex cachehead- Just cacheheadnode ->- chainBlockMain cachehead ch >>= \case- False -> do- $(logDebugS) "Cache" $- "Reverting cache head not in main chain: " <>- blockHashToHex cachehead- removeHeadC >> f- True ->- chainGetBlock newhead ch >>= \case- Nothing -> do- $(logErrorS) "Cache" $- "Cache head node not found: " <>- blockHashToHex newhead- throwIO $- LogicError $- "Cache head node not found: " <>- cs (blockHashToHex newhead)- Just newheadnode ->- next newheadnode cacheheadnode+ removeHeadC >> f+ True -> chainGetBlock newhead ch >>= \case+ Nothing -> do+ $(logErrorS) "Cache" $+ "Cache head node not found: " <>+ blockHashToHex newhead+ throwIO $+ LogicError $+ "Cache head node not found: " <>+ cs (blockHashToHex newhead)+ Just newheadnode ->+ next newheadnode cacheheadnode next newheadnode cacheheadnode = do ch <- asks cacheChain chainGetAncestor (nodeHeight cacheheadnode + 1) newheadnode ch >>= \case Nothing -> $(logWarnS) "Cache" $- "Ancestor not found at height " <>- cs (show (nodeHeight cacheheadnode + 1)) <>- " for block: " <>- blockHashToHex (headerHash (nodeHeader newheadnode))+ "Ancestor not found at height " <>+ cs (show (nodeHeight cacheheadnode + 1)) <>+ " for block: " <>+ blockHashToHex (headerHash (nodeHeader newheadnode)) Just newcachenode -> importBlockC (headerHash (nodeHeader newcachenode)) >> f @@ -770,18 +771,18 @@ (MonadUnliftIO m, StoreReadExtra m, MonadLoggerIO m) => [TxData] -> CacheX m () importMultiTxC txs = do- $(logDebugS) "Cache" $ "Processing " <> cs (show (length txs)) <> " txs…"+ $(logDebugS) "Cache" $ "Processing " <> cs (show (length txs)) <> " txs" $(logDebugS) "Cache" $ "Getting address information for " <> cs (show (length alladdrs)) <>- " addresses…"+ " addresses" addrmap <- getaddrmap let addrs = HashMap.keys addrmap $(logDebugS) "Cache" $ "Getting balances for " <> cs (show (HashMap.size addrmap)) <>- " addresses…"+ " addresses" balmap <- getbalances addrs $(logDebugS) "Cache" $- "Getting unspent data for " <> cs (show (length allops)) <> " outputs…"+ "Getting unspent data for " <> cs (show (length allops)) <> " outputs" unspentmap <- getunspents gap <- lift getMaxGap now <- systemSeconds <$> liftIO getSystemTime@@ -792,21 +793,20 @@ "Affected xpub " <> cs (show i) <> "/" <> cs (show (length xpubs)) <> ": " <> xpubtxt- addrs' <-- withLockWait lockKey $ do- $(logDebugS) "Cache" $- "Getting xpub balances for " <> cs (show (length xpubs)) <>- " xpubs…"- xmap <- getxbals xpubs- let addrmap' = faddrmap (HashMap.keysSet xmap) addrmap- $(logDebugS) "Cache" "Starting Redis import pipeline…"- runRedis $ do- x <- redisImportMultiTx addrmap' unspentmap txs- y <- redisUpdateBalances addrmap' balmap- z <- redisTouchKeys now (HashMap.keys xmap)- return $ x >> y >> z >> return ()- $(logDebugS) "Cache" "Completed Redis pipeline"- return $ getNewAddrs gap xmap (HashMap.elems addrmap')+ addrs' <- do+ $(logDebugS) "Cache" $+ "Getting xpub balances for " <> cs (show (length xpubs)) <>+ " xpubs"+ xmap <- getxbals xpubs+ let addrmap' = faddrmap (HashMap.keysSet xmap) addrmap+ $(logDebugS) "Cache" "Starting Redis import pipeline"+ runRedis $ do+ x <- redisImportMultiTx addrmap' unspentmap txs+ y <- redisUpdateBalances addrmap' balmap+ z <- redisTouchKeys now (HashMap.keys xmap)+ return $ x >> y >> z >> return ()+ $(logDebugS) "Cache" "Completed Redis pipeline"+ return $ getNewAddrs gap xmap (HashMap.elems addrmap') cacheAddAddresses addrs' where alladdrsls = HashSet.toList alladdrs@@ -913,22 +913,21 @@ cacheAddAddresses addrs = do $(logDebugS) "Cache" $ "Adding " <> cs (show (length addrs)) <> " new generated addresses"- $(logDebugS) "Cache" "Getting balances…"+ $(logDebugS) "Cache" "Getting balances" balmap <- HashMap.fromListWith (<>) <$> mapM (uncurry getbal) addrs- $(logDebugS) "Cache" "Getting unspent outputs…"+ $(logDebugS) "Cache" "Getting unspent outputs" utxomap <- HashMap.fromListWith (<>) <$> mapM (uncurry getutxo) addrs- $(logDebugS) "Cache" "Getting transactions…"+ $(logDebugS) "Cache" "Getting transactions" txmap <- HashMap.fromListWith (<>) <$> mapM (uncurry gettxmap) addrs- withLockWait lockKey $ do- $(logDebugS) "Cache" "Running Redis pipeline…"- runRedis $ do- a <- forM (HashMap.toList balmap) (uncurry redisAddXPubBalances)- b <- forM (HashMap.toList utxomap) (uncurry redisAddXPubUnspents)- c <- forM (HashMap.toList txmap) (uncurry redisAddXPubTxs)- return $ sequence_ a >> sequence_ b >> sequence_ c- $(logDebugS) "Cache" "Completed Redis pipeline"+ $(logDebugS) "Cache" "Running Redis pipeline"+ runRedis $ do+ a <- forM (HashMap.toList balmap) (uncurry redisAddXPubBalances)+ b <- forM (HashMap.toList utxomap) (uncurry redisAddXPubUnspents)+ c <- forM (HashMap.toList txmap) (uncurry redisAddXPubTxs)+ return $ sequence_ a >> sequence_ b >> sequence_ c+ $(logDebugS) "Cache" "Completed Redis pipeline" let xpubs = HashSet.toList . HashSet.fromList . map addressXPubSpec $ Map.elems amap- $(logDebugS) "Cache" "Getting xpub balances…"+ $(logDebugS) "Cache" "Getting xpub balances" xmap <- getbals xpubs gap <- lift getMaxGap let notnulls = getnotnull balmap@@ -1013,34 +1012,39 @@ cachePrime :: (StoreReadBase m, MonadUnliftIO m, MonadLoggerIO m) => CacheX m (Maybe BlockHash)-cachePrime =- cacheGetHead >>= \case- Nothing -> do- $(logDebugS) "Cache" "Cache has no best block set"- lift getBestBlock >>= \case- Nothing -> do- $(logDebugS) "Cache" "Best block not set yet"- return Nothing- Just newhead -> do- ch <- asks cacheChain- chbest <- chainGetBest ch- if headerHash (nodeHeader chbest) == newhead- then do+cachePrime = cacheGetHead >>= \case+ Nothing -> do+ $(logDebugS) "Cache" "Cache has no best block set"+ lift getBestBlock >>= \case+ Nothing -> do+ $(logDebugS) "Cache" "Best block not set yet"+ return Nothing+ Just newhead -> do+ ch <- asks cacheChain+ chbest <- chainGetBest ch+ if headerHash (nodeHeader chbest) == newhead+ then do+ m <- withLockWait 20 $ do+ $(logInfoS) "Cache" "Priming cache" mem <- lift getMempool- withLockWait lockKey $ do- $(logInfoS) "Cache" "Priming cache…"- runRedis $ do- a <- redisAddToMempool mem- b <- redisSetHead newhead- return $ a >> b- $(logDebugS) "Cache" "Primed"- return (Just newhead)- else do- $(logDebugS)- "Cache"- "Not priming cache because not in sync"- return Nothing- Just cachehead -> return (Just cachehead)+ runRedis $ do+ a <- redisAddToMempool mem+ b <- redisSetHead newhead+ return $ a >> b+ $(logDebugS) "Cache" "Primed"+ return newhead+ case m of+ Nothing -> do+ $(logErrorS) "Cache"+ "Could not get lock to prime cache"+ return Nothing+ Just x -> return (Just x)+ else do+ $(logDebugS)+ "Cache"+ "Not priming cache because not in sync"+ return Nothing+ Just cachehead -> return (Just cachehead) cacheSetHead :: (MonadLoggerIO m, StoreReadBase m) => BlockHash -> CacheX m () cacheSetHead bh = do
src/Haskoin/Store/Database/Reader.hs view
@@ -1,25 +1,23 @@+{-# LANGUAGE FlexibleContexts #-} {-# LANGUAGE FlexibleInstances #-} {-# LANGUAGE LambdaCase #-} module Haskoin.Store.Database.Reader ( -- * RocksDB Database Access DatabaseReader (..) , DatabaseReaderT- , connectRocksDB , withDatabaseReader ) where -import Conduit (mapC, runConduit, runResourceT,- sinkList, (.|))+import Conduit (mapC, runConduit, sinkList, (.|)) import Control.Monad.Except (runExceptT, throwError) import Control.Monad.Reader (ReaderT, ask, asks, runReaderT)+import Data.Default (def) import Data.Function (on) import Data.List (sortBy) import Data.Maybe (fromMaybe) import Data.Word (Word32)-import Database.RocksDB (Compression (..), DB,- Options (..), ReadOptions,- defaultOptions,- defaultReadOptions, open)+import Database.RocksDB (Config (..), DB, withDB,+ withIter) import Database.RocksDB.Query (insert, matching, matchingAsList, matchingSkip, retrieve) import Haskoin (Address, BlockHash, BlockHeight,@@ -37,53 +35,45 @@ UnspentKey (..), VersionKey (..), toUnspent, valToBalance, valToUnspent)-import UnliftIO (MonadIO, liftIO)+import UnliftIO (MonadIO, MonadUnliftIO, liftIO) type DatabaseReaderT = ReaderT DatabaseReader data DatabaseReader = DatabaseReader- { databaseHandle :: !DB- , databaseReadOptions :: !ReadOptions- , databaseMaxGap :: !Word32- , databaseInitialGap :: !Word32- , databaseNetwork :: !Network+ { databaseHandle :: !DB+ , databaseMaxGap :: !Word32+ , databaseInitialGap :: !Word32+ , databaseNetwork :: !Network } dataVersion :: Word32 dataVersion = 16 -connectRocksDB ::- MonadIO m => Network -> Word32 -> Word32 -> FilePath -> m DatabaseReader-connectRocksDB net igap gap dir = do- db <-- open- dir- defaultOptions- { createIfMissing = True- , compression = SnappyCompression- , maxOpenFiles = -1- , writeBufferSize = 2 ^ (30 :: Integer)- }+withDatabaseReader :: MonadUnliftIO m+ => Network+ -> Word32+ -> Word32+ -> FilePath+ -> DatabaseReaderT m a+ -> m a+withDatabaseReader net igap gap dir f =+ withDB dir def{createIfMissing = True, maxFiles = Just (-1)} $ \db -> do let bdb = DatabaseReader- { databaseReadOptions = defaultReadOptions- , databaseHandle = db+ { databaseHandle = db , databaseMaxGap = gap , databaseNetwork = net , databaseInitialGap = igap } initRocksDB bdb- return bdb--withDatabaseReader :: MonadIO m => DatabaseReader -> DatabaseReaderT m a -> m a-withDatabaseReader = flip runReaderT+ runReaderT f bdb initRocksDB :: MonadIO m => DatabaseReader -> m ()-initRocksDB bdb@DatabaseReader {databaseReadOptions = opts, databaseHandle = db} = do+initRocksDB bdb@DatabaseReader{databaseHandle = db} = do e <- runExceptT $- retrieve db opts VersionKey >>= \case+ retrieve db VersionKey >>= \case Just v | v == dataVersion -> return () | v == 15 -> migrate15to16 bdb >> initRocksDB bdb@@ -94,8 +84,8 @@ Right () -> return () migrate15to16 :: MonadIO m => DatabaseReader -> m ()-migrate15to16 DatabaseReader {databaseReadOptions = opts, databaseHandle = db} = do- xs <- liftIO $ matchingAsList db opts OldMemKeyS+migrate15to16 DatabaseReader{databaseHandle = db} = do+ xs <- liftIO $ matchingAsList db OldMemKeyS let ys = map (\(OldMemKey t h, ()) -> (t, h)) xs insert db MemKey ys insert db VersionKey (16 :: Word32)@@ -104,37 +94,35 @@ setInitRocksDB db = insert db VersionKey dataVersion getBestDatabaseReader :: MonadIO m => DatabaseReader -> m (Maybe BlockHash)-getBestDatabaseReader DatabaseReader {databaseReadOptions = opts, databaseHandle = db} =- retrieve db opts BestKey+getBestDatabaseReader DatabaseReader{databaseHandle = db} =+ retrieve db BestKey getBlocksAtHeightDB :: MonadIO m => BlockHeight -> DatabaseReader -> m [BlockHash]-getBlocksAtHeightDB h DatabaseReader {databaseReadOptions = opts, databaseHandle = db} =- retrieve db opts (HeightKey h) >>= \case+getBlocksAtHeightDB h DatabaseReader{databaseHandle = db} =+ retrieve db (HeightKey h) >>= \case Nothing -> return [] Just ls -> return ls getDatabaseReader :: MonadIO m => BlockHash -> DatabaseReader -> m (Maybe BlockData)-getDatabaseReader h DatabaseReader {databaseReadOptions = opts, databaseHandle = db} =- retrieve db opts (BlockKey h)+getDatabaseReader h DatabaseReader{databaseHandle = db} =+ retrieve db (BlockKey h) getTxDataDB :: MonadIO m => TxHash -> DatabaseReader -> m (Maybe TxData)-getTxDataDB th DatabaseReader {databaseReadOptions = opts, databaseHandle = db} =- retrieve db opts (TxKey th)+getTxDataDB th DatabaseReader{databaseHandle = db} =+ retrieve db (TxKey th) getSpenderDB :: MonadIO m => OutPoint -> DatabaseReader -> m (Maybe Spender)-getSpenderDB op DatabaseReader {databaseReadOptions = opts, databaseHandle = db} =- retrieve db opts $ SpenderKey op+getSpenderDB op DatabaseReader{databaseHandle = db} =+ retrieve db $ SpenderKey op getBalanceDB :: MonadIO m => Address -> DatabaseReader -> m (Maybe Balance)-getBalanceDB a DatabaseReader { databaseReadOptions = opts- , databaseHandle = db- } =- fmap (valToBalance a) <$> retrieve db opts (BalKey a)+getBalanceDB a DatabaseReader{databaseHandle = db} =+ fmap (valToBalance a) <$> retrieve db (BalKey a) getMempoolDB :: MonadIO m => DatabaseReader -> m [TxRef]-getMempoolDB DatabaseReader {databaseReadOptions = opts, databaseHandle = db} =- fmap f . fromMaybe [] <$> retrieve db opts MemKey+getMempoolDB DatabaseReader{databaseHandle = db} =+ fmap f . fromMaybe [] <$> retrieve db MemKey where f (t, h) = TxRef {txRefBlock = MemRef t, txRefHash = h} @@ -155,34 +143,29 @@ -> Limits -> DatabaseReader -> m [TxRef]-getAddressTxsDB a limits bdb@DatabaseReader { databaseReadOptions = opts- , databaseHandle = db- } =- liftIO . runResourceT . runConduit $- x .| applyLimitsC limits .| mapC (uncurry f) .| sinkList+getAddressTxsDB a limits bdb@DatabaseReader{databaseHandle = db} =+ liftIO $ withIter db $ \it -> runConduit $+ x it .| applyLimitsC limits .| mapC (uncurry f) .| sinkList where- x =+ x it = case start limits of- Nothing -> matching db opts (AddrTxKeyA a)+ Nothing -> matching it (AddrTxKeyA a) Just (AtTx txh) -> getTxDataDB txh bdb >>= \case Just TxData {txDataBlock = b@BlockRef {}} ->- matchingSkip db opts (AddrTxKeyA a) (AddrTxKeyB a b)- _ -> matching db opts (AddrTxKeyA a)+ matchingSkip it (AddrTxKeyA a) (AddrTxKeyB a b)+ _ -> matching it (AddrTxKeyA a) Just (AtBlock bh) -> matchingSkip- db- opts+ it (AddrTxKeyA a) (AddrTxKeyB a (BlockRef bh maxBound)) f AddrTxKey {addrTxKeyT = t} () = t f _ _ = undefined getUnspentDB :: MonadIO m => OutPoint -> DatabaseReader -> m (Maybe Unspent)-getUnspentDB p DatabaseReader { databaseReadOptions = opts- , databaseHandle = db- } =- fmap (valToUnspent p) <$> retrieve db opts (UnspentKey p)+getUnspentDB p DatabaseReader{databaseHandle = db} =+ fmap (valToUnspent p) <$> retrieve db (UnspentKey p) getAddressesUnspentsDB :: MonadIO m@@ -203,25 +186,22 @@ -> Limits -> DatabaseReader -> m [Unspent]-getAddressUnspentsDB a limits bdb@DatabaseReader { databaseReadOptions = opts- , databaseHandle = db- } =- liftIO . runResourceT . runConduit $- x .| applyLimitsC limits .| mapC (uncurry toUnspent) .| sinkList+getAddressUnspentsDB a limits bdb@DatabaseReader{databaseHandle = db} =+ liftIO $ withIter db $ \it -> runConduit $+ x it .| applyLimitsC limits .| mapC (uncurry toUnspent) .| sinkList where- x = case start limits of- Nothing -> matching db opts (AddrOutKeyA a)+ x it = case start limits of+ Nothing -> matching it (AddrOutKeyA a) Just (AtBlock h) -> matchingSkip- db- opts+ it (AddrOutKeyA a) (AddrOutKeyB a (BlockRef h maxBound)) Just (AtTx txh) -> getTxDataDB txh bdb >>= \case Just TxData {txDataBlock = b@BlockRef {}} ->- matchingSkip db opts (AddrOutKeyA a) (AddrOutKeyB a b)- _ -> matching db opts (AddrOutKeyA a)+ matchingSkip it (AddrOutKeyA a) (AddrOutKeyB a b)+ _ -> matching it (AddrOutKeyA a) instance MonadIO m => StoreReadBase (DatabaseReaderT m) where getNetwork = asks databaseNetwork
src/Haskoin/Store/Database/Writer.hs view
@@ -285,7 +285,7 @@ runMaybeT $ MaybeT f <|> MaybeT g where f = getBestBlockH <$> readTVarIO hm- g = withDatabaseReader db getBestBlock+ g = runReaderT getBestBlock db getBlocksAtHeightI :: MonadIO m => BlockHeight@@ -294,7 +294,7 @@ getBlocksAtHeightI bh Writer {getState = hm, getReader = db} = getBlocksAtHeightH bh <$> readTVarIO hm >>= \case Just bs -> return bs- Nothing -> withDatabaseReader db $ getBlocksAtHeight bh+ Nothing -> runReaderT (getBlocksAtHeight bh) db getBlockI :: MonadIO m => BlockHash@@ -304,7 +304,7 @@ runMaybeT $ MaybeT f <|> MaybeT g where f = getBlockH bh <$> readTVarIO hm- g = withDatabaseReader db $ getBlock bh+ g = runReaderT (getBlock bh) db getTxDataI :: MonadIO m => TxHash@@ -314,20 +314,20 @@ runMaybeT $ MaybeT f <|> MaybeT g where f = getTxDataH th <$> readTVarIO hm- g = withDatabaseReader db $ getTxData th+ g = runReaderT (getTxData th) db getSpenderI :: MonadIO m => OutPoint -> Writer -> m (Maybe Spender) getSpenderI op Writer {getReader = db, getState = hm} = getSpenderH op <$> readTVarIO hm >>= \case Just (Modified s) -> return (Just s) Just Deleted -> return Nothing- Nothing -> withDatabaseReader db (getSpender op)+ Nothing -> runReaderT (getSpender op) db getBalanceI :: MonadIO m => Address -> Writer -> m (Maybe Balance) getBalanceI a Writer {getReader = db, getState = hm} = getBalanceH a <$> readTVarIO hm >>= \case Just b -> return $ Just (valToBalance a b)- Nothing -> withDatabaseReader db $ getBalance a+ Nothing -> runReaderT (getBalance a) db getUnspentI :: MonadIO m => OutPoint@@ -337,13 +337,13 @@ getUnspentH op <$> readTVarIO hm >>= \case Just (Modified u) -> return (Just (valToUnspent op u)) Just Deleted -> return Nothing- Nothing -> withDatabaseReader db (getUnspent op)+ Nothing -> runReaderT (getUnspent op) db getMempoolI :: MonadIO m => Writer -> m [TxRef] getMempoolI Writer {getState = hm, getReader = db} = getMempoolH <$> readTVarIO hm >>= \case Just xs -> return xs- Nothing -> withDatabaseReader db getMempool+ Nothing -> runReaderT getMempool db runTx :: MonadIO m => MemoryTx a -> WriterT m a runTx f = ReaderT $ atomically . runReaderT f . getState
src/Haskoin/Store/Logic.hs view
@@ -30,7 +30,7 @@ import Data.Either (rights) import qualified Data.HashSet as S import qualified Data.IntMap.Strict as I-import Data.List (nub, sortOn)+import Data.List (delete, nub, sortOn) import Data.Maybe (catMaybes, fromMaybe, isJust, isNothing, mapMaybe) import Data.Ord (Down (Down))@@ -54,7 +54,9 @@ TxRef (..), UnixTime, Unspent (..), confirmed) import Haskoin.Store.Database.Writer (MemoryTx, WriterT, runTx)-import UnliftIO (Exception, MonadIO, liftIO)+import UnliftIO (Exception, MonadIO,+ MonadUnliftIO, async, liftIO,+ waitAny) type MonadImport m = ( MonadError ImportException m@@ -125,6 +127,21 @@ runMonadMemory $ importTx (MemRef w) w rbf tx return True +multiAsync :: MonadUnliftIO m => Int -> [m a] -> m [a]+multiAsync n = go []+ where+ go [] [] = return []+ go acc xs+ | length acc >= n || null xs = do+ (a, x) <- waitAny acc+ (x :) <$> go (delete a acc) xs+ | otherwise =+ case xs of+ [] -> undefined+ (x : xs') -> do+ a <- async x+ go (a : acc) xs'+ preLoadMemory :: MonadLoggerIO m => [Tx] -> WriterT m () preLoadMemory txs = do $(logDebugS) "BlockStore" "Pre-loading memory"@@ -133,6 +150,8 @@ liftIO (runLoggingT (runReaderT go r) l) where go = do+ -- _ <- multiAsync 32 $ map loadPrevOut (concatMap prevOuts txs)+ -- _ <- multiAsync 32 $ map loadOutputBalances txs mapM_ loadPrevOut (concatMap prevOuts txs) mapM_ loadOutputBalances txs runTx . setMempool =<< getMempool
src/Haskoin/Store/Manager.hs view
@@ -7,6 +7,7 @@ import Control.Monad (forever, unless, when) import Control.Monad.Logger (MonadLoggerIO)+import Control.Monad.Reader (ReaderT (ReaderT), runReaderT) import Data.Serialize (decode) import Data.Time.Clock (NominalDiffTime) import Data.Word (Word32)@@ -39,7 +40,7 @@ connectRedis) import Haskoin.Store.Common (StoreEvent (..)) import Haskoin.Store.Database.Reader (DatabaseReader (..),- connectRocksDB,+ DatabaseReaderT, withDatabaseReader) import Network.Socket (SockAddr (..)) import NQE (Inbox, Process (..), Publisher,@@ -100,7 +101,7 @@ withStore :: (MonadLoggerIO m, MonadUnliftIO m) => StoreConfig -> (Store -> m a) -> m a withStore cfg action =- connectDB cfg >>= \db ->+ connectDB cfg $ ReaderT $ \db -> withPublisher $ \pub -> withPublisher $ \node_pub -> withSubscription node_pub $ \node_sub ->@@ -117,9 +118,9 @@ , storeNetwork = storeConfNetwork cfg } -connectDB :: MonadIO m => StoreConfig -> m DatabaseReader+connectDB :: MonadUnliftIO m => StoreConfig -> DatabaseReaderT m a -> m a connectDB cfg =- connectRocksDB+ withDatabaseReader (storeConfNetwork cfg) (storeConfInitialGap cfg) (storeConfGap cfg)@@ -181,8 +182,7 @@ action (Just (c conn)) where f conn cwinbox =- withDatabaseReader db $- cacheWriter (c conn) cwinbox+ runReaderT (cacheWriter (c conn) cwinbox) db c conn = CacheConfig { cacheConn = conn , cacheMin = storeConfCacheMin cfg
src/Haskoin/Store/Web.hs view
@@ -18,7 +18,8 @@ import Conduit () import Control.Applicative ((<|>)) import Control.Monad (forever, unless, when, (<=<))-import Control.Monad.Logger (MonadLoggerIO, logInfoS)+import Control.Monad.Logger (MonadLoggerIO, logErrorS,+ logInfoS) import Control.Monad.Reader (ReaderT, ask, asks, local, runReaderT) import Control.Monad.Trans (lift)@@ -27,6 +28,7 @@ import Data.Aeson.Encode.Pretty (Config (..), defConfig, encodePretty') import Data.Aeson.Encoding (encodingToLazyByteString, list)+import Data.Aeson.Text (encodeToLazyText) import Data.ByteString.Builder (lazyByteString) import qualified Data.ByteString.Lazy as L import qualified Data.ByteString.Lazy.Char8 as C@@ -42,6 +44,7 @@ import Data.String.Conversions (cs) import Data.Text (Text) import qualified Data.Text.Encoding as T+import Data.Text.Lazy (toStrict) import Data.Time.Clock (NominalDiffTime, diffUTCTime, getCurrentTime) import Data.Time.Clock.System (getSystemTime, systemSeconds)@@ -101,14 +104,15 @@ } data WebConfig = WebConfig- { webHost :: !String- , webPort :: !Int- , webStore :: !Store- , webMaxPending :: !Int- , webMaxLimits :: !WebLimits- , webReqLog :: !Bool- , webTimeouts :: !WebTimeouts- , webVersion :: !String+ { webHost :: !String+ , webPort :: !Int+ , webStore :: !Store+ , webMaxDiff :: !Int+ , webMaxPending :: !Int+ , webMaxLimits :: !WebLimits+ , webReqLog :: !Bool+ , webTimeouts :: !WebTimeouts+ , webVersion :: !String } data WebTimeouts = WebTimeouts@@ -874,9 +878,10 @@ return h blockHealthCheck :: (MonadUnliftIO m, MonadLoggerIO m, StoreReadBase m)- => Chain -> m BlockHealth-blockHealthCheck ch = do- let blockHealthMaxDiff = 2+ => WebConfig -> m BlockHealth+blockHealthCheck cfg = do+ let ch = storeChain $ webStore cfg+ blockHealthMaxDiff = webMaxDiff cfg blockHealthHeaders <- H.nodeHeight <$> chainGetBest ch blockHealthBlocks <-@@ -923,14 +928,18 @@ healthCheck :: (MonadUnliftIO m, MonadLoggerIO m, StoreReadBase m) => WebConfig -> m HealthCheck healthCheck cfg@WebConfig {..} = do- healthBlocks <- blockHealthCheck (storeChain webStore)+ healthBlocks <- blockHealthCheck cfg healthLastBlock <- lastBlockHealthCheck (storeChain webStore) webTimeouts healthLastTx <- lastTxHealthCheck (storeChain webStore) webTimeouts healthPendingTxs <- pendingTxsHealthCheck cfg healthPeers <- peerHealthCheck (storeManager webStore) let healthNetwork = getNetworkName (storeNetwork webStore) healthVersion = webVersion- return HealthCheck {..}+ hc = HealthCheck {..}+ unless (isOK hc) $ do+ let t = toStrict $ encodeToLazyText hc+ $(logErrorS) "Web" $ "Health check failed: " <> t+ return hc scottyDbStats :: MonadLoggerIO m => WebT m () scottyDbStats = do@@ -1062,7 +1071,7 @@ runInWebReader f = do bdb <- asks (storeDB . webStore) mc <- asks (storeCache . webStore)- lift $ withDatabaseReader bdb (withCache mc f)+ lift $ runReaderT (withCache mc f) bdb runNoCache :: MonadIO m => Bool -> ReaderT WebConfig m a -> ReaderT WebConfig m a runNoCache False f = f
test/Haskoin/StoreSpec.hs view
@@ -7,6 +7,7 @@ import Conduit import Control.Monad import Control.Monad.Logger+import Control.Monad.Reader import Data.ByteString (ByteString) import qualified Data.ByteString as B import Data.ByteString.Base64@@ -47,7 +48,7 @@ bestHeight `shouldBe` 8 it "get a block and its transactions" $ withTestStore bchRegTest "get-block-txs" $ \TestStore {..} ->- withDatabaseReader testStoreDB $ do+ flip runReaderT testStoreDB $ do let h1 = "5369ef2386c72acdf513ffd80aeba2a1774e2f004d120761e54a8bf614173f3e" get_the_block h = receive testStoreEvents >>= \case