packages feed

eventsource-stub-store 1.0.2 → 1.0.3

raw patch · 5 files changed

+134/−182 lines, 5 filesdep +asyncdep +transformers-basedep ~protoludePVP: major bump suggested

API removals or changes: PVP suggests a major version bump

Dependencies added: async, transformers-base

Dependency ranges changed: protolude

API changes (from Hackage documentation)

- EventSource.Store.Stub: [streamNextNumber] :: Stream -> EventNumber
+ EventSource.Store.Stub: [streamNumber] :: Stream -> EventNumber

Files

CHANGELOG.md view
@@ -1,6 +1,10 @@+1.0.3+=====+  * Fix `ExpectedVersion` logic.+ 1.0.2 =====-  = Fix Stackkage build+  * Fix Stackkage build  1.0.1 =====
eventsource-stub-store.cabal view
@@ -1,9 +1,11 @@--- This file has been generated from package.yaml by hpack version 0.17.1.+-- This file has been generated from package.yaml by hpack version 0.20.0. -- -- see: https://github.com/sol/hpack+--+-- hash: 6a21808b5f7d73e5d2c75c4c4a57bf948ececa55aac722e20385c26e3c37cf75  name:           eventsource-stub-store-version:        1.0.2+version:        1.0.3 synopsis:       An in-memory stub store implementation. description:    An in-memory stub store implementation. category:       Eventsourcing@@ -21,7 +23,6 @@     LICENSE.md     package.yaml     README.md-    stack.yaml  source-repository head   type: git@@ -30,17 +31,19 @@ library   hs-source-dirs:       library-  default-extensions: NoImplicitPrelude   ghc-options: -Wall   build-depends:-      base >=4.9 && <5-    , eventsource-api >=1.1-    , protolude >=0.1.10 && <0.3+      async+    , base >=4.9 && <5     , containers+    , eventsource-api >=1.1     , mtl     , stm+    , transformers-base   exposed-modules:       EventSource.Store.Stub+  other-modules:+      Paths_eventsource_stub_store   default-language: Haskell2010  test-suite eventsource-stub-store-test-suite@@ -50,15 +53,16 @@       test-suite   ghc-options: -Wall -rtsopts -threaded -with-rtsopts=-N   build-depends:-      base+      aeson+    , base+    , eventsource-api     , eventsource-store-specs ==1.*     , eventsource-stub-store+    , protolude     , tasty     , tasty-hspec-    , protolude-    , aeson-    , eventsource-api   other-modules:       Test.EventSource.Event       Test.EventSource.Store.Stub+      Paths_eventsource_stub_store   default-language: Haskell2010
library/EventSource/Store/Stub.hs view
@@ -1,3 +1,4 @@+{-# LANGUAGE FlexibleContexts  #-} {-# LANGUAGE OverloadedStrings #-} {-# LANGUAGE RecordWildCards   #-} --------------------------------------------------------------------------------@@ -23,11 +24,18 @@   ) where  ---------------------------------------------------------------------------------import Control.Concurrent.STM-import qualified Data.Map.Strict as M-import Data.Sequence (Seq, (|>))+import           Control.Concurrent.STM (STM)+import qualified Control.Concurrent.STM as STM+import Data.Foldable (toList, for_, foldl')++--------------------------------------------------------------------------------+import           Control.Concurrent.Async (async)+import           Control.Monad.Base (MonadBase, liftBase)+import           Control.Monad.State.Strict (execState)+import           Data.Map.Strict (Map)+import qualified Data.Map.Strict as Map+import           Data.Sequence (Seq, (|>)) import qualified Data.Sequence as S-import Protolude  -------------------------------------------------------------------------------- import EventSource.Store@@ -36,36 +44,85 @@ -------------------------------------------------------------------------------- -- | Holds stream state data. data Stream =-  Stream { streamNextNumber :: EventNumber+  Stream { streamNumber :: EventNumber          , streamEvents :: Seq SavedEvent          }  ---------------------------------------------------------------------------------type Sub = TChan SavedEvent+type Sub = STM.TChan SavedEvent type Subs = Map SubscriptionId Sub  -------------------------------------------------------------------------------- data StubStore =-  StubStore { _streams :: TVar (Map StreamName Stream)-            , _subs :: TVar (Map StreamName Subs)+  StubStore { _streams :: STM.TVar (Map StreamName Stream)+            , _subs :: STM.TVar (Map StreamName Subs)             }  --------------------------------------------------------------------------------+data Version+  = NoStreamYet+  | Version EventNumber++--------------------------------------------------------------------------------+getCurrentVersion :: StreamName -> Map StreamName Stream -> Version+getCurrentVersion name m =+  case Map.lookup name m of+    Nothing -> NoStreamYet+    Just s  -> Version $ streamNumber s++--------------------------------------------------------------------------------+data Registered =+  Registered { _regSavedEvents :: !(Seq SavedEvent)+             , _regNumber      :: !EventNumber+             , _regNewMap      :: !(Map StreamName Stream)+             }++--------------------------------------------------------------------------------+registerEvents :: StreamName+               -> [Event]+               -> Map StreamName Stream+               -> Registered+registerEvents name xs m =+  case Map.lookup name m of+    Nothing     -> appendStream (Stream (-1) mempty)+    Just stream -> appendStream stream+  where+    appendStream stream =+      let+          cur       = streamNumber stream+          nextNum   = cur + fromIntegral (length xs)+          saved     = appEvents cur (streamEvents stream)+          newStream = Stream nextNum saved+          regd      = Registered saved nextNum (Map.insert name newStream m)++       in regd++    appEvents ver xss =+      let go (acc, cur) evt =+            let next   = cur + 1+                saved  = SavedEvent next evt Nothing+                newAcc = acc |> saved++             in (newAcc, next)++       in fst $ foldl' go (xss,ver) xs++-------------------------------------------------------------------------------- -- | Creates a new stub event store. newStub :: IO StubStore-newStub = StubStore <$> newTVarIO mempty <*> newTVarIO mempty+newStub = StubStore <$> STM.newTVarIO mempty <*> STM.newTVarIO mempty  -------------------------------------------------------------------------------- -- | Returns current 'StubStore' streams state. streams :: StubStore -> IO (Map StreamName Stream)-streams StubStore{..} = readTVarIO _streams+streams StubStore{..} = STM.readTVarIO _streams  -------------------------------------------------------------------------------- -- | Returns the last event of stream. lastStreamEvent :: StubStore -> StreamName -> IO (Maybe SavedEvent) lastStreamEvent stub name = do   streamMap <- streams stub-  return (go =<< M.lookup name streamMap)+  return (go =<< Map.lookup name streamMap)     where       go stream =         case S.viewr $ streamEvents stream of@@ -76,37 +133,22 @@ -- | Returns all subscriptions a stream has. subscriptionIds :: StubStore -> StreamName -> IO [SubscriptionId] subscriptionIds StubStore{..} name = do-  subMap <- readTVarIO _subs-  case M.lookup name subMap of+  subMap <- STM.readTVarIO _subs+  case Map.lookup name subMap of     Nothing -> return []-    Just subs -> return $ M.keys subs-----------------------------------------------------------------------------------appendStream :: [Event] -> Stream -> Stream-appendStream = flip $ foldl' go-  where-    go s e =-      let num = streamNextNumber s-          evts = streamEvents s in-      s { streamNextNumber = num + 1-        , streamEvents = evts |> SavedEvent num e Nothing-        }-----------------------------------------------------------------------------------newStream :: [Event] -> Stream-newStream xs = appendStream xs (Stream 0 mempty)+    Just subs -> return $ Map.keys subs  ---------------------------------------------------------------------------------notifySubs :: StubStore -> StreamName -> [SavedEvent] -> STM ()+notifySubs :: Foldable f => StubStore -> StreamName -> f SavedEvent -> STM () notifySubs StubStore{..} name events = do-  subMap <- readTVar _subs-  for_ (M.lookup name subMap) $ \subs ->+  subMap <- STM.readTVar _subs+  for_ (Map.lookup name subMap) $ \subs ->     for_ subs $ \sub ->       for_ events $ \e ->-        writeTChan sub e+        STM.writeTChan sub e  ---------------------------------------------------------------------------------buildEvent :: (EncodeEvent a, MonadIO m) => a -> m Event+buildEvent :: (EncodeEvent a, MonadBase IO m) => a -> m Event buildEvent a = do   eid <- freshEventId   let start = Event { eventType = ""@@ -121,87 +163,58 @@ instance Store StubStore where   appendEvents self@StubStore{..} name ver xs = do     events <- traverse buildEvent xs-    liftIO $ async $ atomically $ do-      streamMap <- readTVar _streams--      case M.lookup name streamMap of-        Nothing -> do-          case ver of-            StreamExists ->-              throwSTM $ ExpectedVersionException ver NoStream-            ExactVersion v ->-              unless (v == 0) $ throwSTM-                              $ ExpectedVersionException ver NoStream-            _ -> return ()--          let _F Nothing  = Just $ newStream events-              _F (Just s) = Just $ appendStream events s--              newStreamMap = M.alter _F name streamMap--          writeTVar _streams newStreamMap+    liftBase $ async $ STM.atomically $ do+      streamMap <- STM.readTVar _streams -          -- This part is already performed in 'appendStream' but difficult-          -- to take its logic apart from building 'SavedEvent's.-          let saved = (\(num, evt) -> SavedEvent num evt Nothing) <$> zip [0..] events-          notifySubs self name saved-          let Just last = getLast $ foldMap (Last . Just) saved-              nextNum = eventNumber last + 1-          return nextNum+      let persistEvents =+            do let regd = registerEvents name events streamMap +               notifySubs self name $ _regSavedEvents regd+               (_regNumber regd) <$ STM.writeTVar _streams (_regNewMap regd) -        Just stream -> do-          let currentNumber = streamNextNumber stream+      case getCurrentVersion name streamMap of+        NoStreamYet ->           case ver of-            NoStream ->-              throwSTM $ ExpectedVersionException ver StreamExists-            ExactVersion v ->-              unless (v == streamNextNumber stream - 1)-                $ throwSTM-                $ ExpectedVersionException ver (ExactVersion currentNumber)--            _ -> return ()--          let nextStream = appendStream events stream-              newStreamMap = M.adjust (const nextStream) name streamMap--          writeTVar _streams newStreamMap--          -- This part is already performed in 'appendStream' but difficult-          -- to take its logic apart from building 'SavedEvent's.-          let saved = (\(num, evt) -> SavedEvent num evt Nothing) <$> zip [currentNumber..] events-          notifySubs self name saved-          let Just last = getLast $ foldMap (Last . Just) saved-              nextNum = eventNumber last + 1-          return nextNum-+            NoStream   -> persistEvents+            AnyVersion -> persistEvents+            _          -> STM.throwSTM $ ExpectedVersionException ver NoStream+        Version num ->+          case ver of+            AnyVersion   -> persistEvents+            StreamExists -> persistEvents+            NoStream     -> STM.throwSTM $ ExpectedVersionException ver+                                         $ ExactVersion num+            ExactVersion expVer+              | expVer == num -> persistEvents+              | otherwise     -> STM.throwSTM $ ExpectedVersionException ver+                                              $ ExactVersion num -  readBatch StubStore{..} name (Batch start _) = liftIO $ async $ atomically $ do-    streamMap <- readTVar _streams-    case M.lookup name streamMap of-      Nothing -> return $ ReadFailure StreamNotFound+  readBatch StubStore{..} name (Batch start _) = liftBase $ async $ STM.atomically $ do+    streamMap <- STM.readTVar _streams+    case Map.lookup name streamMap of+      Nothing -> return $ ReadFailure $ StreamNotFound name       Just stream -> do         let events = S.filter ((>= start) . eventNumber) $ streamEvents stream             slice = Slice { sliceEvents = toList events                           , sliceEndOfStream = True-                          , sliceNextEventNumber = streamNextNumber stream+                          , sliceNextEventNumber = streamNumber stream                           }          return $ ReadSuccess slice    subscribe StubStore{..} name = do     sid <- freshSubscriptionId-    liftIO $ atomically $ do-      chan <- newTChan-      let sub = Subscription sid $ liftIO $ atomically $ do-            saved <- readTChan chan+    liftBase $ STM.atomically $ do+      chan <- STM.newTChan+      let sub = Subscription sid $ liftBase $ STM.atomically $ do+            saved <- STM.readTChan chan             return $ Right saved -      subMap <- readTVar _subs-      let _F Nothing  = Just $ M.singleton sid  chan-          _F (Just m) = Just $ M.insert sid chan m+      subMap <- STM.readTVar _subs+      let _F Nothing  = Just $ Map.singleton sid  chan+          _F (Just m) = Just $ Map.insert sid chan m -          nextSubMap = M.alter _F name subMap+          nextSubMap = Map.alter _F name subMap -      writeTVar _subs nextSubMap+      STM.writeTVar _subs nextSubMap       return sub
package.yaml view
@@ -8,19 +8,17 @@ - LICENSE.md - package.yaml - README.md-- stack.yaml ghc-options: -Wall github: YoEight/eventsource-api library:-  default-extensions:-    - NoImplicitPrelude   dependencies:   - base >=4.9 && <5   - eventsource-api >=1.1-  - protolude >=0.1.10 && <0.3   - containers   - mtl   - stm+  - transformers-base+  - async   source-dirs: library license: BSD3 license-file: LICENSE.md@@ -45,4 +43,4 @@     - -with-rtsopts=-N     main: Main.hs     source-dirs: test-suite-version: '1.0.2'+version: '1.0.3'
− stack.yaml
@@ -1,67 +0,0 @@-# This file was automatically generated by 'stack init'-#-# Some commonly used options have been documented as comments in this file.-# For advanced use and comprehensive documentation of the format, please see:-# http://docs.haskellstack.org/en/stable/yaml_configuration/--# Resolver to choose a 'specific' stackage snapshot or a compiler version.-# A snapshot resolver dictates the compiler version and the set of packages-# to be used for project dependencies. For example:-#-# resolver: lts-3.5-# resolver: nightly-2015-09-21-# resolver: ghc-7.10.2-# resolver: ghcjs-0.1.0_ghc-7.10.2-# resolver:-#  name: custom-snapshot-#  location: "./custom-snapshot.yaml"-resolver: nightly-2017-08-15-compiler: ghc-8.2.1--# User packages to be built.-# Various formats can be used as shown in the example below.-#-# packages:-# - some-directory-# - https://example.com/foo/bar/baz-0.0.2.tar.gz-# - location:-#    git: https://github.com/commercialhaskell/stack.git-#    commit: e7b331f14bcffb8367cd58fbfc8b40ec7642100a-# - location: https://github.com/commercialhaskell/stack/commit/e7b331f14bcffb8367cd58fbfc8b40ec7642100a-#   extra-dep: true-#  subdirs:-#  - auto-update-#  - wai-#-# A package marked 'extra-dep: true' will only be built if demanded by a-# non-dependency (i.e. a user package), and its test suites and benchmarks-# will not be run. This is useful for tweaking upstream packages.-packages:-- '.'-# Dependency packages to be pulled from upstream that are not in the resolver-# (e.g., acme-missiles-0.3)-extra-deps: []--# Override default flag values for local packages and extra-deps-flags: {}--# Extra package databases containing global packages-extra-package-dbs: []--# Control whether we use the GHC we find on the path-# system-ghc: true-#-# Require a specific version of stack, using version ranges-# require-stack-version: -any # Default-# require-stack-version: ">=1.2"-#-# Override the architecture used by stack, especially useful on Windows-# arch: i386-# arch: x86_64-#-# Extra directories used by stack for building-# extra-include-dirs: [/path/to/dir]-# extra-lib-dirs: [/path/to/dir]-#-# Allow a newer minor version of GHC than the snapshot specifies-# compiler-check: newer-minor