cqrs-sqlite3 0.7.1 → 0.8.0
raw patch · 3 files changed
+48/−49 lines, 3 filesdep +conduitdep −enumeratordep ~cqrsPVP ok
version bump matches the API change (PVP)
Dependencies added: conduit
Dependencies removed: enumerator
Dependency ranges changed: cqrs
API changes (from Hackage documentation)
- Data.CQRS.EventStore.Backend.Sqlite3Utils: enumQueryResult :: Database -> String -> [SQLData] -> Enumerator [SQLData] IO b
+ Data.CQRS.EventStore.Backend.Sqlite3Utils: instance Eq State
+ Data.CQRS.EventStore.Backend.Sqlite3Utils: sourceQuery :: Database -> String -> [SQLData] -> Source IO [SQLData]
Files
- cqrs-sqlite3.cabal +3/−3
- src/Data/CQRS/EventStore/Backend/Sqlite3.hs +18/−24
- src/Data/CQRS/EventStore/Backend/Sqlite3Utils.hs +27/−22
cqrs-sqlite3.cabal view
@@ -1,5 +1,5 @@ Name: cqrs-sqlite3-Version: 0.7.1+Version: 0.8.0 Synopsis: SQLite3 backend for the cqrs package. Description: SQLite3 backend for the cqrs package. License: MIT@@ -14,8 +14,8 @@ build-depends: base == 4.* , bytestring >= 0.9.0.1 , cereal >= 0.3.3 && < 0.4- , cqrs >= 0.7.1 && < 0.8- , enumerator >= 0.4.15 && < 0.5+ , conduit >= 0.1 && < 0.2+ , cqrs >= 0.8.0 && < 0.9 , direct-sqlite >= 1.1 && < 1.2 extensions: ScopedTypeVariables ghc-options: -Wall
src/Data/CQRS/EventStore/Backend/Sqlite3.hs view
@@ -5,14 +5,13 @@ import Control.Monad (when, forM_, liftM) import Data.ByteString (ByteString)-import Data.CQRS.EventStore.Backend (EventStoreBackend(..), RawEvent)-import Data.CQRS.EventStore.Backend.Sqlite3Utils (withTransaction, execSql, enumQueryResult)+import Data.Conduit (Source, ($=), ($$), runResourceT)+import qualified Data.Conduit.List as CL+import Data.CQRS.EventStore.Backend (EventStoreBackend(..), RawEvent, RawSnapshot(..))+import Data.CQRS.EventStore.Backend.Sqlite3Utils (withTransaction, execSql, sourceQuery) import Data.CQRS.GUID (GUID) import Data.CQRS.PersistedEvent (PersistedEvent(..))-import Data.CQRS.Serialize (decode')-import Data.Enumerator (Enumerator, (>>==), ($=), run_)-import qualified Data.Enumerator.List as EL-import Data.Serialize (encode)+import Data.CQRS.Serialize (decode', encode) import qualified Database.SQLite3 as SQL import Database.SQLite3 (Database, SQLData(..)) import Prelude hiding (catch)@@ -79,8 +78,7 @@ let unpackColumns [ SQLInteger v ] = v unpackColumns columns = error $ badQueryResultMsg [show guid] columns -- Get the current version number of the aggregate.- curVer <- run_ $ EL.fold (\x -> max x . unpackColumns) 0 >>==- (enumQueryResult database getCurrentVersionSql [toSQLData guid])+ curVer <- runResourceT $ (sourceQuery database getCurrentVersionSql [toSQLData guid]) $$ CL.fold (\x -> max x . unpackColumns) 0 -- Sanity check current version number. when (fromIntegral curVer /= originatingVersion) $@@ -101,19 +99,17 @@ , toSQLData $ peGlobalVer e ] -retrieveEvents :: Database -> GUID -> Int -> IO [RawEvent]+retrieveEvents :: Database -> GUID -> Int -> Source IO RawEvent retrieveEvents database guid v0 = do -- Unpack the columns into tuples. let unpackColumns [SQLInteger version, SQLInteger gversion, SQLBlob eventData] = PersistedEvent guid eventData (fromIntegral version) (fromIntegral gversion) unpackColumns columns = error $ badQueryResultMsg [show guid, show v0] columns -- Find events with version numbers.- run_ $ EL.consume >>==- (enumQueryResult database selectEventsSql [toSQLData guid, toSQLData v0] $=- (EL.map unpackColumns))+ fmap unpackColumns $ sourceQuery database selectEventsSql [toSQLData guid, toSQLData v0] -enumerateAllEvents :: Database -> Int -> Enumerator RawEvent IO a+enumerateAllEvents :: Database -> Int -> Source IO RawEvent enumerateAllEvents database minVersion = do- enumQueryResult database enumerateAllEventsSql [toSQLData minVersion] $= EL.map+ sourceQuery database enumerateAllEventsSql [toSQLData minVersion] $= CL.map (\columns -> do case columns of [ SQLInteger gv, SQLBlob g, SQLInteger v, SQLBlob ed ] ->@@ -121,23 +117,21 @@ _ -> error $ badQueryResultMsg [show minVersion] columns) -writeSnapshot :: Database -> GUID -> (Int, ByteString) -> IO ()-writeSnapshot database guid (v,a) = do+writeSnapshot :: Database -> GUID -> RawSnapshot -> IO ()+writeSnapshot database guid (RawSnapshot v d) = do execSql database writeSnapshotSql [ toSQLData guid- , toSQLData a+ , toSQLData d , toSQLData v ] -getLatestSnapshot :: Database -> GUID -> IO (Maybe (Int, ByteString))+getLatestSnapshot :: Database -> GUID -> IO (Maybe RawSnapshot) getLatestSnapshot database guid = do -- Unpack columns from result.- let unpackColumns :: [SQLData] -> Maybe (Int,ByteString)- unpackColumns [SQLBlob a, SQLInteger v] = Just (fromIntegral v, a)- unpackColumns columns = error $ badQueryResultMsg [show guid] columns+ let unpackColumns [SQLBlob d, SQLInteger v] = Just $ RawSnapshot (fromIntegral v) d+ unpackColumns columns = error $ badQueryResultMsg [show guid] columns -- Run the query.- run_ $ EL.fold const Nothing >>==- (enumQueryResult database selectSnapshotSql [toSQLData guid] $= (EL.map unpackColumns))+ runResourceT $ (sourceQuery database selectSnapshotSql [toSQLData guid] $= (CL.map unpackColumns)) $$ CL.fold const Nothing getLatestVersion :: Database -> IO Int getLatestVersion database = do@@ -146,7 +140,7 @@ unpackColumns [SQLInteger v] = fromIntegral v unpackColumns columns = error $ badQueryResultMsg [] columns -- Run the query.- liftM head $ run_ $ EL.consume >>== (enumQueryResult database getLatestVersionSql [] $= (EL.map unpackColumns))+ liftM head $ runResourceT $ (sourceQuery database getLatestVersionSql [] $= (CL.map unpackColumns)) $$ CL.consume -- | Open an SQLite3-based event store using the named SQLite database file. -- The database file is created if it does not exist.
src/Data/CQRS/EventStore/Backend/Sqlite3Utils.hs view
@@ -1,13 +1,15 @@ {-| Implementation of an SQLite3-based event store. -} module Data.CQRS.EventStore.Backend.Sqlite3Utils- ( enumQueryResult+ ( sourceQuery , execSql , withTransaction ) where -import Control.Exception (catch, bracket, finally, onException, SomeException)-import Data.Enumerator (Enumerator, Iteratee(..), tryIO, continue, (>>==), Stream(..))-import Data.Enumerator.Internal (checkContinue0)+import Control.Exception (catch, bracket, onException, SomeException)+import Control.Monad (liftM, when)+import Data.Conduit (Source)+import qualified Data.Conduit as C+import Data.IORef (newIORef, readIORef, writeIORef) import qualified Database.SQLite3 as SQL import Database.SQLite3 (Database, Statement, SQLData(..), StepResult(..)) import Prelude hiding (catch)@@ -34,26 +36,29 @@ _ <- SQL.step stmt return () -enumQueryResults :: Statement -> Enumerator [SQLData] IO b-enumQueryResults stmt = checkContinue0 $ \loop k ->- do- nextResult <- tryIO $ SQL.step stmt- case nextResult of- Done -> continue k- Row -> do- cols <- tryIO $ SQL.columns stmt- k (Chunks [cols]) >>== loop+data State = Unbound+ | Bound+ deriving (Eq) --- | Perform a query and enumerate the results. Each result is a list--- of returned columns.-enumQueryResult :: Database -> String -> [SQLData] -> Enumerator [SQLData] IO b-enumQueryResult database sql parameters step = do- stmt <- tryIO $ SQL.prepare database sql- Iteratee $ finally- (do+sourceQuery :: Database -> String -> [SQLData] -> Source IO [SQLData]+sourceQuery database sql parameters =+ C.sourceIO+ (do+ stateRef <- newIORef Unbound+ stmt <- SQL.prepare database sql+ return (stateRef, stmt))+ (\(_,stmt) -> SQL.finalize stmt)+ (\(stateRef,stmt) -> do+ -- Bind parameters if necessary.+ state <- readIORef stateRef+ when (state == Unbound) $ do SQL.bind stmt parameters- runIteratee $ enumQueryResults stmt step)- (SQL.finalize stmt)+ writeIORef stateRef $ Bound+ -- Fetch results.+ nextResult <- SQL.step stmt+ case nextResult of+ Done -> return C.Closed+ Row -> liftM C.Open $ SQL.columns stmt) -- | Execute an IO action with an active transaction. withTransaction :: Database -> IO a -> IO a