superevent (empty) → 0.1.0.0
raw patch · 10 files changed
+751/−0 lines, 10 filesdep +QuickCheckdep +aesondep +asyncsetup-changed
Dependencies added: QuickCheck, aeson, async, base, bytestring, conduit, containers, contravariant, hashable, hasql, hasql-migration, hasql-pool, hasql-transaction, hspec, mtl, random, stm, string-qq, superevent, temporary, text, time, timespan, transformers, unordered-containers, uuid, vector
Files
- ChangeLog.md +5/−0
- LICENSE +30/−0
- README.md +3/−0
- Setup.hs +2/−0
- src/SuperEvent/Store/Hasql.hs +364/−0
- src/SuperEvent/Store/Hasql/Utils.hs +20/−0
- src/SuperEvent/Store/Types.hs +108/−0
- superevent.cabal +90/−0
- test/Spec.hs +1/−0
- test/SuperEvent/Store/HasqlSpec.hs +128/−0
+ ChangeLog.md view
@@ -0,0 +1,5 @@+# Changelog for superevent++## 0.1.0.0++* First release.
+ LICENSE view
@@ -0,0 +1,30 @@+Copyright Alexander Thiemann (c) 2018++All rights reserved.++Redistribution and use in source and binary forms, with or without+modification, are permitted provided that the following conditions are met:++ * Redistributions of source code must retain the above copyright+ notice, this list of conditions and the following disclaimer.++ * Redistributions in binary form must reproduce the above+ copyright notice, this list of conditions and the following+ disclaimer in the documentation and/or other materials provided+ with the distribution.++ * Neither the name of Alexander Thiemann nor the names of other+ contributors may be used to endorse or promote products derived+ from this software without specific prior written permission.++THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS+"AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT+LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR+A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT+OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,+SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT+LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,+DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY+THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT+(INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE+OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
+ README.md view
@@ -0,0 +1,3 @@+# superevent++A simple opinionated event store implementation. In the future it might become backend independent, but currently it ships with a `hasql` backend for PostgreSQL.
+ Setup.hs view
@@ -0,0 +1,2 @@+import Distribution.Simple+main = defaultMain
+ src/SuperEvent/Store/Hasql.hs view
@@ -0,0 +1,364 @@+{-# LANGUAGE MultiParamTypeClasses #-}+{-# LANGUAGE QuasiQuotes #-}+{-# LANGUAGE StrictData #-}+{-# LANGUAGE BangPatterns #-}+{-# LANGUAGE OverloadedStrings #-}+module SuperEvent.Store.Hasql+ ( newPgSqlStore+ , withTempStore+ , DbStore+ )+where++import SuperEvent.Store.Hasql.Utils+import SuperEvent.Store.Types++import Control.Exception+import Control.Monad.Trans+import Data.Conduit+import Data.Functor.Contravariant+import Data.Maybe+import Data.Monoid+import Data.String.QQ+import Data.Time.TimeSpan+import Hasql.Query+import System.Random+import qualified Data.ByteString as BS+import qualified Data.ByteString.Char8 as BSC+import qualified Data.Vector as V+import qualified Hasql.Connection as C+import qualified Hasql.Decoders as D+import qualified Hasql.Encoders as E+import qualified Hasql.Migration as M+import qualified Hasql.Pool as P+import qualified Hasql.Session as S+import qualified Hasql.Transaction as Tx++data DbStore+ = DbStore+ { db_store :: Store+ }++newPgSqlStore :: BS.ByteString -> IO DbStore+newPgSqlStore connStr =+ do pool <- P.acquire (20, 60 * 5, connStr)+ let store = Store pool+ let migs = [M.MigrationScript "schemav1" schemaV1]+ withPool store $+ let loop [] = pure ()+ loop (mig : more) =+ do res <-+ dbTx Tx.ReadCommitted Tx.Write $ M.runMigration mig+ case res of+ M.MigrationError err -> fail err+ M.MigrationSuccess -> loop more+ in loop (M.MigrationInitialization : migs)+ pure $ DbStore store++-- | Temporary postgres store for tests+withTempStore :: (DbStore -> IO a) -> IO a+withTempStore run =+ bracket allocDb removeDb $ \(_, dbname) ->+ do putStrLn ("TempDB is: " <> show dbname)+ bracket (newPgSqlStore $ "dbname=" <> dbname) (\_ -> pure ()) run+ where+ assertRight y =+ case y of+ Right x -> pure x+ Left errMsg -> fail (show errMsg)+ removeDb (globalConn, dbname) =+ do runRes2 <-+ flip S.run globalConn $ S.sql $ "DROP DATABASE IF EXISTS " <> dbname+ assertRight runRes2+ C.release globalConn+ allocDb =+ do globalConnE <- C.acquire ""+ globalConn <- assertRight globalConnE+ dbnameSuffix <-+ BSC.pack . take 10 . randomRs ('a', 'z') <$>+ newStdGen+ let dbname = "eventstore_temp_" <> dbnameSuffix+ runRes <-+ flip S.run globalConn $+ do S.sql $ "DROP DATABASE IF EXISTS " <> dbname+ S.sql $ "CREATE DATABASE " <> dbname+ assertRight runRes+ localConnE <- C.acquire $ "dbname=" <> dbname+ localConn <- assertRight localConnE+ runRes' <-+ flip S.run localConn $+ S.sql $ "CREATE EXTENSION hstore"+ assertRight runRes'+ C.release localConn+ pure (globalConn, dbname)++encStreamId :: E.Value StreamId+encStreamId = contramap unStreamId E.text++encEventType :: E.Value EventType+encEventType = contramap unEventType E.text++encEventNumber :: E.Value EventNumber+encEventNumber = contramap unEventNumber E.int8++encGlobalPosition :: E.Value GlobalPosition+encGlobalPosition = contramap unGlobalPosition E.int8++decEventNumber :: D.Row EventNumber+decEventNumber = EventNumber <$> D.value D.int8++decRecordedEvent :: D.Row RecordedEvent+decRecordedEvent =+ RecordedEvent+ <$> (StreamId <$> D.value D.text)+ <*> D.value D.uuid+ <*> decEventNumber+ <*> (EventType <$> D.value D.text)+ <*> D.value D.jsonb+ <*> D.value D.jsonb+ <*> D.value D.timestamptz++qStreamVersion :: Query StreamId (Maybe EventNumber)+qStreamVersion =+ statement sql encoder decoder True+ where+ sql =+ "SELECT MAX(version) FROM events WHERE stream = $1"+ encoder =+ E.value encStreamId+ decoder =+ fmap EventNumber <$> D.singleRow (D.nullableValue D.int8)++data WriteEvent+ = WriteEvent+ { we_stream :: StreamId+ , we_number :: EventNumber+ , we_data :: EventData+ }++qWriteEvent :: Query WriteEvent ()+qWriteEvent =+ statement sql encoder D.unit True+ where+ sql =+ "INSERT INTO events "+ <> "(id, stream, version, type, data, meta_data)"+ <> " VALUES "+ <> "($1, $2, $3, $4, $5, $6)"+ encoder =+ contramap (ed_guid . we_data) (E.value E.uuid)+ <> contramap we_stream (E.value encStreamId)+ <> contramap we_number (E.value encEventNumber)+ <> contramap (ed_type . we_data) (E.value encEventType)+ <> contramap (ed_data . we_data) (E.value E.jsonb)+ <> contramap (ed_metadata . we_data) (E.value E.jsonb)++data SingleEventQuery+ = SingleEventQuery+ { seq_stream :: StreamId+ , seq_number :: EventNumber+ }++qSingleEvent :: Query SingleEventQuery (Maybe RecordedEvent)+qSingleEvent =+ statement sql encoder decoder True+ where+ sql =+ "SELECT "+ <> "stream, id, version, type, data, meta_data, created "+ <> "FROM "+ <> "events "+ <> "WHERE stream = $1 AND version = $2 LIMIT 1"+ encoder =+ contramap seq_stream (E.value encStreamId)+ <> contramap seq_number (E.value encEventNumber)+ decoder =+ D.maybeRow decRecordedEvent++data EventStreamQuery+ = EventStreamQuery+ { esq_stream :: StreamId+ , esq_number :: EventNumber+ , esq_limit :: Int+ }++sqlEventStream :: ReadDirection -> BS.ByteString+sqlEventStream readDir =+ "SELECT "+ <> "stream, id, version, type, data, meta_data, created "+ <> "FROM "+ <> "events "+ <> "WHERE stream = $1 AND version "+ <> (if readDir == RdForward then ">=" else "<=")+ <> " $2 "+ <> "ORDER BY version "+ <> (if readDir == RdForward then "ASC" else "DESC")+ <> " LIMIT $3"++encEventStreamQuery :: E.Params EventStreamQuery+encEventStreamQuery =+ contramap esq_stream (E.value encStreamId)+ <> contramap esq_number (E.value encEventNumber)+ <> contramap (fromIntegral . esq_limit) (E.value E.int8)++qEventStream :: ReadDirection -> Query EventStreamQuery (V.Vector RecordedEvent)+qEventStream readDir =+ statement (sqlEventStream readDir) encEventStreamQuery decoder True+ where+ decoder =+ D.rowsVector decRecordedEvent++data GlobalEventQuery+ = GlobalEventQuery+ { geq_position :: GlobalPosition+ , geq_limit :: Int+ }++qGlobalEvent ::+ ReadDirection+ -> Query GlobalEventQuery (V.Vector (GlobalPosition, RecordedEvent))+qGlobalEvent readDir =+ statement sql encoder decoder True+ where+ sql =+ "SELECT "+ <> "position, stream, id, version, type, data, meta_data, created "+ <> "FROM "+ <> "events "+ <> "WHERE position "+ <> (if readDir == RdForward then ">=" else "<=")+ <> " $1 "+ <> "ORDER BY position "+ <> (if readDir == RdForward then "ASC" else "DESC")+ <> " LIMIT $2 "+ encoder =+ contramap geq_position (E.value encGlobalPosition)+ <> contramap (fromIntegral . geq_limit) (E.value E.int8)+ decoder =+ D.rowsVector $+ (,)+ <$> (GlobalPosition <$> D.value D.int8)+ <*> decRecordedEvent++dbWriteToStream ::+ DbStore+ -> StreamId+ -> ExpectedVersion+ -> V.Vector EventData+ -> IO WriteResult+dbWriteToStream db streamId ev events =+ withPool (db_store db) $+ dbTx Tx.Serializable Tx.Write $+ do myVersion <- Tx.query streamId qStreamVersion+ case ev of+ EvAny -> continueIf True myVersion+ EvNoStream -> continueIf (isNothing myVersion) myVersion+ EvStreamExists -> continueIf (isJust myVersion) myVersion+ EvExact expected -> continueIf (myVersion == Just expected) myVersion+ where+ continue vers =+ flip V.imapM_ events $ \idx event ->+ do let we =+ WriteEvent+ { we_stream = streamId+ , we_number =+ incrementTimes idx $+ nextEventNumber $ fromMaybe firstEventNumber vers+ , we_data = event+ }+ Tx.query we qWriteEvent+ continueIf cond vers =+ if cond+ then do continue vers+ pure WrSuccess+ else pure WrWrongExpectedVersion++instance EventStoreWriter IO DbStore where+ writeToStream = dbWriteToStream++dbReadEvent ::+ DbStore+ -> StreamId+ -> EventNumber+ -> IO EventReadResult+dbReadEvent store streamId eventNumber =+ withPool (db_store store) $+ do res <- S.query (SingleEventQuery streamId eventNumber) qSingleEvent+ case res of+ Nothing -> pure ErrFailed+ Just v -> pure (ErrValue v)++dbReadStreamEvents ::+ DbStore+ -> StreamId -> EventNumber -> Int -> ReadDirection+ -> IO (V.Vector RecordedEvent)+dbReadStreamEvents store streamId eventNumber size readDir =+ withPool (db_store store) $+ S.query (EventStreamQuery streamId eventNumber size) (qEventStream readDir)++dbReadAllEvents ::+ DbStore+ -> GlobalPosition -> Int -> ReadDirection+ -> IO (V.Vector (GlobalPosition, RecordedEvent))+dbReadAllEvents store eventNumber size readDir =+ withPool (db_store store) $+ S.query (GlobalEventQuery eventNumber size) (qGlobalEvent readDir)++instance EventStoreReader IO DbStore where+ readEvent = dbReadEvent+ readStreamEvents = dbReadStreamEvents+ readAllEvents = dbReadAllEvents++-- | Poor mans subscriber implementation as 'hasql' does not support+-- LISTEN/NOTIFY. Could replace with REDIS?+dbSubscribeTo ::+ DbStore+ -> SubscriptionConfig+ -> ConduitM () RecordedEvent IO ()+dbSubscribeTo store config =+ do startNumber <-+ case startPosition of+ SspBeginning -> pure firstEventNumber+ SspFrom x -> pure x+ SspCurrent ->+ liftIO $ withPool (db_store store) $+ fromMaybe firstEventNumber <$> S.query streamId qStreamVersion+ innerLoop startNumber+ where+ innerLoop !pos =+ do batch <-+ liftIO $+ dbReadStreamEvents store streamId pos 1000 RdForward+ if V.null batch+ then do liftIO $ sleepTS (milliseconds 500)+ innerLoop pos+ else do V.mapM_ yield batch+ let lastEl = V.last batch+ innerLoop (nextEventNumber (re_number lastEl))++ streamId = sc_stream config+ startPosition = sc_startPosition config++instance EventStoreSubscriber IO DbStore where+ subscribeTo = dbSubscribeTo++schemaV1 :: BS.ByteString+schemaV1 = [s|+CREATE TABLE events (+ id UUID NOT NULL,+ stream TEXT NOT NULL,+ version INT8 NOT NULL,+ position SERIAL8 NOT NULL,+ type TEXT NOT NULL,+ data JSONB NOT NULL,+ meta_data JSONB NOT NULL,+ created timestamptz NOT NULL DEFAULT NOW(),+ CONSTRAINT "event_id" PRIMARY KEY (id, stream),+ UNIQUE (stream, version),+ CONSTRAINT valid_version CHECK (version >= position)+);++CREATE INDEX event_stream ON events (stream);+CREATE INDEX event_stream_version ON events (stream, version);+CREATE INDEX event_position ON events (position);+|]
+ src/SuperEvent/Store/Hasql/Utils.hs view
@@ -0,0 +1,20 @@+module SuperEvent.Store.Hasql.Utils where++import qualified Hasql.Pool as P+import qualified Hasql.Session as S+import qualified Hasql.Transaction as Tx+import qualified Hasql.Transaction.Sessions as Tx++newtype Store+ = Store+ { unStore :: P.Pool }++dbTx :: Tx.IsolationLevel -> Tx.Mode -> Tx.Transaction a -> S.Session a+dbTx = Tx.transaction++withPool :: Store -> S.Session a -> IO a+withPool pss sess =+ do res <- P.use (unStore pss) sess+ case res of+ Left err -> fail (show err)+ Right ok -> pure ok
+ src/SuperEvent/Store/Types.hs view
@@ -0,0 +1,108 @@+{-# LANGUAGE StrictData #-}+{-# LANGUAGE MultiParamTypeClasses #-}+{-# LANGUAGE FunctionalDependencies #-}+{-# LANGUAGE GeneralizedNewtypeDeriving #-}+module SuperEvent.Store.Types where++import Data.Aeson+import Data.Conduit+import Data.Hashable+import Data.Int+import Data.Text (Text)+import Data.Time+import Data.UUID (UUID)+import Data.Vector (Vector)++newtype EventType+ = EventType { unEventType :: Text }+ deriving (Show, Eq, Ord)++newtype EventNumber+ = EventNumber { unEventNumber :: Int64 }+ deriving (Show, Eq, Ord)++firstEventNumber ::EventNumber+firstEventNumber = EventNumber 0++nextEventNumber :: EventNumber -> EventNumber+nextEventNumber (EventNumber x) = EventNumber (x + 1)++incrementTimes :: Int -> EventNumber -> EventNumber+incrementTimes n (EventNumber x) = EventNumber (x + fromIntegral n)++newtype GlobalPosition+ = GlobalPosition { unGlobalPosition :: Int64 }+ deriving (Show, Eq, Ord)++data EventData+ = EventData+ { ed_guid :: UUID+ , ed_type :: EventType+ , ed_data :: Value+ , ed_metadata :: Value+ } deriving (Show, Eq)++data ExpectedVersion+ = EvAny+ | EvNoStream+ | EvStreamExists+ | EvExact EventNumber+ deriving (Show, Eq)++newtype StreamId+ = StreamId { unStreamId :: Text }+ deriving (Show, Eq, Ord, Hashable)++data WriteResult+ = WrSuccess+ | WrWrongExpectedVersion+ deriving (Show, Eq)++class EventStoreWriter m es | es -> m where+ writeToStream ::+ es -> StreamId -> ExpectedVersion -> Vector EventData -> m WriteResult++data RecordedEvent+ = RecordedEvent+ { re_stream :: StreamId+ , re_guid :: UUID+ , re_number :: EventNumber+ , re_type :: EventType+ , re_data :: Value+ , re_metadata :: Value+ , re_created :: UTCTime+ } deriving (Show, Eq)++data ReadDirection+ = RdForward+ | RdBackward+ deriving (Show, Eq, Ord, Enum, Bounded)++data EventReadResult+ = ErrFailed+ | ErrValue RecordedEvent+ deriving (Show, Eq)++class EventStoreReader m es | es -> m where+ readEvent :: es -> StreamId -> EventNumber -> m EventReadResult+ readStreamEvents ::+ es -> StreamId -> EventNumber -> Int -> ReadDirection+ -> m (Vector RecordedEvent)+ readAllEvents ::+ es -> GlobalPosition -> Int -> ReadDirection+ -> m (Vector (GlobalPosition, RecordedEvent))++data SubscriptionStartPosition+ = SspBeginning+ | SspFrom EventNumber+ | SspCurrent+ deriving (Show, Eq)++data SubscriptionConfig+ = SubscriptionConfig+ { sc_startPosition :: SubscriptionStartPosition+ , sc_stream :: StreamId+ } deriving (Show, Eq)++class EventStoreSubscriber m es | es -> m where+ subscribeTo :: es -> SubscriptionConfig -> ConduitM () RecordedEvent m ()
+ superevent.cabal view
@@ -0,0 +1,90 @@+-- This file has been generated from package.yaml by hpack version 0.20.0.+--+-- see: https://github.com/sol/hpack+--+-- hash: 1636a007f3d46981133cef8a5a092ad267cf35b4404cb0ce81c84ee90cc74890++name: superevent+version: 0.1.0.0+synopsis: A simple opinionated event store implementation+description: A simple opinionated event store implementation+category: Database+homepage: https://github.com/agrafix/superevent#readme+bug-reports: https://github.com/agrafix/superevent/issues+author: Alexander Thiemann+maintainer: mail@athiemann.net+copyright: 2018 Alexander Thiemann <mail@athiemann.net>+license: BSD3+license-file: LICENSE+build-type: Simple+cabal-version: >= 1.10++extra-source-files:+ ChangeLog.md+ README.md++source-repository head+ type: git+ location: https://github.com/agrafix/superevent++library+ hs-source-dirs:+ src+ ghc-options: -Wall+ build-depends:+ aeson+ , async+ , base >=4.7 && <5+ , bytestring+ , conduit+ , containers+ , contravariant+ , hashable+ , hasql+ , hasql-migration+ , hasql-pool+ , hasql-transaction+ , mtl+ , random+ , stm+ , string-qq+ , text+ , time+ , timespan+ , transformers+ , unordered-containers+ , uuid+ , vector+ exposed-modules:+ SuperEvent.Store.Hasql+ SuperEvent.Store.Hasql.Utils+ SuperEvent.Store.Types+ other-modules:+ Paths_superevent+ default-language: Haskell2010++test-suite superevent-test+ type: exitcode-stdio-1.0+ main-is: Spec.hs+ hs-source-dirs:+ test+ ghc-options: -threaded -rtsopts -with-rtsopts=-N+ build-depends:+ QuickCheck+ , aeson+ , async+ , base >=4.7 && <5+ , conduit+ , hspec+ , mtl+ , stm+ , superevent+ , temporary+ , text+ , transformers+ , uuid+ , vector+ other-modules:+ SuperEvent.Store.HasqlSpec+ Paths_superevent+ default-language: Haskell2010
+ test/Spec.hs view
@@ -0,0 +1,1 @@+{-# OPTIONS_GHC -F -pgmF hspec-discover #-}
+ test/SuperEvent/Store/HasqlSpec.hs view
@@ -0,0 +1,128 @@+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE OverloadedStrings #-}+module SuperEvent.Store.HasqlSpec (spec) where++import SuperEvent.Store.Hasql+import SuperEvent.Store.Types++import Control.Concurrent.Async+import Control.Concurrent.STM+import Control.Monad+import Control.Monad.Trans+import Data.Aeson+import Data.Conduit+import Test.Hspec+import qualified Data.Text as T+import qualified Data.UUID.V4 as UUID+import qualified Data.Vector as V++spec :: Spec+spec =+ around withTempStore $+ do it "should work for simple reads" simpleWriteRead+ it "should work for simple stream reads" simpleWriteReadStream+ it "should work for simple global reads" simpleWriteReadGlobal+ it "should work for simple subscriptions" simpleWriteSubStream++simpleWriteRead ::+ (EventStoreReader IO es, EventStoreWriter IO es)+ => es -> IO ()+simpleWriteRead store =+ do let payload :: T.Text+ payload = "Hello"+ entry guid =+ EventData guid (EventType "foo") (toJSON payload) (toJSON ())+ stream = StreamId "text"+ guid <- UUID.nextRandom+ writeRes <- writeToStream store stream EvAny (V.singleton $ entry guid)+ writeRes `shouldBe` WrSuccess+ evt <- readEvent store stream (nextEventNumber firstEventNumber)+ case evt of+ ErrFailed -> expectationFailure "Read failed"+ ErrValue re -> re_guid re `shouldBe` guid++simpleWriteReadStream ::+ (EventStoreReader IO es, EventStoreWriter IO es)+ => es -> IO ()+simpleWriteReadStream store =+ do let payload :: T.Text+ payload = "Hello"+ entry guid =+ EventData guid (EventType "foo") (toJSON payload) (toJSON ())+ stream = StreamId "text"+ events <-+ V.forM (V.fromList [0..99]) $ \(_ :: Int) ->+ do guid <- UUID.nextRandom+ pure (entry guid)+ writeRes <- writeToStream store stream EvAny events+ writeRes `shouldBe` WrSuccess+ let writtenGuids = V.map ed_guid events+ readHelper evtNum lim =+ V.map re_guid <$> readStreamEvents store stream evtNum lim RdForward+ evts <- readHelper firstEventNumber 10+ evts `shouldBe` V.take 10 writtenGuids++ evts2 <- readHelper (incrementTimes 11 firstEventNumber) 100+ evts2 `shouldBe` V.take 90 (V.drop 10 writtenGuids)++simpleWriteReadGlobal ::+ (EventStoreReader IO es, EventStoreWriter IO es)+ => es -> IO ()+simpleWriteReadGlobal store =+ do let payload :: T.Text+ payload = "Hello"+ entry guid =+ EventData guid (EventType "foo") (toJSON payload) (toJSON ())+ stream = StreamId "text"+ events <-+ V.forM (V.fromList [0..99]) $ \(_ :: Int) ->+ do guid <- UUID.nextRandom+ pure (entry guid)+ writeRes <- writeToStream store stream EvAny events+ writeRes `shouldBe` WrSuccess+ let writtenGuids = V.map ed_guid events+ readHelper pos lim =+ V.map (re_guid . snd) <$> readAllEvents store pos lim RdForward+ evts <- readHelper (GlobalPosition 0) 10+ evts `shouldBe` V.take 10 writtenGuids++ evts2 <- readHelper (GlobalPosition 11) 100+ evts2 `shouldBe` V.take 90 (V.drop 10 writtenGuids)++simpleWriteSubStream ::+ (EventStoreSubscriber IO es, EventStoreReader IO es, EventStoreWriter IO es)+ => es -> IO ()+simpleWriteSubStream store =+ do let payload :: T.Text+ payload = "Hello"+ entry guid =+ EventData guid (EventType "foo") (toJSON payload) (toJSON ())+ stream = StreamId "text"+ outVar <- newTVarIO []+ poller <-+ async $+ do let consumer =+ do val <- await+ case val of+ Nothing ->+ do liftIO $ putStrLn "Consumer was terminated"+ pure ()+ Just v ->+ do liftIO $ atomically $ modifyTVar' outVar ((:) v)+ consumer+ subscribeTo store (SubscriptionConfig SspBeginning stream) $$ consumer+ events <-+ V.forM (V.fromList [0..9999]) $ \(_ :: Int) ->+ do guid <- UUID.nextRandom+ pure (entry guid)+ writeRes <- writeToStream store stream EvAny events+ writeRes `shouldBe` WrSuccess+ let writtenGuids = V.map ed_guid events+ finalResult <-+ atomically $+ do vals <- readTVar outVar+ when (length vals < 9999) retry+ pure (V.map re_guid $ V.reverse $ V.fromList vals)+ uninterruptibleCancel poller+ finalResult `shouldBe` writtenGuids