eventstore 0.8.0.0 → 0.9.0.0
raw patch · 49 files changed
+5196/−3536 lines, 49 filesdep +eventstoredep +tastydep +tasty-hunitdep −sodiumdep ~aesondep ~basedep ~textPVP ok
version bump matches the API change (PVP)
Dependencies added: eventstore, tasty, tasty-hunit
Dependencies removed: sodium
Dependency ranges changed: aeson, base, text, time
API changes (from Hackage documentation)
- Database.EventStore: AllEventsSlice :: !ReadAllResult -> !Position -> !Position -> !Bool -> ![ResolvedEvent] -> !ReadDirection -> AllEventsSlice
- Database.EventStore: CatchupStreamDeleted :: Text -> CatchupError
- Database.EventStore: CatchupSubscriptionDropReason :: Text -> DropReason -> CatchupError
- Database.EventStore: CatchupUnexpectedStreamStatus :: Text -> ReadStreamResult -> CatchupError
- Database.EventStore: MaxAttempt :: HostName -> Int -> Int -> ConnectionException
- Database.EventStore: RA_ACCESS_DENIED :: ReadAllResult
- Database.EventStore: RA_ERROR :: ReadAllResult
- Database.EventStore: RA_NOT_MODIFIED :: ReadAllResult
- Database.EventStore: RA_SUCCESS :: ReadAllResult
- Database.EventStore: RE_ACCESS_DENIED :: ReadEventResult
- Database.EventStore: RE_ERROR :: ReadEventResult
- Database.EventStore: RE_NOT_FOUND :: ReadEventResult
- Database.EventStore: RE_NO_STREAM :: ReadEventResult
- Database.EventStore: RE_STREAM_DELETED :: ReadEventResult
- Database.EventStore: RE_SUCCESS :: ReadEventResult
- Database.EventStore: RS_ACCESS_DENIED :: ReadStreamResult
- Database.EventStore: RS_ERROR :: ReadStreamResult
- Database.EventStore: RS_NOT_MODIFIED :: ReadStreamResult
- Database.EventStore: RS_NO_STREAM :: ReadStreamResult
- Database.EventStore: RS_STREAM_DELETED :: ReadStreamResult
- Database.EventStore: RS_SUCCESS :: ReadStreamResult
- Database.EventStore: ReadResult :: !ReadEventResult -> !Text -> !Int32 -> !(Maybe ResolvedEvent) -> ReadResult
- Database.EventStore: StreamEventsSlice :: !ReadStreamResult -> !Text -> !Int32 -> !Int32 -> !Int32 -> !Bool -> ![ResolvedEvent] -> !ReadDirection -> StreamEventsSlice
- Database.EventStore: allEventsSliceDirection :: AllEventsSlice -> !ReadDirection
- Database.EventStore: allEventsSliceEvents :: AllEventsSlice -> ![ResolvedEvent]
- Database.EventStore: allEventsSliceFrom :: AllEventsSlice -> !Position
- Database.EventStore: allEventsSliceIsEOS :: AllEventsSlice -> !Bool
- Database.EventStore: allEventsSliceNext :: AllEventsSlice -> !Position
- Database.EventStore: allEventsSliceResult :: AllEventsSlice -> !ReadAllResult
- Database.EventStore: anyStream :: ExpectedVersion
- Database.EventStore: class Identifiable a
- Database.EventStore: data AllEventsSlice
- Database.EventStore: data CatchupError
- Database.EventStore: data OperationException
- Database.EventStore: data ReadAllResult
- Database.EventStore: data ReadEventResult
- Database.EventStore: data ReadStreamResult
- Database.EventStore: data StreamEventsSlice
- Database.EventStore: deleteStreamPosition :: DeleteResult -> Position
- Database.EventStore: emptyStream :: ExpectedVersion
- Database.EventStore: eventResolved :: ResolvedEvent -> Bool
- Database.EventStore: exactStream :: Int32 -> ExpectedVersion
- Database.EventStore: noStream :: ExpectedVersion
- Database.EventStore: positionCommit :: Position -> !Int64
- Database.EventStore: positionPrepare :: Position -> !Int64
- 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: readResultEventNumber :: ReadResult -> !Int32
- Database.EventStore: readResultResolvedEvent :: ReadResult -> !(Maybe ResolvedEvent)
- Database.EventStore: readResultStatus :: ReadResult -> !ReadEventResult
- Database.EventStore: readResultStreamId :: ReadResult -> !Text
- Database.EventStore: recordedEventCreated :: RecordedEvent -> !(Maybe UTCTime)
- Database.EventStore: recordedEventData :: RecordedEvent -> !ByteString
- Database.EventStore: recordedEventId :: RecordedEvent -> !UUID
- Database.EventStore: recordedEventIsJson :: RecordedEvent -> !Bool
- Database.EventStore: recordedEventMetadata :: RecordedEvent -> !(Maybe ByteString)
- Database.EventStore: recordedEventNumber :: RecordedEvent -> !Int32
- Database.EventStore: recordedEventStreamId :: RecordedEvent -> !Text
- Database.EventStore: recordedEventType :: RecordedEvent -> !Text
- Database.EventStore: resolvedEventLink :: ResolvedEvent -> !(Maybe RecordedEvent)
- Database.EventStore: resolvedEventRecord :: ResolvedEvent -> !(Maybe RecordedEvent)
- Database.EventStore: s_credentials :: Settings -> Maybe Credentials
- Database.EventStore: s_heartbeatInterval :: Settings -> NominalDiffTime
- Database.EventStore: s_heartbeatTimeout :: Settings -> NominalDiffTime
- Database.EventStore: s_logger :: Settings -> Maybe (Log -> IO ())
- Database.EventStore: s_reconnect_delay_secs :: Settings -> Int
- Database.EventStore: s_requireMaster :: Settings -> Bool
- Database.EventStore: s_retry :: Settings -> Retry
- Database.EventStore: streamACLDeleteRoles :: StreamACL -> ![Text]
- Database.EventStore: streamACLMetaReadRoles :: StreamACL -> ![Text]
- Database.EventStore: streamACLMetaWriteRoles :: StreamACL -> ![Text]
- Database.EventStore: streamACLReadRoles :: StreamACL -> ![Text]
- Database.EventStore: streamACLWriteRoles :: StreamACL -> ![Text]
- Database.EventStore: streamEventsSliceDirection :: StreamEventsSlice -> !ReadDirection
- Database.EventStore: streamEventsSliceEvents :: StreamEventsSlice -> ![ResolvedEvent]
- Database.EventStore: streamEventsSliceIsEOS :: StreamEventsSlice -> !Bool
- Database.EventStore: streamEventsSliceLast :: StreamEventsSlice -> !Int32
- Database.EventStore: streamEventsSliceNext :: StreamEventsSlice -> !Int32
- Database.EventStore: streamEventsSliceResult :: StreamEventsSlice -> !ReadStreamResult
- Database.EventStore: streamEventsSliceStart :: StreamEventsSlice -> !Int32
- Database.EventStore: streamEventsSliceStreamId :: StreamEventsSlice -> !Text
- Database.EventStore: streamMetaResultData :: StreamMetadataResult -> !StreamMetadata
- Database.EventStore: streamMetaResultStream :: StreamMetadataResult -> !Text
- Database.EventStore: streamMetaResultVersion :: StreamMetadataResult -> !Int32
- Database.EventStore: streamMetadataACL :: StreamMetadata -> !StreamACL
- Database.EventStore: streamMetadataCacheControl :: StreamMetadata -> !(Maybe TimeSpan)
- Database.EventStore: streamMetadataCustom :: StreamMetadata -> !Object
- Database.EventStore: streamMetadataGetCustomProperty :: FromJSON a => StreamMetadata -> Text -> Maybe a
- Database.EventStore: streamMetadataGetCustomPropertyValue :: StreamMetadata -> Text -> Maybe Value
- Database.EventStore: streamMetadataMaxAge :: StreamMetadata -> !(Maybe TimeSpan)
- Database.EventStore: streamMetadataMaxCount :: StreamMetadata -> !(Maybe Int32)
- Database.EventStore: streamMetadataTruncateBefore :: StreamMetadata -> !(Maybe Int32)
- Database.EventStore: subId :: Identifiable a => Subscription a -> UUID
- Database.EventStore: subIsSubscribedToAll :: Subscription a -> Bool
- Database.EventStore: subLastCommitPos :: Identifiable a => Subscription a -> Int64
- Database.EventStore: subLastEventNumber :: Identifiable a => Subscription a -> Maybe Int32
- Database.EventStore: subNextEvent :: Subscription a -> (IO (NextEvent a))
- Database.EventStore: subResolveLinkTos :: Subscription Regular -> Bool
- Database.EventStore: subStreamId :: Subscription a -> Text
- Database.EventStore: subUnsubscribe :: Subscription a -> (IO ())
- Database.EventStore: transactionSendEvents :: Transaction -> [Event] -> IO (Async ())
- Database.EventStore: transactionStart :: Connection -> Text -> ExpectedVersion -> IO (Async Transaction)
- Database.EventStore: writeNextExpectedVersion :: WriteResult -> !Int32
- Database.EventStore: writePosition :: WriteResult -> !Position
- Database.EventStore.Logging: MaxAttemptConnectionReached :: Int -> ErrorMessage
- Database.EventStore.Logging: PackageParsingError :: String -> ErrorMessage
- Database.EventStore.Logging: WrongPackageFraming :: ErrorMessage
- Database.EventStore.Logging: instance Show ErrorMessage
- Database.EventStore.Logging: instance Show InfoMessage
- Database.EventStore.Logging: instance Show Log
+ Database.EventStore: Aborted :: OperationError
+ Database.EventStore: All :: StreamType
+ Database.EventStore: AllStream :: StreamName
+ Database.EventStore: ClosedConnection :: ConnectionException
+ Database.EventStore: MaxAttemptConnectionReached :: HostName -> Int -> Int -> ConnectionException
+ Database.EventStore: PackageParsingError :: String -> ServerConnectionError
+ Database.EventStore: PersistActionAborted :: PersistActionException
+ Database.EventStore: PersistActionAccessDenied :: PersistActionException
+ Database.EventStore: PersistActionAlreadyExist :: PersistActionException
+ Database.EventStore: PersistActionDoesNotExist :: PersistActionException
+ Database.EventStore: PersistActionFail :: PersistActionException
+ Database.EventStore: ReadAccessDenied :: StreamName -> ReadResult t a
+ Database.EventStore: ReadError :: Maybe Text -> ReadResult t a
+ Database.EventStore: ReadEvent :: !Text -> !Int32 -> !ResolvedEvent -> ReadEvent
+ Database.EventStore: ReadEventNotFound :: !Text -> !Int32 -> ReadEvent
+ Database.EventStore: ReadNoStream :: ReadResult RegularStream a
+ Database.EventStore: ReadNotModified :: ReadResult t a
+ Database.EventStore: ReadStreamDeleted :: Text -> ReadResult RegularStream a
+ Database.EventStore: ReadSuccess :: a -> ReadResult t a
+ Database.EventStore: RegularStream :: StreamType
+ Database.EventStore: RunningPersist :: UUID -> Text -> Text -> Int32 -> Text -> Int64 -> (Maybe Int32) -> Running
+ Database.EventStore: RunningReg :: UUID -> Text -> Bool -> Int64 -> (Maybe Int32) -> Running
+ Database.EventStore: StreamName :: Text -> StreamName
+ Database.EventStore: SubAborted :: SubDropReason
+ Database.EventStore: SubAccessDenied :: SubDropReason
+ Database.EventStore: SubNotFound :: SubDropReason
+ Database.EventStore: SubPersistDeleted :: SubDropReason
+ Database.EventStore: SubUnsubscribed :: SubDropReason
+ Database.EventStore: SubscriptionClosed :: Running -> SubDropReason -> SubscriptionClosed
+ Database.EventStore: WrongPackageFraming :: ServerConnectionError
+ Database.EventStore: [positionCommit] :: Position -> !Int64
+ Database.EventStore: [positionPrepare] :: Position -> !Int64
+ 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: [readEventNumber] :: ReadEvent -> !Int32
+ Database.EventStore: [readEventResolved] :: ReadEvent -> !ResolvedEvent
+ Database.EventStore: [readEventStream] :: ReadEvent -> !Text
+ Database.EventStore: [recordedEventCreated] :: RecordedEvent -> !(Maybe UTCTime)
+ Database.EventStore: [recordedEventData] :: RecordedEvent -> !ByteString
+ Database.EventStore: [recordedEventId] :: RecordedEvent -> !UUID
+ Database.EventStore: [recordedEventIsJson] :: RecordedEvent -> !Bool
+ Database.EventStore: [recordedEventMetadata] :: RecordedEvent -> !(Maybe ByteString)
+ Database.EventStore: [recordedEventNumber] :: RecordedEvent -> !Int32
+ Database.EventStore: [recordedEventStreamId] :: RecordedEvent -> !Text
+ Database.EventStore: [recordedEventType] :: RecordedEvent -> !Text
+ Database.EventStore: [resolvedEventLink] :: ResolvedEvent -> !(Maybe RecordedEvent)
+ Database.EventStore: [resolvedEventPosition] :: ResolvedEvent -> !(Maybe Position)
+ Database.EventStore: [resolvedEventRecord] :: ResolvedEvent -> !(Maybe RecordedEvent)
+ Database.EventStore: [s_credentials] :: Settings -> Maybe Credentials
+ Database.EventStore: [s_heartbeatInterval] :: Settings -> NominalDiffTime
+ Database.EventStore: [s_heartbeatTimeout] :: Settings -> NominalDiffTime
+ Database.EventStore: [s_logger] :: Settings -> Maybe (Log -> IO ())
+ Database.EventStore: [s_reconnect_delay_secs] :: Settings -> Int
+ Database.EventStore: [s_requireMaster] :: Settings -> Bool
+ Database.EventStore: [s_retry] :: Settings -> Retry
+ Database.EventStore: [streamACLDeleteRoles] :: StreamACL -> ![Text]
+ Database.EventStore: [streamACLMetaReadRoles] :: StreamACL -> ![Text]
+ Database.EventStore: [streamACLMetaWriteRoles] :: StreamACL -> ![Text]
+ Database.EventStore: [streamACLReadRoles] :: StreamACL -> ![Text]
+ Database.EventStore: [streamACLWriteRoles] :: StreamACL -> ![Text]
+ Database.EventStore: [streamMetaResultData] :: StreamMetadataResult -> !StreamMetadata
+ Database.EventStore: [streamMetaResultStream] :: StreamMetadataResult -> !Text
+ Database.EventStore: [streamMetaResultVersion] :: StreamMetadataResult -> !Int32
+ Database.EventStore: [streamMetadataACL] :: StreamMetadata -> !StreamACL
+ Database.EventStore: [streamMetadataCacheControl] :: StreamMetadata -> !(Maybe TimeSpan)
+ Database.EventStore: [streamMetadataCustom] :: StreamMetadata -> !Object
+ Database.EventStore: [streamMetadataMaxAge] :: StreamMetadata -> !(Maybe TimeSpan)
+ Database.EventStore: [streamMetadataMaxCount] :: StreamMetadata -> !(Maybe Int32)
+ Database.EventStore: [streamMetadataTruncateBefore] :: StreamMetadata -> !(Maybe Int32)
+ Database.EventStore: [writeNextExpectedVersion] :: WriteResult -> !Int32
+ Database.EventStore: [writePosition] :: WriteResult -> !Position
+ Database.EventStore: anyVersion :: ExpectedVersion
+ Database.EventStore: class Slice a where type family Loc a
+ Database.EventStore: data AllSlice
+ Database.EventStore: data OperationError
+ Database.EventStore: data PersistActionException
+ Database.EventStore: data ReadEvent
+ Database.EventStore: data Running
+ Database.EventStore: data ServerConnectionError
+ Database.EventStore: data StreamName
+ Database.EventStore: data StreamSlice
+ Database.EventStore: data StreamType
+ Database.EventStore: data SubDropReason
+ Database.EventStore: data SubscriptionClosed
+ Database.EventStore: data SubscriptionId
+ Database.EventStore: data TransactionId
+ Database.EventStore: emptyStreamVersion :: ExpectedVersion
+ Database.EventStore: exactEventVersion :: Int32 -> ExpectedVersion
+ Database.EventStore: getCustomProperty :: FromJSON a => StreamMetadata -> Text -> Maybe a
+ Database.EventStore: getCustomPropertyValue :: StreamMetadata -> Text -> Maybe Value
+ Database.EventStore: getSubId :: Subscription a -> IO SubscriptionId
+ Database.EventStore: getSubLastCommitPos :: Subscription a -> IO Int64
+ Database.EventStore: getSubLastEventNumber :: Subscription a -> IO (Maybe Int32)
+ Database.EventStore: getSubResolveLinkTos :: Subscription Regular -> Bool
+ Database.EventStore: getSubStream :: Subscription a -> Text
+ Database.EventStore: instance GHC.Classes.Eq Database.EventStore.SubscriptionId
+ Database.EventStore: instance GHC.Classes.Eq Database.EventStore.TransactionId
+ Database.EventStore: instance GHC.Classes.Ord Database.EventStore.SubscriptionId
+ Database.EventStore: instance GHC.Classes.Ord Database.EventStore.TransactionId
+ Database.EventStore: instance GHC.Exception.Exception Database.EventStore.SubscriptionClosed
+ Database.EventStore: instance GHC.Show.Show Database.EventStore.SubscriptionClosed
+ Database.EventStore: instance GHC.Show.Show Database.EventStore.SubscriptionId
+ Database.EventStore: instance GHC.Show.Show Database.EventStore.TransactionId
+ Database.EventStore: isEventResolvedLink :: ResolvedEvent -> Bool
+ Database.EventStore: isSubscribedToAll :: Subscription a -> Bool
+ Database.EventStore: nextEvent :: Subscription a -> IO ResolvedEvent
+ Database.EventStore: nextEventMaybe :: Subscription a -> IO (Maybe ResolvedEvent)
+ Database.EventStore: noStreamVersion :: ExpectedVersion
+ Database.EventStore: recordedEventDataAsJson :: FromJSON a => RecordedEvent -> Maybe a
+ Database.EventStore: resolvedEventDataAsJson :: FromJSON a => ResolvedEvent -> Maybe a
+ Database.EventStore: sliceDirection :: Slice a => a -> ReadDirection
+ Database.EventStore: sliceEOS :: Slice a => a -> Bool
+ Database.EventStore: sliceEvents :: Slice a => a -> [ResolvedEvent]
+ Database.EventStore: sliceFrom :: Slice a => a -> Loc a
+ Database.EventStore: sliceNext :: Slice a => a -> Loc a
+ Database.EventStore: startTransaction :: Connection -> Text -> ExpectedVersion -> IO (Async Transaction)
+ Database.EventStore: transactionId :: Transaction -> TransactionId
+ Database.EventStore: transactionWrite :: Transaction -> [Event] -> IO (Async ())
+ Database.EventStore: unsubscribe :: Subscription a -> IO ()
+ Database.EventStore: waitTillClosed :: Connection -> IO ()
+ Database.EventStore.Logging: PackageReceived :: Word8 -> UUID -> InfoMessage
+ Database.EventStore.Logging: PackageSent :: Word8 -> UUID -> InfoMessage
+ Database.EventStore.Logging: UnexpectedException :: SomeException -> ErrorMessage
+ Database.EventStore.Logging: instance GHC.Show.Show Database.EventStore.Logging.ErrorMessage
+ Database.EventStore.Logging: instance GHC.Show.Show Database.EventStore.Logging.InfoMessage
+ Database.EventStore.Logging: instance GHC.Show.Show Database.EventStore.Logging.Log
- Database.EventStore: AccessDenied :: Text -> OperationException
+ Database.EventStore: AccessDenied :: StreamName -> OperationError
- Database.EventStore: InvalidOperation :: Text -> OperationException
+ Database.EventStore: InvalidOperation :: Text -> OperationError
- Database.EventStore: InvalidServerResponse :: Word8 -> Word8 -> OperationException
+ Database.EventStore: InvalidServerResponse :: Word8 -> Word8 -> OperationError
- Database.EventStore: InvalidTransaction :: OperationException
+ Database.EventStore: InvalidTransaction :: OperationError
- Database.EventStore: ProtobufDecodingError :: String -> OperationException
+ Database.EventStore: ProtobufDecodingError :: String -> OperationError
- Database.EventStore: ResolvedEvent :: !(Maybe RecordedEvent) -> !(Maybe RecordedEvent) -> ResolvedEvent
+ Database.EventStore: ResolvedEvent :: !(Maybe RecordedEvent) -> !(Maybe RecordedEvent) -> !(Maybe Position) -> ResolvedEvent
- Database.EventStore: ServerError :: (Maybe Text) -> OperationException
+ Database.EventStore: ServerError :: (Maybe Text) -> OperationError
- Database.EventStore: StreamDeleted :: Text -> OperationException
+ Database.EventStore: StreamDeleted :: Text -> OperationError
- Database.EventStore: WrongExpectedVersion :: Text -> ExpectedVersion -> OperationException
+ Database.EventStore: WrongExpectedVersion :: Text -> ExpectedVersion -> OperationError
- Database.EventStore: connectToPersistentSubscription :: Connection -> Text -> Text -> Int32 -> IO (Async (Subscription Persistent))
+ Database.EventStore: connectToPersistentSubscription :: Connection -> Text -> Text -> Int32 -> IO (Subscription Persistent)
- Database.EventStore: createPersistentSubscription :: Connection -> Text -> Text -> PersistentSubscriptionSettings -> IO (Async ())
+ Database.EventStore: createPersistentSubscription :: Connection -> Text -> Text -> PersistentSubscriptionSettings -> IO (Async (Maybe PersistActionException))
- Database.EventStore: data ReadResult
+ Database.EventStore: data ReadResult :: StreamType -> * -> *
- Database.EventStore: deletePersistentSubscription :: Connection -> Text -> Text -> IO (Async ())
+ Database.EventStore: deletePersistentSubscription :: Connection -> Text -> Text -> IO (Async (Maybe PersistActionException))
- Database.EventStore: readAllEventsBackward :: Connection -> Position -> Int32 -> Bool -> IO (Async AllEventsSlice)
+ Database.EventStore: readAllEventsBackward :: Connection -> Position -> Int32 -> Bool -> IO (Async AllSlice)
- Database.EventStore: readAllEventsForward :: Connection -> Position -> Int32 -> Bool -> IO (Async AllEventsSlice)
+ Database.EventStore: readAllEventsForward :: Connection -> Position -> Int32 -> Bool -> IO (Async AllSlice)
- Database.EventStore: readEvent :: Connection -> Text -> Int32 -> Bool -> IO (Async ReadResult)
+ Database.EventStore: readEvent :: Connection -> Text -> Int32 -> Bool -> IO (Async (ReadResult RegularStream ReadEvent))
- Database.EventStore: readStreamEventsBackward :: Connection -> Text -> Int32 -> Int32 -> Bool -> IO (Async StreamEventsSlice)
+ Database.EventStore: readStreamEventsBackward :: Connection -> Text -> Int32 -> Int32 -> Bool -> IO (Async (ReadResult RegularStream StreamSlice))
- Database.EventStore: readStreamEventsForward :: Connection -> Text -> Int32 -> Int32 -> Bool -> IO (Async StreamEventsSlice)
+ Database.EventStore: readStreamEventsForward :: Connection -> Text -> Int32 -> Int32 -> Bool -> IO (Async (ReadResult RegularStream StreamSlice))
- Database.EventStore: resolvedEventOriginal :: ResolvedEvent -> Maybe RecordedEvent
+ Database.EventStore: resolvedEventOriginal :: ResolvedEvent -> RecordedEvent
- Database.EventStore: resolvedEventOriginalId :: ResolvedEvent -> Maybe UUID
+ Database.EventStore: resolvedEventOriginalId :: ResolvedEvent -> UUID
- Database.EventStore: resolvedEventOriginalStreamId :: ResolvedEvent -> Maybe Text
+ Database.EventStore: resolvedEventOriginalStreamId :: ResolvedEvent -> Text
- Database.EventStore: subscribe :: Connection -> Text -> Bool -> IO (Async (Subscription Regular))
+ Database.EventStore: subscribe :: Connection -> Text -> Bool -> IO (Subscription Regular)
- Database.EventStore: subscribeToAll :: Connection -> Bool -> IO (Async (Subscription Regular))
+ Database.EventStore: subscribeToAll :: Connection -> Bool -> IO (Subscription Regular)
- Database.EventStore: updatePersistentSubscription :: Connection -> Text -> Text -> PersistentSubscriptionSettings -> IO (Async ())
+ Database.EventStore: updatePersistentSubscription :: Connection -> Text -> Text -> PersistentSubscriptionSettings -> IO (Async (Maybe PersistActionException))
- Database.EventStore: withJson :: Value -> EventData
+ Database.EventStore: withJson :: ToJSON a => a -> EventData
- Database.EventStore: withJsonAndMetadata :: Value -> Value -> EventData
+ Database.EventStore: withJsonAndMetadata :: (ToJSON a, ToJSON b) => a -> b -> EventData
Files
- CHANGELOG.markdown +9/−0
- Database/EventStore.hs +434/−197
- Database/EventStore/Catchup.hs +0/−186
- Database/EventStore/Internal/Connection.hs +125/−30
- Database/EventStore/Internal/Execution/Production.hs +636/−0
- Database/EventStore/Internal/Generator.hs +42/−0
- Database/EventStore/Internal/Manager/Operation.hs +0/−197
- Database/EventStore/Internal/Manager/Operation/Model.hs +178/−0
- Database/EventStore/Internal/Manager/Subscription.hs +220/−1112
- Database/EventStore/Internal/Manager/Subscription/Driver.hs +481/−0
- Database/EventStore/Internal/Manager/Subscription/Message.hs +395/−0
- Database/EventStore/Internal/Manager/Subscription/Model.hs +318/−0
- Database/EventStore/Internal/Manager/Subscription/Packages.hs +137/−0
- Database/EventStore/Internal/Operation.hs +193/−0
- Database/EventStore/Internal/Operation/Catchup.hs +102/−0
- Database/EventStore/Internal/Operation/DeleteStream.hs +63/−0
- Database/EventStore/Internal/Operation/DeleteStream/Message.hs +64/−0
- Database/EventStore/Internal/Operation/DeleteStreamOperation.hs +0/−137
- Database/EventStore/Internal/Operation/Read/Common.hs +145/−0
- Database/EventStore/Internal/Operation/ReadAllEvents.hs +62/−0
- Database/EventStore/Internal/Operation/ReadAllEvents/Message.hs +84/−0
- Database/EventStore/Internal/Operation/ReadAllEventsOperation.hs +0/−200
- Database/EventStore/Internal/Operation/ReadEvent.hs +67/−0
- Database/EventStore/Internal/Operation/ReadEvent/Message.hs +75/−0
- Database/EventStore/Internal/Operation/ReadEventOperation.hs +0/−170
- Database/EventStore/Internal/Operation/ReadStreamEvents.hs +66/−0
- Database/EventStore/Internal/Operation/ReadStreamEvents/Message.hs +80/−0
- Database/EventStore/Internal/Operation/ReadStreamEventsOperation.hs +0/−205
- Database/EventStore/Internal/Operation/StreamMetadata.hs +89/−0
- Database/EventStore/Internal/Operation/Transaction.hs +109/−0
- Database/EventStore/Internal/Operation/Transaction/Message.hs +136/−0
- Database/EventStore/Internal/Operation/TransactionStartOperation.hs +0/−283
- Database/EventStore/Internal/Operation/Write/Common.hs +49/−0
- Database/EventStore/Internal/Operation/WriteEvents.hs +61/−0
- Database/EventStore/Internal/Operation/WriteEvents/Message.hs +71/−0
- Database/EventStore/Internal/Operation/WriteEventsOperation.hs +0/−161
- Database/EventStore/Internal/Operations.hs +32/−0
- Database/EventStore/Internal/Packages.hs +5/−16
- Database/EventStore/Internal/Processor.hs +256/−216
- Database/EventStore/Internal/Reader.hs +0/−95
- Database/EventStore/Internal/Stream.hs +31/−0
- Database/EventStore/Internal/TimeSpan.hs +33/−3
- Database/EventStore/Internal/Types.hs +139/−222
- Database/EventStore/Internal/Util/Sodium.hs +0/−35
- Database/EventStore/Internal/Writer.hs +0/−31
- Database/EventStore/Logging.hs +10/−8
- README.md +95/−18
- eventstore.cabal +45/−14
- tests/integration.hs +59/−0
CHANGELOG.markdown view
@@ -1,3 +1,12 @@+0.9.0.0+-------+* Rewrite entirely the internals.+* Implement integration tests.+* Rename every `ExpectedVersion` smart constructors.+* Improve internal and public documentation.+* Improve failure reports when the connection dropped.+* Implement more robust internal connection.+ 0.8.0.0 ------- * Implement competing consumers.
Database/EventStore.hs view
@@ -1,5 +1,7 @@-{-# LANGUAGE OverloadedStrings #-}-{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE DeriveDataTypeable #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE RecordWildCards #-} -------------------------------------------------------------------------------- -- | -- Module : Database.EventStore@@ -12,15 +14,10 @@ -- -------------------------------------------------------------------------------- module Database.EventStore- ( -- * Event- Event- , EventData- , createEvent- , withJson- , withJsonAndMetadata- -- * Connection- , Connection+ ( -- * Connection+ Connection , ConnectionException(..)+ , ServerConnectionError(..) , Credentials , Settings(..) , Retry@@ -30,6 +27,13 @@ , defaultSettings , connect , shutdown+ , waitTillClosed+ -- * Event+ , Event+ , EventData+ , createEvent+ , withJson+ , withJsonAndMetadata -- * Read Operations , StreamMetadataResult(..) , readEvent@@ -41,8 +45,8 @@ -- * Write Operations , StreamACL(..) , StreamMetadata(..)- , streamMetadataGetCustomPropertyValue- , streamMetadataGetCustomProperty+ , getCustomPropertyValue+ , getCustomProperty , emptyStreamACL , emptyStreamMetadata , deleteStream@@ -95,38 +99,43 @@ , timeSpanTotalMillis -- * Transaction , Transaction- , transactionStart+ , TransactionId+ , startTransaction+ , transactionId , transactionCommit , transactionRollback- , transactionSendEvents- -- * Volatile Subscription- , DropReason(..)- , Identifiable+ , transactionWrite+ -- * Subscription+ , SubscriptionClosed(..)+ , SubscriptionId , Subscription- , NextEvent- , Regular- , Catchup- , Persistent+ , S.Running(..)+ , S.SubDropReason(..)+ -- * Volatile Subscription+ , S.Regular , subscribe , subscribeToAll- , subNextEvent- , subId- , subStreamId- , subIsSubscribedToAll- , subResolveLinkTos- , subLastCommitPos- , subLastEventNumber- , subUnsubscribe+ , getSubId+ , getSubStream+ , isSubscribedToAll+ , unsubscribe+ , nextEvent+ , nextEventMaybe+ , getSubResolveLinkTos+ , getSubLastCommitPos+ , getSubLastEventNumber -- * Catch-up Subscription- , CatchupError(..)+ , S.Catchup , subscribeFrom , subscribeToAllFrom , waitTillCatchup , hasCaughtUp -- * Persistent Subscription+ , S.Persistent , PersistentSubscriptionSettings(..) , SystemConsumerStrategy(..) , NakAction(..)+ , S.PersistActionException(..) , notifyEventsProcessed , notifyEventsFailed , defaultPersistentSubscriptionSettings@@ -135,31 +144,35 @@ , deletePersistentSubscription , connectToPersistentSubscription -- * Results- , AllEventsSlice(..)- , DeleteResult(..)+ , Slice(..)+ , AllSlice+ , Op.DeleteResult(..) , WriteResult(..) , ReadResult(..) , RecordedEvent(..)- , StreamEventsSlice(..)+ , Op.ReadEvent(..)+ , StreamType(..)+ , StreamSlice , Position(..) , ReadDirection(..)- , ReadAllResult(..)- , ReadEventResult(..) , ResolvedEvent(..)- , ReadStreamResult(..)- , OperationException(..)- , eventResolved+ , OperationError(..)+ , StreamName(..)+ , isEventResolvedLink , resolvedEventOriginal+ , resolvedEventDataAsJson , resolvedEventOriginalStreamId , resolvedEventOriginalId+ , recordedEventDataAsJson , positionStart , positionEnd -- * Misc+ , DropReason(..) , ExpectedVersion- , anyStream- , noStream- , emptyStream- , exactStream+ , anyVersion+ , noStreamVersion+ , emptyStreamVersion+ , exactEventVersion -- * Re-export , module Control.Concurrent.Async , (<>)@@ -167,29 +180,31 @@ -------------------------------------------------------------------------------- import Control.Concurrent-import Control.Concurrent.STM (atomically)+import Control.Concurrent.STM import Control.Exception-import Data.ByteString.Lazy (fromStrict)+import Control.Monad (when) import Data.Int+import Data.Maybe import Data.Monoid ((<>))+import Data.Typeable -------------------------------------------------------------------------------- import Control.Concurrent.Async-import Data.Aeson (decode) import Data.Text hiding (group)+import Data.UUID ---------------------------------------------------------------------------------import Database.EventStore.Catchup-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+import Database.EventStore.Internal.Connection hiding (Connection)+import qualified Database.EventStore.Internal.Manager.Subscription as S+import Database.EventStore.Internal.Manager.Subscription.Message+import Database.EventStore.Internal.Operation (OperationError(..))+import qualified Database.EventStore.Internal.Operations as Op+import Database.EventStore.Internal.Operation.Read.Common+import Database.EventStore.Internal.Operation.Write.Common+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.TimeSpan+import Database.EventStore.Internal.Types+import Database.EventStore.Internal.Execution.Production -------------------------------------------------------------------------------- -- Connection@@ -197,37 +212,40 @@ -- | Represents a connection to a single EventStore node. data Connection = Connection- { _runCmd :: Cmd -> IO ()+ { _prod :: Production , _settings :: Settings } -------------------------------------------------------------------------------- -- | Creates a new 'Connection' to a single node. It maintains a full duplex -- connection to the EventStore. An EventStore 'Connection' operates quite--- differently than say a SQL connection. Normally when you use a SQL+-- differently than say a SQL connection. Normally when you use an EventStore -- connection you want to keep the connection open for a much longer of time -- than when you use a SQL connection. ----- Another difference is that with the EventStore 'Connection' all operation+-- Another difference is that with the EventStore 'Connection' all operations -- are handled in a full async manner (even if you call the synchronous -- behaviors). Many threads can use an EvenStore 'Connection' at the same time -- or a single thread can make many asynchronous requests. To get the most--- performance out of the connection it is generally recommend to use it in+-- performance out of the connection it is generally recommended to use it in -- this way. connect :: Settings -> String -- ^ HostName -> Int -- ^ Port -> IO Connection connect settings host port = do- processor <- newProcessor settings- processor (DoConnect host port)+ prod <- newExecutionModel settings host port+ return $ Connection prod settings - return $ Connection processor settings+--------------------------------------------------------------------------------+-- | Waits the 'Connection' to be closed.+waitTillClosed :: Connection -> IO ()+waitTillClosed Connection{..} = prodWaitTillClosed _prod -------------------------------------------------------------------------------- -- | Asynchronously closes the 'Connection'. shutdown :: Connection -> IO ()-shutdown Connection{..} = _runCmd DoShutdown+shutdown Connection{..} = shutdownExecutionModel _prod -------------------------------------------------------------------------------- -- | Sends a single 'Event' to given stream.@@ -240,6 +258,163 @@ sendEvents mgr evt_stream exp_ver [evt] --------------------------------------------------------------------------------+-- | Represents a subscription id.+newtype SubscriptionId = SubId UUID deriving (Eq, Ord, Show)++--------------------------------------------------------------------------------+-- | Determines whether or not any link events encontered in the stream will be+-- resolved.+getSubResolveLinkTos :: Subscription S.Regular -> Bool+getSubResolveLinkTos = S._subTos . _subInner++--------------------------------------------------------------------------------+-- | Non blocking version of `waitTillCatchup`.+hasCaughtUp :: Subscription S.Catchup -> IO Bool+hasCaughtUp sub = atomically $ _hasCaughtUp sub++--------------------------------------------------------------------------------+-- | Waits until 'CatchupSubscription' subscription catch-up its stream.+waitTillCatchup :: Subscription S.Catchup -> IO ()+waitTillCatchup sub = atomically $ do+ caughtUp <- _hasCaughtUp sub+ when (not caughtUp) retry++--------------------------------------------------------------------------------+_hasCaughtUp :: Subscription S.Catchup -> STM Bool+_hasCaughtUp Subscription{..} = do+ SubState sm _ <- readTVar _subVar+ return $ S.hasCaughtUp sm++--------------------------------------------------------------------------------+-- | Tracks a 'Subcription' lifecycle. It holds a 'Subscription' state machine+-- and `SubDropReason` if any.+data SubState a = SubState (S.Subscription a) (Maybe S.SubDropReason)++--------------------------------------------------------------------------------+-- | It's possible to subscribe to a stream and be notified when new events are+-- written to that stream. There are three types of subscription which are+-- available, all of which can be useful in different situations.+--+-- * 'S.Regular'+--+-- * 'S.Catchup'+--+-- * 'S.Persistent'+data Subscription a =+ Subscription+ { _subVar :: TVar (SubState a)+ , _subRun :: TMVar S.Running+ , _subStream :: Text+ , _subProd :: Production+ , _subInner :: a+ }++--------------------------------------------------------------------------------+-- | Gets the ID of the subscription.+getSubId :: Subscription a -> IO SubscriptionId+getSubId Subscription{..} = atomically $ do+ run <- readTMVar _subRun+ return $ SubId $ S.runningUUID run++--------------------------------------------------------------------------------+-- | Gets the subscription stream name.+getSubStream :: Subscription a -> Text+getSubStream = _subStream++--------------------------------------------------------------------------------+-- | Asynchronously unsubscribe from the the stream.+unsubscribe :: Subscription a -> IO ()+unsubscribe Subscription{..} = do+ run <- atomically $ readTMVar _subRun+ pushUnsubscribe _subProd run++--------------------------------------------------------------------------------+-- | If the subscription is on the $all stream.+isSubscribedToAll :: Subscription a -> Bool+isSubscribedToAll = (== "") . getSubStream++--------------------------------------------------------------------------------+-- | The last commit position seen on the subscription (if this a subscription+-- to $all stream).+getSubLastCommitPos :: Subscription a -> IO Int64+getSubLastCommitPos Subscription{..} = atomically $ do+ run <- readTMVar _subRun+ return $ S.runningLastCommitPosition run++--------------------------------------------------------------------------------+-- | The last event number seen on the subscription (if this is a subscription+-- to a single stream).+getSubLastEventNumber :: Subscription a -> IO (Maybe Int32)+getSubLastEventNumber Subscription{..} = atomically $ do+ run <- readTMVar _subRun+ return $ S.runningLastEventNumber run++--------------------------------------------------------------------------------+-- | Awaits for the next event.+nextEvent :: Subscription a -> IO ResolvedEvent+nextEvent sub = atomically $ do+ m <- _nextEventMaybe sub+ case m of+ Nothing -> retry+ Just e -> return e++--------------------------------------------------------------------------------+-- | Non blocking version of 'nextEvent'.+nextEventMaybe :: Subscription a -> IO (Maybe ResolvedEvent)+nextEventMaybe = atomically . _nextEventMaybe++--------------------------------------------------------------------------------+_nextEventMaybe :: Subscription a -> STM (Maybe ResolvedEvent)+_nextEventMaybe Subscription{..} = do+ SubState sub close <- readTVar _subVar+ run <- readTMVar _subRun+ let (res, nxt) = S.readNext sub+ case res of+ Nothing -> do+ case close of+ Nothing -> return Nothing+ Just err -> throwSTM $ SubscriptionClosed run err+ Just e -> do+ writeTVar _subVar $ SubState nxt close+ return $ Just e++--------------------------------------------------------------------------------+-- | Acknowledges those event ids have been successfully processed.+notifyEventsProcessed :: Subscription S.Persistent -> [UUID] -> IO ()+notifyEventsProcessed Subscription{..} evts = do+ run <- atomically $ readTMVar _subRun+ pushAckPersist _subProd run evts++--------------------------------------------------------------------------------+-- | Acknowledges those event ids have failed to be processed successfully.+notifyEventsFailed :: Subscription S.Persistent+ -> NakAction+ -> Maybe Text+ -> [UUID]+ -> IO ()+notifyEventsFailed Subscription{..} act res evts = do+ run <- atomically $ readTMVar _subRun+ pushNakPersist _subProd run act res evts++--------------------------------------------------------------------------------+-- | Modifies 'SubState' internal state machine, letting any 'SubDropReason'+-- untouched.+modifySubSM :: (S.Subscription a -> S.Subscription a)+ -> SubState a+ -> SubState a+modifySubSM k (SubState sm r) = SubState (k sm) r++--------------------------------------------------------------------------------+-- | This exception is raised when the user tries to get the next event from a+-- 'Subscription' that is already closed.+data SubscriptionClosed =+ SubscriptionClosed S.Running S.SubDropReason+ deriving (Show, Typeable)++--------------------------------------------------------------------------------++instance Exception SubscriptionClosed+-------------------------------------------------------------------------------- -- | Sends a list of 'Event' to given stream. sendEvents :: Connection -> Text -- ^ Stream name@@ -247,11 +422,9 @@ -> [Event] -> IO (Async WriteResult) sendEvents Connection{..} evt_stream exp_ver evts = do- (as, mvar) <- createAsync-- let op = writeEventsOperation _settings mvar evt_stream exp_ver evts-- _runCmd (NewOperation op)+ (k, as) <- createOpAsync+ let op = Op.writeEvents _settings evt_stream exp_ver evts+ pushOperation _prod k op return as --------------------------------------------------------------------------------@@ -260,46 +433,92 @@ -> Text -- ^ Stream name -> ExpectedVersion -> Maybe Bool -- ^ Hard delete- -> IO (Async DeleteResult)+ -> IO (Async Op.DeleteResult) deleteStream Connection{..} evt_stream exp_ver hard_del = do- (as, mvar) <- createAsync+ (k, as) <- createOpAsync+ let op = Op.deleteStream _settings evt_stream exp_ver hard_del+ pushOperation _prod k op+ return as - let op = deleteStreamOperation _settings mvar evt_stream exp_ver hard_del+--------------------------------------------------------------------------------+-- | Represents a multi-request transaction with the EventStore.+data Transaction =+ Transaction+ { _tStream :: Text+ , _tTransId :: TransactionId+ , _tExpVer :: ExpectedVersion+ , _tConn :: Connection+ } - _runCmd (NewOperation op)- return as+--------------------------------------------------------------------------------+-- | The id of a 'Transaction'.+newtype TransactionId =+ TransactionId { _unTransId :: Int64 }+ deriving (Eq, Ord, Show) --------------------------------------------------------------------------------+-- | Gets the id of a 'Transaction'.+transactionId :: Transaction -> TransactionId+transactionId = _tTransId++-------------------------------------------------------------------------------- -- | Starts a transaction on given stream.-transactionStart :: Connection+startTransaction :: Connection -> Text -- ^ Stream name -> ExpectedVersion -> IO (Async Transaction)-transactionStart Connection{..} evt_stream exp_ver = do- (as, mvar) <- createAsync+startTransaction conn@Connection{..} evt_stream exp_ver = do+ (k, as) <- createOpAsync+ let op = Op.transactionStart _settings evt_stream exp_ver+ pushOperation _prod k op+ let _F trans_id =+ Transaction+ { _tStream = evt_stream+ , _tTransId = TransactionId trans_id+ , _tExpVer = exp_ver+ , _tConn = conn+ }+ return $ fmap _F as - let op = transactionStartOperation _settings- _runCmd- mvar- evt_stream- exp_ver+--------------------------------------------------------------------------------+-- | Asynchronously writes to a transaction in the EventStore.+transactionWrite :: Transaction -> [Event] -> IO (Async ())+transactionWrite Transaction{..} evts = do+ (k, as) <- createOpAsync+ let Connection{..} = _tConn+ raw_id = _unTransId _tTransId+ op = Op.transactionWrite _settings _tStream _tExpVer raw_id evts+ pushOperation _prod k op+ return as - _runCmd (NewOperation op)+--------------------------------------------------------------------------------+-- | Asynchronously commits this transaction.+transactionCommit :: Transaction -> IO (Async WriteResult)+transactionCommit Transaction{..} = do+ (k, as) <- createOpAsync+ let Connection{..} = _tConn+ raw_id = _unTransId _tTransId+ op = Op.transactionCommit _settings _tStream _tExpVer raw_id+ pushOperation _prod k op return as --------------------------------------------------------------------------------+-- | There isn't such of thing in EventStore parlance. Basically, if you want to+-- rollback, you just have to not 'transactionCommit' a 'Transaction'.+transactionRollback :: Transaction -> IO ()+transactionRollback _ = return ()++-------------------------------------------------------------------------------- -- | Reads a single event from given stream. readEvent :: Connection -> Text -- ^ Stream name -> Int32 -- ^ Event number -> Bool -- ^ Resolve Link Tos- -> IO (Async ReadResult)+ -> IO (Async (ReadResult 'RegularStream Op.ReadEvent)) readEvent Connection{..} stream_id evt_num res_link_tos = do- (as, mvar) <- createAsync-- let op = readEventOperation _settings mvar stream_id evt_num res_link_tos-- _runCmd (NewOperation op)+ (k, as) <- createOpAsync+ let op = Op.readEvent _settings stream_id evt_num res_link_tos+ pushOperation _prod k op return as --------------------------------------------------------------------------------@@ -309,7 +528,7 @@ -> Int32 -- ^ From event number -> Int32 -- ^ Batch size -> Bool -- ^ Resolve Link Tos- -> IO (Async StreamEventsSlice)+ -> IO (Async (ReadResult 'RegularStream StreamSlice)) readStreamEventsForward mgr = readStreamEventsCommon mgr Forward @@ -320,7 +539,7 @@ -> Int32 -- ^ From event number -> Int32 -- ^ Batch size -> Bool -- ^ Resolve Link Tos- -> IO (Async StreamEventsSlice)+ -> IO (Async (ReadResult 'RegularStream StreamSlice)) readStreamEventsBackward mgr = readStreamEventsCommon mgr Backward @@ -331,19 +550,11 @@ -> Int32 -> Int32 -> Bool- -> IO (Async StreamEventsSlice)+ -> IO (Async (ReadResult 'RegularStream StreamSlice)) readStreamEventsCommon Connection{..} dir stream_id start cnt res_link_tos = do- (as, mvar) <- createAsync-- let op = readStreamEventsOperation _settings- dir- mvar- stream_id- start- cnt- res_link_tos-- _runCmd (NewOperation op)+ (k, as) <- createOpAsync+ let op = Op.readStreamEvents _settings dir stream_id start cnt res_link_tos+ pushOperation _prod k op return as --------------------------------------------------------------------------------@@ -352,7 +563,7 @@ -> Position -> Int32 -- ^ Batch size -> Bool -- ^ Resolve Link Tos- -> IO (Async AllEventsSlice)+ -> IO (Async AllSlice) readAllEventsForward mgr = readAllEventsCommon mgr Forward @@ -362,7 +573,7 @@ -> Position -> Int32 -- ^ Batch size -> Bool -- ^ Resolve Link Tos- -> IO (Async AllEventsSlice)+ -> IO (Async AllSlice) readAllEventsBackward mgr = readAllEventsCommon mgr Backward @@ -372,19 +583,11 @@ -> Position -> Int32 -> Bool- -> IO (Async AllEventsSlice)+ -> IO (Async AllSlice) readAllEventsCommon Connection{..} dir pos max_c res_link_tos = do- (as, mvar) <- createAsync-- let op = readAllEventsOperation _settings- dir- mvar- c_pos- p_pos- max_c- res_link_tos-- _runCmd (NewOperation op)+ (k, as) <- createOpAsync+ let op = Op.readAllEvents _settings c_pos p_pos max_c res_link_tos dir+ pushOperation _prod k op return as where Position c_pos p_pos = pos@@ -394,21 +597,26 @@ subscribe :: Connection -> Text -- ^ Stream name -> Bool -- ^ Resolve Link Tos- -> IO (Async (Subscription Regular))+ -> IO (Subscription S.Regular) subscribe Connection{..} stream_id res_lnk_tos = do- tmp <- newEmptyMVar- _runCmd (NewSub stream_id res_lnk_tos (putMVar tmp))- async $ readMVar tmp+ mvar <- newEmptyTMVarIO+ var <- newTVarIO $ SubState S.regularSubscription Nothing+ let mk r = putTMVar mvar r+ recv = readTVar var+ send = writeTVar var+ dropped r = do+ SubState sm _ <- readTVar var+ writeTVar var $ SubState sm (Just r)+ cb = createSubAsync mk recv send dropped+ pushConnectStream _prod cb stream_id res_lnk_tos+ return $ Subscription var mvar stream_id _prod (S.Regular res_lnk_tos) -------------------------------------------------------------------------------- -- | Subcribes to $all stream. subscribeToAll :: Connection -> Bool -- ^ Resolve Link Tos- -> IO (Async (Subscription Regular))-subscribeToAll Connection{..} res_lnk_tos = do- tmp <- newEmptyMVar- _runCmd (NewSub "" res_lnk_tos (putMVar tmp))- async $ readMVar tmp+ -> IO (Subscription S.Regular)+subscribeToAll conn res_lnk_tos = subscribe conn "" res_lnk_tos -------------------------------------------------------------------------------- -- | Subscribes to given stream. If last checkpoint is defined, this will@@ -420,14 +628,11 @@ -> Bool -- ^ Resolve Link Tos -> Maybe Int32 -- ^ Last checkpoint -> Maybe Int32 -- ^ Batch size- -> 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+ -> IO (Subscription S.Catchup)+subscribeFrom conn stream_id res_lnk_tos last_chk_pt batch_m =+ subscribeFromCommon conn stream_id res_lnk_tos batch_m tpe where- evts_fwd cur_num batch_size =- readStreamEventsForward conn stream_id cur_num batch_size res_lnk_tos-- get_sub = subscribe conn stream_id res_lnk_tos+ tpe = Op.RegularCatchup stream_id (fromMaybe 0 last_chk_pt) -------------------------------------------------------------------------------- -- | Same as 'subscribeFrom' but applied to $all stream.@@ -435,15 +640,43 @@ -> Bool -- ^ Resolve Link Tos -> Maybe Position -- ^ Last checkpoint -> Maybe Int32 -- ^ Batch size- -> IO (Subscription Catchup)-subscribeToAllFrom conn res_lnk_tos last_chk_pt batch_m = do- catchupAllStart evts_fwd get_sub last_chk_pt batch_m+ -> IO (Subscription S.Catchup)+subscribeToAllFrom conn res_lnk_tos last_chk_pt batch_m =+ subscribeFromCommon conn "" res_lnk_tos batch_m tpe where- evts_fwd pos batch_size =- readAllEventsForward conn pos batch_size res_lnk_tos+ Position c_pos p_pos = fromMaybe positionStart last_chk_pt+ tpe = Op.AllCatchup c_pos p_pos - get_sub = subscribeToAll conn res_lnk_tos+--------------------------------------------------------------------------------+subscribeFromCommon :: Connection+ -> Text+ -> Bool+ -> Maybe Int32+ -> Op.CatchupState+ -> IO (Subscription S.Catchup)+subscribeFromCommon Connection{..} stream_id res_lnk_tos batch_m tpe = do+ mvar <- newEmptyTMVarIO+ var <- newTVarIO $ SubState S.catchupSubscription Nothing+ let readFrom res =+ case res of+ Left _ -> return ()+ Right (xs, eos, chk) -> atomically $ do+ s <- readTVar var+ let nxt_s = modifySubSM (S.batchRead xs eos chk) s+ writeTVar var nxt_s+ mk = putTMVar mvar+ rcv = readTVar var+ send = writeTVar var+ dropped r = do+ SubState sm _ <- readTVar var+ writeTVar var $ SubState sm (Just r)+ op = Op.catchup _settings tpe res_lnk_tos batch_m+ cb = createSubAsync mk rcv send dropped + pushOperation _prod readFrom op+ pushConnectStream _prod cb stream_id res_lnk_tos+ return $ Subscription var mvar stream_id _prod S.Catchup+ -------------------------------------------------------------------------------- -- | Asynchronously sets the metadata for a stream. setStreamMetadata :: Connection@@ -451,46 +684,20 @@ -> ExpectedVersion -> StreamMetadata -> IO (Async WriteResult)-setStreamMetadata conn evt_stream exp_ver metadata =- let dat = withJson $ streamMetadataJSON metadata- evt = createEvent "$metadata" Nothing dat in- sendEvent conn (metaStreamOf evt_stream) exp_ver evt+setStreamMetadata Connection{..} evt_stream exp_ver metadata = do+ (k, as) <- createOpAsync+ let op = Op.setMetaStream _settings evt_stream exp_ver metadata+ pushOperation _prod k op+ return as -------------------------------------------------------------------------------- -- | 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- async $ atomically $ waitSTM as >>= extractStreamMetadataResult evt_stream-----------------------------------------------------------------------------------extractStreamMetadataResult :: Monad m- => Text- -> ReadResult- -> m StreamMetadataResult-extractStreamMetadataResult stream rres =- case readResultStatus rres of- RE_SUCCESS ->- case action of- Just orig ->- case decode $ fromStrict $ recordedEventData orig of- Just s ->- let res = StreamMetadataResult- { streamMetaResultStream = stream- , streamMetaResultVersion = evt_number- , streamMetaResultData = s- } in- return res- Nothing -> fail "StreamMetadata: wrong format."- Nothing -> fail "impossible: extractStreamMetadataResult"- RE_STREAM_DELETED -> return $ DeletedStreamMetadataResult stream- RE_NOT_FOUND -> return $ NotFoundStreamMetadataResult stream- RE_NO_STREAM -> return $ NotFoundStreamMetadataResult stream- _ -> fail "unexpected ReadEventResult"-- where- action = readResultResolvedEvent rres >>= resolvedEventOriginal- evt_number = readResultEventNumber rres+getStreamMetadata Connection{..} evt_stream = do+ (k, as) <- createOpAsync+ let op = Op.readMetaStream _settings evt_stream+ pushOperation _prod k op+ return as -------------------------------------------------------------------------------- -- | Asynchronously create a persistent subscription group on a stream.@@ -498,11 +705,15 @@ -> Text -> Text -> PersistentSubscriptionSettings- -> IO (Async ())+ -> IO (Async (Maybe S.PersistActionException)) createPersistentSubscription Connection{..} group stream sett = do- (as, mvar) <- createAsync- _runCmd (CreatePersist group stream sett (putMVar mvar))- return as+ mvar <- newEmptyTMVarIO+ let _F res = atomically $+ case res of+ Left e -> putTMVar mvar (Just e)+ _ -> putTMVar mvar Nothing+ pushCreatePersist _prod _F group stream sett+ async $ atomically $ readTMVar mvar -------------------------------------------------------------------------------- -- | Asynchronously update a persistent subscription group on a stream.@@ -510,22 +721,30 @@ -> Text -> Text -> PersistentSubscriptionSettings- -> IO (Async ())+ -> IO (Async (Maybe S.PersistActionException)) updatePersistentSubscription Connection{..} group stream sett = do- (as, mvar) <- createAsync- _runCmd (UpdatePersist group stream sett (putMVar mvar))- return as+ mvar <- newEmptyTMVarIO+ let _F res = atomically $+ case res of+ Left e -> putTMVar mvar (Just e)+ _ -> putTMVar mvar Nothing+ pushUpdatePersist _prod _F group stream sett+ async $ atomically $ readTMVar mvar -------------------------------------------------------------------------------- -- | Asynchronously delete a persistent subscription group on a stream. deletePersistentSubscription :: Connection -> Text -> Text- -> IO (Async ())+ -> IO (Async (Maybe S.PersistActionException)) deletePersistentSubscription Connection{..} group stream = do- (as, mvar) <- createAsync- _runCmd (DeletePersist group stream (putMVar mvar))- return as+ mvar <- newEmptyTMVarIO+ let _F res = atomically $+ case res of+ Left e -> putTMVar mvar (Just e)+ _ -> putTMVar mvar Nothing+ pushDeletePersist _prod _F group stream+ async $ atomically $ readTMVar mvar -------------------------------------------------------------------------------- -- | Asynchronously connect to a persistent subscription given a group on a@@ -534,22 +753,40 @@ -> Text -> Text -> Int32- -> IO (Async (Subscription Persistent))+ -> IO (Subscription S.Persistent) connectToPersistentSubscription Connection{..} group stream bufSize = do- mvar <- newEmptyMVar- _runCmd (ConnectPersist group stream bufSize (putMVar mvar))- async $ readMVar mvar+ mvar <- newEmptyTMVarIO+ var <- newTVarIO $ SubState S.persistentSubscription Nothing+ let mk r = putTMVar mvar r+ recv = readTVar var+ send = writeTVar var+ dropped r = do+ SubState sm _ <- readTVar var+ writeTVar var $ SubState sm (Just r)+ cb = createSubAsync mk recv send dropped+ pushConnectPersist _prod cb group stream bufSize+ return $ Subscription var mvar stream _prod (S.Persistent group) ---------------------------------------------------------------------------------createAsync :: IO (Async a, MVar (OperationExceptional a))-createAsync = do+createOpAsync :: IO (Either OperationError a -> IO (), Async a)+createOpAsync = do mvar <- newEmptyMVar as <- async $ do res <- readMVar mvar either throwIO return res-- return (as, mvar)+ return (putMVar mvar, as) ---------------------------------------------------------------------------------metaStreamOf :: Text -> Text-metaStreamOf s = "$$" <> s+createSubAsync :: (S.Running -> STM ())+ -> STM (SubState a)+ -> (SubState a -> STM ())+ -> (S.SubDropReason -> STM ())+ -> (S.SubConnectEvent -> IO ())+createSubAsync mk rcv send quit = go+ where+ go (S.SubConfirmed run) = atomically $ mk run+ go (S.EventAppeared e) = atomically $ do+ SubState sm close <- rcv+ let nxt = S.eventArrived e sm+ send $ SubState nxt close+ go (S.Dropped r) = atomically $ quit r
− Database/EventStore/Catchup.hs
@@ -1,186 +0,0 @@-{-# LANGUAGE DeriveDataTypeable #-}-{-# LANGUAGE OverloadedStrings #-}-{-# LANGUAGE TypeFamilies #-}------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Catchup--- 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.Catchup where-----------------------------------------------------------------------------------import Control.Concurrent-import Control.Exception-import Control.Monad-import Data.Foldable (traverse_)-import Data.Int-import Data.Maybe-----------------------------------------------------------------------------------import Control.Concurrent.Async-import Data.Text-----------------------------------------------------------------------------------import Database.EventStore.Internal.Manager.Subscription-import Database.EventStore.Internal.Operation.ReadStreamEventsOperation-import Database.EventStore.Internal.Operation.ReadAllEventsOperation-import Database.EventStore.Internal.Types-----------------------------------------------------------------------------------defaultBatchSize :: Int32-defaultBatchSize = 500-----------------------------------------------------------------------------------catchupStart :: (Int32 -> Int32 -> IO (Async StreamEventsSlice))- -> IO (Async (Subscription Regular))- -> Text- -> Maybe Int32- -> Maybe Int32- -> IO (Subscription Catchup)-catchupStart evt_fwd get_sub stream_id batch_size_m last_m = do- chan <- newChan- var <- newEmptyMVar- let batch_size = fromMaybe defaultBatchSize batch_size_m- nxt_read_evt = fromMaybe 0 last_m-- as <- async $ do- res_m <- readEventsTill evt_fwd- (writeChan chan)- stream_id- nxt_read_evt- batch_size-- maybe (return ()) throwIO res_m- action <- get_sub- sub <- wait action- putMVar var sub- keepAwaitingSubEvent stream_id chan sub-- 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 sub-----------------------------------------------------------------------------------catchupAllStart :: ( Position -> Int32 -> IO (Async AllEventsSlice))- -> IO (Async (Subscription Regular))- -> Maybe Position- -> Maybe Int32- -> IO (Subscription Catchup)-catchupAllStart evt_fwd get_sub last_chk_pt_m batch_size_m = do- chan <- newChan- var <- newEmptyMVar- let batch_size = fromMaybe defaultBatchSize batch_size_m- start_pos = fromMaybe positionStart last_chk_pt_m-- as <- async $ do- res_m <- readAllTill evt_fwd- (writeChan chan)- start_pos- batch_size-- maybe (return ()) throwIO res_m- action <- get_sub- sub <- wait action- putMVar var sub- keepAwaitingSubEvent "" chan sub-- 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 sub-----------------------------------------------------------------------------------readEventsTill :: (Int32 -> Int32 -> IO (Async StreamEventsSlice))- -> (Either CatchupError ResolvedEvent -> IO ())- -> Text- -> Int32- -> Int32- -> IO (Maybe CatchupError)-readEventsTill evts_fwd proc_evt stream_id start batch_size =- loop False start- where- loop done cur_evt_num- | done = return Nothing- | otherwise = do- action <- evts_fwd cur_evt_num batch_size- slice <- wait action- case streamEventsSliceResult slice of- RS_SUCCESS -> do- let nxt = streamEventsSliceNext slice- n_done = streamEventsSliceIsEOS slice- evts = streamEventsSliceEvents slice-- traverse_ (proc_evt . Right) evts- loop n_done nxt- RS_NO_STREAM -> loop True cur_evt_num- RS_STREAM_DELETED -> reportError deletedError- s -> reportError $ unexpectedError s-- deletedError = CatchupStreamDeleted stream_id-- unexpectedError s = CatchupUnexpectedStreamStatus stream_id s-- reportError e = do- proc_evt $ Left e- return $ Just e-----------------------------------------------------------------------------------readAllTill :: (Position -> Int32 -> IO (Async AllEventsSlice))- -> (Either CatchupError ResolvedEvent -> IO ())- -> Position- -> Int32- -> IO (Maybe CatchupError)-readAllTill evts_fwd proc_evt start batch_size =- loop False start- where- loop done pos- | done = return Nothing- | otherwise = do- action <- evts_fwd pos batch_size- slice <- wait action- let evts = allEventsSliceEvents slice- nxt = allEventsSliceNext slice- n_done = allEventsSliceIsEOS slice-- traverse_ (proc_evt . Right) evts- loop n_done nxt-----------------------------------------------------------------------------------keepAwaitingSubEvent :: Text- -> Chan (Either CatchupError ResolvedEvent)- -> Subscription Regular- -> IO ()-keepAwaitingSubEvent stream_id chan sub = forever $ do- evt_e <- subNextEvent sub- case evt_e of- Right evt -> writeChan chan (Right evt)- Left r -> do- let e = CatchupSubscriptionDropReason stream_id r-- writeChan chan (Left e)- throwIO e
Database/EventStore/Internal/Connection.hs view
@@ -1,4 +1,6 @@ {-# LANGUAGE DeriveDataTypeable #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE RecordWildCards #-} {-# LANGUAGE ScopedTypeVariables #-} -------------------------------------------------------------------------------- -- |@@ -14,16 +16,18 @@ module Database.EventStore.Internal.Connection ( Connection , ConnectionException(..)+ , HostName , connUUID , connClose- , connFlush , connSend , connRecv+ , connIsClosed , newConnection ) where -------------------------------------------------------------------------------- import Control.Concurrent+import Control.Concurrent.STM import Control.Exception import qualified Data.ByteString as B import Data.Typeable@@ -31,79 +35,170 @@ -------------------------------------------------------------------------------- import Data.UUID+import Data.UUID.V4 import Network-import System.Random -------------------------------------------------------------------------------- import Database.EventStore.Internal.Types import Database.EventStore.Logging ---------------------------------------------------------------------------------data ConnectionException =- MaxAttempt HostName Int Int -- ^ HostName Port MaxAttempt's value+-- | Type of connection issue that can arise during the communication with the+-- server.+data ConnectionException+ = MaxAttemptConnectionReached HostName Int Int+ -- ^ The max reconnection attempt threshold has been reached. Holds a+ -- 'HostName', the port used and the given threshold.+ | ClosedConnection+ -- ^ Use of a close 'Connection'. deriving (Show, Typeable) -------------------------------------------------------------------------------- instance Exception ConnectionException --------------------------------------------------------------------------------+data In a where+ Id :: In UUID+ Close :: In ()+ Send :: B.ByteString -> In ()+ Recv :: Int -> In B.ByteString++--------------------------------------------------------------------------------+-- | Internal representation of a connection with the server. data Connection = Connection- { connUUID :: UUID- , connClose :: IO ()- , connFlush :: IO ()- , connSend :: B.ByteString -> IO ()- , connRecv :: Int -> IO B.ByteString+ { _var :: TMVar State+ , _host :: HostName+ , _port :: Int+ , _setts :: Settings } --------------------------------------------------------------------------------+data State+ = Offline+ | Online !UUID !Handle+ | Closed++--------------------------------------------------------------------------------+-- | Creates a new 'Connection'. newConnection :: Settings -> HostName -> Int -> IO Connection-newConnection sett host port =+newConnection setts host port = do+ var <- newTMVarIO Offline+ return $ Connection var host port setts++--------------------------------------------------------------------------------+-- | Gets current 'Connection' 'UUID'.+connUUID :: Connection -> IO UUID+connUUID conn = execute conn Id++--------------------------------------------------------------------------------+-- | Closes the 'Connection'. It will not retry to reconnect after that call. it+-- means a new 'Connection' has to be created. 'ClosedConnection' exception+-- will be raised if the same 'Connection' object is used after a 'connClose'+-- call.+connClose :: Connection -> IO ()+connClose conn = execute conn Close++--------------------------------------------------------------------------------+-- | Writes 'ByteString' into the buffer.+connSend :: Connection -> B.ByteString -> IO ()+connSend conn b = execute conn (Send b)++--------------------------------------------------------------------------------+-- | Asks the requested amount of bytes from the 'handle'.+connRecv :: Connection -> Int -> IO B.ByteString+connRecv conn i = execute conn (Recv i)++--------------------------------------------------------------------------------+-- | Returns True if the connection is in closed state.+connIsClosed :: Connection -> STM Bool+connIsClosed Connection{..} = do+ r <- readTMVar _var+ case r of+ Closed -> return True+ _ -> return False++--------------------------------------------------------------------------------+-- | Main connection logic. It will automatically reconnect to the server when+-- a exception occured while the 'Handle' is accessed.+execute :: forall a. Connection -> In a -> IO a+execute Connection{..} i = do+ res <- atomically $ do+ s <- takeTMVar _var+ case s of+ Offline -> return $ Right Nothing+ Online u hdl -> return $ Right $ Just (u, hdl)+ Closed -> return $ Left ClosedConnection+ case i of+ Close ->+ case res of+ Left _ -> atomically $ putTMVar _var Closed+ Right Nothing -> atomically $ putTMVar _var Closed+ Right (Just (_, h)) -> do+ hClose h+ atomically $ putTMVar _var Closed+ other ->+ case res of+ Left e -> do+ atomically $ putTMVar _var Closed+ throwIO e+ Right alt -> do+ sres <- case alt of+ Nothing -> newState _setts _host _port+ Just (u, h) -> return $ Right $ Online u h+ case sres of+ Left e -> do+ atomically $ putTMVar _var Closed+ throwIO e+ Right s -> do+ atomically $ putTMVar _var s+ let Online u h = s+ case other of+ Id -> return u+ Send b -> B.hPut h b >> hFlush h+ Recv siz -> B.hGet h siz+ Close -> error "impossible execute"++--------------------------------------------------------------------------------+newState :: Settings -> HostName -> Int -> IO (Either ConnectionException State)+newState sett host port = case s_retry sett of AtMost n -> let loop i = do _settingsLog sett (Info $ Connecting i)- catch (connect sett host port) $ \(_ :: SomeException) -> do+ let action = fmap Right $ connect sett host port+ catch action $ \(_ :: SomeException) -> do threadDelay delay if n <= i- then do- _settingsLog sett- $ Error- $ MaxAttemptConnectionReached i- throwIO $ MaxAttempt host port n+ then return $+ Left $+ MaxAttemptConnectionReached host port n else loop (i + 1) in loop 1 KeepRetrying -> let endlessly i = do _settingsLog sett (Info $ Connecting i)- catch (connect sett host port) $ \(_ :: SomeException) -> do+ let action = fmap Right $ connect sett host port+ catch action $ \(_ :: SomeException) -> threadDelay delay >> endlessly (i + 1) in endlessly (1 :: Int) where- delay = (s_reconnect_delay_secs sett) * secs+ delay = s_reconnect_delay_secs sett * secs -------------------------------------------------------------------------------- secs :: Int secs = 1000000 ---------------------------------------------------------------------------------connect :: Settings -> HostName -> Int -> IO Connection+connect :: Settings -> HostName -> Int -> IO State connect sett host port = do hdl <- connectTo host (PortNumber $ fromIntegral port) hSetBuffering hdl NoBuffering- uuid <- randomIO+ uuid <- nextRandom regularConnection sett hdl uuid ---------------------------------------------------------------------------------regularConnection :: Settings -> Handle -> UUID -> IO Connection+regularConnection :: Settings -> Handle -> UUID -> IO State 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- }+ return $ Online uuid h
+ Database/EventStore/Internal/Execution/Production.hs view
@@ -0,0 +1,636 @@+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE DeriveDataTypeable #-}+{-# LANGUAGE ExistentialQuantification #-}+{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE ScopedTypeVariables #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Execution.Production+-- Copyright : (C) 2015 Yorick Laupa+-- License : (see the file LICENSE)+--+-- Maintainer : Yorick Laupa <yo.eight@gmail.com>+-- Stability : provisional+-- Portability : non-portable+--+-- Production execution model. It's striving for robustness. The model consists+-- on 4 threads. The Reader thread that reads 'Package' from the connection, the+-- Runner thread which executes finalizers submitted by the user (typically what+-- to do on operation completion or when a event has arrived for a subscription)+-- , the writer thread that sends 'Package' to the server, and the Manager+-- thread that handles requests coming both from the user and the Reader+-- thread. If the Reader or Runner threads die, it will be restarted by the+-- Manager thread if the connection hasn't been closed by user in the meantime.+--------------------------------------------------------------------------------+module Database.EventStore.Internal.Execution.Production+ ( Production+ , ServerConnectionError(..)+ , newExecutionModel+ , pushOperation+ , shutdownExecutionModel+ , pushConnectStream+ , pushConnectPersist+ , pushCreatePersist+ , pushUpdatePersist+ , pushDeletePersist+ , pushAckPersist+ , pushNakPersist+ , pushUnsubscribe+ , prodWaitTillClosed+ ) where++--------------------------------------------------------------------------------+import Prelude hiding (take)+import Control.Concurrent+import Control.Concurrent.STM+import Control.Exception+import Control.Monad+import Control.Monad.Fix+import Data.IORef+import Data.Int+import Data.Foldable+import Data.Typeable+import Text.Printf++--------------------------------------------------------------------------------+import Data.Serialize.Get hiding (Done)+import Data.Serialize.Put+import Data.Text+import Data.UUID++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Connection+import Database.EventStore.Internal.Generator+import Database.EventStore.Internal.Manager.Subscription hiding+ ( submitPackage+ , unsubscribe+ , ackPersist+ , nakPersist+ , abort+ )+import Database.EventStore.Internal.Operation hiding (retry)+import Database.EventStore.Internal.Packages+import Database.EventStore.Internal.Processor+import Database.EventStore.Internal.Types+import Database.EventStore.Logging++--------------------------------------------------------------------------------+data Worker+ = Reader ThreadId+ | Runner ThreadId+ | Writer ThreadId+ deriving Show++--------------------------------------------------------------------------------+-- | Raised when the server responded in an unexpected way.+data ServerConnectionError+ = WrongPackageFraming+ -- ^ TCP package sent by the server had a wrong framing.+ | PackageParsingError String+ -- ^ Server sent a malformed TCP package.+ deriving (Show, Typeable)++--------------------------------------------------------------------------------+instance Exception ServerConnectionError++--------------------------------------------------------------------------------+-- | Used to determine if we hit the end of the queue.+data Slot a = Slot !a | End++--------------------------------------------------------------------------------+-- | A 'TQueue' that can be cycled.+newtype CycleQueue a = CycleQueue (TQueue (Slot a))++--------------------------------------------------------------------------------+-- | Creates an empty 'CycleQueue'.+newCycleQueue :: IO (CycleQueue a)+newCycleQueue = fmap CycleQueue newTQueueIO++--------------------------------------------------------------------------------+-- | Gets an element from the 'CycleQueue'.+readCycleQueue :: CycleQueue a -> STM a+readCycleQueue (CycleQueue q) = do+ Slot a <- readTQueue q+ return a++--------------------------------------------------------------------------------+-- | Writes an element to the 'CycleQueue'.+writeCycleQueue :: CycleQueue a -> a -> STM ()+writeCycleQueue (CycleQueue q) a = writeTQueue q (Slot a)++--------------------------------------------------------------------------------+-- | Empties a 'CycleQueue'.+emptyCycleQueue :: CycleQueue a -> STM ()+emptyCycleQueue (CycleQueue q) = writeTQueue q End >> go+ where+ go = do+ s <- readTQueue q+ case s of+ End -> return ()+ _ -> go++--------------------------------------------------------------------------------+-- | Updates a 'CycleQueue'.+updateCycleQueue :: CycleQueue a -> (a -> STM (Maybe a)) -> STM ()+updateCycleQueue (CycleQueue q) k = writeTQueue q End >> go+ where+ go = do+ s <- readTQueue q+ case s of+ End -> return ()+ Slot a -> do+ r <- k a+ case r of+ Nothing -> go+ Just a' -> writeTQueue q (Slot a') >> go++--------------------------------------------------------------------------------+-- | Indicates if a 'CycleQueue' is empty.+isEmptyCycleQueue :: CycleQueue a -> STM Bool+isEmptyCycleQueue (CycleQueue q) = isEmptyTQueue q++--------------------------------------------------------------------------------+wkUpdState :: Worker -> State -> State+wkUpdState (Reader tid) s = s { _reader = Just tid }+wkUpdState (Runner tid) s = s { _runner = Just tid }+wkUpdState (Writer tid) s = s { _writer = Just tid }++--------------------------------------------------------------------------------+-- | Holds the execution model state.+data Production =+ Prod+ { _submit :: TVar (Msg -> IO ())+ -- ^ The action to call when pushing new command.+ , _waitClosed :: STM ()+ -- ^ Action that attests the execution model has been closed successfully.+ -- It doesn't mean the execution model hasn't been shutdown because of+ -- some random exception.+ }++--------------------------------------------------------------------------------+-- | Main execution environment used among different transitions.+data Env =+ Env+ { _setts :: Settings+ -- ^ Global settings reference.+ , _queue :: CycleQueue Msg+ -- ^ That queue ties the user, the reader thread and the manager thread.+ -- The user and the reader push new messages onto the queue while the+ -- manager dequeue and handles one message at the time.+ , _pkgQueue :: CycleQueue Package+ -- ^ That queue ties the writer thread with the manager thread. The writer+ -- dequeue packages from that queue and sends those to the server. While+ -- the manager pushes new packages on every new submitted operation.+ , _jobQueue :: CycleQueue Job+ -- ^ That queue ties the runner thread with the manager thread. The runner+ -- dequeues IO action from it while the manager pushes new command+ -- finalizers as those arrived.+ , _state :: TVar State+ -- ^ Holds manager thread state.+ , _nextSubmit :: TVar (Msg -> IO ())+ -- ^ Indicates the action to call in order to push new commands.+ , _connRef :: IORef Connection+ -- ^ Connection to the server.+ , _disposed :: TMVar ()+ -- ^ Indicates when the production execution model has been shutdown and+ -- disposed any ongoing operations.+ }++--------------------------------------------------------------------------------+data Msg+ = Stopped Worker SomeException+ | Arrived Package+ | Shutdown+ | forall a.+ NewOperation (Either OperationError a -> IO ()) (Operation a)+ | ConnectStream (SubConnectEvent -> IO ()) Text Bool+ | ConnectPersist (SubConnectEvent -> IO ()) Text Text Int32+ | Unsubscribe Running+ | CreatePersist (Either PersistActionException ConfirmedAction -> IO ())+ Text Text PersistentSubscriptionSettings+ | UpdatePersist (Either PersistActionException ConfirmedAction -> IO ())+ Text Text PersistentSubscriptionSettings+ | DeletePersist (Either PersistActionException ConfirmedAction -> IO ())+ Text Text+ | AckPersist Running [UUID]+ | NakPersist Running NakAction (Maybe Text) [UUID]++--------------------------------------------------------------------------------+pushCmd :: Production -> Msg -> IO ()+pushCmd (Prod _sender _) msg = do+ push <- readTVarIO _sender+ push msg++--------------------------------------------------------------------------------+-- | Asks to shutdown the connection to the server asynchronously.+shutdownExecutionModel :: Production -> IO ()+shutdownExecutionModel prod = pushCmd prod Shutdown++--------------------------------------------------------------------------------+-- | Pushes a new 'Operation' asynchronously.+pushOperation :: Production+ -> (Either OperationError a -> IO ())+ -> Operation a+ -> IO ()+pushOperation prod k op = pushCmd prod (NewOperation k op)++--------------------------------------------------------------------------------+-- | Subscribes to a regular stream.+pushConnectStream :: Production+ -> (SubConnectEvent -> IO ())+ -> Text+ -> Bool+ -> IO ()+pushConnectStream prod k n tos = pushCmd prod (ConnectStream k n tos)++--------------------------------------------------------------------------------+-- | Subscribes to a persistent subscription.+pushConnectPersist :: Production+ -> (SubConnectEvent -> IO ())+ -> Text+ -> Text+ -> Int32+ -> IO ()+pushConnectPersist prod k g n buf = pushCmd prod (ConnectPersist k g n buf)++--------------------------------------------------------------------------------+-- | Creates a persistent subscription.+pushCreatePersist :: Production+ -> (Either PersistActionException ConfirmedAction -> IO ())+ -> Text+ -> Text+ -> PersistentSubscriptionSettings+ -> IO ()+pushCreatePersist prod k g n setts = pushCmd prod (CreatePersist k g n setts)++--------------------------------------------------------------------------------+-- | Updates a persistent subscription.+pushUpdatePersist :: Production+ -> (Either PersistActionException ConfirmedAction -> IO ())+ -> Text+ -> Text+ -> PersistentSubscriptionSettings+ -> IO ()+pushUpdatePersist prod k g n setts = pushCmd prod (UpdatePersist k g n setts)++--------------------------------------------------------------------------------+-- | Deletes a persistent subscription.+pushDeletePersist :: Production+ -> (Either PersistActionException ConfirmedAction -> IO ())+ -> Text+ -> Text+ -> IO ()+pushDeletePersist prod k g n = pushCmd prod (DeletePersist k g n)++--------------------------------------------------------------------------------+-- | Acknowledges a set of events has been successfully handled.+pushAckPersist :: Production -> Running -> [UUID] -> IO ()+pushAckPersist prod run evts = pushCmd prod (AckPersist run evts)++--------------------------------------------------------------------------------+-- | Acknowledges a set of events hasn't been handled successfully.+pushNakPersist :: Production+ -> Running+ -> NakAction+ -> Maybe Text+ -> [UUID]+ -> IO ()+pushNakPersist prod run act res evts =+ pushCmd prod (NakPersist run act res evts)++--------------------------------------------------------------------------------+-- | Unsubscribe from a subscription.+pushUnsubscribe :: Production -> Running -> IO ()+pushUnsubscribe prod r = pushCmd prod (Unsubscribe r)++--------------------------------------------------------------------------------+-- | Waits the execution model to close properly.+prodWaitTillClosed :: Production -> IO ()+prodWaitTillClosed (Prod _ disposed) = atomically disposed++--------------------------------------------------------------------------------+newtype Job = Job (IO ())++--------------------------------------------------------------------------------+-- Internal Production state.+--------------------------------------------------------------------------------+data State =+ State+ { _proc :: !(Processor (IO ()))+ , _reader :: !(Maybe ThreadId)+ , _runner :: !(Maybe ThreadId)+ , _writer :: !(Maybe ThreadId)+ }++--------------------------------------------------------------------------------+emptyState :: Settings -> Generator -> State+emptyState setts gen = State (newProcessor setts gen) Nothing Nothing Nothing++--------------------------------------------------------------------------------+updateProc :: Processor (IO ()) -> State -> State+updateProc p s = s { _proc = p }++--------------------------------------------------------------------------------+-- | Reader thread. Keeps reading 'Package' from the connection.+--------------------------------------------------------------------------------+reader :: Settings -> CycleQueue Msg -> Connection -> IO ()+reader sett queue c = forever $ do+ header_bs <- connRecv c 4+ case runGet getLengthPrefix header_bs of+ Left _ -> throwIO WrongPackageFraming+ Right length_prefix -> connRecv c length_prefix >>= parsePackage+ where+ parsePackage bs =+ case runGet getPackage bs of+ Left e -> throwIO $ PackageParsingError e+ Right pkg -> do+ atomically $ writeCycleQueue queue (Arrived pkg)+ let cmd = packageCmd pkg+ uuid = packageCorrelation pkg+ _settingsLog sett $ Info $ PackageReceived cmd uuid++--------------------------------------------------------------------------------+-- | Writer thread, writes incoming 'Package's+--------------------------------------------------------------------------------+writer :: Settings -> CycleQueue Package -> Connection -> IO ()+writer setts pkg_queue conn = forever $ do+ pkg <- atomically $ readCycleQueue pkg_queue+ connSend conn $ runPut $ putPackage pkg+ let cmd = packageCmd pkg+ uuid = packageCorrelation pkg+ _settingsLog setts $ Info $ PackageSent cmd uuid++--------------------------------------------------------------------------------+getLengthPrefix :: Get Int+getLengthPrefix = fmap fromIntegral getWord32le++--------------------------------------------------------------------------------+getPackage :: Get Package+getPackage = do+ cmd <- getWord8+ flg <- getFlag+ col <- getUUID+ cred <- getCredentials flg+ rest <- remaining+ dta <- getBytes rest++ let pkg = Package+ { packageCmd = cmd+ , packageCorrelation = col+ , packageData = dta+ , packageCred = cred+ }++ return pkg++--------------------------------------------------------------------------------+getFlag :: Get Flag+getFlag = do+ wd <- getWord8+ case wd of+ 0x00 -> return None+ 0x01 -> return Authenticated+ _ -> fail $ printf "TCP: Unhandled flag value 0x%x" wd++--------------------------------------------------------------------------------+getCredEntryLength :: Get Int+getCredEntryLength = fmap fromIntegral getWord8++--------------------------------------------------------------------------------+getCredentials :: Flag -> Get (Maybe Credentials)+getCredentials None = return Nothing+getCredentials _ = do+ loginLen <- getCredEntryLength+ login <- getBytes loginLen+ passwLen <- getCredEntryLength+ passw <- getBytes passwLen+ return $ Just $ credentials login passw++--------------------------------------------------------------------------------+getUUID :: Get UUID+getUUID = do+ bs <- getLazyByteString 16+ case fromByteString bs of+ Just uuid -> return uuid+ _ -> fail "TCP: Wrong UUID format"++--------------------------------------------------------------------------------+-- Runner thread. Keeps running job comming from the Manager thread.+--------------------------------------------------------------------------------+runner :: CycleQueue Job -> IO ()+runner job_queue = forever $ do+ Job j <- atomically $ readCycleQueue job_queue+ j++--------------------------------------------------------------------------------+-- | Spawns a new thread worker.+spawn :: Env -> (ThreadId -> Worker) -> IO Worker+spawn Env{..} mk = do+ conn <- readIORef _connRef+ tid <- mfix $ \tid ->+ let worker = mk tid+ action =+ case worker of+ Reader _ -> reader _setts _queue conn+ Runner _ -> runner _jobQueue+ Writer _ -> writer _setts _pkgQueue conn in+ forkFinally action $ \r ->+ case r of+ Left e ->+ case asyncExceptionFromException e of+ Just ThreadKilled -> return ()+ _ -> atomically $ writeCycleQueue _queue+ $ Stopped worker e+ _ -> return ()+ return $ mk tid++--------------------------------------------------------------------------------+-- | Loops over a 'Processor''s 'Transition' state machine, returning an updated+-- 'Processor' model at the end.+runTransition :: Env -> Transition (IO ()) -> STM (Processor (IO ()))+runTransition Env{..} = go+ where+ go (Produce j nxt) = do+ let job = Job j+ writeCycleQueue _jobQueue job+ go nxt+ go (Transmit pkg nxt) = do+ writeCycleQueue _pkgQueue pkg+ go nxt+ go (Await new_proc) = return new_proc++--------------------------------------------------------------------------------+-- | First execution mode. It spawns initial reader, runner and writer threads.+-- Then it switches to 'cruising' mode.+bootstrap :: Env -> IO ()+bootstrap env@Env{..} = do+ rew <- spawn env Reader+ ruw <- spawn env Runner+ wrw <- spawn env Writer+ let _F = wkUpdState rew .+ wkUpdState ruw .+ wkUpdState wrw+ atomically $ modifyTVar' _state _F+ cruising env++--------------------------------------------------------------------------------+-- | Crusing execution mode. Reads and handle message coming from the channel as+-- those are arrived. That mode is used when the connection to the server is+-- still live. We might have deconnection once in a while but at the end, if+-- we managed to reconnect to it, we consider everything is fine.+cruising :: Env -> IO ()+cruising env@Env{..} = do+ msg <- atomically $ readCycleQueue _queue+ s <- readTVarIO _state+ case msg of+ Stopped _ e -> throwIO e+ Arrived pkg -> do+ let sm = submitPackage pkg $ _proc s+ atomically $ do+ new_proc <- runTransition env sm+ modifyTVar' _state $ updateProc new_proc+ cruising env+ Shutdown -> throwIO ClosedConnection+ NewOperation k op -> do+ let sm = newOperation k op $ _proc s+ atomically $ do+ new_proc <- runTransition env sm+ modifyTVar' _state $ updateProc new_proc+ cruising env+ ConnectStream k n tos -> do+ let sm = connectRegularStream k n tos $ _proc s+ atomically $ do+ new_proc <- runTransition env sm+ modifyTVar' _state $ updateProc new_proc+ cruising env+ ConnectPersist k g n b -> do+ let sm = connectPersistent k g n b $ _proc s+ atomically $ do+ new_proc <- runTransition env sm+ modifyTVar' _state $ updateProc new_proc+ cruising env+ Unsubscribe r -> do+ let sm = unsubscribe r $ _proc s+ atomically $ do+ new_proc <- runTransition env sm+ modifyTVar' _state $ updateProc new_proc+ cruising env+ CreatePersist k g n psetts -> do+ let sm = createPersistent k g n psetts $ _proc s+ atomically $ do+ new_proc <- runTransition env sm+ modifyTVar' _state $ updateProc new_proc+ cruising env+ UpdatePersist k g n psetts -> do+ let sm = updatePersistent k g n psetts $ _proc s+ atomically $ do+ new_proc <- runTransition env sm+ modifyTVar' _state $ updateProc new_proc+ cruising env+ DeletePersist k g n -> do+ let sm = deletePersistent k g n $ _proc s+ atomically $ do+ new_proc <- runTransition env sm+ modifyTVar' _state $ updateProc new_proc+ cruising env+ AckPersist run evts -> do+ let sm = ackPersist (return ()) run evts $ _proc s+ atomically $ do+ new_proc <- runTransition env sm+ modifyTVar' _state $ updateProc new_proc+ cruising env+ NakPersist run act res evts -> do+ let sm = nakPersist (return ()) run act res evts $ _proc s+ atomically $ do+ new_proc <- runTransition env sm+ modifyTVar' _state $ updateProc new_proc+ cruising env++--------------------------------------------------------------------------------+-- | That mode is triggered either because the user asks to shutdown the+-- connection or because the connection to server has been dropped and we+-- can't reconnect.+closing :: Env -> IO ()+closing env@Env{..} = do+ State _ retid rutid wutid <- readTVarIO _state+ -- We kill reader and writer threads to avoid in fly package in the cleaning+ -- phase.+ traverse_ killThread retid+ traverse_ killThread wutid++ -- Discards every 'Package' that was about to be sent.+ atomically $ emptyCycleQueue _pkgQueue++ -- Takes care of 'Package's that have already arrived. Just in case those+ -- are completing or moving forward ongoing operations. Every ongoing+ -- request is kept for later reconnection. Some transient operations like+ -- Ack, Nak or Unsubscribe are just discard.+ atomically $ updateCycleQueue _queue $ \nxt ->+ case nxt of+ Arrived pkg -> do+ s <- readTVar _state+ let sm = submitPackage pkg $ _proc s+ nxt_proc <- runTransition env sm+ modifyTVar' _state $ updateProc nxt_proc+ return Nothing+ Shutdown -> return Nothing+ AckPersist _ _ -> return Nothing+ NakPersist _ _ _ _ -> return Nothing+ Unsubscribe _ -> return Nothing+ Stopped _ _ -> return Nothing+ _ -> return $ Just nxt++ -- If the connection is already closed, it will throw an exception. We just+ -- make sure it doesn't interfere with the cleaning process.+ conn <- readIORef _connRef+ _ <- try $ connClose conn :: (IO (Either ConnectionException ()))+ atomically $ do+ s <- readTVar _state+ _ <- runTransition env $ abort $ _proc s+ return ()++ -- Waits the runner thread to deal with its jobs list.+ atomically $ do+ end <- isEmptyCycleQueue _jobQueue+ unless end retry++ traverse_ killThread rutid++--------------------------------------------------------------------------------+raiseException :: Exception e => e -> Msg -> IO ()+raiseException e _ = throwIO e++--------------------------------------------------------------------------------+-- | Main Production execution model entry point.+newExecutionModel :: Settings -> HostName -> Int -> IO Production+newExecutionModel setts host port = do+ gen <- newGenerator+ queue <- newCycleQueue+ pkg_queue <- newCycleQueue+ job_queue <- newCycleQueue+ conn <- newConnection setts host port+ conn_ref <- newIORef conn+ var <- newTVarIO $ emptyState setts gen+ nxt_sub <- newTVarIO (atomically . writeCycleQueue queue)+ disposed <- newEmptyTMVarIO+ let env = Env setts queue pkg_queue job_queue var nxt_sub conn_ref disposed+ handler res = do+ closing env+ case res of+ Left e -> do+ _settingsLog setts (Error $ UnexpectedException e)+ case fromException e of+ Just (_ :: ConnectionException) -> atomically $ do+ writeTVar nxt_sub (raiseException e)+ putTMVar disposed ()+ _ -> do new_conn <- newConnection setts host port+ writeIORef conn_ref new_conn+ _ <- forkFinally (bootstrap env) handler+ return ()++ _ -> atomically $ putTMVar disposed ()+ _ <- forkFinally (bootstrap env) handler+ return $ Prod nxt_sub $ do+ closed <- connIsClosed conn+ unless closed retry+ readTMVar disposed
+ Database/EventStore/Internal/Generator.hs view
@@ -0,0 +1,42 @@+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Generator+-- Copyright : (C) 2015 Yorick Laupa+-- License : (see the file LICENSE)+--+-- Maintainer : Yorick Laupa <yo.eight@gmail.com>+-- Stability : provisional+-- Portability : non-portable+--+-- Pure UUID generator.+--------------------------------------------------------------------------------+module Database.EventStore.Internal.Generator+ ( Generator+ , nextUUID+ , newGenerator+ , splitGenerator+ ) where++--------------------------------------------------------------------------------+import Data.UUID+import System.Random++--------------------------------------------------------------------------------+-- | Pure 'UUID' generator.+newtype Generator = Generator StdGen++--------------------------------------------------------------------------------+-- | Gets the next fresh 'UUID'.+nextUUID :: Generator -> (UUID, Generator)+nextUUID (Generator g) = let (u, nxt) = random g in (u, Generator nxt)++--------------------------------------------------------------------------------+-- | Builds 2 new 'Generator's out of one.+splitGenerator :: Generator -> (Generator, Generator)+splitGenerator (Generator g) =+ let (g1, g2) = split g in (Generator g1, Generator g2)++--------------------------------------------------------------------------------+-- | Creates a new 'Generator'.+newGenerator :: IO Generator+newGenerator = fmap Generator getStdGen
− Database/EventStore/Internal/Manager/Operation.hs
@@ -1,197 +0,0 @@-{-# LANGUAGE RecordWildCards #-}-{-# LANGUAGE ExistentialQuantification #-}------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Manager.Operation--- 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.Operation- ( Decision(..)- , OperationParams(..)- , operationNetwork- ) where-----------------------------------------------------------------------------------import qualified Data.Map.Strict as M-import Data.Monoid ((<>))-import Data.Word-----------------------------------------------------------------------------------import Data.ProtocolBuffers-import Data.Serialize-import Data.UUID-import FRP.Sodium-import System.Random-----------------------------------------------------------------------------------import Database.EventStore.Internal.Types hiding (Event, newEvent)-import Database.EventStore.Internal.Util.Sodium-----------------------------------------------------------------------------------newtype Manager = Manager (M.Map UUID Operation)-----------------------------------------------------------------------------------initManager :: Manager-initManager = Manager M.empty------------------------------------------------------------------------------------- Operation----------------------------------------------------------------------------------data Decision- = DoNothing- | EndOperation- | Retry- | Reconnection- | Subscribed-----------------------------------------------------------------------------------data Operation- = Operation- { operationCreatePackage :: UUID -> IO Package- , operationInspect :: Package -> IO Decision- }-----------------------------------------------------------------------------------data OperationParams- = forall req resp. (Encode req, Decode resp) =>- OperationParams- { opSettings :: !Settings- , opRequestCmd :: !Word8- , opResponseCmd :: !Word8-- , opRequest :: IO req- , opSuccess :: resp -> IO Decision- , opFailure :: OperationException -> IO Decision- }-----------------------------------------------------------------------------------createOperation :: Settings -> OperationParams -> Operation-createOperation sett params =- Operation- { operationCreatePackage = createPackage sett params- , operationInspect = inspection params- }-----------------------------------------------------------------------------------createPackage :: Settings -> OperationParams -> UUID -> IO Package-createPackage Settings{..} OperationParams{..} uuid = do- req <- opRequest-- let pack = Package- { packageCmd = opRequestCmd- , packageCorrelation = uuid- , packageData = runPut $ encodeMessage req- , packageCred = s_credentials- }-- return pack-----------------------------------------------------------------------------------inspection :: OperationParams -> Package -> IO Decision-inspection params@OperationParams{..} pack- | found == exp_v = deeperInspection params pack- | otherwise = failed (InvalidServerResponse exp_v found)- where- exp_v = opResponseCmd- failed = opFailure- found = packageCmd pack-----------------------------------------------------------------------------------deeperInspection :: OperationParams -> Package -> IO Decision-deeperInspection OperationParams{..} pack =- case runGet decodeMessage bytes of- Left e -> failed (ProtobufDecodingError e)- Right msg -> succeed msg- where- failed = opFailure- succeed = opSuccess- bytes = packageData pack------------------------------------------------------------------------------------- Event----------------------------------------------------------------------------------data Register = Register UUID Operation--newtype Remove = Remove UUID--data Response = Response !Package !Operation-----------------------------------------------------------------------------------operationNetwork :: Settings- -> (Package -> Reactive ())- -> Reactive ()- -> Event Package- -> Reactive (OperationParams -> Reactive ())-operationNetwork sett push_pkg push_reco e_pkg = do- (on_new, push_new) <- newEvent- (on_reg, push_reg) <- newEvent- (on_rem, push_rem) <- newEvent- (on_ret, push_ret) <- newEvent-- let mgr_e = fmap register on_reg <>- fmap remove on_rem-- mgr_b <- accum initManager mgr_e-- let resp_e = filterJust $ snapshot response e_pkg mgr_b-- on_new_op = fmap (createOperation sett) on_new <> on_ret-- push_reg_io = pushAsync2 $ \uuid op -> push_reg $ Register uuid op- push_rem_io = pushAsync (push_rem . Remove)- push_retry_io = pushAsync2 $ \uuid op -> do- push_rem $ Remove uuid- push_ret op- push_reco_io = pushAsync2 $ \uuid op -> do- push_reco- push_rem $ Remove uuid- push_ret op- push_send_io = pushAsync push_pkg-- _ <- listen on_new_op $ \op -> do- uuid <- randomIO- push_reg_io uuid op-- _ <- listen resp_e $ \(Response pkg op) -> do- decision <- operationInspect op pkg- let corr_id = packageCorrelation pkg-- case decision of- DoNothing -> return ()- EndOperation -> push_rem_io corr_id- Retry -> push_retry_io corr_id op- Reconnection -> push_reco_io corr_id op- _ -> fail unexpectedDecision-- _ <- listen on_reg $ \(Register uuid op) ->- operationCreatePackage op uuid >>= push_send_io-- return push_new-----------------------------------------------------------------------------------unexpectedDecision :: String-unexpectedDecision = "Unexpected decision Processor.handlingOperation"------------------------------------------------------------------------------------- Model----------------------------------------------------------------------------------register :: Register -> Manager -> Manager-register (Register uuid op) (Manager m) = Manager $ M.insert uuid op m-----------------------------------------------------------------------------------remove :: Remove -> Manager -> Manager-remove (Remove uuid) (Manager m) = Manager $ M.delete uuid m------------------------------------------------------------------------------------- Snapshot----------------------------------------------------------------------------------response :: Package -> Manager -> Maybe Response-response pkg (Manager m) = fmap (Response pkg) $ M.lookup corr_id m- where- corr_id = packageCorrelation pkg
+ Database/EventStore/Internal/Manager/Operation/Model.hs view
@@ -0,0 +1,178 @@+{-# LANGUAGE ExistentialQuantification #-}+{-# LANGUAGE RecordWildCards #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Manager.Operation.Model+-- Copyright : (C) 2015 Yorick Laupa+-- License : (see the file LICENSE)+--+-- Maintainer : Yorick Laupa <yo.eight@gmail.com>+-- Stability : provisional+-- Portability : non-portable+--+-- Main operation bookkeeping structure.+--------------------------------------------------------------------------------+module Database.EventStore.Internal.Manager.Operation.Model+ ( Model+ , Transition(..)+ , newModel+ , pushOperation+ , submitPackage+ , abort+ ) where++--------------------------------------------------------------------------------+import Data.Word++--------------------------------------------------------------------------------+import qualified Data.HashMap.Strict as H+import Data.ProtocolBuffers+import Data.Serialize+import Data.UUID++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Generator+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Entry of a running 'Operation'.+data Elem r =+ forall a resp. Decode resp =>+ Elem+ { _opOp :: Operation a+ , _opCmd :: Word8+ , _opCont :: resp -> SM a ()+ , _opCb :: Either OperationError a -> r+ }++--------------------------------------------------------------------------------+-- | Operation internal state.+data State r =+ State+ { _gen :: Generator+ -- ^ 'UUID' generator.+ , _pending :: H.HashMap UUID (Elem r)+ -- ^ Contains all running 'Operation's.+ }++--------------------------------------------------------------------------------+initState :: Generator -> State r+initState g = State g H.empty++--------------------------------------------------------------------------------+-- | Type of requests handled by the model.+data Request r+ = forall a. New (Operation a) (Either OperationError a -> r)+ -- ^ Register a new 'Operation'.+ | Pkg Package+ -- ^ Submit a package.+ | Abort+ -- ^ Aborts every pending operation.++--------------------------------------------------------------------------------+-- | Output produces by the interpretation of an 'Operation'.+data Transition r+ = Produce r (Transition r)+ -- ^ Produces an intermediary value.+ | Transmit Package (Transition r)+ -- ^ Asks for sending the given 'Package'.+ | Await (Model r)+ -- ^ waits for more input.++--------------------------------------------------------------------------------+-- | Main 'Operation' bookkeeping state machine.+newtype Model r = Model (Request r -> Maybe (Transition r))++--------------------------------------------------------------------------------+-- | Pushes a new 'Operation' to model. The given 'Operation' state-machine is+-- initialized and produces a 'Package'.+pushOperation :: (Either OperationError a -> r)+ -> Operation a+ -> Model r+ -> Transition r+pushOperation cb op (Model k) = let Just t = k (New op cb) in t++--------------------------------------------------------------------------------+-- | Submits a 'Package' to the model. If the model isn't concerned by the+-- 'Package', it will returns 'Nothing'. Because 'Operation' can implement+-- complex logic (retry for instance), it returns a 'Step'.+submitPackage :: Package -> Model r -> Maybe (Transition r)+submitPackage pkg (Model k) = k (Pkg pkg)++--------------------------------------------------------------------------------+-- | Aborts every pending operation.+abort :: Model r -> Transition r+abort (Model k) = let Just t = k Abort in t++--------------------------------------------------------------------------------+runOperation :: Settings+ -> (Either OperationError a -> r)+ -> Operation a+ -> SM a ()+ -> State r+ -> Transition r+runOperation setts cb op start init_st = go init_st start+ where+ go st (Return _) = Await $ Model $ handle setts st+ go st (Yield a n) = Produce (cb $ Right a) (go st n)+ go st (FreshId k) =+ let (new_id, nxt_gen) = nextUUID $ _gen st+ nxt_st = st { _gen = nxt_gen } in+ go nxt_st $ k new_id+ go st (SendPkg ci co rq k) =+ let (new_uuid, nxt_gen) = nextUUID $ _gen st+ pkg = Package+ { packageCmd = ci+ , packageCorrelation = new_uuid+ , packageData = runPut $ encodeMessage rq+ , packageCred = s_credentials setts+ }+ elm = Elem op co k cb+ ps = H.insert new_uuid elm $ _pending st+ nxt_st = st { _pending = ps+ , _gen = nxt_gen+ } in+ Transmit pkg (Await $ Model $ handle setts nxt_st)+ go st (Failure m) =+ case m of+ Just e -> Produce (cb $ Left e) (Await $ Model $ handle setts st)+ _ -> runOperation setts cb op op st++--------------------------------------------------------------------------------+runPackage :: Settings -> State r -> Package -> Maybe (Transition r)+runPackage setts st Package{..} = do+ Elem op resp_cmd cont cb <- H.lookup packageCorrelation $ _pending st+ let nxt_ps = H.delete packageCorrelation $ _pending st+ nxt_st = st { _pending = nxt_ps }+ if resp_cmd /= packageCmd+ then+ let r = cb $ Left $ InvalidServerResponse resp_cmd packageCmd in+ return $ Produce r (Await $ Model $ handle setts nxt_st)+ else+ case runGet decodeMessage packageData of+ Left e ->+ let r = cb $ Left $ ProtobufDecodingError e in+ return $ Produce r (Await $ Model $ handle setts nxt_st)+ Right m -> return $ runOperation setts cb op (cont m) nxt_st++--------------------------------------------------------------------------------+abortOperations :: Settings -> State r -> Transition r+abortOperations setts init_st = go init_st $ H.toList $ _pending init_st+ where+ go st ((key, Elem _ _ _ k):xs) =+ let ps = H.delete key $ _pending st+ nxt_st = st { _pending = ps } in+ Produce (k $ Left Aborted) $ go nxt_st xs+ go st [] = Await $ Model $ handle setts st++--------------------------------------------------------------------------------+-- | Creates a new 'Operation' model state-machine.+newModel :: Settings -> Generator -> Model r+newModel setts g = Model $ handle setts $ initState g++--------------------------------------------------------------------------------+handle :: Settings -> State r -> Request r -> Maybe (Transition r)+handle setts st (New op cb) = Just $ runOperation setts cb op op st+handle setts st (Pkg pkg) = runPackage setts st pkg+handle setts st Abort = Just $ abortOperations setts st
Database/EventStore/Internal/Manager/Subscription.hs view
@@ -1,1112 +1,220 @@-{-# 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+{-# LANGUAGE GADTs #-}+{-# LANGUAGE Rank2Types #-}+{-# LANGUAGE ScopedTypeVariables #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Manager.Subscription+-- Copyright : (C) 2015 Yorick Laupa+-- License : (see the file LICENSE)+--+-- Maintainer : Yorick Laupa <yo.eight@gmail.com>+-- Stability : provisional+-- Portability : non-portable+--+-- Main subscription state machine declaration module. It also declares every+-- functions required to drive a 'Subscription'.+--------------------------------------------------------------------------------+module Database.EventStore.Internal.Manager.Subscription+ ( module Database.EventStore.Internal.Manager.Subscription.Driver+ , Regular(..)+ , Persistent(..)+ , Catchup(..)+ , Subscription+ , Running(..)+ , Checkpoint(..)+ , regularSubscription+ , catchupSubscription+ , persistentSubscription+ , eventArrived+ , readNext+ , batchRead+ , hasCaughtUp+ , runningUUID+ , runningLastEventNumber+ , runningLastCommitPosition+ ) where++--------------------------------------------------------------------------------+import Data.Int++--------------------------------------------------------------------------------+import Data.Sequence+import Data.Text (Text)++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Manager.Subscription.Driver+import Database.EventStore.Internal.Manager.Subscription.Model+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Also referred as volatile subscription. For example, if a stream has 100+-- events in it when a subscriber connects, the subscriber can expect to see+-- event number 101 onwards until the time the subscription is closed or+-- dropped.+data Regular = Regular { _subTos :: Bool }++--------------------------------------------------------------------------------+-- | This kind of subscription specifies a starting point, in the form of an+-- event number or transaction file position. The given function will be+-- called for events from the starting point until the end of the stream, and+-- then for subsequently written events.+--+-- For example, if a starting point of 50 is specified when a stream has 100+-- events in it, the subscriber can expect to see events 51 through 100, and+-- then any events subsequently written until such time as the subscription is+-- dropped or closed.+data Catchup = Catchup++--------------------------------------------------------------------------------+-- | The server remembers the state of the subscription. This allows for many+-- different modes of operations compared to a regular or catchup subscription+-- where the client holds the subscription state.+-- (Need EventStore >= v3.1.0).+data Persistent = Persistent { _perGroup :: Text }++--------------------------------------------------------------------------------+-- | Represents the different type of inputs a subscription state-machine can+-- handle.+data Input t a where+ -- A event has written to the stream. Subscription state machine should+ -- store that event withing its state.+ Arrived :: ResolvedEvent -> Input t (Subscription t)+ -- The user asks for the next event coming from the server.+ ReadNext :: Input t (Maybe ResolvedEvent, Subscription t)+ -- A batch read has been made. It's only use for 'Catchup' subscription+ -- type. It gives the list of read events and indicates if it reaches the+ -- end of the stream along with the next checkpoint to point at.+ BatchRead :: [ResolvedEvent]+ -> Bool+ -> Checkpoint+ -> Input Catchup (Subscription Catchup)+ -- Used only for 'Catchup' subscription type. Asks if the subscription+ -- read every events up to the checkpoint given by the user.+ CaughtUp :: Input Catchup Bool++--------------------------------------------------------------------------------+-- | Main subscription state machine.+newtype Subscription t = Subscription (forall a. Input t a -> a)++--------------------------------------------------------------------------------+-- | Submit a new event to the subscription state machine. Internally,+-- that event should be stored into the subscription buffer.+eventArrived :: ResolvedEvent -> Subscription t -> Subscription t+eventArrived e (Subscription k) = k (Arrived e)++--------------------------------------------------------------------------------+-- | Reads the next available event. Returns 'Nothing' it there is any. When+-- returning an event, it will be removed from the subscription buffer.+readNext :: Subscription t -> (Maybe ResolvedEvent, Subscription t)+readNext (Subscription k) = k ReadNext++--------------------------------------------------------------------------------+-- | Submits a list of events read from a stream. It's only used by a 'Catchup'+-- subscription.+batchRead :: [ResolvedEvent]+ -> Bool -- ^ If it reaches the end of the stream.+ -> Checkpoint+ -> Subscription Catchup+ -> Subscription Catchup+batchRead es eos nxt (Subscription k) = k (BatchRead es eos nxt)++--------------------------------------------------------------------------------+-- | Indicates if the subscription caught up the end of the stream, meaning the+-- subscription is actually live. Only used by 'Catchup' subscription.+hasCaughtUp :: Subscription Catchup -> Bool+hasCaughtUp (Subscription k) = k CaughtUp++--------------------------------------------------------------------------------+-- | Main 'Regular' subscription state machine.+regularSubscription :: Subscription Regular+regularSubscription = baseSubscription++--------------------------------------------------------------------------------+-- | Main 'Persistent' subscription state machine.+persistentSubscription :: Subscription Persistent+persistentSubscription = baseSubscription++--------------------------------------------------------------------------------+-- | Represents the next checkpoint to reach on a catchup subscription. Wheither+-- it's a regular stream or the $all stream, it either point to an 'Int32' or+-- a 'Position'.+data Checkpoint = CheckpointNumber Int32 | CheckpointPosition Position++--------------------------------------------------------------------------------+-- | Depending either if the subscription concerns a regular stream or $all,+-- indicates if an event number (or 'Position') is lesser that the current the+-- given 'CheckPoint'.+beforeChk :: Checkpoint -> ResolvedEvent -> Bool+beforeChk (CheckpointNumber num) re =+ recordedEventNumber (resolvedEventOriginal re) < num+beforeChk (CheckpointPosition pos) re =+ maybe False (< pos) $ resolvedEventPosition re++--------------------------------------------------------------------------------+-- | That subscription state machine accumulates events coming from batch read+-- and any real time change made on a stream. That state machine will not+-- served any recent change made on the stream until it reaches the end of the+-- stream. On every batch read, it makes sure events contained in that batch+-- are deleted from the subscription buffer in order to avoid duplicates. That+-- implemention has been chosen to avoid potential message lost between the+-- moment with reach the end of the stream and the delay required by asking+-- for a subscription.+catchupSubscription :: Subscription Catchup+catchupSubscription = Subscription $ catchingUp empty empty+ where+ catchingUp :: forall a. Seq ResolvedEvent+ -> Seq ResolvedEvent+ -> Input Catchup a+ -> a+ catchingUp b s (Arrived e) =+ Subscription $ catchingUp b (s |> e)+ catchingUp b s ReadNext =+ case viewl b of+ EmptyL -> (Nothing, Subscription $ catchingUp b s)+ e :< rest -> (Just e, Subscription $ catchingUp rest s)+ catchingUp b s (BatchRead es eos nxt_pt) =+ let nxt_b = foldl (|>) b es+ nxt_s = dropWhileL (beforeChk nxt_pt) s+ nxt = if eos+ then Subscription $ caughtUp nxt_b nxt_s+ else Subscription $ catchingUp nxt_b nxt_s in+ nxt+ catchingUp _ _ CaughtUp = False++ caughtUp :: forall a. Seq ResolvedEvent+ -> Seq ResolvedEvent+ -> Input Catchup a+ -> a+ caughtUp b s (Arrived e) = Subscription $ caughtUp b (s |> e)+ caughtUp b s ReadNext =+ case viewl b of+ EmptyL -> live s ReadNext+ e :< rest ->+ case viewl rest of+ EmptyL -> (Just e, Subscription $ live s)+ _ -> (Just e, Subscription $ caughtUp rest s)+ caughtUp b s (BatchRead _ _ _) = Subscription $ caughtUp b s+ caughtUp _ _ CaughtUp = False++ live :: forall a. Seq ResolvedEvent -> Input Catchup a -> a+ live s (Arrived e) = Subscription $ live (s |> e)+ live s ReadNext =+ case viewl s of+ EmptyL -> (Nothing, Subscription $ live s)+ e :< rest -> (Just e, Subscription $ live rest)+ live s (BatchRead _ _ _) = Subscription $ live s+ live _ CaughtUp = True++--------------------------------------------------------------------------------+-- | Base subscription used for 'Regular' or 'Persistent' subscription.+baseSubscription :: forall t. Subscription t+baseSubscription = Subscription $ go empty+ where+ go :: forall a. Seq ResolvedEvent -> Input t a -> a+ go s (Arrived e) = Subscription $ go (s |> e)+ go s ReadNext =+ case viewl s of+ EmptyL -> (Nothing, Subscription $ go s)+ e :< rest -> (Just e, Subscription $ go rest)+ go _ BatchRead{} = error "impossible: base subscription"+ go _ CaughtUp = error "impossible: base subscription"
+ Database/EventStore/Internal/Manager/Subscription/Driver.hs view
@@ -0,0 +1,481 @@+{-# LANGUAGE DeriveDataTypeable #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE Rank2Types #-}+{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE ScopedTypeVariables #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Manager.Subscription.Driver+-- Copyright : (C) 2015 Yorick Laupa+-- License : (see the file LICENSE)+--+-- Maintainer : Yorick Laupa <yo.eight@gmail.com>+-- Stability : provisional+-- Portability : non-portable+--+-- Subscription model driver. It drivers the model accordingly depending on the+-- 'Package' or commands submitted to it.+--------------------------------------------------------------------------------+module Database.EventStore.Internal.Manager.Subscription.Driver+ ( SubDropReason(..)+ , SubConnectEvent(..)+ , PersistActionException(..)+ , ConfirmedAction(..)+ , NakAction(..)+ , Driver+ , newDriver+ , submitPackage+ , connectToStream+ , connectToPersist+ , createPersist+ , updatePersist+ , deletePersist+ , ackPersist+ , nakPersist+ , unsubscribe+ , abort+ ) where++--------------------------------------------------------------------------------+import Control.Exception+import Data.Int+import Data.Maybe+import Data.Typeable++--------------------------------------------------------------------------------+import Data.ByteString+import qualified Data.HashMap.Strict as H+import Data.Serialize+import Data.ProtocolBuffers+import Data.Text+import Data.UUID++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Generator+import Database.EventStore.Internal.Manager.Subscription.Message+import Database.EventStore.Internal.Manager.Subscription.Model+import Database.EventStore.Internal.Manager.Subscription.Packages+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Set of events that can occurs during a subscription lifetime.+data SubConnectEvent+ = EventAppeared ResolvedEvent+ -- ^ A wild event appeared !+ | Dropped SubDropReason+ -- ^ The subscription connection dropped.+ | SubConfirmed Running+ -- ^ Subscription connection is confirmed. It means that subscription can+ -- receive events from the server.++--------------------------------------------------------------------------------+-- | Indicates why a subscription has been dropped.+data SubDropReason+ = SubUnsubscribed+ -- ^ Subscription connection has been closed by the user.+ | SubAccessDenied+ -- ^ The current user is not allowed to operate on the supplied stream.+ | SubNotFound+ -- ^ Given stream name doesn't exist.+ | SubPersistDeleted+ -- ^ Given stream is deleted.+ | SubAborted+ -- ^ Occurs when the user shutdown the connection from the server or if+ -- the connection to the server is no longer possible.+ deriving (Show, Eq)++--------------------------------------------------------------------------------+-- | Enumerates all persistent action exceptions.+data PersistActionException+ = PersistActionFail+ -- ^ The action failed.+ | PersistActionAlreadyExist+ -- ^ Happens when creating a persistent subscription on a stream with a+ -- group name already taken.+ | PersistActionDoesNotExist+ -- ^ An operation tried to do something on a persistent subscription or a+ -- stream that don't exist.+ | PersistActionAccessDenied+ -- ^ The current user is not allowed to operate on the supplied stream or+ -- persistent subscription.+ | PersistActionAborted+ -- ^ That action has been aborted because the user shutdown the connection+ -- to the server or the connection to the server is no longer possible.+ deriving (Show, Typeable)++--------------------------------------------------------------------------------+instance Exception PersistActionException++--------------------------------------------------------------------------------+-- | Emitted when a persistent action has been carried out successfully.+data ConfirmedAction =+ ConfirmedAction+ { caId :: !UUID+ -- ^ Action id.+ , caGroup :: !Text+ -- ^ Subscription group name.+ , caStream :: !Text+ -- ^ Stream name.+ , caAction :: !PersistAction+ -- ^ Persistent action type.+ }++--------------------------------------------------------------------------------+-- | Submits a 'Package' to a subscription driver. If the 'Package' was+-- processed by the driver, it will return a final value and a new driver with+-- its internal state updated accordingly.+submitPackage :: Package -> Driver r -> Maybe (r, Driver r)+submitPackage pkg (Driver k) = k (Pkg pkg)++--------------------------------------------------------------------------------+-- | Starts a regular subscription connection. It returns the associated+-- 'Package' and updates driver internal state.+connectToStream :: (SubConnectEvent -> r)+ -> Text -- ^ Stream name.+ -> Bool -- ^ Resolve Link TOS+ -> Driver r+ -> (Package, Driver r)+connectToStream c s t (Driver k) = k (Cmd $ ConnectReg c s t)++--------------------------------------------------------------------------------+-- | Starts a persistent subscription connection. It returns the associated+-- 'Package' and updates driver internal state.+connectToPersist :: (SubConnectEvent -> r)+ -> Text -- ^ Group name.+ -> Text -- ^ Stream name.+ -> Int32 -- ^ Buffer size.+ -> Driver r+ -> (Package, Driver r)+connectToPersist c g s b (Driver k) = k (Cmd $ ConnectPersist c g s b)++--------------------------------------------------------------------------------+-- | Creates a persistent subscription. It returns the associated 'Package' and+-- updates driver internal state.+createPersist :: (Either PersistActionException ConfirmedAction -> r)+ -> Text -- ^ Group name.+ -> Text -- ^ Stream name.+ -> PersistentSubscriptionSettings+ -> Driver r+ -> (Package, Driver r)+createPersist c g s ss (Driver k) =+ k (Cmd $ ApplyPersistAction c g s (PersistCreate ss))++--------------------------------------------------------------------------------+-- | Updates a persistent subscription. It returns the associated 'Package' and+-- updates driver internal state.+updatePersist :: (Either PersistActionException ConfirmedAction -> r)+ -> Text -- ^ Group name.+ -> Text -- ^ Stream name.+ -> PersistentSubscriptionSettings+ -> Driver r+ -> (Package, Driver r)+updatePersist c g s ss (Driver k) =+ k (Cmd $ ApplyPersistAction c g s (PersistUpdate ss))++--------------------------------------------------------------------------------+-- | Deletes a persistent subscription. It returns the associated 'Package' and+-- updates driver internal state.+deletePersist :: (Either PersistActionException ConfirmedAction -> r)+ -> Text -- ^ Group name.+ -> Text -- ^ Stream name.+ -> Driver r+ -> (Package, Driver r)+deletePersist c g s (Driver k) =+ k (Cmd $ ApplyPersistAction c g s PersistDelete)++--------------------------------------------------------------------------------+-- | Given a persistent subscription, acknowledges a set of events have been+-- successfully processed. It returns the associated 'Package' and updates+-- driver internal state.+ackPersist :: r+ -> Running+ -> [UUID] -- ^ Event ids.+ -> Driver r+ -> (Package, Driver r)+ackPersist r i evts (Driver k) = k (Cmd $ PersistAck r i evts)++--------------------------------------------------------------------------------+-- | Given a persistent subscription, indicates a set of event haven't been+-- processed correctly. It returns the associated 'Package' and updates driver+-- internal state.+nakPersist :: r+ -> Running+ -> NakAction+ -> Maybe Text -- ^ Reason.+ -> [UUID] -- ^ Event ids.+ -> Driver r+ -> (Package, Driver r)+nakPersist r i na mt evts (Driver k) =+ k (Cmd $ PersistNak r i na mt evts)++--------------------------------------------------------------------------------+-- | Unsubscribe from a subscription.+unsubscribe :: Running -> Driver r -> (Package, Driver r)+unsubscribe r (Driver k) = k (Cmd $ Unsubscribe r)++--------------------------------------------------------------------------------+-- | Aborts every pending action.+abort :: Driver r -> [r]+abort (Driver k) = k Abort++--------------------------------------------------------------------------------+-- EventStore result mappers:+-- =========================+-- EventStore protocol has several values that means the exact same thing. Those+-- functions convert a specific EventStore to uniform result type common to all+-- persistent actions.+--------------------------------------------------------------------------------+createRException :: CreatePersistentSubscriptionResult+ -> Maybe PersistActionException+createRException CPS_Success = Nothing+createRException CPS_AlreadyExists = Just PersistActionAlreadyExist+createRException CPS_Fail = Just PersistActionFail+createRException CPS_AccessDenied = Just PersistActionAccessDenied++--------------------------------------------------------------------------------+deleteRException :: DeletePersistentSubscriptionResult+ -> Maybe PersistActionException+deleteRException DPS_Success = Nothing+deleteRException DPS_DoesNotExist = Just PersistActionDoesNotExist+deleteRException DPS_Fail = Just PersistActionFail+deleteRException DPS_AccessDenied = Just PersistActionAccessDenied++--------------------------------------------------------------------------------+updateRException :: UpdatePersistentSubscriptionResult+ -> Maybe PersistActionException+updateRException UPS_Success = Nothing+updateRException UPS_DoesNotExist = Just PersistActionDoesNotExist+updateRException UPS_Fail = Just PersistActionFail+updateRException UPS_AccessDenied = Just PersistActionAccessDenied++--------------------------------------------------------------------------------+toSubDropReason :: DropReason -> SubDropReason+toSubDropReason D_Unsubscribed = SubUnsubscribed+toSubDropReason D_NotFound = SubNotFound+toSubDropReason D_AccessDenied = SubAccessDenied+toSubDropReason D_PersistentSubscriptionDeleted = SubPersistDeleted++--------------------------------------------------------------------------------+--------------------------------------------------------------------------------+-- | Type of inputs handled by the 'Subscription' driver.+data In r a where+ -- A command consists of receiving some parameters, updating the+ -- 'Subscription' model accordingly and thus modifying driver internal+ -- state.+ Cmd :: Cmd r -> In r (Package, Driver r)+ -- A 'Package' has been submitted to the 'Subscription' driver. If the+ -- driver recognize that 'Package', it returns a final value and update+ -- the driver internal state.+ Pkg :: Package -> In r (Maybe (r, Driver r))+ -- Aborts every pending action.+ Abort :: In r [r]++--------------------------------------------------------------------------------+-- | Set of commands handled by the driver.+data Cmd r+ = ConnectReg (SubConnectEvent -> r) Text Bool+ -- ^ Creates a regular 'Subscription' connection. When a 'SubConnectEvent'+ -- has arrived, the driver will use the provided callback and emit a+ -- final value. It holds a stream name and `Resolve Link TOS` setting.+ | ConnectPersist (SubConnectEvent -> r)+ Text+ Text+ Int32+ -- ^ Creates a persistent 'Subscription' connection. When a+ -- 'SubConnectEvent' has arrived, the driver will use the provided+ -- callback and emit a final value. It holds a group name, stream+ -- name and a buffer size.++ | Unsubscribe Running+ -- ^ Unsubscribe from a subscription.++ | ApplyPersistAction (Either PersistActionException ConfirmedAction -> r)+ Text+ Text+ PersistAction+ -- ^ Creates a persistent action. Depending of the failure or the success+ -- of that action, the driver will use the provided callback to emit a+ -- final value. It hols a group name, a stream name and a persistent+ -- action.++ | PersistAck r Running [UUID]+ -- ^ Acks a set of Event 'UUID' to notify those events have been correctly+ -- handled. It holds a 'Running' subscription and a set of `UUID`+ -- representing event id. When the ack would be confirmed the driver+ -- will return the supplied final value.+ | PersistNak r+ Running+ NakAction+ (Maybe Text)+ [UUID]+ -- ^ Naks a set of Event 'UUID' to notify those events haven't been+ -- handled correctly. it holds a 'Running' subscription, a 'NakAction',+ -- an optional reason and a set of event ids. When the nak would be+ -- confirmed, the driver will return the provided final value.++--------------------------------------------------------------------------------+-- | Driver internal state.+data State r =+ State+ { _model :: !Model+ -- ^ Subscription model.+ , _gen :: !Generator+ -- ^ 'UUID' generator.+ , _reg :: !(H.HashMap UUID (Cmd r))+ -- ^ Holds ongoing commands. When stored, it means an action hasn't been+ -- confirmed yet.+ }++--------------------------------------------------------------------------------+initState :: Generator -> State r+initState gen = State newModel gen H.empty++--------------------------------------------------------------------------------+-- | Subscription driver state machine.+newtype Driver r = Driver (forall a. In r a -> a)++--------------------------------------------------------------------------------+-- | Creates a new subscription 'Driver' state machine.+newDriver :: forall r. Settings -> Generator -> Driver r+newDriver setts gen = Driver $ go (initState gen)+ where+ go :: forall a. State r -> In r a -> a+ go st@State{..} (Pkg Package{..}) = do+ elm <- H.lookup packageCorrelation _reg+ case packageCmd of+ 0xC2 -> do+ _ <- querySubscription packageCorrelation _model+ msg <- maybeDecodeMessage packageData+ let e = getField $ streamResolvedEvent msg+ evt = newResolvedEventFromBuf e+ app = EventAppeared evt+ ConnectReg k _ _ = elm+ return (k app, Driver $ go st)++ 0xC7 -> do+ _ <- querySubscription packageCorrelation _model+ msg <- maybeDecodeMessage packageData+ let e = getField $ psseaEvt msg+ evt = newResolvedEvent e+ app = EventAppeared evt+ ConnectPersist k _ _ _ = elm+ return (k app, Driver $ go st)++ 0xC1 -> do+ msg <- maybeDecodeMessage packageData+ let lcp = getField $ subscribeLastCommitPos msg+ len = getField $ subscribeLastEventNumber msg+ meta = RegularMeta lcp len+ ConnectReg k _ _ = elm+ nxt_m = confirmedSubscription packageCorrelation meta _model+ run <- querySubscription packageCorrelation nxt_m+ let nxt_st = st { _model = nxt_m }+ evt = SubConfirmed run+ return (k evt, Driver $ go nxt_st)++ 0xC6 -> do+ msg <- maybeDecodeMessage packageData+ let lcp = getField $ pscLastCommitPos msg+ sid = getField $ pscId msg+ len = getField $ pscLastEvtNumber msg+ meta = PersistMeta sid lcp len+ ConnectPersist k _ _ _ = elm+ nxt_m = confirmedSubscription packageCorrelation meta _model+ run <- querySubscription packageCorrelation nxt_m+ let nxt_st = st { _model = nxt_m }+ evt = SubConfirmed run+ return (k evt, Driver $ go nxt_st)++ 0xC9 -> confirmPAction elm (getField . cpscResult) createRException+ 0xCF -> confirmPAction elm (getField . upscResult) updateRException+ 0xCB -> confirmPAction elm (getField . dpscResult) deleteRException++ 0xC4 -> do+ run <- querySubscription packageCorrelation _model+ msg <- maybeDecodeMessage packageData+ let reason = fromMaybe D_Unsubscribed $ getField+ $ dropReason msg+ nxt_m = unsubscribed run _model+ dreason = toSubDropReason reason+ evt = Dropped dreason+ nxt_reg = H.delete packageCorrelation _reg+ nxt_st = st { _model = nxt_m+ , _reg = nxt_reg }+ case elm of+ ConnectReg k _ _ -> return (k evt, Driver $ go nxt_st)+ ConnectPersist k _ _ _ -> return (k evt, Driver $ go nxt_st)+ _ -> Nothing++ _ -> Nothing+ where+ confirmPAction :: Decode m+ => Cmd r+ -> (m -> e)+ -> (e -> Maybe PersistActionException)+ -> Maybe (r, Driver r)+ confirmPAction (ApplyPersistAction k g n c) fd em = do+ msg <- maybeDecodeMessage packageData+ _ <- queryPersistentAction packageCorrelation _model+ let nxt_m = confirmedAction packageCorrelation _model+ nxt_rg = H.delete packageCorrelation _reg+ nxt_st = st { _model = nxt_m+ , _reg = nxt_rg+ }+ evt = ConfirmedAction packageCorrelation g n c+ case em $ fd msg of+ Just e -> return (k $ Left e, Driver $ go nxt_st)+ Nothing -> return (k $ Right evt, Driver $ go nxt_st)+ confirmPAction _ _ _ = Nothing++ go st@State{..} (Cmd cmd) =+ case cmd of+ ConnectReg _ s tos ->+ let (u, nxt_g) = nextUUID _gen+ pkg = createConnectRegularPackage setts u s tos+ nxt_m = connectReg s tos u _model+ nxt_st = st { _model = nxt_m+ , _gen = nxt_g+ , _reg = H.insert u cmd _reg } in+ (pkg, Driver $ go nxt_st)+ ConnectPersist _ gn n b ->+ let (u, nxt_g) = nextUUID _gen+ pkg = createConnectPersistPackage setts u gn n b+ nxt_m = connectPersist gn n b u _model+ nxt_st = st { _model = nxt_m+ , _gen = nxt_g+ , _reg = H.insert u cmd _reg } in+ (pkg, Driver $ go nxt_st)+ Unsubscribe r ->+ let pkg = createUnsubscribePackage setts $ runningUUID r in+ (pkg, Driver $ go st)+ ApplyPersistAction _ gn n a ->+ let (u, nxt_g) = nextUUID _gen+ pkg = createPersistActionPackage setts u gn n a+ nxt_m = persistAction gn n u a _model+ nxt_st = st { _model = nxt_m+ , _gen = nxt_g+ , _reg = H.insert u cmd _reg } in+ (pkg, Driver $ go nxt_st)+ PersistAck _ run evts ->+ let RunningPersist _ _ _ _ sid _ _ = run+ u = runningUUID run+ pkg = createAckPackage setts u sid evts in+ (pkg, Driver $ go st)+ PersistNak _ run na r evts ->+ let RunningPersist _ _ _ _ sid _ _ = run+ u = runningUUID run+ pkg = createNakPackage setts u sid na r evts in+ (pkg, Driver $ go st)+ go st Abort = (H.elems $ _reg st) >>= _F+ where+ _F (ConnectReg k _ _) = [k $ Dropped SubAborted]+ _F (ConnectPersist k _ _ _) = [k $ Dropped SubAborted]+ _F (ApplyPersistAction k _ _ _) = [k $ Left PersistActionAborted]+ _F _ = []++--------------------------------------------------------------------------------+maybeDecodeMessage :: Decode a => ByteString -> Maybe a+maybeDecodeMessage bytes =+ case runGet decodeMessage bytes of+ Right a -> Just a+ _ -> Nothing
+ Database/EventStore/Internal/Manager/Subscription/Message.hs view
@@ -0,0 +1,395 @@+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE DeriveGeneric #-}+{-# OPTIONS_GHC -fcontext-stack=26 #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Manager.Subscription.Message+-- 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.Internal.Manager.Subscription.Message where++--------------------------------------------------------------------------------+import Data.ByteString (ByteString)+import Data.Int+import GHC.Generics (Generic)++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text (Text)++--------------------------------------------------------------------------------+import Database.EventStore.Internal.TimeSpan+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Stream subscription connection request.+data SubscribeToStream+ = SubscribeToStream+ { subscribeStreamId :: Required 1 (Value Text)+ , subscribeResolveLinkTos :: Required 2 (Value Bool)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode SubscribeToStream++--------------------------------------------------------------------------------+-- | 'SubscribeToStream' smart constructor.+subscribeToStream :: Text -> Bool -> SubscribeToStream+subscribeToStream stream_id res_link_tos =+ SubscribeToStream+ { subscribeStreamId = putField stream_id+ , subscribeResolveLinkTos = putField res_link_tos+ }++--------------------------------------------------------------------------------+-- | Stream subscription connection response.+data SubscriptionConfirmation+ = SubscriptionConfirmation+ { subscribeLastCommitPos :: Required 1 (Value Int64)+ , subscribeLastEventNumber :: Optional 2 (Value Int32)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode SubscriptionConfirmation++--------------------------------------------------------------------------------+-- | Serialized event sent by the server when a new event has been appended to a+-- stream.+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)++--------------------------------------------------------------------------------+-- | A message sent by the server when a subscription has been dropped.+data SubscriptionDropped+ = SubscriptionDropped+ { dropReason :: Optional 1 (Enumeration DropReason) }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode SubscriptionDropped++--------------------------------------------------------------------------------+-- | A message sent to the server to indicate the user asked to end a+-- subscription.+data UnsubscribeFromStream = UnsubscribeFromStream deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode UnsubscribeFromStream++--------------------------------------------------------------------------------+-- | Create persistent subscription request.+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' smart constructor.+_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++--------------------------------------------------------------------------------+-- | Create persistent subscription outcome.+data CreatePersistentSubscriptionResult+ = CPS_Success+ | CPS_AlreadyExists+ | CPS_Fail+ | CPS_AccessDenied+ deriving (Enum, Eq, Show)++--------------------------------------------------------------------------------+-- | Create persistent subscription response.+data CreatePersistentSubscriptionCompleted =+ CreatePersistentSubscriptionCompleted+ { cpscResult :: Required 1 (Enumeration CreatePersistentSubscriptionResult)+ , cpscReason :: Optional 2 (Value Text)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode CreatePersistentSubscriptionCompleted++--------------------------------------------------------------------------------+-- | Delete persistent subscription request.+data DeletePersistentSubscription =+ DeletePersistentSubscription+ { dpsGroupName :: Required 1 (Value Text)+ , dpsStreamId :: Required 2 (Value Text)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode DeletePersistentSubscription++--------------------------------------------------------------------------------+-- | 'DeletePersistentSubscription' smart construction.+_deletePersistentSubscription :: Text -> Text -> DeletePersistentSubscription+_deletePersistentSubscription group_name stream_id =+ DeletePersistentSubscription+ { dpsGroupName = putField group_name+ , dpsStreamId = putField stream_id+ }++--------------------------------------------------------------------------------+-- | Delete persistent subscription outcome.+data DeletePersistentSubscriptionResult+ = DPS_Success+ | DPS_DoesNotExist+ | DPS_Fail+ | DPS_AccessDenied+ deriving (Enum, Eq, Show)++--------------------------------------------------------------------------------+-- | Delete persistent subscription response.+data DeletePersistentSubscriptionCompleted =+ DeletePersistentSubscriptionCompleted+ { dpscResult :: Required 1 (Enumeration DeletePersistentSubscriptionResult)+ , dpscReason :: Optional 2 (Value Text)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode DeletePersistentSubscriptionCompleted++--------------------------------------------------------------------------------+-- | Update persistent subscription request.+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' smart constructor.+_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++--------------------------------------------------------------------------------+-- | Update persistent subscription outcome.+data UpdatePersistentSubscriptionResult+ = UPS_Success+ | UPS_DoesNotExist+ | UPS_Fail+ | UPS_AccessDenied+ deriving (Enum, Eq, Show)++--------------------------------------------------------------------------------+-- | Update persistent subscription response.+data UpdatePersistentSubscriptionCompleted =+ UpdatePersistentSubscriptionCompleted+ { upscResult :: Required 1 (Enumeration UpdatePersistentSubscriptionResult)+ , upscReason :: Optional 2 (Value Text)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode UpdatePersistentSubscriptionCompleted++--------------------------------------------------------------------------------+-- | Connect to a persistent subscription request.+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' smart constructor.+_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+ }++--------------------------------------------------------------------------------+-- | Ack processed events request.+data PersistentSubscriptionAckEvents =+ PersistentSubscriptionAckEvents+ { psaeId :: Required 1 (Value Text)+ , psaeProcessedEvtIds :: Repeated 2 (Value ByteString)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode PersistentSubscriptionAckEvents++--------------------------------------------------------------------------------+-- | 'PersistentSubscriptionAckEvents' smart constructor.+persistentSubscriptionAckEvents :: Text+ -> [ByteString]+ -> PersistentSubscriptionAckEvents+persistentSubscriptionAckEvents sub_id evt_ids =+ PersistentSubscriptionAckEvents+ { psaeId = putField sub_id+ , psaeProcessedEvtIds = putField evt_ids+ }++--------------------------------------------------------------------------------+-- | Gathers every possible Nak actions.+data NakAction+ = NA_Unknown+ | NA_Park+ | NA_Retry+ | NA_Skip+ | NA_Stop+ deriving (Enum, Eq, Show)++--------------------------------------------------------------------------------+-- | Nak processed events request.+data PersistentSubscriptionNakEvents =+ PersistentSubscriptionNakEvents+ { psneId :: Required 1 (Value Text)+ , psneProcessedEvtIds :: Repeated 2 (Value ByteString)+ , psneMsg :: Optional 3 (Value Text)+ , psneAction :: Required 4 (Enumeration NakAction)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode PersistentSubscriptionNakEvents++--------------------------------------------------------------------------------+-- | 'PersistentSubscriptionNakEvents' smart constructor.+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+ }++--------------------------------------------------------------------------------+-- | Connection to persistent subscription response.+data PersistentSubscriptionConfirmation =+ PersistentSubscriptionConfirmation+ { pscLastCommitPos :: Required 1 (Value Int64)+ , pscId :: Required 2 (Value Text)+ , pscLastEvtNumber :: Optional 3 (Value Int32)+ } deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode PersistentSubscriptionConfirmation++--------------------------------------------------------------------------------+-- | Avalaible event sent by the server in the context of a persistent+-- subscription..+data PersistentSubscriptionStreamEventAppeared =+ PersistentSubscriptionStreamEventAppeared+ { psseaEvt :: Required 1 (Message ResolvedIndexedEvent) }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode PersistentSubscriptionStreamEventAppeared
+ Database/EventStore/Internal/Manager/Subscription/Model.hs view
@@ -0,0 +1,318 @@+{-# LANGUAGE GADTs #-}+{-# LANGUAGE Rank2Types #-}+{-# LANGUAGE RecordWildCards #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Manager.Subscription.Model+-- Copyright : (C) 2015 Yorick Laupa+-- License : (see the file LICENSE)+--+-- Maintainer : Yorick Laupa <yo.eight@gmail.com>+-- Stability : provisional+-- Portability : non-portable+--+-- Main Subscription bookkeeping structure.+--------------------------------------------------------------------------------+module Database.EventStore.Internal.Manager.Subscription.Model+ ( PersistAction(..)+ , PendingAction(..)+ , Running(..)+ , Meta(..)+ , Model+ , runningUUID+ , runningLastEventNumber+ , runningLastCommitPosition+ , querySubscription+ , queryPersistentAction+ , confirmedSubscription+ , confirmedAction+ , newModel+ , unsubscribed+ , connectReg+ , connectPersist+ , persistAction+ ) where++--------------------------------------------------------------------------------+import Data.Int++--------------------------------------------------------------------------------+import qualified Data.HashMap.Strict as H+import Data.Text+import Data.UUID++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Type of persistent action.+data PersistAction+ = PersistCreate PersistentSubscriptionSettings+ | PersistUpdate PersistentSubscriptionSettings+ | PersistDelete++--------------------------------------------------------------------------------+-- | Represents an persistent action that hasn't been completed yet.+data PendingAction =+ PendingAction+ { _paGroup :: !Text+ , _paStream :: !Text+ , _paTpe :: !PersistAction+ }++--------------------------------------------------------------------------------+type Register a = H.HashMap UUID a++--------------------------------------------------------------------------------+-- | Represents a 'Subscription' which is about to be confirmed.+data Pending+ = PendingReg Text Bool+ -- ^ Related to regular subscription. In order of appearance:+ --+ -- * Stream name.+ --+ -- * Resolve Link TOS.+ | PendingPersist Text Text Int32+ -- ^ Related to persistent subscription. In order of appearance:+ --+ -- * Group name.+ --+ -- * Stream name.+ --+ -- * Buffer size.+ deriving Show++--------------------------------------------------------------------------------+-- | Represents a running subscription. Gathers useful information.+data Running+ = RunningReg UUID Text Bool Int64 (Maybe Int32)+ -- ^ Related regular subscription. In order of appearance:+ --+ -- * Subscription id.+ --+ -- * Stream name.+ --+ -- * Resolve Link TOS.+ --+ -- * Last commit position.+ --+ -- * Last event number.+ | RunningPersist UUID Text Text Int32 Text Int64 (Maybe Int32)+ -- ^ Related to persistent subscription. In order of appearance:+ --+ -- * Subscription id.+ --+ -- * Group name.+ --+ -- * Stream name.+ --+ -- * Buffer size.+ --+ -- * Persistence subscription id.+ --+ -- * Last commit position.+ --+ -- * Last event number.+ deriving Show++--------------------------------------------------------------------------------+-- | Gets the event number of a running subscription.+runningLastEventNumber :: Running -> Maybe Int32+runningLastEventNumber (RunningReg _ _ _ _ i) = i+runningLastEventNumber (RunningPersist _ _ _ _ _ _ i) = i++--------------------------------------------------------------------------------+-- | Gets the commit position of a running subscription.+runningLastCommitPosition :: Running -> Int64+runningLastCommitPosition (RunningReg _ _ _ i _) = i+runningLastCommitPosition (RunningPersist _ _ _ _ _ i _) = i++--------------------------------------------------------------------------------+-- | Gets the 'UUID' of a running subscription.+runningUUID :: Running -> UUID+runningUUID (RunningReg i _ _ _ _) = i+runningUUID (RunningPersist i _ _ _ _ _ _) = i++--------------------------------------------------------------------------------+-- | Type of requests handled by the model.+data Request a where+ -- Read request.+ Query :: Query a -> Request a+ -- Write request.+ Execute :: Action -> Request Model++--------------------------------------------------------------------------------+-- | Set of a piece of information we can query from the 'Subscription' model.+data Query a where+ -- Query a running 'Subscription'.+ QuerySub :: UUID -> Query (Maybe Running)+ -- Query a pending persistent action.+ QueryAction :: UUID -> Query (Maybe PendingAction)++--------------------------------------------------------------------------------+-- | Set of actions handled by the 'Subscription' model.+data Action+ = Connect UUID Connect+ -- ^ Subscription connection.+ | Confirmed Confirmed+ -- ^ Subscription action confirmation.+ | Unsubscribed UUID+ -- ^ Subscription no longer exist.+ | PersistAction Text Text UUID PersistAction+ -- ^ Add a new persist action.++--------------------------------------------------------------------------------+-- | Subscription connection information.+data Connect+ = ConnectReg Text Bool+ -- | |---- Resolve TOS link.+ -- |--------- Stream name.+ | ConnectPersist Text Text Int32+ -- | | |---- Buffer size.+ -- | |--------- Stream name.+ -- |-------------- Group name.++--------------------------------------------------------------------------------+-- | Information related to a confirmed 'Subscription'.+data Meta+ = RegularMeta Int64 (Maybe Int32)+ -- | |------------- Last commit position.+ -- |------------------- Last event number.+ | PersistMeta Text Int64 (Maybe Int32)+ -- | | |------------- Subscription Id.+ -- | |------------------- Last commit position.+ -- |------------------------ Last event number.++--------------------------------------------------------------------------------+-- | Subscription action confirmation.+data Confirmed+ = ConfirmedConnection UUID Meta+ -- ^ Confirms a 'Subscription' connection has handled successfully.+ | ConfirmedPersistAction UUID+ -- ^ Confirms a persist action has been handled successfully.++--------------------------------------------------------------------------------+-- | Retrieves a running 'Subscription'.+querySubscription :: UUID -> Model -> Maybe Running+querySubscription u (Model k) = k $ Query $ QuerySub u++--------------------------------------------------------------------------------+-- | Retrieves an ongoing persistent action.+queryPersistentAction :: UUID -> Model -> Maybe PendingAction+queryPersistentAction u (Model k) = k $ Query $ QueryAction u++--------------------------------------------------------------------------------+-- | Registers a regular 'Subscription' request.+connectReg :: Text -> Bool -> UUID -> Model -> Model+connectReg n t u (Model k) = k $ Execute $ Connect u (ConnectReg n t)++--------------------------------------------------------------------------------+-- | Registers a persistent 'Subscription' request.+connectPersist :: Text -> Text -> Int32 -> UUID -> Model -> Model+connectPersist g n b u (Model k) =+ k $ Execute $ Connect u (ConnectPersist g n b)++--------------------------------------------------------------------------------+-- | Registers a persistent action.+persistAction :: Text -> Text -> UUID -> PersistAction -> Model -> Model+persistAction g n u a (Model k) = k $ Execute $ PersistAction g n u a++--------------------------------------------------------------------------------+-- | Confirms a subscription.+confirmedSubscription :: UUID -> Meta -> Model -> Model+confirmedSubscription u m (Model k) =+ k $ Execute $ Confirmed $ ConfirmedConnection u m++--------------------------------------------------------------------------------+-- | Confirms a persistent action. It doesn't assume if the action went well.+confirmedAction :: UUID -> Model -> Model+confirmedAction u (Model k) = k $ Execute $ Confirmed $ ConfirmedPersistAction u++--------------------------------------------------------------------------------+-- | Remove a 'Subscription'.+unsubscribed :: Running -> Model -> Model+unsubscribed r (Model k) = k $ Execute $ Unsubscribed $ runningUUID r++--------------------------------------------------------------------------------+-- | 'Subscription' model internal state.+data State =+ State+ { _stPending :: !(Register Pending)+ -- ^ Holds all pending 'Subscription's+ , _stRunning :: !(Register Running)+ -- ^ Holds all 'Subscription's that are currently running.+ , _stAction :: !(Register PendingAction)+ -- ^ Holds all pending persistent actions.+ }++--------------------------------------------------------------------------------+emptyState :: State+emptyState = State H.empty H.empty H.empty++--------------------------------------------------------------------------------+-- | Subscription operations state machine. Keeps every information related to+-- subscription updated.+newtype Model = Model (forall a. Request a -> a)++--------------------------------------------------------------------------------+-- | Creates a new 'Subscription' model.+newModel :: Model+newModel = Model $ modelHandle emptyState++--------------------------------------------------------------------------------+-- | Main model handler.+modelHandle :: State -> Request a -> a+modelHandle s (Execute e) =+ case e of+ Connect u c ->+ case c of+ ConnectReg n tos ->+ let p = PendingReg n tos+ nxt_ps = H.insert u p $ _stPending s+ nxt_s = s { _stPending = nxt_ps } in+ Model $ modelHandle nxt_s+ ConnectPersist g n b ->+ let p = PendingPersist g n b+ nxt_ps = H.insert u p $ _stPending s+ nxt_s = s { _stPending = nxt_ps } in+ Model $ modelHandle nxt_s+ Confirmed c ->+ case c of+ ConfirmedConnection u tpe ->+ case tpe of+ RegularMeta lc le ->+ case H.lookup u $ _stPending s of+ Just (PendingReg n tos) ->+ let r = RunningReg u n tos lc le+ nxt_rs = H.insert u r $ _stRunning s+ nxt_s = s { _stRunning = nxt_rs } in+ Model $ modelHandle nxt_s+ _ -> Model $ modelHandle s+ PersistMeta sb lc le ->+ case H.lookup u $ _stPending s of+ Just (PendingPersist g n b) ->+ let r = RunningPersist u g n b sb lc le+ nxt_rs = H.insert u r $ _stRunning s+ nxt_s = s { _stRunning = nxt_rs } in+ Model $ modelHandle nxt_s+ _ -> Model $ modelHandle s+ ConfirmedPersistAction u ->+ case H.lookup u $ _stAction s of+ Just (PendingAction{}) ->+ let nxt_as = H.delete u $ _stAction s+ nxt_s = s { _stAction = nxt_as } in+ Model $ modelHandle nxt_s+ _ -> Model $ modelHandle s+ Unsubscribed u ->+ let nxt_ps = H.delete u $ _stRunning s+ nxt_s = s { _stRunning = nxt_ps } in+ Model $ modelHandle nxt_s+ PersistAction g n u t ->+ let a = PendingAction g n t+ nxt_as = H.insert u a $ _stAction s+ nxt_s = s { _stAction = nxt_as } in+ Model $ modelHandle nxt_s+modelHandle s (Query q) =+ case q of+ QuerySub u -> H.lookup u $ _stRunning s+ QueryAction u -> H.lookup u $ _stAction s
+ Database/EventStore/Internal/Manager/Subscription/Packages.hs view
@@ -0,0 +1,137 @@+{-# LANGUAGE RecordWildCards #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Manager.Subscription.Packages+-- 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.Internal.Manager.Subscription.Packages where++--------------------------------------------------------------------------------+import Data.Int++--------------------------------------------------------------------------------+import Data.ByteString.Lazy (toStrict)+import Data.ProtocolBuffers+import Data.Serialize+import Data.Text+import Data.UUID++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Manager.Subscription.Message+import Database.EventStore.Internal.Manager.Subscription.Model+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+import Prelude++--------------------------------------------------------------------------------+-- | Creates a regular subscription connection 'Package'.+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++--------------------------------------------------------------------------------+-- | Creates a persistent subscription connection 'Package'.+createConnectPersistPackage :: Settings+ -> UUID+ -> Text+ -> Text+ -> Int32+ -> Package+createConnectPersistPackage Settings{..} uuid grp stream bufSize =+ Package+ { packageCmd = 0xC5+ , packageCorrelation = uuid+ , packageData = runPut $ encodeMessage msg+ , packageCred = s_credentials+ }+ where+ msg = _connectToPersistentSubscription grp stream bufSize++--------------------------------------------------------------------------------+-- | Creates a persistent subscription 'Package'.+createPersistActionPackage :: Settings+ -> UUID+ -> Text+ -> Text+ -> PersistAction+ -> Package+createPersistActionPackage Settings{..} u grp strm tpe =+ Package+ { packageCmd = cmd+ , packageCorrelation = u+ , packageData = runPut msg+ , packageCred = s_credentials+ }+ where+ msg =+ case tpe of+ PersistCreate sett ->+ encodeMessage $ _createPersistentSubscription grp strm sett+ PersistUpdate sett ->+ encodeMessage $ _updatePersistentSubscription grp strm sett+ PersistDelete ->+ encodeMessage $ _deletePersistentSubscription grp strm+ cmd =+ case tpe of+ PersistCreate _ -> 0xC8+ PersistUpdate _ -> 0xCE+ PersistDelete -> 0xCA++--------------------------------------------------------------------------------+-- | Creates Ack 'Package'.+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 = fmap (toStrict . toByteString) eids+ msg = persistentSubscriptionAckEvents sid bytes++--------------------------------------------------------------------------------+-- | Create Nak 'Package'.+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 = fmap (toStrict . toByteString) eids+ msg = persistentSubscriptionNakEvents sid bytes txt act++--------------------------------------------------------------------------------+-- | Create an unsubscribe 'Package'.+createUnsubscribePackage :: Settings -> UUID -> Package+createUnsubscribePackage Settings{..} uuid =+ Package+ { packageCmd = 0xC3+ , packageCorrelation = uuid+ , packageData = runPut $ encodeMessage UnsubscribeFromStream+ , packageCred = s_credentials+ }
+ Database/EventStore/Internal/Operation.hs view
@@ -0,0 +1,193 @@+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE DeriveDataTypeable #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE KindSignatures #-}+{-# LANGUAGE Rank2Types #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation+-- 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.Internal.Operation where++--------------------------------------------------------------------------------+import Control.Applicative+import Control.Exception+import Control.Monad+import Data.Typeable++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text+import Data.UUID+import Data.Word++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+import Prelude++--------------------------------------------------------------------------------+-- | Operation result sent by the server.+data OpResult+ = OP_SUCCESS+ | OP_PREPARE_TIMEOUT+ | OP_COMMIT_TIMEOUT+ | OP_FORWARD_TIMEOUT+ | OP_WRONG_EXPECTED_VERSION+ | OP_STREAM_DELETED+ | OP_INVALID_TRANSACTION+ | OP_ACCESS_DENIED+ deriving (Eq, Enum, Show)++--------------------------------------------------------------------------------+-- | Operation exception that can occurs on an operation response.+data OperationError+ = WrongExpectedVersion Text ExpectedVersion -- ^ Stream and Expected Version+ | StreamDeleted Text -- ^ Stream+ | InvalidTransaction+ | AccessDenied StreamName -- ^ Stream+ | InvalidServerResponse Word8 Word8 -- ^ Expected, Found+ | ProtobufDecodingError String+ | ServerError (Maybe Text) -- ^ Reason+ | InvalidOperation Text+ -- ^ Invalid operation state. If happens, it's a driver bug.+ | Aborted+ -- ^ Occurs when the user asked to close the connection or if the+ -- connection can't reconnect anymore.+ deriving (Show, Typeable)++--------------------------------------------------------------------------------+instance Exception OperationError++--------------------------------------------------------------------------------+-- | Main operation state machine instruction.+data SM o a+ = Return a+ -- ^ Lifts a pure value into the intruction tree. Also marks the end of+ -- an instruction tree.+ | Yield o (SM o a)+ -- ^ Emits an operation return value.+ | FreshId (UUID -> SM o a)+ -- ^ Asks for an unused 'UUID'.+ | forall rq rp. (Encode rq, Decode rp) =>+ SendPkg Word8 Word8 rq (rp -> SM o a)+ -- ^ Send a request message given a command and an expected command.+ -- response. It also carries a callback to call when response comes in.+ | Failure (Maybe OperationError)+ -- ^ Ends the instruction interpretation. If holds Nothing, the+ -- interpretation should resume from the beginning. Otherwise it ends+ -- by indicating what went wrong.++--------------------------------------------------------------------------------+instance Functor (SM o) where+ fmap f (Return a) = Return (f a)+ fmap f (Yield o n) = Yield o (fmap f n)+ fmap f (FreshId k) = FreshId (fmap f . k)+ fmap f (SendPkg ci co p k) = SendPkg ci co p (fmap f . k)+ fmap _ (Failure e) = Failure e++--------------------------------------------------------------------------------+instance Applicative (SM o) where+ pure = return+ (<*>) = ap++--------------------------------------------------------------------------------+instance Monad (SM o) where+ return = Return++ Return a >>= f = f a+ Yield o n >>= f = Yield o (n >>= f)+ FreshId k >>= f = FreshId ((f =<<) . k)+ SendPkg ci co p k >>= f = SendPkg ci co p ((f =<<) . k)+ Failure e >>= _ = Failure e++--------------------------------------------------------------------------------+-- | Asks for a unused 'UUID'.+freshId :: SM o UUID+freshId = FreshId Return++--------------------------------------------------------------------------------+-- | Raises an 'OperationError'.+failure :: OperationError -> SM o a+failure e = Failure $ Just e++--------------------------------------------------------------------------------+-- | Asks to resume the interpretation from the beginning.+retry :: SM o a+retry = Failure Nothing++--------------------------------------------------------------------------------+-- | Sends a request to the server given a command request and response. It+-- returns the expected deserialized message.+send :: (Encode rq, Decode rp) => Word8 -> Word8 -> rq -> SM o rp+send ci co rq = SendPkg ci co rq Return++--------------------------------------------------------------------------------+-- | Emits operation return value.+yield :: o -> SM o ()+yield o = Yield o (Return ())++--------------------------------------------------------------------------------+-- | Replaces every emitted value, via 'yield' function by calling the given+-- callback.+foreach :: SM a x -> (a -> SM b x) -> SM b x+foreach start k = go start+ where+ go (Return x) = Return x+ go (Yield a n) = k a >> go n+ go (FreshId ki) = FreshId (go . ki)+ go (SendPkg ci co p kp) = SendPkg ci co p (go . kp)+ go (Failure e) = Failure e++--------------------------------------------------------------------------------+-- | Maps every emitted value, via 'yield', using given function.+mapOp :: (a -> b) -> SM a () -> SM b ()+mapOp k sm = foreach sm (yield . k)++--------------------------------------------------------------------------------+-- | An operation is just a 'SM' tree.+type Operation a = SM a ()++--------------------------------------------------------------------------------+-- | Raises 'WrongExpectedVersion' exception.+wrongVersion :: Text -> ExpectedVersion -> SM o a+wrongVersion stream ver = failure (WrongExpectedVersion stream ver)++--------------------------------------------------------------------------------+-- | Raises 'StreamDeleted' exception.+streamDeleted :: Text -> SM o a+streamDeleted stream = failure (StreamDeleted stream)++--------------------------------------------------------------------------------+-- | Raises 'InvalidTransaction' exception.+invalidTransaction :: SM o a+invalidTransaction = failure InvalidTransaction++--------------------------------------------------------------------------------+-- | Raises 'AccessDenied' exception.+accessDenied :: StreamName -> SM o a+accessDenied = failure . AccessDenied++--------------------------------------------------------------------------------+-- | Raises 'ProtobufDecodingError' exception.+protobufDecodingError :: String -> SM o a+protobufDecodingError = failure . ProtobufDecodingError++--------------------------------------------------------------------------------+-- | Raises 'ServerError' exception.+serverError :: Maybe Text -> SM o a+serverError = failure . ServerError++--------------------------------------------------------------------------------+-- | Raises 'InvalidServerResponse' exception.+invalidServerResponse :: Word8 -> Word8 -> SM o a+invalidServerResponse expe got = failure $ InvalidServerResponse expe got
+ Database/EventStore/Internal/Operation/Catchup.hs view
@@ -0,0 +1,102 @@+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE OverloadedStrings #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.Catchup+-- 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.Internal.Operation.Catchup+ ( CatchupState(..)+ , catchup+ ) where++--------------------------------------------------------------------------------+import Control.Monad+import Data.Int+import Data.Maybe++--------------------------------------------------------------------------------+import Data.Text (Text)++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Manager.Subscription (Checkpoint(..))+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Operation.Read.Common+import Database.EventStore.Internal.Operation.ReadAllEvents+import Database.EventStore.Internal.Operation.ReadStreamEvents+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+defaultBatchSize :: Int32+defaultBatchSize = 500++--------------------------------------------------------------------------------+streamNotFound :: OperationError+streamNotFound = InvalidOperation "Catchup. inexistant stream"++--------------------------------------------------------------------------------+-- | Catchup operation state.+data CatchupState+ = RegularCatchup Text Int32+ -- ^ Indicates the stream name and the next event number to start from.+ | AllCatchup Int64 Int64+ -- ^ Indicates the commit and prepare position. Used when catching up from+ -- the $all stream.++--------------------------------------------------------------------------------+-- | Stream catching up operation.+catchup :: Settings+ -> CatchupState+ -> Bool+ -> Maybe Int32+ -> Operation ([ResolvedEvent], Bool, Checkpoint)+catchup setts init_tpe tos bat_siz = go init_tpe+ where+ batch = fromMaybe defaultBatchSize bat_siz+ go tpe = do+ let action =+ case tpe of+ RegularCatchup stream cur_evt ->+ let op = readStreamEvents setts Forward stream cur_evt+ batch tos in+ mapOp Left op+ AllCatchup c_pos p_pos ->+ let op = readAllEvents setts c_pos p_pos batch+ tos Forward in+ mapOp Right op++ foreach action $ \res -> do+ (eos, evts, nchk, nxt_tpe) <- case res of+ Right as -> do+ let Position nxt_c nxt_p = sliceNext as+ tmp_tpe = AllCatchup nxt_c nxt_p+ chk = CheckpointPosition $ sliceNext as+ return (sliceEOS as, sliceEvents as, chk, tmp_tpe)+ Left rr -> fromReadResult rr $ \as ->+ let RegularCatchup s _ = tpe+ nxt = sliceNext as+ tmp_tpe = RegularCatchup s nxt+ chk = CheckpointNumber nxt in+ return (sliceEOS as, sliceEvents as, chk, tmp_tpe)++ yield (evts, eos, nchk)+ when (not eos) $ go nxt_tpe++--------------------------------------------------------------------------------+fromReadResult :: ReadResult 'RegularStream a -> (a -> SM b x) -> SM b x+fromReadResult res k =+ case res of+ ReadNoStream -> failure streamNotFound+ ReadStreamDeleted s -> failure $ StreamDeleted s+ ReadNotModified -> failure $ ServerError Nothing+ ReadError e -> failure $ ServerError e+ ReadAccessDenied s -> failure $ AccessDenied s+ ReadSuccess ss -> k ss
+ Database/EventStore/Internal/Operation/DeleteStream.hs view
@@ -0,0 +1,63 @@+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE Rank2Types #-}+{-# LANGUAGE RecordWildCards #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.DeleteStream+-- 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.Internal.Operation.DeleteStream+ ( DeleteResult(..)+ , deleteStream+ ) where++--------------------------------------------------------------------------------+import Data.Maybe++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Operation.DeleteStream.Message+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Returned after deleting a stream. 'Position' of the write.+newtype DeleteResult = DeleteResult Position deriving (Eq, Show)++--------------------------------------------------------------------------------+-- | Delete a regular stream operation.+deleteStream :: Settings+ -> Text+ -> ExpectedVersion+ -> Maybe Bool+ -> Operation DeleteResult+deleteStream Settings{..} s v hard = do+ let msg = newRequest s (expVersionInt32 v) s_requireMaster hard+ resp <- send 0x8A 0x8B msg+ let r = getField $ _result resp+ com_pos = getField $ _commitPosition resp+ prep_pos = getField $ _preparePosition resp+ com_pos_int = fromMaybe (-1) com_pos+ prep_pos_int = fromMaybe (-1) prep_pos+ pos = Position com_pos_int prep_pos_int+ res = DeleteResult pos+ case r of+ OP_SUCCESS -> yield res+ OP_PREPARE_TIMEOUT -> retry+ OP_FORWARD_TIMEOUT -> retry+ OP_COMMIT_TIMEOUT -> retry+ OP_WRONG_EXPECTED_VERSION -> wrongVersion s v+ OP_STREAM_DELETED -> streamDeleted s+ OP_INVALID_TRANSACTION -> invalidTransaction+ OP_ACCESS_DENIED -> accessDenied (StreamName s)
+ Database/EventStore/Internal/Operation/DeleteStream/Message.hs view
@@ -0,0 +1,64 @@+{-# LANGUAGE DeriveGeneric #-}+{-# LANGUAGE DataKinds #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.DeleteStream.Message+-- 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.Internal.Operation.DeleteStream.Message where++--------------------------------------------------------------------------------+import Data.Int+import GHC.Generics++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation++--------------------------------------------------------------------------------+-- | Delete stream request.+data Request+ = Request+ { _streamId :: Required 1 (Value Text)+ , _expectedVersion :: Required 2 (Value Int32)+ , _requireMaster :: Required 3 (Value Bool)+ , _hardDelete :: Optional 4 (Value Bool)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode Request++--------------------------------------------------------------------------------+-- | 'Request' smart constructor.+newRequest :: Text -> Int32 -> Bool -> Maybe Bool -> Request+newRequest stream_id exp_ver req_master hard_delete =+ Request+ { _streamId = putField stream_id+ , _expectedVersion = putField exp_ver+ , _requireMaster = putField req_master+ , _hardDelete = putField hard_delete+ }++--------------------------------------------------------------------------------+-- | Delete stream response.+data Response+ = Response+ { _result :: Required 1 (Enumeration OpResult)+ , _message :: Optional 2 (Value Text)+ , _preparePosition :: Optional 3 (Value Int64)+ , _commitPosition :: Optional 4 (Value Int64)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode Response
− Database/EventStore/Internal/Operation/DeleteStreamOperation.hs
@@ -1,137 +0,0 @@-{-# LANGUAGE DeriveGeneric #-}-{-# LANGUAGE DataKinds #-}------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Operation.DeleteStreamOperation--- 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.Operation.DeleteStreamOperation- ( deleteStreamOperation ) where-----------------------------------------------------------------------------------import Control.Concurrent-import Data.Int-import Data.Maybe-import GHC.Generics (Generic)-----------------------------------------------------------------------------------import Data.ProtocolBuffers-import Data.Text-----------------------------------------------------------------------------------import Database.EventStore.Internal.Manager.Operation-import Database.EventStore.Internal.Types-----------------------------------------------------------------------------------data DeleteStream- = DeleteStream- { _deleteStreamId :: Required 1 (Value Text)- , _deleteStreamExpectedVersion :: Required 2 (Value Int32)- , _deleteStreamRequireMaster :: Required 3 (Value Bool)- , _deleteStreamHardDelete :: Optional 4 (Value Bool)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Encode DeleteStream-----------------------------------------------------------------------------------newDeleteStream :: Text- -> Int32- -> Bool- -> Maybe Bool- -> DeleteStream-newDeleteStream stream_id exp_ver req_master hard_delete =- DeleteStream- { _deleteStreamId = putField stream_id- , _deleteStreamExpectedVersion = putField exp_ver- , _deleteStreamRequireMaster = putField req_master- , _deleteStreamHardDelete = putField hard_delete- }-----------------------------------------------------------------------------------data DeleteStreamCompleted- = DeleteStreamCompleted- { _deleteCompletedResult :: Required 1 (Enumeration OpResult)- , _deleteCompletedMessage :: Optional 2 (Value Text)- , _deleteCompletedPreparePosition :: Optional 3 (Value Int64)- , _deleteCompletedCommitPosition :: Optional 4 (Value Int64)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode DeleteStreamCompleted-----------------------------------------------------------------------------------deleteStreamOperation :: Settings- -> MVar (OperationExceptional DeleteResult)- -> Text- -> ExpectedVersion- -> Maybe Bool- -> OperationParams-deleteStreamOperation settings mvar stream_id exp_ver hard_del =- OperationParams- { opSettings = settings- , opRequestCmd = 0x8A- , opResponseCmd = 0x8B-- , opRequest =- let req_master = s_requireMaster settings- exp_ver_int = expVersionInt32 exp_ver- request = newDeleteStream stream_id- exp_ver_int- req_master- hard_del in- return request-- , opSuccess = inspect mvar stream_id exp_ver- , opFailure = failed mvar- }-----------------------------------------------------------------------------------inspect :: MVar (OperationExceptional DeleteResult)- -> Text- -> ExpectedVersion- -> DeleteStreamCompleted- -> IO Decision-inspect mvar stream exp_ver dsc = go (getField $ _deleteCompletedResult dsc)- where- go OP_SUCCESS = succeed mvar dsc- go OP_PREPARE_TIMEOUT = return Retry- go OP_FORWARD_TIMEOUT = return Retry- go OP_COMMIT_TIMEOUT = return Retry- go OP_WRONG_EXPECTED_VERSION = failed mvar wrong_version- go OP_STREAM_DELETED = failed mvar (StreamDeleted stream)- go OP_INVALID_TRANSACTION = failed mvar InvalidTransaction- go OP_ACCESS_DENIED = failed mvar (AccessDenied stream)-- wrong_version = WrongExpectedVersion stream exp_ver-----------------------------------------------------------------------------------succeed :: MVar (OperationExceptional DeleteResult)- -> DeleteStreamCompleted- -> IO Decision-succeed mvar wec = do- putMVar mvar (Right wr)- return EndOperation- where- com_pos = getField $ _deleteCompletedCommitPosition wec- pre_pos = getField $ _deleteCompletedPreparePosition wec- com_pos_int = fromMaybe (-1) com_pos- pre_pos_int = fromMaybe (-1) pre_pos- pos = Position com_pos_int pre_pos_int- wr = DeleteResult pos-----------------------------------------------------------------------------------failed :: MVar (OperationExceptional DeleteResult)- -> OperationException- -> IO Decision-failed mvar e = do- putMVar mvar (Left e)- return EndOperation
+ Database/EventStore/Internal/Operation/Read/Common.hs view
@@ -0,0 +1,145 @@+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE KindSignatures #-}+{-# LANGUAGE TypeFamilies #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.Read.Common+-- 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.Internal.Operation.Read.Common where++--------------------------------------------------------------------------------+import Control.Applicative+import Data.Foldable+import Data.Int+import Data.Monoid+import Data.Traversable++--------------------------------------------------------------------------------+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+import Prelude++--------------------------------------------------------------------------------+-- | Enumeration detailing the possible outcomes of reading a stream.+data ReadResult :: StreamType -> * -> * where+ ReadSuccess :: a -> ReadResult t a+ ReadNoStream :: ReadResult 'RegularStream a+ ReadStreamDeleted :: Text -> ReadResult 'RegularStream a+ ReadNotModified :: ReadResult t a+ ReadError :: Maybe Text -> ReadResult t a+ ReadAccessDenied :: StreamName -> ReadResult t a++--------------------------------------------------------------------------------+instance Eq a => Eq (ReadResult t a) where+ ReadSuccess a == ReadSuccess b = a == b+ ReadNoStream == ReadNoStream = True+ ReadStreamDeleted s == ReadStreamDeleted v = s == v+ ReadNotModified == ReadNotModified = True+ ReadError e == ReadError u = e == u+ ReadAccessDenied s == ReadAccessDenied v = s == v+ _ == _ = False++--------------------------------------------------------------------------------+instance Show a => Show (ReadResult t a) where+ show (ReadSuccess a) = "ReadSuccess " ++ show a+ show ReadNoStream = "ReadNoStream"+ show (ReadStreamDeleted s) = "ReadStreamDeleted" ++ show s+ show ReadNotModified = "ReadNoModified"+ show (ReadError e) = "ReadError" ++ show e+ show (ReadAccessDenied s) = "ReadAccessDenied " ++ show s++--------------------------------------------------------------------------------+instance Functor (ReadResult t) where+ fmap f (ReadSuccess a) = ReadSuccess (f a)+ fmap _ ReadNoStream = ReadNoStream+ fmap _ (ReadStreamDeleted s) = ReadStreamDeleted s+ fmap _ ReadNotModified = ReadNotModified+ fmap _ (ReadError e) = ReadError e+ fmap _ (ReadAccessDenied s) = ReadAccessDenied s++--------------------------------------------------------------------------------+instance Foldable (ReadResult t) where+ foldMap f (ReadSuccess a) = f a+ foldMap _ _ = mempty++--------------------------------------------------------------------------------+instance Traversable (ReadResult t) where+ traverse f (ReadSuccess a) = fmap ReadSuccess $ f a+ traverse _ ReadNoStream = pure ReadNoStream+ traverse _ (ReadStreamDeleted s) = pure $ ReadStreamDeleted s+ traverse _ ReadNotModified = pure ReadNotModified+ traverse _ (ReadError e) = pure $ ReadError e+ traverse _ (ReadAccessDenied s) = pure $ ReadAccessDenied s++--------------------------------------------------------------------------------+-- | Gathers common slice operations.+class Slice a where+ type Loc a++ sliceEvents :: a -> [ResolvedEvent]+ -- ^ Gets slice's 'ResolvedEvent's.+ sliceDirection :: a -> ReadDirection+ -- ^ Gets slice's reading direction.+ sliceEOS :: a -> Bool+ -- ^ If the slice reaches the end of the stream.+ sliceFrom :: a -> Loc a+ -- ^ Gets the starting location of this slice.+ sliceNext :: a -> Loc a+ -- ^ Gets the next location of this slice.++--------------------------------------------------------------------------------+-- | Regular stream slice.+data StreamSlice =+ StreamSlice+ { sliceStream :: !Text+ , sliceLast :: !Int32+ , _ssDir :: !ReadDirection+ , _ssFrom :: !Int32+ , _ssNext :: !Int32+ , _ssEvents :: ![ResolvedEvent]+ , _ssEOS :: !Bool+ } deriving Show++--------------------------------------------------------------------------------+instance Slice StreamSlice where+ type Loc StreamSlice = Int32++ sliceEvents = _ssEvents+ sliceDirection = _ssDir+ sliceEOS = _ssEOS+ sliceFrom = _ssFrom+ sliceNext = _ssNext++--------------------------------------------------------------------------------+-- | Represents a slice of the $all stream.+data AllSlice =+ AllSlice+ { _saFrom :: !Position+ , _saNext :: !Position+ , _saDir :: !ReadDirection+ , _saEvents :: ![ResolvedEvent]+ , _saEOS :: !Bool+ } deriving Show++--------------------------------------------------------------------------------+instance Slice AllSlice where+ type Loc AllSlice = Position++ sliceEvents = _saEvents+ sliceDirection = _saDir+ sliceEOS = _saEOS+ sliceFrom = _saFrom+ sliceNext = _saNext
+ Database/EventStore/Internal/Operation/ReadAllEvents.hs view
@@ -0,0 +1,62 @@+{-# LANGUAGE RecordWildCards #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.ReadAllEvents+-- 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.Internal.Operation.ReadAllEvents+ ( readAllEvents ) where++--------------------------------------------------------------------------------+import Data.Int+import Data.Maybe++--------------------------------------------------------------------------------+import Data.ProtocolBuffers++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Operation.Read.Common+import Database.EventStore.Internal.Operation.ReadAllEvents.Message+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Batch read on $all stream operation.+readAllEvents :: Settings+ -> Int64+ -> Int64+ -> Int32+ -> Bool+ -> ReadDirection+ -> Operation AllSlice+readAllEvents Settings{..} c_pos p_pos max_c tos dir = do+ let msg = newRequest c_pos p_pos max_c tos s_requireMaster+ cmd = case dir of+ Forward -> 0xB6+ Backward -> 0xB8++ resp_cmd = case dir of+ Forward -> 0xB7+ Backward -> 0xB9+ resp <- send cmd resp_cmd msg+ let r = getField $ _Result resp+ err = getField $ _Error resp+ nc_pos = getField $ _NextCommitPosition resp+ np_pos = getField $ _NextPreparePosition resp+ es = getField $ _Events resp+ evts = fmap newResolvedEventFromBuf es+ eos = null evts+ f_pos = Position c_pos p_pos+ n_pos = Position nc_pos np_pos+ slice = AllSlice f_pos n_pos dir evts eos+ case fromMaybe SUCCESS r of+ ERROR -> serverError err+ ACCESS_DENIED -> accessDenied AllStream+ _ -> yield slice
+ Database/EventStore/Internal/Operation/ReadAllEvents/Message.hs view
@@ -0,0 +1,84 @@+{-# LANGUAGE DeriveGeneric #-}+{-# LANGUAGE DataKinds #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.ReadAllEvents.Message+-- 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.Internal.Operation.ReadAllEvents.Message where++--------------------------------------------------------------------------------+import Data.Int+import GHC.Generics++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Batch read on $all stream request.+data Request+ = Request+ { _commitPosition :: Required 1 (Value Int64)+ , _preparePosition :: Required 2 (Value Int64)+ , _maxCount :: Required 3 (Value Int32)+ , _resolveLinkTos :: Required 4 (Value Bool)+ , _requireMaster :: Required 5 (Value Bool)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode Request++--------------------------------------------------------------------------------+-- | 'Request' smart constructor.+newRequest :: Int64+ -> Int64+ -> Int32+ -> Bool+ -> Bool+ -> Request+newRequest c_pos p_pos max_c res_link_tos req_master =+ Request+ { _commitPosition = putField c_pos+ , _preparePosition = putField p_pos+ , _maxCount = putField max_c+ , _resolveLinkTos = putField res_link_tos+ , _requireMaster = putField req_master+ }++--------------------------------------------------------------------------------+-- | Enumeration detailing the possible outcomes of reading a slice of $all+-- stream.+data Result+ = SUCCESS+ | NOT_MODIFIED+ | ERROR+ | ACCESS_DENIED+ deriving (Eq, Enum, Show)++--------------------------------------------------------------------------------+-- | Batch read on $all stream response.+data Response+ = Response+ { _CommitPosition :: Required 1 (Value Int64)+ , _PreparePosition :: Required 2 (Value Int64)+ , _Events :: Repeated 3 (Message ResolvedEventBuf)+ , _NextCommitPosition :: Required 4 (Value Int64)+ , _NextPreparePosition :: Required 5 (Value Int64)+ , _Result :: Optional 6 (Enumeration Result)+ , _Error :: Optional 7 (Value Text)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode Response
− Database/EventStore/Internal/Operation/ReadAllEventsOperation.hs
@@ -1,200 +0,0 @@-{-# LANGUAGE DeriveGeneric #-}-{-# LANGUAGE DataKinds #-}-{-# LANGUAGE OverloadedStrings #-}------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Operation.ReadAllEventsOperation--- 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.Operation.ReadAllEventsOperation- ( AllEventsSlice(..)- , ReadAllResult(..)- , readAllEventsOperation- ) where-----------------------------------------------------------------------------------import Control.Concurrent-import Data.Int-import Data.Maybe-import GHC.Generics (Generic)-----------------------------------------------------------------------------------import Data.Text hiding (null)-import Data.ProtocolBuffers-----------------------------------------------------------------------------------import Database.EventStore.Internal.Manager.Operation-import Database.EventStore.Internal.Types-----------------------------------------------------------------------------------data ReadAllEvents- = ReadAllEvents- { readAllEventsCommitPosition :: Required 1 (Value Int64)- , readAllEventsPreparePosition :: Required 2 (Value Int64)- , readAllEventsMaxCount :: Required 3 (Value Int32)- , readAllEventsResolveLinkTos :: Required 4 (Value Bool)- , readAllEventsRequireMaster :: Required 5 (Value Bool)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Encode ReadAllEvents-----------------------------------------------------------------------------------newReadAllEvents :: Int64- -> Int64- -> Int32- -> Bool- -> Bool- -> ReadAllEvents-newReadAllEvents c_pos p_pos max_c res_link_tos req_master =- ReadAllEvents- { readAllEventsCommitPosition = putField c_pos- , readAllEventsPreparePosition = putField p_pos- , readAllEventsMaxCount = putField max_c- , readAllEventsResolveLinkTos = putField res_link_tos- , readAllEventsRequireMaster = putField req_master- }------------------------------------------------------------------------------------- | Enumeration detailing the possible outcomes of reading a slice of $all--- stream.-data ReadAllResult- = RA_SUCCESS- | RA_NOT_MODIFIED- | RA_ERROR- | RA_ACCESS_DENIED- deriving (Eq, Enum, Show)-----------------------------------------------------------------------------------data ReadAllEventsCompleted- = ReadAllEventsCompleted- { readAECCommitPosition :: Required 1 (Value Int64)- , readAECPreparePosition :: Required 2 (Value Int64)- , readAECEvents :: Repeated 3 (Message ResolvedEventBuf)- , readAECNextCommitPosition :: Required 4 (Value Int64)- , readAECNextPreparePosition :: Required 5 (Value Int64)- , readAECResult :: Optional 6 (Enumeration ReadAllResult)- , readAECError :: Optional 7 (Value Text)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode ReadAllEventsCompleted------------------------------------------------------------------------------------- | The result of a read operation from the $all stream.-data AllEventsSlice- = AllEventsSlice- { allEventsSliceResult :: !ReadAllResult- -- ^ Representing the status of the read attempt.- , allEventsSliceFrom :: !Position- -- ^ Representing the position where the next slice should be read- -- from.- , allEventsSliceNext :: !Position- -- ^ Representing the position where the next slice should be read from.- , allEventsSliceIsEOS :: !Bool- -- ^ Representing whether or not this is the end of the $all stream.- , allEventsSliceEvents :: ![ResolvedEvent]- -- ^ The events read.- , allEventsSliceDirection :: !ReadDirection- -- ^ The direction of read request.- }- deriving Show-----------------------------------------------------------------------------------newAllEventsSlice :: ReadDirection -> ReadAllEventsCompleted -> AllEventsSlice-newAllEventsSlice dir raec = aes- where- res = fromMaybe RA_SUCCESS (getField $ readAECResult raec)- evts = fmap newResolvedEventFromBuf (getField $ readAECEvents raec)- r_com = getField $ readAECCommitPosition raec- r_pre = getField $ readAECPreparePosition raec- r_n_com = getField $ readAECNextCommitPosition raec- r_n_pre = getField $ readAECNextPreparePosition raec- from_pos = Position r_com r_pre- next_pos = Position r_n_com r_n_pre-- aes = AllEventsSlice- { allEventsSliceResult = res- , allEventsSliceFrom = from_pos- , allEventsSliceNext = next_pos- , allEventsSliceIsEOS = null evts- , allEventsSliceEvents = evts- , allEventsSliceDirection = dir- }-----------------------------------------------------------------------------------readAllEventsOperation :: Settings- -> ReadDirection- -> MVar (OperationExceptional AllEventsSlice)- -> Int64- -> Int64- -> Int32- -> Bool- -> OperationParams-readAllEventsOperation settings dir mvar c_pos p_pos max_c res_link_tos =- OperationParams- { opSettings = settings- , opRequestCmd = req- , opResponseCmd = resp-- , opRequest =- let req_master = s_requireMaster settings- request = newReadAllEvents c_pos- p_pos- max_c- res_link_tos- req_master in- return request-- , opSuccess = inspect mvar dir- , opFailure = failed mvar- }- where- req = case dir of- Forward -> 0xB6- Backward -> 0xB8-- resp = case dir of- Forward -> 0xB7- Backward -> 0xB9-----------------------------------------------------------------------------------inspect :: MVar (OperationExceptional AllEventsSlice)- -> ReadDirection- -> ReadAllEventsCompleted- -> IO Decision-inspect mvar dir raec = go res- where- res = fromMaybe RA_SUCCESS (getField $ readAECResult raec)- may_err = getField $ readAECError raec-- go RA_ERROR = failed mvar (ServerError may_err)- go RA_ACCESS_DENIED = failed mvar (AccessDenied "$all")- go _ = succeed mvar dir raec-----------------------------------------------------------------------------------succeed :: MVar (OperationExceptional AllEventsSlice)- -> ReadDirection- -> ReadAllEventsCompleted- -> IO Decision-succeed mvar dir raec = do- putMVar mvar (Right ses)- return EndOperation- where- ses = newAllEventsSlice dir raec-----------------------------------------------------------------------------------failed :: MVar (OperationExceptional AllEventsSlice)- -> OperationException- -> IO Decision-failed mvar e = do- putMVar mvar (Left e)- return EndOperation
+ Database/EventStore/Internal/Operation/ReadEvent.hs view
@@ -0,0 +1,67 @@+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE RecordWildCards #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.DeleteStream+-- 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.Internal.Operation.ReadEvent+ ( ReadEvent(..)+ , readEvent+ ) where++--------------------------------------------------------------------------------+import Data.Int++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Operation.ReadEvent.Message+import Database.EventStore.Internal.Operation.Read.Common+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Represents the result of looking up a specific event number from a stream.+data ReadEvent+ = ReadEventNotFound+ { readEventStream :: !Text+ , readEventNumber :: !Int32+ }+ | ReadEvent+ { readEventStream :: !Text+ , readEventNumber :: !Int32+ , readEventResolved :: !ResolvedEvent+ } deriving Show++--------------------------------------------------------------------------------+-- | Read a specific event given event number operation.+readEvent :: Settings+ -> Text+ -> Int32+ -> Bool+ -> Operation (ReadResult 'RegularStream ReadEvent)+readEvent Settings{..} s evtn tos = do+ let msg = newRequest s evtn tos s_requireMaster+ resp <- send 0xB0 0xB1 msg+ let r = getField $ _result resp+ evt = newResolvedEvent $ getField $ _indexedEvent resp+ err = getField $ _error resp+ not_found = ReadSuccess $ ReadEventNotFound s evtn+ found = ReadSuccess $ ReadEvent s evtn evt+ case r of+ NOT_FOUND -> yield not_found+ NO_STREAM -> yield ReadNoStream+ STREAM_DELETED -> yield $ ReadStreamDeleted s+ ERROR -> yield (ReadError err)+ ACCESS_DENIED -> yield $ ReadAccessDenied $ StreamName s+ SUCCESS -> yield found
+ Database/EventStore/Internal/Operation/ReadEvent/Message.hs view
@@ -0,0 +1,75 @@+{-# LANGUAGE DeriveGeneric #-}+{-# LANGUAGE DataKinds #-}+{-# OPTIONS_GHC -fcontext-stack=26 #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.ReadEvent.Message+-- 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.Internal.Operation.ReadEvent.Message where++--------------------------------------------------------------------------------+import Data.Int+import GHC.Generics++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Read event on a regular stream request.+data Request+ = Request+ { _streamId :: Required 1 (Value Text)+ , _eventNumber :: Required 2 (Value Int32)+ , _resolveLinkTos :: Required 3 (Value Bool)+ , _requireMaster :: Required 4 (Value Bool)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode Request++--------------------------------------------------------------------------------+-- | 'Request' smart constructor.+newRequest :: Text -> Int32 -> Bool -> Bool -> Request+newRequest stream_id evt_num res_link_tos req_master =+ Request+ { _streamId = putField stream_id+ , _eventNumber = putField evt_num+ , _resolveLinkTos = putField res_link_tos+ , _requireMaster = putField req_master+ }++--------------------------------------------------------------------------------+-- | Enumeration representing the status of a single event read operation.+data Result+ = SUCCESS+ | NOT_FOUND+ | NO_STREAM+ | STREAM_DELETED+ | ERROR+ | ACCESS_DENIED+ deriving (Eq, Enum, Show)++--------------------------------------------------------------------------------+-- | Read event on a regular stream response.+data Response+ = Response+ { _result :: Required 1 (Enumeration Result)+ , _indexedEvent :: Required 2 (Message ResolvedIndexedEvent)+ , _error :: Optional 3 (Value Text)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode Response
− Database/EventStore/Internal/Operation/ReadEventOperation.hs
@@ -1,170 +0,0 @@-{-# LANGUAGE DeriveGeneric #-}-{-# LANGUAGE DataKinds #-}-{-# OPTIONS_GHC -fcontext-stack=26 #-}------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Operation.ReadEventOperation--- 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.Operation.ReadEventOperation- ( ReadResult(..)- , ReadEventResult(..)- , readEventOperation- ) where-----------------------------------------------------------------------------------import Control.Concurrent-import Data.Int-import GHC.Generics (Generic)-----------------------------------------------------------------------------------import Data.ProtocolBuffers-import Data.Text-----------------------------------------------------------------------------------import Database.EventStore.Internal.Manager.Operation-import Database.EventStore.Internal.Types-----------------------------------------------------------------------------------data ReadEvent- = ReadEvent- { readEventStreamId :: Required 1 (Value Text)- , readEventNumber :: Required 2 (Value Int32)- , readEventResolveLinkTos :: Required 3 (Value Bool)- , readEventRequireMaster :: Required 4 (Value Bool)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Encode ReadEvent-----------------------------------------------------------------------------------newReadEvent :: Text -> Int32 -> Bool -> Bool -> ReadEvent-newReadEvent stream_id evt_num res_link_tos req_master =- ReadEvent- { readEventStreamId = putField stream_id- , readEventNumber = putField evt_num- , readEventResolveLinkTos = putField res_link_tos- , readEventRequireMaster = putField req_master- }------------------------------------------------------------------------------------- | Enumeration representing the status of a single event read operation.-data ReadEventResult- = RE_SUCCESS- | RE_NOT_FOUND- | RE_NO_STREAM- | RE_STREAM_DELETED- | RE_ERROR- | RE_ACCESS_DENIED- deriving (Eq, Enum, Show)-----------------------------------------------------------------------------------data ReadEventCompleted- = ReadEventCompleted- { readCompletedResult :: Required 1 (Enumeration ReadEventResult)- , readCompletedIndexedEvent :: Required 2 (Message ResolvedIndexedEvent)- , readCompletedError :: Optional 3 (Value Text)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode ReadEventCompleted------------------------------------------------------------------------------------- | Result of a single event read operation to the EventStore.-data ReadResult- = ReadResult- { readResultStatus :: !ReadEventResult -- ^ Attempt status- , readResultStreamId :: !Text -- ^ Stream name- , readResultEventNumber :: !Int32 -- ^ Event number- , readResultResolvedEvent :: !(Maybe ResolvedEvent)- }- deriving Show-----------------------------------------------------------------------------------newReadResult :: ReadEventResult- -> Text- -> Int32- -> ResolvedIndexedEvent- -> ReadResult-newReadResult status stream_id evt_num rie = rr- where- may_re =- case status of- RE_SUCCESS -> Just $ newResolvedEvent rie- _ -> Nothing-- rr = ReadResult- { readResultStatus = status- , readResultStreamId = stream_id- , readResultEventNumber = evt_num- , readResultResolvedEvent = may_re- }-----------------------------------------------------------------------------------readEventOperation :: Settings- -> MVar (OperationExceptional ReadResult)- -> Text- -> Int32- -> Bool -- ^ Resolve link TOS- -> OperationParams-readEventOperation settings mvar stream_id evt_num res_link_tos =- OperationParams- { opSettings = settings- , opRequestCmd = 0xB0- , opResponseCmd = 0xB1-- , opRequest =- let req_master = s_requireMaster settings- request = newReadEvent stream_id- evt_num- res_link_tos- req_master in- return request-- , opSuccess = inspect mvar stream_id evt_num- , opFailure = failed mvar- }-----------------------------------------------------------------------------------inspect :: MVar (OperationExceptional ReadResult)- -> Text- -> Int32- -> ReadEventCompleted- -> IO Decision-inspect mvar stream_id evt_num reco = go (getField $ readCompletedResult reco)- where- may_err = getField $ readCompletedError reco-- go RE_ERROR = failed mvar (ServerError may_err)- go RE_ACCESS_DENIED = failed mvar (AccessDenied stream_id)- go _ = succeed mvar stream_id evt_num reco-----------------------------------------------------------------------------------succeed :: MVar (OperationExceptional ReadResult)- -> Text- -> Int32- -> ReadEventCompleted- -> IO Decision-succeed mvar stream_id evt_num reco = do- putMVar mvar (Right rr)- return EndOperation- where- status = getField $ readCompletedResult reco- rie = getField $ readCompletedIndexedEvent reco- rr = newReadResult status stream_id evt_num rie-----------------------------------------------------------------------------------failed :: MVar (OperationExceptional ReadResult)- -> OperationException- -> IO Decision-failed mvar e = do- putMVar mvar (Left e)- return EndOperation
+ Database/EventStore/Internal/Operation/ReadStreamEvents.hs view
@@ -0,0 +1,66 @@+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE RecordWildCards #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.ReadStreamEvents+-- 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.Internal.Operation.ReadStreamEvents+ ( readStreamEvents ) where++--------------------------------------------------------------------------------+import Data.Int++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Operation.Read.Common+import Database.EventStore.Internal.Operation.ReadStreamEvents.Message+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Batch read from a regular stream operation.+readStreamEvents :: Settings+ -> ReadDirection+ -> Text+ -> Int32+ -> Int32+ -> Bool+ -> Operation (ReadResult 'RegularStream StreamSlice)+readStreamEvents Settings{..} dir s st cnt tos = do+ let req_cmd =+ case dir of+ Forward -> 0xB2+ Backward -> 0xB4+ resp_cmd =+ case dir of+ Forward -> 0xB3+ Backward -> 0xB5++ msg = newRequest s st cnt tos s_requireMaster+ resp <- send req_cmd resp_cmd msg+ let r = getField $ _result resp+ es = getField $ _events resp+ evts = fmap newResolvedEvent es+ err = getField $ _error resp+ eos = getField $ _endOfStream resp+ nxt = getField $ _nextNumber resp+ lst = getField $ _lastNumber resp+ found = StreamSlice s lst dir st nxt evts eos+ case r of+ NO_STREAM -> yield ReadNoStream+ STREAM_DELETED -> yield $ ReadStreamDeleted s+ NOT_MODIFIED -> yield ReadNotModified+ ERROR -> yield (ReadError err)+ ACCESS_DENIED -> yield $ ReadAccessDenied $ StreamName s+ SUCCESS -> yield (ReadSuccess found)
+ Database/EventStore/Internal/Operation/ReadStreamEvents/Message.hs view
@@ -0,0 +1,80 @@+{-# LANGUAGE DeriveGeneric #-}+{-# LANGUAGE DataKinds #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.ReadStreamEvents.Message+-- 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.Internal.Operation.ReadStreamEvents.Message where++--------------------------------------------------------------------------------+import Data.Int+import GHC.Generics++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Batch read on regular stream request.+data Request+ = Request+ { _streamId :: Required 1 (Value Text)+ , _eventNumber :: Required 2 (Value Int32)+ , _maxCount :: Required 3 (Value Int32)+ , _resolveLinkTos :: Required 4 (Value Bool)+ , _requireMaster :: Required 5 (Value Bool)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+-- | 'Request' smart constructor.+newRequest :: Text -> Int32 -> Int32 -> Bool -> Bool -> Request+newRequest stream_id evt_num max_c res_link_tos req_master =+ Request+ { _streamId = putField stream_id+ , _eventNumber = putField evt_num+ , _maxCount = putField max_c+ , _resolveLinkTos = putField res_link_tos+ , _requireMaster = putField req_master+ }++--------------------------------------------------------------------------------+instance Encode Request++--------------------------------------------------------------------------------+-- | Enumeration detailing the possible outcomes of reading a slice of a stream+data Result+ = SUCCESS+ | NO_STREAM+ | STREAM_DELETED+ | NOT_MODIFIED+ | ERROR+ | ACCESS_DENIED+ deriving (Eq, Enum, Show)++--------------------------------------------------------------------------------+-- | Batch read on regular stream response.+data Response+ = Response+ { _events :: Repeated 1 (Message ResolvedIndexedEvent)+ , _result :: Required 2 (Enumeration Result)+ , _nextNumber :: Required 3 (Value Int32)+ , _lastNumber :: Required 4 (Value Int32)+ , _endOfStream :: Required 5 (Value Bool)+ , _lastCommitPosition :: Required 6 (Value Int64)+ , _error :: Optional 7 (Value Text)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode Response
− Database/EventStore/Internal/Operation/ReadStreamEventsOperation.hs
@@ -1,205 +0,0 @@-{-# LANGUAGE DeriveGeneric #-}-{-# LANGUAGE DataKinds #-}------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Operation.ReadStreamEventsOperation--- 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.Operation.ReadStreamEventsOperation- ( StreamEventsSlice(..)- , ReadStreamResult(..)- , readStreamEventsOperation- ) where-----------------------------------------------------------------------------------import Control.Concurrent-import Data.Int-import GHC.Generics (Generic)-----------------------------------------------------------------------------------import Data.ProtocolBuffers-import Data.Text-----------------------------------------------------------------------------------import Database.EventStore.Internal.Manager.Operation-import Database.EventStore.Internal.Types-----------------------------------------------------------------------------------data ReadStreamEvents- = ReadStreamEvents- { _readStreamId :: Required 1 (Value Text)- , _readStreamEventNumber :: Required 2 (Value Int32)- , _readStreamMaxCount :: Required 3 (Value Int32)- , _readStreamResolveLinkTos :: Required 4 (Value Bool)- , _readStreamRequireMaster :: Required 5 (Value Bool)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------newReadStreamEvents :: Text- -> Int32- -> Int32- -> Bool- -> Bool- -> ReadStreamEvents-newReadStreamEvents stream_id evt_num max_c res_link_tos req_master =- ReadStreamEvents- { _readStreamId = putField stream_id- , _readStreamEventNumber = putField evt_num- , _readStreamMaxCount = putField max_c- , _readStreamResolveLinkTos = putField res_link_tos- , _readStreamRequireMaster = putField req_master- }-----------------------------------------------------------------------------------instance Encode ReadStreamEvents------------------------------------------------------------------------------------- | Enumeration detailing the possible outcomes of reading a slice of a stream-data ReadStreamResult- = RS_SUCCESS- | RS_NO_STREAM- | RS_STREAM_DELETED- | RS_NOT_MODIFIED- | RS_ERROR- | RS_ACCESS_DENIED- deriving (Eq, Enum, Show)-----------------------------------------------------------------------------------data ReadStreamEventsCompleted- = ReadStreamEventsCompleted- { _readSECEvents :: Repeated 1 (Message ResolvedIndexedEvent)- , _readSECResult :: Required 2 (Enumeration ReadStreamResult)- , _readSECNextNumber :: Required 3 (Value Int32)- , _readSECLastNumber :: Required 4 (Value Int32)- , _readSECEndOfStream :: Required 5 (Value Bool)- , _readSECLastCommitPosition :: Required 6 (Value Int64)- , _readSECError :: Optional 7 (Value Text)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode ReadStreamEventsCompleted------------------------------------------------------------------------------------- | Represents the result of a single read operation to the EventStore.-data StreamEventsSlice- = StreamEventsSlice- { streamEventsSliceResult :: !ReadStreamResult- -- ^ Representing the status of the read attempt.- , streamEventsSliceStreamId :: !Text- -- ^ The name of the stream read.- , streamEventsSliceStart :: !Int32- -- ^ The starting point (represented as a sequence number) of the read- -- operation.- , streamEventsSliceNext :: !Int32- -- ^ The next event number that can be read.- , streamEventsSliceLast :: !Int32- -- ^ The last event number in the stream.- , streamEventsSliceIsEOS :: !Bool- -- ^ Representing whether or not this is the end of the stream.- , streamEventsSliceEvents :: ![ResolvedEvent]- -- ^ The events read represented as 'ResolvedEvent'- , streamEventsSliceDirection :: !ReadDirection- -- ^ The direction of the read request.- }- deriving Show-----------------------------------------------------------------------------------newStreamEventsSlice :: Text- -> Int32- -> ReadDirection- -> ReadStreamEventsCompleted- -> StreamEventsSlice-newStreamEventsSlice stream_id start dir reco = ses- where- evts = getField $ _readSECEvents reco-- ses = StreamEventsSlice- { streamEventsSliceResult = getField $ _readSECResult reco- , streamEventsSliceStreamId = stream_id- , streamEventsSliceStart = start- , streamEventsSliceNext = getField $ _readSECNextNumber reco- , streamEventsSliceLast = getField $ _readSECLastNumber reco- , streamEventsSliceIsEOS = getField $ _readSECEndOfStream reco- , streamEventsSliceEvents = fmap newResolvedEvent evts- , streamEventsSliceDirection = dir- }-----------------------------------------------------------------------------------readStreamEventsOperation :: Settings- -> ReadDirection- -> MVar (OperationExceptional StreamEventsSlice)- -> Text- -> Int32- -> Int32- -> Bool- -> OperationParams-readStreamEventsOperation settings dir mvar stream_id start cnt res_link_tos =- OperationParams- { opSettings = settings- , opRequestCmd = req- , opResponseCmd = resp-- , opRequest =- let req_master = s_requireMaster settings- request = newReadStreamEvents stream_id- start- cnt- res_link_tos- req_master in- return request-- , opSuccess = inspect mvar dir stream_id start- , opFailure = failed mvar- }- where- req = case dir of- Forward -> 0xB2- Backward -> 0xB4-- resp = case dir of- Forward -> 0xB3- Backward -> 0xB5-----------------------------------------------------------------------------------inspect :: MVar (OperationExceptional StreamEventsSlice)- -> ReadDirection- -> Text- -> Int32- -> ReadStreamEventsCompleted- -> IO Decision-inspect mvar dir stream_id start rsec = go (getField $ _readSECResult rsec)- where- may_err = getField $ _readSECError rsec-- go RS_ERROR = failed mvar (ServerError may_err)- go RS_ACCESS_DENIED = failed mvar (AccessDenied stream_id)- go _ = succeed mvar dir stream_id start rsec-----------------------------------------------------------------------------------succeed :: MVar (OperationExceptional StreamEventsSlice)- -> ReadDirection- -> Text- -> Int32- -> ReadStreamEventsCompleted- -> IO Decision-succeed mvar dir stream_id start rsec = do- putMVar mvar (Right ses)- return EndOperation- where- ses = newStreamEventsSlice stream_id start dir rsec-----------------------------------------------------------------------------------failed :: MVar (OperationExceptional StreamEventsSlice)- -> OperationException- -> IO Decision-failed mvar e = do- putMVar mvar (Left e)- return EndOperation
+ Database/EventStore/Internal/Operation/StreamMetadata.hs view
@@ -0,0 +1,89 @@+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE Rank2Types #-}+{-# LANGUAGE RecordWildCards #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.StreamMetadata+-- 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.Internal.Operation.StreamMetadata+ ( readMetaStream+ , setMetaStream+ ) where++--------------------------------------------------------------------------------+import Data.Int+import Data.Monoid ((<>))++--------------------------------------------------------------------------------+import Data.Aeson (decode)+import Data.ByteString.Lazy (fromStrict)+import Data.Text (Text)++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Operation.Read.Common+import Database.EventStore.Internal.Operation.ReadEvent+import Database.EventStore.Internal.Operation.Write.Common+import Database.EventStore.Internal.Operation.WriteEvents+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+metaStream :: Text -> Text+metaStream s = "$$" <> s++--------------------------------------------------------------------------------+-- | Read stream metadata operation.+readMetaStream :: Settings -> Text -> Operation StreamMetadataResult+readMetaStream setts s =+ foreach (readEvent setts (metaStream s) (-1) False) $ \tmp -> do+ onReadResult tmp $ \n e_num evt -> do+ let bytes = recordedEventData $ resolvedEventOriginal evt+ case decode $ fromStrict bytes of+ Just pv -> yield $ StreamMetadataResult n e_num pv+ Nothing -> failure invalidFormat++--------------------------------------------------------------------------------+-- | Set stream metadata operation.+setMetaStream :: Settings+ -> Text+ -> ExpectedVersion+ -> StreamMetadata+ -> Operation WriteResult+setMetaStream setts s v meta =+ let stream = metaStream s+ json = streamMetadataJSON meta+ evt = createEvent "$metadata" Nothing (withJson json)+ inner = writeEvents setts stream v [evt] in+ foreach inner yield++--------------------------------------------------------------------------------+invalidFormat :: OperationError+invalidFormat = InvalidOperation "Invalid metadata format"++--------------------------------------------------------------------------------+streamNotFound :: OperationError+streamNotFound = InvalidOperation "Read metadata on an inexistant stream"++--------------------------------------------------------------------------------+onReadResult :: ReadResult 'RegularStream ReadEvent+ -> (Text -> Int32 -> ResolvedEvent -> SM a b)+ -> SM a b+onReadResult (ReadSuccess r) k =+ case r of+ ReadEvent s n e -> k s n e+ _ -> failure streamNotFound+onReadResult ReadNoStream _ = failure streamNotFound+onReadResult (ReadStreamDeleted s) _ = failure $ StreamDeleted s+onReadResult ReadNotModified _ = failure $ ServerError Nothing+onReadResult (ReadError e) _ = failure $ ServerError e+onReadResult (ReadAccessDenied s) _ = failure $ AccessDenied s
+ Database/EventStore/Internal/Operation/Transaction.hs view
@@ -0,0 +1,109 @@+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE KindSignatures #-}+{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE TypeFamilies #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.Transaction+-- 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.Internal.Operation.Transaction+ ( transactionStart+ , transactionWrite+ , transactionCommit+ ) where++--------------------------------------------------------------------------------+import Data.Int+import Data.Maybe+import Data.Traversable++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Operation.Transaction.Message+import Database.EventStore.Internal.Operation.Write.Common+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+import Prelude++--------------------------------------------------------------------------------+-- | Start transaction operation.+transactionStart :: Settings -> Text -> ExpectedVersion -> Operation Int64+transactionStart Settings{..} stream exp_v = do+ let msg = newStart stream (expVersionInt32 exp_v) s_requireMaster+ resp <- send 0x84 0x85 msg+ let tid = getField $ _transId resp+ r = getField $ _result resp+ case r of+ OP_PREPARE_TIMEOUT -> retry+ OP_FORWARD_TIMEOUT -> retry+ OP_COMMIT_TIMEOUT -> retry+ OP_WRONG_EXPECTED_VERSION -> wrongVersion stream exp_v+ OP_STREAM_DELETED -> streamDeleted stream+ OP_INVALID_TRANSACTION -> invalidTransaction+ OP_ACCESS_DENIED -> accessDenied $ StreamName stream+ OP_SUCCESS -> yield tid++--------------------------------------------------------------------------------+-- | Transactional write operation.+transactionWrite :: Settings+ -> Text+ -> ExpectedVersion+ -> Int64+ -> [Event]+ -> Operation ()+transactionWrite Settings{..} stream exp_v trans_id evts = do+ nevts <- traverse eventToNewEvent evts+ let msg = newWrite trans_id nevts s_requireMaster+ resp <- send 0x86 0x87 msg+ let r = getField $ _wwResult resp+ case r of+ OP_PREPARE_TIMEOUT -> retry+ OP_FORWARD_TIMEOUT -> retry+ OP_COMMIT_TIMEOUT -> retry+ OP_WRONG_EXPECTED_VERSION -> wrongVersion stream exp_v+ OP_STREAM_DELETED -> streamDeleted stream+ OP_INVALID_TRANSACTION -> invalidTransaction+ OP_ACCESS_DENIED -> accessDenied $ StreamName stream+ OP_SUCCESS -> yield ()++--------------------------------------------------------------------------------+-- | Transactional commit operation.+transactionCommit :: Settings+ -> Text+ -> ExpectedVersion+ -> Int64+ -> Operation WriteResult+transactionCommit Settings{..} stream exp_v trans_id = do+ let msg = newCommit trans_id s_requireMaster+ resp <- send 0x88 0x89 msg+ let r = getField $ _ccResult resp+ com_pos = getField $ _commitPosition resp+ pre_pos = getField $ _preparePosition resp+ lst_num = getField $ _lastNumber resp+ p_int = fromMaybe (-1) pre_pos+ c_int = fromMaybe (-1) com_pos+ pos = Position c_int p_int+ res = WriteResult lst_num pos+ case r of+ OP_PREPARE_TIMEOUT -> retry+ OP_FORWARD_TIMEOUT -> retry+ OP_COMMIT_TIMEOUT -> retry+ OP_WRONG_EXPECTED_VERSION -> wrongVersion stream exp_v+ OP_STREAM_DELETED -> streamDeleted stream+ OP_INVALID_TRANSACTION -> invalidTransaction+ OP_ACCESS_DENIED -> accessDenied $ StreamName stream+ OP_SUCCESS -> yield res
+ Database/EventStore/Internal/Operation/Transaction/Message.hs view
@@ -0,0 +1,136 @@+{-# LANGUAGE DeriveGeneric #-}+{-# LANGUAGE DataKinds #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.Transaction.Message+-- 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.Internal.Operation.Transaction.Message where++--------------------------------------------------------------------------------+import Data.Int+import GHC.Generics++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Start transaction request.+data Start =+ Start+ { _streamId :: Required 1 (Value Text)+ , _expectedVersion :: Required 2 (Value Int32)+ , _requireMaster :: Required 3 (Value Bool)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode Start++--------------------------------------------------------------------------------+-- | 'Start' smart constructor.+newStart :: Text -> Int32 -> Bool -> Start+newStart stream_id exp_ver req_master =+ Start+ { _streamId = putField stream_id+ , _expectedVersion = putField exp_ver+ , _requireMaster = putField req_master+ }++--------------------------------------------------------------------------------+-- | Start transaction response.+data Started =+ Started+ { _transId :: Required 1 (Value Int64)+ , _result :: Required 2 (Enumeration OpResult)+ , _message :: Optional 3 (Value Text)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode Started++--------------------------------------------------------------------------------+-- | Write transactional events request.+data Write =+ Write+ { _wTransId :: Required 1 (Value Int64)+ , _events :: Repeated 2 (Message NewEvent)+ , _wRequireMaster :: Required 3 (Value Bool)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode Write++--------------------------------------------------------------------------------+-- | 'Write' smart constructor.+newWrite :: Int64 -> [NewEvent] -> Bool -> Write+newWrite trans_id evts req_master =+ Write+ { _wTransId = putField trans_id+ , _events = putField evts+ , _wRequireMaster = putField req_master+ }++--------------------------------------------------------------------------------+-- | Write transactional events response.+data Written =+ Written+ { _wwTransId :: Required 1 (Value Int64)+ , _wwResult :: Required 2 (Enumeration OpResult)+ , _wwMessage :: Optional 3 (Value Text)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode Written++--------------------------------------------------------------------------------+-- | Commit transaction request.+data Commit =+ Commit+ { _cTransId :: Required 1 (Value Int64)+ , _cRequireMaster :: Required 2 (Value Bool)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode Commit++--------------------------------------------------------------------------------+-- | 'Commit' smart constructor.+newCommit :: Int64 -> Bool -> Commit+newCommit trans_id req_master =+ Commit+ { _cTransId = putField trans_id+ , _cRequireMaster = putField req_master+ }++--------------------------------------------------------------------------------+-- | Commit transaction response.+data Committed =+ Committed+ { _ccTransId :: Required 1 (Value Int64)+ , _ccResult :: Required 2 (Enumeration OpResult)+ , _ccMessage :: Optional 3 (Value Text)+ , _firstNumber :: Required 4 (Value Int32)+ , _lastNumber :: Required 5 (Value Int32)+ , _preparePosition :: Optional 6 (Value Int64)+ , _commitPosition :: Optional 7 (Value Int64)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode Committed
− Database/EventStore/Internal/Operation/TransactionStartOperation.hs
@@ -1,283 +0,0 @@-{-# LANGUAGE RecordWildCards #-}------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Operation.TransactionStartOperation--- 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.Operation.TransactionStartOperation- ( transactionStartOperation ) where-----------------------------------------------------------------------------------import Control.Concurrent-import Control.Exception-import Data.Int-import Data.Maybe-import Data.Traversable-import Prelude-----------------------------------------------------------------------------------import Control.Concurrent.Async-import Data.ProtocolBuffers-import Data.Text-----------------------------------------------------------------------------------import Database.EventStore.Internal.Manager.Operation-import Database.EventStore.Internal.Processor-import Database.EventStore.Internal.Types-----------------------------------------------------------------------------------data TransactionEnv- = TransactionEnv- { _transSettings :: Settings- , _transProcessor :: Cmd -> IO ()- , _transStreamId :: Text- , _transExpectedVersion :: ExpectedVersion- }-----------------------------------------------------------------------------------transactionStartOperation :: Settings- -> (Cmd -> IO ())- -> MVar (OperationExceptional Transaction)- -> Text- -> ExpectedVersion- -> OperationParams-transactionStartOperation settings procss mvar stream_id exp_ver =- OperationParams- { opSettings = settings- , opRequestCmd = 0x84- , opResponseCmd = 0x85-- , opRequest =- let req_master = s_requireMaster settings- exp_ver_int = expVersionInt32 exp_ver- request = newTransactionStart stream_id- exp_ver_int- req_master in- return request-- , opSuccess = inspectTrans env mvar- , opFailure = failed mvar- }- where- env = TransactionEnv- { _transSettings = settings- , _transProcessor = procss- , _transStreamId = stream_id- , _transExpectedVersion = exp_ver- }-----------------------------------------------------------------------------------inspectTrans :: TransactionEnv- -> MVar (OperationExceptional Transaction)- -> TransactionStartCompleted- -> IO Decision-inspectTrans env mvar tsc = go (getField $ transactionSCResult tsc)- where- go OP_SUCCESS = succeedTrans env mvar tsc- go OP_PREPARE_TIMEOUT = return Retry- go OP_FORWARD_TIMEOUT = return Retry- go OP_COMMIT_TIMEOUT = return Retry- go OP_WRONG_EXPECTED_VERSION = failed mvar wrong_version- go OP_STREAM_DELETED = failed mvar (StreamDeleted stream_id)- go OP_INVALID_TRANSACTION = failed mvar InvalidTransaction- go OP_ACCESS_DENIED = failed mvar (AccessDenied stream_id)-- exp_ver = _transExpectedVersion env- stream_id = _transStreamId env- wrong_version = WrongExpectedVersion stream_id exp_ver-----------------------------------------------------------------------------------succeedTrans :: TransactionEnv- -> MVar (OperationExceptional Transaction)- -> TransactionStartCompleted- -> IO Decision-succeedTrans env mvar tsc = do- putMVar mvar (Right trans)- return EndOperation- where- trans_id = getField $ transactionSCId tsc- trans = createTransaction env trans_id-----------------------------------------------------------------------------------failed :: MVar (OperationExceptional a)- -> OperationException- -> IO Decision-failed mvar e = do- putMVar mvar (Left e)- return EndOperation-----------------------------------------------------------------------------------createTransaction :: TransactionEnv -> Int64 -> Transaction-createTransaction env@TransactionEnv{..} trans_id = trans- where- trans = Transaction- { transactionId = trans_id- , transactionStreamId = _transStreamId- , transactionExpectedVersion = _transExpectedVersion-- , transactionCommit = do- (as, mvar) <- createAsync-- let op = transactionCommitOperation env trans_id mvar-- _transProcessor (NewOperation op)- return as-- , transactionSendEvents = \evts -> do- (as, mvar) <- createAsync-- let op = transactionWriteOperation env trans_id mvar evts-- _transProcessor (NewOperation op)- return as-- , transactionRollback = return ()- }-----------------------------------------------------------------------------------transactionWriteOperation :: TransactionEnv- -> Int64- -> MVar (OperationExceptional ())- -> [Event]- -> OperationParams-transactionWriteOperation env trans_id mvar evts =- OperationParams- { opSettings = settings- , opRequestCmd = 0x86- , opResponseCmd = 0x87-- , opRequest = do- new_evts <- traverse eventToNewEvent evts-- let request = newTransactionWrite trans_id- new_evts- req_master-- return request-- , opSuccess = inspectWrite env mvar- , opFailure = failed mvar- }- where- settings = _transSettings env- req_master = s_requireMaster settings-----------------------------------------------------------------------------------transactionCommitOperation :: TransactionEnv- -> Int64- -> MVar (OperationExceptional WriteResult)- -> OperationParams-transactionCommitOperation env trans_id mvar =- OperationParams- { opSettings = settings- , opRequestCmd = 0x88- , opResponseCmd = 0x89-- , opRequest =- let request = newTransactionCommit trans_id req_master in-- return request-- , opSuccess = inspectCommit env mvar- , opFailure = failed mvar- }- where- settings = _transSettings env- req_master = s_requireMaster settings-----------------------------------------------------------------------------------inspectWrite :: TransactionEnv- -> MVar (OperationExceptional ())- -> TransactionWriteCompleted- -> IO Decision-inspectWrite env mvar twc = go (getField $ transactionWCResult twc)- where- go OP_SUCCESS = succeedWrite mvar twc- go OP_PREPARE_TIMEOUT = return Retry- go OP_FORWARD_TIMEOUT = return Retry- go OP_COMMIT_TIMEOUT = return Retry- go OP_WRONG_EXPECTED_VERSION = failed mvar wrong_version- go OP_STREAM_DELETED = failed mvar (StreamDeleted stream_id)- go OP_INVALID_TRANSACTION = failed mvar InvalidTransaction- go OP_ACCESS_DENIED = failed mvar (AccessDenied stream_id)-- exp_ver = _transExpectedVersion env- stream_id = _transStreamId env- wrong_version = WrongExpectedVersion stream_id exp_ver-----------------------------------------------------------------------------------succeedWrite :: MVar (OperationExceptional ())- -> TransactionWriteCompleted- -> IO Decision-succeedWrite mvar _ = do- putMVar mvar (Right ())- return EndOperation-----------------------------------------------------------------------------------inspectCommit :: TransactionEnv- -> MVar (OperationExceptional WriteResult)- -> TransactionCommitCompleted- -> IO Decision-inspectCommit env mvar tcc = go (getField $ transactionCCResult tcc)- where- go OP_SUCCESS = succeedCommit mvar tcc- go OP_PREPARE_TIMEOUT = return Retry- go OP_FORWARD_TIMEOUT = return Retry- go OP_COMMIT_TIMEOUT = return Retry- go OP_WRONG_EXPECTED_VERSION = failed mvar wrong_version- go OP_STREAM_DELETED = failed mvar (StreamDeleted stream_id)- go OP_INVALID_TRANSACTION = failed mvar InvalidTransaction- go OP_ACCESS_DENIED = failed mvar (AccessDenied stream_id)-- exp_ver = _transExpectedVersion env- stream_id = _transStreamId env- wrong_version = WrongExpectedVersion stream_id exp_ver-----------------------------------------------------------------------------------succeedCommit :: MVar (OperationExceptional WriteResult)- -> TransactionCommitCompleted- -> IO Decision-succeedCommit mvar tcc = do- putMVar mvar (Right wr)- return EndOperation- where- last_evt_num = getField $ transactionCCLastNumber tcc- com_pos = getField $ transactionCCCommitPosition tcc- pre_pos = getField $ transactionCCPreparePosition tcc- com_pos_int = fromMaybe (-1) com_pos- pre_pos_int = fromMaybe (-1) pre_pos- pos = Position com_pos_int pre_pos_int- wr = WriteResult last_evt_num pos-----------------------------------------------------------------------------------eventToNewEvent :: Event -> IO NewEvent-eventToNewEvent evt =- newEvent evt_type- evt_id- evt_data_type- evt_metadata_type- evt_data_bytes- evt_metadata_bytes- where- evt_type = eventType evt- evt_id = eventId evt- evt_data_bytes = eventDataBytes $ eventData evt- evt_data_type = eventDataType $ eventData evt- evt_metadata_bytes = eventMetadataBytes $ eventData evt- evt_metadata_type = eventMetadataType $ eventData evt-----------------------------------------------------------------------------------createAsync :: IO (Async a, MVar (OperationExceptional a))-createAsync = do- mvar <- newEmptyMVar- as <- async $ do- res <- readMVar mvar- either throwIO return res-- return (as, mvar)
+ Database/EventStore/Internal/Operation/Write/Common.hs view
@@ -0,0 +1,49 @@+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.Write.Common+-- 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.Internal.Operation.Write.Common where++--------------------------------------------------------------------------------+import Data.Int++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Returned after writing to a stream.+data WriteResult+ = WriteResult+ { writeNextExpectedVersion :: !Int32+ -- ^ Next expected version of the stream.+ , writePosition :: !Position+ -- ^ 'Position' of the write.+ }+ deriving (Eq, Show)++--------------------------------------------------------------------------------+-- | Constructs a 'NewEvent' from an 'Event'.+eventToNewEvent :: Event -> SM a NewEvent+eventToNewEvent evt = do+ uuid <- maybe freshId return evt_id+ return $ newEvent evt_type+ uuid+ evt_data_type+ evt_metadata_type+ evt_data_bytes+ evt_metadata_bytes+ where+ evt_type = eventType evt+ evt_id = eventId evt+ evt_data_bytes = eventDataBytes $ eventData evt+ evt_data_type = eventDataType $ eventData evt+ evt_metadata_bytes = eventMetadataBytes $ eventData evt+ evt_metadata_type = eventMetadataType $ eventData evt
+ Database/EventStore/Internal/Operation/WriteEvents.hs view
@@ -0,0 +1,61 @@+{-# LANGUAGE RecordWildCards #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.WriteEvents+-- 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.Internal.Operation.WriteEvents+ ( writeEvents ) where++--------------------------------------------------------------------------------+import Data.Maybe+import Data.Traversable++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Operation.Write.Common+import Database.EventStore.Internal.Operation.WriteEvents.Message+import Database.EventStore.Internal.Stream+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+import Prelude++--------------------------------------------------------------------------------+-- | Write events operation.+writeEvents :: Settings+ -> Text+ -> ExpectedVersion+ -> [Event]+ -> Operation WriteResult+writeEvents Settings{..} s v evts = do+ nevts <- traverse eventToNewEvent evts+ let msg = newRequest s (expVersionInt32 v) nevts s_requireMaster+ resp <- send 0x82 0x83 msg+ let r = getField $ _result resp+ com_pos = getField $ _commitPosition resp+ prep_pos = getField $ _preparePosition resp+ lst_num = getField $ _lastNumber resp+ com_pos_int = fromMaybe (-1) com_pos+ prep_pos_int = fromMaybe (-1) prep_pos+ pos = Position com_pos_int prep_pos_int+ res = WriteResult lst_num pos+ case r of+ OP_SUCCESS -> yield res+ OP_PREPARE_TIMEOUT -> retry+ OP_FORWARD_TIMEOUT -> retry+ OP_COMMIT_TIMEOUT -> retry+ OP_WRONG_EXPECTED_VERSION -> wrongVersion s v+ OP_STREAM_DELETED -> streamDeleted s+ OP_INVALID_TRANSACTION -> invalidTransaction+ OP_ACCESS_DENIED -> accessDenied (StreamName s)
+ Database/EventStore/Internal/Operation/WriteEvents/Message.hs view
@@ -0,0 +1,71 @@+{-# LANGUAGE DeriveGeneric #-}+{-# LANGUAGE DataKinds #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operation.WriteEvents.Message+-- 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.Internal.Operation.WriteEvents.Message where++--------------------------------------------------------------------------------+import Data.Int+import GHC.Generics++--------------------------------------------------------------------------------+import Data.ProtocolBuffers+import Data.Text++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation+import Database.EventStore.Internal.Types++--------------------------------------------------------------------------------+-- | Write events request.+data Request+ = Request+ { _streamId :: Required 1 (Value Text)+ , _expectedVersion :: Required 2 (Value Int32)+ , _events :: Repeated 3 (Message NewEvent)+ , _requireMaster :: Required 4 (Value Bool)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Encode Request++--------------------------------------------------------------------------------+-- | 'Request' smart constructor.+newRequest :: Text -- ^ Stream+ -> Int32 -- ^ Expected version+ -> [NewEvent] -- ^ Events+ -> Bool -- ^ Require master+ -> Request+newRequest stream_id exp_ver evts req_master =+ Request+ { _streamId = putField stream_id+ , _expectedVersion = putField exp_ver+ , _events = putField evts+ , _requireMaster = putField req_master+ }++--------------------------------------------------------------------------------+-- | Write events response.+data Response+ = Response+ { _result :: Required 1 (Enumeration OpResult)+ , _message :: Optional 2 (Value Text)+ , _firstNumber :: Required 3 (Value Int32)+ , _lastNumber :: Required 4 (Value Int32)+ , _preparePosition :: Optional 5 (Value Int64)+ , _commitPosition :: Optional 6 (Value Int64)+ }+ deriving (Generic, Show)++--------------------------------------------------------------------------------+instance Decode Response
− Database/EventStore/Internal/Operation/WriteEventsOperation.hs
@@ -1,161 +0,0 @@-{-# LANGUAGE DeriveGeneric #-}-{-# LANGUAGE DataKinds #-}------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Operation.WriteEventsOperation--- 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.Operation.WriteEventsOperation- ( writeEventsOperation ) where-----------------------------------------------------------------------------------import Control.Concurrent-import Data.Int-import Data.Maybe-import Data.Traversable-import GHC.Generics (Generic)-import Prelude-----------------------------------------------------------------------------------import Data.ProtocolBuffers-import Data.Text-----------------------------------------------------------------------------------import Database.EventStore.Internal.Manager.Operation-import Database.EventStore.Internal.Types-----------------------------------------------------------------------------------data WriteEvents- = WriteEvents- { _writeStreamId :: Required 1 (Value Text)- , _writeExpectedVersion :: Required 2 (Value Int32)- , _writeEvents :: Repeated 3 (Message NewEvent)- , _writeRequireMaster :: Required 4 (Value Bool)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Encode WriteEvents-----------------------------------------------------------------------------------newWriteEvents :: Text -- ^ Stream- -> Int32 -- ^ Expected version- -> [NewEvent] -- ^ Events- -> Bool -- ^ Require master- -> WriteEvents-newWriteEvents stream_id exp_ver evts req_master =- WriteEvents- { _writeStreamId = putField stream_id- , _writeExpectedVersion = putField exp_ver- , _writeEvents = putField evts- , _writeRequireMaster = putField req_master- }-----------------------------------------------------------------------------------data WriteEventsCompleted- = WriteEventsCompleted- { _writeCompletedResult :: Required 1 (Enumeration OpResult)- , _writeCompletedMessage :: Optional 2 (Value Text)- , _writeCompletedFirstNumber :: Required 3 (Value Int32)- , _writeCompletedLastNumber :: Required 4 (Value Int32)- , _writeCompletedPreparePosition :: Optional 5 (Value Int64)- , _writeCompletedCommitPosition :: Optional 6 (Value Int64)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode WriteEventsCompleted-----------------------------------------------------------------------------------writeEventsOperation :: Settings- -> MVar (OperationExceptional WriteResult)- -> Text- -> ExpectedVersion- -> [Event]- -> OperationParams-writeEventsOperation settings mvar evt_stream exp_ver evts =- OperationParams- { opSettings = settings- , opRequestCmd = 0x82- , opResponseCmd = 0x83-- , opRequest = do- new_evts <- traverse eventToNewEvent evts-- let require_master = s_requireMaster settings- exp_ver_int32 = expVersionInt32 exp_ver- request = newWriteEvents evt_stream- exp_ver_int32- new_evts- require_master- return request-- , opSuccess = inspect mvar evt_stream exp_ver- , opFailure = failed mvar- }-----------------------------------------------------------------------------------inspect :: MVar (OperationExceptional WriteResult)- -> Text- -> ExpectedVersion- -> WriteEventsCompleted- -> IO Decision-inspect mvar stream exp_ver wec = go (getField $ _writeCompletedResult wec)- where- go OP_SUCCESS = succeed mvar wec- go OP_PREPARE_TIMEOUT = return Retry- go OP_FORWARD_TIMEOUT = return Retry- go OP_COMMIT_TIMEOUT = return Retry- go OP_WRONG_EXPECTED_VERSION = failed mvar wrong_version- go OP_STREAM_DELETED = failed mvar (StreamDeleted stream)- go OP_INVALID_TRANSACTION = failed mvar InvalidTransaction- go OP_ACCESS_DENIED = failed mvar (AccessDenied stream)-- wrong_version = WrongExpectedVersion stream exp_ver-----------------------------------------------------------------------------------succeed :: MVar (OperationExceptional WriteResult)- -> WriteEventsCompleted- -> IO Decision-succeed mvar wec = do- putMVar mvar (Right wr)- return EndOperation- where- last_evt_num = getField $ _writeCompletedLastNumber wec- com_pos = getField $ _writeCompletedCommitPosition wec- pre_pos = getField $ _writeCompletedPreparePosition wec- com_pos_int = fromMaybe (-1) com_pos- pre_pos_int = fromMaybe (-1) pre_pos- pos = Position com_pos_int pre_pos_int- wr = WriteResult last_evt_num pos-----------------------------------------------------------------------------------failed :: MVar (OperationExceptional WriteResult)- -> OperationException- -> IO Decision-failed mvar e = do- putMVar mvar (Left e)- return EndOperation-----------------------------------------------------------------------------------eventToNewEvent :: Event -> IO NewEvent-eventToNewEvent evt =- newEvent evt_type- evt_id- evt_data_type- evt_metadata_type- evt_data_bytes- evt_metadata_bytes- where- evt_type = eventType evt- evt_id = eventId evt- evt_data_bytes = eventDataBytes $ eventData evt- evt_data_type = eventDataType $ eventData evt- evt_metadata_bytes = eventMetadataBytes $ eventData evt- evt_metadata_type = eventMetadataType $ eventData evt
+ Database/EventStore/Internal/Operations.hs view
@@ -0,0 +1,32 @@+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Operations+-- Copyright : (C) 2015 Yorick Laupa+-- License : (see the file LICENSE)+--+-- Maintainer : Yorick Laupa <yo.eight@gmail.com>+-- Stability : provisional+-- Portability : non-portable+--+-- Mega module to easily import operation in the main EventStore module.+--------------------------------------------------------------------------------+module Database.EventStore.Internal.Operations+ ( module Database.EventStore.Internal.Operation.Catchup+ , module Database.EventStore.Internal.Operation.DeleteStream+ , module Database.EventStore.Internal.Operation.ReadAllEvents+ , module Database.EventStore.Internal.Operation.ReadEvent+ , module Database.EventStore.Internal.Operation.ReadStreamEvents+ , module Database.EventStore.Internal.Operation.StreamMetadata+ , module Database.EventStore.Internal.Operation.Transaction+ , module Database.EventStore.Internal.Operation.WriteEvents+ ) where++--------------------------------------------------------------------------------+import Database.EventStore.Internal.Operation.Catchup+import Database.EventStore.Internal.Operation.DeleteStream+import Database.EventStore.Internal.Operation.ReadAllEvents+import Database.EventStore.Internal.Operation.ReadEvent+import Database.EventStore.Internal.Operation.ReadStreamEvents+import Database.EventStore.Internal.Operation.StreamMetadata+import Database.EventStore.Internal.Operation.Transaction+import Database.EventStore.Internal.Operation.WriteEvents
Database/EventStore/Internal/Packages.hs view
@@ -11,8 +11,7 @@ -------------------------------------------------------------------------------- module Database.EventStore.Internal.Packages ( -- * Package Smart Contructors- heartbeatPackage- , heartbeatResponsePackage+ heartbeatResponsePackage -- * Cereal Put , putPackage ) where@@ -24,7 +23,6 @@ -------------------------------------------------------------------------------- import Data.Serialize.Put import Data.UUID-import System.Random -------------------------------------------------------------------------------- import Database.EventStore.Internal.Types@@ -32,19 +30,7 @@ -------------------------------------------------------------------------------- -- Encode ---------------------------------------------------------------------------------heartbeatPackage :: IO Package-heartbeatPackage = do- uuid <- randomIO- let pack = Package- { packageCmd = 0x01- , packageCorrelation = uuid- , packageData = B.empty- , packageCred = Nothing- }-- return pack----------------------------------------------------------------------------------+-- | Constructs a heartbeat response given the 'UUID' of heartbeat request. heartbeatResponsePackage :: UUID -> Package heartbeatResponsePackage uuid = Package@@ -55,6 +41,7 @@ } --------------------------------------------------------------------------------+-- | Serializes a 'Package' into raw bytes. putPackage :: Package -> Put putPackage pack = do putWord32le length_prefix@@ -80,5 +67,7 @@ credSize (Credentials login passw) = B.length login + B.length passw + 2 --------------------------------------------------------------------------------+-- | The minimun size a 'Package' should have. It's basically a command byte,+-- correlation bytes ('UUID') and a 'Flag' byte. mandatorySize :: Int mandatorySize = 18
Database/EventStore/Internal/Processor.hs view
@@ -1,5 +1,6 @@-{-# LANGUAGE RecordWildCards #-}-{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE ExistentialQuantification #-}+{-# LANGUAGE RecordWildCards #-} -------------------------------------------------------------------------------- -- | -- Module : Database.EventStore.Internal.Processor@@ -10,274 +11,313 @@ -- Stability : provisional -- Portability : non-portable --+-- Top level operation and subscription logic of EventStore driver. -------------------------------------------------------------------------------- module Database.EventStore.Internal.Processor- ( ConnectionException(..)- , InternalException(..)- , Cmd(..)+ ( Processor+ , Transition(..) , newProcessor+ , connectRegularStream+ , connectPersistent+ , createPersistent+ , updatePersistent+ , deletePersistent+ , ackPersist+ , nakPersist+ , newOperation+ , submitPackage+ , unsubscribe+ , abort ) where ---------------------------------------------------------------------------------import Control.Concurrent-import Control.Exception-import Data.Functor (void) import Data.Int-import Data.Monoid ((<>))-import Data.Word ---------------------------------------------------------------------------------import Data.Text (Text)+import Data.Text import Data.UUID-import FRP.Sodium-import Network ---------------------------------------------------------------------------------import Database.EventStore.Internal.Connection-import Database.EventStore.Internal.Manager.Operation-import Database.EventStore.Internal.Manager.Subscription+import Database.EventStore.Internal.Generator+import Database.EventStore.Internal.Operation hiding (SM(..)) import Database.EventStore.Internal.Packages-import Database.EventStore.Internal.Reader-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----------------------------------------------------------------------------------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+import Database.EventStore.Internal.Types ----------------------------------------------------------------------------------- State----------------------------------------------------------------------------------data State- = Offline- | Online- { _uuidCon :: !UUID- , _packageCount :: !Int- , _host :: !HostName- , _port :: !Int- , _cleanup :: !(IO ())- }+import qualified Database.EventStore.Internal.Manager.Operation.Model as Op+import qualified Database.EventStore.Internal.Manager.Subscription as Sub ----------------------------------------------------------------------------------- Event----------------------------------------------------------------------------------data Connect = Connect HostName Int-data Connected = Connected HostName Int UUID (IO ())-data Cleanup = Cleanup-data Reconnect = Reconnect-data Reconnected = Reconnected UUID (IO ())+-- | Type of inputs handled by the 'Processor' driver.+data In r+ = Cmd (Cmd r)+ -- ^ A command can be an 'Operation' or a 'Subscription' actions.+ | Pkg Package+ -- ^ Handle a 'Package' coming from the server. ---------------------------------------------------------------------------------heartbeatRequestCmd :: Word8-heartbeatRequestCmd = 0x01+-- | Type of commmand a 'Processor' can handle.+data Cmd r+ = SubscriptionCmd (SubscriptionCmd r)+ -- ^ Subcription related commands.+ | forall a. NewOp (Operation a) (Either OperationError a -> r)+ -- ^ Register a new 'Operation'.+ | Abort+ -- ^ Aborts every pending operation. ---------------------------------------------------------------------------------network :: Settings -> Chan Package -> Reactive Processor-network sett chan = do- (onConnect, pushConnect) <- newEvent- (onConnected, pushConnected) <- newEvent- (onCleanup, pushCleanup) <- newEvent- (onReconnect, pushReconnect) <- newEvent- (onReconnected, pushReconnected) <- newEvent- (onReceived, pushReceived) <- newEvent- (onSend, pushSend) <- newEvent+-- | Supported subscription command.+data SubscriptionCmd r+ = ConnectStream (Sub.SubConnectEvent -> r) Text Bool+ -- ^ Creates a regular subscription connection.+ | ConnectPersist (Sub.SubConnectEvent -> r) Text Text Int32+ -- ^ Creates a persistent subscription connection. - push_new_op <- operationNetwork sett- pushSend- (pushReconnect Reconnect)- onReceived+ | CreatePersist (Either Sub.PersistActionException Sub.ConfirmedAction -> r)+ Text+ Text+ PersistentSubscriptionSettings+ -- ^ Creates a persistent subscription. - runSubCmd <- subscriptionNetwork sett pushSend onReceived+ | Unsubscribe Sub.Running+ -- ^ Unsubscribes a subscription. - let stateE = fmap connected onConnected <>- fmap reconnected onReconnected <>- fmap received onReceived+ | UpdatePersist (Either Sub.PersistActionException Sub.ConfirmedAction -> r)+ Text+ Text+ PersistentSubscriptionSettings+ -- ^ Updates a persistent subscription. - stateB <- accum Offline stateE+ | DeletePersist (Either Sub.PersistActionException Sub.ConfirmedAction -> r)+ Text+ Text+ -- ^ Deletes a persistent subscription. - let heartbeatP pkg = packageCmd pkg == heartbeatRequestCmd- onlyHeartbeats = filterE heartbeatP onReceived+ | AckPersist r Sub.Running [UUID]+ -- ^ Acknowledges a set of events has been successfully handled. - con_snap = fmap connectSnapshot onConnect- reco_snap = snapshot reconnectSnapshot onReconnect stateB- clean_snap = snapshot cleanupSnapshot onCleanup stateB+ | NakPersist r Sub.Running Sub.NakAction (Maybe Text) [UUID]+ -- ^ Acknowledges a set of events hasn't been handled successfully. - full_reco c = do- pushCleanup Cleanup- pushReconnect c+--------------------------------------------------------------------------------+-- | Creates a regular subscription connection.+connectRegularStream :: (Sub.SubConnectEvent -> r)+ -> Text -- ^ Stream name.+ -> Bool -- ^ Resolve Link TOS.+ -> Processor r+ -> Transition r+connectRegularStream c s tos (Processor k) =+ k $ Cmd $ SubscriptionCmd $ ConnectStream c s tos - push_reco_io = pushAsync full_reco Reconnect- push_recv_io = \pkg -> sync $ pushReceived pkg- push_recod_io = pushAsync2 $ \u c -> pushReconnected $ Reconnected u c- push_send_io = pushAsync pushSend- push_con_io = pushAsync4 $ \h p u c ->- pushConnected $ Connected h p u c+--------------------------------------------------------------------------------+-- | Creates a persistent subscription connection.+connectPersistent :: (Sub.SubConnectEvent -> r)+ -> Text -- ^ Group name.+ -> Text -- ^ Stream name.+ -> Int32 -- ^ Buffer size.+ -> Processor r+ -> Transition r+connectPersistent c g s siz (Processor k) =+ k $ Cmd $ SubscriptionCmd $ ConnectPersist c g s siz - _ <- listen con_snap $ \(ConnectionSnapshot host port) ->- connection sett- chan- push_recv_io- (push_con_io host port)- push_reco_io- host- port+--------------------------------------------------------------------------------+-- | Creates a persistent subscription.+createPersistent :: (Either Sub.PersistActionException Sub.ConfirmedAction -> r)+ -> Text -- ^ Group name.+ -> Text -- ^ Stream name.+ -> PersistentSubscriptionSettings+ -> Processor r+ -> Transition r+createPersistent c g s sett (Processor k) =+ k $ Cmd $ SubscriptionCmd $ CreatePersist c g s sett - _ <- listen reco_snap $ \(ConnectionSnapshot host port) ->- connection sett- chan- push_recv_io- push_recod_io- push_reco_io- host- port+--------------------------------------------------------------------------------+-- | Updates a persistent subscription.+updatePersistent :: (Either Sub.PersistActionException Sub.ConfirmedAction -> r)+ -> Text -- ^ Group name.+ -> Text -- ^ Stream name.+ -> PersistentSubscriptionSettings+ -> Processor r+ -> Transition r+updatePersistent c g s sett (Processor k) =+ k $ Cmd $ SubscriptionCmd $ UpdatePersist c g s sett - _ <- listen clean_snap $ \(CleanupSnapshot finalizer) -> finalizer+--------------------------------------------------------------------------------+-- | Deletes a persistent subscription.+deletePersistent :: (Either Sub.PersistActionException Sub.ConfirmedAction -> r)+ -> Text -- ^ Group name.+ -> Text -- ^ Stream name.+ -> Processor r+ -> Transition r+deletePersistent c g s (Processor k) =+ k $ Cmd $ SubscriptionCmd $ DeletePersist c g s - _ <- listen onlyHeartbeats $ \pkg ->- push_send_io $ heartbeatResponsePackage (packageCorrelation pkg)+--------------------------------------------------------------------------------+-- | Acknowledges a set of events has been successfully handled.+ackPersist :: r -> Sub.Running -> [UUID] -> Processor r -> Transition r+ackPersist r run evts (Processor k) =+ k $ Cmd $ SubscriptionCmd $ AckPersist r run evts - 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)+--------------------------------------------------------------------------------+-- | Acknowledges a set of events hasn't been handled successfully.+nakPersist :: r+ -> Sub.Running+ -> Sub.NakAction+ -> Maybe Text+ -> [UUID]+ -> Processor r+ -> Transition r+nakPersist r run act res evts (Processor k) =+ k $ Cmd $ SubscriptionCmd $ NakPersist r run act res evts - _ <- listen onSend (writeChan chan)+--------------------------------------------------------------------------------+-- | Registers a new 'Operation'.+newOperation :: (Either OperationError a -> r)+ -> Operation a+ -> Processor r+ -> Transition r+newOperation c op (Processor k) = k $ Cmd $ NewOp op c - return runCmd+--------------------------------------------------------------------------------+-- | Submits a 'Package'.+submitPackage :: Package -> Processor r -> Transition r+submitPackage pkg (Processor k) = k $ Pkg pkg ----------------------------------------------------------------------------------- Observer+-- | Unsubscribes a subscription.+unsubscribe :: Sub.Running -> Processor r -> Transition r+unsubscribe r (Processor k) = k $ Cmd $ SubscriptionCmd $ Unsubscribe r+ ---------------------------------------------------------------------------------data ConnectionSnapshot- = ConnectionSnapshot- { _conHost :: !HostName- , _conPort :: !Int- }+-- | Aborts every pending operation.+abort :: Processor r -> Transition r+abort (Processor k) = k $ Cmd Abort ---------------------------------------------------------------------------------connectSnapshot :: Connect -> ConnectionSnapshot-connectSnapshot (Connect host port) =- ConnectionSnapshot- { _conHost = host- , _conPort = port+-- | 'Processor' internal state.+data State r =+ State+ { _subDriver :: Sub.Driver r+ -- ^ Subscription driver.+ , _opModel :: Op.Model r+ -- ^ Operation model. } ---------------------------------------------------------------------------------connection :: Settings- -> Chan Package- -> (Package -> IO ())- -> (UUID -> IO () -> IO ())- -> IO ()- -> HostName- -> Int- -> IO ()-connection sett chan push_pkg push_con push_reco host port = do- conn <- newConnection sett host port- 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- _settingsLog sett (Info $ Disconnected $ connUUID conn)+initState :: Settings -> Generator -> State r+initState setts g = State (Sub.newDriver setts g1) (Op.newModel setts g2)+ where+ (g1, g2) = splitGenerator g ---------------------------------------------------------------------------------recovering :: IO () -> Either SomeException () -> IO ()-recovering recover (Left some_ex) = do- case fromException some_ex of- Just e ->- case e of- ConnectionClosedByServer- -> recover- Stopped- -> return ()- _ -> recover-recovering _ _ = return ()+-- | Represents the state transition of 'Processor' state machine.+data Transition r+ = Produce r (Transition r)+ -- ^ Produces a final value.+ | Transmit Package (Transition r)+ -- ^ Indicates to send the given 'Package'.+ | Await (Processor r)+ -- ^ Waits for more input. ---------------------------------------------------------------------------------reconnectSnapshot :: Reconnect -> State -> ConnectionSnapshot-reconnectSnapshot _ s =- ConnectionSnapshot- { _conHost = _host s- , _conPort = _port s- }+-- | Processor state-machine.+newtype Processor r = Processor (In r -> Transition r) ---------------------------------------------------------------------------------newtype CleanupSnapshot = CleanupSnapshot (IO ())+loopOpTransition :: State r -> Op.Transition r -> Transition r+loopOpTransition st (Op.Produce r nxt) =+ Produce r (loopOpTransition st nxt)+loopOpTransition st (Op.Transmit pkg nxt) =+ Transmit pkg (loopOpTransition st nxt)+loopOpTransition st (Op.Await m) =+ let nxt_st = st { _opModel = m } in Await $ Processor $ handle nxt_st ---------------------------------------------------------------------------------cleanupSnapshot :: Cleanup -> State -> CleanupSnapshot-cleanupSnapshot _ s =- case s of- Offline {}- -> CleanupSnapshot (return ())- Online {}- -> CleanupSnapshot $ _cleanup s+abortTransition :: State r -> Op.Transition r -> [r] -> Transition r+abortTransition st init_op init_rs = abortOp init_op+ where+ abortOp (Op.Produce r nxt) = Produce r (abortOp nxt)+ abortOp (Op.Transmit _ nxt) = abortOp nxt+ abortOp _ = abortSub init_rs ------------------------------------------------------------------------------------ Model----------------------------------------------------------------------------------connected :: Connected -> State -> State-connected (Connected host port uuid cl) s =- case s of- Offline- -> Online- { _uuidCon = uuid- , _packageCount = 0- , _host = host- , _port = port- , _cleanup = cl- }- _ -> s+ abortSub [] = Await $ Processor $ handle st+ abortSub (r:rs) = Produce r (abortSub rs) ---------------------------------------------------------------------------------reconnected :: Reconnected -> State -> State-reconnected (Reconnected uuid cl) s =- case s of- Online {}- -> s { _uuidCon = uuid- , _cleanup = cl- }- _ -> s+handle :: State r -> In r -> Transition r+handle = go+ where+ go st (Cmd tpe) =+ case tpe of+ NewOp op cb ->+ let sm = Op.pushOperation cb op $ _opModel st in+ loopOpTransition st sm+ SubscriptionCmd cmd -> subCmd st cmd+ Abort ->+ let sm = Op.abort $ _opModel st+ rs = Sub.abort $ _subDriver st in+ abortTransition st sm rs + go st (Pkg pkg)+ | packageCmd pkg == 0x01 =+ let r_pkg = heartbeatResponsePackage $ packageCorrelation pkg in+ Transmit r_pkg $ Await $ Processor $ go st+ | otherwise =+ let sm_m = Op.submitPackage pkg $ _opModel st in+ case fmap (loopOpTransition st) sm_m of+ Just nxt -> nxt+ Nothing ->+ case Sub.submitPackage pkg $ _subDriver st of+ Nothing -> Await $ Processor $ go st+ Just (r, nxt_drv) ->+ let nxt_st = st { _subDriver = nxt_drv } in+ Produce r $ Await $ Processor $ go nxt_st++ subCmd st@State{..} cmd =+ case cmd of+ ConnectStream k s tos ->+ let (pkg, nxt_drv) = Sub.connectToStream k s tos _subDriver+ nxt_st = st { _subDriver = nxt_drv }+ nxt = Processor $ go nxt_st in+ Transmit pkg $ Await nxt+ ConnectPersist k g s b ->+ let (pkg, nxt_drv) = Sub.connectToPersist k g s b _subDriver+ nxt_st = st { _subDriver = nxt_drv }+ nxt = Processor $ go nxt_st in+ Transmit pkg $ Await nxt+ Unsubscribe r ->+ let (pkg, nxt_drv) = Sub.unsubscribe r _subDriver+ nxt_st = st { _subDriver = nxt_drv }+ nxt = Processor $ go nxt_st in+ Transmit pkg $ Await nxt+ CreatePersist k g s ss ->+ let (pkg, nxt_drv) = Sub.createPersist k g s ss _subDriver+ nxt_st = st { _subDriver = nxt_drv }+ nxt = Processor $ go nxt_st in+ Transmit pkg $ Await nxt+ UpdatePersist k g s ss ->+ let (pkg, nxt_drv) = Sub.updatePersist k g s ss _subDriver+ nxt_st = st { _subDriver = nxt_drv }+ nxt = Processor $ go nxt_st in+ Transmit pkg $ Await nxt+ DeletePersist k g s ->+ let (pkg, nxt_drv) = Sub.deletePersist k g s _subDriver+ nxt_st = st { _subDriver = nxt_drv }+ nxt = Processor $ go nxt_st in+ Transmit pkg $ Await nxt+ AckPersist r run evts ->+ let (pkg, nxt_drv) = Sub.ackPersist r run evts _subDriver+ nxt_st = st { _subDriver = nxt_drv }+ nxt = Processor $ go nxt_st in+ Transmit pkg $ Await nxt+ NakPersist r run act res evts ->+ let (pkg, nxt_drv) = Sub.nakPersist r run act res evts+ _subDriver+ nxt_st = st { _subDriver = nxt_drv }+ nxt = Processor $ go nxt_st in+ Transmit pkg $ Await nxt+ ---------------------------------------------------------------------------------received :: Package -> State -> State-received _ s =- case s of- Online {}- -> let cnt = _packageCount s in s { _packageCount = cnt + 1 }- _ -> s+-- | Creates a new 'Processor' state-machine.+newProcessor :: Settings -> Generator -> Processor r+newProcessor setts gen = Processor $ handle $ initState setts gen
− Database/EventStore/Internal/Reader.hs
@@ -1,95 +0,0 @@------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Reader--- 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.Reader (readerThread) where-----------------------------------------------------------------------------------import Prelude hiding (take)-import Control.Monad-import Text.Printf-----------------------------------------------------------------------------------import Data.Serialize.Get-import Data.UUID-----------------------------------------------------------------------------------import Database.EventStore.Internal.Connection-import Database.EventStore.Internal.Types-import Database.EventStore.Logging-----------------------------------------------------------------------------------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 _ -> _settingsLog sett (Error WrongPackageFraming)- Right length_prefix -> connRecv c length_prefix >>= parsePackage- where- parsePackage bs =- case runGet getPackage bs of- Left e -> _settingsLog sett (Error $ PackageParsingError e)- Right pack -> push_p pack------------------------------------------------------------------------------------- Parsers----------------------------------------------------------------------------------getLengthPrefix :: Get Int-getLengthPrefix = fmap fromIntegral getWord32le-----------------------------------------------------------------------------------getPackage :: Get Package-getPackage = do- cmd <- getWord8- flg <- getFlag- col <- getUUID- cred <- getCredentials flg- rest <- remaining- dta <- getBytes rest-- let pack = Package- { packageCmd = cmd- , packageCorrelation = col- , packageData = dta- , packageCred = cred- }-- return pack-----------------------------------------------------------------------------------getFlag :: Get Flag-getFlag = do- wd <- getWord8- case wd of- 0x00 -> return None- 0x01 -> return Authenticated- _ -> fail $ printf "TCP: Unhandled flag value 0x%x" wd-----------------------------------------------------------------------------------getCredEntryLength :: Get Int-getCredEntryLength = fmap fromIntegral getWord8-----------------------------------------------------------------------------------getCredentials :: Flag -> Get (Maybe Credentials)-getCredentials None = return Nothing-getCredentials _ = do- loginLen <- getCredEntryLength- login <- getBytes loginLen- passwLen <- getCredEntryLength- passw <- getBytes passwLen- return $ Just $ credentials login passw-----------------------------------------------------------------------------------getUUID :: Get UUID-getUUID = do- bs <- getLazyByteString 16- case fromByteString bs of- Just uuid -> return uuid- _ -> fail "TCP: Wrong UUID format"
+ Database/EventStore/Internal/Stream.hs view
@@ -0,0 +1,31 @@+{-# LANGUAGE GADTs #-}+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE KindSignatures #-}+--------------------------------------------------------------------------------+-- |+-- Module : Database.EventStore.Internal.Stream+-- 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.Internal.Stream where++--------------------------------------------------------------------------------+import Data.Text++--------------------------------------------------------------------------------+-- | A stream can either point to $all or a regular one.+data StreamType = All | RegularStream deriving (Eq, Ord)++--------------------------------------------------------------------------------+-- | Represents a regular stream name or $all stream.+data StreamName = StreamName Text | AllStream deriving Eq++--------------------------------------------------------------------------------+instance Show StreamName where+ show (StreamName t) = show t+ show AllStream = "$all"
Database/EventStore/Internal/TimeSpan.hs view
@@ -9,7 +9,7 @@ -- Stability : provisional -- Portability : non-portable ----- Sorry but had no choice.+-- .NET TimeSpan implemented in Haskell. -------------------------------------------------------------------------------- module Database.EventStore.Internal.TimeSpan ( TimeSpan@@ -43,8 +43,7 @@ import Data.Text.Lazy.Builder ----------------------------------------------------------------------------------- | .NET TimeSpan: TimeSpan represents a duration of time. A TimeSpan can be--- negative or positive. Sorry+-- | .NET TimeSpan: Represents a time interval. newtype TimeSpan = TimeSpan Int64 deriving (Eq, Ord) --------------------------------------------------------------------------------@@ -64,6 +63,7 @@ parseJSON _ = empty --------------------------------------------------------------------------------+-- | Determines weither a number is positive or negative. parseFormatLiteral :: Parser FormatLiteral parseFormatLiteral = do c <- peekChar'@@ -165,21 +165,29 @@ / realToFrac ticksPerMillisecond) :: Double) --------------------------------------------------------------------------------+-- | Initializes a new instance of the TimeSpan structure to the specified+-- number of ticks. timeSpanTicks :: Int64 -> TimeSpan timeSpanTicks = TimeSpan --------------------------------------------------------------------------------+-- | Initializes a new instance of the TimeSpan structure to a specified number+-- of hours, minutes, and seconds. timeSpanHoursMinsSecs :: Int64 -> Int64 -> Int64 -> TimeSpan timeSpanHoursMinsSecs hh mm ss = TimeSpan $ totalSecs * ticksPerSecond where totalSecs = (hh * 3600) + (mm * 60) + ss --------------------------------------------------------------------------------+-- | Initializes a new instance of the TimeSpan structure to a specified number+-- of days, hours, minutes, and seconds. timeSpanDaysHoursMinsSecs :: Int64 -> Int64 -> Int64 -> Int64 -> TimeSpan timeSpanDaysHoursMinsSecs dd hh mm ss = timeSpanDaysHoursMinsSecsMillis dd hh mm ss 0 --------------------------------------------------------------------------------+-- | Initializes a new instance of the TimeSpan structure to a specified number+-- of days, hours, minutes, seconds, and milliseconds. timeSpanDaysHoursMinsSecsMillis :: Int64 -> Int64 -> Int64@@ -195,56 +203,78 @@ ss) * 1000 + ms --------------------------------------------------------------------------------+-- | Gets the number of ticks that represent the value of the current 'TimeSpan'+-- structure. timeSpanGetTicks :: TimeSpan -> Int64 timeSpanGetTicks (TimeSpan i) = i --------------------------------------------------------------------------------+-- | Gets the days component of the time interval represented by the current+-- 'TimeSpan' structure. timeSpanGetDays :: TimeSpan -> Int64 timeSpanGetDays (TimeSpan i) = truncate $ (realToFrac i :: Double) / (realToFrac ticksPerDay) --------------------------------------------------------------------------------+-- | Gets the hours component of the time interval represented by the current+-- 'TimeSpan' structure. timeSpanGetHours :: TimeSpan -> Int64 timeSpanGetHours (TimeSpan i) = mod (truncate $ (realToFrac i :: Double) / (realToFrac ticksPerHour)) 24 --------------------------------------------------------------------------------+-- | Gets the minutes component of the time interval represented by the current+-- 'TimeSpan' structure. timeSpanGetMinutes :: TimeSpan -> Int64 timeSpanGetMinutes (TimeSpan i) = mod (truncate $ (realToFrac i :: Double) / (realToFrac ticksPerMinute)) 60 --------------------------------------------------------------------------------+-- | Gets the seconds component of the time interval represented by the current+-- 'TimeSpan' structure. timeSpanGetSeconds :: TimeSpan -> Int64 timeSpanGetSeconds (TimeSpan i) = mod (truncate $ (realToFrac i :: Double) / (realToFrac ticksPerSecond)) 60 --------------------------------------------------------------------------------+-- | Gets the milliseconds component of the time interval represented by the+-- current 'TimeSpan' structure. timeSpanGetMillis :: TimeSpan -> Int64 timeSpanGetMillis (TimeSpan i) = mod (truncate $ (realToFrac i :: Double) / (realToFrac ticksPerMillisecond)) 1000 --------------------------------------------------------------------------------+-- | Returns a 'TimeSpan' that represents a specified number of seconds, where+-- the specification is accurate to the nearest millisecond. timeSpanFromSeconds :: Double -> TimeSpan timeSpanFromSeconds i = interval i millisPerSecond --------------------------------------------------------------------------------+-- | Returns a 'TimeSpan' that represents a specified number of minutes, where+-- the specification is accurate to the nearest millisecond. timeSpanFromMinutes :: Double -> TimeSpan timeSpanFromMinutes i = interval i millisPerMinute --------------------------------------------------------------------------------+-- | Returns a 'TimeSpan' that represents a specified number of hours, where the+-- specification is accurate to the nearest millisecond. timeSpanFromHours :: Double -> TimeSpan timeSpanFromHours i = interval i millisPerHour --------------------------------------------------------------------------------+-- | Returns a 'TimeSpan' that represents a specified number of days, where the+-- specification is accurate to the nearest millisecond. timeSpanFromDays :: Double -> TimeSpan timeSpanFromDays i = interval i millisPerDay --------------------------------------------------------------------------------+-- | Gets the value of the current 'TimeSpan' structure expressed in whole and+-- fractional milliseconds. timeSpanTotalMillis :: TimeSpan -> Int64 timeSpanTotalMillis (TimeSpan i) = let tmp = (realToFrac i) * millisPerTick in
Database/EventStore/Internal/Types.hs view
@@ -22,7 +22,7 @@ -------------------------------------------------------------------------------- import Control.Applicative import Control.Exception-import Control.Monad+import Control.Monad (mzero) import Data.ByteString (ByteString) import Data.ByteString.Lazy (fromStrict, toStrict) import Data.Int@@ -34,7 +34,6 @@ import GHC.Generics (Generic) ---------------------------------------------------------------------------------import Control.Concurrent.Async hiding (link) import qualified Data.Aeson as A import Data.Aeson.Types (Object, ToJSON(..), Pair, Parser, (.=)) import qualified Data.HashMap.Strict as H@@ -43,7 +42,6 @@ import Data.Time import Data.Time.Clock.POSIX import Data.UUID (UUID, fromByteString, toByteString)-import System.Random -------------------------------------------------------------------------------- import Database.EventStore.Logging@@ -52,33 +50,20 @@ -------------------------------------------------------------------------------- -- Exceptions --------------------------------------------------------------------------------+-- | Represent a class of error where the user is not at fault. It could be+-- either the client or the server. data InternalException = ConnectionClosedByServer- | Stopped+ -- ^ Happens when the server deliberately close the connection. This+ -- probably happens if the client didn't respect EventStore+ -- communication error. For instance, the client takes too much time to+ -- respond to a heartbeat request. deriving (Show, Typeable) -------------------------------------------------------------------------------- instance Exception InternalException ---------------------------------------------------------------------------------data OperationException- = WrongExpectedVersion Text ExpectedVersion -- ^ Stream and Expected Version- | StreamDeleted Text -- ^ Stream- | InvalidTransaction- | AccessDenied Text -- ^ Stream- | InvalidServerResponse Word8 Word8 -- ^ Expected, Found- | ProtobufDecodingError String- | ServerError (Maybe Text) -- ^ Reason- | InvalidOperation Text- deriving (Show, Typeable)-----------------------------------------------------------------------------------instance Exception OperationException-----------------------------------------------------------------------------------type OperationExceptional a = Either OperationException a---------------------------------------------------------------------------------- -- Event -------------------------------------------------------------------------------- -- | Contains event information like its type and data. Only used for write@@ -91,6 +76,7 @@ } deriving (Eq, Show) --------------------------------------------------------------------------------+-- | Create an 'Event' meant to be persisted. createEvent :: Text -- ^ Event type -> Maybe UUID -- ^ Event ID, generated if 'Nothing' -> EventData -- ^ Event data@@ -104,28 +90,33 @@ deriving (Eq, Show) --------------------------------------------------------------------------------+-- | Maps 'Event' inner data type to an 'Int32' understandable by the server. eventDataType :: EventData -> Int32 eventDataType (Json _ _) = 1 --------------------------------------------------------------------------------+-- | Maps 'Event' inner metadata type to an 'Int32' understandable by the server. eventMetadataType :: EventData -> Int32 eventMetadataType _ = 0 -------------------------------------------------------------------------------- -- | Creates a event using JSON format-withJson :: A.Value -> EventData-withJson value = Json value Nothing+withJson :: ToJSON a => a -> EventData+withJson value = Json (toJSON value) Nothing -------------------------------------------------------------------------------- -- | Create a event with metadata using JSON format-withJsonAndMetadata :: A.Value -> A.Value -> EventData-withJsonAndMetadata value metadata = Json value (Just metadata)+withJsonAndMetadata :: (ToJSON a, ToJSON b) => a -> b -> EventData+withJsonAndMetadata value metadata =+ Json (toJSON value) (Just $ toJSON metadata) --------------------------------------------------------------------------------+-- | Serializes 'EventData''s data to a raw 'ByteString'. eventDataBytes :: EventData -> ByteString eventDataBytes (Json value _) = toStrict $ A.encode value --------------------------------------------------------------------------------+-- | Serializes 'EventData' metadata to a raw 'ByteString'. eventMetadataBytes :: EventData -> Maybe ByteString eventMetadataBytes (Json _ meta_m) = fmap (toStrict . A.encode) meta_m @@ -149,6 +140,7 @@ deriving (Eq, Show) --------------------------------------------------------------------------------+-- | Maps a 'ExpectedVersion' to an 'Int32' understandable by the server. expVersionInt32 :: ExpectedVersion -> Int32 expVersionInt32 Any = -2 expVersionInt32 NoStream = -1@@ -157,44 +149,33 @@ -------------------------------------------------------------------------------- -- | This write should not conflict with anything and should always succeed.-anyStream :: ExpectedVersion-anyStream = Any+anyVersion :: ExpectedVersion+anyVersion = Any -------------------------------------------------------------------------------- -- | The stream being written to should not yet exist. If it does exist -- treat that as a concurrency problem.-noStream :: ExpectedVersion-noStream = NoStream+noStreamVersion :: ExpectedVersion+noStreamVersion = NoStream -------------------------------------------------------------------------------- -- | The stream should exist and should be empty. If it does not exist or -- is not empty, treat that as a concurrency problem.-emptyStream :: ExpectedVersion-emptyStream = EmptyStream+emptyStreamVersion :: ExpectedVersion+emptyStreamVersion = EmptyStream -------------------------------------------------------------------------------- -- | States that the last event written to the stream should have a -- sequence number matching your expected value.-exactStream :: Int32 -> ExpectedVersion-exactStream i+exactEventVersion :: Int32 -> ExpectedVersion+exactEventVersion i | i < 0 = error $ "expected version must be >= 0, but is " ++ show i | otherwise = Exact i -------------------------------------------------------------------------------- -- EventStore Messages ---------------------------------------------------------------------------------data OpResult- = OP_SUCCESS- | OP_PREPARE_TIMEOUT- | OP_COMMIT_TIMEOUT- | OP_FORWARD_TIMEOUT- | OP_WRONG_EXPECTED_VERSION- | OP_STREAM_DELETED- | OP_INVALID_TRANSACTION- | OP_ACCESS_DENIED- deriving (Eq, Enum, Show)----------------------------------------------------------------------------------+-- | Serializes form of an 'Event'. data NewEvent = NewEvent { newEventId :: Required 1 (Value ByteString)@@ -210,16 +191,16 @@ instance Encode NewEvent --------------------------------------------------------------------------------+-- | 'NewEvent' smart constructor. newEvent :: Text -- ^ Event type- -> Maybe UUID -- ^ Event ID+ -> UUID -- ^ Event ID -> Int32 -- ^ Data content type -> Int32 -- ^ Metadata content type -> ByteString -- ^ Event data -> Maybe ByteString -- ^ Metadata- -> IO NewEvent-newEvent evt_type evt_id data_type meta_type evt_data evt_meta = do- new_uuid <- maybe randomIO return evt_id- let uuid_bytes = toStrict $ toByteString new_uuid+ -> NewEvent+newEvent evt_type evt_id data_type meta_type evt_data evt_meta =+ let uuid_bytes = toStrict $ toByteString evt_id new_evt = NewEvent { newEventId = putField uuid_bytes , newEventType = putField evt_type@@ -227,115 +208,11 @@ , newEventMetadataType = putField meta_type , newEventData = putField evt_data , newEventMetadata = putField evt_meta- }-- return new_evt-----------------------------------------------------------------------------------data TransactionStart- = TransactionStart- { transactionStartStreamId :: Required 1 (Value Text)- , transactionStartExpectedVersion :: Required 2 (Value Int32)- , transactionStartRequireMaster :: Required 3 (Value Bool)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------newTransactionStart :: Text- -> Int32- -> Bool- -> TransactionStart-newTransactionStart stream_id exp_ver req_master =- TransactionStart- { transactionStartStreamId = putField stream_id- , transactionStartExpectedVersion = putField exp_ver- , transactionStartRequireMaster = putField req_master- }-----------------------------------------------------------------------------------instance Encode TransactionStart-----------------------------------------------------------------------------------data TransactionStartCompleted- = TransactionStartCompleted- { transactionSCId :: Required 1 (Value Int64)- , transactionSCResult :: Required 2 (Enumeration OpResult)- , transactionSCMessage :: Optional 3 (Value Text)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode TransactionStartCompleted-----------------------------------------------------------------------------------data TransactionWrite- = TransactionWrite- { transactionWriteId :: Required 1 (Value Int64)- , transactionWriteEvents :: Repeated 2 (Message NewEvent)- , transactionWriteRequireMaster :: Required 3 (Value Bool)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Encode TransactionWrite-----------------------------------------------------------------------------------newTransactionWrite :: Int64 -> [NewEvent] -> Bool -> TransactionWrite-newTransactionWrite trans_id evts req_master =- TransactionWrite- { transactionWriteId = putField trans_id- , transactionWriteEvents = putField evts- , transactionWriteRequireMaster = putField req_master- }-----------------------------------------------------------------------------------data TransactionWriteCompleted- = TransactionWriteCompleted- { transactionWCId :: Required 1 (Value Int64)- , transactionWCResult :: Required 2 (Enumeration OpResult)- , transactionWCMessage :: Optional 3 (Value Text)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode TransactionWriteCompleted-----------------------------------------------------------------------------------data TransactionCommit- = TransactionCommit- { transactionCommitId :: Required 1 (Value Int64)- , transactionCommitRequireMaster :: Required 2 (Value Bool)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Encode TransactionCommit-----------------------------------------------------------------------------------newTransactionCommit :: Int64 -> Bool -> TransactionCommit-newTransactionCommit trans_id req_master =- TransactionCommit- { transactionCommitId = putField trans_id- , transactionCommitRequireMaster = putField req_master- }-----------------------------------------------------------------------------------data TransactionCommitCompleted- = TransactionCommitCompleted- { transactionCCId :: Required 1 (Value Int64)- , transactionCCResult :: Required 2 (Enumeration OpResult)- , transactionCCMessage :: Optional 3 (Value Text)- , transactionCCFirstNumber :: Required 4 (Value Int32)- , transactionCCLastNumber :: Required 5 (Value Int32)- , transactionCCPreparePosition :: Optional 6 (Value Int64)- , transactionCCCommitPosition :: Optional 7 (Value Int64)- }- deriving (Generic, Show)-----------------------------------------------------------------------------------instance Decode TransactionCommitCompleted+ } in+ new_evt --------------------------------------------------------------------------------+-- | Represents a serialized event coming from the server. data EventRecord = EventRecord { eventRecordStreamId :: Required 1 (Value Text)@@ -355,6 +232,8 @@ instance Decode EventRecord --------------------------------------------------------------------------------+-- | Represents a serialized event representiong either an event or a link+-- event. data ResolvedIndexedEvent = ResolvedIndexedEvent { resolvedIndexedRecord :: Optional 1 (Message EventRecord)@@ -366,6 +245,7 @@ instance Decode ResolvedIndexedEvent --------------------------------------------------------------------------------+-- | Represents a serialized event sent by the server in a subscription context. data ResolvedEventBuf = ResolvedEventBuf { resolvedEventBufEvent :: Required 1 (Message EventRecord)@@ -388,9 +268,22 @@ { positionCommit :: !Int64 -- ^ Commit position of the record , positionPrepare :: !Int64 -- ^ Prepare position of the record }- deriving (Eq, Show)+ deriving Show --------------------------------------------------------------------------------+instance Eq Position where+ Position ac ap == Position bc bp = ac == bc && ap == bp++--------------------------------------------------------------------------------+instance Ord Position where+ compare (Position ac ap) (Position bc bp) =+ if ac < bc || (ac == bc && ap < bp)+ then LT+ else if ac > bc || (ac == bc && ap > bp)+ then GT+ else EQ++-------------------------------------------------------------------------------- -- | Representing the start of the transaction file. positionStart :: Position positionStart = Position 0 0@@ -401,23 +294,6 @@ positionEnd = Position (-1) (-1) ----------------------------------------------------------------------------------- | Returned after writing to a stream.-data WriteResult- = WriteResult- { writeNextExpectedVersion :: !Int32- -- ^ Next expected version of the stream.- , writePosition :: !Position- -- ^ 'Position' of the write.- }- deriving (Eq, Show)------------------------------------------------------------------------------------- | Returned after deleting a stream. 'Position' of the write.-newtype DeleteResult- = DeleteResult { deleteStreamPosition :: Position }- deriving (Eq, Show)---------------------------------------------------------------------------------- -- | Represents a previously written event. data RecordedEvent = RecordedEvent@@ -441,10 +317,17 @@ deriving Show --------------------------------------------------------------------------------+-- | Tries to parse JSON object from the given 'RecordedEvent'.+recordedEventDataAsJson :: A.FromJSON a => RecordedEvent -> Maybe a+recordedEventDataAsJson = A.decode . fromStrict . recordedEventData++--------------------------------------------------------------------------------+-- | Converts a raw 'Int64' into an 'UTCTime' toUTC :: Int64 -> UTCTime toUTC = posixSecondsToUTCTime . (/1000) . realToFrac . CTime --------------------------------------------------------------------------------+-- | Constructs a 'RecordedEvent' from an 'EventRecord'. newRecordedEvent :: EventRecord -> RecordedEvent newRecordedEvent er = re where@@ -474,29 +357,38 @@ -- link event. , resolvedEventLink :: !(Maybe RecordedEvent) -- ^ The link event if this 'ResolvedEvent' is a link event.+ , resolvedEventPosition :: !(Maybe Position)+ -- ^ Possible 'Position' of that event. } deriving Show --------------------------------------------------------------------------------+-- | Constructs a 'ResolvedEvent' from a 'ResolvedIndexedEvent'. newResolvedEvent :: ResolvedIndexedEvent -> ResolvedEvent newResolvedEvent rie = re where record = getField $ resolvedIndexedRecord rie link = getField $ resolvedIndexedLink rie re = ResolvedEvent- { resolvedEventRecord = fmap newRecordedEvent record- , resolvedEventLink = fmap newRecordedEvent link+ { resolvedEventRecord = fmap newRecordedEvent record+ , resolvedEventLink = fmap newRecordedEvent link+ , resolvedEventPosition = Nothing } --------------------------------------------------------------------------------+-- | Constructs a 'ResolvedEvent' from a 'ResolvedEventBuf'. newResolvedEventFromBuf :: ResolvedEventBuf -> ResolvedEvent newResolvedEventFromBuf reb = re where record = Just $ newRecordedEvent $ getField $ resolvedEventBufEvent reb link = getField $ resolvedEventBufLink reb+ com = getField $ resolvedEventBufCommitPosition reb+ pre = getField $ resolvedEventBufPreparePosition reb+ pos = Position com pre re = ResolvedEvent- { resolvedEventRecord = record- , resolvedEventLink = fmap newRecordedEvent link+ { resolvedEventRecord = record+ , resolvedEventLink = fmap newRecordedEvent link+ , resolvedEventPosition = Just pos } --------------------------------------------------------------------------------@@ -504,25 +396,29 @@ -- -- If this 'ResolvedEvent' represents a link event, the link will be the -- original event, otherwise it will be the event.-resolvedEventOriginal :: ResolvedEvent -> Maybe RecordedEvent-resolvedEventOriginal (ResolvedEvent record link) =- link <|> record+resolvedEventOriginal :: ResolvedEvent -> RecordedEvent+resolvedEventOriginal (ResolvedEvent record link _) =+ let Just evt = link <|> record in evt --------------------------------------------------------------------------------+-- | Tries to desarialize 'resolvedEventOriginal' data as JSON.+resolvedEventDataAsJson :: A.FromJSON a => ResolvedEvent -> Maybe a+resolvedEventDataAsJson = recordedEventDataAsJson . resolvedEventOriginal++-------------------------------------------------------------------------------- -- | Indicates whether this 'ResolvedEvent' is a resolved link event.-eventResolved :: ResolvedEvent -> Bool-eventResolved = isJust . resolvedEventOriginal+isEventResolvedLink :: ResolvedEvent -> Bool+isEventResolvedLink = isJust . resolvedEventLink -------------------------------------------------------------------------------- -- | The stream name of the original event.-resolvedEventOriginalStreamId :: ResolvedEvent -> Maybe Text-resolvedEventOriginalStreamId =- fmap recordedEventStreamId . resolvedEventOriginal+resolvedEventOriginalStreamId :: ResolvedEvent -> Text+resolvedEventOriginalStreamId = recordedEventStreamId . resolvedEventOriginal -------------------------------------------------------------------------------- -- | The ID of the original event.-resolvedEventOriginalId :: ResolvedEvent -> Maybe UUID-resolvedEventOriginalId = fmap recordedEventId . resolvedEventOriginal+resolvedEventOriginalId :: ResolvedEvent -> UUID+resolvedEventOriginalId = recordedEventId . resolvedEventOriginal -------------------------------------------------------------------------------- -- | Represents the direction of read operation (both from $all an usual@@ -533,35 +429,16 @@ deriving (Eq, Show) ----------------------------------------------------------------------------------- Transaction------------------------------------------------------------------------------------ | Represents a multi-request transaction with the EventStore.-data Transaction- = Transaction- { transactionId :: Int64- -- ^ The ID of the transaction. This can be used to recover a- -- transaction later.- , transactionStreamId :: Text- -- ^ The name of the stream.- , transactionExpectedVersion :: ExpectedVersion- -- ^ Expected version of the stream.- , transactionCommit :: IO (Async WriteResult)- -- ^ Asynchronously commits this transaction.- , transactionSendEvents :: [Event] -> IO (Async ())- -- ^ Asynchronously writes to a transaction in the EventStore.- , transactionRollback :: IO ()- -- ^ Rollback this transaction.- }---------------------------------------------------------------------------------- -- Flag --------------------------------------------------------------------------------+-- | Indicates either a 'Package' contains 'Credentials' data or not. data Flag = None | Authenticated deriving Show --------------------------------------------------------------------------------+-- | Maps a 'Flag' into a 'Word8' understandable by the server. flagWord8 :: Flag -> Word8 flagWord8 None = 0x00 flagWord8 Authenticated = 0x01@@ -578,6 +455,7 @@ deriving (Eq, Show) --------------------------------------------------------------------------------+-- | Creates a 'Credentials' given a login and a password. credentials :: ByteString -- ^ Login -> ByteString -- ^ Password -> Credentials@@ -586,6 +464,7 @@ -------------------------------------------------------------------------------- -- Package --------------------------------------------------------------------------------+-- | Represents a package exchanged between the client and the server. data Package = Package { packageCmd :: !Word8@@ -642,6 +521,7 @@ } --------------------------------------------------------------------------------+-- | Triggers the logger callback if it has been set. _settingsLog :: Settings -> Log -> IO () _settingsLog Settings{..} l = case s_logger of@@ -669,6 +549,14 @@ -- ^ Roles and users permitted to write stream metadata. } deriving Show +-------------------------------------------------------------------------------+instance A.FromJSON StreamACL where+ parseJSON = parseStreamACL++-------------------------------------------------------------------------------+instance A.ToJSON StreamACL where+ toJSON = streamACLJSON+ -------------------------------------------------------------------------------- -- | 'StreamACL' with no role or users whatsoever. emptyStreamACL :: StreamACL@@ -703,19 +591,16 @@ -------------------------------------------------------------------------------- -- | Gets a custom property value from metadata.-streamMetadataGetCustomPropertyValue :: StreamMetadata -> Text -> Maybe A.Value-streamMetadataGetCustomPropertyValue s k = H.lookup k obj+getCustomPropertyValue :: StreamMetadata -> Text -> Maybe A.Value+getCustomPropertyValue 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+getCustomProperty :: A.FromJSON a => StreamMetadata -> Text -> Maybe a+getCustomProperty s k = do+ v <- getCustomPropertyValue s k case A.fromJSON v of A.Error _ -> Nothing A.Success a -> return a@@ -724,6 +609,10 @@ instance A.FromJSON StreamMetadata where parseJSON = parseStreamMetadata +-------------------------------------------------------------------------------+instance A.ToJSON StreamMetadata where+ toJSON = streamMetadataJSON+ -------------------------------------------------------------------------------- -- | 'StreamMetadata' with everything set to 'Nothing', using 'emptyStreamACL' -- and an empty 'Object'.@@ -738,12 +627,14 @@ } --------------------------------------------------------------------------------+-- | Maps an 'Object' to a list of 'Pair' to ease the 'StreamMetadata'. customMetaToPairs :: Object -> [Pair] customMetaToPairs = fmap go . H.toList where go (k,v) = k .= v --------------------------------------------------------------------------------+-- | Serialized a 'StreamACL' to 'Value' for serialization purpose. streamACLJSON :: StreamACL -> A.Value streamACLJSON StreamACL{..} = A.object [ p_readRoles .= streamACLReadRoles@@ -754,6 +645,7 @@ ] --------------------------------------------------------------------------------+-- | Serialized a 'StreamMetadata' to 'Value' for serialization purpose. streamMetadataJSON :: StreamMetadata -> A.Value streamMetadataJSON StreamMetadata{..} = A.object $ [ p_maxAge .= streamMetadataMaxAge@@ -768,48 +660,60 @@ -------------------------------------------------------------------------------- -- Stream ACL Properties --------------------------------------------------------------------------------+-- | Read ACL property. p_readRoles :: Text p_readRoles = "$r" --------------------------------------------------------------------------------+-- | Write ACL property. p_writeRoles :: Text p_writeRoles = "$w" --------------------------------------------------------------------------------+-- | Delete ACL property. p_deleteRoles :: Text p_deleteRoles = "$d" --------------------------------------------------------------------------------+-- | Metadata read ACL property. p_metaReadRoles :: Text p_metaReadRoles = "$mr" --------------------------------------------------------------------------------+-- | Metadata write ACL property. p_metaWriteRoles :: Text p_metaWriteRoles = "$mw" -------------------------------------------------------------------------------- -- Internal MetaData Properties --------------------------------------------------------------------------------+-- | Max age metadata property. p_maxAge :: Text p_maxAge = "$maxAge" --------------------------------------------------------------------------------+-- | Max count metadata property. p_maxCount :: Text p_maxCount = "$maxCount" --------------------------------------------------------------------------------+-- | truncated before metadata property. p_truncateBefore :: Text p_truncateBefore = "$tb" --------------------------------------------------------------------------------+-- | Cache control metadata property. p_cacheControl :: Text p_cacheControl = "$cacheControl" --------------------------------------------------------------------------------+-- | ACL metadata property. p_acl :: Text p_acl = "$acl" --------------------------------------------------------------------------------+-- | Gathers every internal metadata properties into a 'Set'. It used to safely+-- 'StreamMetadata' in JSON. internalMetaProperties :: S.Set Text internalMetaProperties = S.fromList [ p_maxAge@@ -820,19 +724,21 @@ ] --------------------------------------------------------------------------------+-- | Only keeps the properties the users has set. keepUserProperties :: Object -> Object keepUserProperties = H.filterWithKey go where go k _ = not $ S.member k internalMetaProperties -----------------------------------------------------------------------------------------------------------------------------------------------------------------+-- | Parses a 'NominalDiffTime' from an 'Object' given a JSON property. parseNominalDiffTime :: Text -> Object -> Parser (Maybe NominalDiffTime) parseNominalDiffTime k m = fmap (fmap go) (m A..: k) where go i = (realToFrac $ CTime i) --------------------------------------------------------------------------------+-- | Parses 'StreamACL'. parseStreamACL :: A.Value -> Parser StreamACL parseStreamACL (A.Object m) = StreamACL <$>@@ -844,6 +750,7 @@ parseStreamACL _ = mzero --------------------------------------------------------------------------------+-- | Parses 'StreamMetadata'. parseStreamMetadata :: A.Value -> Parser StreamMetadata parseStreamMetadata (A.Object m) = StreamMetadata <$>@@ -858,13 +765,16 @@ -------------------------------------------------------------------------------- -- Builder --------------------------------------------------------------------------------+-- | Allows to build a structure using 'Monoid' functions. type Builder a = Endo a --------------------------------------------------------------------------------+-- | Build a structure given a 'Builder' and an initial value. build :: a -> Builder a -> a build a (Endo k) = k a --------------------------------------------------------------------------------+-- | A 'Builder' applies to 'StreamACL'. type StreamACLBuilder = Builder StreamACL --------------------------------------------------------------------------------@@ -928,6 +838,7 @@ modifyStreamACL b acl = build acl b --------------------------------------------------------------------------------+-- | A 'Builder' applies to 'StreamMetadata'. type StreamMetadataBuilder = Builder StreamMetadata --------------------------------------------------------------------------------@@ -951,10 +862,13 @@ setCacheControl d = Endo $ \s -> s { streamMetadataCacheControl = Just d } --------------------------------------------------------------------------------+-- | Overwrites any previous 'StreamACL' by the given one in a+-- 'StreamMetadataBuilder'. setACL :: StreamACL -> StreamMetadataBuilder setACL a = Endo $ \s -> s { streamMetadataACL = a } --------------------------------------------------------------------------------+-- | Updates a 'StreamMetadata''s 'StreamACL' given a 'StreamACLBuilder'. modifyACL :: StreamACLBuilder -> StreamMetadataBuilder modifyACL b = Endo $ \s -> let old = streamMetadataACL s@@ -1005,21 +919,24 @@ -- ^ 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.+ -- ^ Distributes events to each client in a round robin fashion. deriving (Show, Eq) --------------------------------------------------------------------------------+-- | Maps a 'SystemConsumerStrategy' to a 'Text' understandable by the server. strategyText :: SystemConsumerStrategy -> Text strategyText DispatchToSingle = "DispatchToSingle" strategyText RoundRobin = "RoundRobin" --------------------------------------------------------------------------------+-- | Tries to parse a 'SystemConsumerStrategy' given a raw 'Text'. strategyFromText :: Text -> Maybe SystemConsumerStrategy strategyFromText "DispatchToSingle" = Just DispatchToSingle strategyFromText "RoundRobin" = Just RoundRobin strategyFromText _ = Nothing --------------------------------------------------------------------------------+-- | Gathers every persistent subscription property. data PersistentSubscriptionSettings = PersistentSubscriptionSettings { psSettingsResolveLinkTos :: !Bool
− Database/EventStore/Internal/Util/Sodium.hs
@@ -1,35 +0,0 @@------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Util.Sodium--- 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.Util.Sodium where-----------------------------------------------------------------------------------import Control.Concurrent (forkIO)-import Data.Functor (void)-----------------------------------------------------------------------------------import FRP.Sodium-----------------------------------------------------------------------------------pushAsync :: (a -> Reactive ()) -> a -> IO ()-pushAsync push a = void $ forkIO $ sync $ push a-----------------------------------------------------------------------------------pushAsync2 :: (a -> b -> Reactive ()) -> a -> b -> IO ()-pushAsync2 push a b = void $ forkIO $ sync $ push a b-----------------------------------------------------------------------------------pushAsync3 :: (a -> b -> c -> Reactive ()) -> a -> b -> c -> IO ()-pushAsync3 push a b c = void $ forkIO $ sync $ push a b c-----------------------------------------------------------------------------------pushAsync4 :: (a -> b -> c -> d -> Reactive ()) -> a -> b -> c -> d -> IO ()-pushAsync4 push a b c d = void $ forkIO $ sync $ push a b c d
− Database/EventStore/Internal/Writer.hs
@@ -1,31 +0,0 @@------------------------------------------------------------------------------------ |--- Module : Database.EventStore.Internal.Writer--- 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.Writer (writerThread) where-----------------------------------------------------------------------------------import Control.Concurrent-import Control.Monad-----------------------------------------------------------------------------------import Data.Serialize.Put-----------------------------------------------------------------------------------import Database.EventStore.Internal.Connection-import Database.EventStore.Internal.Packages-import Database.EventStore.Internal.Types-----------------------------------------------------------------------------------writerThread :: Chan Package -> Connection -> IO ()-writerThread chan c = forever $ do- pkg <- readChan chan- connSend c (runPut $ putPackage pkg)- connFlush c
Database/EventStore/Logging.hs view
@@ -12,6 +12,10 @@ module Database.EventStore.Logging where --------------------------------------------------------------------------------+import Control.Exception+import Data.Word++-------------------------------------------------------------------------------- import Data.UUID --------------------------------------------------------------------------------@@ -23,13 +27,7 @@ -------------------------------------------------------------------------------- -- | Classifies error-like log messages.-data ErrorMessage- = MaxAttemptConnectionReached Int- -- ^ Indicates max attempt value.- | WrongPackageFraming- | PackageParsingError String- -- ^ Indicates parsing error message.- deriving Show+data ErrorMessage = UnexpectedException SomeException deriving Show -------------------------------------------------------------------------------- -- | Classifies info-like log messages.@@ -38,8 +36,12 @@ -- ^ Indicates current attempt. | ConnectionClosed UUID -- ^ Indicates connection 'UUID'.- | Connected UUID + | Connected UUID -- ^ Indicates connection 'UUID'. | Disconnected UUID -- ^ Indicates connection 'UUID'+ | PackageSent Word8 UUID+ -- ^ Indicates a package has been sent.+ | PackageReceived Word8 UUID+ -- ^ Indicates the client's received a package from the server. deriving Show
README.md view
@@ -6,27 +6,104 @@ 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).+ * Read event(s) from regular or $all stream (forward or backward).+ * Write event(s) to regular stream.+ * Delete regular stream.+ * Transactional writes to regular stream.+ * Volatile subscriptions to regular or $all stream.+ * Catch-up subscriptions to regular or $all stream.+ * Competing consumers (a.k.a Persistent subscriptions) to regular stream.+ * Authenticated communication with EventStore server.+ * Read stream metadata (ACL and custom properties).+ * Write stream metadata (ACL and custom properties). -TODO-====- 1. SSL+Not implemented yet+===================+ * Secured connection with the server (SSL). Requirements ============- 1. GHC >= 7.8.3- 2. Cabal >= 1.18- 3. EventStore >= 3.0.0 (>= 3.1.0 if you want competing consumers)+ * GHC >= 7.8.3+ * Cabal >= 1.18+ * EventStore >= 3.0.0 (>= 3.1.0 if you want competing consumers) -Tested on Linux and OSX Yosemite.+Install+======= -BSD3 License+* Using [Hackage](https://hackage.haskell.org/package/eventstore)+```+$ cabal update+$ cabal install eventstore+```++* From source+```+$ git clone https://github.com/YoEight/eventstore.git+$ cd eventstore+$ cabal install --only-dependencies+$ cabal configure +$ cabal install+```++How to test+===========+Tests are available. Those assume a server is running on `127.0.0.1` and `1113` port.+```+$ cabal install --only-dependencies --enable-tests+$ cabal configure --enable-tests+$ cabal test+```++How to use+==========++```haskell+{-# LANGUAGE OverloadedStrings #-} -- That library uses `Text` pervasively. This pragma permits to use+ -- String literal when a Text is needed.+module Main where ++import Data.Aeson+-- It requires to have `aeson` package installed. Note that EventStore doesn't constraint you to JSON+-- format but putting common use aside, by doing so you'll be able to use some interesting EventStore+-- features like its Complex Event Processing (CEP) capabality.+ +import Database.EventStore+-- Note that import also re-exports 'Control.Concurrent.Async' module, allowing the use of 'wait'+-- function for instance.++main :: IO ()+main = do+ -- A common pattern with an EventStore connection is to create a single instance only and pass it + -- wherever you need it (it's threadsafe). It's very important to not consider an EventStore connection like + -- its regular SQL counterpart. An EventStore connection will try its best to reconnect+ -- automatically to the server if the connection dropped. Of course that behavior can be tuned+ -- through some settings.+ conn <- connect defaultSettings "127.0.0.1" 1113+ let js = "isHaskellTheBest" .= True -- (.=) comes from Data.Aeson module.+ evt = createEvent "programming" Nothing (withJson js)+ + -- Appends an event to a stream named `languages`. + as <- sendEvent conn "languages" anyVersion evt+ + -- EventStore interactions are fundamentally asynchronous. Nothing requires you to wait + -- for the completion of an operation, but it's good to know if something went wrong.+ _ <- wait as+ + -- Again, if you decide to `shutdown` an EventStore connection, it means your application is + -- about to terminate.+ shutdown conn+ + -- Make sure the EventStore connection completes every ongoing operation. For instance, if + -- at the moment we call `shutdown` and some operations (or subscriptions) were still pending,+ -- the connection aborted all of them.+ waitTillClosed conn+```+Notes+=====+That library was tested on Linux and OSX Yosemite.++Contributions and bug reports are welcome!++BSD3 License ++-Yorick Laupa
eventstore.cabal view
@@ -10,7 +10,7 @@ -- PVP summary: +-+------- breaking API changes -- | | +----- non-breaking API additions -- | | | +--- code changes with no API change-version: 0.8.0.0+version: 0.9.0.0 -- A short (one-line) description of the package. synopsis: EventStore TCP Client@@ -57,23 +57,38 @@ exposed-modules: Database.EventStore Database.EventStore.Logging -- Modules included in this library but not exported.- other-modules: Database.EventStore.Catchup- Database.EventStore.Internal.Connection+ other-modules: Database.EventStore.Internal.Connection+ Database.EventStore.Internal.Execution.Production+ Database.EventStore.Internal.Generator+ Database.EventStore.Internal.Operation Database.EventStore.Internal.Packages Database.EventStore.Internal.Processor- Database.EventStore.Internal.Reader+ Database.EventStore.Internal.Stream Database.EventStore.Internal.TimeSpan Database.EventStore.Internal.Types- Database.EventStore.Internal.Manager.Operation+ Database.EventStore.Internal.Manager.Operation.Model Database.EventStore.Internal.Manager.Subscription- Database.EventStore.Internal.Operation.DeleteStreamOperation- Database.EventStore.Internal.Operation.ReadAllEventsOperation- Database.EventStore.Internal.Operation.ReadEventOperation- Database.EventStore.Internal.Operation.ReadStreamEventsOperation- Database.EventStore.Internal.Operation.TransactionStartOperation- Database.EventStore.Internal.Operation.WriteEventsOperation- Database.EventStore.Internal.Util.Sodium- Database.EventStore.Internal.Writer+ Database.EventStore.Internal.Manager.Subscription.Driver+ Database.EventStore.Internal.Manager.Subscription.Message+ Database.EventStore.Internal.Manager.Subscription.Model+ Database.EventStore.Internal.Manager.Subscription.Packages+ Database.EventStore.Internal.Operations+ Database.EventStore.Internal.Operation.Catchup+ Database.EventStore.Internal.Operation.DeleteStream+ Database.EventStore.Internal.Operation.DeleteStream.Message+ Database.EventStore.Internal.Operation.ReadAllEvents+ Database.EventStore.Internal.Operation.ReadAllEvents.Message+ Database.EventStore.Internal.Operation.ReadEvent+ Database.EventStore.Internal.Operation.ReadEvent.Message+ Database.EventStore.Internal.Operation.ReadStreamEvents+ Database.EventStore.Internal.Operation.ReadStreamEvents.Message+ Database.EventStore.Internal.Operation.Read.Common+ Database.EventStore.Internal.Operation.StreamMetadata+ Database.EventStore.Internal.Operation.Transaction+ Database.EventStore.Internal.Operation.Transaction.Message+ Database.EventStore.Internal.Operation.WriteEvents+ Database.EventStore.Internal.Operation.WriteEvents.Message+ Database.EventStore.Internal.Operation.Write.Common -- LANGUAGE extensions used by modules in this package. -- other-extensions:@@ -88,7 +103,6 @@ , network ==2.6.* , protobuf >=0.2 && <0.3 , random ==1.*- , sodium ==0.11.* , text >=1.1.1 && <1.3 , time >=1.4 && <1.6 , uuid ==1.3.*@@ -103,3 +117,20 @@ ghc-options: -Wall default-language: Haskell2010++test-suite integration-tests+ type: exitcode-stdio-1.0+ default-language: Haskell2010+ hs-source-dirs: tests+ main-is: integration.hs++ ghc-options: -Wall++ build-depends: base,+ eventstore,+ tasty,+ tasty-hunit,+ aeson,+ text,+ stm,+ time
+ tests/integration.hs view
@@ -0,0 +1,59 @@+{-# LANGUAGE OverloadedStrings #-}+--------------------------------------------------------------------------------+-- |+-- Module : Main+-- Copyright : (C) 2015 Yorick Laupa+-- License : (see the file LICENSE)+--+-- Maintainer : Yorick Laupa <yo.eight@gmail.com>+-- Stability : provisional+-- Portability : non-portable+--+-- Main integration entry point.+--------------------------------------------------------------------------------+module Main where++--------------------------------------------------------------------------------+import Data.Time+import Database.EventStore+import Database.EventStore.Logging+import Numeric+import Test.Tasty+import Test.Tasty.Ingredients.Basic++--------------------------------------------------------------------------------+import Tests++--------------------------------------------------------------------------------+main :: IO ()+main = do+ let setts = defaultSettings+ { s_credentials = Just $ credentials "admin" "changeit"+ , s_reconnect_delay_secs = 1+ -- , s_logger = Just logger+ }+ conn <- connect setts "127.0.0.1" 1113+ let tree = tests conn+ defaultMainWithIngredients [consoleTestReporter] tree++--------------------------------------------------------------------------------+logger :: Log -> IO ()+logger l = do+ t <- getCurrentTime+ putStr "["+ putStr $ show t+ putStr "] "+ showLog l+ where+ showLog (Info m) = do+ putStr "[INFO] "+ case m of+ PackageSent cmd uuid ->+ putStrLn $ "Sent 0x" ++ showHex cmd "" ++ " over " ++ show uuid+ PackageReceived cmd uuid ->+ putStrLn $ "Received 0x" ++ showHex cmd "" ++ " over "+ ++ show uuid+ _ -> print m+ showLog (Error m) = do+ putStr "!!! ERROR !!! "+ print m