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 +5/−1
- eventsource-stub-store.cabal +15/−11
- library/EventSource/Store/Stub.hs +111/−98
- package.yaml +3/−5
- stack.yaml +0/−67
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