hdbi-conduit 1.2.0 → 1.3.0
raw patch · 3 files changed
+186/−98 lines, 3 filesdep ~hdbidep ~hdbi-sqlite
Dependency ranges changed: hdbi, hdbi-sqlite
Files
- Data/Conduit/HDBI.hs +135/−83
- hdbi-conduit.cabal +4/−4
- testsrc/runtests.hs +47/−11
Data/Conduit/HDBI.hs view
@@ -1,35 +1,40 @@ {-# LANGUAGE BangPatterns , TypeFamilies+, ScopedTypeVariables #-} module Data.Conduit.HDBI ( -- * Conduit functions selectAll- , selectAllRows- , selectRawAll- , selectRawAllRows , insertAll- , insertAllRows , insertAllCount- , insertAllRowsCount+ , insertAllTrans+ , flushAt+ , flushBy -- * Auxiliary conduit functions , statementSource , statementSink , statementSinkCount+ , statementSinkTrans -- * ResourceT functions , allocConnection , allocStmt , executeStmt- , executeStmtRow- , executeStmtRaw+ -- * Stream typing helpers+ , asSqlVals+ , asThisType ) where +import Control.Exception (try, throw, SomeException(..))+import Control.Monad (when) import Control.Monad.IO.Class+import Control.Monad.Trans.Class (lift) import Control.Monad.Trans.Resource import Data.Conduit import Database.HDBI+import qualified Data.Conduit.List as L allocConnection :: (Connection con, MonadResource m) => IO con -> m (ReleaseKey, con) allocConnection con = allocate con disconnect@@ -37,89 +42,55 @@ allocStmt :: (Statement stmt, MonadResource m) => IO stmt -> m (ReleaseKey, stmt) allocStmt stmt = allocate stmt finish -executeStmt :: (Connection con, (ConnStatement con) ~ stmt, MonadResource m) => con -> Query -> [SqlValue] -> m (ReleaseKey, stmt)-executeStmt con query vals = do+executeStmt :: (Connection con, (ConnStatement con) ~ stmt, ToRow row, MonadResource m)+ => con -> Query -> row -> m (ReleaseKey, stmt)+executeStmt con query row = do (key, stmt) <- allocStmt $ prepare con query- liftIO $ execute stmt vals+ liftIO $ execute stmt row return (key, stmt) -executeStmtRow :: (Connection con, (ConnStatement con) ~ stmt, ToRow row, MonadResource m) => con -> Query -> row -> m (ReleaseKey, stmt)-executeStmtRow con query row = do- (key, stmt) <- allocStmt $ prepare con query- liftIO $ executeRow stmt row- return (key, stmt) -executeStmtRaw :: (Connection con, (ConnStatement con) ~ stmt, MonadResource m) => con -> Query -> m (ReleaseKey, stmt)-executeStmtRaw con query = do- (key, stmt) <- allocStmt $ prepare con query- liftIO $ executeRaw stmt- return (key, stmt)---- | fetch all results of query-selectAll :: (Connection con, MonadResource m)+-- | Execute query and stream result+selectAll :: (Connection con, MonadResource m, FromRow row, ToRow params) => con- -> Query -- query to execute- -> [SqlValue] -- query parameters- -> Source m [SqlValue]-selectAll con query params = statementSource fetch $ do- st <- prepare con query- execute st params- return st---- | same as `selectAll` but reburn stream of `FromRow` instances-selectAllRows :: (Connection con, MonadResource m, FromRow a)- => con- -> Query- -> [SqlValue]- -> Source m a-selectAllRows con query params = statementSource fetchRow $ do- st <- prepare con query- execute st params- return st---- | same as `selectAll` but without query parameters-selectRawAll :: (Connection con, MonadResource m)- => con- -> Query- -> Source m [SqlValue]-selectRawAll con query = statementSource fetch $ do- st <- prepare con query- executeRaw st- return st---- | same as `selectRawAll` but return stream of `FromRow` instances-selectRawAllRows :: (Connection con, MonadResource m, FromRow a) => con -> Query -> Source m a-selectRawAllRows con query = statementSource fetchRow $ do- st <- prepare con query- executeRaw st- return st+ -> Query+ -> params+ -> Source m row+selectAll con query params = statementSource (liftIO . fetch) execStmt+ where+ execStmt = do+ st <- prepare con query+ (r :: Either SomeException ()) <- try $ execute st params+ case r of+ Left e -> do+ finish st+ throw e+ Right _ -> return st -- | same as `insertAll` but also count executed rows-insertAllCount :: (Connection con, MonadResource m, Num count) => con -> Query -> Sink [SqlValue] m count-insertAllCount con query = statementSinkCount execute $ prepare con query+insertAllCount :: (Connection con, MonadResource m, Num count, ToRow a) => con -> Query -> Sink a m count+insertAllCount con query = statementSinkCount rowPutter $ prepare con query --- | same as `insertAllRows` but also count executed rows-insertAllRowsCount :: (Connection con, MonadResource m, Num count, ToRow a) => con -> Query -> Sink a m count-insertAllRowsCount con query = statementSinkCount executeRow $ prepare con query --- | perform `execute` for each bunch of values-insertAll :: (Connection con, MonadResource m)+-- | perform `executeRow` for each input row+insertAll :: (Connection con, MonadResource m, ToRow a) => con -> Query- -> Sink [SqlValue] m ()-insertAll con query = statementSink execute $ prepare con query+ -> Sink a m ()+insertAll con query = statementSink rowPutter $ prepare con query --- | perfor `executeRow` for each input row-insertAllRows :: (Connection con, MonadResource m, ToRow a)- => con- -> Query- -> Sink a m ()-insertAllRows con query = statementSink executeRow $ prepare con query+-- | Execute query on each (Chunk row) and commit on each Flush. The last query+-- is always commit, so be carefull.+insertAllTrans :: (Connection con, MonadResource m, ToRow a)+ => con+ -> Query+ -> Sink (Flush a) m ()+insertAllTrans con query = statementSinkTrans con rowPutter $ prepare con query -- | Get all values from the statement until action return ''Just a'' statementSource :: (Statement stmt, MonadResource m)- => (stmt -> IO (Maybe a)) -- action to execute until it return+ => (stmt -> m (Maybe a)) -- action to execute until it return -- Nothing -> IO stmt -- statement constructor -> Source m a@@ -129,7 +100,7 @@ statementSource' where statementSource' st = do- row <- liftIO $ getter st+ row <- lift $ getter st case row of Nothing -> return () Just r -> do@@ -139,7 +110,7 @@ -- | Execute action many times with given thread of values, return the count -- of executions statementSinkCount :: (Statement stmt, MonadResource m, Num count)- => (stmt -> a -> IO ()) -- action to execute each time+ => (stmt -> a -> m ()) -- action to execute each time -> IO stmt -- statement constructor -> Sink a m count statementSinkCount putter stmt = bracketP@@ -152,14 +123,12 @@ case next of Nothing -> return ac Just n -> do- liftIO $ do- reset st- putter st n+ lift $ putter st n statementSinkCount' (ac+1) st -- | Same as `statementSinkCount` but without counting, just return () statementSink :: (Statement stmt, MonadResource m)- => (stmt -> a -> IO ())+ => (stmt -> a -> m ()) -> IO stmt -> Sink a m () statementSink putter stmt = bracketP@@ -172,7 +141,90 @@ case next of Nothing -> return () Just n -> do- liftIO $ do- reset st- putter st n+ lift $ putter st n statementSink' st++-- | Execute each chunk with putter function and commit transaction on each+-- flush. The last action is always commit.+statementSinkTrans :: (Connection con, (ConnStatement con) ~ stmt, MonadResource m, MonadIO m)+ => con+ -> (stmt -> a -> m ())+ -> IO stmt+ -> Sink (Flush a) m ()+statementSinkTrans con putter stmt = do+ intrans <- liftIO $ inTransaction con+ bracketP+ stmt+ finish+ $ statementSinkTrans' intrans+ where+ statementSinkTrans' intrans st = do+ next <- await+ case next of+ Nothing -> do+ when intrans $ liftIO $ commit con+ return ()+ Just n -> case n of+ Flush -> do+ when intrans $ liftIO $ commit con+ statementSinkTrans' False st+ Chunk val -> do+ when (not intrans) $ liftIO $ begin con+ lift $ putter st val+ statementSinkTrans' True st++-- | The behaviour is the same as `groupBy` function. Each time when prefix+-- return False the conduit yields Flush.+flushBy :: (Monad m) => (a -> a -> Bool) -> Conduit a m (Flush a)+flushBy pref = flushBy' Nothing+ where+ flushBy' !lst = case lst of+ Nothing -> do+ n <- await+ case n of+ Nothing -> return ()+ Just x -> do+ yield $ Chunk x+ flushBy' $ Just x+ Just l -> do+ n <- await+ case n of+ Nothing -> return ()+ Just x -> do+ when (not $ pref l x) $ yield Flush+ yield $ Chunk x+ flushBy' $ Just x++-- | separate each `i` chunks with Flush+flushAt :: (Monad m, Integral i) => i -> Conduit a m (Flush a)+flushAt cnt = flushAt' c+ where+ c = max 1 cnt+ flushAt' 0 = do+ yield Flush+ flushAt' c+ flushAt' !count = do+ n <- await+ case n of+ Nothing -> return ()+ Just x -> do+ yield $ Chunk x+ flushAt' $ count - 1+++rowPutter :: (MonadIO m, Statement stmt, ToRow row)+ => stmt+ -> row+ -> m ()+rowPutter st row = liftIO $ do+ reset st+ execute st row++-- | Function to fuse when no data convertion is needed+asSqlVals :: (Monad m) => Conduit [SqlValue] m [SqlValue]+asSqlVals = L.map id++-- | To specify actual stream type by fusing with it. The value of argument is+-- not used+asThisType :: (Monad m) => a -> Conduit a m a+asThisType _ = L.map id
hdbi-conduit.cabal view
@@ -2,7 +2,7 @@ -- documentation, see http://haskell.org/cabal/users-guide/ name: hdbi-conduit-version: 1.2.0+version: 1.3.0 synopsis: Conduit glue for HDBI -- description: license: BSD3@@ -29,7 +29,7 @@ ghc-prof-options: -fprof-auto build-depends: base < 5 , conduit- , hdbi >= 1.2.0+ , hdbi >= 1.3.0 , resourcet , transformers @@ -43,8 +43,8 @@ build-depends: base < 5 , QuickCheck , conduit- , hdbi >= 1.2.0- , hdbi-sqlite >= 1.2.0+ , hdbi >= 1.3.0+ , hdbi-sqlite >= 1.3.0 , quickcheck-assertions , resourcet , test-framework
testsrc/runtests.hs view
@@ -32,18 +32,54 @@ allTests c = testGroup "All tests" [ testProperty "Insert + fold" $ insertFold c , testProperty "Insert + copy" $ insertCopy c- , testProperty "Insert + copy + sum" $ insertCopySum c]+ , testProperty "Insert + copy + sum" $ insertCopySum c+ , testProperty "Insert trans fluahAt" $ insertTransFlushAt c+ , testProperty "Insert trans flushBy" $ insertTransFlushBy c+ ] sumPairs :: (Num a, Num b) => (a, b) -> (a, b) -> (a, b) sumPairs (!a, !b) (!x, !y) = (a+x, b+y) ++insertTransFlushAt :: SQliteConnection -> Positive Int -> [(Integer, Integer)] -> Property+insertTransFlushAt c count vals = M.monadicIO $ do+ (Just res, tr) <- M.run $ do+ runRaw c "delete from values1"+ runResourceT+ $ L.sourceList vals+ $= (flushAt $ getPositive count)+ $$ insertAllTrans c "insert into values1 (val1, val2) values (?,?)"+ r <- runFetchOne c "select count(*) from values1" ()+ tr <- inTransaction c+ return (r, tr)+ _ <- M.stop $ res ?== (length vals)+ M.stop $ tr ?== False++insertTransFlushBy :: SQliteConnection -> NonEmptyList Integer -> Property+insertTransFlushBy con vals = M.monadicIO $ do+ (tr, r) <- M.run $ do+ runRaw con "delete from values1"+ runResourceT+ $ L.sourceList nvals+ $= flushBy signFlush+ $= L.map (fmap one) -- flush is the functor+ $$ insertAllTrans con "insert into values1 (val1) values (?)"+ tr <- inTransaction con+ Just r <- runFetchOne con "select sum(val1) from values1" ()+ return (tr, r)+ _ <- M.stop $ tr ?== False+ M.stop $ r ?== (sum nvals)+ where+ nvals = getNonEmpty vals+ signFlush a b = (signum a) == (signum b)+ insertFold :: SQliteConnection -> [(Integer, Integer)] -> Property insertFold c vals = M.monadicIO $ do res <- M.run $ withTransaction c $ do runRaw c "delete from values1"- runManyRows c "insert into values1(val1, val2) values (?,?)" vals+ runMany c "insert into values1(val1, val2) values (?,?)" vals runResourceT- $ selectRawAllRows c "select val1, val2 from values1"+ $ selectAll c "select val1, val2 from values1" () $$ L.fold sumPairs (0 :: Integer, 0 :: Integer) M.stop $ res ?== (foldl' sumPairs (0, 0) vals) @@ -54,9 +90,9 @@ runRaw c "delete from values2" runResourceT $ L.sourceList vals- $$ insertAllRows c "insert into values1(val1, val2) values (?,?)"+ $$ insertAll c "insert into values1(val1, val2) values (?,?)" runResourceT- $ selectRawAll c "select val1, val2 from values1"+ $ selectAll c "select val1, val2 from values1" () $= asThisType (undefined :: (Int, Int)) $$ insertAllCount c "insert into values2(val1, val2) values (?,?)" M.stop $ res == (length vals) @@ -67,19 +103,19 @@ "values2", "values3"] runResourceT- $ L.sourceList (map (\(a, b) -> [toSql a, toSql b]) vals)+ $ L.sourceList vals $$ insertAll c "insert into values1(val1, val2) values (?,?)" runResourceT- $ selectRawAll c "select val1, val2 from values1"+ $ selectAll c "select val1, val2 from values1" () $= asSqlVals $$ insertAll c "insert into values2(val1, val2) values (?,?)" runResourceT $ (U.zip- (selectRawAllRows c "select val1, val2 from values1")- (selectRawAllRows c "select val1, val2 from values2"))+ (selectAll c "select val1, val2 from values1" ())+ (selectAll c "select val1, val2 from values2" ())) $= L.map (\(a, b :: (Integer, Integer)) -> sumPairs a b)- $$ insertAllRows c "insert into values3(val1, val2) values (?,?)"+ $$ insertAll c "insert into values3(val1, val2) values (?,?)" runResourceT- $ selectRawAllRows c "select val1, val2 from values3"+ $ selectAll c "select val1, val2 from values3" () $$ L.fold sumPairs (0 :: Integer, 0 :: Integer) let (a, b) = foldl' sumPairs (0, 0) vals M.stop $ (a*2, b*2) ==? res