haskoin-store 0.15.2 → 0.16.0
raw patch · 17 files changed
+2610/−2133 lines, 17 filesdep +resourcetdep −directory
Dependencies added: resourcet
Dependencies removed: directory
Files
- CHANGELOG.md +13/−0
- README.md +9/−0
- app/Main.hs +84/−508
- haskoin-store.cabal +8/−4
- src/Haskoin/Store.hs +11/−440
- src/Network/Haskoin/Store/Block.hs +102/−86
- src/Network/Haskoin/Store/Data.hs +129/−34
- src/Network/Haskoin/Store/Data/Cached.hs +198/−0
- src/Network/Haskoin/Store/Data/ImportDB.hs +108/−133
- src/Network/Haskoin/Store/Data/KeyValue.hs +23/−38
- src/Network/Haskoin/Store/Data/Memory.hs +404/−0
- src/Network/Haskoin/Store/Data/RocksDB.hs +72/−97
- src/Network/Haskoin/Store/Data/STM.hs +0/−423
- src/Network/Haskoin/Store/Logic.hs +331/−330
- src/Network/Haskoin/Store/Messages.hs +63/−34
- src/Network/Haskoin/Store/Web.hs +1040/−0
- test/Spec.hs +15/−6
CHANGELOG.md view
@@ -4,6 +4,18 @@ The format is based on [Keep a Changelog](http://keepachangelog.com/en/1.0.0/) and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0.html). +## 0.16.0+### Added+- Orphan transaction support.+- Full address balance cache in RocksDB.+- Full unspent output cache in RocksDB.++### Changed+- Significantly refactor code.+- Move web stuff to its own module.+- Change types related to databases.+- Make xpub balance, transaction and unspent queries fetch data in parallel.+ ## 0.15.2 ### Added - Internal data types to support orphan transactions.@@ -12,6 +24,7 @@ - Do not spam block actor with pings. - Fix balance/unspent cache not reverting when importing fails. - Fix transaction sorting algorithm not including transaction position information.+- Fix conflicting mempool transaction preventing block from importing. ## 0.15.1 ### Changed
README.md view
@@ -20,6 +20,15 @@ * From the root of this repository run `stack --nix build --copy-bins`. * File will usually be installed in `~/.local/bin/haskoin-store`. +## Cache++A memory-based RocksDB database can be used as a cache to store:++* Address balances.+* Unspent outputs.++Give `haskoin-store` the path to a directory mapped to RAM, and it will populate a RockDB database for caching. Needs around 25 GB at the moment (May 2019).+ ## API Documentation
app/Main.hs view
@@ -35,12 +35,12 @@ import NQE import Options.Applicative import Paths_haskoin_store as P-import System.Directory import System.Exit import System.FilePath import System.IO.Unsafe import Text.Read (readMaybe) import UnliftIO+import UnliftIO.Directory import Web.Scotty.Trans as S data Config = Config@@ -50,6 +50,7 @@ , configDiscover :: !Bool , configPeers :: ![(Host, Maybe Port)] , configVersion :: !Bool+ , configCache :: !FilePath } defPort :: Int@@ -58,14 +59,6 @@ defNetwork :: Network defNetwork = btc -instance Parsable BlockHash where- parseParam =- maybe (Left "could not decode block hash") Right . hexToBlockHash . cs--instance Parsable TxHash where- parseParam =- maybe (Left "could not decode tx hash") Right . hexToTxHash . cs- netNames :: String netNames = intercalate "|" (map getNetworkName allNets) @@ -94,6 +87,10 @@ many . option (eitherReader peerReader) $ metavar "HOST" <> long "peer" <> short 'p' <> help "Network peer (as many as required)"+ configCache <-+ option str $+ long "cache" <> short 'c' <> help "Memory mapped disk for cache" <>+ value "" configVersion <- switch $ long "version" <> short 'v' <> help "Show version" return Config {..}@@ -122,22 +119,6 @@ _ -> Left "Peer information could not be parsed" return (host, port) -defHandler :: Monad m => Network -> Except -> ActionT Except m ()-defHandler net e = do- proto <- setupBin- case e of- ThingNotFound -> status status404- BadRequest -> status status400- UserError _ -> status status400- StringError _ -> status status400- ServerError -> status status500- S.raw $ serialAny net proto e--maybeSerial :: (Monad m, JsonSerial a, BinSerial a) => Network -> Bool -- ^ protobuf- -> Maybe a -> ActionT Except m ()-maybeSerial _ _ Nothing = raise ThingNotFound-maybeSerial net proto (Just x) = S.raw $ serialAny net proto x- myDirectory :: FilePath myDirectory = unsafePerformIO $ getAppUserDataDirectory "haskoin-store" {-# NOINLINE myDirectory #-}@@ -151,496 +132,91 @@ exitSuccess when (null (configPeers conf) && not (configDiscover conf)) . liftIO $ die "ERROR: Specify peers to connect or enable peer discovery."- let net = configNetwork conf- let wdir = configDir conf </> getNetworkName net- liftIO $ createDirectoryIfMissing True wdir- db <-- open- (wdir </> "db")- R.defaultOptions- { createIfMissing = True- , compression = SnappyCompression- , maxOpenFiles = -1- , writeBufferSize = 2 `shift` 30- }- withPublisher $ \pub ->- withStore (scfg conf db pub) $ \st -> runWeb conf st db pub+ run conf where- scfg conf db pub =- StoreConfig- { storeConfMaxPeers = 20- , storeConfInitPeers =- map- (second (fromMaybe (getDefaultPort (configNetwork conf))))- (configPeers conf)- , storeConfDiscover = configDiscover conf- , storeConfDB = db- , storeConfNetwork = configNetwork conf- , storeConfListen = (`sendSTM` pub) . Event- } opts = info (helper <*> config) $ fullDesc <> progDesc "Blockchain store and API" <> Options.Applicative.header ("haskoin-store version " <> showVersion P.version) -runWeb ::- (MonadUnliftIO m, MonadLoggerIO m)- => Config- -> Store- -> DB- -> Publisher StoreEvent- -> m ()-runWeb conf st db pub = do- l <- askLoggerIO- scottyT (configPort conf) (runner l) $ do- defaultHandler (defHandler net)- S.get "/block/best" $ do- cors- n <- parse_no_tx- proto <- setupBin- res <-- runMaybeT $ do- bh <-- MaybeT $ withBlockDB defaultReadOptions db getBestBlock- b <-- MaybeT . withBlockDB defaultReadOptions db $ getBlock bh- if n- then return b {blockDataTxs = take 1 (blockDataTxs b)}- else return b- maybeSerial net proto res- S.get "/block/:block" $ do- cors- block <- param "block"- n <- parse_no_tx- proto <- setupBin- res <-- runMaybeT $ do- b <-- MaybeT . withBlockDB defaultReadOptions db $- getBlock block- if n- then return b {blockDataTxs = take 1 (blockDataTxs b)}- else return b- maybeSerial net proto res- S.get "/block/height/:height" $ do- cors- height <- param "height"- no_tx <- parse_no_tx- proto <- setupBin- res <-- do bs <-- withBlockDB defaultReadOptions db $- getBlocksAtHeight height- fmap catMaybes . forM bs $ \bh ->- runMaybeT $ do- b <-- MaybeT . withBlockDB defaultReadOptions db $- getBlock bh- return- b- { blockDataTxs =- if no_tx- then take 1 (blockDataTxs b)- else blockDataTxs b- }- S.raw $ serialAny net proto res- S.get "/block/heights" $ do- cors- heights <- param "heights"- no_tx <- parse_no_tx- proto <- setupBin- res <-- withBlockDB defaultReadOptions db $ do- bs <- concat <$> mapM getBlocksAtHeight (nub heights)- fmap catMaybes . forM bs $ \bh ->- runMaybeT $ do- b <- MaybeT $ getBlock bh- return- b- { blockDataTxs =- if no_tx- then take 1 (blockDataTxs b)- else blockDataTxs b- }- S.raw $ serialAny net proto res- S.get "/blocks" $ do- cors- blocks <- param "blocks"- no_tx <- parse_no_tx- proto <- setupBin- res <-- withBlockDB defaultReadOptions db $- fmap catMaybes . forM blocks $ \bh ->- runMaybeT $ do- b <- MaybeT $ getBlock bh- return- b- { blockDataTxs =- if no_tx- then take 1 (blockDataTxs b)- else blockDataTxs b- }- S.raw $ serialAny net proto res- S.get "/mempool" $ do- cors- (l, s) <- parse_limits- proto <- setupBin- stream $ \io flush' ->- runResourceT . withBlockDB defaultReadOptions db $ do- runConduit $ getMempoolLimit l s .| streamAny net proto io- liftIO flush'- S.get "/transaction/:txid" $ do- cors- txid <- param "txid"- proto <- setupBin- res <- withBlockDB defaultReadOptions db $ getTransaction txid- maybeSerial net proto res- S.get "/transaction/:txid/hex" $ do- cors- txid <- param "txid"- res <- withBlockDB defaultReadOptions db $ getTransaction txid- case res of- Nothing -> raise ThingNotFound- Just x ->- text . cs . encodeHex $ Serialize.encode (transactionData x)- S.get "/transaction/:txid/bin" $ do- cors- txid <- param "txid"- res <- withBlockDB defaultReadOptions db $ getTransaction txid- case res of- Nothing -> raise ThingNotFound- Just x -> do- S.setHeader "Content-Type" "application/octet-stream"- S.raw $ Serialize.encodeLazy (transactionData x)- S.get "/transaction/:txid/after/:height" $ do- cors- txid <- param "txid"- height <- param "height"- proto <- setupBin- res <-- withBlockDB defaultReadOptions db $- cbAfterHeight 10000 height txid- S.raw $ serialAny net proto res- S.get "/transactions" $ do- cors- txids <- param "txids"- proto <- setupBin- res <-- withBlockDB defaultReadOptions db $- catMaybes <$> mapM getTransaction (nub txids)- S.raw $ serialAny net proto res- S.get "/transactions/hex" $ do- cors- txids <- param "txids"- res <-- withBlockDB defaultReadOptions db $- catMaybes <$> mapM getTransaction (nub txids)- S.json $ map (encodeHex . Serialize.encode . transactionData) res- S.get "/transactions/bin" $ do- cors- txids <- param "txids"- res <-- withBlockDB defaultReadOptions db $- catMaybes <$> mapM getTransaction (nub txids)- S.setHeader "Content-Type" "application/octet-stream"- S.raw . L.concat $ map (Serialize.encodeLazy . transactionData) res- S.get "/address/:address/transactions" $ do- cors- a <- parse_address- (l, s) <- parse_limits- proto <- setupBin- stream $ \io flush' ->- runResourceT . withBlockDB defaultReadOptions db $ do- runConduit $- getAddressTxsLimit l s a .| streamAny net proto io- liftIO flush'- S.get "/address/:address/transactions/full" $ do- cors- a <- parse_address- (l, s) <- parse_limits- proto <- setupBin- stream $ \io flush' ->- runResourceT . withBlockDB defaultReadOptions db $ do- runConduit $- getAddressTxsFull l s a .| streamAny net proto io- liftIO flush'- S.get "/address/transactions" $ do- cors- as <- parse_addresses- (l, s) <- parse_limits- proto <- setupBin- stream $ \io flush' ->- runResourceT . withBlockDB defaultReadOptions db $ do- runConduit $- getAddressesTxsLimit l s as .| streamAny net proto io- liftIO flush'- S.get "/address/transactions/full" $ do- cors- as <- parse_addresses- (l, s) <- parse_limits- proto <- setupBin- stream $ \io flush' ->- runResourceT . withBlockDB defaultReadOptions db $ do- runConduit $- getAddressesTxsFull l s as .| streamAny net proto io- liftIO flush'- S.get "/address/:address/unspent" $ do- cors- a <- parse_address- (l, s) <- parse_limits- proto <- setupBin- stream $ \io flush' ->- runResourceT . withBlockDB defaultReadOptions db $ do- runConduit $- getAddressUnspentsLimit l s a .| streamAny net proto io- liftIO flush'- S.get "/address/unspent" $ do- cors- as <- parse_addresses- (l, s) <- parse_limits- proto <- setupBin- stream $ \io flush' ->- runResourceT . withBlockDB defaultReadOptions db $ do- runConduit $- getAddressesUnspentsLimit l s as .|- streamAny net proto io- liftIO flush'- S.get "/address/:address/balance" $ do- cors- address <- parse_address- proto <- setupBin- res <-- withBlockDB defaultReadOptions db $- getBalance address >>= \case- Just b -> return b- Nothing ->- return- Balance- { balanceAddress = address- , balanceAmount = 0- , balanceUnspentCount = 0- , balanceZero = 0- , balanceTxCount = 0- , balanceTotalReceived = 0+cacheDir :: Network -> FilePath -> Maybe FilePath+cacheDir net "" = Nothing+cacheDir net ch = Just (ch </> getNetworkName net </> "cache")++run :: (MonadLoggerIO m, MonadUnliftIO m) => Config -> m ()+run Config { configPort = port+ , configNetwork = net+ , configDiscover = disc+ , configPeers = peers+ , configCache = cache_path+ , configDir = db_dir+ } =+ flip finally clear $ do+ $(logInfoS) "Main" $+ "Creating working directory if not found: " <> cs wd+ createDirectoryIfMissing True wd+ db <-+ do dbh <-+ open+ (wd </> "db")+ R.defaultOptions+ { createIfMissing = True+ , compression = SnappyCompression+ , maxOpenFiles = -1+ , writeBufferSize = 2 `shift` 30+ }+ return BlockDB {blockDB = dbh, blockDBopts = defaultReadOptions}+ cdb <-+ case cd of+ Nothing -> return Nothing+ Just ch -> do+ $(logInfoS) "Main" $ "Deleting cache directory: " <> cs ch+ removePathForcibly ch+ $(logInfoS) "Main" $ "Creating cache directory: " <> cs ch+ createDirectoryIfMissing True ch+ dbh <- open ch R.defaultOptions {createIfMissing = True}+ return $+ Just+ BlockDB+ { blockDB = dbh+ , blockDBopts = defaultReadOptions }- S.raw $ serialAny net proto res- S.get "/address/balances" $ do- cors- addresses <- parse_addresses- proto <- setupBin- res <-- withBlockDB defaultReadOptions db $ do- let f a Nothing =- Balance- { balanceAddress = a- , balanceAmount = 0- , balanceUnspentCount = 0- , balanceZero = 0- , balanceTxCount = 0- , balanceTotalReceived = 0+ $(logInfoS) "Main" "Populating cache (if active)..."+ ldb <- newLayeredDB db cdb+ $(logInfoS) "Main" "Finished populating cache"+ withPublisher $ \pub ->+ let scfg =+ StoreConfig+ { storeConfMaxPeers = 20+ , storeConfInitPeers =+ map+ (second (fromMaybe (getDefaultPort net)))+ peers+ , storeConfDiscover = disc+ , storeConfDB = ldb+ , storeConfNetwork = net+ , storeConfListen = (`sendSTM` pub) . Event+ }+ in withStore scfg $ \str ->+ let wcfg =+ WebConfig+ { webPort = port+ , webNetwork = net+ , webDB = ldb+ , webPublisher = pub+ , webStore = str }- f _ (Just b) = b- mapM (\a -> f a <$> getBalance a) addresses- S.raw $ serialAny net proto res- S.get "/xpub/:xpub/balances" $ do- cors- xpub <- parse_xpub- proto <- setupBin- res <- withBlockDB defaultReadOptions db $ xpubBals xpub- S.raw $ serialAny net proto res- S.get "/xpub/:xpub/transactions" $ do- cors- x <- parse_xpub- (l, s) <- parse_limits- proto <- setupBin- bs <- withBlockDB defaultReadOptions db $ xpubBals x- stream $ \io flush' ->- runResourceT . withBlockDB defaultReadOptions db $ do- runConduit $ xpubTxsLimit l s bs .| streamAny net proto io- liftIO flush'- S.get "/xpub/:xpub/transactions/full" $ do- cors- xpub <- parse_xpub- (l, s) <- parse_limits- proto <- setupBin- bs <- withBlockDB defaultReadOptions db $ xpubBals xpub- stream $ \io flush' ->- runResourceT . withBlockDB defaultReadOptions db $ do- runConduit $ xpubTxsFull l s bs .| streamAny net proto io- liftIO flush'- S.get "/xpub/:xpub/unspent" $ do- cors- x <- parse_xpub- proto <- setupBin- (l, s) <- parse_limits- stream $ \io flush' ->- runResourceT . withBlockDB defaultReadOptions db $ do- runConduit $- xpubUnspentLimit l s x .| streamAny net proto io- liftIO flush'- S.get "/xpub/:xpub" $ do- cors- x <- parse_xpub- (l, s) <- parse_limits- proto <- setupBin- res <-- lift . runResourceT $- withBlockDB defaultReadOptions db $ xpubSummary l s x- S.raw $ serialAny net proto res- S.post "/transactions" $ do- cors- proto <- setupBin- b <- body- let bin = eitherToMaybe . Serialize.decode- hex = bin <=< decodeHex . cs . C.filter (not . isSpace)- tx <-- case hex b <|> bin (L.toStrict b) of- Nothing -> raise (UserError "decode tx fail")- Just x -> return x- lift (publishTx net pub st db tx) >>= \case- Right () -> do- S.raw $ serialAny net proto (TxId (txHash tx))- lift $- $(logDebugS) "Main" $- "Success publishing tx " <> txHashToHex (txHash tx)- Left e -> do- case e of- PubNoPeers -> status status500- PubTimeout -> status status500- PubPeerDisconnected -> status status500- PubNotFound -> status status500- PubReject _ -> status status400- S.raw $ serialAny net proto (UserError (show e))- lift $- $(logErrorS) "Main" $- "Error publishing tx " <> txHashToHex (txHash tx) <>- ": " <>- cs (show e)- finish- S.get "/dbstats" $ do- cors- getProperty db Stats >>= text . cs . fromJust- S.get "/events" $ do- cors- proto <- setupBin- stream $ \io flush' ->- withSubscription pub $ \sub ->- forever $- flush' >> receive sub >>= \se -> do- let me =- case se of- StoreBestBlock block_hash ->- Just (EventBlock block_hash)- StoreMempoolNew tx_hash ->- Just (EventTx tx_hash)- _ -> Nothing- case me of- Nothing -> return ()- Just e -> do- let bs =- serialAny net proto e <>- if proto- then mempty- else "\n"- io (lazyByteString bs)- S.get "/peers" $ do- cors- proto <- setupBin- ps <- getPeersInformation (storeManager st)- S.raw $ serialAny net proto ps- S.get "/health" $ do- cors- proto <- setupBin- h <-- liftIO . withBlockDB defaultReadOptions db $- healthCheck net (storeManager st) (storeChain st)- when (not (healthOK h) || not (healthSynced h)) $ status status503- S.raw $ serialAny net proto h- notFound $ raise ThingNotFound- where- parse_limits = do- let b = do- height <- param "height"- pos <- param "pos" `rescue` const (return maxBound)- return $ StartBlock height pos- m = do- time <- param "time"- return $ StartMem time- o = do- o <- param "offset" `rescue` const (return 0)- return $ StartOffset o- l <- (Just <$> param "limit") `rescue` const (return Nothing)- s <- b <|> m <|> o- return (l, s)- parse_address = do- address <- param "address"- case stringToAddr net address of- Nothing -> next- Just a -> return a- parse_addresses = do- addresses <- param "addresses"- let as = mapMaybe (stringToAddr net) addresses- unless (length as == length addresses) next- return as- parse_xpub = do- t <- param "xpub"- case xPubImport net t of- Nothing -> next- Just x -> return x- net = configNetwork conf- parse_no_tx = param "notx" `rescue` const (return False)- runner f l = do- u <- askUnliftIO- unliftIO u (runLoggingT l f)- cors = setHeader "Access-Control-Allow-Origin" "*"--serialAny ::- (JsonSerial a, BinSerial a)- => Network- -> Bool -- ^ binary- -> a- -> L.ByteString-serialAny net True = runPutLazy . binSerial net-serialAny net False = encodingToLazyByteString . jsonSerial net--streamAny ::- (JsonSerial i, BinSerial i, MonadIO m)- => Network- -> Bool -- ^ protobuf- -> (Builder -> IO ())- -> ConduitT i o m ()-streamAny net True io = binConduit net .| mapC lazyByteString .| streamConduit io-streamAny net False io = jsonListConduit net .| streamConduit io--jsonListConduit :: (JsonSerial a, Monad m) => Network -> ConduitT a Builder m ()-jsonListConduit net =- yield "[" >> mapC (fromEncoding . jsonSerial net) .| intersperseC "," >> yield "]"--binConduit :: (BinSerial i, Monad m) => Network -> ConduitT i L.ByteString m ()-binConduit net = mapC (runPutLazy . binSerial net)--streamConduit :: MonadIO m => (i -> IO ()) -> ConduitT i o m ()-streamConduit io = mapM_C (liftIO . io)--setupBin :: Monad m => ActionT Except m Bool-setupBin =- let p = do- setHeader "Content-Type" "application/octet-stream"- return True- j = do- setHeader "Content-Type" "application/json"- return False- in S.header "accept" >>= \case- Nothing -> j- Just x ->- if is_binary x- then p- else j+ in runWeb wcfg where- is_binary x =- let ts =- map- (T.takeWhile (/= ';'))- (T.splitOn "," (T.filter (not . isSpace) x))- in elem "application/octet-stream" ts+ clear =+ case cd of+ Nothing -> return ()+ Just ch -> do+ $(logInfoS) "Main" $ "Deleting cache directory: " <> cs ch+ removePathForcibly ch+ wd = db_dir </> getNetworkName net+ cd =+ case cache_path of+ "" -> Nothing+ ch -> Just (ch </> getNetworkName net </> "cache")
haskoin-store.cabal view
@@ -4,10 +4,10 @@ -- -- see: https://github.com/sol/hpack ----- hash: f789909b15a88fd2b0d4a03016ad5b760ac4a7b2a246a8c19e133598c5940dc0+-- hash: a806db1bfb6b706e998c48d71994d5fe905b80350b267bcdfcb11398f97bee3e name: haskoin-store-version: 0.15.2+version: 0.16.0 synopsis: Storage and index for Bitcoin and Bitcoin Cash description: Store blocks, transactions, and balances for Bitcoin or Bitcoin Cash, and make that information via REST API. category: Bitcoin, Finance, Network@@ -33,12 +33,14 @@ other-modules: Network.Haskoin.Store.Block Network.Haskoin.Store.Data+ Network.Haskoin.Store.Data.Cached Network.Haskoin.Store.Data.ImportDB Network.Haskoin.Store.Data.KeyValue+ Network.Haskoin.Store.Data.Memory Network.Haskoin.Store.Data.RocksDB- Network.Haskoin.Store.Data.STM Network.Haskoin.Store.Logic Network.Haskoin.Store.Messages+ Network.Haskoin.Store.Web autogen-modules: Paths_haskoin_store hs-source-dirs:@@ -60,6 +62,7 @@ , network , nqe , random+ , resourcet , rocksdb-haskell , rocksdb-query , scotty@@ -86,7 +89,6 @@ , conduit , containers , data-default- , directory , filepath , hashable , haskoin-core@@ -99,6 +101,7 @@ , nqe , optparse-applicative , random+ , resourcet , rocksdb-haskell , rocksdb-query , scotty@@ -137,6 +140,7 @@ , network , nqe , random+ , resourcet , rocksdb-haskell , rocksdb-query , scotty
src/Haskoin/Store.hs view
@@ -35,7 +35,12 @@ , StartFrom(..) , UnixTime , BlockPos+ , BlockDB(..)+ , LayeredDB(..)+ , WebConfig(..)+ , newLayeredDB , withStore+ , runWeb , store , getBestBlock , getBlocksAtHeight@@ -58,9 +63,6 @@ , getAddressesTxsFull , getAddressesTxsLimit , getPeersInformation- , xpubTxs- , xpubTxsLimit- , xpubTxsFull , xpubBals , xpubUnspent , xpubUnspentLimit@@ -71,11 +73,8 @@ , confirmed , cbAfterHeight , healthCheck- , mergeSourcesBy- , withBlockDB- , withBlockSTM- , withBalanceSTM- , withUnspentSTM+ , withBlockMem+ , withLayeredDB ) where import Conduit@@ -96,41 +95,16 @@ import Haskoin.Node import Network.Haskoin.Store.Block import Network.Haskoin.Store.Data+import Network.Haskoin.Store.Data.Cached+import Network.Haskoin.Store.Data.Memory import Network.Haskoin.Store.Data.RocksDB-import Network.Haskoin.Store.Data.STM import Network.Haskoin.Store.Messages+import Network.Haskoin.Store.Web import Network.Socket (SockAddr (..)) import NQE import System.Random import UnliftIO -data PubExcept- = PubNoPeers- | PubReject RejectCode- | PubTimeout- | PubNotFound- | PubPeerDisconnected- deriving Eq--instance Show PubExcept where- show PubNoPeers = "no peers"- show (PubReject c) =- "rejected: " <>- case c of- RejectMalformed -> "malformed"- RejectInvalid -> "invalid"- RejectObsolete -> "obsolete"- RejectDuplicate -> "duplicate"- RejectNonStandard -> "not standard"- RejectDust -> "dust"- RejectInsufficientFee -> "insufficient fee"- RejectCheckpoint -> "checkpoint"- show PubTimeout = "timeout"- show PubNotFound = "not found"- show PubPeerDisconnected = "peer disconnected"--instance Exception PubExcept- withStore :: (MonadLoggerIO m, MonadUnliftIO m) => StoreConfig@@ -161,7 +135,7 @@ let ncfg = NodeConfig { nodeConfMaxPeers = storeConfMaxPeers cfg- , nodeConfDB = storeConfDB cfg+ , nodeConfDB = blockDB . layeredDB $ storeConfDB cfg , nodeConfPeers = storeConfInitPeers cfg , nodeConfDiscover = storeConfDiscover cfg , nodeConfEvents = storeDispatch b l@@ -231,406 +205,3 @@ StoreTxReject p th (rejectCode r) (getVarString (rejectReason r)) storeDispatch _ _ (PeerEvent _) = return ()--healthCheck ::- (MonadUnliftIO m, StoreRead m)- => Network- -> Manager- -> Chain- -> m HealthCheck-healthCheck net mgr ch = do- n <- timeout (5 * 1000 * 1000) $ chainGetBest ch- b <-- runMaybeT $ do- h <- MaybeT getBestBlock- MaybeT $ getBlock h- p <- timeout (5 * 1000 * 1000) $ managerGetPeers mgr- let k = isNothing n || isNothing b || maybe False (not . null) p- s =- isJust $ do- x <- n- y <- b- guard $ nodeHeight x - blockDataHeight y <= 1- return- HealthCheck- { healthBlockBest = headerHash . blockDataHeader <$> b- , healthBlockHeight = blockDataHeight <$> b- , healthHeaderBest = headerHash . nodeHeader <$> n- , healthHeaderHeight = nodeHeight <$> n- , healthPeers = length <$> p- , healthNetwork = getNetworkName net- , healthOK = k- , healthSynced = s- }---- | Publish a new transaction to the network.-publishTx ::- (MonadUnliftIO m, MonadLoggerIO m)- => Network- -> Publisher StoreEvent- -> Store- -> DB- -> Tx- -> m (Either PubExcept ())-publishTx net pub st db tx = do- $(logDebugS) "PubTx" $- "Preparing to publish tx: " <> txHashToHex (txHash tx)- e <- withSubscription pub $ \s ->- withBlockDB defaultReadOptions db $- getTransaction (txHash tx) >>= \case- Just _ -> do- $(logErrorS) "PubTx" $- "Tx already in DB: " <> txHashToHex (txHash tx)- return $ Right ()- Nothing ->- E.runExceptT $ do- $(logDebugS) "PubTx" $ "Getting peers from manager..."- managerGetPeers (storeManager st) >>= \case- [] -> do- $(logErrorS) "PubTx" $ "No peers connected."- E.throwError PubNoPeers- OnlinePeer { onlinePeerMailbox = p- , onlinePeerAddress = a- }:_ -> do- $(logDebugS) "PubTx" $- "Sending tx " <> txHashToHex (txHash tx) <>- " to peer " <>- T.pack (show a)- MTx tx `sendMessage` p- let t =- if getSegWit net- then InvWitnessTx- else InvTx- sendMessage- (MGetData- (GetData- [InvVector t (getTxHash (txHash tx))]))- p- f p s- $(logDebugS) "PubTx" $ "Finished for tx: " <> txHashToHex (txHash tx)- return e- where- t = 15 * 1000 * 1000- f p s = do- $(logDebugS) "PubTx" $- "Waiting for peer to relay tx " <> txHashToHex (txHash tx)- liftIO (timeout t (E.runExceptT (g p s))) >>= \case- Nothing -> do- $(logErrorS) "PubTx" $- "Peer did not relay tx " <> txHashToHex (txHash tx)- E.throwError PubTimeout- Just (Left e) -> do- $(logErrorS) "PubTx" $- "Error publishing tx " <> txHashToHex (txHash tx) <>- T.pack (show e)- E.throwError e- Just (Right ()) -> do- $(logDebugS) "PubTx" $- "Success publishing tx " <> txHashToHex (txHash tx)- g p s =- receive s >>= \case- StoreTxReject p' h' c _- | p == p' && h' == txHash tx -> E.throwError $ PubReject c- StorePeerDisconnected p' _- | p == p' -> E.throwError PubPeerDisconnected- StoreMempoolNew h'- | h' == txHash tx -> return ()- _ -> g p s---- | Obtain information about connected peers from peer manager process.-getPeersInformation :: MonadIO m => Manager -> m [PeerInformation]-getPeersInformation mgr = mapMaybe toInfo <$> managerGetPeers mgr- where- toInfo op = do- ver <- onlinePeerVersion op- let as = onlinePeerAddress op- ua = getVarString $ userAgent ver- vs = version ver- sv = services ver- rl = relay ver- return- PeerInformation- { peerUserAgent = ua- , peerAddress = as- , peerVersion = vs- , peerServices = sv- , peerRelay = rl- }--xpubBals :: (Monad m, BalanceRead m) => XPubKey -> m [XPubBal]-xpubBals xpub = (<>) <$> go 0 0 <*> go 1 0- where- go m n = do- xs <- catMaybes <$> mapM (uncurry b) (as m n)- case xs of- [] -> return []- _ -> (xs <>) <$> go m (n + 20)- b a p =- getBalance a >>= \case- Nothing -> return Nothing- Just b' -> return $ Just XPubBal {xPubBalPath = p, xPubBal = b'}- as m n =- map- (\(a, _, n') -> (a, [m, n']))- (take 20 (deriveAddrs (pubSubKey xpub m) n))--xpubTxs ::- (Monad m, BalanceRead m, StoreStream m)- => Maybe BlockRef- -> [XPubBal]- -> ConduitT () BlockTx m ()-xpubTxs m bs = do- xs <-- forM bs $ \XPubBal {xPubBal = b} ->- return $ getAddressTxs (balanceAddress b) m- mergeSourcesBy (flip compare `on` blockTxBlock) xs .| dedup--xpubTxsLimit ::- (Monad m, BalanceRead m, StoreStream m)- => Maybe Word32- -> StartFrom- -> [XPubBal]- -> ConduitT () BlockTx m ()-xpubTxsLimit l s bs = do- xpubTxs (mbr s) bs .| (offset s >> limit l)--xpubTxsFull ::- (Monad m, BalanceRead m, StoreStream m, StoreRead m)- => Maybe Word32- -> StartFrom- -> [XPubBal]- -> ConduitT () Transaction m ()-xpubTxsFull l s bs =- xpubTxsLimit l s bs .| concatMapMC (getTransaction . blockTxHash)--xpubUnspent ::- (Monad m, StoreStream m, BalanceRead m, StoreRead m)- => Maybe BlockRef- -> XPubKey- -> ConduitT () XPubUnspent m ()-xpubUnspent mbr xpub = do- bals <- lift $ xpubBals xpub- xs <-- forM bals $ \XPubBal {xPubBalPath = p, xPubBal = b} ->- return $ getAddressUnspents (balanceAddress b) mbr .| mapC (f p)- mergeSourcesBy (flip compare `on` (unspentBlock . xPubUnspent)) xs- where- f p t = XPubUnspent {xPubUnspentPath = p, xPubUnspent = t}--xpubUnspentLimit ::- (Monad m, StoreStream m, BalanceRead m, StoreRead m)- => Maybe Word32- -> StartFrom- -> XPubKey- -> ConduitT () XPubUnspent m ()-xpubUnspentLimit l s x =- xpubUnspent (mbr s) x .| (offset s >> limit l)--xpubSummary ::- (Monad m, StoreStream m, BalanceRead m, StoreRead m)- => Maybe Word32- -> StartFrom- -> XPubKey- -> m XPubSummary-xpubSummary l s x = do- bs <- xpubBals x- let f XPubBal {xPubBalPath = p, xPubBal = Balance {balanceAddress = a}} =- (a, p)- pm = H.fromList $ map f bs- txs <- runConduit $ xpubTxsFull l s bs .| sinkList- let as =- nub- [ a- | t <- txs- , let is = transactionInputs t- , let os = transactionOutputs t- , let ais =- mapMaybe- (eitherToMaybe . scriptToAddressBS . inputPkScript)- is- , let aos =- mapMaybe- (eitherToMaybe . scriptToAddressBS . outputScript)- os- , a <- ais ++ aos- ]- ps = H.fromList $ mapMaybe (\a -> (a, ) <$> H.lookup a pm) as- ex = foldl max 0 [i | XPubBal {xPubBalPath = [x, i]} <- bs, x == 0]- ch = foldl max 0 [i | XPubBal {xPubBalPath = [x, i]} <- bs, x == 1]- return- XPubSummary- { xPubSummaryReceived =- sum (map (balanceTotalReceived . xPubBal) bs)- , xPubSummaryConfirmed = sum (map (balanceAmount . xPubBal) bs)- , xPubSummaryZero = sum (map (balanceZero . xPubBal) bs)- , xPubSummaryPaths = ps- , xPubSummaryTxs = txs- , xPubChangeIndex = ch- , xPubExternalIndex = ex- }---- | Check if any of the ancestors of this transaction is a coinbase after the--- specified height. Returns 'Nothing' if answer cannot be computed before--- hitting limits.-cbAfterHeight ::- (Monad m, StoreRead m)- => Int -- ^ how many ancestors to test before giving up- -> BlockHeight- -> TxHash- -> m TxAfterHeight-cbAfterHeight d h t- | d <= 0 = return $ TxAfterHeight Nothing- | otherwise = TxAfterHeight <$> runMaybeT (snd <$> tst d t)- where- tst e x- | e <= 0 = MaybeT $ return Nothing- | otherwise = do- let e' = e - 1- tx <- MaybeT $ getTransaction x- if any isCoinbase (transactionInputs tx)- then return (e', blockRefHeight (transactionBlock tx) > h)- else case transactionBlock tx of- BlockRef {blockRefHeight = b}- | b <= h -> return (e', False)- _ ->- r e' . nub $- map- (outPointHash . inputPoint)- (transactionInputs tx)- r e [] = return (e, False)- r e (n:ns) = do- (e', s) <- tst e n- if s- then return (e', True)- else r e' ns---- Snatched from:--- https://github.com/cblp/conduit-merge/blob/master/src/Data/Conduit/Merge.hs-mergeSourcesBy ::- (Foldable f, Monad m)- => (a -> a -> Ordering)- -> f (ConduitT () a m ())- -> ConduitT i a m ()-mergeSourcesBy f = mergeSealed . fmap sealConduitT . toList- where- mergeSealed sources = do- prefetchedSources <- lift $ traverse ($$++ await) sources- go [(a, s) | (s, Just a) <- prefetchedSources]- go [] = pure ()- go sources = do- let (a, src1):sources1 = sortBy (f `on` fst) sources- yield a- (src2, mb) <- lift $ src1 $$++ await- let sources2 =- case mb of- Nothing -> sources1- Just b -> (b, src2) : sources1- go sources2--getMempoolLimit ::- (Monad m, StoreStream m)- => Maybe Word32- -> StartFrom- -> ConduitT () TxHash m ()-getMempoolLimit _ StartBlock {} = return ()-getMempoolLimit l (StartMem t) =- getMempool (Just t) .| mapC snd .| limit l-getMempoolLimit l s =- getMempool Nothing .| mapC snd .| (offset s >> limit l)--getAddressTxsLimit ::- (Monad m, StoreStream m)- => Maybe Word32- -> StartFrom- -> Address- -> ConduitT () BlockTx m ()-getAddressTxsLimit l s a =- getAddressTxs a (mbr s) .| (offset s >> limit l)--getAddressTxsFull ::- (Monad m, StoreStream m, StoreRead m)- => Maybe Word32- -> StartFrom- -> Address- -> ConduitT () Transaction m ()-getAddressTxsFull l s a =- getAddressTxsLimit l s a .| concatMapMC (getTransaction . blockTxHash)--getAddressesTxsLimit ::- (Monad m, StoreStream m)- => Maybe Word32- -> StartFrom- -> [Address]- -> ConduitT () BlockTx m ()-getAddressesTxsLimit l s as =- mergeSourcesBy- (flip compare `on` blockTxBlock)- (map (`getAddressTxs` mbr s) as) .|- dedup .|- (offset s >> limit l)--getAddressesTxsFull ::- (Monad m, StoreStream m, StoreRead m)- => Maybe Word32- -> StartFrom- -> [Address]- -> ConduitT () Transaction m ()-getAddressesTxsFull l s as =- mergeSourcesBy- (flip compare `on` blockTxBlock)- (map (`getAddressTxs` mbr s) as) .|- dedup .|- (offset s >> limit l) .|- concatMapMC (getTransaction . blockTxHash)--getAddressUnspentsLimit ::- (Monad m, StoreStream m)- => Maybe Word32- -> StartFrom- -> Address- -> ConduitT () Unspent m ()-getAddressUnspentsLimit l s a =- getAddressUnspents a (mbr s) .| (offset s >> limit l)--getAddressesUnspentsLimit ::- (Monad m, StoreStream m)- => Maybe Word32- -> StartFrom- -> [Address]- -> ConduitT () Unspent m ()-getAddressesUnspentsLimit l s as =- mergeSourcesBy- (flip compare `on` unspentBlock)- (map (`getAddressUnspents` mbr s) as) .|- (offset s >> limit l)--offset :: Monad m => StartFrom -> ConduitT i i m ()-offset (StartOffset o) = dropC (fromIntegral o)-offset _ = return ()--limit :: Monad m => Maybe Word32 -> ConduitT i i m ()-limit Nothing = mapC id-limit (Just n) = takeC (fromIntegral n)--mbr :: StartFrom -> Maybe BlockRef-mbr (StartBlock h p) = Just (BlockRef h p)-mbr (StartMem t) = Just (MemRef t)-mbr (StartOffset _) = Nothing--dedup :: (Eq i, Monad m) => ConduitT i i m ()-dedup =- let dd Nothing =- await >>= \case- Just x -> do- yield x- dd (Just x)- Nothing -> return ()- dd (Just x) =- await >>= \case- Just y- | x == y -> dd (Just x)- | otherwise -> do- yield y- dd (Just y)- Nothing -> return ()- in dd Nothing
src/Network/Haskoin/Store/Block.hs view
@@ -1,5 +1,6 @@ {-# LANGUAGE DeriveAnyClass #-} {-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE FlexibleInstances #-} {-# LANGUAGE LambdaCase #-} {-# LANGUAGE MultiWayIf #-} {-# LANGUAGE OverloadedStrings #-}@@ -9,6 +10,7 @@ ( blockStore ) where +import Conduit import Control.Arrow import Control.Monad.Except import Control.Monad.Logger@@ -24,8 +26,6 @@ import Haskoin.Node import Network.Haskoin.Store.Data import Network.Haskoin.Store.Data.ImportDB-import Network.Haskoin.Store.Data.RocksDB-import Network.Haskoin.Store.Data.STM import Network.Haskoin.Store.Logic import Network.Haskoin.Store.Messages import NQE@@ -50,13 +50,44 @@ -- | Block store process state. data BlockRead = BlockRead- { mySelf :: !BlockStore- , myConfig :: !BlockConfig- , myPeer :: !(TVar (Maybe Syncing))- , myUnspent :: !(TVar UnspentMap)- , myBalances :: !(TVar BalanceMap)+ { mySelf :: !BlockStore+ , myConfig :: !BlockConfig+ , myPeer :: !(TVar (Maybe Syncing)) } +type BlockT m = ReaderT BlockRead m++runImport ::+ MonadLoggerIO m+ => ReaderT ImportDB (ExceptT ImportException m) a+ -> ReaderT BlockRead m (Either ImportException a)+runImport f =+ ReaderT $ \r -> runExceptT (runImportDB (blockConfDB (myConfig r)) f)++runLayered :: ReaderT LayeredDB m a -> ReaderT BlockRead m a+runLayered f = ReaderT $ \r -> runReaderT f (blockConfDB (myConfig r))++instance MonadIO m => StoreRead (ReaderT BlockRead m) where+ isInitialized = runLayered isInitialized+ getBestBlock = runLayered getBestBlock+ getBlocksAtHeight = runLayered . getBlocksAtHeight+ getBlock = runLayered . getBlock+ getTxData = runLayered . getTxData+ getSpender = runLayered . getSpender+ getSpenders = runLayered . getSpenders+ getOrphanTx = runLayered . getOrphanTx+ getUnspent = runLayered . getUnspent+ getBalance = runLayered . getBalance++instance (MonadResource m, MonadUnliftIO m) =>+ StoreStream (ReaderT BlockRead m) where+ getMempool = transPipe runLayered . getMempool+ getOrphans = transPipe runLayered getOrphans+ getAddressUnspents a x = transPipe runLayered $ getAddressUnspents a x+ getAddressTxs a x = transPipe runLayered $ getAddressTxs a x+ getAddressBalances = transPipe runLayered getAddressBalances+ getUnspents = transPipe runLayered getUnspents+ -- | Run block store process. blockStore :: (MonadUnliftIO m, MonadLoggerIO m)@@ -66,22 +97,13 @@ blockStore cfg inbox = do $(logInfoS) "Block" "Initializing block store..." pb <- newTVarIO Nothing- um <- newTVarIO M.empty- bm <- newTVarIO (M.empty, []) runReaderT (ini >> run)- BlockRead- { mySelf = inboxToMailbox inbox- , myConfig = cfg- , myPeer = pb- , myUnspent = um- , myBalances = bm- }+ BlockRead {mySelf = inboxToMailbox inbox, myConfig = cfg, myPeer = pb} where ini = do- (db, net) <- (blockConfDB &&& blockConfNet) <$> asks myConfig- (um, bm) <- asks (myUnspent &&& myBalances)- runExceptT (initDB net db um bm) >>= \case+ net <- asks (blockConfNet . myConfig)+ runImport (initDB net) >>= \case Left e -> do $(logErrorS) "Block" $ "Could not initialize block store: " <> fromString (show e)@@ -90,56 +112,34 @@ run = withAsync (pingMe (inboxToMailbox inbox)) . const . forever $ do $(logDebugS) "Block" "Awaiting message..."- receive inbox >>= processBlockMessage+ receive inbox >>= \msg ->+ ReaderT $ \r ->+ runResourceT (runReaderT (processBlockMessage msg) r) isSynced :: (MonadLoggerIO m, MonadUnliftIO m) => ReaderT BlockRead m Bool isSynced = do- (db, ch) <- (blockConfDB &&& blockConfChain) <$> asks myConfig+ ch <- asks (blockConfChain . myConfig) $(logDebugS) "Block" "Testing if synced with header chain..."- withBlockDB defaultReadOptions db $- getBestBlock >>= \case- Nothing -> do- $(logErrorS) "Block" "Block database uninitialized"- throwIO Uninitialized- Just bb -> do- $(logDebugS) "Block" $ "Best block: " <> blockHashToHex bb- chainGetBest ch >>= \cb -> do- $(logDebugS) "Block" $- "Best chain block " <>- blockHashToHex (headerHash (nodeHeader cb)) <>- " at height " <>- cs (show (nodeHeight cb))- let s = headerHash (nodeHeader cb) == bb- $(logDebugS) "Block" $ "Synced: " <> cs (show s)- return s+ getBestBlock >>= \case+ Nothing -> do+ $(logErrorS) "Block" "Block database uninitialized"+ throwIO Uninitialized+ Just bb -> do+ $(logDebugS) "Block" $ "Best block: " <> blockHashToHex bb+ chainGetBest ch >>= \cb -> do+ $(logDebugS) "Block" $+ "Best chain block " <>+ blockHashToHex (headerHash (nodeHeader cb)) <>+ " at height " <>+ cs (show (nodeHeight cb))+ let s = headerHash (nodeHeader cb) == bb+ $(logDebugS) "Block" $ "Synced: " <> cs (show s)+ return s mempool :: (MonadUnliftIO m, MonadLoggerIO m) => Peer -> ReaderT BlockRead m () mempool p = MMempool `sendMessage` p -pruneCache :: (MonadUnliftIO m, MonadLoggerIO m) => ReaderT BlockRead m ()-pruneCache = do- um <- asks myUnspent- bm <- asks myBalances- do u <- readTVarIO um- b <- readTVarIO bm- $(logDebugS) "Block" $- "Unspent output cache pre-prune tx count: " <>- fromString (show (M.size u))- $(logDebugS) "Block" $- "Address cache pre-prune count: " <>- fromString (show (M.size (fst b)))- atomically $- withUnspentSTM um pruneUnspent >> withBalanceSTM bm pruneBalance- do u <- readTVarIO um- b <- readTVarIO bm- $(logDebugS) "Block" $- "Unspent output cache post-prune tx count: " <>- fromString (show (M.size u))- $(logDebugS) "Block" $- "Address cache post-prune count: " <>- fromString (show (M.size (fst b)))- processBlock :: (MonadUnliftIO m, MonadLoggerIO m) => Peer@@ -154,13 +154,10 @@ "Block" ("Cannot accept block " <> hex <> " from non-syncing peer") mzero- db <- blockConfDB <$> asks myConfig n <- cbn upr net <- blockConfNet <$> asks myConfig- um <- asks myUnspent- bm <- asks myBalances- runExceptT (runImportDB db um bm $ importBlock net b n) >>= \case+ lift (runImport (importBlock net b n)) >>= \case Right () -> do l <- blockConfListener <$> asks myConfig $(logInfoS) "Block" $@@ -170,7 +167,6 @@ cs (show (nodeHeight n)) atomically $ l (StoreBestBlock (headerHash (blockHeader b))) lift $ isSynced >>= \x -> when x (mempool p)- when (nodeHeight n `mod` 1000 == 0) (lift pruneCache) Left e -> do $(logErrorS) "Block" $ "Error importing block " <>@@ -220,11 +216,7 @@ (PeerMisbehaving "We do not like peers that cannot find them blocks") p -processTx ::- (MonadUnliftIO m, MonadLoggerIO m)- => Peer- -> Tx- -> ReaderT BlockRead m ()+processTx :: (MonadUnliftIO m, MonadLoggerIO m) => Peer -> Tx -> BlockT m () processTx _p tx = isSynced >>= \case False ->@@ -233,10 +225,8 @@ True -> do $(logInfoS) "Block" $ "Incoming tx: " <> txHashToHex (txHash tx) now <- fromIntegral . systemSeconds <$> liftIO getSystemTime- (net, db) <- (blockConfNet &&& blockConfDB) <$> asks myConfig- um <- asks myUnspent- bm <- asks myBalances- runExceptT (runImportDB db um bm $ newMempoolTx net tx now) >>= \case+ net <- asks (blockConfNet . myConfig)+ runImport (newMempoolTx net tx now) >>= \case Left e -> $(logErrorS) "Block" $ "Error importing tx: " <> txHashToHex (txHash tx) <> ": " <>@@ -250,6 +240,36 @@ $(logDebugS) "Block" $ "Not importing mempool tx: " <> txHashToHex (txHash tx) +processOrphans ::+ (MonadResource m, MonadUnliftIO m, MonadLoggerIO m) => BlockT m ()+processOrphans =+ isSynced >>= \case+ False -> $(logDebugS) "Block" "Not importing orphans as not yet in sync"+ True -> do+ now <- fromIntegral . systemSeconds <$> liftIO getSystemTime+ (ldb, net) <- asks ((blockConfDB &&& blockConfNet) . myConfig)+ $(logDebugS) "Block" "Getting expired orphan transactions..."+ old <- runConduit $ getOldOrphans now .| sinkList+ case old of+ [] ->+ $(logDebugS) "Block" "No old orphan transactions to remove"+ _ -> do+ $(logDebugS) "Block" $+ "Removing " <> cs (show (length old)) <>+ "expired orphan transactions..."+ runImport $ mapM_ deleteOrphanTx old+ return ()+ $(logDebugS) "Block" "Selecting orphan transactions to import..."+ orphans <- runConduit $ getOrphans .| sinkList+ case orphans of+ [] -> $(logDebugS) "Block" "No orphan tranasctions to import"+ _ ->+ $(logDebugS) "Block" $+ "Importing " <> cs (show (length orphans)) <>+ " orphan transactions"+ forM_ orphans $ runImport . uncurry (importOrphan net)+ $(logDebugS) "Block" $ "Finished importing orphans"+ processTxs :: (MonadUnliftIO m, MonadLoggerIO m) => Peer@@ -260,14 +280,13 @@ False -> $(logDebugS) "Block" "Ignoring incoming tx inv (not synced yet)" True -> do- db <- blockConfDB <$> asks myConfig $(logDebugS) "Block" $ "Received " <> fromString (show (length hs)) <> " transaction inventory" xs <- fmap catMaybes . forM hs $ \h -> runMaybeT $ do- t <- withBlockDB defaultReadOptions db $ getTxData h+ t <- lift $ getTxData h guard (isNothing t) return (getTxHash h) unless (null xs) $ do@@ -354,9 +373,8 @@ OnlinePeer {onlinePeerMailbox = p}:_ -> return (Just p) cbn = chainGetBest =<< blockConfChain <$> asks myConfig dbn = do- db <- blockConfDB <$> asks myConfig bb <-- withBlockDB defaultReadOptions db getBestBlock >>= \case+ lift getBestBlock >>= \case Nothing -> do $(logErrorS) "Block" "Best block not in database" throwIO Uninitialized@@ -413,11 +431,8 @@ "Reverting best block " <> blockHashToHex d <> " as it is not in main chain..." resetPeer- db <- blockConfDB <$> asks myConfig- um <- asks myUnspent- bm <- asks myBalances- net <- blockConfNet <$> asks myConfig- runExceptT (runImportDB db um bm $ revertBlock net d) >>= \case+ net <- asks (blockConfNet . myConfig)+ lift (runImport (revertBlock net d)) >>= \case Left e -> do $(logErrorS) "Block" $ "Could not revert best block: " <>@@ -439,9 +454,9 @@ Just Syncing {syncingPeer = p, syncingHead = b, syncingTime = now} processBlockMessage ::- (MonadUnliftIO m, MonadLoggerIO m)+ (MonadResource m, MonadUnliftIO m, MonadLoggerIO m) => BlockMessage- -> ReaderT BlockRead m ()+ -> BlockT m () processBlockMessage (BlockNewBest bn) = do $(logDebugS) "Block" $ "New best block header " <> fromString (show (nodeHeight bn)) <> ": " <>@@ -455,7 +470,8 @@ processBlockMessage (BlockNotFound p bs) = processNoBlocks p bs processBlockMessage (BlockTxReceived p tx) = processTx p tx processBlockMessage (BlockTxAvailable p ts) = processTxs p ts-processBlockMessage (BlockPing r) = checkTime >> atomically (r ())+processBlockMessage (BlockPing r) =+ processOrphans >> checkTime >> atomically (r ()) processBlockMessage PurgeMempool = purgeMempool pingMe :: MonadLoggerIO m => Mailbox BlockMessage -> m ()
src/Network/Haskoin/Store/Data.hs view
@@ -16,9 +16,10 @@ import qualified Data.ByteString as B import Data.ByteString.Short (ShortByteString) import qualified Data.ByteString.Short as B.Short+import Data.Default import Data.Hashable import Data.HashMap.Strict (HashMap)-import qualified Data.HashMap.Strict as H+import qualified Data.HashMap.Strict as M import Data.Int import qualified Data.IntMap as I import Data.IntMap.Strict (IntMap)@@ -29,33 +30,40 @@ import qualified Data.Text.Encoding as T import qualified Data.Text.Lazy as T.Lazy import Data.Word+import Database.RocksDB (DB, ReadOptions) import GHC.Generics import Haskoin as H import Network.Socket (SockAddr) import Paths_haskoin_store as P+import UnliftIO import UnliftIO.Exception import qualified Web.Scotty.Trans as Scotty -type UnixTime = Word64-type BlockPos = Word32+encodeShort :: Serialize a => a -> ShortByteString+encodeShort = B.Short.toShort . S.encode -newtype InitException = IncorrectVersion Word32- deriving (Show, Read, Eq, Ord, Exception)+decodeShort :: Serialize a => ShortByteString -> a+decodeShort bs = case S.decode (B.Short.fromShort bs) of+ Left e -> error e+ Right a -> a -class UnspentWrite m where- addUnspent :: Unspent -> m ()- delUnspent :: OutPoint -> m ()- pruneUnspent :: m ()+data BlockDB =+ BlockDB+ { blockDB :: !DB+ , blockDBopts :: !ReadOptions+ } -class UnspentRead m where- getUnspent :: OutPoint -> m (Maybe Unspent)+data LayeredDB =+ LayeredDB+ { layeredDB :: !BlockDB+ , layeredCache :: !(Maybe BlockDB)+ } -class BalanceWrite m where- setBalance :: Balance -> m ()- pruneBalance :: m ()+type UnixTime = Word64+type BlockPos = Word32 -class BalanceRead m where- getBalance :: Address -> m (Maybe Balance)+newtype InitException = IncorrectVersion Word32+ deriving (Show, Read, Eq, Ord, Exception) class StoreRead m where isInitialized :: m (Either InitException Bool)@@ -66,19 +74,8 @@ getOrphanTx :: TxHash -> m (Maybe (UnixTime, Tx)) getSpenders :: TxHash -> m (IntMap Spender) getSpender :: OutPoint -> m (Maybe Spender)--getTransaction ::- (Monad m, StoreRead m) => TxHash -> m (Maybe Transaction)-getTransaction h = runMaybeT $ do- d <- MaybeT $ getTxData h- sm <- lift $ getSpenders h- return $ toTransaction d sm--class StoreStream m where- getMempool :: Maybe UnixTime -> ConduitT () (UnixTime, TxHash) m ()- getOrphans :: ConduitT () (UnixTime, Tx) m ()- getAddressUnspents :: Address -> Maybe BlockRef -> ConduitT () Unspent m ()- getAddressTxs :: Address -> Maybe BlockRef -> ConduitT () BlockTx m ()+ getBalance :: Address -> m (Maybe Balance)+ getUnspent :: OutPoint -> m (Maybe Unspent) class StoreWrite m where setInit :: m ()@@ -89,14 +86,32 @@ insertSpender :: OutPoint -> Spender -> m () deleteSpender :: OutPoint -> m () insertAddrTx :: Address -> BlockTx -> m ()- removeAddrTx :: Address -> BlockTx -> m ()+ deleteAddrTx :: Address -> BlockTx -> m () insertAddrUnspent :: Address -> Unspent -> m ()- removeAddrUnspent :: Address -> Unspent -> m ()+ deleteAddrUnspent :: Address -> Unspent -> m () insertMempoolTx :: TxHash -> UnixTime -> m () deleteMempoolTx :: TxHash -> UnixTime -> m () insertOrphanTx :: Tx -> UnixTime -> m () deleteOrphanTx :: TxHash -> m ()+ setBalance :: Balance -> m ()+ insertUnspent :: Unspent -> m ()+ deleteUnspent :: OutPoint -> m () +getTransaction ::+ (Monad m, StoreRead m) => TxHash -> m (Maybe Transaction)+getTransaction h = runMaybeT $ do+ d <- MaybeT $ getTxData h+ sm <- lift $ getSpenders h+ return $ toTransaction d sm++class StoreStream m where+ getMempool :: Maybe UnixTime -> ConduitT () (UnixTime, TxHash) m ()+ getOrphans :: ConduitT () (UnixTime, Tx) m ()+ getAddressUnspents :: Address -> Maybe BlockRef -> ConduitT () Unspent m ()+ getAddressTxs :: Address -> Maybe BlockRef -> ConduitT () BlockTx m ()+ getAddressBalances :: ConduitT () Balance m ()+ getUnspents :: ConduitT () Unspent m ()+ -- | Serialize such that ordering is inverted. putUnixTime w = putWord64be $ maxBound - w getUnixTime = (maxBound -) <$> getWord64be@@ -835,7 +850,7 @@ [ "balance" .= object ["received" .= r, "confirmed" .= c, "unconfirmed" .= z] , "indices" .= object ["change" .= ch, "external" .= ext]- , "paths" .= object (mapMaybe (uncurry f) (H.toList ps))+ , "paths" .= object (mapMaybe (uncurry f) (M.toList ps)) , "txs" .= map (transactionToJSON net) ts ] where@@ -865,8 +880,8 @@ put z put ext put ch- putWord64be (fromIntegral $ H.size ps)- forM_ (H.toList ps) $ \(a, p) -> do+ putWord64be (fromIntegral $ M.size ps)+ forM_ (M.toList ps) $ \(a, p) -> do binSerial net a put p putWord64be (fromIntegral $ length ts)@@ -1002,3 +1017,83 @@ | StartMem !UnixTime | StartOffset !Word32 deriving (Show, Eq, Generic)++data BalVal = BalVal+ { balValAmount :: !Word64+ , balValZero :: !Word64+ , balValUnspentCount :: !Word64+ , balValTxCount :: !Word64+ , balValTotalReceived :: !Word64+ } deriving (Show, Read, Eq, Ord, Generic, Hashable, Serialize)++balValToBalance :: Address -> BalVal -> Balance+balValToBalance a BalVal { balValAmount = v+ , balValZero = z+ , balValUnspentCount = u+ , balValTxCount = t+ , balValTotalReceived = r+ } =+ Balance+ { balanceAddress = a+ , balanceAmount = v+ , balanceZero = z+ , balanceUnspentCount = u+ , balanceTxCount = t+ , balanceTotalReceived = r+ }++balanceToBalVal :: Balance -> (Address, BalVal)+balanceToBalVal Balance { balanceAddress = a+ , balanceAmount = v+ , balanceZero = z+ , balanceUnspentCount = u+ , balanceTxCount = t+ , balanceTotalReceived = r+ } =+ ( a+ , BalVal+ { balValAmount = v+ , balValZero = z+ , balValUnspentCount = u+ , balValTxCount = t+ , balValTotalReceived = r+ })++-- | Default balance for an address.+instance Default BalVal where+ def =+ BalVal+ { balValAmount = 0+ , balValZero = 0+ , balValUnspentCount = 0+ , balValTxCount = 0+ , balValTotalReceived = 0+ }++data UnspentVal = UnspentVal+ { unspentValBlock :: !BlockRef+ , unspentValAmount :: !Word64+ , unspentValScript :: !ShortByteString+ } deriving (Show, Read, Eq, Ord, Generic, Hashable, Serialize)++unspentToUnspentVal :: Unspent -> (OutPoint, UnspentVal)+unspentToUnspentVal Unspent { unspentBlock = b+ , unspentPoint = p+ , unspentAmount = v+ , unspentScript = s+ } =+ ( p+ , UnspentVal+ {unspentValBlock = b, unspentValAmount = v, unspentValScript = s})++unspentValToUnspent :: OutPoint -> UnspentVal -> Unspent+unspentValToUnspent p UnspentVal { unspentValBlock = b+ , unspentValAmount = v+ , unspentValScript = s+ } =+ Unspent+ { unspentBlock = b+ , unspentPoint = p+ , unspentAmount = v+ , unspentScript = s+ }
+ src/Network/Haskoin/Store/Data/Cached.hs view
@@ -0,0 +1,198 @@+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE FlexibleInstances #-}+{-# LANGUAGE LambdaCase #-}+module Network.Haskoin.Store.Data.Cached where++import Conduit+import Control.Applicative+import Control.Monad.Except+import Control.Monad.Logger+import Control.Monad.Reader (MonadReader, ReaderT)+import qualified Control.Monad.Reader as R+import Control.Monad.Trans.Maybe+import qualified Data.ByteString as B+import qualified Data.ByteString.Short as B.Short+import Data.IntMap.Strict (IntMap)+import Data.List+import Data.Maybe+import Data.Serialize (Serialize, encode)+import Data.String.Conversions (cs)+import Database.RocksDB as R+import Database.RocksDB.Query as R+import Haskoin+import Network.Haskoin.Store.Data+import Network.Haskoin.Store.Data.KeyValue+import Network.Haskoin.Store.Data.RocksDB+import NQE (query)+import UnliftIO++newLayeredDB :: MonadUnliftIO m => BlockDB -> Maybe BlockDB -> m LayeredDB+newLayeredDB blocks Nothing =+ return LayeredDB {layeredDB = blocks, layeredCache = Nothing}+newLayeredDB blocks (Just cache) = do+ bulkCopy opts db cdb BalKeyS+ bulkCopy opts db cdb UnspentKeyB+ return LayeredDB {layeredDB = blocks, layeredCache = Just cache}+ where+ BlockDB {blockDBopts = opts, blockDB = db} = blocks+ BlockDB {blockDB = cdb} = cache++withLayeredDB :: LayeredDB -> ReaderT LayeredDB m a -> m a+withLayeredDB = flip R.runReaderT++isInitializedC :: MonadIO m => LayeredDB -> m (Either InitException Bool)+isInitializedC LayeredDB {layeredDB = db} = isInitializedDB db++getBestBlockC :: MonadIO m => LayeredDB -> m (Maybe BlockHash)+getBestBlockC LayeredDB {layeredDB = db} = getBestBlockDB db++getBlocksAtHeightC :: MonadIO m => BlockHeight -> LayeredDB -> m [BlockHash]+getBlocksAtHeightC h LayeredDB {layeredDB = db} = getBlocksAtHeightDB h db++getBlockC :: MonadIO m => BlockHash -> LayeredDB -> m (Maybe BlockData)+getBlockC bh LayeredDB {layeredDB = db} = getBlockDB bh db++getTxDataC :: MonadIO m => TxHash -> LayeredDB -> m (Maybe TxData)+getTxDataC th LayeredDB {layeredDB = db} = getTxDataDB th db++getOrphanTxC :: MonadIO m => TxHash -> LayeredDB -> m (Maybe (UnixTime, Tx))+getOrphanTxC h LayeredDB {layeredDB = db} = getOrphanTxDB h db++getSpenderC :: MonadIO m => OutPoint -> LayeredDB -> m (Maybe Spender)+getSpenderC p LayeredDB {layeredDB = db} = getSpenderDB p db++getSpendersC :: MonadIO m => TxHash -> LayeredDB -> m (IntMap Spender)+getSpendersC t LayeredDB {layeredDB = db} = getSpendersDB t db++getBalanceC :: MonadIO m => Address -> LayeredDB -> m (Maybe Balance)+getBalanceC a LayeredDB {layeredCache = Just db} = getBalanceDB a db+getBalanceC a LayeredDB {layeredDB = db} = getBalanceDB a db++getUnspentC :: MonadIO m => OutPoint -> LayeredDB -> m (Maybe Unspent)+getUnspentC op LayeredDB {layeredCache = Just db} = getUnspentDB op db+getUnspentC op LayeredDB {layeredDB = db} = getUnspentDB op db++getUnspentsC ::+ (MonadResource m, MonadIO m) => LayeredDB -> ConduitT () Unspent m ()+getUnspentsC LayeredDB {layeredDB = db} = getUnspentsDB db++getMempoolC ::+ (MonadResource m, MonadUnliftIO m)+ => Maybe UnixTime+ -> LayeredDB+ -> ConduitT () (UnixTime, TxHash) m ()+getMempoolC mpu LayeredDB {layeredDB = db} = getMempoolDB mpu db++getOrphansC ::+ (MonadUnliftIO m, MonadResource m)+ => LayeredDB+ -> ConduitT () (UnixTime, Tx) m ()+getOrphansC LayeredDB {layeredDB = db} = getOrphansDB db++getAddressBalancesC ::+ (MonadUnliftIO m, MonadResource m)+ => LayeredDB+ -> ConduitT () Balance m ()+getAddressBalancesC LayeredDB {layeredDB = db} = getAddressBalancesDB db++getAddressUnspentsC ::+ (MonadUnliftIO m, MonadResource m)+ => Address+ -> Maybe BlockRef+ -> LayeredDB+ -> ConduitT () Unspent m ()+getAddressUnspentsC addr mbr LayeredDB {layeredDB = db} =+ getAddressUnspentsDB addr mbr db++getAddressTxsC ::+ (MonadUnliftIO m, MonadResource m)+ => Address+ -> Maybe BlockRef+ -> LayeredDB+ -> ConduitT () BlockTx m ()+getAddressTxsC addr mbr LayeredDB {layeredDB = db} =+ getAddressTxsDB addr mbr db++bulkCopy ::+ (Serialize k, MonadUnliftIO m) => ReadOptions -> DB -> DB -> k -> m ()+bulkCopy opts db cdb k =+ runResourceT $ do+ ch <- newTBQueueIO 1000000+ withAsync (iterate ch) $ \a -> write_batch ch [] 0+ where+ iterate ch =+ withIterator db opts $ \it -> do+ iterSeek it (encode k)+ recurse it ch+ write_batch ch acc l+ | l >= 10000 = do+ write cdb defaultWriteOptions acc+ write_batch ch [] 0+ | otherwise =+ atomically (readTBQueue ch) >>= \case+ Just (key, val) -> write_batch ch (Put key val : acc) (l + 1)+ Nothing -> write cdb defaultWriteOptions acc+ recurse it ch =+ iterEntry it >>= \case+ Nothing -> atomically $ writeTBQueue ch Nothing+ Just (key, val) ->+ let pfx = B.take (B.length (encode k)) key+ in if pfx == encode k+ then do+ atomically . writeTBQueue ch $ Just (key, val)+ iterNext it+ recurse it ch+ else atomically $ writeTBQueue ch Nothing++instance (MonadUnliftIO m, MonadResource m) =>+ StoreStream (ReaderT LayeredDB m) where+ getMempool x = do+ c <- R.ask+ getMempoolC x c+ getOrphans = do+ c <- R.ask+ getOrphansC c+ getAddressUnspents a x = do+ c <- R.ask+ getAddressUnspentsC a x c+ getAddressTxs a x = do+ c <- R.ask+ getAddressTxsC a x c+ getAddressBalances = do+ c <- R.ask+ getAddressBalancesC c+ getUnspents = do+ c <- R.ask+ getUnspentsC c++instance MonadIO m => StoreRead (ReaderT LayeredDB m) where+ isInitialized = do+ c <- R.ask+ isInitializedC c+ getBestBlock = do+ c <- R.ask+ getBestBlockC c+ getBlocksAtHeight h = do+ c <- R.ask+ getBlocksAtHeightC h c+ getBlock b = do+ c <- R.ask+ getBlockC b c+ getTxData t = do+ c <- R.ask+ getTxDataC t c+ getSpender p = do+ c <- R.ask+ getSpenderC p c+ getSpenders t = do+ c <- R.ask+ getSpendersC t c+ getOrphanTx h = do+ c <- R.ask+ getOrphanTxC h c+ getUnspent a = do+ c <- R.ask+ getUnspentC a c+ getBalance a = do+ c <- R.ask+ getBalanceC a c
src/Network/Haskoin/Store/Data/ImportDB.hs view
@@ -1,11 +1,15 @@ {-# LANGUAGE FlexibleContexts #-} {-# LANGUAGE FlexibleInstances #-} {-# LANGUAGE LambdaCase #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE TupleSections #-} module Network.Haskoin.Store.Data.ImportDB where import Conduit import Control.Applicative import Control.Monad.Except+import Control.Monad.Logger import Control.Monad.Reader (ReaderT) import qualified Control.Monad.Reader as R import Control.Monad.Trans.Maybe@@ -16,50 +20,48 @@ import qualified Data.IntMap.Strict as I import Data.List import Data.Maybe+import Data.String.Conversions (cs) import Database.RocksDB as R import Database.RocksDB.Query as R import Haskoin import Network.Haskoin.Store.Data+import Network.Haskoin.Store.Data.Cached import Network.Haskoin.Store.Data.KeyValue+import Network.Haskoin.Store.Data.Memory import Network.Haskoin.Store.Data.RocksDB-import Network.Haskoin.Store.Data.STM import UnliftIO data ImportDB = ImportDB- { importRocksDB :: !(ReadOptions, DB)- , importHashMap :: !(TVar HashMapDB)- , importUnspentMap :: !(TVar UnspentMap)- , importBalanceMap :: !(TVar BalanceMap)+ { importLayeredDB :: !LayeredDB+ , importHashMap :: !(TVar BlockMem) } runImportDB ::- (MonadError e m, MonadIO m)- => DB- -> TVar UnspentMap- -> TVar BalanceMap+ (MonadError e m, MonadLoggerIO m)+ => LayeredDB -> ReaderT ImportDB m a -> m a-runImportDB db um bm f = do- hm <- newTVarIO emptyHashMapDB- um' <- atomically $ readTVar um >>= newTVar- bm' <- atomically $ readTVar bm >>= newTVar- x <-- R.runReaderT- f- ImportDB- { importRocksDB = (defaultReadOptions, db)- , importHashMap = hm- , importUnspentMap = um'- , importBalanceMap = bm'- }+runImportDB ldb f = do+ hm <- newTVarIO emptyBlockMem+ x <- R.runReaderT f ImportDB {importLayeredDB = ldb, importHashMap = hm} ops <- hashMapOps <$> readTVarIO hm- writeBatch db ops- atomically $ do- readTVar um' >>= writeTVar um- readTVar bm' >>= writeTVar bm+ $(logDebugS) "ImportDB" "Committing changes to database and cache..."+ case cache of+ Just BlockDB {blockDB = cdb} -> do+ cops <- cacheMapOps <$> readTVarIO hm+ let del Put {} = False+ del Del {} = True+ (delcops, addcops) = partition del cops+ writeBatch cdb delcops+ writeBatch db ops+ writeBatch cdb addcops+ Nothing -> writeBatch db ops+ $(logDebugS) "ImportDB" "Finished committing changes to database and cache" return x+ where+ LayeredDB {layeredDB = BlockDB {blockDB = db}, layeredCache = cache} = ldb -hashMapOps :: HashMapDB -> [BatchOp]+hashMapOps :: BlockMem -> [BatchOp] hashMapOps db = bestBlockOp (hBest db) <> blockHashOps (hBlock db) <>@@ -73,6 +75,11 @@ orphanOps (hOrphans db) <> unspentOps (hUnspent db) +cacheMapOps :: BlockMem -> [BatchOp]+cacheMapOps db =+ balOps (hBalance db) <>+ unspentOps (hUnspent db)+ bestBlockOp :: Maybe BlockHash -> [BatchOp] bestBlockOp Nothing = [] bestBlockOp (Just b) = [insertOp BestKey b]@@ -158,179 +165,157 @@ orphanOps = map (uncurry f) . M.toList where f h (Just x) = insertOp (OrphanKey h) x- f h Nothing = deleteOp (OrphanKey h)+ f h Nothing = deleteOp (OrphanKey h) -unspentOps :: HashMap TxHash (IntMap (Maybe Unspent)) -> [BatchOp]+unspentOps :: HashMap TxHash (IntMap (Maybe UnspentVal)) -> [BatchOp] unspentOps = concatMap (uncurry f) . M.toList where f h = map (uncurry (g h)) . I.toList- g h i (Just u) =- insertOp- (UnspentKey (OutPoint h (fromIntegral i)))- UnspentVal- { unspentValAmount = unspentAmount u- , unspentValBlock = unspentBlock u- , unspentValScript = B.Short.fromShort (unspentScript u)- }- g h i Nothing = deleteOp (UnspentKey (OutPoint h (fromIntegral i)))+ g h i (Just u) = insertOp (UnspentKey (OutPoint h (fromIntegral i))) u+ g h i Nothing = deleteOp (UnspentKey (OutPoint h (fromIntegral i))) isInitializedI :: MonadIO m => ImportDB -> m (Either InitException Bool)-isInitializedI ImportDB {importRocksDB = db} =- uncurry withBlockDB db isInitialized+isInitializedI ImportDB {importLayeredDB = ldb} =+ withLayeredDB ldb isInitialized setInitI :: MonadIO m => ImportDB -> m ()-setInitI ImportDB {importRocksDB = (_, db), importHashMap = hm} = do- atomically $ withBlockSTM hm setInit+setInitI ImportDB { importLayeredDB = LayeredDB {layeredDB = BlockDB {blockDB = db}}+ , importHashMap = hm+ } = do+ withBlockMem hm setInit setInitDB db setBestI :: MonadIO m => BlockHash -> ImportDB -> m () setBestI bh ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ setBest bh+ withBlockMem hm $ setBest bh insertBlockI :: MonadIO m => BlockData -> ImportDB -> m () insertBlockI b ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ insertBlock b+ withBlockMem hm $ insertBlock b insertAtHeightI :: MonadIO m => BlockHash -> BlockHeight -> ImportDB -> m () insertAtHeightI b h ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ insertAtHeight b h+ withBlockMem hm $ insertAtHeight b h insertTxI :: MonadIO m => TxData -> ImportDB -> m () insertTxI t ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ insertTx t+ withBlockMem hm $ insertTx t insertSpenderI :: MonadIO m => OutPoint -> Spender -> ImportDB -> m () insertSpenderI p s ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ insertSpender p s+ withBlockMem hm $ insertSpender p s deleteSpenderI :: MonadIO m => OutPoint -> ImportDB -> m () deleteSpenderI p ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ deleteSpender p+ withBlockMem hm $ deleteSpender p insertAddrTxI :: MonadIO m => Address -> BlockTx -> ImportDB -> m () insertAddrTxI a t ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ insertAddrTx a t+ withBlockMem hm $ insertAddrTx a t -removeAddrTxI :: MonadIO m => Address -> BlockTx -> ImportDB -> m ()-removeAddrTxI a t ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ removeAddrTx a t+deleteAddrTxI :: MonadIO m => Address -> BlockTx -> ImportDB -> m ()+deleteAddrTxI a t ImportDB {importHashMap = hm} =+ withBlockMem hm $ deleteAddrTx a t insertAddrUnspentI :: MonadIO m => Address -> Unspent -> ImportDB -> m () insertAddrUnspentI a u ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ insertAddrUnspent a u+ withBlockMem hm $ insertAddrUnspent a u -removeAddrUnspentI :: MonadIO m => Address -> Unspent -> ImportDB -> m ()-removeAddrUnspentI a u ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ removeAddrUnspent a u+deleteAddrUnspentI :: MonadIO m => Address -> Unspent -> ImportDB -> m ()+deleteAddrUnspentI a u ImportDB {importHashMap = hm} =+ withBlockMem hm $ deleteAddrUnspent a u insertMempoolTxI :: MonadIO m => TxHash -> UnixTime -> ImportDB -> m () insertMempoolTxI t p ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ insertMempoolTx t p+ withBlockMem hm $ insertMempoolTx t p deleteMempoolTxI :: MonadIO m => TxHash -> UnixTime -> ImportDB -> m () deleteMempoolTxI t p ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ deleteMempoolTx t p+ withBlockMem hm $ deleteMempoolTx t p insertOrphanTxI :: MonadIO m => Tx -> UnixTime -> ImportDB -> m () insertOrphanTxI t p ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ insertOrphanTx t p+ withBlockMem hm $ insertOrphanTx t p deleteOrphanTxI :: MonadIO m => TxHash -> ImportDB -> m () deleteOrphanTxI t ImportDB {importHashMap = hm} =- atomically . withBlockSTM hm $ deleteOrphanTx t+ withBlockMem hm $ deleteOrphanTx t getBestBlockI :: MonadIO m => ImportDB -> m (Maybe BlockHash)-getBestBlockI ImportDB {importHashMap = hm, importRocksDB = db} =+getBestBlockI ImportDB {importHashMap = hm, importLayeredDB = db} = runMaybeT $ MaybeT f <|> MaybeT g where- f = atomically $ withBlockSTM hm getBestBlock- g = uncurry withBlockDB db getBestBlock+ f = withBlockMem hm getBestBlock+ g = withLayeredDB db getBestBlock getBlocksAtHeightI :: MonadIO m => BlockHeight -> ImportDB -> m [BlockHash]-getBlocksAtHeightI bh ImportDB {importHashMap = hm, importRocksDB = db} = do- xs <- atomically . withBlockSTM hm $ getBlocksAtHeight bh- ys <- uncurry withBlockDB db $ getBlocksAtHeight bh+getBlocksAtHeightI bh ImportDB {importHashMap = hm, importLayeredDB = db} = do+ xs <- withBlockMem hm $ getBlocksAtHeight bh+ ys <- withLayeredDB db $ getBlocksAtHeight bh return . nub $ xs <> ys getBlockI :: MonadIO m => BlockHash -> ImportDB -> m (Maybe BlockData)-getBlockI bh ImportDB {importRocksDB = db, importHashMap = hm} =+getBlockI bh ImportDB {importLayeredDB = db, importHashMap = hm} = runMaybeT $ MaybeT f <|> MaybeT g where- f = atomically . withBlockSTM hm $ getBlock bh- g = uncurry withBlockDB db $ getBlock bh+ f = withBlockMem hm $ getBlock bh+ g = withLayeredDB db $ getBlock bh getTxDataI :: MonadIO m => TxHash -> ImportDB -> m (Maybe TxData)-getTxDataI th ImportDB {importRocksDB = db, importHashMap = hm} =+getTxDataI th ImportDB {importLayeredDB = db, importHashMap = hm} = runMaybeT $ MaybeT f <|> MaybeT g where- f = atomically . withBlockSTM hm $ getTxData th- g = uncurry withBlockDB db $ getTxData th+ f = withBlockMem hm $ getTxData th+ g = withLayeredDB db $ getTxData th getOrphanTxI :: MonadIO m => TxHash -> ImportDB -> m (Maybe (UnixTime, Tx))-getOrphanTxI h ImportDB {importRocksDB = db, importHashMap = hm} =+getOrphanTxI h ImportDB {importLayeredDB = db, importHashMap = hm} = fmap join . runMaybeT $ MaybeT f <|> MaybeT g where f = getOrphanTxH h <$> readTVarIO hm- g = Just <$> uncurry withBlockDB db (getOrphanTx h)+ g = Just <$> withLayeredDB db (getOrphanTx h) getSpenderI :: MonadIO m => OutPoint -> ImportDB -> m (Maybe Spender)-getSpenderI op ImportDB {importRocksDB = db, importHashMap = hm} =- getSpenderH op <$> readTVarIO hm >>= \case- Just s -> return s- Nothing -> uncurry withBlockDB db $ getSpender op+getSpenderI op ImportDB {importLayeredDB = db, importHashMap = hm} =+ fmap join . runMaybeT $ MaybeT f <|> MaybeT g+ where+ f = getSpenderH op <$> readTVarIO hm+ g = Just <$> withLayeredDB db (getSpender op) getSpendersI :: MonadIO m => TxHash -> ImportDB -> m (IntMap Spender)-getSpendersI t ImportDB {importRocksDB = db, importHashMap = hm} = do+getSpendersI t ImportDB {importLayeredDB = db, importHashMap = hm} = do hsm <- getSpendersH t <$> readTVarIO hm- dsm <- I.map Just <$> uncurry withBlockDB db (getSpenders t)+ dsm <- I.map Just <$> withLayeredDB db (getSpenders t) return . I.map fromJust . I.filter isJust $ hsm <> dsm getBalanceI :: MonadIO m => Address -> ImportDB -> m (Maybe Balance)-getBalanceI a ImportDB { importRocksDB = db- , importHashMap = hm- , importBalanceMap = bm- } =- runMaybeT $- MaybeT (atomically . runMaybeT $ cachemap <|> hashmap) <|> database+getBalanceI a ImportDB {importLayeredDB = db, importHashMap = hm} =+ runMaybeT $ MaybeT f <|> MaybeT g where- cachemap = MaybeT . withBalanceSTM bm $ getBalance a- hashmap = MaybeT . withBlockSTM hm $ getBalance a- database = MaybeT . uncurry withBlockDB db $ getBalance a+ f = withBlockMem hm $ getBalance a+ g = withLayeredDB db $ getBalance a setBalanceI :: MonadIO m => Balance -> ImportDB -> m ()-setBalanceI b ImportDB {importHashMap = hm, importBalanceMap = bm} =- atomically $ do- withBlockSTM hm $ setBalance b- withBalanceSTM bm $ setBalance b+setBalanceI b ImportDB {importHashMap = hm} =+ withBlockMem hm $ setBalance b getUnspentI :: MonadIO m => OutPoint -> ImportDB -> m (Maybe Unspent)-getUnspentI op ImportDB { importRocksDB = db- , importHashMap = hm- , importUnspentMap = um- } = do- u <-- atomically . runMaybeT $ do- let x = withUnspentSTM um (getUnspent op)- y = getUnspentH op <$> readTVar hm- Just <$> MaybeT x <|> MaybeT y- case u of- Nothing -> uncurry withBlockDB db $ getUnspent op- Just x -> return x+getUnspentI op ImportDB {importLayeredDB = db, importHashMap = hm} =+ fmap join . runMaybeT $ MaybeT f <|> MaybeT g+ where+ f = getUnspentH op <$> readTVarIO hm+ g = Just <$> withLayeredDB db (getUnspent op) -addUnspentI :: MonadIO m => Unspent -> ImportDB -> m ()-addUnspentI u ImportDB {importHashMap = hm, importUnspentMap = um} =- atomically $ do- withBlockSTM hm $ addUnspent u- withUnspentSTM um $ addUnspent u+insertUnspentI :: MonadIO m => Unspent -> ImportDB -> m ()+insertUnspentI u ImportDB {importHashMap = hm} =+ withBlockMem hm $ insertUnspent u -delUnspentI :: MonadIO m => OutPoint -> ImportDB -> m ()-delUnspentI p ImportDB {importHashMap = hm, importUnspentMap = um} =- atomically $ do- withUnspentSTM um $ delUnspent p- withBlockSTM hm $ delUnspent p+deleteUnspentI :: MonadIO m => OutPoint -> ImportDB -> m ()+deleteUnspentI p ImportDB {importHashMap = hm} =+ withBlockMem hm $ deleteUnspent p -instance (MonadIO m) => StoreRead (ReaderT ImportDB m) where+instance MonadIO m => StoreRead (ReaderT ImportDB m) where isInitialized = R.ask >>= isInitializedI getBestBlock = R.ask >>= getBestBlockI getBlocksAtHeight h = R.ask >>= getBlocksAtHeightI h@@ -339,8 +324,10 @@ getSpender p = R.ask >>= getSpenderI p getSpenders t = R.ask >>= getSpendersI t getOrphanTx h = R.ask >>= getOrphanTxI h+ getUnspent a = R.ask >>= getUnspentI a+ getBalance a = R.ask >>= getBalanceI a -instance (MonadIO m) => StoreWrite (ReaderT ImportDB m) where+instance MonadIO m => StoreWrite (ReaderT ImportDB m) where setInit = R.ask >>= setInitI setBest h = R.ask >>= setBestI h insertBlock b = R.ask >>= insertBlockI b@@ -349,25 +336,13 @@ insertSpender p s = R.ask >>= insertSpenderI p s deleteSpender p = R.ask >>= deleteSpenderI p insertAddrTx a t = R.ask >>= insertAddrTxI a t- removeAddrTx a t = R.ask >>= removeAddrTxI a t+ deleteAddrTx a t = R.ask >>= deleteAddrTxI a t insertAddrUnspent a u = R.ask >>= insertAddrUnspentI a u- removeAddrUnspent a u = R.ask >>= removeAddrUnspentI a u+ deleteAddrUnspent a u = R.ask >>= deleteAddrUnspentI a u insertMempoolTx t p = R.ask >>= insertMempoolTxI t p deleteMempoolTx t p = R.ask >>= deleteMempoolTxI t p insertOrphanTx t p = R.ask >>= insertOrphanTxI t p deleteOrphanTx t = R.ask >>= deleteOrphanTxI t--instance (MonadIO m) => UnspentRead (ReaderT ImportDB m) where- getUnspent a = R.ask >>= getUnspentI a--instance (MonadIO m) => UnspentWrite (ReaderT ImportDB m) where- addUnspent u = R.ask >>= addUnspentI u- delUnspent p = R.ask >>= delUnspentI p- pruneUnspent = return ()--instance (MonadIO m) => BalanceRead (ReaderT ImportDB m) where- getBalance a = R.ask >>= getBalanceI a--instance (MonadIO m) => BalanceWrite (ReaderT ImportDB m) where+ insertUnspent u = R.ask >>= insertUnspentI u+ deleteUnspent p = R.ask >>= deleteUnspentI p setBalance b = R.ask >>= setBalanceI b- pruneBalance = return ()
src/Network/Haskoin/Store/Data/KeyValue.hs view
@@ -27,6 +27,7 @@ | AddrTxKeyB { addrTxKeyA :: !Address , addrTxKeyB :: !BlockRef }+ | AddrTxKeyS deriving (Show, Eq, Ord, Generic, Hashable) instance Serialize AddrTxKey@@ -50,6 +51,8 @@ putWord8 0x05 put a put b+ -- 0x05+ put AddrTxKeyS = putWord8 0x05 get = do guard . (== 0x05) =<< getWord8 a <- get@@ -80,6 +83,7 @@ | AddrOutKeyB { addrOutKeyA :: !Address , addrOutKeyB :: !BlockRef }+ | AddrOutKeyS deriving (Show, Read, Eq, Ord, Generic, Hashable) instance Serialize AddrOutKey@@ -90,15 +94,17 @@ put a put b put p- -- 0x06 · StoreAddr- put AddrOutKeyA {addrOutKeyA = a} = do- putWord8 0x06- put a -- 0x06 · StoreAddr · BlockRef put AddrOutKeyB {addrOutKeyA = a, addrOutKeyB = b} = do putWord8 0x06 put a put b+ -- 0x06 · StoreAddr+ put AddrOutKeyA {addrOutKeyA = a} = do+ putWord8 0x06+ put a+ -- 0x06+ put AddrOutKeyS = putWord8 0x06 get = do guard . (== 0x06) =<< getWord8 AddrOutKey <$> get <*> get <*> get@@ -157,29 +163,28 @@ data UnspentKey = UnspentKey { unspentKey :: !OutPoint } | UnspentKeyS { unspentKeyS :: !TxHash }+ | UnspentKeyB deriving (Show, Read, Eq, Ord, Generic, Hashable) -instance Serialize UnspentKey where+instance Serialize UnspentKey -- 0x09 · TxHash · Index+ where put UnspentKey {unspentKey = OutPoint {outPointHash = h, outPointIndex = i}} = do putWord8 0x09 put h put i+ -- 0x09 · TxHash put UnspentKeyS {unspentKeyS = t} = do putWord8 0x09 put t+ -- 0x09+ put UnspentKeyB = putWord8 0x09 get = do guard . (== 0x09) =<< getWord8 h <- get i <- get return $ UnspentKey OutPoint {outPointHash = h, outPointIndex = i} -data UnspentVal = UnspentVal- { unspentValBlock :: !BlockRef- , unspentValAmount :: !Word64- , unspentValScript :: !ByteString- } deriving (Show, Read, Eq, Ord, Generic, Hashable, Serialize)- instance R.Key UnspentKey instance R.KeyValue UnspentKey UnspentVal@@ -269,8 +274,11 @@ instance R.KeyValue HeightKey [BlockHash] -- | Address balance database key.-newtype BalKey- = BalKey { balanceKey :: Address }+data BalKey+ = BalKey+ { balanceKey :: !Address+ }+ | BalKeyS deriving (Show, Read, Eq, Ord, Generic, Hashable) instance Serialize BalKey where@@ -278,36 +286,13 @@ put BalKey {balanceKey = a} = do putWord8 0x04 put a+ -- 0x04+ put BalKeyS = putWord8 0x04 get = do guard . (== 0x04) =<< getWord8 BalKey <$> get instance R.Key BalKey---- | Address balance database value.-data BalVal = BalVal- { balValAmount :: !Word64- -- ^ balance in satoshi- , balValZero :: !Word64- -- ^ unconfirmed balance in satoshi- , balValUnspentCount :: !Word64- -- ^ number of unspent outputs- , balValTxCount :: !Word64- -- ^ number of transactions- , balValTotalReceived :: !Word64- -- ^ total amount received by this address- } deriving (Show, Read, Eq, Ord, Generic, Hashable, Serialize)---- | Default balance for an address.-instance Default BalVal where- def =- BalVal- { balValAmount = 0- , balValZero = 0- , balValUnspentCount = 0- , balValTxCount = 0- , balValTotalReceived = 0- } instance R.KeyValue BalKey BalVal
+ src/Network/Haskoin/Store/Data/Memory.hs view
@@ -0,0 +1,404 @@+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE FlexibleInstances #-}+{-# LANGUAGE TupleSections #-}+{-# OPTIONS_GHC -Wno-orphans #-}+module Network.Haskoin.Store.Data.Memory where++import Conduit+import Control.Monad+import Control.Monad.Reader (MonadReader, ReaderT)+import qualified Control.Monad.Reader as R+import qualified Data.ByteString.Short as B.Short+import Data.Function+import Data.HashMap.Strict (HashMap)+import qualified Data.HashMap.Strict as M+import Data.IntMap.Strict (IntMap)+import qualified Data.IntMap.Strict as I+import Data.List+import Data.Maybe+import Haskoin+import Network.Haskoin.Store.Data+import Network.Haskoin.Store.Data.KeyValue+import Network.Haskoin.Store.Messages+import UnliftIO++withBlockMem :: MonadIO m => TVar BlockMem -> ReaderT (TVar BlockMem) m a -> m a+withBlockMem = flip R.runReaderT++data BlockMem = BlockMem+ { hBest :: !(Maybe BlockHash)+ , hBlock :: !(HashMap BlockHash BlockData)+ , hHeight :: !(HashMap BlockHeight [BlockHash])+ , hTx :: !(HashMap TxHash TxData)+ , hSpender :: !(HashMap TxHash (IntMap (Maybe Spender)))+ , hUnspent :: !(HashMap TxHash (IntMap (Maybe UnspentVal)))+ , hBalance :: !(HashMap Address BalVal)+ , hAddrTx :: !(HashMap Address (HashMap BlockRef (HashMap TxHash Bool)))+ , hAddrOut :: !(HashMap Address (HashMap BlockRef (HashMap OutPoint (Maybe OutVal))))+ , hMempool :: !(HashMap UnixTime (HashMap TxHash Bool))+ , hOrphans :: !(HashMap TxHash (Maybe (UnixTime, Tx)))+ , hInit :: !Bool+ } deriving (Eq, Show)++emptyBlockMem :: BlockMem+emptyBlockMem =+ BlockMem+ { hBest = Nothing+ , hBlock = M.empty+ , hHeight = M.empty+ , hTx = M.empty+ , hSpender = M.empty+ , hUnspent = M.empty+ , hBalance = M.empty+ , hAddrTx = M.empty+ , hAddrOut = M.empty+ , hMempool = M.empty+ , hOrphans = M.empty+ , hInit = False+ }++isInitializedH :: BlockMem -> Either InitException Bool+isInitializedH = Right . hInit++getBestBlockH :: BlockMem -> Maybe BlockHash+getBestBlockH = hBest++getBlocksAtHeightH :: BlockHeight -> BlockMem -> [BlockHash]+getBlocksAtHeightH h = M.lookupDefault [] h . hHeight++getBlockH :: BlockHash -> BlockMem -> Maybe BlockData+getBlockH h = M.lookup h . hBlock++getTxDataH :: TxHash -> BlockMem -> Maybe TxData+getTxDataH t = M.lookup t . hTx++getSpenderH :: OutPoint -> BlockMem -> Maybe (Maybe Spender)+getSpenderH op db = do+ m <- M.lookup (outPointHash op) (hSpender db)+ I.lookup (fromIntegral (outPointIndex op)) m++getSpendersH :: TxHash -> BlockMem -> IntMap (Maybe Spender)+getSpendersH t = M.lookupDefault I.empty t . hSpender++getBalanceH :: Address -> BlockMem -> Maybe Balance+getBalanceH a = fmap (balValToBalance a) . M.lookup a . hBalance++getMempoolH ::+ Monad m+ => Maybe UnixTime+ -> BlockMem+ -> ConduitT () (UnixTime, TxHash) m ()+getMempoolH mpu db =+ let f ts =+ case mpu of+ Nothing -> False+ Just pu -> ts > pu+ ls =+ dropWhile (f . fst) .+ sortBy (flip compare) . M.toList . M.map (M.keys . M.filter id) $+ hMempool db+ in yieldMany [(u, h) | (u, hs) <- ls, h <- hs]++getOrphansH :: Monad m => BlockMem -> ConduitT () (UnixTime, Tx) m ()+getOrphansH = yieldMany . catMaybes . M.elems . hOrphans++getOrphanTxH :: TxHash -> BlockMem -> Maybe (Maybe (UnixTime, Tx))+getOrphanTxH h = M.lookup h . hOrphans++getUnspentsH :: Monad m => BlockMem -> ConduitT () Unspent m ()+getUnspentsH BlockMem {hUnspent = us} =+ yieldMany+ [ u+ | (h, m) <- M.toList us+ , (i, mv) <- I.toList m+ , v <- maybeToList mv+ , let p = OutPoint h (fromIntegral i)+ , let u = unspentValToUnspent p v+ ]++getAddressTxsH :: Address -> Maybe BlockRef -> BlockMem -> [BlockTx]+getAddressTxsH a mbr db =+ dropWhile h .+ sortBy (flip compare) . catMaybes . concatMap (uncurry f) . M.toList $+ M.lookupDefault M.empty a (hAddrTx db)+ where+ f b hm = map (uncurry (g b)) $ M.toList hm+ g b h' True =+ Just+ BlockTx+ {blockTxBlock = b, blockTxHash = h'}+ g _ _ False = Nothing+ h BlockTx {blockTxBlock = b} =+ case mbr of+ Nothing -> False+ Just br -> b > br++getAddressBalancesH :: Monad m => BlockMem -> ConduitT () Balance m ()+getAddressBalancesH BlockMem {hBalance = bm} =+ yieldMany (M.toList bm) .| mapC (uncurry balValToBalance)++getAddressUnspentsH ::+ Address -> Maybe BlockRef -> BlockMem -> [Unspent]+getAddressUnspentsH a mbr db =+ dropWhile h .+ sortBy (flip compare) . catMaybes . concatMap (uncurry f) . M.toList $+ M.lookupDefault M.empty a (hAddrOut db)+ where+ f b hm = map (uncurry (g b)) $ M.toList hm+ g b p (Just u) =+ Just+ Unspent+ { unspentBlock = b+ , unspentAmount = outValAmount u+ , unspentScript = B.Short.toShort (outValScript u)+ , unspentPoint = p+ }+ g _ _ Nothing = Nothing+ h Unspent {unspentBlock = b} =+ case mbr of+ Nothing -> False+ Just br -> b > br++setInitH :: BlockMem -> BlockMem+setInitH db = db {hInit = True}++setBestH :: BlockHash -> BlockMem -> BlockMem+setBestH h db = db {hBest = Just h}++insertBlockH :: BlockData -> BlockMem -> BlockMem+insertBlockH bd db =+ db {hBlock = M.insert (headerHash (blockDataHeader bd)) bd (hBlock db)}++insertAtHeightH :: BlockHash -> BlockHeight -> BlockMem -> BlockMem+insertAtHeightH h g db = db {hHeight = M.insertWith f g [h] (hHeight db)}+ where+ f xs ys = nub $ xs <> ys++insertTxH :: TxData -> BlockMem -> BlockMem+insertTxH tx db = db {hTx = M.insert (txHash (txData tx)) tx (hTx db)}++insertSpenderH :: OutPoint -> Spender -> BlockMem -> BlockMem+insertSpenderH op s db =+ db+ { hSpender =+ M.insertWith+ (<>)+ (outPointHash op)+ (I.singleton (fromIntegral (outPointIndex op)) (Just s))+ (hSpender db)+ }++deleteSpenderH :: OutPoint -> BlockMem -> BlockMem+deleteSpenderH op db =+ db+ { hSpender =+ M.insertWith+ (<>)+ (outPointHash op)+ (I.singleton (fromIntegral (outPointIndex op)) Nothing)+ (hSpender db)+ }++setBalanceH :: Balance -> BlockMem -> BlockMem+setBalanceH bal db = db {hBalance = M.insert a b (hBalance db)}+ where+ (a, b) = balanceToBalVal bal++insertAddrTxH :: Address -> BlockTx -> BlockMem -> BlockMem+insertAddrTxH a btx db =+ let s =+ M.singleton+ a+ (M.singleton+ (blockTxBlock btx)+ (M.singleton (blockTxHash btx) True))+ in db {hAddrTx = M.unionWith (M.unionWith M.union) s (hAddrTx db)}++deleteAddrTxH :: Address -> BlockTx -> BlockMem -> BlockMem+deleteAddrTxH a btx db =+ let s =+ M.singleton+ a+ (M.singleton+ (blockTxBlock btx)+ (M.singleton (blockTxHash btx) False))+ in db {hAddrTx = M.unionWith (M.unionWith M.union) s (hAddrTx db)}++insertAddrUnspentH :: Address -> Unspent -> BlockMem -> BlockMem+insertAddrUnspentH a u db =+ let uns =+ OutVal+ { outValAmount = unspentAmount u+ , outValScript = B.Short.fromShort (unspentScript u)+ }+ s =+ M.singleton+ a+ (M.singleton+ (unspentBlock u)+ (M.singleton (unspentPoint u) (Just uns)))+ in db {hAddrOut = M.unionWith (M.unionWith M.union) s (hAddrOut db)}++deleteAddrUnspentH :: Address -> Unspent -> BlockMem -> BlockMem+deleteAddrUnspentH a u db =+ let s =+ M.singleton+ a+ (M.singleton+ (unspentBlock u)+ (M.singleton (unspentPoint u) Nothing))+ in db {hAddrOut = M.unionWith (M.unionWith M.union) s (hAddrOut db)}++insertMempoolTxH :: TxHash -> UnixTime -> BlockMem -> BlockMem+insertMempoolTxH h u db =+ let s = M.singleton u (M.singleton h True)+ in db {hMempool = M.unionWith M.union s (hMempool db)}++deleteMempoolTxH :: TxHash -> UnixTime -> BlockMem -> BlockMem+deleteMempoolTxH h u db =+ let s = M.singleton u (M.singleton h False)+ in db {hMempool = M.unionWith M.union s (hMempool db)}++insertOrphanTxH :: Tx -> UnixTime -> BlockMem -> BlockMem+insertOrphanTxH tx u db =+ db {hOrphans = M.insert (txHash tx) (Just (u, tx)) (hOrphans db)}++deleteOrphanTxH :: TxHash -> BlockMem -> BlockMem+deleteOrphanTxH h db = db {hOrphans = M.insert h Nothing (hOrphans db)}++getUnspentH :: OutPoint -> BlockMem -> Maybe (Maybe Unspent)+getUnspentH op db = do+ m <- M.lookup (outPointHash op) (hUnspent db)+ fmap (unspentValToUnspent op) <$> I.lookup (fromIntegral (outPointIndex op)) m++insertUnspentH :: Unspent -> BlockMem -> BlockMem+insertUnspentH u db =+ db+ { hUnspent =+ M.insertWith+ (<>)+ (outPointHash (unspentPoint u))+ (I.singleton+ (fromIntegral (outPointIndex (unspentPoint u)))+ (Just (snd (unspentToUnspentVal u))))+ (hUnspent db)+ }++deleteUnspentH :: OutPoint -> BlockMem -> BlockMem+deleteUnspentH op db =+ db+ { hUnspent =+ M.insertWith+ (<>)+ (outPointHash op)+ (I.singleton (fromIntegral (outPointIndex op)) Nothing)+ (hUnspent db)+ }++instance MonadIO m => StoreRead (ReaderT (TVar BlockMem) m) where+ isInitialized = do+ v <- R.ask >>= readTVarIO+ return $ isInitializedH v+ getBestBlock = do+ v <- R.ask >>= readTVarIO+ return $ getBestBlockH v+ getBlocksAtHeight h = do+ v <- R.ask >>= readTVarIO+ return $ getBlocksAtHeightH h v+ getBlock b = do+ v <- R.ask >>= readTVarIO+ return $ getBlockH b v+ getTxData t = do+ v <- R.ask >>= readTVarIO+ return $ getTxDataH t v+ getSpender t = do+ v <- R.ask >>= readTVarIO+ return . join $ getSpenderH t v+ getSpenders t = do+ v <- R.ask >>= readTVarIO+ return . I.map fromJust . I.filter isJust $ getSpendersH t v+ getOrphanTx h = do+ v <- R.ask >>= readTVarIO+ return . join $ getOrphanTxH h v+ getUnspent p = do+ v <- R.ask >>= readTVarIO+ return . join $ getUnspentH p v+ getBalance a = do+ v <- R.ask >>= readTVarIO+ return $ getBalanceH a v++instance MonadIO m => StoreStream (ReaderT (TVar BlockMem) m) where+ getMempool m = do+ v <- R.ask >>= readTVarIO+ getMempoolH m v+ getOrphans = do+ v <- R.ask >>= readTVarIO+ getOrphansH v+ getAddressTxs a m = do+ v <- R.ask >>= readTVarIO+ yieldMany $ getAddressTxsH a m v+ getAddressUnspents a m = do+ v <- R.ask >>= readTVarIO+ yieldMany $ getAddressUnspentsH a m v+ getAddressBalances = do+ v <- R.ask >>= readTVarIO+ getAddressBalancesH v+ getUnspents = do+ v <- R.ask >>= readTVarIO+ getUnspentsH v++instance (MonadIO m) => StoreWrite (ReaderT (TVar BlockMem) m) where+ setInit = do+ v <- R.ask+ atomically $ modifyTVar v setInitH+ setBest h = do+ v <- R.ask+ atomically $ modifyTVar v (setBestH h)+ insertBlock b = do+ v <- R.ask+ atomically $ modifyTVar v (insertBlockH b)+ insertAtHeight h g = do+ v <- R.ask+ atomically $ modifyTVar v (insertAtHeightH h g)+ insertTx t = do+ v <- R.ask+ atomically $ modifyTVar v (insertTxH t)+ insertSpender p s = do+ v <- R.ask+ atomically $ modifyTVar v (insertSpenderH p s)+ deleteSpender p = do+ v <- R.ask+ atomically $ modifyTVar v (deleteSpenderH p)+ insertAddrTx a t = do+ v <- R.ask+ atomically $ modifyTVar v (insertAddrTxH a t)+ deleteAddrTx a t = do+ v <- R.ask+ atomically $ modifyTVar v (deleteAddrTxH a t)+ insertAddrUnspent a u = do+ v <- R.ask+ atomically $ modifyTVar v (insertAddrUnspentH a u)+ deleteAddrUnspent a u = do+ v <- R.ask+ atomically $ modifyTVar v (deleteAddrUnspentH a u)+ insertMempoolTx h t = do+ v <- R.ask+ atomically $ modifyTVar v (insertMempoolTxH h t)+ deleteMempoolTx h t = do+ v <- R.ask+ atomically $ modifyTVar v (deleteMempoolTxH h t)+ insertOrphanTx t u = do+ v <- R.ask+ atomically $ modifyTVar v (insertOrphanTxH t u)+ deleteOrphanTx h = do+ v <- R.ask+ atomically $ modifyTVar v (deleteOrphanTxH h)+ setBalance b = do+ v <- R.ask+ atomically $ modifyTVar v (setBalanceH b)+ insertUnspent h = do+ v <- R.ask+ atomically $ modifyTVar v (insertUnspentH h)+ deleteUnspent p = do+ v <- R.ask+ atomically $ modifyTVar v (deleteUnspentH p)
src/Network/Haskoin/Store/Data/RocksDB.hs view
@@ -1,11 +1,10 @@-{-# LANGUAGE DeriveAnyClass #-}-{-# LANGUAGE FlexibleInstances #-}-{-# LANGUAGE LambdaCase #-}+{-# LANGUAGE DeriveAnyClass #-}+{-# LANGUAGE LambdaCase #-} {-# OPTIONS_GHC -Wno-orphans #-} module Network.Haskoin.Store.Data.RocksDB where import Conduit-import Control.Monad.Reader (ReaderT)+import Control.Monad.Reader (MonadReader, ReaderT) import qualified Control.Monad.Reader as R import qualified Data.ByteString.Short as B.Short import Data.IntMap (IntMap)@@ -18,80 +17,62 @@ import Network.Haskoin.Store.Data.KeyValue import UnliftIO -type BlockDB = (ReadOptions, DB)- dataVersion :: Word32 dataVersion = 15 -withBlockDB :: ReadOptions -> DB -> ReaderT BlockDB m a -> m a-withBlockDB opts db f = R.runReaderT f (opts, db)--data ExceptRocksDB =- MempoolTxNotFound- deriving (Eq, Show, Read, Exception)--isInitializedDB ::- MonadIO m => ReadOptions -> DB -> m (Either InitException Bool)-isInitializedDB opts db =+isInitializedDB :: MonadIO m => BlockDB -> m (Either InitException Bool)+isInitializedDB BlockDB {blockDBopts = opts, blockDB = db} = retrieve db opts VersionKey >>= \case Just v | v == dataVersion -> return (Right True) | otherwise -> return (Left (IncorrectVersion v)) Nothing -> return (Right False) -getBestBlockDB :: MonadIO m => ReadOptions -> DB -> m (Maybe BlockHash)-getBestBlockDB opts db = retrieve db opts BestKey+setInitDB :: MonadIO m => DB -> m ()+setInitDB db = insert db VersionKey dataVersion -getBlocksAtHeightDB ::- MonadIO m => BlockHeight -> ReadOptions -> DB -> m [BlockHash]-getBlocksAtHeightDB h opts db =+getBestBlockDB :: MonadIO m => BlockDB -> m (Maybe BlockHash)+getBestBlockDB BlockDB {blockDBopts = opts, blockDB = db} =+ retrieve db opts BestKey++getBlocksAtHeightDB :: MonadIO m => BlockHeight -> BlockDB -> m [BlockHash]+getBlocksAtHeightDB h BlockDB {blockDBopts = opts, blockDB = db} = retrieve db opts (HeightKey h) >>= \case Nothing -> return [] Just ls -> return ls -getBlockDB :: MonadIO m => BlockHash -> ReadOptions -> DB -> m (Maybe BlockData)-getBlockDB h opts db = retrieve db opts (BlockKey h)+getBlockDB :: MonadIO m => BlockHash -> BlockDB -> m (Maybe BlockData)+getBlockDB h BlockDB {blockDBopts = opts, blockDB = db} =+ retrieve db opts (BlockKey h) getTxDataDB ::- MonadIO m => TxHash -> ReadOptions -> DB -> m (Maybe TxData)-getTxDataDB th opts db = retrieve db opts (TxKey th)+ MonadIO m => TxHash -> BlockDB -> m (Maybe TxData)+getTxDataDB th BlockDB {blockDBopts = opts, blockDB = db} =+ retrieve db opts (TxKey th) -getSpenderDB :: MonadIO m => OutPoint -> ReadOptions -> DB -> m (Maybe Spender)-getSpenderDB op opts db = retrieve db opts $ SpenderKey op+getSpenderDB :: MonadIO m => OutPoint -> BlockDB -> m (Maybe Spender)+getSpenderDB op BlockDB {blockDBopts = opts, blockDB = db} =+ retrieve db opts $ SpenderKey op -getSpendersDB :: MonadIO m => TxHash -> ReadOptions -> DB -> m (IntMap Spender)-getSpendersDB th opts db =+getSpendersDB :: MonadIO m => TxHash -> BlockDB -> m (IntMap Spender)+getSpendersDB th BlockDB {blockDBopts = opts, blockDB = db} = I.fromList . map (uncurry f) <$> liftIO (matchingAsList db opts (SpenderKeyS th)) where f (SpenderKey op) s = (fromIntegral (outPointIndex op), s) f _ _ = undefined -getBalanceDB :: MonadIO m => Address -> ReadOptions -> DB -> m (Maybe Balance)-getBalanceDB a opts db = fmap f <$> retrieve db opts (BalKey a)- where- f BalVal { balValAmount = v- , balValZero = z- , balValUnspentCount = c- , balValTxCount = t- , balValTotalReceived = r- } =- Balance- { balanceAddress = a- , balanceAmount = v- , balanceZero = z- , balanceUnspentCount = c- , balanceTxCount = t- , balanceTotalReceived = r- }+getBalanceDB :: MonadIO m => Address -> BlockDB -> m (Maybe Balance)+getBalanceDB a BlockDB {blockDBopts = opts, blockDB = db} =+ fmap (balValToBalance a) <$> retrieve db opts (BalKey a) getMempoolDB :: (MonadIO m, MonadResource m) => Maybe UnixTime- -> ReadOptions- -> DB+ -> BlockDB -> ConduitT () (UnixTime, TxHash) m ()-getMempoolDB mpu opts db = x .| mapC (uncurry f)+getMempoolDB mpu BlockDB {blockDBopts = opts, blockDB = db} =+ x .| mapC (uncurry f) where x = case mpu of@@ -102,23 +83,23 @@ getOrphansDB :: (MonadIO m, MonadResource m)- => ReadOptions- -> DB+ => BlockDB -> ConduitT () (UnixTime, Tx) m ()-getOrphansDB opts db = matching db opts OrphanKeyS .| mapC snd+getOrphansDB BlockDB {blockDBopts = opts, blockDB = db} =+ matching db opts OrphanKeyS .| mapC snd -getOrphanTxDB ::- MonadIO m => TxHash -> ReadOptions -> DB -> m (Maybe (UnixTime, Tx))-getOrphanTxDB h opts db = retrieve db opts (OrphanKey h)+getOrphanTxDB :: MonadIO m => TxHash -> BlockDB -> m (Maybe (UnixTime, Tx))+getOrphanTxDB h BlockDB {blockDBopts = opts, blockDB = db} =+ retrieve db opts (OrphanKey h) getAddressTxsDB :: (MonadIO m, MonadResource m) => Address -> Maybe BlockRef- -> ReadOptions- -> DB+ -> BlockDB -> ConduitT () BlockTx m ()-getAddressTxsDB a mbr opts db = x .| mapC (uncurry f)+getAddressTxsDB a mbr BlockDB {blockDBopts = opts, blockDB = db} =+ x .| mapC (uncurry f) where x = case mbr of@@ -127,14 +108,33 @@ f AddrTxKey {addrTxKeyT = t} () = t f _ _ = undefined +getAddressBalancesDB ::+ (MonadIO m, MonadResource m)+ => BlockDB+ -> ConduitT () Balance m ()+getAddressBalancesDB BlockDB {blockDBopts = opts, blockDB = db} =+ matching db opts BalKeyS .| mapC (\(BalKey a, b) -> balValToBalance a b)++getUnspentsDB ::+ (MonadIO m, MonadResource m)+ => BlockDB+ -> ConduitT () Unspent m ()+getUnspentsDB BlockDB {blockDBopts = opts, blockDB = db} =+ matching db opts UnspentKeyB .|+ mapC (\(UnspentKey k, v) -> unspentFromDB k v)++getUnspentDB :: MonadIO m => OutPoint -> BlockDB -> m (Maybe Unspent)+getUnspentDB p BlockDB {blockDBopts = opts, blockDB = db} =+ fmap (unspentValToUnspent p) <$> retrieve db opts (UnspentKey p)+ getAddressUnspentsDB :: (MonadIO m, MonadResource m) => Address -> Maybe BlockRef- -> ReadOptions- -> DB+ -> BlockDB -> ConduitT () Unspent m ()-getAddressUnspentsDB a mbr opts db = x .| mapC (uncurry f)+getAddressUnspentsDB a mbr BlockDB {blockDBopts = opts, blockDB = db} =+ x .| mapC (uncurry f) where x = case mbr of@@ -151,39 +151,14 @@ } f _ _ = undefined -getUnspentDB :: MonadIO m => OutPoint -> ReadOptions -> DB -> m (Maybe Unspent)-getUnspentDB op opts db = fmap f <$> retrieve db opts (UnspentKey op)- where- f u =- Unspent- { unspentBlock = unspentValBlock u- , unspentPoint = op- , unspentAmount = unspentValAmount u- , unspentScript = B.Short.toShort (unspentValScript u)- }--setInitDB :: MonadIO m => DB -> m ()-setInitDB db = insert db VersionKey dataVersion--instance MonadIO m => StoreRead (ReaderT BlockDB m) where- isInitialized = R.ask >>= uncurry isInitializedDB- getBestBlock = R.ask >>= uncurry getBestBlockDB- getBlocksAtHeight h = R.ask >>= uncurry (getBlocksAtHeightDB h)- getBlock h = R.ask >>= uncurry (getBlockDB h)- getTxData t = R.ask >>= uncurry (getTxDataDB t)- getSpenders p = R.ask >>= uncurry (getSpendersDB p)- getSpender p = R.ask >>= uncurry (getSpenderDB p)- getOrphanTx h = R.ask >>= uncurry (getOrphanTxDB h)--instance (MonadIO m, MonadResource m) =>- StoreStream (ReaderT BlockDB m) where- getMempool p = lift R.ask >>= uncurry (getMempoolDB p)- getOrphans = lift R.ask >>= uncurry getOrphansDB- getAddressTxs a b = R.ask >>= uncurry (getAddressTxsDB a b)- getAddressUnspents a b = R.ask >>= uncurry (getAddressUnspentsDB a b)--instance (MonadIO m) => BalanceRead (ReaderT BlockDB m) where- getBalance a = R.ask >>= uncurry (getBalanceDB a)--instance (MonadIO m) => UnspentRead (ReaderT BlockDB m) where- getUnspent p = R.ask >>= uncurry (getUnspentDB p)+unspentFromDB :: OutPoint -> UnspentVal -> Unspent+unspentFromDB p UnspentVal { unspentValBlock = b+ , unspentValAmount = v+ , unspentValScript = s+ } =+ Unspent+ { unspentBlock = b+ , unspentAmount = v+ , unspentPoint = p+ , unspentScript = s+ }
− src/Network/Haskoin/Store/Data/STM.hs
@@ -1,423 +0,0 @@-{-# LANGUAGE FlexibleInstances #-}-{-# LANGUAGE TupleSections #-}-{-# OPTIONS_GHC -Wno-orphans #-}-module Network.Haskoin.Store.Data.STM where--import Conduit-import Control.Monad-import Control.Monad.Reader (ReaderT)-import qualified Control.Monad.Reader as R-import qualified Data.ByteString.Short as B.Short-import Data.Function-import Data.HashMap.Strict (HashMap)-import qualified Data.HashMap.Strict as M-import Data.IntMap.Strict (IntMap)-import qualified Data.IntMap.Strict as I-import Data.List-import Data.Maybe-import Haskoin-import Network.Haskoin.Store.Data-import Network.Haskoin.Store.Data.KeyValue-import UnliftIO--type BlockSTM = ReaderT (TVar HashMapDB) STM-type UnspentSTM = ReaderT (TVar UnspentMap) STM-type BalanceSTM = ReaderT (TVar BalanceMap) STM-type UnspentMap = HashMap TxHash (IntMap Unspent)-type BalanceMap = (HashMap Address Balance, [Address])--withBlockSTM :: TVar HashMapDB -> ReaderT (TVar HashMapDB) STM a -> STM a-withBlockSTM = flip R.runReaderT--withUnspentSTM :: TVar UnspentMap -> ReaderT (TVar UnspentMap) STM a -> STM a-withUnspentSTM = flip R.runReaderT--withBalanceSTM :: TVar BalanceMap -> ReaderT (TVar BalanceMap) STM a -> STM a-withBalanceSTM = flip R.runReaderT--data HashMapDB = HashMapDB- { hBest :: !(Maybe BlockHash)- , hBlock :: !(HashMap BlockHash BlockData)- , hHeight :: !(HashMap BlockHeight [BlockHash])- , hTx :: !(HashMap TxHash TxData)- , hSpender :: !(HashMap TxHash (IntMap (Maybe Spender)))- , hUnspent :: !(HashMap TxHash (IntMap (Maybe Unspent)))- , hBalance :: !(HashMap Address BalVal)- , hAddrTx :: !(HashMap Address (HashMap BlockRef (HashMap TxHash Bool)))- , hAddrOut :: !(HashMap Address (HashMap BlockRef (HashMap OutPoint (Maybe OutVal))))- , hMempool :: !(HashMap UnixTime (HashMap TxHash Bool))- , hOrphans :: !(HashMap TxHash (Maybe (UnixTime, Tx)))- , hInit :: !Bool- } deriving (Eq, Show)--emptyHashMapDB :: HashMapDB-emptyHashMapDB =- HashMapDB- { hBest = Nothing- , hBlock = M.empty- , hHeight = M.empty- , hTx = M.empty- , hSpender = M.empty- , hUnspent = M.empty- , hBalance = M.empty- , hAddrTx = M.empty- , hAddrOut = M.empty- , hMempool = M.empty- , hOrphans = M.empty- , hInit = False- }--isInitializedH :: HashMapDB -> Either InitException Bool-isInitializedH = Right . hInit--getBestBlockH :: HashMapDB -> Maybe BlockHash-getBestBlockH = hBest--getBlocksAtHeightH ::- BlockHeight -> HashMapDB -> [BlockHash]-getBlocksAtHeightH h = M.lookupDefault [] h . hHeight--getBlockH :: BlockHash -> HashMapDB -> Maybe BlockData-getBlockH h = M.lookup h . hBlock--getTxDataH :: TxHash -> HashMapDB -> Maybe TxData-getTxDataH t = M.lookup t . hTx--getSpenderH :: OutPoint -> HashMapDB -> Maybe (Maybe Spender)-getSpenderH op db = do- m <- M.lookup (outPointHash op) (hSpender db)- I.lookup (fromIntegral (outPointIndex op)) m--getSpendersH :: TxHash -> HashMapDB -> IntMap (Maybe Spender)-getSpendersH t = M.lookupDefault I.empty t . hSpender--getBalanceH :: Address -> HashMapDB -> Maybe Balance-getBalanceH a = fmap f . M.lookup a . hBalance- where- f b =- Balance- { balanceAddress = a- , balanceAmount = balValAmount b- , balanceZero = balValZero b- , balanceUnspentCount = balValUnspentCount b- , balanceTxCount = balValTxCount b- , balanceTotalReceived = balValTotalReceived b- }--getMempoolH ::- Monad m- => Maybe UnixTime- -> HashMapDB- -> ConduitT () (UnixTime, TxHash) m ()-getMempoolH mpu db =- let f ts =- case mpu of- Nothing -> False- Just pu -> ts > pu- ls =- dropWhile (f . fst) .- sortBy (flip compare) . M.toList . M.map (M.keys . M.filter id) $- hMempool db- in yieldMany [(u, h) | (u, hs) <- ls, h <- hs]--getOrphansH :: Monad m => HashMapDB -> ConduitT () (UnixTime, Tx) m ()-getOrphansH = yieldMany . catMaybes . M.elems . hOrphans--getOrphanTxH :: TxHash -> HashMapDB -> Maybe (Maybe (UnixTime, Tx))-getOrphanTxH h = M.lookup h . hOrphans--getAddressTxsH :: Address -> Maybe BlockRef -> HashMapDB -> [BlockTx]-getAddressTxsH a mbr db =- dropWhile h .- sortBy (flip compare) . catMaybes . concatMap (uncurry f) . M.toList $- M.lookupDefault M.empty a (hAddrTx db)- where- f b hm = map (uncurry (g b)) $ M.toList hm- g b h' True =- Just- BlockTx- {blockTxBlock = b, blockTxHash = h'}- g _ _ False = Nothing- h BlockTx {blockTxBlock = b} =- case mbr of- Nothing -> False- Just br -> b > br--getAddressUnspentsH ::- Address -> Maybe BlockRef -> HashMapDB -> [Unspent]-getAddressUnspentsH a mbr db =- dropWhile h .- sortBy (flip compare) . catMaybes . concatMap (uncurry f) . M.toList $- M.lookupDefault M.empty a (hAddrOut db)- where- f b hm = map (uncurry (g b)) $ M.toList hm- g b p (Just u) =- Just- Unspent- { unspentBlock = b- , unspentAmount = outValAmount u- , unspentScript = B.Short.toShort (outValScript u)- , unspentPoint = p- }- g _ _ Nothing = Nothing- h Unspent {unspentBlock = b} =- case mbr of- Nothing -> False- Just br -> b > br--setInitH :: HashMapDB -> HashMapDB-setInitH db = db {hInit = True}--setBestH :: BlockHash -> HashMapDB -> HashMapDB-setBestH h db = db {hBest = Just h}--insertBlockH :: BlockData -> HashMapDB -> HashMapDB-insertBlockH bd db =- db {hBlock = M.insert (headerHash (blockDataHeader bd)) bd (hBlock db)}--insertAtHeightH :: BlockHash -> BlockHeight -> HashMapDB -> HashMapDB-insertAtHeightH h g db = db {hHeight = M.insertWith f g [h] (hHeight db)}- where- f xs ys = nub $ xs <> ys--insertTxH :: TxData -> HashMapDB -> HashMapDB-insertTxH tx db = db {hTx = M.insert (txHash (txData tx)) tx (hTx db)}--insertSpenderH :: OutPoint -> Spender -> HashMapDB -> HashMapDB-insertSpenderH op s db =- db- { hSpender =- M.insertWith- (<>)- (outPointHash op)- (I.singleton (fromIntegral (outPointIndex op)) (Just s))- (hSpender db)- }--deleteSpenderH :: OutPoint -> HashMapDB -> HashMapDB-deleteSpenderH op db =- db- { hSpender =- M.insertWith- (<>)- (outPointHash op)- (I.singleton (fromIntegral (outPointIndex op)) Nothing)- (hSpender db)- }--setBalanceH :: Balance -> HashMapDB -> HashMapDB-setBalanceH b db = db {hBalance = M.insert (balanceAddress b) x (hBalance db)}- where- x =- BalVal- { balValAmount = balanceAmount b- , balValZero = balanceZero b- , balValUnspentCount = balanceUnspentCount b- , balValTxCount = balanceTxCount b- , balValTotalReceived = balanceTotalReceived b- }--insertAddrTxH :: Address -> BlockTx -> HashMapDB -> HashMapDB-insertAddrTxH a btx db =- let s =- M.singleton- a- (M.singleton- (blockTxBlock btx)- (M.singleton (blockTxHash btx) True))- in db {hAddrTx = M.unionWith (M.unionWith M.union) s (hAddrTx db)}--removeAddrTxH :: Address -> BlockTx -> HashMapDB -> HashMapDB-removeAddrTxH a btx db =- let s =- M.singleton- a- (M.singleton- (blockTxBlock btx)- (M.singleton (blockTxHash btx) False))- in db {hAddrTx = M.unionWith (M.unionWith M.union) s (hAddrTx db)}--insertAddrUnspentH :: Address -> Unspent -> HashMapDB -> HashMapDB-insertAddrUnspentH a u db =- let uns =- OutVal- { outValAmount = unspentAmount u- , outValScript = B.Short.fromShort (unspentScript u)- }- s =- M.singleton- a- (M.singleton- (unspentBlock u)- (M.singleton (unspentPoint u) (Just uns)))- in db {hAddrOut = M.unionWith (M.unionWith M.union) s (hAddrOut db)}--removeAddrUnspentH :: Address -> Unspent -> HashMapDB -> HashMapDB-removeAddrUnspentH a u db =- let s =- M.singleton- a- (M.singleton- (unspentBlock u)- (M.singleton (unspentPoint u) Nothing))- in db {hAddrOut = M.unionWith (M.unionWith M.union) s (hAddrOut db)}--insertMempoolTxH :: TxHash -> UnixTime -> HashMapDB -> HashMapDB-insertMempoolTxH h u db =- let s = M.singleton u (M.singleton h True)- in db {hMempool = M.unionWith M.union s (hMempool db)}--deleteMempoolTxH :: TxHash -> UnixTime -> HashMapDB -> HashMapDB-deleteMempoolTxH h u db =- let s = M.singleton u (M.singleton h False)- in db {hMempool = M.unionWith M.union s (hMempool db)}--insertOrphanTxH :: Tx -> UnixTime -> HashMapDB -> HashMapDB-insertOrphanTxH tx u db =- db {hOrphans = M.insert (txHash tx) (Just (u, tx)) (hOrphans db)}--deleteOrphanTxH :: TxHash -> HashMapDB -> HashMapDB-deleteOrphanTxH h db = db {hOrphans = M.insert h Nothing (hOrphans db)}--getUnspentH :: OutPoint -> HashMapDB -> Maybe (Maybe Unspent)-getUnspentH op db = do- m <- M.lookup (outPointHash op) (hUnspent db)- I.lookup (fromIntegral (outPointIndex op)) m--addUnspentH :: Unspent -> HashMapDB -> HashMapDB-addUnspentH u db =- db- { hUnspent =- M.insertWith- (<>)- (outPointHash (unspentPoint u))- (I.singleton- (fromIntegral (outPointIndex (unspentPoint u)))- (Just u))- (hUnspent db)- }--delUnspentH :: OutPoint -> HashMapDB -> HashMapDB-delUnspentH op db =- db- { hUnspent =- M.insertWith- (<>)- (outPointHash op)- (I.singleton (fromIntegral (outPointIndex op)) Nothing)- (hUnspent db)- }--instance StoreRead BlockSTM where- isInitialized = fmap isInitializedH . lift . readTVar =<< R.ask- getBestBlock = fmap getBestBlockH . lift . readTVar =<< R.ask- getBlocksAtHeight h =- fmap (getBlocksAtHeightH h) . lift . readTVar =<< R.ask- getBlock b = fmap (getBlockH b) . lift . readTVar =<< R.ask- getTxData t = fmap (getTxDataH t) . lift . readTVar =<< R.ask- getSpender t = fmap (join . getSpenderH t) . lift . readTVar =<< R.ask- getSpenders t =- fmap (I.map fromJust . I.filter isJust . getSpendersH t) .- lift . readTVar =<<- R.ask- getOrphanTx h = fmap (join . getOrphanTxH h) . lift . readTVar =<< R.ask--instance BalanceRead BlockSTM where- getBalance a = fmap (getBalanceH a) . lift . readTVar =<< R.ask--instance UnspentRead BlockSTM where- getUnspent op = fmap (join . getUnspentH op) . lift . readTVar =<< R.ask--instance BalanceWrite BlockSTM where- setBalance b = lift . (`modifyTVar` setBalanceH b) =<< R.ask- pruneBalance = return ()--instance StoreStream BlockSTM where- getMempool m = getMempoolH m =<< lift . lift . readTVar =<< lift R.ask- getOrphans = getOrphansH =<< lift . lift . readTVar =<< lift R.ask- getAddressTxs a m =- yieldMany . getAddressTxsH a m =<< lift . lift . readTVar =<< lift R.ask- getAddressUnspents a m =- yieldMany . getAddressUnspentsH a m =<<- lift . lift . readTVar =<< lift R.ask--instance StoreWrite BlockSTM where- setInit = lift . (`modifyTVar` setInitH) =<< R.ask- setBest h = lift . (`modifyTVar` setBestH h) =<< R.ask- insertBlock b = lift . (`modifyTVar` insertBlockH b) =<< R.ask- insertAtHeight h g = lift . (`modifyTVar` insertAtHeightH h g) =<< R.ask- insertTx t = lift . (`modifyTVar` insertTxH t) =<< R.ask- insertSpender p s = lift . (`modifyTVar` insertSpenderH p s) =<< R.ask- deleteSpender p = lift . (`modifyTVar` deleteSpenderH p) =<< R.ask- insertAddrTx a t = lift . (`modifyTVar` insertAddrTxH a t) =<< R.ask- removeAddrTx a t = lift . (`modifyTVar` removeAddrTxH a t) =<< R.ask- insertAddrUnspent a u =- lift . (`modifyTVar` insertAddrUnspentH a u) =<< R.ask- removeAddrUnspent a u =- lift . (`modifyTVar` removeAddrUnspentH a u) =<< R.ask- insertMempoolTx h t = lift . (`modifyTVar` insertMempoolTxH h t) =<< R.ask- deleteMempoolTx h t = lift . (`modifyTVar` deleteMempoolTxH h t) =<< R.ask- deleteOrphanTx h = lift . (`modifyTVar` deleteOrphanTxH h) =<< R.ask- insertOrphanTx t u = lift . (`modifyTVar` insertOrphanTxH t u) =<< R.ask--instance UnspentWrite BlockSTM where- addUnspent h = lift . (`modifyTVar` addUnspentH h) =<< R.ask- delUnspent p = lift . (`modifyTVar` delUnspentH p) =<< R.ask- pruneUnspent = return ()--instance UnspentRead UnspentSTM where- getUnspent op = do- um <- lift . readTVar =<< R.ask- return $ do- m <- M.lookup (outPointHash op) um- I.lookup (fromIntegral (outPointIndex op)) m--instance UnspentWrite UnspentSTM where- addUnspent u = do- v <- R.ask- lift . modifyTVar v $- M.insertWith- (<>)- (outPointHash (unspentPoint u))- (I.singleton (fromIntegral (outPointIndex (unspentPoint u))) u)- delUnspent op = lift . (`modifyTVar` M.update g (outPointHash op)) =<< R.ask- where- g m =- let n = I.delete (fromIntegral (outPointIndex op)) m- in if I.null n- then Nothing- else Just n- pruneUnspent = do- v <- R.ask- lift . modifyTVar v $ \um ->- if M.size um > 2 ^ (21 :: Int)- then let g is = unspentBlock (head (I.elems is))- ls =- sortBy- (compare `on` (g . snd))- (filter (not . I.null . snd) (M.toList um))- in M.fromList (drop (2 ^ (20 :: Int)) ls)- else um--instance BalanceRead BalanceSTM where- getBalance a = do- b <- fmap fst $ lift . readTVar =<< R.ask- return $ M.lookup a b--instance BalanceWrite BalanceSTM where- setBalance b = do- v <- R.ask- lift . modifyTVar v $ \(m, s) ->- let m' = M.insert (balanceAddress b) b m- s' = balanceAddress b : s- in (m', s')- pruneBalance = do- v <- R.ask- lift . modifyTVar v $ \(m, s) ->- if length s > 2 ^ (21 :: Int)- then let s' = take (2 ^ (20 :: Int)) s- m' = M.fromList (mapMaybe (g m) s')- in (m', s')- else (m, s)- where- g m a = (a, ) <$> M.lookup a m
src/Network/Haskoin/Store/Logic.hs view
@@ -12,7 +12,7 @@ import Control.Monad.Logger import qualified Data.ByteString as B import qualified Data.ByteString.Short as B.Short-import Data.Either+import Data.Either (rights) import qualified Data.IntMap.Strict as I import Data.List import Data.Maybe@@ -25,82 +25,119 @@ import Haskoin import Network.Haskoin.Block.Headers (computeSubsidy) import Network.Haskoin.Store.Data-import Network.Haskoin.Store.Data.ImportDB-import Network.Haskoin.Store.Data.STM import UnliftIO data ImportException- = PrevBlockNotBest !BlockHash- | UnconfirmedCoinbase !TxHash+ = PrevBlockNotBest !Text+ | UnconfirmedCoinbase !Text | BestBlockUnknown- | BestBlockNotFound !BlockHash- | BlockNotBest !BlockHash- | OrphanTx !TxHash- | TxNotFound !TxHash- | NoUnspent !OutPoint- | TxInvalidOp !TxHash- | TxDeleted !TxHash- | TxDoubleSpend !TxHash- | AlreadyUnspent !OutPoint- | TxConfirmed !TxHash- | OutputOutOfRange !OutPoint- | BalanceNotFound !Address- | InsufficientBalance !Address- | InsufficientZeroBalance !Address- | InsufficientOutputs !Address- | InsufficientFunds !TxHash+ | BestBlockNotFound !Text+ | BlockNotBest !Text+ | OrphanTx !Text+ | TxNotFound !Text+ | NoUnspent !Text+ | TxInvalidOp !Text+ | TxDeleted !Text+ | TxDoubleSpend !Text+ | AlreadyUnspent !Text+ | TxConfirmed !Text+ | OutputOutOfRange !Text+ | BalanceNotFound !Text+ | InsufficientBalance !Text+ | InsufficientZeroBalance !Text+ | InsufficientOutputs !Text+ | InsufficientFunds !Text | InitException !InitException- | DuplicatePrevOutput !TxHash+ | DuplicatePrevOutput !Text deriving (Show, Read, Eq, Ord, Exception) initDB ::- (MonadIO m, MonadError ImportException m, MonadLoggerIO m)+ ( StoreRead m+ , StoreWrite m+ , MonadLogger m+ , MonadError ImportException m+ ) => Network- -> DB- -> TVar UnspentMap- -> TVar BalanceMap -> m ()-initDB net db um bm =- runImportDB db um bm $- isInitialized >>= \case- Left e -> do- $(logErrorS) "BlockLogic" $- "Initialization exception: " <> fromString (show e)- throwError (InitException e)- Right True -> do- $(logDebugS) "BlockLogic" "Database is already initialized"- return ()- Right False -> do- $(logDebugS)- "BlockLogic"- "Initializing database by importing genesis block"- importBlock net (genesisBlock net) (genesisNode net)- setInit+initDB net =+ isInitialized >>= \case+ Left e -> do+ $(logErrorS) "BlockLogic" $+ "Initialization exception: " <> fromString (show e)+ throwError (InitException e)+ Right True -> do+ $(logDebugS) "BlockLogic" "Database is already initialized"+ return ()+ Right False -> do+ $(logDebugS)+ "BlockLogic"+ "Initializing database by importing genesis block"+ importBlock net (genesisBlock net) (genesisNode net)+ setInit +getOldOrphans ::+ (StoreStream m, MonadResource m)+ => UnixTime+ -> ConduitT () TxHash m ()+getOldOrphans now =+ getOrphans .| filterC ((< now - 600) . fst) .| mapC (txHash . snd)++importOrphan ::+ ( StoreRead m+ , StoreWrite m+ , MonadLogger m+ , MonadError ImportException m+ )+ => Network+ -> UnixTime+ -> Tx+ -> m ()+importOrphan net t tx = do+ $(logDebugS) "Block" $+ "Attempting to import orphan tx " <> txHashToHex (txHash tx)+ go `catchError` ex+ where+ go = do+ newMempoolTx net tx t >>= \case+ True ->+ $(logDebugS) "BlockLogic" $+ "Succesfully imported orphan transaction: " <>+ txHashToHex (txHash tx)+ False ->+ $(logDebugS) "BlockLogic" $+ "Orphan transaction already imported: " <>+ txHashToHex (txHash tx)+ deleteOrphanTx (txHash tx)+ ex (OrphanTx _) = do+ $(logDebugS) "BlockLogic" $+ "Transaction still orphan: " <> txHashToHex (txHash tx)+ ex e = do+ $(logErrorS) "BlockLogic" $+ "Error importing orphan tx: " <> txHashToHex (txHash tx) <> ": " <>+ cs (show e)+ deleteOrphanTx (txHash tx)+ newMempoolTx ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , UnspentRead m- , UnspentWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m+ , MonadError ImportException m ) => Network -> Tx -> UnixTime -> m Bool newMempoolTx net tx w = do- $(logInfoS) "BlockLogic" $- "Adding transaction to mempool: " <> txHashToHex (txHash tx)- getTxData (txHash tx) >>= \case- Just x- | not (txDataDeleted x) -> do- $(logWarnS) "BlockLogic" $- "Transaction already exists: " <> txHashToHex (txHash tx)- return False- _ -> go+ $(logInfoS) "BlockLogic" $+ "Adding transaction to mempool: " <> txHashToHex (txHash tx)+ getTxData (txHash tx) >>= \case+ Just x+ | not (txDataDeleted x) -> do+ $(logWarnS) "BlockLogic" $+ "Transaction already exists: " <>+ txHashToHex (txHash tx)+ return False+ _ -> go where go = do orp <-@@ -108,9 +145,10 @@ mapM (getTxData . outPointHash . prevOutput) (txIn tx) if orp then do- $(logErrorS) "BlockLogic" $+ $(logWarnS) "BlockLogic" $ "Transaction is orphan: " <> txHashToHex (txHash tx)- throwError $ OrphanTx (txHash tx)+ insertOrphanTx tx w+ throwError $ OrphanTx (txHashToHex (txHash tx)) else f f = do us <-@@ -148,96 +186,93 @@ isrbf th = transactionRBF <$> getImportTx th revertBlock ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , UnspentRead m- , UnspentWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m+ , MonadError ImportException m ) => Network -> BlockHash -> m () revertBlock net bh = do- bd <-- getBestBlock >>= \case- Nothing -> do- $(logErrorS) "BlockLogic" "Best block unknown"- throwError BestBlockUnknown- Just h ->- getBlock h >>= \case- Nothing -> do- $(logErrorS) "BlockLogic" "Best block not found"- throwError (BestBlockNotFound h)- Just b- | h == bh -> return b- | otherwise -> do- $(logErrorS) "BlockLogic" $- "Attempted to delete block that isn't best: " <>- blockHashToHex h- throwError (BlockNotBest bh)- txs <- mapM (fmap transactionData . getImportTx) (blockDataTxs bd)- mapM_ (deleteTx net False . txHash . snd) (reverse (sortTxs txs))- setBest (prevBlock (blockDataHeader bd))- insertBlock bd {blockDataMainChain = False}+ bd <-+ getBestBlock >>= \case+ Nothing -> do+ $(logErrorS) "BlockLogic" "Best block unknown"+ throwError BestBlockUnknown+ Just h ->+ getBlock h >>= \case+ Nothing -> do+ $(logErrorS) "BlockLogic" "Best block not found"+ throwError (BestBlockNotFound (blockHashToHex h))+ Just b+ | h == bh -> return b+ | otherwise -> do+ $(logErrorS) "BlockLogic" $+ "Attempted to delete block that isn't best: " <>+ blockHashToHex h+ throwError (BlockNotBest (blockHashToHex bh))+ txs <-+ mapM (fmap transactionData . getImportTx) (blockDataTxs bd)+ mapM_+ (deleteTx net False . txHash . snd)+ (reverse (sortTxs txs))+ setBest (prevBlock (blockDataHeader bd))+ insertBlock bd {blockDataMainChain = False} importBlock ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , UnspentRead m- , UnspentWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m+ , MonadError ImportException m ) => Network -> Block -> BlockNode -> m () importBlock net b n = do- getBestBlock >>= \case- Nothing- | isGenesis n -> do- $(logInfoS) "BlockLogic" $- "Importing genesis block: " <>- blockHashToHex (headerHash (nodeHeader n))- return ()- | otherwise -> do- $(logErrorS) "BlockLogic" $- "Importing non-genesis block when best block unknown: " <>- blockHashToHex (headerHash (blockHeader b))- throwError BestBlockUnknown- Just h- | prevBlock (blockHeader b) == h -> return ()- | otherwise -> do- $(logErrorS) "BlockLogic" $- "Block " <> blockHashToHex (headerHash (blockHeader b)) <>- " does not build on current best " <>- blockHashToHex h- throwError (PrevBlockNotBest (prevBlock (nodeHeader n)))- insertBlock- BlockData- { blockDataHeight = nodeHeight n- , blockDataMainChain = True- , blockDataWork = nodeWork n- , blockDataHeader = nodeHeader n- , blockDataSize = fromIntegral (B.length (encode b))- , blockDataTxs = map txHash (blockTxns b)- , blockDataWeight = fromIntegral w- , blockDataSubsidy = subsidy (nodeHeight n)- , blockDataFees = cb_out_val - subsidy (nodeHeight n)- , blockDataOutputs = ts_out_val- }- insertAtHeight (headerHash (nodeHeader n)) (nodeHeight n)- setBest (headerHash (nodeHeader n))- $(logDebugS) "Block" $ "Importing or confirming block transactions..."- mapM_ (uncurry import_or_confirm) (sortTxs (blockTxns b))- $(logDebugS) "Block" $- "Done importing transactions for block " <>- blockHashToHex (headerHash (nodeHeader n))+ getBestBlock >>= \case+ Nothing+ | isGenesis n -> do+ $(logInfoS) "BlockLogic" $+ "Importing genesis block: " <>+ blockHashToHex (headerHash (nodeHeader n))+ return ()+ | otherwise -> do+ $(logErrorS) "BlockLogic" $+ "Importing non-genesis block when best block unknown: " <>+ blockHashToHex (headerHash (blockHeader b))+ throwError BestBlockUnknown+ Just h+ | prevBlock (blockHeader b) == h -> return ()+ | otherwise -> do+ $(logErrorS) "BlockLogic" $+ "Block " <> blockHashToHex (headerHash (blockHeader b)) <>+ " does not build on current best " <>+ blockHashToHex h+ throwError+ (PrevBlockNotBest+ (blockHashToHex (prevBlock (nodeHeader n))))+ insertBlock+ BlockData+ { blockDataHeight = nodeHeight n+ , blockDataMainChain = True+ , blockDataWork = nodeWork n+ , blockDataHeader = nodeHeader n+ , blockDataSize = fromIntegral (B.length (encode b))+ , blockDataTxs = map txHash (blockTxns b)+ , blockDataWeight = fromIntegral w+ , blockDataSubsidy = subsidy (nodeHeight n)+ , blockDataFees = cb_out_val - subsidy (nodeHeight n)+ , blockDataOutputs = ts_out_val+ }+ insertAtHeight (headerHash (nodeHeader n)) (nodeHeight n)+ setBest (headerHash (nodeHeader n))+ $(logDebugS) "Block" $ "Importing or confirming block transactions..."+ mapM_ (uncurry import_or_confirm) (sortTxs (blockTxns b))+ $(logDebugS) "Block" $+ "Done importing transactions for block " <>+ blockHashToHex (headerHash (nodeHeader n)) where import_or_confirm x tx = getTxData (txHash tx) >>= \case@@ -278,14 +313,10 @@ in is <> go ds importTx ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , UnspentRead m- , UnspentWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m+ , MonadError ImportException m ) => Network -> BlockRef@@ -293,46 +324,54 @@ -> Tx -> m () importTx net br tt tx = do- when (length (nub (map prevOutput (txIn tx))) < length (txIn tx)) $ do- $(logErrorS) "BlockLogic" $- "Transaction spends same output twice: " <> txHashToHex (txHash tx)- throwError (DuplicatePrevOutput (txHash tx))- when (iscb && not (confirmed br)) $ do- $(logErrorS) "BlockLogic" $- "Attempting to import coinbase to the mempool: " <>- txHashToHex (txHash tx)- throwError (UnconfirmedCoinbase (txHash tx))- us <-- if iscb- then return []- else forM (txIn tx) $ \TxIn {prevOutput = op} -> uns op- when- (not (confirmed br) &&- sum (map unspentAmount us) < sum (map outValue (txOut tx))) $ do- $(logErrorS) "BlockLogic" $- "Insufficient funds: " <> txHashToHex (txHash tx)- throwError (InsufficientFunds th)- zipWithM_ (spendOutput net br (txHash tx)) [0 ..] us- zipWithM_ (newOutput net br . OutPoint (txHash tx)) [0 ..] (txOut tx)- rbf <- getrbf- let t =- Transaction- { transactionBlock = br- , transactionVersion = txVersion tx- , transactionLockTime = txLockTime tx- , transactionInputs =- if iscb- then zipWith mkcb (txIn tx) ws- else zipWith3 mkin us (txIn tx) ws- , transactionOutputs = map mkout (txOut tx)- , transactionDeleted = False- , transactionRBF = rbf- , transactionTime = tt- }- let (d, _) = fromTransaction t- insertTx d- updateAddressCounts (txAddresses t) (+ 1)- unless (confirmed br) $ insertMempoolTx (txHash tx) (memRefTime br)+ when (length (nub (map prevOutput (txIn tx))) < length (txIn tx)) $ do+ $(logErrorS) "BlockLogic" $+ "Transaction spends same output twice: " <>+ txHashToHex (txHash tx)+ throwError (DuplicatePrevOutput (txHashToHex (txHash tx)))+ when (iscb && not (confirmed br)) $ do+ $(logErrorS) "BlockLogic" $+ "Attempting to import coinbase to the mempool: " <>+ txHashToHex (txHash tx)+ throwError (UnconfirmedCoinbase (txHashToHex (txHash tx)))+ us <-+ if iscb+ then return []+ else forM (txIn tx) $ \TxIn {prevOutput = op} -> uns op+ when+ (not (confirmed br) &&+ sum (map unspentAmount us) < sum (map outValue (txOut tx))) $ do+ $(logErrorS) "BlockLogic" $+ "Insufficient funds: " <> txHashToHex (txHash tx)+ throwError (InsufficientFunds (txHashToHex th))+ zipWithM_+ (\i u -> spendOutput net br (txHash tx) i u)+ [0 ..]+ us+ zipWithM_+ (\i o -> newOutput net br (OutPoint (txHash tx) i) o)+ [0 ..]+ (txOut tx)+ rbf <- getrbf+ let t =+ Transaction+ { transactionBlock = br+ , transactionVersion = txVersion tx+ , transactionLockTime = txLockTime tx+ , transactionInputs =+ if iscb+ then zipWith mkcb (txIn tx) ws+ else zipWith3 mkin us (txIn tx) ws+ , transactionOutputs = map mkout (txOut tx)+ , transactionDeleted = False+ , transactionRBF = rbf+ , transactionTime = tt+ }+ let (d, _) = fromTransaction t+ insertTx d+ updateAddressCounts net (txAddresses t) (+ 1)+ unless (confirmed br) $+ insertMempoolTx (txHash tx) (memRefTime br) where uns op = getUnspent op >>= \case@@ -349,7 +388,7 @@ txHashToHex (outPointHash op) <> " " <> fromString (show (outPointIndex op))- throwError (NoUnspent op)+ throwError (NoUnspent (cs (show op))) Just Spender {spenderHash = s} -> do deleteTx net True s getUnspent op >>= \case@@ -359,7 +398,7 @@ txHashToHex (outPointHash op) <> " " <> fromString (show (outPointIndex op))- throwError (NoUnspent op)+ throwError (NoUnspent (cs (show op))) Just u -> return u th = txHash tx iscb = all (== nullOutPoint) (map prevOutput (txIn tx))@@ -372,7 +411,7 @@ let hs = nub $ map (outPointHash . prevOutput) (txIn tx) in fmap or . forM hs $ \h -> getTxData h >>= \case- Nothing -> throwError (TxNotFound h)+ Nothing -> throwError (TxNotFound (txHashToHex h)) Just t | confirmed (txDataBlock t) -> return False | txDataRBF t -> return True@@ -401,14 +440,10 @@ } confirmTx ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , BalanceRead m- , BalanceWrite m- , UnspentRead m- , UnspentWrite m , MonadLogger m+ , MonadError ImportException m ) => Network -> TxData@@ -420,7 +455,7 @@ case scriptToAddressBS (prevScript p) of Left _ -> return () Right a -> do- removeAddrTx+ deleteAddrTx a BlockTx {blockTxBlock = txDataBlock t, blockTxHash = txHash tx}@@ -431,8 +466,8 @@ let op = OutPoint (txHash tx) n s <- getSpender (OutPoint (txHash tx) n) when (isNothing s) $ do- delUnspent op- addUnspent+ deleteUnspent op+ insertUnspent Unspent { unspentBlock = br , unspentPoint = op@@ -442,7 +477,7 @@ case scriptToAddressBS (scriptOutput o) of Left _ -> return () Right a -> do- removeAddrTx+ deleteAddrTx a BlockTx {blockTxBlock = txDataBlock t, blockTxHash = txHash tx}@@ -450,7 +485,7 @@ a BlockTx {blockTxBlock = br, blockTxHash = txHash tx} when (isNothing s) $ do- removeAddrUnspent+ deleteAddrUnspent a Unspent { unspentBlock = txDataBlock t@@ -467,7 +502,7 @@ , unspentScript = B.Short.toShort (scriptOutput o) } reduceBalance net False False a (outValue o)- increaseBalance net True False a (outValue o)+ increaseBalance True False a (outValue o) insertTx t {txDataBlock = br} deleteMempoolTx (txHash tx) (memRefTime (txDataBlock t)) @@ -484,14 +519,10 @@ concat <$> mapM getRecursiveTx ss deleteTx ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , UnspentRead m- , UnspentWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m+ , MonadError ImportException m ) => Network -> Bool -- ^ only delete transaction if unconfirmed@@ -503,7 +534,7 @@ Nothing -> do $(logErrorS) "BlockLogic" $ "Transaciton not found: " <> txHashToHex h- throwError (TxNotFound h)+ throwError (TxNotFound (txHashToHex h)) Just t | txDataDeleted t -> do $(logWarnS) "BlockLogic" $@@ -512,7 +543,7 @@ | mo && confirmed (txDataBlock t) -> do $(logErrorS) "BlockLogic" $ "Will not delete confirmed transaction: " <> txHashToHex h- throwError (TxConfirmed h)+ throwError (TxConfirmed (txHashToHex h)) | otherwise -> go t where go t = do@@ -525,37 +556,38 @@ unless (confirmed (txDataBlock t)) $ deleteMempoolTx h (memRefTime (txDataBlock t)) insertTx t {txDataDeleted = True}- updateAddressCounts (txDataAddresses t) (subtract 1)+ updateAddressCounts net (txDataAddresses t) (subtract 1) insertDeletedMempoolTx ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m , MonadLogger m+ , MonadError ImportException m ) => Tx -> UnixTime -> m () insertDeletedMempoolTx tx w = do- us <-- forM (txIn tx) $ \TxIn {prevOutput = op} ->- getImportTx (outPointHash op) >>= getTxOutput (outPointIndex op)- rbf <- getrbf- let (d, _) =- fromTransaction- Transaction- { transactionBlock = MemRef w- , transactionVersion = txVersion tx- , transactionLockTime = txLockTime tx- , transactionInputs = zipWith3 mkin us (txIn tx) ws- , transactionOutputs = map mkout (txOut tx)- , transactionDeleted = True- , transactionRBF = rbf- , transactionTime = w- }- $(logWarnS) "BlockLogic" $- "Inserting deleted mempool transaction: " <> txHashToHex (txHash tx)- insertTx d+ us <-+ forM (txIn tx) $ \TxIn {prevOutput = op} ->+ getImportTx (outPointHash op) >>=+ getTxOutput (outPointIndex op)+ rbf <- getrbf+ let (d, _) =+ fromTransaction+ Transaction+ { transactionBlock = MemRef w+ , transactionVersion = txVersion tx+ , transactionLockTime = txLockTime tx+ , transactionInputs = zipWith3 mkin us (txIn tx) ws+ , transactionOutputs = map mkout (txOut tx)+ , transactionDeleted = True+ , transactionRBF = rbf+ , transactionTime = w+ }+ $(logWarnS) "BlockLogic" $+ "Inserting deleted mempool transaction: " <> txHashToHex (txHash tx)+ insertTx d where ws = map Just (txWitness tx) <> repeat Nothing getrbf@@ -567,7 +599,7 @@ Nothing -> do $(logErrorS) "BlockLogic" $ "Transaction not found: " <> txHashToHex h- throwError (TxNotFound h)+ throwError (TxNotFound (txHashToHex h)) Just t | confirmed (txDataBlock t) -> return False | txDataRBF t -> return True@@ -589,13 +621,8 @@ } newOutput ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , UnspentRead m- , UnspentWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m ) => Network@@ -604,18 +631,15 @@ -> TxOut -> m () newOutput net br op to = do- addUnspent u+ insertUnspent u case scriptToAddressBS (scriptOutput to) of Left _ -> return () Right a -> do insertAddrUnspent a u insertAddrTx a- BlockTx- { blockTxHash = outPointHash op- , blockTxBlock = br- }- increaseBalance net (confirmed br) True a (outValue to)+ BlockTx {blockTxHash = outPointHash op, blockTxBlock = br}+ increaseBalance (confirmed br) True a (outValue to) where u = Unspent@@ -626,14 +650,10 @@ } delOutput ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , UnspentRead m- , UnspentWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m+ , MonadError ImportException m ) => Network -> OutPoint@@ -641,11 +661,11 @@ delOutput net op = do t <- getImportTx (outPointHash op) u <- getTxOutput (outPointIndex op) t- delUnspent op+ deleteUnspent op case scriptToAddressBS (outputScript u) of Left _ -> return () Right a -> do- removeAddrUnspent+ deleteAddrUnspent a Unspent { unspentScript = B.Short.toShort (outputScript u)@@ -653,7 +673,7 @@ , unspentPoint = op , unspentAmount = outputAmount u }- removeAddrTx+ deleteAddrTx a BlockTx { blockTxHash = outPointHash op@@ -667,7 +687,7 @@ (outputAmount u) getImportTx ::- (MonadError ImportException m, StoreRead m, MonadLogger m)+ (StoreRead m, MonadLogger m, MonadError ImportException m) => TxHash -> m Transaction getImportTx th =@@ -675,18 +695,18 @@ Nothing -> do $(logErrorS) "BlockLogic" $ "Tranasction not found: " <> txHashToHex th- throwError $ TxNotFound th+ throwError $ TxNotFound (txHashToHex th) Just d | txDataDeleted d -> do $(logErrorS) "BlockLogic" $ "Transaction deleted: " <> txHashToHex th- throwError $ TxDeleted th+ throwError $ TxDeleted (txHashToHex th) | otherwise -> do sm <- getSpenders th return $ toTransaction d sm getTxOutput ::- (MonadError ImportException m, MonadLogger m)+ (MonadLogger m, MonadError ImportException m) => Word32 -> Transaction -> m StoreOutput@@ -696,8 +716,8 @@ "Output out of range " <> txHashToHex (txHash (transactionData tx)) <> " " <> fromString (show i)- throwError $- OutputOutOfRange+ throwError . OutputOutOfRange . cs $+ show OutPoint { outPointHash = txHash (transactionData tx) , outPointIndex = i@@ -705,14 +725,10 @@ return $ transactionOutputs tx !! fromIntegral i spendOutput ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , UnspentRead m- , UnspentWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m+ , MonadError ImportException m ) => Network -> BlockRef@@ -721,9 +737,7 @@ -> Unspent -> m () spendOutput net br th ix u = do- insertSpender- (unspentPoint u)- Spender {spenderHash = th, spenderIndex = ix}+ insertSpender (unspentPoint u) Spender {spenderHash = th, spenderIndex = ix} case scriptToAddressBS (B.Short.fromShort (unspentScript u)) of Left _ -> return () Right a -> do@@ -733,74 +747,63 @@ False a (unspentAmount u)- removeAddrUnspent a u- insertAddrTx- a- BlockTx- { blockTxHash = th- , blockTxBlock = br- }- delUnspent (unspentPoint u)+ deleteAddrUnspent a u+ insertAddrTx a BlockTx {blockTxHash = th, blockTxBlock = br}+ deleteUnspent (unspentPoint u) unspendOutput ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , UnspentRead m- , UnspentWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m+ , MonadError ImportException m ) => Network -> OutPoint -> m () unspendOutput net op = do- t <- getImportTx (outPointHash op)- o <- getTxOutput (outPointIndex op) t- s <-- case outputSpender o of- Nothing -> do- $(logErrorS) "BlockLogic" $- "Output already unspent: " <> txHashToHex (outPointHash op) <>- " " <>- fromString (show (outPointIndex op))- throwError (AlreadyUnspent op)- Just s -> return s- x <- getImportTx (spenderHash s)- deleteSpender op- let u =- Unspent- { unspentAmount = outputAmount o- , unspentBlock = transactionBlock t- , unspentScript = B.Short.toShort (outputScript o)- , unspentPoint = op- }- addUnspent u- case scriptToAddressBS (outputScript o) of- Left _ -> return ()- Right a -> do- insertAddrUnspent a u- removeAddrTx- a- BlockTx- { blockTxHash = spenderHash s- , blockTxBlock = transactionBlock x+ t <- getImportTx (outPointHash op)+ o <- getTxOutput (outPointIndex op) t+ s <-+ case outputSpender o of+ Nothing -> do+ $(logErrorS) "BlockLogic" $+ "Output already unspent: " <>+ txHashToHex (outPointHash op) <>+ " " <>+ fromString (show (outPointIndex op))+ throwError (AlreadyUnspent (cs (show op)))+ Just s -> return s+ x <- getImportTx (spenderHash s)+ deleteSpender op+ let u =+ Unspent+ { unspentAmount = outputAmount o+ , unspentBlock = transactionBlock t+ , unspentScript = B.Short.toShort (outputScript o)+ , unspentPoint = op }- increaseBalance- net- (confirmed (unspentBlock u))- False- a- (outputAmount o)+ insertUnspent u+ case scriptToAddressBS (outputScript o) of+ Left _ -> return ()+ Right a -> do+ insertAddrUnspent a u+ deleteAddrTx+ a+ BlockTx+ { blockTxHash = spenderHash s+ , blockTxBlock = transactionBlock x+ }+ increaseBalance+ (confirmed (unspentBlock u))+ False+ a+ (outputAmount o) reduceBalance ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m+ , MonadError ImportException m ) => Network -> Bool -- ^ spend or delete confirmed output@@ -813,7 +816,7 @@ Nothing -> do $(logErrorS) "BlockLogic" $ "Balance not found for address " <> addrText net a- throwError (BalanceNotFound a)+ throwError (BalanceNotFound (addrText net a)) Just b -> do when (v > amnt b) $ do $(logErrorS) "BlockLogic" $@@ -821,11 +824,12 @@ " (needs: " <> cs (show v) <> ", has: " <>- cs (show (amnt b)) <> ")"+ cs (show (amnt b)) <>+ ")" throwError $ if c- then InsufficientBalance a- else InsufficientZeroBalance a+ then InsufficientBalance (addrText net a)+ else InsufficientZeroBalance (addrText net a) setBalance b { balanceAmount =@@ -857,23 +861,19 @@ addr = case addrToString net a of Nothing -> "???"- Just x -> x+ Just x -> x increaseBalance ::- ( MonadError ImportException m- , StoreRead m+ ( StoreRead m , StoreWrite m- , BalanceRead m- , BalanceWrite m , MonadLogger m )- => Network- -> Bool -- ^ add confirmed output+ => Bool -- ^ add confirmed output -> Bool -- ^ increase total received -> Address -> Word64 -> m ()-increaseBalance _net c t a v = do+increaseBalance c t a v = do b <- getBalance a >>= \case Nothing ->@@ -913,15 +913,16 @@ else "unconfirmed" updateAddressCounts ::- (MonadError ImportException m, BalanceWrite m, BalanceRead m)- => [Address]+ (StoreWrite m, StoreRead m, Monad m, MonadError ImportException m)+ => Network+ -> [Address] -> (Word64 -> Word64) -> m ()-updateAddressCounts as f =+updateAddressCounts net as f = forM_ as $ \a -> do b <- getBalance a >>= \case- Nothing -> throwError (BalanceNotFound a)+ Nothing -> throwError (BalanceNotFound (addrText net a)) Just b -> return b setBalance b {balanceTxCount = f (balanceTxCount b)}
src/Network/Haskoin/Store/Messages.hs view
@@ -4,77 +4,79 @@ {-# LANGUAGE MultiParamTypeClasses #-} module Network.Haskoin.Store.Messages where -import Data.ByteString (ByteString)+import Data.ByteString (ByteString) import Data.Word-import Database.RocksDB (DB)+import Database.RocksDB (DB) import Haskoin import Haskoin.Node+import Network.Haskoin.Store.Data import Network.Socket import NQE+import UnliftIO.Exception+import UnliftIO.STM (TVar) + -- | Mailbox for block store. type BlockStore = Mailbox BlockMessage -- | Store mailboxes.-data Store = Store- { storeManager :: !Manager+data Store =+ Store+ { storeManager :: !Manager -- ^ peer manager mailbox- , storeChain :: !Chain+ , storeChain :: !Chain -- ^ chain header process mailbox- , storeBlock :: !BlockStore+ , storeBlock :: !BlockStore -- ^ block storage mailbox- }+ } -- | Configuration for a 'Store'.-data StoreConfig = StoreConfig- { storeConfMaxPeers :: !Int+data StoreConfig =+ StoreConfig+ { storeConfMaxPeers :: !Int -- ^ max peers to connect to- , storeConfInitPeers :: ![HostPort]+ , storeConfInitPeers :: ![HostPort] -- ^ static set of peers to connect to- , storeConfDiscover :: !Bool+ , storeConfDiscover :: !Bool -- ^ discover new peers?- , storeConfDB :: !DB+ , storeConfDB :: !LayeredDB -- ^ RocksDB database handler- , storeConfNetwork :: !Network+ , storeConfNetwork :: !Network -- ^ network constants- , storeConfListen :: !(Listen StoreEvent)- }+ , storeConfListen :: !(Listen StoreEvent)+ -- ^ listen to store events+ } -- | Configuration for a block store.-data BlockConfig = BlockConfig- { blockConfManager :: !Manager+data BlockConfig =+ BlockConfig+ { blockConfManager :: !Manager -- ^ peer manager from running node- , blockConfChain :: !Chain+ , blockConfChain :: !Chain -- ^ chain from a running node- , blockConfListener :: !(Listen StoreEvent)+ , blockConfListener :: !(Listen StoreEvent) -- ^ listener for store events- , blockConfDB :: !DB+ , blockConfDB :: !LayeredDB -- ^ RocksDB database handle- , blockConfNet :: !Network+ , blockConfNet :: !Network -- ^ network constants- }+ } -- | Messages that a 'BlockStore' can accept. data BlockMessage = BlockNewBest !BlockNode -- ^ new block header in chain- | BlockPeerConnect !Peer- !SockAddr+ | BlockPeerConnect !Peer !SockAddr -- ^ new peer connected- | BlockPeerDisconnect !Peer- !SockAddr+ | BlockPeerDisconnect !Peer !SockAddr -- ^ peer disconnected- | BlockReceived !Peer- !Block+ | BlockReceived !Peer !Block -- ^ new block received from a peer- | BlockNotFound !Peer- ![BlockHash]+ | BlockNotFound !Peer ![BlockHash] -- ^ block not found- | BlockTxReceived !Peer- !Tx+ | BlockTxReceived !Peer !Tx -- ^ transaction received from peer- | BlockTxAvailable !Peer- ![TxHash]+ | BlockTxAvailable !Peer ![TxHash] -- ^ peer has transactions available | BlockPing !(Listen ()) -- ^ internal housekeeping ping@@ -104,3 +106,30 @@ !RejectCode !ByteString -- ^ peer rejected transaction++data PubExcept+ = PubNoPeers+ | PubReject RejectCode+ | PubTimeout+ | PubNotFound+ | PubPeerDisconnected+ deriving Eq++instance Show PubExcept where+ show PubNoPeers = "no peers"+ show (PubReject c) =+ "rejected: " <>+ case c of+ RejectMalformed -> "malformed"+ RejectInvalid -> "invalid"+ RejectObsolete -> "obsolete"+ RejectDuplicate -> "duplicate"+ RejectNonStandard -> "not standard"+ RejectDust -> "dust"+ RejectInsufficientFee -> "insufficient fee"+ RejectCheckpoint -> "checkpoint"+ show PubTimeout = "timeout"+ show PubNotFound = "not found"+ show PubPeerDisconnected = "peer disconnected"++instance Exception PubExcept
+ src/Network/Haskoin/Store/Web.hs view
@@ -0,0 +1,1040 @@+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE FlexibleInstances #-}+{-# LANGUAGE LambdaCase #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE RankNTypes #-}+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE TupleSections #-}+module Network.Haskoin.Store.Web where+import Conduit hiding (runResourceT)+import Control.Applicative ((<|>))+import Control.Arrow+import Control.Exception ()+import Control.Monad+import Control.Monad.Logger+import Control.Monad.Reader (MonadReader, ReaderT)+import qualified Control.Monad.Reader as R+import Control.Monad.Trans.Maybe+import Data.Aeson.Encoding (encodingToLazyByteString,+ fromEncoding)+import Data.Bits+import Data.ByteString.Builder+import qualified Data.ByteString.Lazy as L+import qualified Data.ByteString.Lazy.Char8 as C+import Data.Char+import Data.Foldable+import Data.Function+import qualified Data.HashMap.Strict as H+import Data.List+import Data.Maybe+import Data.Serialize as Serialize+import Data.String.Conversions+import qualified Data.Text.Lazy as T+import Data.Version+import Data.Word (Word32)+import Database.RocksDB as R+import Haskoin+import Haskoin.Node+import Network.Haskoin.Store.Data+import Network.Haskoin.Store.Data.Cached+import Network.Haskoin.Store.Messages+import Network.HTTP.Types+import NQE+import qualified Paths_haskoin_store as P+import Text.Read (readMaybe)+import UnliftIO+import UnliftIO.Resource+import Web.Scotty.Internal.Types (ActionT (ActionT, runAM))+import Web.Scotty.Trans as S++type WebT m = ActionT Except (ReaderT LayeredDB m)++data WebConfig =+ WebConfig+ { webPort :: !Int+ , webNetwork :: !Network+ , webDB :: !LayeredDB+ , webPublisher :: !(Publisher StoreEvent)+ , webStore :: !Store+ }++instance Parsable BlockHash where+ parseParam =+ maybe (Left "could not decode block hash") Right . hexToBlockHash . cs++instance Parsable TxHash where+ parseParam =+ maybe (Left "could not decode tx hash") Right . hexToTxHash . cs++instance MonadIO m => StoreRead (WebT m) where+ isInitialized = lift isInitialized+ getBestBlock = lift getBestBlock+ getBlocksAtHeight = lift . getBlocksAtHeight+ getBlock = lift . getBlock+ getTxData = lift . getTxData+ getSpender = lift . getSpender+ getSpenders = lift . getSpenders+ getOrphanTx = lift . getOrphanTx+ getUnspent = lift . getUnspent+ getBalance = lift . getBalance++instance (MonadResource m, MonadUnliftIO m) =>+ StoreStream (WebT (ReaderT LayeredDB m)) where+ getMempool = transPipe lift . getMempool+ getOrphans = transPipe lift getOrphans+ getAddressUnspents a x = transPipe lift $ getAddressUnspents a x+ getAddressTxs a x = transPipe lift $ getAddressTxs a x+ getAddressBalances = transPipe lift getAddressBalances+ getUnspents = transPipe lift getUnspents++askDB :: Monad m => WebT m LayeredDB+askDB = lift R.ask++defHandler :: Monad m => Network -> Except -> WebT m ()+defHandler net e = do+ proto <- setupBin+ case e of+ ThingNotFound -> status status404+ BadRequest -> status status400+ UserError _ -> status status400+ StringError _ -> status status400+ ServerError -> status status500+ protoSerial net proto e++maybeSerial ::+ (Monad m, JsonSerial a, BinSerial a)+ => Network+ -> Bool -- ^ binary+ -> Maybe a+ -> WebT m ()+maybeSerial _ _ Nothing = raise ThingNotFound+maybeSerial net proto (Just x) = S.raw $ serialAny net proto x++protoSerial ::+ (Monad m, JsonSerial a, BinSerial a)+ => Network+ -> Bool+ -> a+ -> WebT m ()+protoSerial net proto = S.raw . serialAny net proto++scottyBestBlock :: MonadIO m => Network -> WebT m ()+scottyBestBlock net = do+ cors+ n <- parseNoTx+ proto <- setupBin+ res <-+ runMaybeT $ do+ h <- MaybeT getBestBlock+ b <- MaybeT $ getBlock h+ return $ pruneTx n b+ maybeSerial net proto res++scottyBlock :: MonadIO m => Network -> WebT m ()+scottyBlock net = do+ cors+ block <- param "block"+ n <- parseNoTx+ proto <- setupBin+ res <-+ runMaybeT $ do+ b <- MaybeT $ getBlock block+ return $ pruneTx n b+ maybeSerial net proto res++scottyBlockHeight :: MonadIO m => Network -> WebT m ()+scottyBlockHeight net = do+ cors+ height <- param "height"+ n <- parseNoTx+ proto <- setupBin+ res <-+ fmap catMaybes $ do+ hs <- getBlocksAtHeight height+ forM hs $ \h ->+ runMaybeT $ do+ b <- MaybeT $ getBlock h+ return $ pruneTx n b+ protoSerial net proto res++scottyBlockHeights :: MonadIO m => Network -> WebT m ()+scottyBlockHeights net = do+ cors+ heights <- param "heights"+ n <- parseNoTx+ proto <- setupBin+ bs <- concat <$> mapM getBlocksAtHeight (nub heights)+ res <-+ fmap catMaybes . forM bs $ \bh ->+ runMaybeT $ do+ b <- MaybeT $ getBlock bh+ return $ pruneTx n b+ protoSerial net proto res++scottyBlocks :: MonadIO m => Network -> WebT m ()+scottyBlocks net = do+ cors+ blocks <- param "blocks"+ n <- parseNoTx+ proto <- setupBin+ res <-+ fmap catMaybes . forM blocks $ \bh ->+ runMaybeT $ do+ b <- MaybeT $ getBlock bh+ return $ pruneTx n b+ protoSerial net proto res++scottyMempool :: MonadUnliftIO m => Network -> WebT m ()+scottyMempool net = do+ cors+ (l, s) <- parseLimits+ proto <- setupBin+ db <- askDB+ stream $ \io flush' -> do+ runResourceT . withLayeredDB db $+ runConduit $ getMempoolLimit l s .| streamAny net proto io+ flush'++scottyTransaction :: MonadIO m => Network -> WebT m ()+scottyTransaction net = do+ cors+ txid <- param "txid"+ proto <- setupBin+ res <- getTransaction txid+ maybeSerial net proto res++scottyRawTransaction :: MonadIO m => Bool -> WebT m ()+scottyRawTransaction hex = do+ cors+ txid <- param "txid"+ res <- getTransaction txid+ case res of+ Nothing -> raise ThingNotFound+ Just x ->+ if hex+ then text . cs . encodeHex . Serialize.encode $+ transactionData x+ else do+ S.setHeader "Content-Type" "application/octet-stream"+ S.raw $ Serialize.encodeLazy (transactionData x)++scottyTxAfterHeight :: MonadIO m => Network -> WebT m ()+scottyTxAfterHeight net = do+ cors+ txid <- param "txid"+ height <- param "height"+ proto <- setupBin+ res <- cbAfterHeight 10000 height txid+ protoSerial net proto res++scottyTransactions :: MonadIO m => Network -> WebT m ()+scottyTransactions net = do+ cors+ txids <- param "txids"+ proto <- setupBin+ res <- catMaybes <$> mapM getTransaction (nub txids)+ protoSerial net proto res++scottyRawTransactions :: MonadIO m => Bool -> WebT m ()+scottyRawTransactions hex = do+ cors+ txids <- param "txids"+ res <- catMaybes <$> mapM getTransaction (nub txids)+ if hex+ then S.json $ map (encodeHex . Serialize.encode . transactionData) res+ else do+ S.setHeader "Content-Type" "application/octet-stream"+ S.raw . L.concat $ map (Serialize.encodeLazy . transactionData) res++scottyAddressTxs :: MonadUnliftIO m => Network -> Bool -> WebT m ()+scottyAddressTxs net full = do+ cors+ a <- parseAddress net+ (l, s) <- parseLimits+ proto <- setupBin+ db <- askDB+ stream $ \io flush' -> do+ runResourceT . withLayeredDB db . runConduit $ f proto l s a io+ flush'+ where+ f proto l s a io+ | full = getAddressTxsFull l s a .| streamAny net proto io+ | otherwise = getAddressTxsLimit l s a .| streamAny net proto io++scottyAddressesTxs :: MonadUnliftIO m => Network -> Bool -> WebT m ()+scottyAddressesTxs net full = do+ cors+ as <- parseAddresses net+ (l, s) <- parseLimits+ proto <- setupBin+ db <- askDB+ stream $ \io flush' -> do+ runResourceT . withLayeredDB db . runConduit $ f proto l s as io+ flush'+ where+ f proto l s as io+ | full = getAddressesTxsFull l s as .| streamAny net proto io+ | otherwise = getAddressesTxsLimit l s as .| streamAny net proto io++scottyAddressUnspent :: MonadUnliftIO m => Network -> WebT m ()+scottyAddressUnspent net = do+ cors+ a <- parseAddress net+ (l, s) <- parseLimits+ proto <- setupBin+ db <- askDB+ stream $ \io flush' -> do+ runResourceT . withLayeredDB db . runConduit $+ getAddressUnspentsLimit l s a .| streamAny net proto io+ flush'++scottyAddressesUnspent :: MonadUnliftIO m => Network -> WebT m ()+scottyAddressesUnspent net = do+ cors+ as <- parseAddresses net+ (l, s) <- parseLimits+ proto <- setupBin+ db <- askDB+ stream $ \io flush' -> do+ runResourceT . withLayeredDB db . runConduit $+ getAddressesUnspentsLimit l s as .| streamAny net proto io+ flush'++scottyAddressBalance :: MonadIO m => Network -> WebT m ()+scottyAddressBalance net = do+ cors+ address <- parseAddress net+ proto <- setupBin+ res <-+ getBalance address >>= \case+ Just b -> return b+ Nothing ->+ return+ Balance+ { balanceAddress = address+ , balanceAmount = 0+ , balanceUnspentCount = 0+ , balanceZero = 0+ , balanceTxCount = 0+ , balanceTotalReceived = 0+ }+ protoSerial net proto res++scottyAddressesBalances :: MonadIO m => Network -> WebT m ()+scottyAddressesBalances net = do+ cors+ as <- parseAddresses net+ proto <- setupBin+ let f a Nothing =+ Balance+ { balanceAddress = a+ , balanceAmount = 0+ , balanceUnspentCount = 0+ , balanceZero = 0+ , balanceTxCount = 0+ , balanceTotalReceived = 0+ }+ f _ (Just b) = b+ res <- mapM (\a -> f a <$> getBalance a) as+ protoSerial net proto res++scottyXpubBalances :: MonadUnliftIO m => Network -> WebT m ()+scottyXpubBalances net = do+ cors+ xpub <- parseXpub net+ proto <- setupBin+ db <- askDB+ res <- liftIO . runResourceT . withLayeredDB db $ xpubBals xpub+ protoSerial net proto res++scottyXpubTxs :: MonadUnliftIO m => Network -> Bool -> WebT m ()+scottyXpubTxs net full = do+ cors+ x <- parseXpub net+ (l, s) <- parseLimits+ proto <- setupBin+ db <- askDB+ bs <- liftIO . runResourceT . withLayeredDB db $ xpubBals x+ stream $ \io flush' -> do+ runResourceT . withLayeredDB db . runConduit $ f proto l s bs io+ flush'+ where+ f proto l s bs io+ | full =+ getAddressesTxsFull l s (map (balanceAddress . xPubBal) bs) .|+ streamAny net proto io+ | otherwise =+ getAddressesTxsLimit l s (map (balanceAddress . xPubBal) bs) .|+ streamAny net proto io++scottyXpubUnspents :: MonadIO m => Network -> WebT m ()+scottyXpubUnspents net = do+ cors+ x <- parseXpub net+ proto <- setupBin+ (l, s) <- parseLimits+ db <- askDB+ stream $ \io flush' -> do+ runResourceT . withLayeredDB db . runConduit $+ xpubUnspentLimit net l s x .| streamAny net proto io+ flush'++scottyXpubSummary :: MonadUnliftIO m => Network -> WebT m ()+scottyXpubSummary net = do+ cors+ x <- parseXpub net+ (l, s) <- parseLimits+ proto <- setupBin+ db <- askDB+ res <- liftIO . runResourceT . withLayeredDB db $ xpubSummary l s x+ protoSerial net proto res++scottyPostTx ::+ (MonadUnliftIO m, MonadLoggerIO m)+ => Network+ -> Store+ -> Publisher StoreEvent+ -> WebT m ()+scottyPostTx net st pub = do+ cors+ proto <- setupBin+ b <- body+ let bin = eitherToMaybe . Serialize.decode+ hex = bin <=< decodeHex . cs . C.filter (not . isSpace)+ tx <-+ case hex b <|> bin (L.toStrict b) of+ Nothing -> raise (UserError "decode tx fail")+ Just x -> return x+ lift (publishTx net pub st tx) >>= \case+ Right () -> do+ protoSerial net proto (TxId (txHash tx))+ $(logDebugS) "Web" $+ "Success publishing tx " <> txHashToHex (txHash tx)+ Left e -> do+ case e of+ PubNoPeers -> status status500+ PubTimeout -> status status500+ PubPeerDisconnected -> status status500+ PubNotFound -> status status500+ PubReject _ -> status status400+ protoSerial net proto (UserError (show e))+ $(logErrorS) "Web" $+ "Error publishing tx " <> txHashToHex (txHash tx) <> ": " <>+ cs (show e)+ finish++scottyDbStats :: MonadIO m => WebT m ()+scottyDbStats = do+ cors+ LayeredDB {layeredDB = BlockDB {blockDB = db}} <- askDB+ lift (getProperty db Stats) >>= text . cs . fromJust++scottyEvents :: MonadUnliftIO m => Network -> Publisher StoreEvent -> WebT m ()+scottyEvents net pub = do+ cors+ proto <- setupBin+ stream $ \io flush' ->+ withSubscription pub $ \sub ->+ forever $+ flush' >> receive sub >>= \se -> do+ let me =+ case se of+ StoreBestBlock block_hash ->+ Just (EventBlock block_hash)+ StoreMempoolNew tx_hash -> Just (EventTx tx_hash)+ _ -> Nothing+ case me of+ Nothing -> return ()+ Just e ->+ let bs =+ serialAny net proto e <>+ if proto+ then mempty+ else "\n"+ in io (lazyByteString bs)++scottyPeers :: MonadIO m => Network -> Store -> WebT m ()+scottyPeers net st = do+ cors+ proto <- setupBin+ ps <- getPeersInformation (storeManager st)+ protoSerial net proto ps++scottyHealth :: MonadUnliftIO m => Network -> Store -> WebT m ()+scottyHealth net st = do+ cors+ proto <- setupBin+ h <- lift $ healthCheck net (storeManager st) (storeChain st)+ when (not (healthOK h) || not (healthSynced h)) $ status status503+ protoSerial net proto h++runWeb :: (MonadLoggerIO m, MonadUnliftIO m) => WebConfig -> m ()+runWeb WebConfig { webDB = db+ , webPort = port+ , webNetwork = net+ , webStore = st+ , webPublisher = pub+ } = do+ runner <- askRunInIO+ scottyT port (runner . withLayeredDB db) $ do+ defaultHandler (defHandler net)+ S.get "/block/best" $ scottyBestBlock net+ S.get "/block/:block" $ scottyBlock net+ S.get "/block/height/:height" $ scottyBlockHeight net+ S.get "/block/heights" $ scottyBlockHeights net+ S.get "/blocks" $ scottyBlocks net+ S.get "/mempool" $ scottyMempool net+ S.get "/transaction/:txid" $ scottyTransaction net+ S.get "/transaction/:txid/hex" $ scottyRawTransaction True+ S.get "/transaction/:txid/bin" $ scottyRawTransaction False+ S.get "/transaction/:txid/after/:height" $ scottyTxAfterHeight net+ S.get "/transactions" $ scottyTransactions net+ S.get "/transactions/hex" $ scottyRawTransactions True+ S.get "/transactions/bin" $ scottyRawTransactions False+ S.get "/address/:address/transactions" $ scottyAddressTxs net False+ S.get "/address/:address/transactions/full" $ scottyAddressTxs net True+ S.get "/address/transactions" $ scottyAddressesTxs net False+ S.get "/address/transactions/full" $ scottyAddressesTxs net True+ S.get "/address/:address/unspent" $ scottyAddressUnspent net+ S.get "/address/unspent" $ scottyAddressesUnspent net+ S.get "/address/:address/balance" $ scottyAddressBalance net+ S.get "/address/balances" $ scottyAddressesBalances net+ S.get "/xpub/:xpub/balances" $ scottyXpubBalances net+ S.get "/xpub/:xpub/transactions" $ scottyXpubTxs net False+ S.get "/xpub/:xpub/transactions/full" $ scottyXpubTxs net True+ S.get "/xpub/:xpub/unspent" $ scottyXpubUnspents net+ S.get "/xpub/:xpub" $ scottyXpubSummary net+ S.post "/transactions" $ scottyPostTx net st pub+ S.get "/dbstats" $ scottyDbStats+ S.get "/events" $ scottyEvents net pub+ S.get "/peers" $ scottyPeers net st+ S.get "/health" $ scottyHealth net st+ notFound $ raise ThingNotFound++parseLimits :: (ScottyError e, Monad m) => ActionT e m (Maybe Word32, StartFrom)+parseLimits = do+ let b = do+ height <- param "height"+ pos <- param "pos" `rescue` const (return maxBound)+ return $ StartBlock height pos+ m = do+ time <- param "time"+ return $ StartMem time+ o = do+ o <- param "offset" `rescue` const (return 0)+ return $ StartOffset o+ l <- (Just <$> param "limit") `rescue` const (return Nothing)+ s <- b <|> m <|> o+ return (l, s)++parseAddress net = do+ address <- param "address"+ case stringToAddr net address of+ Nothing -> next+ Just a -> return a++parseAddresses net = do+ addresses <- param "addresses"+ let as = mapMaybe (stringToAddr net) addresses+ unless (length as == length addresses) next+ return as++parseXpub :: (Monad m, ScottyError e) => Network -> ActionT e m XPubKey+parseXpub net = do+ t <- param "xpub"+ case xPubImport net t of+ Nothing -> next+ Just x -> return x++parseNoTx :: (Monad m, ScottyError e) => ActionT e m Bool+parseNoTx = param "notx" `rescue` const (return False)++pruneTx False b = b+pruneTx True b = b {blockDataTxs = take 1 (blockDataTxs b)}++cors :: Monad m => ActionT e m ()+cors = setHeader "Access-Control-Allow-Origin" "*"++serialAny ::+ (JsonSerial a, BinSerial a)+ => Network+ -> Bool -- ^ binary+ -> a+ -> L.ByteString+serialAny net True = runPutLazy . binSerial net+serialAny net False = encodingToLazyByteString . jsonSerial net++streamAny ::+ (JsonSerial i, BinSerial i, MonadIO m)+ => Network+ -> Bool -- ^ protobuf+ -> (Builder -> IO ())+ -> ConduitT i o m ()+streamAny net True io = binConduit net .| mapC lazyByteString .| streamConduit io+streamAny net False io = jsonListConduit net .| streamConduit io++jsonListConduit :: (JsonSerial a, Monad m) => Network -> ConduitT a Builder m ()+jsonListConduit net =+ yield "[" >> mapC (fromEncoding . jsonSerial net) .| intersperseC "," >> yield "]"++binConduit :: (BinSerial i, Monad m) => Network -> ConduitT i L.ByteString m ()+binConduit net = mapC (runPutLazy . binSerial net)++streamConduit :: MonadIO m => (i -> IO ()) -> ConduitT i o m ()+streamConduit io = mapM_C (liftIO . io)++setupBin :: Monad m => ActionT Except m Bool+setupBin =+ let p = do+ setHeader "Content-Type" "application/octet-stream"+ return True+ j = do+ setHeader "Content-Type" "application/json"+ return False+ in S.header "accept" >>= \case+ Nothing -> j+ Just x ->+ if is_binary x+ then p+ else j+ where+ is_binary = (== "application/octet-stream")++instance MonadLoggerIO m => MonadLoggerIO (WebT m) where+ askLoggerIO = lift askLoggerIO++instance MonadLogger m => MonadLogger (WebT m) where+ monadLoggerLog loc src lvl = lift . monadLoggerLog loc src lvl++healthCheck ::+ (MonadUnliftIO m, StoreRead m)+ => Network+ -> Manager+ -> Chain+ -> m HealthCheck+healthCheck net mgr ch = do+ n <- timeout (5 * 1000 * 1000) $ chainGetBest ch+ b <-+ runMaybeT $ do+ h <- MaybeT getBestBlock+ MaybeT $ getBlock h+ p <- timeout (5 * 1000 * 1000) $ managerGetPeers mgr+ let k = isNothing n || isNothing b || maybe False (not . null) p+ s =+ isJust $ do+ x <- n+ y <- b+ guard $ nodeHeight x - blockDataHeight y <= 1+ return+ HealthCheck+ { healthBlockBest = headerHash . blockDataHeader <$> b+ , healthBlockHeight = blockDataHeight <$> b+ , healthHeaderBest = headerHash . nodeHeader <$> n+ , healthHeaderHeight = nodeHeight <$> n+ , healthPeers = length <$> p+ , healthNetwork = getNetworkName net+ , healthOK = k+ , healthSynced = s+ }++-- | Obtain information about connected peers from peer manager process.+getPeersInformation :: MonadIO m => Manager -> m [PeerInformation]+getPeersInformation mgr = mapMaybe toInfo <$> managerGetPeers mgr+ where+ toInfo op = do+ ver <- onlinePeerVersion op+ let as = onlinePeerAddress op+ ua = getVarString $ userAgent ver+ vs = version ver+ sv = services ver+ rl = relay ver+ return+ PeerInformation+ { peerUserAgent = ua+ , peerAddress = as+ , peerVersion = vs+ , peerServices = sv+ , peerRelay = rl+ }++xpubBals ::+ (MonadResource m, MonadUnliftIO m, StoreRead m) => XPubKey -> m [XPubBal]+xpubBals xpub = do+ (rk, ss) <- allocate (newTVarIO []) (\as -> readTVarIO as >>= mapM_ cancel)+ stp0 <- newTVarIO False+ stp1 <- newTVarIO False+ q0 <- newTBQueueIO 20+ q1 <- newTBQueueIO 20+ xs <-+ withAsync (go stp0 ss q0 0) $ \_ ->+ withAsync (go stp1 ss q1 1) $ \_ ->+ withAsync (red ss stp0 q0) $ \r0 ->+ withAsync (red ss stp1 q1) $ \r1 -> do+ xs0 <- wait r0+ xs1 <- wait r1+ return $ xs0 <> xs1+ release rk+ return xs+ where+ stp e =+ readTVarIO e >>= \s ->+ if s+ then return ()+ else await >>= \case+ Nothing -> return ()+ Just x -> yield x >> stp e+ go e ss q m =+ runConduit $+ yieldMany (as m) .| stp e .| mapMC (uncurry (b ss)) .| conduitToQueue q+ red ss e q = runConduit $ queueToConduit q .| f ss e 0 .| sinkList+ b ss a p = mask_ $ do+ s <-+ async $+ getBalance a >>= \case+ Nothing -> return Nothing+ Just b' -> return $ Just XPubBal {xPubBalPath = p, xPubBal = b'}+ atomically $ modifyTVar ss (s :)+ return s+ as m = map (\(a, _, n') -> (a, [m, n'])) (deriveAddrs (pubSubKey xpub m) 0)+ f ss e n+ | n < 20 =+ await >>= \case+ Just a ->+ wait a >>= \case+ Nothing -> f ss e (n + 1)+ Just b -> yield b >> f ss e 0+ Nothing -> return ()+ | otherwise = do+ atomically $ writeTVar e True+ await >>= \case+ Just a -> do+ cancel a+ atomically $ modifyTVar ss (Data.List.delete a)+ f ss e n+ Nothing -> return ()++xpubUnspent ::+ ( MonadResource m+ , MonadUnliftIO m+ , StoreStream m+ , StoreRead m+ )+ => Network+ -> Maybe BlockRef+ -> XPubKey+ -> ConduitT () XPubUnspent m ()+xpubUnspent net mbr xpub = do+ (_, as) <-+ lift $ allocate (newTVarIO []) (\as -> readTVarIO as >>= mapM_ cancel)+ xs <-+ lift $ do+ bals <- xpubBals xpub+ forM bals $ \XPubBal {xPubBalPath = p, xPubBal = b} ->+ mask_ $ do+ q <- newTBQueueIO 10+ a <-+ async . runConduit $+ getAddressUnspents (balanceAddress b) mbr .| mapC (f p) .|+ conduitToQueue q+ atomically $ modifyTVar as (a :)+ return $ queueToConduit q+ mergeSourcesBy (flip compare `on` (unspentBlock . xPubUnspent)) xs+ where+ f p t = XPubUnspent {xPubUnspentPath = p, xPubUnspent = t}++xpubUnspentLimit ::+ ( MonadResource m+ , MonadUnliftIO m+ , StoreStream m+ , StoreRead m+ )+ => Network+ -> Maybe Word32+ -> StartFrom+ -> XPubKey+ -> ConduitT () XPubUnspent m ()+xpubUnspentLimit net l s x =+ xpubUnspent net (mbr s) x .| (offset s >> limit l)++xpubSummary ::+ (MonadResource m, MonadUnliftIO m, StoreStream m, StoreRead m)+ => Maybe Word32+ -> StartFrom+ -> XPubKey+ -> m XPubSummary+xpubSummary l s x = do+ bs <- xpubBals x+ let f XPubBal {xPubBalPath = p, xPubBal = Balance {balanceAddress = a}} =+ (a, p)+ pm = H.fromList $ map f bs+ txs <-+ runConduit $+ getAddressesTxsFull l s (map (balanceAddress . xPubBal) bs) .| sinkList+ let as =+ nub+ [ a+ | t <- txs+ , let is = transactionInputs t+ , let os = transactionOutputs t+ , let ais =+ mapMaybe+ (eitherToMaybe . scriptToAddressBS . inputPkScript)+ is+ , let aos =+ mapMaybe+ (eitherToMaybe . scriptToAddressBS . outputScript)+ os+ , a <- ais ++ aos+ ]+ ps = H.fromList $ mapMaybe (\a -> (a, ) <$> H.lookup a pm) as+ ex = foldl max 0 [i | XPubBal {xPubBalPath = [x, i]} <- bs, x == 0]+ ch = foldl max 0 [i | XPubBal {xPubBalPath = [x, i]} <- bs, x == 1]+ return+ XPubSummary+ { xPubSummaryReceived =+ sum (map (balanceTotalReceived . xPubBal) bs)+ , xPubSummaryConfirmed = sum (map (balanceAmount . xPubBal) bs)+ , xPubSummaryZero = sum (map (balanceZero . xPubBal) bs)+ , xPubSummaryPaths = ps+ , xPubSummaryTxs = txs+ , xPubChangeIndex = ch+ , xPubExternalIndex = ex+ }++-- | Check if any of the ancestors of this transaction is a coinbase after the+-- specified height. Returns 'Nothing' if answer cannot be computed before+-- hitting limits.+cbAfterHeight ::+ (MonadIO m, StoreRead m)+ => Int -- ^ how many ancestors to test before giving up+ -> BlockHeight+ -> TxHash+ -> m TxAfterHeight+cbAfterHeight d h t+ | d <= 0 = return $ TxAfterHeight Nothing+ | otherwise = do+ x <- fmap snd <$> tst d t+ return $ TxAfterHeight x+ where+ tst e x+ | e <= 0 = return Nothing+ | otherwise = do+ let e' = e - 1+ getTransaction x >>= \case+ Nothing -> return Nothing+ Just tx ->+ if any isCoinbase (transactionInputs tx)+ then return $+ Just (e', blockRefHeight (transactionBlock tx) > h)+ else case transactionBlock tx of+ BlockRef {blockRefHeight = b}+ | b <= h -> return $ Just (e', False)+ _ ->+ r e' . nub $+ map+ (outPointHash . inputPoint)+ (transactionInputs tx)+ r e [] = return $ Just (e, False)+ r e (n:ns) =+ tst e n >>= \case+ Nothing -> return Nothing+ Just (e', s) ->+ if s+ then return $ Just (e', True)+ else r e' ns++-- Snatched from:+-- https://github.com/cblp/conduit-merge/blob/master/src/Data/Conduit/Merge.hs+mergeSourcesBy ::+ (Foldable f, Monad m)+ => (a -> a -> Ordering)+ -> f (ConduitT () a m ())+ -> ConduitT i a m ()+mergeSourcesBy f = mergeSealed . fmap sealConduitT . toList+ where+ mergeSealed sources = do+ prefetchedSources <- lift $ traverse ($$++ await) sources+ go [(a, s) | (s, Just a) <- prefetchedSources]+ go [] = pure ()+ go sources = do+ let (a, src1):sources1 = sortBy (f `on` fst) sources+ yield a+ (src2, mb) <- lift $ src1 $$++ await+ let sources2 =+ case mb of+ Nothing -> sources1+ Just b -> (b, src2) : sources1+ go sources2++getMempoolLimit ::+ (Monad m, StoreStream m)+ => Maybe Word32+ -> StartFrom+ -> ConduitT () TxHash m ()+getMempoolLimit _ StartBlock {} = return ()+getMempoolLimit l (StartMem t) =+ getMempool (Just t) .| mapC snd .| limit l+getMempoolLimit l s =+ getMempool Nothing .| mapC snd .| (offset s >> limit l)++getAddressTxsLimit ::+ (Monad m, StoreStream m)+ => Maybe Word32+ -> StartFrom+ -> Address+ -> ConduitT () BlockTx m ()+getAddressTxsLimit l s a =+ getAddressTxs a (mbr s) .| (offset s >> limit l)++getAddressTxsFull ::+ (Monad m, StoreStream m, StoreRead m)+ => Maybe Word32+ -> StartFrom+ -> Address+ -> ConduitT () Transaction m ()+getAddressTxsFull l s a =+ getAddressTxsLimit l s a .| concatMapMC (getTransaction . blockTxHash)++getAddressesTxsLimit ::+ (MonadResource m, MonadUnliftIO m, StoreStream m)+ => Maybe Word32+ -> StartFrom+ -> [Address]+ -> ConduitT () BlockTx m ()+getAddressesTxsLimit l s addrs = do+ (_, ss) <-+ lift $ allocate (newTVarIO []) (\ss -> readTVarIO ss >>= mapM_ cancel)+ xs <-+ lift $ do+ forM addrs $ \addr -> mask_ $ do+ q <- newTBQueueIO 10+ a <-+ async . runConduit $+ getAddressTxs addr (mbr s) .| conduitToQueue q+ atomically $ modifyTVar ss (a :)+ return $ queueToConduit q+ mergeSourcesBy (flip compare `on` blockTxBlock) xs .| dedup .|+ (offset s >> limit l)++getAddressesTxsFull ::+ (MonadResource m, MonadUnliftIO m, StoreStream m, StoreRead m)+ => Maybe Word32+ -> StartFrom+ -> [Address]+ -> ConduitT () Transaction m ()+getAddressesTxsFull l s as =+ getAddressesTxsLimit l s as .| concatMapMC (getTransaction . blockTxHash)++getAddressUnspentsLimit ::+ (Monad m, StoreStream m)+ => Maybe Word32+ -> StartFrom+ -> Address+ -> ConduitT () Unspent m ()+getAddressUnspentsLimit l s a =+ getAddressUnspents a (mbr s) .| (offset s >> limit l)++getAddressesUnspentsLimit ::+ (Monad m, StoreStream m)+ => Maybe Word32+ -> StartFrom+ -> [Address]+ -> ConduitT () Unspent m ()+getAddressesUnspentsLimit l s as =+ mergeSourcesBy+ (flip compare `on` unspentBlock)+ (map (`getAddressUnspents` mbr s) as) .|+ (offset s >> limit l)++offset :: Monad m => StartFrom -> ConduitT i i m ()+offset (StartOffset o) = dropC (fromIntegral o)+offset _ = return ()++limit :: Monad m => Maybe Word32 -> ConduitT i i m ()+limit Nothing = mapC id+limit (Just n) = takeC (fromIntegral n)++mbr :: StartFrom -> Maybe BlockRef+mbr (StartBlock h p) = Just (BlockRef h p)+mbr (StartMem t) = Just (MemRef t)+mbr (StartOffset _) = Nothing++conduitToQueue :: MonadIO m => TBQueue (Maybe a) -> ConduitT a Void m ()+conduitToQueue q =+ await >>= \case+ Just x -> atomically (writeTBQueue q (Just x)) >> conduitToQueue q+ Nothing -> atomically $ writeTBQueue q Nothing++queueToConduit :: MonadIO m => TBQueue (Maybe a) -> ConduitT () a m ()+queueToConduit q =+ atomically (readTBQueue q) >>= \case+ Just x -> yield x >> queueToConduit q+ Nothing -> return ()++dedup :: (Eq i, Monad m) => ConduitT i i m ()+dedup =+ let dd Nothing =+ await >>= \case+ Just x -> do+ yield x+ dd (Just x)+ Nothing -> return ()+ dd (Just x) =+ await >>= \case+ Just y+ | x == y -> dd (Just x)+ | otherwise -> do+ yield y+ dd (Just y)+ Nothing -> return ()+ in dd Nothing++-- | Publish a new transaction to the network.+publishTx ::+ (MonadUnliftIO m, StoreRead m)+ => Network+ -> Publisher StoreEvent+ -> Store+ -> Tx+ -> m (Either PubExcept ())+publishTx net pub st tx = do+ e <-+ withSubscription pub $ \s ->+ getTransaction (txHash tx) >>= \case+ Just _ -> do+ return $ Right ()+ Nothing -> go s+ return e+ where+ go s = do+ managerGetPeers (storeManager st) >>= \case+ [] -> do+ return $ Left PubNoPeers+ OnlinePeer {onlinePeerMailbox = p, onlinePeerAddress = a}:_ -> do+ MTx tx `sendMessage` p+ let t =+ if getSegWit net+ then InvWitnessTx+ else InvTx+ sendMessage+ (MGetData (GetData [InvVector t (getTxHash (txHash tx))]))+ p+ f p s+ t = 15 * 1000 * 1000+ f p s = do+ liftIO (timeout t (g p s)) >>= \case+ Nothing -> do+ return $ Left PubTimeout+ Just (Left e) -> do+ return $ Left e+ Just (Right ()) -> do+ return $ Right ()+ g p s =+ receive s >>= \case+ StoreTxReject p' h' c _+ | p == p' && h' == txHash tx -> return . Left $ PubReject c+ StorePeerDisconnected p' _+ | p == p' -> return $ Left PubPeerDisconnected+ StoreMempoolNew h'+ | h' == txHash tx -> return $ Right ()+ _ -> g p s
test/Spec.hs view
@@ -15,7 +15,7 @@ import UnliftIO data TestStore = TestStore- { testStoreDB :: !DB+ { testStoreDB :: !LayeredDB , testStoreBlockStore :: !BlockStore , testStoreChain :: !Chain , testStoreEvents :: !(Inbox StoreEvent)@@ -39,7 +39,7 @@ bestHeight `shouldBe` 8 it "get a block and its transactions" $ withTestStore net "get-block-txs" $ \TestStore {..} ->- withBlockDB defaultReadOptions testStoreDB $ do+ withLayeredDB testStoreDB $ do let h1 = "e8588129e146eeb0aa7abdc3590f8c5920cc5ff42daf05c23b29d4ae5b51fc22" h2 =@@ -69,9 +69,9 @@ withTestStore :: MonadUnliftIO m => Network -> String -> (TestStore -> m a) -> m a withTestStore net t f =- withSystemTempDirectory ("haskoin-store-test-" <> t <> "-") $ \w -> do- x <- newInbox+ withSystemTempDirectory ("haskoin-store-test-" <> t <> "-") $ \w -> runNoLoggingT $ do+ x <- newInbox db <- open w@@ -80,12 +80,21 @@ , errorIfExists = True , compression = SnappyCompression }+ let ldb =+ LayeredDB+ { layeredDB =+ BlockDB+ { blockDB = db+ , blockDBopts = defaultReadOptions+ }+ , layeredCache = Nothing+ } let cfg = StoreConfig { storeConfMaxPeers = 20 , storeConfInitPeers = [] , storeConfDiscover = True- , storeConfDB = db+ , storeConfDB = ldb , storeConfNetwork = net , storeConfListen = (`sendSTM` x) }@@ -93,7 +102,7 @@ lift $ f TestStore- { testStoreDB = db+ { testStoreDB = ldb , testStoreBlockStore = storeBlock , testStoreChain = storeChain , testStoreEvents = x