eventstore 0.2.0.1 → 0.3.0.0
raw patch · 4 files changed
+48/−95 lines, 4 filesPVP ok
version bump matches the API change (PVP)
API changes (from Hackage documentation)
+ Database.EventStore: subChan :: Subscription -> Chan (Either DropReason ResolvedEvent)
- Database.EventStore: Subscription :: !UUID -> !Text -> !Bool -> !Int64 -> !(Maybe Int32) -> IO () -> Subscription
+ Database.EventStore: Subscription :: !UUID -> !Text -> !Bool -> !Int64 -> !(Maybe Int32) -> Chan (Either DropReason ResolvedEvent) -> IO () -> Subscription
- Database.EventStore: subscribe :: Connection -> Text -> Bool -> (Subscription -> Either DropReason ResolvedEvent -> IO ()) -> IO (Async Subscription)
+ Database.EventStore: subscribe :: Connection -> Text -> Bool -> IO (Async Subscription)
Files
- Database/EventStore.hs +1/−3
- Database/EventStore/Internal/Manager/Subscription.hs +43/−84
- Database/EventStore/Internal/Processor.hs +3/−7
- eventstore.cabal +1/−1
Database/EventStore.hs view
@@ -277,13 +277,11 @@ subscribe :: Connection -> Text -> Bool- -> (Subscription -> Either DropReason ResolvedEvent -> IO ()) -> IO (Async Subscription)-subscribe Connection{..} stream_id res_lnk_tos cb = do+subscribe Connection{..} stream_id res_lnk_tos = do tmp <- newEmptyMVar processorNewSubcription conProcessor (putMVar tmp)- cb stream_id res_lnk_tos async $ readMVar tmp
Database/EventStore/Internal/Manager/Subscription.hs view
@@ -16,13 +16,13 @@ -------------------------------------------------------------------------------- module Database.EventStore.Internal.Manager.Subscription ( DropReason (..)+ , NewSubscriptionCB , Subscription(..) , subscriptionNetwork ) where -------------------------------------------------------------------------------- import Control.Concurrent-import Control.Concurrent.Async import Control.Monad import Control.Monad.Fix import Data.ByteString (ByteString)@@ -112,9 +112,8 @@ -------------------------------------------------------------------------------- data Pending = Pending- { penConf :: Int64 -> Maybe Int32 -> Subscription- , penCb :: Subscription -> IO ()- , penEvtCb :: Subscription -> Either DropReason ResolvedEvent -> IO ()+ { penConf :: Int64 -> Maybe Int32 -> IO Subscription+ , penCb :: Subscription -> IO () } --------------------------------------------------------------------------------@@ -132,23 +131,15 @@ , subResolveLinkTos :: !Bool , subLastCommitPos :: !Int64 , subLastEventNumber :: !(Maybe Int32)+ , subChan :: Chan (Either DropReason ResolvedEvent) , subUnsubscribe :: IO () } ---------------------------------------------------------------------------------data Sub- = Sub- { _subSub :: Subscription- , _subCb :: Subscription -> Either DropReason ResolvedEvent -> IO ()- , _subChan :: Chan (IO ())- , _subAS :: Async ()- }---------------------------------------------------------------------------------- data Manager = Manager { _pendings :: !(M.Map UUID Pending)- , _subscriptions :: !(M.Map UUID Sub)+ , _subscriptions :: !(M.Map UUID Subscription) } --------------------------------------------------------------------------------@@ -190,7 +181,6 @@ = Appeared { _appSub :: !Subscription , _appEvt :: !ResolvedEvent- , _appCb :: Subscription -> ResolvedEvent -> IO () } --------------------------------------------------------------------------------@@ -200,36 +190,23 @@ sub <- M.lookup packageCorrelation _subscriptions sea <- maybeDecodeMessage packageData let res_evt = getField $ streamResolvedEvent sea- chan = _subChan sub- evt_cb = \s evt -> writeChan chan ((_subCb sub) s $ Right evt)- app = Appeared- { _appSub = _subSub sub- , _appEvt = newResolvedEventFromBuf res_evt- , _appCb = evt_cb- } - return app+ return Appeared+ { _appSub = sub+ , _appEvt = newResolvedEventFromBuf res_evt+ } | otherwise = Nothing ---------------------------------------------------------------------------------confirmSub :: Confirmation -> Manager -> IO (Maybe Sub)+confirmSub :: Confirmation -> Manager -> IO (Maybe Subscription) confirmSub (Confirmation uuid sc) Manager{..} = for (M.lookup uuid _pendings) $ \p -> do let last_com_pos = getField $ subscribeLastCommitPos sc last_evt_num = getField $ subscribeLastEventNumber sc - let sub = penConf p last_com_pos last_evt_num+ sub <- penConf p last_com_pos last_evt_num penCb p sub- chan <- newChan- as <- async $ forever $ do- action <- readChan chan- action- return Sub- { _subSub = sub- , _subCb = penEvtCb p- , _subChan = chan- , _subAS = as- }+ return sub -------------------------------------------------------------------------------- -- Events@@ -240,25 +217,19 @@ , _subCallback :: Subscription -> IO () , _subStream :: !Text , _subResolveLinkTos :: !Bool- , _subEventAppeared :: Subscription- -> Either DropReason ResolvedEvent- -> IO () } -------------------------------------------------------------------------------- data Confirmation = Confirmation !UUID !SubscriptionConfirmation ---------------------------------------------------------------------------------type NewSubscription =- (Subscription -> IO ())- -> (Subscription -> Either DropReason ResolvedEvent -> IO ())- -> Text -> Bool -> IO ()+type NewSubscriptionCB = (Subscription -> IO ()) -> Text -> Bool -> IO () -------------------------------------------------------------------------------- subscriptionNetwork :: Settings -> (Package -> Reactive ()) -> Event Package- -> Reactive NewSubscription+ -> Reactive NewSubscriptionCB subscriptionNetwork sett push_pkg e_pkg = do (on_sub, push_sub) <- newEvent (on_rem, push_rem) <- newEvent@@ -282,20 +253,21 @@ mgr_b push_pkg_io = pushAsync push_pkg - push_sub_io cb evt_cb stream res_lnk_tos = randomIO >>= \uuid -> sync $+ push_sub_io cb stream res_lnk_tos = randomIO >>= \uuid -> sync $ push_sub Subscribe { _subId = uuid , _subCallback = cb , _subStream = stream , _subResolveLinkTos = res_lnk_tos- , _subEventAppeared = evt_cb } _ <- listen on_sub (push_pkg_io . createSubscribePackage sett) - _ <- listen on_app $ \(Appeared sub evt cb) -> cb sub evt+ _ <- listen on_app $ \(Appeared sub evt) ->+ writeChan (subChan sub) (Right evt) - _ <- listen on_drop $ \(Dropped reason sub _ cb) -> cb sub reason+ _ <- listen on_drop $ \(Dropped reason sub _) ->+ writeChan (subChan sub) (Left reason) return push_sub_io @@ -327,7 +299,6 @@ { droppedReason :: !DropReason , droppedSub :: !Subscription , droppedId :: !UUID- , droppedCb :: Subscription -> DropReason -> IO () } --------------------------------------------------------------------------------@@ -343,25 +314,13 @@ go | packageCmd == subscriptionDropped = do sub <- M.lookup packageCorrelation _subscriptions msg <- maybeDecodeMessage packageData- let reason = fromMaybe D_Unsubscribed $ getField $ dropReason msg- chan = _subChan sub-- drop_cb = \s r -> do- var <- newEmptyMVar- writeChan chan $ do- (_subCb sub) s (Left r)- putMVar var ()- readMVar var- cancel $ _subAS sub-- dropped = Dropped- { droppedReason = reason- , droppedSub = _subSub sub- , droppedId = packageCorrelation- , droppedCb = drop_cb- }+ let reason = fromMaybe D_Unsubscribed $ getField $ dropReason msg - return dropped+ return Dropped+ { droppedReason = reason+ , droppedSub = sub+ , droppedId = packageCorrelation+ } | otherwise = Nothing --------------------------------------------------------------------------------@@ -373,29 +332,29 @@ where pending = Pending- { penConf = new_sub- , penCb = _subCallback- , penEvtCb = _subEventAppeared+ { penConf = new_sub+ , penCb = _subCallback } - new_sub com_pos last_evt =- Subscription- { subId = _subId- , subStream = _subStream- , subResolveLinkTos = _subResolveLinkTos- , subLastCommitPos = com_pos- , subLastEventNumber = last_evt- , subUnsubscribe = sync $ unsub _subId- }+ new_sub com_pos last_evt = do+ chan <- newChan + return Subscription+ { subId = _subId+ , subStream = _subStream+ , subResolveLinkTos = _subResolveLinkTos+ , subLastCommitPos = com_pos+ , subLastEventNumber = last_evt+ , subChan = chan+ , subUnsubscribe = sync $ unsub _subId+ }+ ---------------------------------------------------------------------------------confirmed :: Sub -> Manager -> Manager-confirmed ss s@Manager{..} =- s { _pendings = M.delete sub_id _pendings- , _subscriptions = M.insert sub_id ss _subscriptions+confirmed :: Subscription -> Manager -> Manager+confirmed sub@Subscription{..} s@Manager{..} =+ s { _pendings = M.delete subId _pendings+ , _subscriptions = M.insert subId sub _subscriptions }- where- sub_id = subId $ _subSub ss -------------------------------------------------------------------------------- remove :: UUID -> Manager -> Manager
Database/EventStore/Internal/Processor.hs view
@@ -17,6 +17,7 @@ , Processor(..) , DropReason(..) , Subscription(..)+ , NewSubscriptionCB , newProcessor ) where @@ -54,12 +55,7 @@ { processorConnect :: HostName -> Int -> IO () , processorShutdown :: IO () , processorNewOperation :: OperationParams -> IO ()- , processorNewSubcription ::- (Subscription -> IO ())- -> (Subscription -> Either DropReason ResolvedEvent -> IO ())- -> Text- -> Bool- -> IO ()+ , processorNewSubcription :: NewSubscriptionCB } --------------------------------------------------------------------------------@@ -177,7 +173,7 @@ let processor = Processor { processorConnect = \h p -> sync $ pushConnect $ Connect h p- , processorShutdown = pushAsync pushCleanup Cleanup+ , processorShutdown = sync $ pushCleanup Cleanup , processorNewOperation = \o -> sync $ push_new_op o , processorNewSubcription = push_sub }
eventstore.cabal view
@@ -10,7 +10,7 @@ -- PVP summary: +-+------- breaking API changes -- | | +----- non-breaking API additions -- | | | +--- code changes with no API change-version: 0.2.0.1+version: 0.3.0.0 -- A short (one-line) description of the package. synopsis: EventStore Haskell TCP Client