eventstore 0.7.2.1 → 0.8.0.0
raw patch · 13 files changed
+1586/−581 lines, 13 filesdep ~aesonPVP ok
version bump matches the API change (PVP)
Dependency ranges changed: aeson
API changes (from Hackage documentation)
- Database.EventStore: catchupAwait :: Catchup -> IO (Either CatchupError ResolvedEvent)
- Database.EventStore: catchupStream :: Catchup -> Text
- Database.EventStore: catchupUnsubscribe :: Catchup -> IO ()
- Database.EventStore: subAwait :: Subscription -> IO (Either DropReason ResolvedEvent)
- Database.EventStore: subStream :: Subscription -> Text
+ Database.EventStore: DispatchToSingle :: SystemConsumerStrategy
+ Database.EventStore: InvalidOperation :: Text -> OperationException
+ Database.EventStore: NA_Park :: NakAction
+ Database.EventStore: NA_Retry :: NakAction
+ Database.EventStore: NA_Skip :: NakAction
+ Database.EventStore: NA_Stop :: NakAction
+ Database.EventStore: NA_Unknown :: NakAction
+ Database.EventStore: PersistentSubscriptionSettings :: !Bool -> !Int32 -> !Bool -> !TimeSpan -> !Int32 -> !Int32 -> !Int32 -> !Int32 -> !TimeSpan -> !Int32 -> !Int32 -> !Int32 -> !SystemConsumerStrategy -> PersistentSubscriptionSettings
+ Database.EventStore: RoundRobin :: SystemConsumerStrategy
+ Database.EventStore: class Identifiable a
+ Database.EventStore: connectToPersistentSubscription :: Connection -> Text -> Text -> Int32 -> IO (Async (Subscription Persistent))
+ Database.EventStore: createPersistentSubscription :: Connection -> Text -> Text -> PersistentSubscriptionSettings -> IO (Async ())
+ Database.EventStore: data NakAction
+ Database.EventStore: data Persistent
+ Database.EventStore: data PersistentSubscriptionSettings
+ Database.EventStore: data Regular
+ Database.EventStore: data SystemConsumerStrategy
+ Database.EventStore: defaultPersistentSubscriptionSettings :: PersistentSubscriptionSettings
+ Database.EventStore: deletePersistentSubscription :: Connection -> Text -> Text -> IO (Async ())
+ Database.EventStore: notifyEventsFailed :: Subscription Persistent -> NakAction -> Maybe Text -> [UUID] -> IO ()
+ Database.EventStore: notifyEventsProcessed :: Subscription Persistent -> [UUID] -> IO ()
+ Database.EventStore: psSettingsCheckPointAfter :: PersistentSubscriptionSettings -> !TimeSpan
+ Database.EventStore: psSettingsExtraStats :: PersistentSubscriptionSettings -> !Bool
+ Database.EventStore: psSettingsHistoryBufSize :: PersistentSubscriptionSettings -> !Int32
+ Database.EventStore: psSettingsLiveBufSize :: PersistentSubscriptionSettings -> !Int32
+ Database.EventStore: psSettingsMaxCheckPointCount :: PersistentSubscriptionSettings -> !Int32
+ Database.EventStore: psSettingsMaxRetryCount :: PersistentSubscriptionSettings -> !Int32
+ Database.EventStore: psSettingsMaxSubsCount :: PersistentSubscriptionSettings -> !Int32
+ Database.EventStore: psSettingsMinCheckPointCount :: PersistentSubscriptionSettings -> !Int32
+ Database.EventStore: psSettingsMsgTimeout :: PersistentSubscriptionSettings -> !TimeSpan
+ Database.EventStore: psSettingsNamedConsumerStrategy :: PersistentSubscriptionSettings -> !SystemConsumerStrategy
+ Database.EventStore: psSettingsReadBatchSize :: PersistentSubscriptionSettings -> !Int32
+ Database.EventStore: psSettingsResolveLinkTos :: PersistentSubscriptionSettings -> !Bool
+ Database.EventStore: psSettingsStartFrom :: PersistentSubscriptionSettings -> !Int32
+ Database.EventStore: resolvedEventOriginalId :: ResolvedEvent -> Maybe UUID
+ Database.EventStore: s_logger :: Settings -> Maybe (Log -> IO ())
+ Database.EventStore: streamMetadataGetCustomProperty :: FromJSON a => StreamMetadata -> Text -> Maybe a
+ Database.EventStore: streamMetadataGetCustomPropertyValue :: StreamMetadata -> Text -> Maybe Value
+ Database.EventStore: subIsSubscribedToAll :: Subscription a -> Bool
+ Database.EventStore: subNextEvent :: Subscription a -> (IO (NextEvent a))
+ Database.EventStore: subStreamId :: Subscription a -> Text
+ Database.EventStore: timeSpanFromDays :: Double -> TimeSpan
+ Database.EventStore: timeSpanFromHours :: Double -> TimeSpan
+ Database.EventStore: timeSpanFromMinutes :: Double -> TimeSpan
+ Database.EventStore: timeSpanFromSeconds :: Double -> TimeSpan
+ Database.EventStore: timeSpanTotalMillis :: TimeSpan -> Int64
+ Database.EventStore: updatePersistentSubscription :: Connection -> Text -> Text -> PersistentSubscriptionSettings -> IO (Async ())
+ Database.EventStore.Logging: Connected :: UUID -> InfoMessage
+ Database.EventStore.Logging: Connecting :: Int -> InfoMessage
+ Database.EventStore.Logging: ConnectionClosed :: UUID -> InfoMessage
+ Database.EventStore.Logging: Disconnected :: UUID -> InfoMessage
+ Database.EventStore.Logging: Error :: ErrorMessage -> Log
+ Database.EventStore.Logging: Info :: InfoMessage -> Log
+ Database.EventStore.Logging: MaxAttemptConnectionReached :: Int -> ErrorMessage
+ Database.EventStore.Logging: PackageParsingError :: String -> ErrorMessage
+ Database.EventStore.Logging: WrongPackageFraming :: ErrorMessage
+ Database.EventStore.Logging: data ErrorMessage
+ Database.EventStore.Logging: data InfoMessage
+ Database.EventStore.Logging: data Log
+ Database.EventStore.Logging: instance Show ErrorMessage
+ Database.EventStore.Logging: instance Show InfoMessage
+ Database.EventStore.Logging: instance Show Log
- Database.EventStore: Settings :: NominalDiffTime -> NominalDiffTime -> Bool -> Maybe Credentials -> Retry -> Int -> Settings
+ Database.EventStore: Settings :: NominalDiffTime -> NominalDiffTime -> Bool -> Maybe Credentials -> Retry -> Int -> Maybe (Log -> IO ()) -> Settings
- Database.EventStore: data Subscription
+ Database.EventStore: data Subscription a
- Database.EventStore: hasCaughtUp :: Catchup -> IO Bool
+ Database.EventStore: hasCaughtUp :: Subscription Catchup -> IO Bool
- Database.EventStore: subId :: Subscription -> UUID
+ Database.EventStore: subId :: Identifiable a => Subscription a -> UUID
- Database.EventStore: subLastCommitPos :: Subscription -> Int64
+ Database.EventStore: subLastCommitPos :: Identifiable a => Subscription a -> Int64
- Database.EventStore: subLastEventNumber :: Subscription -> (Maybe Int32)
+ Database.EventStore: subLastEventNumber :: Identifiable a => Subscription a -> Maybe Int32
- Database.EventStore: subResolveLinkTos :: Subscription -> Bool
+ Database.EventStore: subResolveLinkTos :: Subscription Regular -> Bool
- Database.EventStore: subUnsubscribe :: Subscription -> IO ()
+ Database.EventStore: subUnsubscribe :: Subscription a -> (IO ())
- Database.EventStore: subscribe :: Connection -> Text -> Bool -> IO (Async Subscription)
+ Database.EventStore: subscribe :: Connection -> Text -> Bool -> IO (Async (Subscription Regular))
- Database.EventStore: subscribeFrom :: Connection -> Text -> Bool -> Maybe Int32 -> Maybe Int32 -> IO Catchup
+ Database.EventStore: subscribeFrom :: Connection -> Text -> Bool -> Maybe Int32 -> Maybe Int32 -> IO (Subscription Catchup)
- Database.EventStore: subscribeToAll :: Connection -> Bool -> IO (Async Subscription)
+ Database.EventStore: subscribeToAll :: Connection -> Bool -> IO (Async (Subscription Regular))
- Database.EventStore: subscribeToAllFrom :: Connection -> Bool -> Maybe Position -> Maybe Int32 -> IO Catchup
+ Database.EventStore: subscribeToAllFrom :: Connection -> Bool -> Maybe Position -> Maybe Int32 -> IO (Subscription Catchup)
- Database.EventStore: waitTillCatchup :: Catchup -> IO ()
+ Database.EventStore: waitTillCatchup :: Subscription Catchup -> IO ()
Files
- CHANGELOG.markdown +9/−0
- Database/EventStore.hs +104/−39
- Database/EventStore/Catchup.hs +28/−77
- Database/EventStore/Internal/Connection.hs +24/−16
- Database/EventStore/Internal/Manager/Subscription.hs +1112/−378
- Database/EventStore/Internal/Operation/TransactionStartOperation.hs +4/−4
- Database/EventStore/Internal/Processor.hs +40/−33
- Database/EventStore/Internal/Reader.hs +6/−7
- Database/EventStore/Internal/TimeSpan.hs +71/−0
- Database/EventStore/Internal/Types.hs +120/−7
- Database/EventStore/Logging.hs +45/−0
- README.md +20/−17
- eventstore.cabal +3/−3
CHANGELOG.markdown view
@@ -1,3 +1,12 @@+0.8.0.0+-------+* Implement competing consumers.+* Expose an uniform API among all kind of subscriptions.+* Rewrite internal subscription management.+* Add missing `Eq` or `Show` instances for exposed datatypes.+* Add `streamMetadataCustomPropertyValue` and `streamMetadataCustomProperty`.+* Add logging capability.+ 0.7.2.1 ------- * Fix compilation issue
Database/EventStore.hs view
@@ -41,6 +41,8 @@ -- * Write Operations , StreamACL(..) , StreamMetadata(..)+ , streamMetadataGetCustomPropertyValue+ , streamMetadataGetCustomProperty , emptyStreamACL , emptyStreamMetadata , deleteStream@@ -86,6 +88,11 @@ , timeSpanGetMinutes , timeSpanGetSeconds , timeSpanGetMillis+ , timeSpanFromSeconds+ , timeSpanFromMinutes+ , timeSpanFromHours+ , timeSpanFromDays+ , timeSpanTotalMillis -- * Transaction , Transaction , transactionStart@@ -94,26 +101,39 @@ , transactionSendEvents -- * Volatile Subscription , DropReason(..)+ , Identifiable , Subscription+ , NextEvent+ , Regular+ , Catchup+ , Persistent , subscribe , subscribeToAll- , subAwait+ , subNextEvent , subId- , subStream+ , subStreamId+ , subIsSubscribedToAll , subResolveLinkTos , subLastCommitPos , subLastEventNumber , subUnsubscribe -- * Catch-up Subscription- , Catchup , CatchupError(..) , subscribeFrom , subscribeToAllFrom- , catchupAwait- , catchupStream- , catchupUnsubscribe , waitTillCatchup , hasCaughtUp+ -- * Persistent Subscription+ , PersistentSubscriptionSettings(..)+ , SystemConsumerStrategy(..)+ , NakAction(..)+ , notifyEventsProcessed+ , notifyEventsFailed+ , defaultPersistentSubscriptionSettings+ , createPersistentSubscription+ , updatePersistentSubscription+ , deletePersistentSubscription+ , connectToPersistentSubscription -- * Results , AllEventsSlice(..) , DeleteResult(..)@@ -131,6 +151,7 @@ , eventResolved , resolvedEventOriginal , resolvedEventOriginalStreamId+ , resolvedEventOriginalId , positionStart , positionEnd -- * Misc@@ -155,19 +176,20 @@ -------------------------------------------------------------------------------- import Control.Concurrent.Async import Data.Aeson (decode)-import Data.Text+import Data.Text hiding (group) -------------------------------------------------------------------------------- import Database.EventStore.Catchup-import Database.EventStore.Internal.Processor-import Database.EventStore.Internal.TimeSpan-import Database.EventStore.Internal.Types+import Database.EventStore.Internal.Manager.Subscription import Database.EventStore.Internal.Operation.DeleteStreamOperation import Database.EventStore.Internal.Operation.ReadAllEventsOperation import Database.EventStore.Internal.Operation.ReadEventOperation import Database.EventStore.Internal.Operation.ReadStreamEventsOperation import Database.EventStore.Internal.Operation.TransactionStartOperation import Database.EventStore.Internal.Operation.WriteEventsOperation+import Database.EventStore.Internal.Processor+import Database.EventStore.Internal.TimeSpan+import Database.EventStore.Internal.Types -------------------------------------------------------------------------------- -- Connection@@ -175,8 +197,8 @@ -- | Represents a connection to a single EventStore node. data Connection = Connection- { conProcessor :: Processor- , conSettings :: Settings+ { _runCmd :: Cmd -> IO ()+ , _settings :: Settings } --------------------------------------------------------------------------------@@ -198,14 +220,14 @@ -> IO Connection connect settings host port = do processor <- newProcessor settings- processorConnect processor host port+ processor (DoConnect host port) return $ Connection processor settings -------------------------------------------------------------------------------- -- | Asynchronously closes the 'Connection'. shutdown :: Connection -> IO ()-shutdown Connection{..} = processorShutdown conProcessor+shutdown Connection{..} = _runCmd DoShutdown -------------------------------------------------------------------------------- -- | Sends a single 'Event' to given stream.@@ -227,9 +249,9 @@ sendEvents Connection{..} evt_stream exp_ver evts = do (as, mvar) <- createAsync - let op = writeEventsOperation conSettings mvar evt_stream exp_ver evts+ let op = writeEventsOperation _settings mvar evt_stream exp_ver evts - processorNewOperation conProcessor op+ _runCmd (NewOperation op) return as --------------------------------------------------------------------------------@@ -242,9 +264,9 @@ deleteStream Connection{..} evt_stream exp_ver hard_del = do (as, mvar) <- createAsync - let op = deleteStreamOperation conSettings mvar evt_stream exp_ver hard_del+ let op = deleteStreamOperation _settings mvar evt_stream exp_ver hard_del - processorNewOperation conProcessor op+ _runCmd (NewOperation op) return as --------------------------------------------------------------------------------@@ -256,13 +278,13 @@ transactionStart Connection{..} evt_stream exp_ver = do (as, mvar) <- createAsync - let op = transactionStartOperation conSettings- conProcessor+ let op = transactionStartOperation _settings+ _runCmd mvar evt_stream exp_ver - processorNewOperation conProcessor op+ _runCmd (NewOperation op) return as --------------------------------------------------------------------------------@@ -275,9 +297,9 @@ readEvent Connection{..} stream_id evt_num res_link_tos = do (as, mvar) <- createAsync - let op = readEventOperation conSettings mvar stream_id evt_num res_link_tos+ let op = readEventOperation _settings mvar stream_id evt_num res_link_tos - processorNewOperation conProcessor op+ _runCmd (NewOperation op) return as --------------------------------------------------------------------------------@@ -313,7 +335,7 @@ readStreamEventsCommon Connection{..} dir stream_id start cnt res_link_tos = do (as, mvar) <- createAsync - let op = readStreamEventsOperation conSettings+ let op = readStreamEventsOperation _settings dir mvar stream_id@@ -321,7 +343,7 @@ cnt res_link_tos - processorNewOperation conProcessor op+ _runCmd (NewOperation op) return as --------------------------------------------------------------------------------@@ -354,7 +376,7 @@ readAllEventsCommon Connection{..} dir pos max_c res_link_tos = do (as, mvar) <- createAsync - let op = readAllEventsOperation conSettings+ let op = readAllEventsOperation _settings dir mvar c_pos@@ -362,7 +384,7 @@ max_c res_link_tos - processorNewOperation conProcessor op+ _runCmd (NewOperation op) return as where Position c_pos p_pos = pos@@ -372,26 +394,20 @@ subscribe :: Connection -> Text -- ^ Stream name -> Bool -- ^ Resolve Link Tos- -> IO (Async Subscription)+ -> IO (Async (Subscription Regular)) subscribe Connection{..} stream_id res_lnk_tos = do tmp <- newEmptyMVar- processorNewSubcription conProcessor- (putMVar tmp)- stream_id- res_lnk_tos+ _runCmd (NewSub stream_id res_lnk_tos (putMVar tmp)) async $ readMVar tmp -------------------------------------------------------------------------------- -- | Subcribes to $all stream. subscribeToAll :: Connection -> Bool -- ^ Resolve Link Tos- -> IO (Async Subscription)+ -> IO (Async (Subscription Regular)) subscribeToAll Connection{..} res_lnk_tos = do tmp <- newEmptyMVar- processorNewSubcription conProcessor- (putMVar tmp)- ""- res_lnk_tos+ _runCmd (NewSub "" res_lnk_tos (putMVar tmp)) async $ readMVar tmp --------------------------------------------------------------------------------@@ -404,7 +420,7 @@ -> Bool -- ^ Resolve Link Tos -> Maybe Int32 -- ^ Last checkpoint -> Maybe Int32 -- ^ Batch size- -> IO Catchup+ -> IO (Subscription Catchup) subscribeFrom conn stream_id res_lnk_tos last_chk_pt batch_m = do catchupStart evts_fwd get_sub stream_id batch_m last_chk_pt where@@ -419,7 +435,7 @@ -> Bool -- ^ Resolve Link Tos -> Maybe Position -- ^ Last checkpoint -> Maybe Int32 -- ^ Batch size- -> IO Catchup+ -> IO (Subscription Catchup) subscribeToAllFrom conn res_lnk_tos last_chk_pt batch_m = do catchupAllStart evts_fwd get_sub last_chk_pt batch_m where@@ -441,6 +457,7 @@ sendEvent conn (metaStreamOf evt_stream) exp_ver evt --------------------------------------------------------------------------------+-- | Asynchronously gets the metadata of a stream. getStreamMetadata :: Connection -> Text -> IO (Async StreamMetadataResult) getStreamMetadata conn evt_stream = do as <- readEvent conn (metaStreamOf evt_stream) (-1) False@@ -474,6 +491,54 @@ where action = readResultResolvedEvent rres >>= resolvedEventOriginal evt_number = readResultEventNumber rres++--------------------------------------------------------------------------------+-- | Asynchronously create a persistent subscription group on a stream.+createPersistentSubscription :: Connection+ -> Text+ -> Text+ -> PersistentSubscriptionSettings+ -> IO (Async ())+createPersistentSubscription Connection{..} group stream sett = do+ (as, mvar) <- createAsync+ _runCmd (CreatePersist group stream sett (putMVar mvar))+ return as++--------------------------------------------------------------------------------+-- | Asynchronously update a persistent subscription group on a stream.+updatePersistentSubscription :: Connection+ -> Text+ -> Text+ -> PersistentSubscriptionSettings+ -> IO (Async ())+updatePersistentSubscription Connection{..} group stream sett = do+ (as, mvar) <- createAsync+ _runCmd (UpdatePersist group stream sett (putMVar mvar))+ return as++--------------------------------------------------------------------------------+-- | Asynchronously delete a persistent subscription group on a stream.+deletePersistentSubscription :: Connection+ -> Text+ -> Text+ -> IO (Async ())+deletePersistentSubscription Connection{..} group stream = do+ (as, mvar) <- createAsync+ _runCmd (DeletePersist group stream (putMVar mvar))+ return as++--------------------------------------------------------------------------------+-- | Asynchronously connect to a persistent subscription given a group on a+-- stream.+connectToPersistentSubscription :: Connection+ -> Text+ -> Text+ -> Int32+ -> IO (Async (Subscription Persistent))+connectToPersistentSubscription Connection{..} group stream bufSize = do+ mvar <- newEmptyMVar+ _runCmd (ConnectPersist group stream bufSize (putMVar mvar))+ async $ readMVar mvar -------------------------------------------------------------------------------- createAsync :: IO (Async a, MVar (OperationExceptional a))
Database/EventStore/Catchup.hs view
@@ -1,5 +1,6 @@ {-# LANGUAGE DeriveDataTypeable #-} {-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE TypeFamilies #-} -------------------------------------------------------------------------------- -- | -- Module : Database.EventStore.Catchup@@ -11,17 +12,7 @@ -- Portability : non-portable -- ---------------------------------------------------------------------------------module Database.EventStore.Catchup- ( Catchup- , CatchupError(..)- , catchupAwait- , catchupStart- , catchupAllStart- , catchupStream- , catchupUnsubscribe- , waitTillCatchup- , hasCaughtUp- ) where+module Database.EventStore.Catchup where -------------------------------------------------------------------------------- import Control.Concurrent@@ -30,7 +21,6 @@ import Data.Foldable (traverse_) import Data.Int import Data.Maybe-import Data.Typeable -------------------------------------------------------------------------------- import Control.Concurrent.Async@@ -43,45 +33,16 @@ import Database.EventStore.Internal.Types ----------------------------------------------------------------------------------- | Errors that could arise during a catch-up subscription. 'Text' value--- represents the stream name.-data CatchupError- = CatchupStreamDeleted Text- | CatchupUnexpectedStreamStatus Text ReadStreamResult- | CatchupSubscriptionDropReason Text DropReason- deriving (Show, Typeable)-----------------------------------------------------------------------------------instance Exception CatchupError------------------------------------------------------------------------------------- | Representing catch-up subscriptions.-data Catchup- = Catchup- { catchupStream :: Text- -- ^ The name of the stream to which the subscription is subscribed.- , catchupChan :: Chan (Either CatchupError ResolvedEvent)- , catchupSubMVar :: MVar Subscription- , catchupUnsubscribe :: IO ()- -- ^ Asynchronously unsubscribes from the stream.- }------------------------------------------------------------------------------------- | Awaits for the next 'ResolvedEvent'.-catchupAwait :: Catchup -> IO (Either CatchupError ResolvedEvent)-catchupAwait c = readChan $ catchupChan c---------------------------------------------------------------------------------- defaultBatchSize :: Int32 defaultBatchSize = 500 -------------------------------------------------------------------------------- catchupStart :: (Int32 -> Int32 -> IO (Async StreamEventsSlice))- -> IO (Async Subscription)+ -> IO (Async (Subscription Regular)) -> Text -> Maybe Int32 -> Maybe Int32- -> IO Catchup+ -> IO (Subscription Catchup) catchupStart evt_fwd get_sub stream_id batch_size_m last_m = do chan <- newChan var <- newEmptyMVar@@ -101,24 +62,26 @@ putMVar var sub keepAwaitingSubEvent stream_id chan sub - let catchup = Catchup- { catchupStream = stream_id- , catchupChan = chan- , catchupSubMVar = var- , catchupUnsubscribe = do+ let catchup = Catchup var+ sub = Subscription+ { subStreamId = stream_id+ , subNextEvent = readChan chan+ , subIsSubscribedToAll = stream_id == ""+ , subUnsubscribe = do cancel as sub_m <- tryTakeMVar var traverse_ subUnsubscribe sub_m+ , _subInternal = catchup } - return catchup+ return sub -------------------------------------------------------------------------------- catchupAllStart :: ( Position -> Int32 -> IO (Async AllEventsSlice))- -> IO (Async Subscription)+ -> IO (Async (Subscription Regular)) -> Maybe Position -> Maybe Int32- -> IO Catchup+ -> IO (Subscription Catchup) catchupAllStart evt_fwd get_sub last_chk_pt_m batch_size_m = do chan <- newChan var <- newEmptyMVar@@ -137,17 +100,19 @@ putMVar var sub keepAwaitingSubEvent "" chan sub - let catchup = Catchup- { catchupStream = ""- , catchupChan = chan- , catchupSubMVar = var- , catchupUnsubscribe = do- cancel as- sub_m <- tryTakeMVar var- traverse_ subUnsubscribe sub_m- }+ let catchup = Catchup var+ sub = Subscription+ { subStreamId = ""+ , subNextEvent = readChan chan+ , subIsSubscribedToAll = True+ , subUnsubscribe = do+ cancel as+ sub_m <- tryTakeMVar var+ traverse_ subUnsubscribe sub_m+ , _subInternal = catchup+ } - return catchup+ return sub -------------------------------------------------------------------------------- readEventsTill :: (Int32 -> Int32 -> IO (Async StreamEventsSlice))@@ -208,10 +173,10 @@ -------------------------------------------------------------------------------- keepAwaitingSubEvent :: Text -> Chan (Either CatchupError ResolvedEvent)- -> Subscription+ -> Subscription Regular -> IO () keepAwaitingSubEvent stream_id chan sub = forever $ do- evt_e <- subAwait sub+ evt_e <- subNextEvent sub case evt_e of Right evt -> writeChan chan (Right evt) Left r -> do@@ -219,17 +184,3 @@ writeChan chan (Left e) throwIO e------------------------------------------------------------------------------------- | Waits until 'Catchup' subscription catch-up its stream.-waitTillCatchup :: Catchup -> IO ()-waitTillCatchup c = do- _ <- readMVar $ catchupSubMVar c- return ()------------------------------------------------------------------------------------- | Non blocking version of `waitTillCatchup`.-hasCaughtUp :: Catchup -> IO Bool-hasCaughtUp c = do- res <- tryReadMVar $ catchupSubMVar c- return $ isJust res
Database/EventStore/Internal/Connection.hs view
@@ -28,7 +28,6 @@ import qualified Data.ByteString as B import Data.Typeable import System.IO-import Text.Printf -------------------------------------------------------------------------------- import Data.UUID@@ -37,6 +36,7 @@ -------------------------------------------------------------------------------- import Database.EventStore.Internal.Types+import Database.EventStore.Logging -------------------------------------------------------------------------------- data ConnectionException =@@ -62,15 +62,20 @@ case s_retry sett of AtMost n -> let loop i = do- printf "Connecting...Attempt %d\n" i+ _settingsLog sett (Info $ Connecting i) catch (connect sett host port) $ \(_ :: SomeException) -> do threadDelay delay- if n <= i then throwIO $ MaxAttempt host port n- else loop (i + 1) in+ if n <= i+ then do+ _settingsLog sett+ $ Error+ $ MaxAttemptConnectionReached i+ throwIO $ MaxAttempt host port n+ else loop (i + 1) in loop 1 KeepRetrying -> let endlessly i = do- printf "Connecting...Attempt %d\n" i+ _settingsLog sett (Info $ Connecting i) catch (connect sett host port) $ \(_ :: SomeException) -> do threadDelay delay >> endlessly (i + 1) in endlessly (1 :: Int)@@ -83,19 +88,22 @@ -------------------------------------------------------------------------------- connect :: Settings -> HostName -> Int -> IO Connection-connect _ host port = do+connect sett host port = do hdl <- connectTo host (PortNumber $ fromIntegral port) hSetBuffering hdl NoBuffering uuid <- randomIO- return $ regularConnection hdl uuid+ regularConnection sett hdl uuid ---------------------------------------------------------------------------------regularConnection :: Handle -> UUID -> Connection-regularConnection h uuid =- Connection- { connUUID = uuid- , connClose = hClose h- , connFlush = hFlush h- , connSend = B.hPut h- , connRecv = B.hGet h- }+regularConnection :: Settings -> Handle -> UUID -> IO Connection+regularConnection sett h uuid = do+ _settingsLog sett (Info $ Connected uuid)+ return Connection+ { connUUID = uuid+ , connClose = do+ _settingsLog sett (Info $ ConnectionClosed uuid)+ hClose h+ , connFlush = hFlush h+ , connSend = B.hPut h+ , connRecv = B.hGet h+ }
Database/EventStore/Internal/Manager/Subscription.hs view
@@ -1,378 +1,1112 @@-{-# LANGUAGE DeriveGeneric #-}-{-# LANGUAGE DataKinds #-}-{-# LANGUAGE MultiWayIf #-}-{-# LANGUAGE RecordWildCards #-}-{-# OPTIONS_GHC -fcontext-stack=26 #-}------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Manager.Subscription--- Copyright : (C) 2014 Yorick Laupa--- License : (see the file LICENSE)------ Maintainer : Yorick Laupa <yo.eight@gmail.com>--- Stability : provisional--- Portability : non-portable-------------------------------------------------------------------------------------module Database.EventStore.Internal.Manager.Subscription- ( DropReason (..)- , NewSubscriptionCB- , Subscription- , subAwait- , subId- , subStream- , subResolveLinkTos- , subLastCommitPos- , subLastEventNumber- , subUnsubscribe- , subscriptionNetwork- ) where-----------------------------------------------------------------------------------import Control.Concurrent-import Control.Monad.Fix-import Data.ByteString (ByteString)-import Data.Functor-import Data.Int-import qualified Data.Map.Strict as M-import Data.Maybe-import Data.Monoid ((<>))-import Data.Traversable (for)-import Data.Word-import GHC.Generics (Generic)-----------------------------------------------------------------------------------import Data.ProtocolBuffers-import Data.Serialize-import Data.Text-import Data.UUID-import FRP.Sodium-import FRP.Sodium.IO-import System.Random-----------------------------------------------------------------------------------import Database.EventStore.Internal.Types hiding (Event, newEvent)-import Database.EventStore.Internal.Util.Sodium-----------------------------------------------------------------------------------data SubscribeToStream- = SubscribeToStream- { subscribeStreamId :: Required 1 (Value Text)- , subscribeResolveLinkTos :: Required 2 (Value Bool)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Encode SubscribeToStream-----------------------------------------------------------------------------------subscribeToStream :: Text -> Bool -> SubscribeToStream-subscribeToStream stream_id res_link_tos =- SubscribeToStream- { subscribeStreamId = putField stream_id- , subscribeResolveLinkTos = putField res_link_tos- }-----------------------------------------------------------------------------------data SubscriptionConfirmation- = SubscriptionConfirmation- { subscribeLastCommitPos :: Required 1 (Value Int64)- , subscribeLastEventNumber :: Optional 2 (Value Int32)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode SubscriptionConfirmation-----------------------------------------------------------------------------------data StreamEventAppeared- = StreamEventAppeared- { streamResolvedEvent :: Required 1 (Message ResolvedEventBuf) }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode StreamEventAppeared------------------------------------------------------------------------------------- | Represents the reason subscription drop happened.-data DropReason- = D_Unsubscribed- | D_AccessDenied- | D_NotFound- | D_PersistentSubscriptionDeleted- deriving (Enum, Eq, Show)-----------------------------------------------------------------------------------data SubscriptionDropped- = SubscriptionDropped- { dropReason :: Optional 1 (Enumeration DropReason) }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode SubscriptionDropped-----------------------------------------------------------------------------------data UnsubscribeFromStream = UnsubscribeFromStream deriving (Generic, Show)-----------------------------------------------------------------------------------instance Encode UnsubscribeFromStream-----------------------------------------------------------------------------------data Pending- = Pending- { penConf :: Int64 -> Maybe Int32 -> IO Subscription- , penCb :: Subscription -> IO ()- }------------------------------------------------------------------------------------- | Represents a subscription to a single stream or $all stream--- in the EventStore.-data Subscription- = Subscription- { subId :: !UUID- -- ^ ID of the subscription.- , subStream :: !Text- -- ^ The name of the stream to which the subscription is subscribed.- , subResolveLinkTos :: !Bool- -- ^ Determines whether or not any link events encontered in the stream- -- will be resolved.- , subLastCommitPos :: !Int64- -- ^ The last commit position seen on the subscription (if this a- -- subscription to $all stream).- , subLastEventNumber :: !(Maybe Int32)- -- ^ The last event number seen on the subscription (if this is a- -- subscription to a single stream).- , subChan :: Chan (Either DropReason ResolvedEvent)- , subUnsubscribe :: IO ()- -- ^ Asynchronously unsubscribe from the the stream.- }------------------------------------------------------------------------------------- | Awaits for the next 'ResolvedEvent'.-subAwait :: Subscription -> IO (Either DropReason ResolvedEvent)-subAwait Subscription{..} = readChan subChan-----------------------------------------------------------------------------------data Manager- = Manager- { _pendings :: !(M.Map UUID Pending)- , _subscriptions :: !(M.Map UUID Subscription)- }-----------------------------------------------------------------------------------initManager :: Manager-initManager = Manager M.empty M.empty------------------------------------------------------------------------------------- Handled Packages----------------------------------------------------------------------------------subscriptionConfirmed :: Word8-subscriptionConfirmed = 0xC1-----------------------------------------------------------------------------------streamEventAppeared :: Word8-streamEventAppeared = 0xC2-----------------------------------------------------------------------------------subscriptionDropped :: Word8-subscriptionDropped = 0xC4--------------------------------------------------------------------------------------------------------------------------------------------------------------------maybeDecodeMessage :: Decode a => ByteString -> Maybe a-maybeDecodeMessage bytes =- case runGet decodeMessage bytes of- Right a -> Just a- _ -> Nothing-----------------------------------------------------------------------------------onConfirmation :: Package -> Maybe Confirmation-onConfirmation Package{..}- | packageCmd == subscriptionConfirmed =- fmap (Confirmation packageCorrelation) $- maybeDecodeMessage packageData- | otherwise = Nothing-----------------------------------------------------------------------------------data Appeared- = Appeared- { _appSub :: !Subscription- , _appEvt :: !ResolvedEvent- }-----------------------------------------------------------------------------------onEventAppeared :: Package -> Manager -> Maybe Appeared-onEventAppeared Package{..} Manager{..}- | packageCmd == streamEventAppeared = do- sub <- M.lookup packageCorrelation _subscriptions- sea <- maybeDecodeMessage packageData- let res_evt = getField $ streamResolvedEvent sea-- return Appeared- { _appSub = sub- , _appEvt = newResolvedEventFromBuf res_evt- }- | otherwise = Nothing-----------------------------------------------------------------------------------confirmSub :: Confirmation -> Manager -> IO (Maybe Subscription)-confirmSub (Confirmation uuid sc) Manager{..} =- for (M.lookup uuid _pendings) $ \p -> do- let last_com_pos = getField $ subscribeLastCommitPos sc- last_evt_num = getField $ subscribeLastEventNumber sc-- sub <- penConf p last_com_pos last_evt_num- penCb p sub- return sub------------------------------------------------------------------------------------- Events----------------------------------------------------------------------------------data Subscribe- = Subscribe- { _subId :: !UUID- , _subCallback :: Subscription -> IO ()- , _subStream :: !Text- , _subResolveLinkTos :: !Bool- }-----------------------------------------------------------------------------------data Confirmation = Confirmation !UUID !SubscriptionConfirmation-----------------------------------------------------------------------------------type NewSubscriptionCB = (Subscription -> IO ()) -> Text -> Bool -> IO ()-----------------------------------------------------------------------------------subscriptionNetwork :: Settings- -> (Package -> Reactive ())- -> Event Package- -> Reactive NewSubscriptionCB-subscriptionNetwork sett push_pkg e_pkg = do- (on_sub, push_sub) <- newEvent- (on_rem, push_rem) <- newEvent-- mgr_b <- mfix $ \mgr_b -> do- let on_con = filterJust $ fmap onConfirmation e_pkg- on_con_sub = filterJust $ executeSyncIO $ snapshot confirmSub- on_con- mgr_b- send_unsub = push_pkg . createUnsubscribePackage sett-- mgr_e = fmap (subscribe send_unsub) on_sub <>- fmap remove on_rem <>- fmap confirmed on_con_sub-- accum initManager mgr_e-- let on_app = filterJust $ snapshot onEventAppeared e_pkg mgr_b- on_drop = filterJust $ execute $ snapshot (dropError push_rem)- e_pkg- mgr_b- push_pkg_io = pushAsync push_pkg-- push_sub_io cb stream res_lnk_tos = randomIO >>= \uuid -> void $- forkIO $ sync $ push_sub Subscribe- { _subId = uuid- , _subCallback = cb- , _subStream = stream- , _subResolveLinkTos = res_lnk_tos- }-- _ <- listen on_sub (push_pkg_io . createSubscribePackage sett)-- _ <- listen on_app $ \(Appeared sub evt) ->- writeChan (subChan sub) (Right evt)-- _ <- listen on_drop $ \(Dropped reason sub _) ->- writeChan (subChan sub) (Left reason)-- return push_sub_io-----------------------------------------------------------------------------------createSubscribePackage :: Settings -> Subscribe -> Package-createSubscribePackage Settings{..} Subscribe{..} =- Package- { packageCmd = 0xC0- , packageCorrelation = _subId- , packageData = runPut $ encodeMessage msg- , packageCred = s_credentials- }- where- msg = subscribeToStream _subStream _subResolveLinkTos-----------------------------------------------------------------------------------createUnsubscribePackage :: Settings -> UUID -> Package-createUnsubscribePackage Settings{..} uuid =- Package- { packageCmd = 0xC3- , packageCorrelation = uuid- , packageData = runPut $ encodeMessage UnsubscribeFromStream- , packageCred = s_credentials- }-----------------------------------------------------------------------------------data Dropped- = Dropped- { droppedReason :: !DropReason- , droppedSub :: !Subscription- , droppedId :: !UUID- }-----------------------------------------------------------------------------------dropError :: (UUID -> Reactive ())- -> Package- -> Manager- -> Reactive (Maybe Dropped)-dropError push_rem Package{..} Manager{..} =- for go $ \d -> do- push_rem $ droppedId d- return d- where- go | packageCmd == subscriptionDropped = do- sub <- M.lookup packageCorrelation _subscriptions- msg <- maybeDecodeMessage packageData- let reason = fromMaybe D_Unsubscribed $ getField $ dropReason msg-- return Dropped- { droppedReason = reason- , droppedSub = sub- , droppedId = packageCorrelation- }- | otherwise = Nothing------------------------------------------------------------------------------------- Model----------------------------------------------------------------------------------subscribe :: (UUID -> Reactive ()) -> Subscribe -> Manager -> Manager-subscribe unsub Subscribe{..} s@Manager{..} =- s { _pendings = M.insert _subId pending _pendings }- where- pending =- Pending- { penConf = new_sub- , penCb = _subCallback- }-- new_sub com_pos last_evt = do- chan <- newChan-- return Subscription- { subId = _subId- , subStream = _subStream- , subResolveLinkTos = _subResolveLinkTos- , subLastCommitPos = com_pos- , subLastEventNumber = last_evt- , subChan = chan- , subUnsubscribe = sync $ unsub _subId- }-----------------------------------------------------------------------------------confirmed :: Subscription -> Manager -> Manager-confirmed sub@Subscription{..} s@Manager{..} =- s { _pendings = M.delete subId _pendings- , _subscriptions = M.insert subId sub _subscriptions- }-----------------------------------------------------------------------------------remove :: UUID -> Manager -> Manager-remove uuid s@Manager{..} = s { _subscriptions = M.delete uuid _subscriptions }+{-# LANGUAGE BangPatterns #-}+{-# LANGUAGE DeriveGeneric #-}+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE ExistentialQuantification #-}+{-# LANGUAGE DeriveDataTypeable #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE MultiWayIf #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE TypeFamilies #-}+{-# OPTIONS_GHC -fcontext-stack=26 #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Manager.Subscription+-- Copyright : (C) 2014 Yorick Laupa+-- License : (see the file LICENSE)+--+-- Maintainer : Yorick Laupa <yo.eight@gmail.com>+-- Stability : provisional+-- Portability : non-portable+--+--------------------------------------------------------------------------------+module Database.EventStore.Internal.Manager.Subscription where++--------------------------------------------------------------------------------+import Control.Concurrent+import Control.Exception+import Control.Monad.Fix+import Data.ByteString (ByteString)+import Data.ByteString.Lazy (toStrict)+import Data.Foldable+import Data.Functor+import Data.Int+import qualified Data.Map.Strict as M+import Data.Maybe+import Data.Monoid ((<>))+import Data.Typeable+import GHC.Generics (Generic)+import Prelude++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Serialize+import Data.Text hiding (group)+import Data.UUID+import FRP.Sodium+import System.Random++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation.ReadStreamEventsOperation+import Database.EventStore.Internal.TimeSpan+import Database.EventStore.Internal.Types hiding (Event, newEvent)+import Database.EventStore.Internal.Util.Sodium++--------------------------------------------------------------------------------+data SubscribeToStream+ = SubscribeToStream+ { subscribeStreamId :: Required 1 (Value Text)+ , subscribeResolveLinkTos :: Required 2 (Value Bool)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode SubscribeToStream++--------------------------------------------------------------------------------+subscribeToStream :: Text -> Bool -> SubscribeToStream+subscribeToStream stream_id res_link_tos =+ SubscribeToStream+ { subscribeStreamId = putField stream_id+ , subscribeResolveLinkTos = putField res_link_tos+ }++--------------------------------------------------------------------------------+data SubscriptionConfirmation+ = SubscriptionConfirmation+ { subscribeLastCommitPos :: Required 1 (Value Int64)+ , subscribeLastEventNumber :: Optional 2 (Value Int32)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode SubscriptionConfirmation++--------------------------------------------------------------------------------+data StreamEventAppeared+ = StreamEventAppeared+ { streamResolvedEvent :: Required 1 (Message ResolvedEventBuf) }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode StreamEventAppeared++--------------------------------------------------------------------------------+-- | Represents the reason subscription drop happened.+data DropReason+ = D_Unsubscribed+ | D_AccessDenied+ | D_NotFound+ | D_PersistentSubscriptionDeleted+ deriving (Enum, Eq, Show)++--------------------------------------------------------------------------------+data SubscriptionDropped+ = SubscriptionDropped+ { dropReason :: Optional 1 (Enumeration DropReason) }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode SubscriptionDropped++--------------------------------------------------------------------------------+data UnsubscribeFromStream = UnsubscribeFromStream deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode UnsubscribeFromStream++--------------------------------------------------------------------------------+data CreatePersistentSubscription =+ CreatePersistentSubscription+ { cpsGroupName :: Required 1 (Value Text)+ , cpsStreamId :: Required 2 (Value Text)+ , cpsResolveLinkTos :: Required 3 (Value Bool)+ , cpsStartFrom :: Required 4 (Value Int32)+ , cpsMsgTimeout :: Required 5 (Value Int32)+ , cpsRecordStats :: Required 6 (Value Bool)+ , cpsLiveBufSize :: Required 7 (Value Int32)+ , cpsReadBatchSize :: Required 8 (Value Int32)+ , cpsBufSize :: Required 9 (Value Int32)+ , cpsMaxRetryCount :: Required 10 (Value Int32)+ , cpsPreferRoundRobin :: Required 11 (Value Bool)+ , cpsChkPtAfterTime :: Required 12 (Value Int32)+ , cpsChkPtMaxCount :: Required 13 (Value Int32)+ , cpsChkPtMinCount :: Required 14 (Value Int32)+ , cpsSubMaxCount :: Required 15 (Value Int32)+ , cpsNamedConsStrategy :: Optional 16 (Value Text)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+_createPersistentSubscription :: Text+ -> Text+ -> PersistentSubscriptionSettings+ -> CreatePersistentSubscription+_createPersistentSubscription group stream sett =+ CreatePersistentSubscription+ { cpsGroupName = putField group+ , cpsStreamId = putField stream+ , cpsResolveLinkTos = putField $ psSettingsResolveLinkTos sett+ , cpsStartFrom = putField $ psSettingsStartFrom sett+ , cpsMsgTimeout = putField $ ms $ psSettingsMsgTimeout sett+ , cpsRecordStats = putField $ psSettingsExtraStats sett+ , cpsLiveBufSize = putField $ psSettingsLiveBufSize sett+ , cpsReadBatchSize = putField $ psSettingsReadBatchSize sett+ , cpsBufSize = putField $ psSettingsHistoryBufSize sett+ , cpsMaxRetryCount = putField $ psSettingsMaxRetryCount sett+ , cpsPreferRoundRobin = putField False+ , cpsChkPtAfterTime = putField $ ms $ psSettingsCheckPointAfter sett+ , cpsChkPtMaxCount = putField $ psSettingsMaxCheckPointCount sett+ , cpsChkPtMinCount = putField $ psSettingsMinCheckPointCount sett+ , cpsSubMaxCount = putField $ psSettingsMaxSubsCount sett+ , cpsNamedConsStrategy = putField $ Just strText+ }+ where+ strText = strategyText $ psSettingsNamedConsumerStrategy sett+ ms = fromIntegral . timeSpanTotalMillis++--------------------------------------------------------------------------------+instance Encode CreatePersistentSubscription++--------------------------------------------------------------------------------+data CreatePersistentSubscriptionResult+ = CPS_Success+ | CPS_AlreadyExists+ | CPS_Fail+ | CPS_AccessDenied+ deriving (Enum, Eq, Show)++--------------------------------------------------------------------------------+data CreatePersistentSubscriptionCompleted =+ CreatePersistentSubscriptionCompleted+ { cpscResult :: Required 1 (Enumeration CreatePersistentSubscriptionResult)+ , cpscReason :: Optional 2 (Value Text)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode CreatePersistentSubscriptionCompleted++--------------------------------------------------------------------------------+data DeletePersistentSubscription =+ DeletePersistentSubscription+ { dpsGroupName :: Required 1 (Value Text)+ , dpsStreamId :: Required 2 (Value Text)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode DeletePersistentSubscription++--------------------------------------------------------------------------------+_deletePersistentSubscription :: Text -> Text -> DeletePersistentSubscription+_deletePersistentSubscription group_name stream_id =+ DeletePersistentSubscription+ { dpsGroupName = putField group_name+ , dpsStreamId = putField stream_id+ }++--------------------------------------------------------------------------------+data DeletePersistentSubscriptionResult+ = DPS_Success+ | DPS_DoesNotExist+ | DPS_Fail+ | DPS_AccessDenied+ deriving (Enum, Eq, Show)++--------------------------------------------------------------------------------+data DeletePersistentSubscriptionCompleted =+ DeletePersistentSubscriptionCompleted+ { dpscResult :: Required 1 (Enumeration DeletePersistentSubscriptionResult)+ , dpscReason :: Optional 2 (Value Text)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode DeletePersistentSubscriptionCompleted++--------------------------------------------------------------------------------+data UpdatePersistentSubscription =+ UpdatePersistentSubscription+ { upsGroupName :: Required 1 (Value Text)+ , upsStreamId :: Required 2 (Value Text)+ , upsResolveLinkTos :: Required 3 (Value Bool)+ , upsStartFrom :: Required 4 (Value Int32)+ , upsMsgTimeout :: Required 5 (Value Int32)+ , upsRecordStats :: Required 6 (Value Bool)+ , upsLiveBufSize :: Required 7 (Value Int32)+ , upsReadBatchSize :: Required 8 (Value Int32)+ , upsBufSize :: Required 9 (Value Int32)+ , upsMaxRetryCount :: Required 10 (Value Int32)+ , upsPreferRoundRobin :: Required 11 (Value Bool)+ , upsChkPtAfterTime :: Required 12 (Value Int32)+ , upsChkPtMaxCount :: Required 13 (Value Int32)+ , upsChkPtMinCount :: Required 14 (Value Int32)+ , upsSubMaxCount :: Required 15 (Value Int32)+ , upsNamedConsStrategy :: Optional 16 (Value Text)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+_updatePersistentSubscription :: Text+ -> Text+ -> PersistentSubscriptionSettings+ -> UpdatePersistentSubscription+_updatePersistentSubscription group stream sett =+ UpdatePersistentSubscription+ { upsGroupName = putField group+ , upsStreamId = putField stream+ , upsResolveLinkTos = putField $ psSettingsResolveLinkTos sett+ , upsStartFrom = putField $ psSettingsStartFrom sett+ , upsMsgTimeout = putField $ ms $ psSettingsMsgTimeout sett+ , upsRecordStats = putField $ psSettingsExtraStats sett+ , upsLiveBufSize = putField $ psSettingsLiveBufSize sett+ , upsReadBatchSize = putField $ psSettingsReadBatchSize sett+ , upsBufSize = putField $ psSettingsHistoryBufSize sett+ , upsMaxRetryCount = putField $ psSettingsMaxRetryCount sett+ , upsPreferRoundRobin = putField False+ , upsChkPtAfterTime = putField $ ms $ psSettingsCheckPointAfter sett+ , upsChkPtMaxCount = putField $ psSettingsMaxCheckPointCount sett+ , upsChkPtMinCount = putField $ psSettingsMinCheckPointCount sett+ , upsSubMaxCount = putField $ psSettingsMaxSubsCount sett+ , upsNamedConsStrategy = putField $ Just strText+ }+ where+ strText = strategyText $ psSettingsNamedConsumerStrategy sett+ ms = fromIntegral . timeSpanTotalMillis++--------------------------------------------------------------------------------+instance Encode UpdatePersistentSubscription++--------------------------------------------------------------------------------+data UpdatePersistentSubscriptionResult+ = UPS_Success+ | UPS_DoesNotExist+ | UPS_Fail+ | UPS_AccessDenied+ deriving (Enum, Eq, Show)++--------------------------------------------------------------------------------+data UpdatePersistentSubscriptionCompleted =+ UpdatePersistentSubscriptionCompleted+ { upscResult :: Required 1 (Enumeration UpdatePersistentSubscriptionResult)+ , upscReason :: Optional 2 (Value Text)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode UpdatePersistentSubscriptionCompleted++--------------------------------------------------------------------------------+data ConnectToPersistentSubscription =+ ConnectToPersistentSubscription+ { ctsId :: Required 1 (Value Text)+ , ctsStreamId :: Required 2 (Value Text)+ , ctsAllowedInFlightMsgs :: Required 3 (Value Int32)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode ConnectToPersistentSubscription++--------------------------------------------------------------------------------+_connectToPersistentSubscription :: Text+ -> Text+ -> Int32+ -> ConnectToPersistentSubscription+_connectToPersistentSubscription sub_id stream_id all_fly_msgs =+ ConnectToPersistentSubscription+ { ctsId = putField sub_id+ , ctsStreamId = putField stream_id+ , ctsAllowedInFlightMsgs = putField all_fly_msgs+ }++--------------------------------------------------------------------------------+data PersistentSubscriptionAckEvents =+ PersistentSubscriptionAckEvents+ { psaeId :: Required 1 (Value Text)+ , psaeProcessedEvtIds :: Required 2 (Value ByteString)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode PersistentSubscriptionAckEvents++--------------------------------------------------------------------------------+persistentSubscriptionAckEvents :: Text+ -> ByteString+ -> PersistentSubscriptionAckEvents+persistentSubscriptionAckEvents sub_id evt_ids =+ PersistentSubscriptionAckEvents+ { psaeId = putField sub_id+ , psaeProcessedEvtIds = putField evt_ids+ }++--------------------------------------------------------------------------------+data NakAction+ = NA_Unknown+ | NA_Park+ | NA_Retry+ | NA_Skip+ | NA_Stop+ deriving (Enum, Eq, Show)++--------------------------------------------------------------------------------+data PersistentSubscriptionNakEvents =+ PersistentSubscriptionNakEvents+ { psneId :: Required 1 (Value Text)+ , psneProcessedEvtIds :: Required 2 (Value ByteString)+ , psneMsg :: Optional 3 (Value Text)+ , psneAction :: Required 4 (Enumeration NakAction)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode PersistentSubscriptionNakEvents++--------------------------------------------------------------------------------+persistentSubscriptionNakEvents :: Text+ -> ByteString+ -> Maybe Text+ -> NakAction+ -> PersistentSubscriptionNakEvents+persistentSubscriptionNakEvents sub_id evt_ids msg action =+ PersistentSubscriptionNakEvents+ { psneId = putField sub_id+ , psneProcessedEvtIds = putField evt_ids+ , psneMsg = putField msg+ , psneAction = putField action+ }++--------------------------------------------------------------------------------+data PersistentSubscriptionConfirmation =+ PersistentSubscriptionConfirmation+ { pscLastCommitPos :: Required 1 (Value Int64)+ , pscId :: Required 2 (Value Text)+ , pscLastEvtNumber :: Optional 3 (Value Int32)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode PersistentSubscriptionConfirmation++--------------------------------------------------------------------------------+data PersistentSubscriptionStreamEventAppeared =+ PersistentSubscriptionStreamEventAppeared+ { psseaEvt :: Required 1 (Message ResolvedIndexedEvent) }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode PersistentSubscriptionStreamEventAppeared++--------------------------------------------------------------------------------+data Sub a where+ RegularSub :: Text -> Bool -> Sub Regular+ PersistentSub :: Text -> Text -> Int32 -> Sub Persistent++--------------------------------------------------------------------------------+data Pending =+ forall s. Push s =>+ Pending+ { _penId :: !UUID+ , _penSub :: !(Sub s)+ , _penCb :: Subscription s -> IO ()+ }++--------------------------------------------------------------------------------+data Confirmed =+ forall s. Push s =>+ Confirmed+ { _conId :: !UUID+ , _conTyp :: !(Sub s)+ , _conCB :: Subscription s -> IO ()+ , _conSub :: !(IO (Subscription s))+ }++--------------------------------------------------------------------------------+data OnGoing =+ forall s. Push s =>+ OnGoing+ { _ongTyp :: !(Sub s)+ , _ongSub :: !(Subscription s)+ }++--------------------------------------------------------------------------------+-- | Value's type returned when calling 'subNextEvent'+type family NextEvent a :: * where+ NextEvent Regular = Either DropReason ResolvedEvent+ NextEvent Persistent = Either DropReason ResolvedEvent+ NextEvent Catchup = Either CatchupError ResolvedEvent++--------------------------------------------------------------------------------+-- | Represents a subscription to a stream.+data Subscription a =+ Subscription+ { subStreamId :: !Text+ -- ^ The name of the stream to which the subscription is subscribed.+ , subUnsubscribe :: !(IO ())+ -- ^ Asynchronously unsubscribe from the the stream.+ , subNextEvent :: !(IO (NextEvent a))+ -- ^ Awaits for the next event.+ , subIsSubscribedToAll :: !Bool+ -- ^ True if this subscription is to $all stream.+ , _subInternal :: !a+ }++--------------------------------------------------------------------------------+-- | Internal use only because we all know that lawless type-classes are bad.+-- But Haskell clearly lacks of a proper module system. So meanwhile, we use+-- type-class as a (hacky) way to have a bit of modularity.+class Push a where+ _pushEvt :: a -> Either DropReason ResolvedEvent -> IO ()++--------------------------------------------------------------------------------+_subPushEvt :: Push a+ => Subscription a+ -> Either DropReason ResolvedEvent+ -> IO ()+_subPushEvt = _pushEvt . _subInternal++--------------------------------------------------------------------------------+-- | Represents a subscription that is directly identifiable. 'Regular' and+-- 'Persistent' fit that description while 'Catchup' doesn't. Because+-- 'Catchup' reads all events from a particular checkpoint and when it's+-- finished, it issues a subscription request.+class Identifiable a where+ _getId :: a -> UUID+ _getLastCommitPos :: a -> Int64+ _getLastEventNumber :: a -> Maybe Int32++--------------------------------------------------------------------------------+-- | Gets the ID of the subscription.+subId :: Identifiable a => Subscription a -> UUID+subId = _getId . _subInternal++--------------------------------------------------------------------------------+-- | The last commit position seen on the subscription (if this a subscription+-- to $all stream).+subLastCommitPos :: Identifiable a => Subscription a -> Int64+subLastCommitPos = _getLastCommitPos . _subInternal++--------------------------------------------------------------------------------+-- | The last event number seen on the subscription (if this is a subscription+-- to a single stream).+subLastEventNumber :: Identifiable a => Subscription a -> Maybe Int32+subLastEventNumber = _getLastEventNumber . _subInternal++--------------------------------------------------------------------------------+-- | Represents a subscription to a single stream or $all stream in the+-- EventStore.+data Regular =+ Regular+ { _regId :: !UUID+ , _regResolveLinkTos :: !Bool+ , _regLastCommitPos :: !Int64+ , _regLastEventNumber :: !(Maybe Int32)+ , _regChan :: !(Chan (Either DropReason ResolvedEvent))+ }++--------------------------------------------------------------------------------+instance Identifiable Regular where+ _getId = _regId+ _getLastCommitPos = _regLastCommitPos+ _getLastEventNumber = _regLastEventNumber++--------------------------------------------------------------------------------+instance Push Regular where+ _pushEvt reg = writeChan (_regChan reg)++--------------------------------------------------------------------------------+-- | Determines whether or not any link events encontered in the stream will be+-- resolved.+subResolveLinkTos :: Subscription Regular -> Bool+subResolveLinkTos Subscription { _subInternal = reg } = _regResolveLinkTos reg++--------------------------------------------------------------------------------+-- | Errors that could arise during a catch-up subscription. 'Text' value+-- represents the stream name.+data CatchupError+ = CatchupStreamDeleted Text+ | CatchupUnexpectedStreamStatus Text ReadStreamResult+ | CatchupSubscriptionDropReason Text DropReason+ deriving (Show, Typeable)++--------------------------------------------------------------------------------+instance Exception CatchupError++--------------------------------------------------------------------------------+-- | Represents catch-up subscription.+data Catchup = Catchup { _catchupSub :: MVar (Subscription Regular) }++--------------------------------------------------------------------------------+-- | Waits until 'Catchup' subscription catch-up its stream.+waitTillCatchup :: Subscription Catchup -> IO ()+waitTillCatchup Subscription { _subInternal = Catchup mvar } = do+ _ <- readMVar mvar+ return ()++--------------------------------------------------------------------------------+-- | Non blocking version of `waitTillCatchup`.+hasCaughtUp :: Subscription Catchup -> IO Bool+hasCaughtUp Subscription { _subInternal = Catchup mvar } =+ fmap isJust $ tryReadMVar mvar++--------------------------------------------------------------------------------+-- | Represents a persistent subscription.+data Persistent =+ Persistent+ { _persistId :: !UUID+ , _persistChan :: !(Chan (Either DropReason ResolvedEvent))+ , _persistSubId :: !Text+ , _persistGroup :: !Text+ , _persistLastCPos :: !Int64+ , _persistLastENum :: !(Maybe Int32)+ , _persistAckCmd :: AckCmd -> IO ()+ }++--------------------------------------------------------------------------------+instance Identifiable Persistent where+ _getId = _persistId+ _getLastCommitPos = _persistLastCPos+ _getLastEventNumber = _persistLastENum++--------------------------------------------------------------------------------+instance Push Persistent where+ _pushEvt p = writeChan (_persistChan p)++--------------------------------------------------------------------------------+-- | Acknowledges those event ids have been successfully processed.+notifyEventsProcessed :: Subscription Persistent -> [UUID] -> IO ()+notifyEventsProcessed sub eids = _persistAckCmd p (AckCmd eids)+ where+ p = _subInternal sub++--------------------------------------------------------------------------------+-- | Acknowledges those event ids have failed to be processed successfully.+notifyEventsFailed :: Subscription Persistent+ -> NakAction+ -> Maybe Text+ -> [UUID]+ -> IO ()+notifyEventsFailed sub act msg eids = _persistAckCmd p (NakCmd act msg eids)+ where+ p = _subInternal sub++--------------------------------------------------------------------------------+data PersistAction+ = PersistCreate PersistentSubscriptionSettings+ | PersistUpdate PersistentSubscriptionSettings+ | PersistDelete++--------------------------------------------------------------------------------+data AckCmd+ = AckCmd [UUID]+ | NakCmd NakAction (Maybe Text) [UUID]++--------------------------------------------------------------------------------+data PendingPersistAction =+ PendingPersistAction+ { _ppaId :: !UUID+ , _ppaGroup :: !Text+ , _ppaStream :: !Text+ , _ppaTyp :: !PersistAction+ , _ppaCB :: Either OperationException () -> IO ()+ }++--------------------------------------------------------------------------------+data PersistActionConfirmed =+ PersistActionConfirmed+ { _pacId :: !UUID+ , _pacResult :: !(Either OperationException ())+ , _pacCB :: Either OperationException () -> IO ()+ }++--------------------------------------------------------------------------------+data Manager+ = Manager+ { _pendings :: !(M.Map UUID Pending)+ , _ongoings :: !(M.Map UUID OnGoing)+ , _pendingPersistActions :: !(M.Map UUID PendingPersistAction)+ }++--------------------------------------------------------------------------------+initManager :: Manager+initManager =+ Manager+ { _pendings = M.empty+ , _ongoings = M.empty+ , _pendingPersistActions = M.empty+ }++--------------------------------------------------------------------------------+--------------------------------------------------------------------------------+maybeDecodeMessage :: Decode a => ByteString -> Maybe a+maybeDecodeMessage bytes =+ case runGet decodeMessage bytes of+ Right a -> Just a+ _ -> Nothing++--------------------------------------------------------------------------------+unsafeDecodeMessage :: Decode a => ByteString -> a+unsafeDecodeMessage bytes =+ case runGet decodeMessage bytes of+ Right a -> a+ Left e -> error $ "decoding error: " ++ e++--------------------------------------------------------------------------------+data Appeared =+ forall s. Push s =>+ Appeared+ { _appSub :: !(Subscription s)+ , _appEvt :: !ResolvedEvent+ }++--------------------------------------------------------------------------------+onEventAppeared :: Package -> Manager -> Maybe Appeared+onEventAppeared Package{..} Manager{..} =+ case M.lookup packageCorrelation _ongoings of+ Just (OnGoing typ sub) ->+ case (packageCmd, typ) of+ (0xC2, RegularSub _ _ ) ->+ let msg = unsafeDecodeMessage packageData+ evt = getField $ streamResolvedEvent msg in+ Just $ Appeared sub (newResolvedEventFromBuf evt)+ (0xC7, PersistentSub _ _ _) ->+ let msg = unsafeDecodeMessage packageData+ evt = getField $ psseaEvt msg in+ Just $ Appeared sub (newResolvedEvent evt)+ _ -> Nothing+ _ -> Nothing++--------------------------------------------------------------------------------+confirmSub :: (UUID -> IO ())+ -> (UUID -> Text -> AckCmd -> IO ())+ -> Package+ -> Manager+ -> Maybe Confirmed+confirmSub unsub ackF Package{..} Manager{..} =+ case M.lookup packageCorrelation _pendings of+ Just (Pending _ typ cb) ->+ case (packageCmd, typ) of+ (0xC1, RegularSub stream tos) ->+ let !msg = unsafeDecodeMessage packageData+ lcp = getField $ subscribeLastCommitPos msg+ len = getField $ subscribeLastEventNumber msg in+ Just $ Confirmed packageCorrelation typ cb $ do+ chan <- newChan+ let reg = Regular+ { _regId = packageCorrelation+ , _regResolveLinkTos = tos+ , _regLastCommitPos = lcp+ , _regLastEventNumber = len+ , _regChan = chan+ }++ return Subscription+ { subStreamId = stream+ , subUnsubscribe = unsub packageCorrelation+ , subNextEvent = readChan chan+ , subIsSubscribedToAll = stream == ""+ , _subInternal = reg+ }+ (0xC6, PersistentSub grp stream _) ->+ let !msg = unsafeDecodeMessage packageData+ lcp = getField $ pscLastCommitPos msg+ sid = getField $ pscId msg+ len = getField $ pscLastEvtNumber msg in+ Just $ Confirmed packageCorrelation typ cb $ do+ chan <- newChan+ let pes = Persistent+ { _persistId = packageCorrelation+ , _persistChan = chan+ , _persistSubId = sid+ , _persistGroup = grp+ , _persistLastCPos = lcp+ , _persistLastENum = len+ , _persistAckCmd = \cmd ->+ ackF packageCorrelation sid cmd+ }++ return Subscription+ { subStreamId = stream+ , subUnsubscribe = unsub packageCorrelation+ , subNextEvent = readChan chan+ , subIsSubscribedToAll = False+ , _subInternal = pes+ }+ _ -> Nothing+ _ -> Nothing++--------------------------------------------------------------------------------+-- Events+--------------------------------------------------------------------------------+data Subscribe+ = Subscribe+ { _subId :: !UUID+ , _subCallback :: Subscription Regular -> IO ()+ , _subStream :: !Text+ , _subResolveLinkTos :: !Bool+ }++--------------------------------------------------------------------------------+data RegisterSub = forall s. Push s => RegisterSub UUID (Sub s) (Subscription s)++--------------------------------------------------------------------------------+-- Commands+--------------------------------------------------------------------------------+data SubCommand+ = forall s. Push s => SubscribeTo (Sub s) (Subscription s -> IO ())+ | SubmitPersistAction Text+ Text+ PersistAction+ (Either OperationException () -> IO ())++--------------------------------------------------------------------------------+subscriptionNetwork :: Settings+ -> (Package -> Reactive ())+ -> Event Package+ -> Reactive (SubCommand -> IO ())+subscriptionNetwork sett push_pkg e_pkg = do+ -- When a subscription request has been submitted by the user.+ (on_sub, push_sub) <- newEvent++ -- When a subscription has been confirmed by EventStore and we succesfully+ -- create a `forall s. Push s => Subscription s` object.+ (on_reg_sub, push_reg_sub) <- newEvent++ -- When a persist action has been emitted by the user.+ (on_persist_action, push_persist_action) <- newEvent++ let push_pkg_io = pushAsync push_pkg+ push_ack_cmd uuid sid cmd =+ push_pkg_io $ createAckCmdPackage sett uuid sid cmd+ mgr_b <- mfix $ \mgr_b -> do+ let send_unsub = push_pkg_io . createUnsubscribePackage sett++ on_con_sub = filterJust $ snapshot (confirmSub send_unsub+ push_ack_cmd)+ e_pkg mgr_b++ on_drop = filterJust $ snapshot dropError e_pkg mgr_b++ on_persist_action_cfrm =+ filterJust $ snapshot onPersistActionConfirmed e_pkg mgr_b++ mgr_e = fmap confirmed on_reg_sub <>+ fmap subscribeRequest on_sub <>+ fmap newPersistAction on_persist_action <>+ fmap dropped on_drop <>+ fmap persistActionConfirmed on_persist_action_cfrm++ _ <- listen on_drop $ \(Dropped reason sub _) ->+ _subPushEvt sub (Left reason)++ _ <- listen on_persist_action_cfrm $ \(PersistActionConfirmed _ res k) ->+ k res++ _ <- listen on_con_sub $ \(Confirmed uuid typ cb action) -> do+ sub <- action+ _ <- forkIO $ sync $ push_reg_sub (RegisterSub uuid typ sub)+ cb sub++ accum initManager mgr_e++ let on_app = filterJust $ snapshot onEventAppeared e_pkg mgr_b+++ runSubCommand (SubscribeTo typ cb) = do+ uuid <- randomIO+ let sub = Pending+ { _penId = uuid+ , _penSub = typ+ , _penCb = cb+ }+ void $ forkIO $ sync $ push_sub sub+ runSubCommand (SubmitPersistAction group stream typ cb) = do+ uuid <- randomIO+ let action = PendingPersistAction+ { _ppaId = uuid+ , _ppaGroup = group+ , _ppaStream = stream+ , _ppaTyp = typ+ , _ppaCB = cb+ }+ void $ forkIO $ sync $ push_persist_action action+ _ <- listen on_sub (push_pkg_io . createSubscriptionPackage sett)++ _ <- listen on_persist_action (push_pkg_io . createPersistActionPkg sett)++ _ <- listen on_app $ \(Appeared sub evt) ->+ _subPushEvt sub (Right evt)+++ return runSubCommand++--------------------------------------------------------------------------------+createSubscriptionPackage :: Settings -> Pending -> Package+createSubscriptionPackage sett (Pending uuid typ _) =+ case typ of+ RegularSub stream tos ->+ createConnectRegularPackage sett uuid stream tos+ PersistentSub grp str bufSize ->+ createConnectPersistPackage sett uuid grp str bufSize++--------------------------------------------------------------------------------+createConnectRegularPackage :: Settings -> UUID -> Text -> Bool -> Package+createConnectRegularPackage Settings{..} uuid stream tos =+ Package+ { packageCmd = 0xC0+ , packageCorrelation = uuid+ , packageData = runPut $ encodeMessage msg+ , packageCred = s_credentials+ }+ where+ msg = subscribeToStream stream tos++--------------------------------------------------------------------------------+createConnectPersistPackage :: Settings+ -> UUID+ -> Text+ -> Text+ -> Int32+ -> Package+createConnectPersistPackage Settings{..} uuid group stream bufSize =+ Package+ { packageCmd = 0xC5+ , packageCorrelation = uuid+ , packageData = runPut $ encodeMessage msg+ , packageCred = s_credentials+ }+ where+ msg = _connectToPersistentSubscription group stream bufSize++--------------------------------------------------------------------------------+createAckCmdPackage :: Settings -> UUID -> Text -> AckCmd -> Package+createAckCmdPackage sett uuid sid cmd =+ case cmd of+ AckCmd eids -> createAckPackage sett uuid sid eids+ NakCmd act msg eids -> createNakPackage sett uuid sid act msg eids++--------------------------------------------------------------------------------+createAckPackage :: Settings -> UUID -> Text -> [UUID] -> Package+createAckPackage Settings{..} corr sid eids =+ Package+ { packageCmd = 0xCC+ , packageCorrelation = corr+ , packageData = runPut $ encodeMessage msg+ , packageCred = s_credentials+ }+ where+ bytes = toStrict $ foldMap toByteString eids+ msg = persistentSubscriptionAckEvents sid bytes++--------------------------------------------------------------------------------+createNakPackage :: Settings+ -> UUID+ -> Text+ -> NakAction+ -> Maybe Text+ -> [UUID]+ -> Package+createNakPackage Settings{..} corr sid act txt eids =+ Package+ { packageCmd = 0xCD+ , packageCorrelation = corr+ , packageData = runPut $ encodeMessage msg+ , packageCred = s_credentials+ }+ where+ bytes = toStrict $ foldMap toByteString eids+ msg = persistentSubscriptionNakEvents sid bytes txt act++--------------------------------------------------------------------------------+createUnsubscribePackage :: Settings -> UUID -> Package+createUnsubscribePackage Settings{..} uuid =+ Package+ { packageCmd = 0xC3+ , packageCorrelation = uuid+ , packageData = runPut $ encodeMessage UnsubscribeFromStream+ , packageCred = s_credentials+ }++--------------------------------------------------------------------------------+createPersistActionPkg :: Settings -> PendingPersistAction -> Package+createPersistActionPkg Settings{..} (PendingPersistAction aId grp strm typ _) =+ Package+ { packageCmd = cmd+ , packageCorrelation = aId+ , packageData = runPut msg+ , packageCred = s_credentials+ }+ where+ msg =+ case typ of+ PersistCreate sett ->+ encodeMessage $ _createPersistentSubscription grp strm sett+ PersistUpdate sett ->+ encodeMessage $ _updatePersistentSubscription grp strm sett+ PersistDelete ->+ encodeMessage $ _deletePersistentSubscription grp strm+ cmd =+ case typ of+ PersistCreate _ -> 0xC8+ PersistUpdate _ -> 0xCE+ PersistDelete -> 0xCA++--------------------------------------------------------------------------------+data Dropped =+ forall s. Push s =>+ Dropped+ { droppedReason :: !DropReason+ , droppedSub :: !(Subscription s)+ , droppedId :: !UUID+ }++--------------------------------------------------------------------------------+dropError :: Package -> Manager -> Maybe Dropped+dropError Package{..} Manager{..}+ | packageCmd == 0xC4 = do+ OnGoing _ sub <- M.lookup packageCorrelation _ongoings+ msg <- maybeDecodeMessage packageData+ let reason = fromMaybe D_Unsubscribed $ getField $ dropReason msg++ return Dropped+ { droppedReason = reason+ , droppedSub = sub+ , droppedId = packageCorrelation+ }+ | otherwise = Nothing++--------------------------------------------------------------------------------+nonEmptyText :: Text -> Maybe Text+nonEmptyText "" = Nothing+nonEmptyText t = Just t++--------------------------------------------------------------------------------+onPersistActionConfirmed :: Package -> Manager -> Maybe PersistActionConfirmed+onPersistActionConfirmed Package{..} Manager{..} =+ case M.lookup packageCorrelation _pendingPersistActions of+ Just (PendingPersistAction _ grp stream typ cb) ->+ case (packageCmd, typ) of+ (0xC9, PersistCreate _) -> do+ msg <- maybeDecodeMessage packageData+ let res = getField $ cpscResult msg+ reason = nonEmptyText =<< getField (cpscReason msg)+ ret =+ case res of+ CPS_Success ->+ Right ()+ CPS_Fail ->+ Left $ peristentCreationFailure grp+ stream+ reason+ CPS_AlreadyExists ->+ Left $ persistentCreationExists grp+ stream+ CPS_AccessDenied ->+ Left $ persistentAccessDenied stream+ pac = PersistActionConfirmed+ { _pacId = packageCorrelation+ , _pacResult = ret+ , _pacCB = cb+ }+ return pac+ (0xCF, PersistUpdate _) -> do+ msg <- maybeDecodeMessage packageData+ let res = getField $ upscResult msg+ reason = nonEmptyText =<< getField (upscReason msg)+ ret =+ case res of+ UPS_Success ->+ Right ()+ UPS_Fail ->+ Left $ peristentCreationFailure grp+ stream+ reason+ UPS_DoesNotExist ->+ Left $ persistentDoesNotExist grp+ stream+ UPS_AccessDenied ->+ Left $ persistentAccessDenied stream+ pac = PersistActionConfirmed+ { _pacId = packageCorrelation+ , _pacResult = ret+ , _pacCB = cb+ }+ return pac+ (0xCB, PersistDelete) -> do+ msg <- maybeDecodeMessage packageData+ let res = getField $ dpscResult msg+ reason = nonEmptyText =<< getField (dpscReason msg)+ ret =+ case res of+ DPS_Success ->+ Right ()+ DPS_Fail ->+ Left $ peristentCreationFailure grp+ stream+ reason+ DPS_DoesNotExist ->+ Left $ persistentDoesNotExist grp+ stream+ DPS_AccessDenied ->+ Left $ persistentAccessDenied stream+ pac = PersistActionConfirmed+ { _pacId = packageCorrelation+ , _pacResult = ret+ , _pacCB = cb+ }+ return pac+ _ -> Nothing+ _ -> Nothing++--------------------------------------------------------------------------------+-- Model+--------------------------------------------------------------------------------+persistActionConfirmed :: PersistActionConfirmed -> Manager -> Manager+persistActionConfirmed pc s@Manager{..} =+ s { _pendingPersistActions = M.delete (_pacId pc) _pendingPersistActions }++--------------------------------------------------------------------------------+subscribeRequest :: Pending -> Manager -> Manager+subscribeRequest p@(Pending uuid _ _) s@Manager{..} =+ s { _pendings = M.insert uuid p _pendings }++--------------------------------------------------------------------------------+dropped :: Dropped -> Manager -> Manager+dropped d s@Manager{..} = s { _ongoings = M.delete (droppedId d) _ongoings }++--------------------------------------------------------------------------------+confirmed :: RegisterSub -> Manager -> Manager+confirmed (RegisterSub uuid typ sub) s@Manager{..} =+ s { _pendings = M.delete uuid _pendings+ , _ongoings = M.insert uuid (OnGoing typ sub) _ongoings+ }++--------------------------------------------------------------------------------+newPersistAction :: PendingPersistAction -> Manager -> Manager+newPersistAction ppa@PendingPersistAction{..} s@Manager{..} =+ s { _pendingPersistActions = M.insert _ppaId ppa _pendingPersistActions }++--------------------------------------------------------------------------------+--------------------------------------------------------------------------------+peristentCreationFailure :: Text+ -> Text+ -> Maybe Text+ -> OperationException+peristentCreationFailure group stream m_reason = InvalidOperation msg+ where+ msg = "Subscription group " <> group <> " on stream " <> stream <>+ " failed" <> reasonTxt++ reasonTxt = foldMap (" reason: " <>) m_reason++--------------------------------------------------------------------------------+persistentCreationExists :: Text -> Text -> OperationException+persistentCreationExists group stream = InvalidOperation msg+ where+ msg = "Subscription group " <> group <> " on stream " <> stream <>+ " already exists."++--------------------------------------------------------------------------------+persistentDoesNotExist :: Text -> Text -> OperationException+persistentDoesNotExist group stream = InvalidOperation msg+ where+ msg = "Subscription group " <> group <> " on stream " <> stream <>+ " doesn't exist."++--------------------------------------------------------------------------------+persistentAccessDenied :: Text -> OperationException+persistentAccessDenied stream = AccessDenied msg+ where+ msg = "Write access denied for stream " <> stream
Database/EventStore/Internal/Operation/TransactionStartOperation.hs view
@@ -35,14 +35,14 @@ data TransactionEnv = TransactionEnv { _transSettings :: Settings- , _transProcessor :: Processor+ , _transProcessor :: Cmd -> IO () , _transStreamId :: Text , _transExpectedVersion :: ExpectedVersion } -------------------------------------------------------------------------------- transactionStartOperation :: Settings- -> Processor+ -> (Cmd -> IO ()) -> MVar (OperationExceptional Transaction) -> Text -> ExpectedVersion@@ -126,7 +126,7 @@ let op = transactionCommitOperation env trans_id mvar - processorNewOperation _transProcessor op+ _transProcessor (NewOperation op) return as , transactionSendEvents = \evts -> do@@ -134,7 +134,7 @@ let op = transactionWriteOperation env trans_id mvar evts - processorNewOperation _transProcessor op+ _transProcessor (NewOperation op) return as , transactionRollback = return ()
Database/EventStore/Internal/Processor.hs view
@@ -14,17 +14,7 @@ module Database.EventStore.Internal.Processor ( ConnectionException(..) , InternalException(..)- , Processor(..)- , DropReason(..)- , NewSubscriptionCB- , Subscription- , subAwait- , subId- , subStream- , subResolveLinkTos- , subLastCommitPos- , subLastEventNumber- , subUnsubscribe+ , Cmd(..) , newProcessor ) where @@ -32,11 +22,12 @@ import Control.Concurrent import Control.Exception import Data.Functor (void)+import Data.Int import Data.Monoid ((<>)) import Data.Word-import Text.Printf --------------------------------------------------------------------------------+import Data.Text (Text) import Data.UUID import FRP.Sodium import Network@@ -50,19 +41,29 @@ import Database.EventStore.Internal.Types hiding (Event, newEvent) import Database.EventStore.Internal.Util.Sodium import Database.EventStore.Internal.Writer+import Database.EventStore.Logging (Log(..), InfoMessage (Disconnected)) -------------------------------------------------------------------------------- -- Processor ---------------------------------------------------------------------------------data Processor- = Processor- { processorConnect :: HostName -> Int -> IO ()- , processorShutdown :: IO ()- , processorNewOperation :: OperationParams -> IO ()- , processorNewSubcription :: NewSubscriptionCB- }+type Result a = a -> IO ()+type EResult a = Result (Either OperationException a) --------------------------------------------------------------------------------+data Cmd+ = DoConnect HostName Int+ | DoShutdown+ | NewOperation OperationParams+ | NewSub Text Bool (Result (Subscription Regular))+ | CreatePersist Text Text PersistentSubscriptionSettings (EResult ())+ | UpdatePersist Text Text PersistentSubscriptionSettings (EResult ())+ | DeletePersist Text Text (EResult ())+ | ConnectPersist Text Text Int32 (Result (Subscription Persistent))++--------------------------------------------------------------------------------+type Processor = Cmd -> IO ()++-------------------------------------------------------------------------------- newProcessor :: Settings -> IO Processor newProcessor sett = sync . network sett =<< newChan @@ -108,7 +109,7 @@ (pushReconnect Reconnect) onReceived - push_sub <- subscriptionNetwork sett pushSend onReceived+ runSubCmd <- subscriptionNetwork sett pushSend onReceived let stateE = fmap connected onConnected <> fmap reconnected onReconnected <>@@ -157,20 +158,26 @@ _ <- listen onlyHeartbeats $ \pkg -> push_send_io $ heartbeatResponsePackage (packageCorrelation pkg) - let processor =- Processor- { processorConnect = \h p -> void $ forkIO $- sync $ pushConnect $ Connect h p- , processorShutdown = void $ forkIO $ sync $- pushCleanup Cleanup- , processorNewOperation = \o -> void $ forkIO $- sync $ push_new_op o- , processorNewSubcription = push_sub- }+ let runCmd (DoConnect h p) =+ void $ forkIO $ sync $ pushConnect $ Connect h p+ runCmd DoShutdown =+ void $ forkIO $ sync $ pushCleanup Cleanup+ runCmd (NewOperation o) =+ void $ forkIO $ sync $ push_new_op o+ runCmd (NewSub stream tos cb) =+ runSubCmd (SubscribeTo (RegularSub stream tos) cb)+ runCmd (CreatePersist g s stgs cb) =+ runSubCmd (SubmitPersistAction g s (PersistCreate stgs) cb)+ runCmd (UpdatePersist g s stgs cb) =+ runSubCmd (SubmitPersistAction g s (PersistUpdate stgs) cb)+ runCmd (DeletePersist g s cb) =+ runSubCmd (SubmitPersistAction g s PersistDelete cb)+ runCmd (ConnectPersist g s b cb) =+ runSubCmd (SubscribeTo (PersistentSub g s b) cb) _ <- listen onSend (writeChan chan) - return processor+ return runCmd -------------------------------------------------------------------------------- -- Observer@@ -200,13 +207,13 @@ -> IO () connection sett chan push_pkg push_con push_reco host port = do conn <- newConnection sett host port- rid <- forkFinally (readerThread push_pkg conn) (recovering push_reco)+ rid <- forkFinally (readerThread sett push_pkg conn) (recovering push_reco) wid <- forkFinally (writerThread chan conn) (recovering push_reco) push_con (connUUID conn) $ do throwTo rid Stopped throwTo wid Stopped connClose conn- printf "Disconnected %s\n" (toString $ connUUID conn)+ _settingsLog sett (Info $ Disconnected $ connUUID conn) -------------------------------------------------------------------------------- recovering :: IO () -> Either SomeException () -> IO ()
Database/EventStore/Internal/Reader.hs view
@@ -23,20 +23,19 @@ -------------------------------------------------------------------------------- import Database.EventStore.Internal.Connection import Database.EventStore.Internal.Types+import Database.EventStore.Logging ---------------------------------------------------------------------------------readerThread :: (Package -> IO ()) -> Connection -> IO ()-readerThread push_p c = forever $ do+readerThread :: Settings -> (Package -> IO ()) -> Connection -> IO ()+readerThread sett push_p c = forever $ do header_bs <- connRecv c 4 case runGet getLengthPrefix header_bs of- Left _- -> error "Wrong package framing"- Right length_prefix- -> connRecv c length_prefix >>= parsePackage+ Left _ -> _settingsLog sett (Error WrongPackageFraming)+ Right length_prefix -> connRecv c length_prefix >>= parsePackage where parsePackage bs = case runGet getPackage bs of- Left e -> error $ printf "Parsing error [%s]" e+ Left e -> _settingsLog sett (Error $ PackageParsingError e) Right pack -> push_p pack --------------------------------------------------------------------------------
Database/EventStore/Internal/TimeSpan.hs view
@@ -23,6 +23,11 @@ , timeSpanGetMinutes , timeSpanGetSeconds , timeSpanGetMillis+ , timeSpanFromSeconds+ , timeSpanFromMinutes+ , timeSpanFromHours+ , timeSpanFromDays+ , timeSpanTotalMillis ) where --------------------------------------------------------------------------------@@ -106,6 +111,22 @@ Positive -> return ts --------------------------------------------------------------------------------+millisPerSecond :: Int64+millisPerSecond = 1000++--------------------------------------------------------------------------------+millisPerMinute :: Int64+millisPerMinute = millisPerSecond * 60++--------------------------------------------------------------------------------+millisPerHour :: Int64+millisPerHour = millisPerMinute * 60++--------------------------------------------------------------------------------+millisPerDay :: Int64+millisPerDay = millisPerHour * 24++-------------------------------------------------------------------------------- ticksPerMillisecond :: Int64 ticksPerMillisecond = 10000 @@ -126,6 +147,24 @@ ticksPerDay = ticksPerHour * 24 --------------------------------------------------------------------------------+millisPerTick :: Double+millisPerTick = 1 / (realToFrac ticksPerMillisecond)++--------------------------------------------------------------------------------+maxMillis :: Int64+maxMillis =+ truncate+ (((realToFrac (maxBound :: Int64) :: Double)+ / realToFrac ticksPerMillisecond) :: Double)++--------------------------------------------------------------------------------+minMillis :: Int64+minMillis =+ truncate+ (((realToFrac (minBound :: Int64) :: Double)+ / realToFrac ticksPerMillisecond) :: Double)++-------------------------------------------------------------------------------- timeSpanTicks :: Int64 -> TimeSpan timeSpanTicks = TimeSpan @@ -190,6 +229,30 @@ (realToFrac ticksPerMillisecond)) 1000 --------------------------------------------------------------------------------+timeSpanFromSeconds :: Double -> TimeSpan+timeSpanFromSeconds i = interval i millisPerSecond++--------------------------------------------------------------------------------+timeSpanFromMinutes :: Double -> TimeSpan+timeSpanFromMinutes i = interval i millisPerMinute++--------------------------------------------------------------------------------+timeSpanFromHours :: Double -> TimeSpan+timeSpanFromHours i = interval i millisPerHour++--------------------------------------------------------------------------------+timeSpanFromDays :: Double -> TimeSpan+timeSpanFromDays i = interval i millisPerDay++--------------------------------------------------------------------------------+timeSpanTotalMillis :: TimeSpan -> Int64+timeSpanTotalMillis (TimeSpan i) =+ let tmp = (realToFrac i) * millisPerTick in+ if tmp > (realToFrac maxMillis) then maxMillis+ else if tmp < (realToFrac minMillis) then minMillis+ else truncate tmp++-------------------------------------------------------------------------------- data FormatLiteral = Positive | Negative --------------------------------------------------------------------------------@@ -255,3 +318,11 @@ if fraction /= 0 then fromText "." <> fromString (padded 7 '0' $ show fraction) else mempty++--------------------------------------------------------------------------------+interval :: Double -> Int64 -> TimeSpan+interval value scale =+ let tmp = value * (realToFrac scale)+ millis = tmp * (if value >= 0 then 0.5 else (-0.5))+ res = truncate (millis * (realToFrac ticksPerMillisecond)) in+ TimeSpan res
Database/EventStore/Internal/Types.hs view
@@ -46,6 +46,7 @@ import System.Random --------------------------------------------------------------------------------+import Database.EventStore.Logging import Database.EventStore.Internal.TimeSpan --------------------------------------------------------------------------------@@ -68,6 +69,7 @@ | InvalidServerResponse Word8 Word8 -- ^ Expected, Found | ProtobufDecodingError String | ServerError (Maybe Text) -- ^ Reason+ | InvalidOperation Text deriving (Show, Typeable) --------------------------------------------------------------------------------@@ -86,7 +88,7 @@ { eventType :: !Text , eventId :: !(Maybe UUID) , eventData :: !EventData- }+ } deriving (Eq, Show) -------------------------------------------------------------------------------- createEvent :: Text -- ^ Event type@@ -99,6 +101,7 @@ -- | Holds event data. data EventData = Json A.Value (Maybe A.Value)+ deriving (Eq, Show) -------------------------------------------------------------------------------- eventDataType :: EventData -> Int32@@ -143,7 +146,7 @@ | NoStream | EmptyStream | Exact Int32- deriving Show+ deriving (Eq, Show) -------------------------------------------------------------------------------- expVersionInt32 :: ExpectedVersion -> Int32@@ -385,7 +388,7 @@ { positionCommit :: !Int64 -- ^ Commit position of the record , positionPrepare :: !Int64 -- ^ Prepare position of the record }- deriving Show+ deriving (Eq, Show) -------------------------------------------------------------------------------- -- | Representing the start of the transaction file.@@ -406,13 +409,13 @@ , writePosition :: !Position -- ^ 'Position' of the write. }- deriving Show+ deriving (Eq, Show) -------------------------------------------------------------------------------- -- | Returned after deleting a stream. 'Position' of the write. newtype DeleteResult = DeleteResult { deleteStreamPosition :: Position }- deriving Show+ deriving (Eq, Show) -------------------------------------------------------------------------------- -- | Represents a previously written event.@@ -517,12 +520,17 @@ fmap recordedEventStreamId . resolvedEventOriginal --------------------------------------------------------------------------------+-- | The ID of the original event.+resolvedEventOriginalId :: ResolvedEvent -> Maybe UUID+resolvedEventOriginalId = fmap recordedEventId . resolvedEventOriginal++-------------------------------------------------------------------------------- -- | Represents the direction of read operation (both from $all an usual -- streams). data ReadDirection = Forward -- ^ From beginning to end | Backward -- ^ From end to beginning- deriving Show+ deriving (Eq, Show) -------------------------------------------------------------------------------- -- Transaction@@ -567,7 +575,7 @@ { credLogin :: !ByteString , credPassword :: !ByteString }- deriving Show+ deriving (Eq, Show) -------------------------------------------------------------------------------- credentials :: ByteString -- ^ Login@@ -617,6 +625,7 @@ , s_credentials :: Maybe Credentials , s_retry :: Retry , s_reconnect_delay_secs :: Int -- ^ In seconds+ , s_logger :: Maybe (Log -> IO ()) } --------------------------------------------------------------------------------@@ -629,9 +638,17 @@ , s_credentials = Nothing , s_retry = atMost 3 , s_reconnect_delay_secs = 3+ , s_logger = Nothing } --------------------------------------------------------------------------------+_settingsLog :: Settings -> Log -> IO ()+_settingsLog Settings{..} l =+ case s_logger of+ Just k -> k l+ _ -> return ()++-------------------------------------------------------------------------------- -- | Millisecond timespan msDiffTime :: Float -> NominalDiffTime msDiffTime i = fromRational $ toRational (i / 1000)@@ -685,6 +702,25 @@ } deriving Show --------------------------------------------------------------------------------+-- | Gets a custom property value from metadata.+streamMetadataGetCustomPropertyValue :: StreamMetadata -> Text -> Maybe A.Value+streamMetadataGetCustomPropertyValue s k = H.lookup k obj+ where+ obj = streamMetadataCustom s++---------------------------------------------------------------------------------+-- | Get a custom property value from metadata.+streamMetadataGetCustomProperty :: A.FromJSON a+ => StreamMetadata+ -> Text+ -> Maybe a+streamMetadataGetCustomProperty s k = do+ v <- streamMetadataGetCustomPropertyValue s k+ case A.fromJSON v of+ A.Error _ -> Nothing+ A.Success a -> return a++------------------------------------------------------------------------------- instance A.FromJSON StreamMetadata where parseJSON = parseStreamMetadata @@ -961,3 +997,80 @@ | DeletedStreamMetadataResult { streamMetaResultStream :: !Text } -- ^ When the stream is soft-deleted. deriving Show++--------------------------------------------------------------------------------+-- | System supported consumer strategies for use with persistent subscriptions.+data SystemConsumerStrategy+ = DispatchToSingle+ -- ^ Distributes events to a single client until it is full. Then round+ -- robin to the next client.+ | RoundRobin+ -- ^ Distribute events to each client in a round robin fashion.+ deriving (Show, Eq)++--------------------------------------------------------------------------------+strategyText :: SystemConsumerStrategy -> Text+strategyText DispatchToSingle = "DispatchToSingle"+strategyText RoundRobin = "RoundRobin"++--------------------------------------------------------------------------------+strategyFromText :: Text -> Maybe SystemConsumerStrategy+strategyFromText "DispatchToSingle" = Just DispatchToSingle+strategyFromText "RoundRobin" = Just RoundRobin+strategyFromText _ = Nothing++--------------------------------------------------------------------------------+data PersistentSubscriptionSettings =+ PersistentSubscriptionSettings+ { psSettingsResolveLinkTos :: !Bool+ -- ^ Whether or not the persistent subscription should resolve linkTo+ -- events to their linked events.+ , psSettingsStartFrom :: !Int32+ -- ^ Where the subscription should start from (position).+ , psSettingsExtraStats :: !Bool+ -- ^ Whether or not in depth latency statistics should be tracked on this+ -- subscription.+ , psSettingsMsgTimeout :: !TimeSpan+ -- ^ The amount of time after which a message should be considered to be+ -- timeout and retried.+ , psSettingsMaxRetryCount :: !Int32+ -- ^ The maximum number of retries (due to timeout) before a message get+ -- considered to be parked.+ , psSettingsLiveBufSize :: !Int32+ -- ^ The size of the buffer listening to live messages as they happen.+ , psSettingsReadBatchSize :: !Int32+ -- ^ The number of events read at a time when paging in history.+ , psSettingsHistoryBufSize :: !Int32+ -- ^ The number of events to cache when paging through history.+ , psSettingsCheckPointAfter :: !TimeSpan+ -- ^ The amount of time to try checkpoint after.+ , psSettingsMinCheckPointCount :: !Int32+ -- ^ The minimum number of messages to checkpoint.+ , psSettingsMaxCheckPointCount :: !Int32+ -- ^ The maximum number of message to checkpoint. If this number is+ -- reached, a checkpoint will be forced.+ , psSettingsMaxSubsCount :: !Int32+ -- ^ The maximum number of subscribers allowed.+ , psSettingsNamedConsumerStrategy :: !SystemConsumerStrategy+ -- ^ The strategy to use for distributing events to client consumers.+ } deriving (Show, Eq)++--------------------------------------------------------------------------------+-- | System default persistent subscription settings.+defaultPersistentSubscriptionSettings :: PersistentSubscriptionSettings+defaultPersistentSubscriptionSettings =+ PersistentSubscriptionSettings+ { psSettingsResolveLinkTos = False+ , psSettingsStartFrom = (-1)+ , psSettingsExtraStats = False+ , psSettingsMsgTimeout = timeSpanFromSeconds 30+ , psSettingsMaxRetryCount = 500+ , psSettingsLiveBufSize = 500+ , psSettingsReadBatchSize = 10+ , psSettingsHistoryBufSize = 20+ , psSettingsCheckPointAfter = timeSpanFromSeconds 2+ , psSettingsMinCheckPointCount = 10+ , psSettingsMaxCheckPointCount = 1000+ , psSettingsMaxSubsCount = 0+ , psSettingsNamedConsumerStrategy = RoundRobin+ }
+ Database/EventStore/Logging.hs view
@@ -0,0 +1,45 @@+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Logging+-- Copyright : (C) 2015 Yorick Laupa+-- License : (see the file LICENSE)+--+-- Maintainer : Yorick Laupa <yo.eight@gmail.com>+-- Stability : provisional+-- Portability : non-portable+--+--------------------------------------------------------------------------------+module Database.EventStore.Logging where++--------------------------------------------------------------------------------+import Data.UUID++--------------------------------------------------------------------------------+-- | Logging main data structure.+data Log+ = Error ErrorMessage+ | Info InfoMessage+ deriving Show++--------------------------------------------------------------------------------+-- | Classifies error-like log messages.+data ErrorMessage+ = MaxAttemptConnectionReached Int+ -- ^ Indicates max attempt value.+ | WrongPackageFraming+ | PackageParsingError String+ -- ^ Indicates parsing error message.+ deriving Show++--------------------------------------------------------------------------------+-- | Classifies info-like log messages.+data InfoMessage+ = Connecting Int+ -- ^ Indicates current attempt.+ | ConnectionClosed UUID+ -- ^ Indicates connection 'UUID'.+ | Connected UUID + -- ^ Indicates connection 'UUID'.+ | Disconnected UUID+ -- ^ Indicates connection 'UUID'+ deriving Show
README.md view
@@ -1,29 +1,32 @@ EventStore Haskell TCP client =============================-[](https://travis-ci.org/YoEight/eventstore) -Basically, all we have for now are:+[](https://gitter.im/YoEight/eventstore?utm_source=badge&utm_medium=badge&utm_campaign=pr-badge&utm_content=badge)+[](https://travis-ci.org/YoEight/eventstore) - 1. NewEvent- 2. DeleteStream- 3. Transaction- 4. ReadEvent- 5. ReadStreamEvents (Forward and Backward)- 6. ReadAllEvents (Forward and Backward)- 7. Volatile subscriptions- 8. Authentication- 9. Catchup subscriptions+That driver supports: + 1. Read event(s) from regular or $all stream (forward or backward).+ 2. Write event(s) to regular stream.+ 3. Delete regular stream.+ 4. Transactional writes to regular stream.+ 5. Volatile subscriptions to regular or $all stream.+ 6. Catch-up subscriptions to regular or $all stream.+ 7. Competing consumers (a.k.a Persistent subscriptions) to regular stream.+ 8. Authenticated communication with EventStore server.+ 9. Read stream metadata (ACL and custom properties).+ 10. Write stream metadata (ACL and custom properties).+ TODO ====-- 1. Persistent Subscriptions- 2. SSL+ 1. SSL Requirements ============ 1. GHC >= 7.8.3- 2. Cabal >= 1.20- 3. EventStore >= 3.0.0+ 2. Cabal >= 1.18+ 3. EventStore >= 3.0.0 (>= 3.1.0 if you want competing consumers) -(Don't know if it works on Windows)+Tested on Linux and OSX Yosemite.++BSD3 License
eventstore.cabal view
@@ -10,7 +10,7 @@ -- PVP summary: +-+------- breaking API changes -- | | +----- non-breaking API additions -- | | | +--- code changes with no API change-version: 0.7.2.1+version: 0.8.0.0 -- A short (one-line) description of the package. synopsis: EventStore TCP Client@@ -55,7 +55,7 @@ library -- Modules exported by the library. exposed-modules: Database.EventStore-+ Database.EventStore.Logging -- Modules included in this library but not exported. other-modules: Database.EventStore.Catchup Database.EventStore.Internal.Connection@@ -80,7 +80,7 @@ -- Other library packages from which modules are imported. build-depends: base >=4.7 && <5- , aeson >=0.8 && <0.9+ , aeson >=0.8 && <0.10 , async >=2.0 && <2.1 , bytestring >=0.10.4 && <0.11 , cereal >=0.4 && <0.5