eventstore 0.13.1.5 → 0.13.1.6
raw patch · 4 files changed
+116/−39 lines, 4 filesdep +asyncPVP ok
version bump matches the API change (PVP)
Dependencies added: async
API changes (from Hackage documentation)
Files
- CHANGELOG.markdown +4/−0
- Database/EventStore.hs +76/−37
- eventstore.cabal +3/−2
- tests/Tests.hs +33/−0
CHANGELOG.markdown view
@@ -1,3 +1,7 @@+0.13.1.6+--------+* Fix a deadlock on user side when asking a catchup subscription on a stream that doesn't exist.+ 0.13.1.5 -------- * Fix an internal connection bug.
Database/EventStore.hs view
@@ -1,7 +1,8 @@-{-# LANGUAGE DataKinds #-}-{-# LANGUAGE DeriveDataTypeable #-}-{-# LANGUAGE OverloadedStrings #-}-{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE DeriveDataTypeable #-}+{-# LANGUAGE ExistentialQuantification #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE RecordWildCards #-} -------------------------------------------------------------------------------- -- | -- Module : Database.EventStore@@ -302,13 +303,20 @@ -------------------------------------------------------------------------------- _hasCaughtUp :: Subscription S.Catchup -> STM Bool _hasCaughtUp Subscription{..} = do- SubState sm _ <- readTVar _subVar- return $ S.hasCaughtUp sm+ res <- _subState+ case res of+ SubState sm _ -> return $ S.hasCaughtUp sm+ SubException e -> throwSTM e -------------------------------------------------------------------------------- -- | Tracks a 'Subcription' lifecycle. It holds a 'Subscription' state machine -- and `SubDropReason` if any.-data SubState a = SubState (S.Subscription a) (Maybe S.SubDropReason)+data SubState a+ = SubState (S.Subscription a) (Maybe S.SubDropReason)+ | forall e. Exception e => SubException e+ -- ^ Hack used to cover a special exception that can arise from+ -- a catchup subscription for instance. One example is when asking a+ -- catchup subscription on stream that doesn't exist yet. -------------------------------------------------------------------------------- -- | It's possible to subscribe to a stream and be notified when new events are@@ -322,11 +330,12 @@ -- * 'S.Persistent' data Subscription a = Subscription- { _subVar :: TVar (SubState a)- , _subRun :: TMVar S.Running- , _subStream :: Text- , _subProd :: Production- , _subInner :: a+ { _subState :: STM (SubState a)+ , _subSetState :: SubState a -> STM ()+ , _subRun :: TMVar S.Running+ , _subStream :: Text+ , _subProd :: Production+ , _subInner :: a } --------------------------------------------------------------------------------@@ -393,17 +402,20 @@ -------------------------------------------------------------------------------- _nextEventMaybe :: Subscription a -> STM (Maybe ResolvedEvent) _nextEventMaybe Subscription{..} = do- SubState sub close <- readTVar _subVar- run <- readTMVar _subRun- let (res, nxt) = S.readNext sub- case res of- Nothing -> do- case close of- Nothing -> return Nothing- Just err -> throwSTM $ SubscriptionClosed run err- Just e -> do- writeTVar _subVar $ SubState nxt close- return $ Just e+ st <- _subState+ case st of+ SubException e -> throwSTM e+ SubState sub close -> do+ run <- readTMVar _subRun+ let (res, nxt) = S.readNext sub+ case res of+ Nothing -> do+ case close of+ Nothing -> return Nothing+ Just err -> throwSTM $ SubscriptionClosed run err+ Just e -> do+ _subSetState $ SubState nxt close+ return $ Just e -------------------------------------------------------------------------------- -- | Acknowledges those event ids have been successfully processed.@@ -461,6 +473,7 @@ -> SubState a -> SubState a modifySubSM k (SubState sm r) = SubState (k sm) r+modifySubSM _ s = s -------------------------------------------------------------------------------- -- | This exception is raised when the user tries to get the next event from a@@ -667,7 +680,10 @@ writeTVar var $ SubState sm (Just r) cb = createSubAsync mk recv send dropped pushConnectStream _prod cb stream_id res_lnk_tos- return $ Subscription var mvar stream_id _prod (S.Regular res_lnk_tos)+ let getState = readTVar var+ setState = writeTVar var+ return $ Subscription getState setState mvar stream_id _prod+ (S.Regular res_lnk_tos) -------------------------------------------------------------------------------- -- | Subcribes to $all stream.@@ -713,27 +729,47 @@ -> Op.CatchupState -> IO (Subscription S.Catchup) subscribeFromCommon Connection{..} stream_id res_lnk_tos batch_m tpe = do- mvar <- newEmptyTMVarIO- var <- newTVarIO $ SubState S.catchupSubscription Nothing+ mvarRun <- newEmptyTMVarIO+ mvarSub <- newEmptyTMVarIO let readFrom res = case res of- Left _ -> return ()+ -- We want to notify the user that something went wrong in the+ -- first phase of a catchup subscription (e.g. reading the+ -- stream forward until we catchup to stream's end). This+ -- prevents a deadlock on user side in case where the user calls+ -- `waitTillCatchup` on a stream that doesn't exist.+ Left e -> atomically $ do+ isEmpty <- isEmptyTMVar mvarSub+ if isEmpty+ then putTMVar mvarSub (SubException e)+ else () <$ swapTMVar mvarSub (SubException e)+ Right (xs, eos, chk) -> atomically $ do- s <- readTVar var- let nxt_s = modifySubSM (S.batchRead xs eos chk) s- writeTVar var nxt_s- mk = putTMVar mvar- rcv = readTVar var- send = writeTVar var+ -- When a catchup subscription receives events for the+ -- first time.+ whenM (isEmptyTMVar mvarSub) $ do+ let initState = SubState S.catchupSubscription Nothing+ putTMVar mvarSub initState++ subState <- takeTMVar mvarSub+ let nxtSubState =+ modifySubSM (S.batchRead xs eos chk) subState++ putTMVar mvarSub nxtSubState+ mk = putTMVar mvarRun+ rcv = readTMVar mvarSub+ send = \x -> () <$ swapTMVar mvarSub x dropped r = do- SubState sm _ <- readTVar var- writeTVar var $ SubState sm (Just r)+ SubState sm _ <- takeTMVar mvarSub+ putTMVar mvarSub $ SubState sm (Just r) op = Op.catchup _settings tpe res_lnk_tos batch_m cb = createSubAsync mk rcv send dropped pushOperation _prod readFrom op pushConnectStream _prod cb stream_id res_lnk_tos- return $ Subscription var mvar stream_id _prod S.Catchup+ let getState = readTMVar mvarSub+ setState = \x -> () <$ swapTMVar mvarSub x+ return $ Subscription getState setState mvarRun stream_id _prod S.Catchup -------------------------------------------------------------------------------- -- | Asynchronously sets the metadata for a stream.@@ -823,7 +859,10 @@ writeTVar var $ SubState sm (Just r) cb = createSubAsync mk recv send dropped pushConnectPersist _prod cb group stream bufSize- return $ Subscription var mvar stream _prod (S.Persistent group)+ let getState = readTVar var+ setState = writeTVar var+ return $ Subscription getState setState mvar stream _prod+ (S.Persistent group) -------------------------------------------------------------------------------- createOpAsync :: IO (Either OperationError a -> IO (), Async a)
eventstore.cabal view
@@ -10,7 +10,7 @@ -- PVP summary: +-+------- breaking API changes -- | | +----- non-breaking API additions -- | | | +--- code changes with no API change-version: 0.13.1.5+version: 0.13.1.6 tested-with: GHC >= 7.8.3 && < 7.11 @@ -146,4 +146,5 @@ stm, time, dotnet-timespan,- connection+ connection,+ async
tests/Tests.hs view
@@ -16,10 +16,12 @@ module Tests where --------------------------------------------------------------------------------+import Control.Concurrent import Control.Exception import Data.Maybe (catMaybes) --------------------------------------------------------------------------------+import Control.Concurrent.Async import Data.Aeson import Data.DotNet.TimeSpan import Test.Tasty@@ -41,6 +43,8 @@ , testCase "Real $all backward" $ readAllEventsBackwardTest conn , testCase "Subscription test" $ subscribeTest conn , testCase "Subscription from test" $ subscribeFromTest conn+ , testCase "Subscription from catchup not blocking" $+ subscribeFromNoStreamTest conn , testCase "Set Stream Metadata" $ setStreamMetadataTest conn , testCase "Get Stream Metadata" $ getStreamMetadataTest conn , testCase "Create persistent sub" $ createPersistentTest conn@@ -224,6 +228,35 @@ return False res <- catch action $ \(_ :: SubscriptionClosed) -> return True assertBool "Should have raised an exception" res++--------------------------------------------------------------------------------+data SubNoStreamTest+ = SubNoStreamTestSuccess+ | SubNoStreamTestWrongException+ | SubNoStreamTestTimeout+ deriving (Eq, Show)++--------------------------------------------------------------------------------+secs :: Int+secs = 1000 * 1000++--------------------------------------------------------------------------------+subscribeFromNoStreamTest :: Connection -> IO ()+subscribeFromNoStreamTest conn = do+ sub <- subscribeFrom conn "non-existent-stream" False Nothing Nothing+ let subAction = do+ res <- try $ waitTillCatchup sub+ case res of+ Left InvalidOperation{} -> return SubNoStreamTestSuccess+ _ -> return SubNoStreamTestWrongException+ timeout = do+ threadDelay (10 * secs)+ return SubNoStreamTestTimeout++ res <- race subAction timeout+ case res of+ Left r -> assertEqual "Wrong test result" SubNoStreamTestSuccess r+ Right r -> assertEqual "Wrong test result" SubNoStreamTestSuccess r -------------------------------------------------------------------------------- setStreamMetadataTest :: Connection -> IO ()