packages feed

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 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