packages feed

sparrow 0.0.1.6 → 0.0.2.0

raw patch · 4 files changed

+79/−65 lines, 4 files

Files

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 ()         }       }