sparrow 0.0.1.6 → 0.0.2.0
raw patch · 4 files changed
+79/−65 lines, 4 files
Files
- sparrow.cabal +2/−2
- src/Web/Dependencies/Sparrow/Server.hs +17/−12
- src/Web/Dependencies/Sparrow/Server/Types.hs +46/−41
- src/Web/Dependencies/Sparrow/Types.hs +14/−10
sparrow.cabal view
@@ -2,10 +2,10 @@ -- -- see: https://github.com/sol/hpack ----- hash: 0ebfaafba5a716832f29a853a6fe6660eafd4c006e46eb404dcef9e832f7bf62+-- hash: 309095bf67849fabbab5152d18d36076dcba4000675a82d0129d9a121f9ad671 name: sparrow-version: 0.0.1.6+version: 0.0.2.0 synopsis: Unified streaming dependency management for web apps description: Please see the README on Github at <https://git.localcooking.com/tooling/sparrow#readme> category: Web
src/Web/Dependencies/Sparrow/Server.hs view
@@ -60,7 +60,9 @@ import qualified Data.UUID as UUID import Data.Singleton.Class (Extractable (runSingleton)) import Data.Proxy (Proxy (..))-import Control.Monad (join, forever, forM_)+import Data.Foldable (Foldable)+import Control.Applicative (Alternative)+import Control.Monad (join, forever) import Control.Monad.Trans (lift) import Control.Monad.State (modify') import Control.Monad.IO.Class (MonadIO (..))@@ -85,7 +87,7 @@ -- | Called per-connection-unpackServer :: forall m stM http initIn initOut deltaIn deltaOut+unpackServer :: forall m f stM http initIn initOut deltaIn deltaOut . MonadIO m => Aligned.MonadBaseControl IO m stM => Extractable stM@@ -93,9 +95,10 @@ => ToJSON initOut => FromJSON deltaIn => ToJSON deltaOut+ => Foldable f => Topic -- ^ Name of Dependency- -> Server m initIn initOut deltaIn deltaOut -- ^ Handler for all clients- -> SparrowServerT http m (MiddlewareT m)+ -> Server m f initIn initOut deltaIn deltaOut -- ^ Handler for all clients+ -> SparrowServerT http f m (MiddlewareT m) unpackServer topic server = do env <- ask' @@ -173,7 +176,7 @@ => Match xs' xs childHttp resultHttp => UrlChunks xs -- ^ Should match the dependency name -> childHttp -- ^ 'Network.Wai.Trans.MiddlewareT', or a function to one- -> SparrowServerT resultHttp m ()+ -> SparrowServerT resultHttp f m () match ts http = tell' (singleton ts http) @@ -188,8 +191,8 @@ matchGroup :: Monad m => MatchGroup xs' xs childHttp resultHttp => UrlChunks xs -- ^ Common 'Topic' prefix- -> SparrowServerT childHttp m () -- ^ Set of handlers- -> SparrowServerT resultHttp m ()+ -> SparrowServerT childHttp f m () -- ^ Set of handlers+ -> SparrowServerT resultHttp f m () matchGroup ts x = do env <- ask' http <- lift (execSparrowServerT' env x)@@ -197,13 +200,15 @@ -- | Host dependencies and websocket-serveDependencies :: forall m stM sec a+serveDependencies :: forall f m stM sec a . MonadBaseControl IO m => Aligned.MonadBaseControl IO m stM => Extractable stM => MonadIO m => MonadCatch m- => SparrowServerT (MiddlewareT m) m a -- ^ Dependencies+ => Foldable f+ => Alternative f+ => SparrowServerT (MiddlewareT m) f m a -- ^ Dependencies -> m (RouterT (MiddlewareT m) sec m ()) serveDependencies server = Aligned.liftBaseWith $ \runInBase -> do let runM :: forall b. m b -> IO b@@ -232,7 +237,7 @@ listener <- async $ forever $ do x <- atomically (TMapChan.lookup envSessionsOutgoing sessionID) runM (send x)- atomically $ putTMVar outgoingListener listener+ atomically (putTMVar outgoingListener listener) initSubs <- liftIO $ atomically $ getCurrentRegisteredTopics env sessionID @@ -250,7 +255,7 @@ -- update client of removed subscription send (WSTopicRemoved topic) - liftIO $ killOnOpenThreads env sessionID topic+ liftIO (killOnOpenThreads env sessionID topic) WSIncoming (WithTopic topic x) -> do mEff <- liftIO $ atomically $ getCallReceive env sessionID topic x@@ -271,7 +276,7 @@ delSubscriberFromAllTopics env sessionID callAllOnUnsubscribe env sessionID- liftIO $ killAllOnOpenThreads env sessionID+ liftIO (killAllOnOpenThreads env sessionID) } wsApp' <- pingPong ((10^6) * 10) wsApp -- every 10 seconds
src/Web/Dependencies/Sparrow/Server/Types.hs view
@@ -55,13 +55,14 @@ import Data.Monoid ((<>)) import Data.Maybe (fromMaybe) import Data.Proxy (Proxy (..))-import Data.Foldable (sequenceA_)+import Data.Foldable (sequenceA_, asum) import Data.Aeson (FromJSON, Value) import qualified Data.Aeson as Aeson import Data.HashMap.Strict (HashMap) import qualified Data.HashMap.Strict as HM import Data.HashSet (HashSet) import qualified Data.HashSet as HS+import Control.Applicative (Alternative) import Control.Monad (forM_) import Control.Monad.Reader (ReaderT (..), runReaderT, ask, MonadReader (..)) import Control.Monad.State (StateT, execStateT, modify', MonadState (..))@@ -81,7 +82,7 @@ type SessionsOutgoing = TMapChan SessionID (WSOutgoing (WithTopic Value)) -sendTo :: Env m -> SessionID -> WSOutgoing (WithTopic Value) -> STM ()+sendTo :: Env f m -> SessionID -> WSOutgoing (WithTopic Value) -> STM () sendTo Env{envSessionsOutgoing} = TMapChan.insert envSessionsOutgoing @@ -92,7 +93,7 @@ -- unsafe because receiver isn't type-caste to the expected Topic unsafeRegisterReceive :: MonadIO m- => Env m -> SessionID -> Topic -> (Value -> Maybe (m ())) -> STM ()+ => Env f m -> SessionID -> Topic -> (Value -> Maybe (m ())) -> STM () unsafeRegisterReceive Env{envRegisteredReceive} sID topic f = do mTopics <- TMapMVar.tryObserve envRegisteredReceive sID case mTopics of@@ -106,26 +107,26 @@ TMapMVar.insertForce topics topic f -unregisterReceive :: Env m -> SessionID -> Topic -> STM ()+unregisterReceive :: Env f m -> SessionID -> Topic -> STM () unregisterReceive Env{envRegisteredReceive} sID topic = do mTopics <- TMapMVar.tryObserve envRegisteredReceive sID case mTopics of Nothing -> pure () Just topics -> TMapMVar.delete topics topic -unregisterSession :: Env m -> SessionID -> STM ()+unregisterSession :: Env f m -> SessionID -> STM () unregisterSession Env{envRegisteredReceive} = TMapMVar.delete envRegisteredReceive getCallReceive :: MonadIO m- => Env m -> SessionID -> Topic -> Value -> STM (Maybe (m ()))+ => Env f m -> SessionID -> Topic -> Value -> STM (Maybe (m ())) getCallReceive Env{envRegisteredReceive} sID topic v = do topics <- TMapMVar.observe envRegisteredReceive sID onReceive <- TMapMVar.observe topics topic pure (onReceive v) -getCurrentRegisteredTopics :: Env m -> SessionID -> STM [Topic]+getCurrentRegisteredTopics :: Env f m -> SessionID -> STM [Topic] getCurrentRegisteredTopics Env{envRegisteredReceive} sID = do mTopics <- TMapMVar.tryObserve envRegisteredReceive sID case mTopics of@@ -134,27 +135,27 @@ type RegisteredTopicInvalidators = TVar (HashMap Topic (Value -> Maybe String)) -registerInvalidator :: forall deltaIn m+registerInvalidator :: forall deltaIn f m . FromJSON deltaIn- => Env m -> Topic -> Proxy deltaIn -> STM ()+ => Env f m -> Topic -> Proxy deltaIn -> STM () registerInvalidator Env{envRegisteredTopicInvalidators} topic Proxy = let go v = case Aeson.fromJSON v of Aeson.Error e -> Just e Aeson.Success (_ :: deltaIn) -> Nothing in modifyTVar' envRegisteredTopicInvalidators (HM.insert topic go) -getValidator :: Env m -> Topic -> STM (Maybe (Value -> Maybe String))+getValidator :: Env f m -> Topic -> STM (Maybe (Value -> Maybe String)) getValidator Env{envRegisteredTopicInvalidators} topic = HM.lookup topic <$> readTVar envRegisteredTopicInvalidators type RegisteredTopicSubscribers = TVar (HashMap Topic (HashSet SessionID)) -addSubscriber :: Env m -> Topic -> SessionID -> STM ()+addSubscriber :: Env f m -> Topic -> SessionID -> STM () addSubscriber Env{envRegisteredTopicSubscribers} topic sID = modifyTVar' envRegisteredTopicSubscribers (HM.alter (Just . maybe (HS.singleton sID) (HS.insert sID)) topic) -delSubscriber :: Env m -> Topic -> SessionID -> STM ()+delSubscriber :: Env f m -> Topic -> SessionID -> STM () delSubscriber Env{envRegisteredTopicSubscribers} topic sID = let go xs | xs == HS.singleton sID = Nothing@@ -162,24 +163,24 @@ in modifyTVar' envRegisteredTopicSubscribers (HM.alter (maybe Nothing go) topic) -delSubscriberFromAllTopics :: Env m -> SessionID -> STM ()+delSubscriberFromAllTopics :: Env f m -> SessionID -> STM () delSubscriberFromAllTopics env@Env{envRegisteredTopicSubscribers} sID = do allTopics <- HM.keys <$> readTVar envRegisteredTopicSubscribers forM_ allTopics (\topic -> delSubscriber env topic sID) -getSubscribers :: Env m -> Topic -> STM [SessionID]+getSubscribers :: Env f m -> Topic -> STM [SessionID] getSubscribers Env{envRegisteredTopicSubscribers} topic = maybe [] HS.toList . HM.lookup topic <$> readTVar envRegisteredTopicSubscribers type RegisteredOnUnsubscribe m = TVar (HashMap SessionID (HashMap Topic (m ()))) -registerOnUnsubscribe :: Env m -> SessionID -> Topic -> m () -> STM ()+registerOnUnsubscribe :: Env f m -> SessionID -> Topic -> m () -> STM () registerOnUnsubscribe Env{envRegisteredOnUnsubscribe} sID topic eff = do xs <- readTVar envRegisteredOnUnsubscribe let topics = fromMaybe HM.empty (HM.lookup sID xs) modifyTVar' envRegisteredOnUnsubscribe (HM.insert sID (HM.insert topic eff topics)) -callOnUnsubscribe :: MonadIO m => Env m -> SessionID -> Topic -> m ()+callOnUnsubscribe :: MonadIO m => Env f m -> SessionID -> Topic -> m () callOnUnsubscribe Env{envRegisteredOnUnsubscribe} sID topic = do mEff <- liftIO $ atomically $ do xs <- readTVar envRegisteredOnUnsubscribe@@ -191,7 +192,7 @@ pure x fromMaybe (pure ()) mEff -callAllOnUnsubscribe :: MonadIO m => Env m -> SessionID -> m ()+callAllOnUnsubscribe :: MonadIO m => Env f m -> SessionID -> m () callAllOnUnsubscribe Env{envRegisteredOnUnsubscribe} sID = do effs <- liftIO $ atomically $ do xs <- readTVar envRegisteredOnUnsubscribe@@ -204,56 +205,60 @@ sequenceA_ effs -type RegisteredOnOpenThreads =- TVar (HashMap SessionID (HashMap Topic [Async ()]))+type RegisteredOnOpenThreads f =+ TVar (HashMap SessionID (HashMap Topic (TVar (f (Async ()))))) -registerOnOpenThreads :: Env m -> SessionID -> Topic -> [Async ()] -> STM ()+registerOnOpenThreads :: Env f m -> SessionID -> Topic -> TVar (f (Async ())) -> STM () registerOnOpenThreads Env{envRegisteredOnOpenThreads} sID topic thread = do xs <- readTVar envRegisteredOnOpenThreads let topics = fromMaybe HM.empty (HM.lookup sID xs) modifyTVar' envRegisteredOnOpenThreads (HM.insert sID (HM.insert topic thread topics)) -killOnOpenThreads :: MonadIO m => Env m -> SessionID -> Topic -> IO ()+killOnOpenThreads :: Foldable f => MonadIO m => Env f m -> SessionID -> Topic -> IO () killOnOpenThreads Env{envRegisteredOnOpenThreads} sID topic = do mThread <- atomically $ do xs <- readTVar envRegisteredOnOpenThreads- case HM.lookup sID xs of+ mThreads' <- case HM.lookup sID xs of Nothing -> pure Nothing Just topics -> do let x = HM.lookup topic topics modifyTVar' envRegisteredOnOpenThreads (HM.adjust (HM.delete topic) sID) pure x+ case mThreads' of+ Nothing -> pure Nothing+ Just threadsRef -> Just <$> readTVar threadsRef case mThread of Nothing -> pure () Just threads -> mapM_ cancel threads -killAllOnOpenThreads :: MonadIO m => Env m -> SessionID -> IO ()+killAllOnOpenThreads :: Foldable f => Alternative f => MonadIO m => Env f m -> SessionID -> IO () killAllOnOpenThreads Env{envRegisteredOnOpenThreads} sID = do threads <- atomically $ do xs <- readTVar envRegisteredOnOpenThreads- case HM.lookup sID xs of+ threadRefs <- case HM.lookup sID xs of Nothing -> pure [] Just topics -> do let x = HM.elems topics modifyTVar' envRegisteredOnOpenThreads (HM.delete sID) pure x- forM_ (concat threads) cancel+ mapM readTVar threadRefs+ forM_ (asum threads) cancel -data Env m = Env+data Env f m = Env { envSessionsOutgoing :: {-# UNPACK #-} !SessionsOutgoing , envRegisteredReceive :: {-# UNPACK #-} !(RegisteredReceive m) , envRegisteredTopicInvalidators :: {-# UNPACK #-} !RegisteredTopicInvalidators , envRegisteredTopicSubscribers :: {-# UNPACK #-} !RegisteredTopicSubscribers , envRegisteredOnUnsubscribe :: {-# UNPACK #-} !(RegisteredOnUnsubscribe m)- , envRegisteredOnOpenThreads :: {-# UNPACK #-} !RegisteredOnOpenThreads+ , envRegisteredOnOpenThreads :: {-# UNPACK #-} !(RegisteredOnOpenThreads f) } -newEnv :: IO (Env m)+newEnv :: IO (Env f m) newEnv = Env <$> atomically newTMapChan <*> atomically newTMapMVar@@ -263,14 +268,14 @@ <*> newTVarIO HM.empty -unsafeBroadcastTopic :: MonadIO m => Env m -> Topic -> Value -> m ()+unsafeBroadcastTopic :: MonadIO m => Env f m -> Topic -> Value -> m () unsafeBroadcastTopic env t v = liftIO $ atomically $ do ss <- getSubscribers env t forM_ ss (\sessionID -> sendTo env sessionID (WSOutgoing (WithTopic t v))) -broadcaster :: MonadIO m => Env m -> Broadcast m+broadcaster :: MonadIO m => Env f m -> Broadcast m broadcaster env = \topic -> do mInvalidator <- liftIO $ atomically $ getValidator env topic case mInvalidator of@@ -284,18 +289,18 @@ type Paper http = RootedPredTrie Text http -newtype SparrowServerT http m a = SparrowServerT- { runSparrowServerT :: ReaderT (Env m) (StateT (Paper http) m) a+newtype SparrowServerT http f m a = SparrowServerT+ { runSparrowServerT :: ReaderT (Env f m) (StateT (Paper http) m) a } deriving (Functor, Applicative, Monad, MonadIO, MonadWriter w, MonadCatch, MonadThrow, MonadMask) -instance MonadTrans (SparrowServerT http) where+instance MonadTrans (SparrowServerT http f) where lift x = SparrowServerT (lift (lift x)) -instance MonadReader r m => MonadReader r (SparrowServerT http m) where+instance MonadReader r m => MonadReader r (SparrowServerT http f m) where ask = lift ask local f (SparrowServerT (ReaderT g)) = SparrowServerT $ ReaderT $ \env -> local f (g env) -instance MonadState s m => MonadState s (SparrowServerT http m) where+instance MonadState s m => MonadState s (SparrowServerT http f m) where get = lift get put x = lift (put x) @@ -310,22 +315,22 @@ execSparrowServerT :: MonadIO m- => SparrowServerT http m a- -> m (Paper http, Env m)+ => SparrowServerT http f m a+ -> m (Paper http, Env f m) execSparrowServerT x = do env <- liftIO newEnv (,env) <$> execSparrowServerT' env x execSparrowServerT' :: Monad m- => Env m- -> SparrowServerT http m a+ => Env f m+ -> SparrowServerT http f m a -> m (Paper http) execSparrowServerT' env (SparrowServerT x) = execStateT (runReaderT x env) mempty -tell' :: Monad m => Paper http -> SparrowServerT http m ()+tell' :: Monad m => Paper http -> SparrowServerT http f m () tell' x = SparrowServerT (lift (modify' (<> x))) -ask' :: Monad m => SparrowServerT http m (Env m)+ask' :: Monad m => SparrowServerT http f m (Env f m) ask' = SparrowServerT ask
src/Web/Dependencies/Sparrow/Types.hs view
@@ -18,9 +18,11 @@ import Data.Aeson.Types (typeMismatch) import Data.Aeson.Attoparsec (attoAeson) import Data.Attoparsec.Text (Parser, takeWhile1, char, sepBy)-import Control.Applicative ((<|>))+import Control.Applicative (Alternative (empty), (<|>)) import Control.DeepSeq (NFData)+import Control.Monad.IO.Class (MonadIO (liftIO)) import Control.Concurrent.Async (Async)+import Control.Concurrent.STM (TVar, newTVarIO) import GHC.Generics (Generic) @@ -33,28 +35,29 @@ , serverSendCurrent :: deltaOut -> m () } -data ServerReturn m initOut deltaIn deltaOut = ServerReturn+data ServerReturn m f initOut deltaIn deltaOut = ServerReturn { serverInitOut :: initOut , serverOnOpen :: ServerArgs m deltaOut- -> m [Async ()]+ -> m (TVar (f (Async ()))) -- ^ invoked once, and should return a 'Control.Concurrent.Async.link'ed long-lived thread -- to kill when the subscription dies , serverOnReceive :: ServerArgs m deltaOut -> deltaIn -> m () -- ^ invoked for each receive } -data ServerContinue m initOut deltaIn deltaOut = ServerContinue- { serverContinue :: Broadcast m -> m (ServerReturn m initOut deltaIn deltaOut)+data ServerContinue m f initOut deltaIn deltaOut = ServerContinue+ { serverContinue :: Broadcast m -> m (ServerReturn m f initOut deltaIn deltaOut) , serverOnUnsubscribe :: m () } -type Server m initIn initOut deltaIn deltaOut =- initIn -> m (Maybe (ServerContinue m initOut deltaIn deltaOut))+type Server m f initIn initOut deltaIn deltaOut =+ initIn -> m (Maybe (ServerContinue m f initOut deltaIn deltaOut)) -staticServer :: Monad m+staticServer :: MonadIO m+ => Alternative f => (initIn -> m (Maybe initOut)) -- ^ Produce an initOut- -> Server m initIn initOut JSONVoid JSONVoid+ -> Server m f initIn initOut JSONVoid JSONVoid staticServer f initIn = do mInitOut <- f initIn case mInitOut of@@ -65,7 +68,8 @@ { serverInitOut = initOut , serverOnOpen = \ServerArgs{serverDeltaReject} -> do serverDeltaReject- pure []+ threadVar <- liftIO (newTVarIO empty)+ pure threadVar , serverOnReceive = \_ _ -> pure () } }