packages feed

mu-grpc-server (empty) → 0.1.0.0

raw patch · 6 files changed

+611/−0 lines, 6 filesdep +asyncdep +basedep +bytestringsetup-changed

Dependencies added: async, base, bytestring, conduit, http2-grpc-proto3-wire, http2-grpc-types, mtl, mu-protobuf, mu-rpc, mu-schema, sop-core, stm, stm-conduit, wai, warp, warp-grpc, warp-tls

Files

+ CHANGELOG.md view
@@ -0,0 +1,5 @@+# Revision history for mu-haskell++## 0.1.0.0 -- YYYY-mm-dd++* First version. Released on an unsuspecting world.
+ LICENSE view
@@ -0,0 +1,202 @@++                                 Apache License+                           Version 2.0, January 2004+                        http://www.apache.org/licenses/++   TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION++   1. Definitions.++      "License" shall mean the terms and conditions for use, reproduction,+      and distribution as defined by Sections 1 through 9 of this document.++      "Licensor" shall mean the copyright owner or entity authorized by+      the copyright owner that is granting the License.++      "Legal Entity" shall mean the union of the acting entity and all+      other entities that control, are controlled by, or are under common+      control with that entity. For the purposes of this definition,+      "control" means (i) the power, direct or indirect, to cause the+      direction or management of such entity, whether by contract or+      otherwise, or (ii) ownership of fifty percent (50%) or more of the+      outstanding shares, or (iii) beneficial ownership of such entity.++      "You" (or "Your") shall mean an individual or Legal Entity+      exercising permissions granted by this License.++      "Source" form shall mean the preferred form for making modifications,+      including but not limited to software source code, documentation+      source, and configuration files.++      "Object" form shall mean any form resulting from mechanical+      transformation or translation of a Source form, including but+      not limited to compiled object code, generated documentation,+      and conversions to other media types.++      "Work" shall mean the work of authorship, whether in Source or+      Object form, made available under the License, as indicated by a+      copyright notice that is included in or attached to the work+      (an example is provided in the Appendix below).++      "Derivative Works" shall mean any work, whether in Source or Object+      form, that is based on (or derived from) the Work and for which the+      editorial revisions, annotations, elaborations, or other modifications+      represent, as a whole, an original work of authorship. For the purposes+      of this License, Derivative Works shall not include works that remain+      separable from, or merely link (or bind by name) to the interfaces of,+      the Work and Derivative Works thereof.++      "Contribution" shall mean any work of authorship, including+      the original version of the Work and any modifications or additions+      to that Work or Derivative Works thereof, that is intentionally+      submitted to Licensor for inclusion in the Work by the copyright owner+      or by an individual or Legal Entity authorized to submit on behalf of+      the copyright owner. For the purposes of this definition, "submitted"+      means any form of electronic, verbal, or written communication sent+      to the Licensor or its representatives, including but not limited to+      communication on electronic mailing lists, source code control systems,+      and issue tracking systems that are managed by, or on behalf of, the+      Licensor for the purpose of discussing and improving the Work, but+      excluding communication that is conspicuously marked or otherwise+      designated in writing by the copyright owner as "Not a Contribution."++      "Contributor" shall mean Licensor and any individual or Legal Entity+      on behalf of whom a Contribution has been received by Licensor and+      subsequently incorporated within the Work.++   2. Grant of Copyright License. Subject to the terms and conditions of+      this License, each Contributor hereby grants to You a perpetual,+      worldwide, non-exclusive, no-charge, royalty-free, irrevocable+      copyright license to reproduce, prepare Derivative Works of,+      publicly display, publicly perform, sublicense, and distribute the+      Work and such Derivative Works in Source or Object form.++   3. Grant of Patent License. Subject to the terms and conditions of+      this License, each Contributor hereby grants to You a perpetual,+      worldwide, non-exclusive, no-charge, royalty-free, irrevocable+      (except as stated in this section) patent license to make, have made,+      use, offer to sell, sell, import, and otherwise transfer the Work,+      where such license applies only to those patent claims licensable+      by such Contributor that are necessarily infringed by their+      Contribution(s) alone or by combination of their Contribution(s)+      with the Work to which such Contribution(s) was submitted. If You+      institute patent litigation against any entity (including a+      cross-claim or counterclaim in a lawsuit) alleging that the Work+      or a Contribution incorporated within the Work constitutes direct+      or contributory patent infringement, then any patent licenses+      granted to You under this License for that Work shall terminate+      as of the date such litigation is filed.++   4. Redistribution. You may reproduce and distribute copies of the+      Work or Derivative Works thereof in any medium, with or without+      modifications, and in Source or Object form, provided that You+      meet the following conditions:++      (a) You must give any other recipients of the Work or+          Derivative Works a copy of this License; and++      (b) You must cause any modified files to carry prominent notices+          stating that You changed the files; and++      (c) You must retain, in the Source form of any Derivative Works+          that You distribute, all copyright, patent, trademark, and+          attribution notices from the Source form of the Work,+          excluding those notices that do not pertain to any part of+          the Derivative Works; and++      (d) If the Work includes a "NOTICE" text file as part of its+          distribution, then any Derivative Works that You distribute must+          include a readable copy of the attribution notices contained+          within such NOTICE file, excluding those notices that do not+          pertain to any part of the Derivative Works, in at least one+          of the following places: within a NOTICE text file distributed+          as part of the Derivative Works; within the Source form or+          documentation, if provided along with the Derivative Works; or,+          within a display generated by the Derivative Works, if and+          wherever such third-party notices normally appear. The contents+          of the NOTICE file are for informational purposes only and+          do not modify the License. You may add Your own attribution+          notices within Derivative Works that You distribute, alongside+          or as an addendum to the NOTICE text from the Work, provided+          that such additional attribution notices cannot be construed+          as modifying the License.++      You may add Your own copyright statement to Your modifications and+      may provide additional or different license terms and conditions+      for use, reproduction, or distribution of Your modifications, or+      for any such Derivative Works as a whole, provided Your use,+      reproduction, and distribution of the Work otherwise complies with+      the conditions stated in this License.++   5. Submission of Contributions. Unless You explicitly state otherwise,+      any Contribution intentionally submitted for inclusion in the Work+      by You to the Licensor shall be under the terms and conditions of+      this License, without any additional terms or conditions.+      Notwithstanding the above, nothing herein shall supersede or modify+      the terms of any separate license agreement you may have executed+      with Licensor regarding such Contributions.++   6. Trademarks. This License does not grant permission to use the trade+      names, trademarks, service marks, or product names of the Licensor,+      except as required for reasonable and customary use in describing the+      origin of the Work and reproducing the content of the NOTICE file.++   7. Disclaimer of Warranty. Unless required by applicable law or+      agreed to in writing, Licensor provides the Work (and each+      Contributor provides its Contributions) on an "AS IS" BASIS,+      WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or+      implied, including, without limitation, any warranties or conditions+      of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A+      PARTICULAR PURPOSE. You are solely responsible for determining the+      appropriateness of using or redistributing the Work and assume any+      risks associated with Your exercise of permissions under this License.++   8. Limitation of Liability. In no event and under no legal theory,+      whether in tort (including negligence), contract, or otherwise,+      unless required by applicable law (such as deliberate and grossly+      negligent acts) or agreed to in writing, shall any Contributor be+      liable to You for damages, including any direct, indirect, special,+      incidental, or consequential damages of any character arising as a+      result of this License or out of the use or inability to use the+      Work (including but not limited to damages for loss of goodwill,+      work stoppage, computer failure or malfunction, or any and all+      other commercial damages or losses), even if such Contributor+      has been advised of the possibility of such damages.++   9. Accepting Warranty or Additional Liability. While redistributing+      the Work or Derivative Works thereof, You may choose to offer,+      and charge a fee for, acceptance of support, warranty, indemnity,+      or other liability obligations and/or rights consistent with this+      License. However, in accepting such obligations, You may act only+      on Your own behalf and on Your sole responsibility, not on behalf+      of any other Contributor, and only if You agree to indemnify,+      defend, and hold each Contributor harmless for any liability+      incurred by, or claims asserted against, such Contributor by reason+      of your accepting any such warranty or additional liability.++   END OF TERMS AND CONDITIONS++   APPENDIX: How to apply the Apache License to your work.++      To apply the Apache License to your work, attach the following+      boilerplate notice, with the fields enclosed by brackets "[]"+      replaced with your own identifying information. (Don't include+      the brackets!)  The text should be enclosed in the appropriate+      comment syntax for the file format. We also recommend that a+      file or class name and description of purpose be included on the+      same "printed page" as the copyright notice for easier+      identification within third-party archives.++   Copyright © 2019-2020 47 Degrees. <http://47deg.com>++   Licensed under the Apache License, Version 2.0 (the "License");+   you may not use this file except in compliance with the License.+   You may obtain a copy of the License at++       http://www.apache.org/licenses/LICENSE-2.0++   Unless required by applicable law or agreed to in writing, software+   distributed under the License is distributed on an "AS IS" BASIS,+   WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.+   See the License for the specific language governing permissions and+   limitations under the License.
+ Setup.hs view
@@ -0,0 +1,2 @@+import Distribution.Simple+main = defaultMain
+ mu-grpc-server.cabal view
@@ -0,0 +1,67 @@+cabal-version:       >=1.10+name:                mu-grpc-server+version:             0.1.0.0+synopsis:            gRPC servers for Mu definitions+description:         With @mu-grpc-server@ you can easily build gRPC servers for mu-haskell!+license:             Apache-2.0+license-file:        LICENSE+author:              Alejandro Serrano, Flavio Corpa+maintainer:          alejandro.serrano@47deg.com+copyright:           Copyright © 2019-2020 <http://47deg.com 47 Degrees>+category:            Network+build-type:          Simple+extra-source-files:  CHANGELOG.md+homepage:            https://higherkindness.io/mu-haskell/+bug-reports:         https://github.com/higherkindness/mu-haskell/issues++source-repository head+  type:     git+  location: https://github.com/higherkindness/mu-haskell++library+  exposed-modules:     Mu.GRpc.Server+  build-depends:       base >=4.12 && <5+                     , async+                     , bytestring+                     , conduit+                     , http2-grpc-proto3-wire+                     , http2-grpc-types+                     , mtl+                     , mu-protobuf+                     , mu-rpc+                     , mu-schema+                     , sop-core+                     , stm+                     , stm-conduit+                     , wai+                     , warp+                     , warp-grpc+                     , warp-tls+  hs-source-dirs:      src+  default-language:    Haskell2010+  ghc-options:         -Wall+                       -fprint-potential-instances++executable grpc-example-server+  main-is:             ExampleServer.hs+  other-modules:       Mu.GRpc.Server+  build-depends:       base >=4.12 && <5+                     , async+                     , bytestring+                     , conduit+                     , http2-grpc-proto3-wire+                     , http2-grpc-types+                     , mtl+                     , mu-protobuf+                     , mu-rpc+                     , mu-schema+                     , sop-core+                     , stm+                     , stm-conduit+                     , wai+                     , warp+                     , warp-grpc+                     , warp-tls+  hs-source-dirs:      src+  default-language:    Haskell2010+  ghc-options:         -Wall
+ src/ExampleServer.hs view
@@ -0,0 +1,19 @@+{-# language DataKinds         #-}+{-# language OverloadedStrings #-}+{-# language TypeFamilies      #-}+module Main where++import           Mu.Adapter.ProtoBuf+import           Mu.GRpc.Server+import           Mu.Rpc.Examples+import           Mu.Schema++type instance AnnotatedSchema ProtoBufAnnotation QuickstartSchema+  = '[ 'AnnField "HelloRequest" "name" ('ProtoBufId 1)+     , 'AnnField "HelloResponse" "message" ('ProtoBufId 1)+     , 'AnnField "HiRequest" "number" ('ProtoBufId 1) ]++main :: IO ()+main = do+  putStrLn "running quickstart application"+  runGRpcApp 8080 quickstartServer
+ src/Mu/GRpc/Server.hs view
@@ -0,0 +1,316 @@+{-# language DataKinds             #-}+{-# language FlexibleContexts      #-}+{-# language FlexibleInstances     #-}+{-# language GADTs                 #-}+{-# language MultiParamTypeClasses #-}+{-# language PolyKinds             #-}+{-# language RankNTypes            #-}+{-# language ScopedTypeVariables   #-}+{-# language TypeApplications      #-}+{-# language TypeOperators         #-}+{-# language UndecidableInstances  #-}+{-|+Description : Execute a Mu 'Server' using gRPC as transport layer++This module allows you to server a Mu 'Server'+as a WAI 'Application' using gRPC as transport layer.++The simples way is to use 'runGRpcApp', all other+variants provide more control over the settings.+-}+module Mu.GRpc.Server+( -- * Run a 'Server' directly+  runGRpcApp, runGRpcAppTrans+, runGRpcAppSettings, Settings+, runGRpcAppTLS, TLSSettings+  -- * Convert a 'Server' into a WAI application+, gRpcApp+  -- * Raise errors as exceptions in IO+, raiseErrors, liftServerConduit+) where++import           Control.Concurrent.Async+import           Control.Concurrent.STM        (atomically)+import           Control.Concurrent.STM.TMVar+import           Control.Monad.Except+import           Data.ByteString               (ByteString)+import qualified Data.ByteString.Char8         as BS+import           Data.Conduit+import           Data.Conduit.TMChan+import           Data.Kind+import           Data.Proxy+import           Network.GRPC.HTTP2.Encoding   (gzip, uncompressed)+import           Network.GRPC.HTTP2.Proto3Wire+import           Network.GRPC.HTTP2.Types      (GRPCStatus (..), GRPCStatusCode (..))+import           Network.GRPC.Server.Handlers+import           Network.GRPC.Server.Wai       as Wai+import           Network.Wai                   (Application)+import           Network.Wai.Handler.Warp      (Port, Settings, run, runSettings)+import           Network.Wai.Handler.WarpTLS   (TLSSettings, runTLS)++import           Mu.Adapter.ProtoBuf.Via+import           Mu.Rpc+import           Mu.Schema+import           Mu.Server++-- | Run a Mu 'Server' on the given port.+runGRpcApp+  :: ( KnownName name, KnownName (FindPackageName anns)+     , GRpcMethodHandlers ServerErrorIO methods handlers )+  => Port+  -> ServerT Maybe ('Service name anns methods) ServerErrorIO handlers+  -> IO ()+runGRpcApp port = runGRpcAppTrans port id++-- | Run a Mu 'Server' on the given port.+runGRpcAppTrans+  :: ( KnownName name, KnownName (FindPackageName anns)+     , GRpcMethodHandlers m methods handlers )+  => Port+  -> (forall a. m a -> ServerErrorIO a)+  -> ServerT Maybe ('Service name anns methods) m handlers+  -> IO ()+runGRpcAppTrans port f svr = run port (gRpcAppTrans f svr)++-- | Run a Mu 'Server' using the given 'Settings'.+--+--   Go to 'Network.Wai.Handler.Warp' to declare 'Settings'.+runGRpcAppSettings+  :: ( KnownName name, KnownName (FindPackageName anns)+     , GRpcMethodHandlers m methods handlers )+  => Settings+  -> (forall a. m a -> ServerErrorIO a)+  -> ServerT Maybe ('Service name anns methods) m handlers+  -> IO ()+runGRpcAppSettings st f svr = runSettings st (gRpcAppTrans f svr)++-- | Run a Mu 'Server' using the given 'TLSSettings' and 'Settings'.+--+--   Go to 'Network.Wai.Handler.WarpTLS' to declare 'TLSSettings'+--   and to 'Network.Wai.Handler.Warp' to declare 'Settings'.+runGRpcAppTLS+  :: ( KnownName name, KnownName (FindPackageName anns)+     , GRpcMethodHandlers m methods handlers )+  => TLSSettings -> Settings+  -> (forall a. m a -> ServerErrorIO a)+  -> ServerT Maybe ('Service name anns methods) m handlers+  -> IO ()+runGRpcAppTLS tls st f svr = runTLS tls st (gRpcAppTrans f svr)++-- | Turn a Mu 'Server' into a WAI 'Application'.+--+--   These 'Application's can be later combined using,+--   for example, @wai-routes@, or you can add middleware+--   from @wai-extra@, among others.+gRpcApp+  :: ( KnownName name, KnownName (FindPackageName anns)+     , GRpcMethodHandlers ServerErrorIO methods handlers )+  => ServerT Maybe ('Service name anns methods) ServerErrorIO handlers+  -> Application+gRpcApp = gRpcAppTrans id++-- | Turn a Mu 'Server' into a WAI 'Application'.+--+--   These 'Application's can be later combined using,+--   for example, @wai-routes@, or you can add middleware+--   from @wai-extra@, among others.+gRpcAppTrans+  :: ( KnownName name, KnownName (FindPackageName anns)+     , GRpcMethodHandlers m methods handlers )+  => (forall a. m a -> ServerErrorIO a)+  -> ServerT Maybe ('Service name anns methods) m handlers+  -> Application+gRpcAppTrans f svr = Wai.grpcApp [uncompressed, gzip]+                                 (gRpcServiceHandlers f svr)++gRpcServiceHandlers+  :: forall name anns methods handlers m.+     (KnownName name, KnownName (FindPackageName anns), GRpcMethodHandlers m methods handlers)+  => (forall a. m a -> ServerErrorIO a)+  -> ServerT Maybe ('Service name anns methods) m handlers+  -> [ServiceHandler]+gRpcServiceHandlers f (Server svr) = gRpcMethodHandlers f packageName serviceName svr+  where packageName = BS.pack (nameVal (Proxy @(FindPackageName anns)))+        serviceName = BS.pack (nameVal (Proxy @name))++class GRpcMethodHandlers (m :: Type -> Type) (ms :: [Method mnm]) (hs :: [Type]) where+  gRpcMethodHandlers :: (forall a. m a -> ServerErrorIO a)+                     -> ByteString -> ByteString+                     -> HandlersT Maybe ms m hs -> [ServiceHandler]++instance GRpcMethodHandlers m '[] '[] where+  gRpcMethodHandlers _ _ _ H0 = []+instance (KnownName name, GRpcMethodHandler m args r h, GRpcMethodHandlers m rest hs)+         => GRpcMethodHandlers m ('Method name anns args r ': rest) (h ': hs) where+  gRpcMethodHandlers f p s (h :<|>: rest)+    = gRpcMethodHandler f (Proxy @args) (Proxy @r) (RPC p s methodName) h+      : gRpcMethodHandlers f p s rest+    where methodName = BS.pack (nameVal (Proxy @name))++class GRpcMethodHandler m args r h where+  gRpcMethodHandler :: (forall a. m a -> ServerErrorIO a)+                    -> Proxy args -> Proxy r -> RPC -> h -> ServiceHandler++-- | Turns a 'Conduit' working on 'ServerErrorIO'+--   into any other base monad which supports 'IO',+--   by raising any error as an exception.+--+--   This function is useful to interoperate with+--   libraries which generate 'Conduit's with other+--   base monads, such as @persistent@.+liftServerConduit+  :: MonadIO m+  => ConduitT a b ServerErrorIO r -> ConduitT a b m r+liftServerConduit = transPipe raiseErrors++-- | Raises errors from 'ServerErrorIO' as exceptions+--   in a monad which supports 'IO'.+--+--   This function is useful to interoperate with other+--   libraries which cannot handle the additional error+--   layer. In particular, with Conduit, as witnessed+--   by 'liftServerConduit'.+raiseErrors :: MonadIO m => ServerErrorIO a -> m a+raiseErrors h+  = liftIO $ do+      h' <- runExceptT h+      case h' of+        Right r -> return r+        Left (ServerError code msg)+          -> closeEarly $ GRPCStatus (serverErrorToGRpcError code)+                                     (BS.pack msg)+  where+    serverErrorToGRpcError :: ServerErrorCode -> GRPCStatusCode+    serverErrorToGRpcError Unknown         = UNKNOWN+    serverErrorToGRpcError Unavailable     = UNAVAILABLE+    serverErrorToGRpcError Unimplemented   = UNIMPLEMENTED+    serverErrorToGRpcError Unauthenticated = UNAUTHENTICATED+    serverErrorToGRpcError Internal        = INTERNAL+    serverErrorToGRpcError NotFound        = NOT_FOUND+    serverErrorToGRpcError Invalid         = INVALID_ARGUMENT++instance GRpcMethodHandler m '[ ] 'RetNothing (m ()) where+  gRpcMethodHandler f _ _ rpc h+    = unary @_ @() @() rpc (\_ _ -> raiseErrors (f h))++instance (ToProtoBufTypeRef rref r)+         => GRpcMethodHandler m '[ ] ('RetSingle rref) (m r) where+  gRpcMethodHandler f _ _ rpc h+    = unary @_ @() @(ViaToProtoBufTypeRef rref r)+            rpc (\_ _ -> ViaToProtoBufTypeRef <$> raiseErrors (f h))++instance (ToProtoBufTypeRef rref r, MonadIO m)+         => GRpcMethodHandler m '[ ] ('RetStream rref)+                              (ConduitT r Void m () -> m ()) where+  gRpcMethodHandler f _ _ rpc h+    = serverStream @_ @() @(ViaToProtoBufTypeRef rref r) rpc sstream+    where sstream :: req -> ()+                  -> IO ((), ServerStream (ViaToProtoBufTypeRef rref r) ())+          sstream _ _ = do+            -- Variable to connect input and output+            var <- newEmptyTMVarIO :: IO (TMVar (Maybe r))+            -- Start executing the handler+            promise <- async (raiseErrors $ ViaToProtoBufTypeRef <$> f (h (toTMVarConduit var)))+              -- Return the information+            let readNext _+                  = do nextOutput <- atomically $ takeTMVar var+                       case nextOutput of+                         Just o  -> return $ Just ((), ViaToProtoBufTypeRef o)+                         Nothing -> do cancel promise+                                       return Nothing+            return ((), ServerStream readNext)++instance (FromProtoBufTypeRef vref v)+         => GRpcMethodHandler m '[ 'ArgSingle vref ] 'RetNothing (v -> m ()) where+  gRpcMethodHandler f _ _ rpc h+    = unary @_ @(ViaFromProtoBufTypeRef vref v) @()+            rpc (\_ -> raiseErrors . f . h . unViaFromProtoBufTypeRef)++instance (FromProtoBufTypeRef vref v, ToProtoBufTypeRef rref r)+         => GRpcMethodHandler m '[ 'ArgSingle vref ] ('RetSingle rref) (v -> m r) where+  gRpcMethodHandler f _ _ rpc h+    = unary @_ @(ViaFromProtoBufTypeRef vref v) @(ViaToProtoBufTypeRef rref r)+            rpc (\_ -> (ViaToProtoBufTypeRef <$>) . raiseErrors . f . h . unViaFromProtoBufTypeRef)++instance (FromProtoBufTypeRef vref v, ToProtoBufTypeRef rref r, MonadIO m)+         => GRpcMethodHandler m '[ 'ArgStream vref ] ('RetSingle rref)+                              (ConduitT () v m () -> m r) where+  gRpcMethodHandler f _ _ rpc h+    = clientStream @_ @(ViaFromProtoBufTypeRef vref v) @(ViaToProtoBufTypeRef rref r)+                   rpc cstream+    where cstream :: req+                  -> IO ((), ClientStream (ViaFromProtoBufTypeRef vref v)+                        (ViaToProtoBufTypeRef rref r) ())+          cstream _ = do+            -- Create a new TMChan+            chan <- newTMChanIO :: IO (TMChan v)+            let producer = sourceTMChan @m chan+            -- Start executing the handler in another thread+            promise <- async (raiseErrors $ ViaToProtoBufTypeRef <$> f (h producer))+            -- Build the actual handler+            let cstreamHandler _ (ViaFromProtoBufTypeRef newInput)+                  = atomically (writeTMChan chan newInput)+                cstreamFinalizer _+                  = atomically (closeTMChan chan) >> wait promise+            -- Return the information+            return ((), ClientStream cstreamHandler cstreamFinalizer)++instance (FromProtoBufTypeRef vref v, ToProtoBufTypeRef rref r, MonadIO m)+         => GRpcMethodHandler m '[ 'ArgSingle vref ] ('RetStream rref)+                              (v -> ConduitT r Void m () -> m ()) where+  gRpcMethodHandler f _ _ rpc h+    = serverStream @_ @(ViaFromProtoBufTypeRef vref v) @(ViaToProtoBufTypeRef rref r)+                   rpc sstream+    where sstream :: req -> ViaFromProtoBufTypeRef vref v+                  -> IO ((), ServerStream (ViaToProtoBufTypeRef rref r) ())+          sstream _ (ViaFromProtoBufTypeRef v) = do+            -- Variable to connect input and output+            var <- newEmptyTMVarIO :: IO (TMVar (Maybe r))+            -- Start executing the handler+            promise <- async (raiseErrors $ ViaToProtoBufTypeRef <$> f (h v (toTMVarConduit var)))+              -- Return the information+            let readNext _+                  = do nextOutput <- atomically $ takeTMVar var+                       case nextOutput of+                         Just o  -> return $ Just ((), ViaToProtoBufTypeRef o)+                         Nothing -> do cancel promise+                                       return Nothing+            return ((), ServerStream readNext)++instance (FromProtoBufTypeRef vref v, ToProtoBufTypeRef rref r, MonadIO m)+         => GRpcMethodHandler m '[ 'ArgStream vref ] ('RetStream rref)+                              (ConduitT () v m () -> ConduitT r Void m () -> m ()) where+  gRpcMethodHandler f _ _ rpc h+    = generalStream @_ @(ViaFromProtoBufTypeRef vref v) @(ViaToProtoBufTypeRef rref r)+                    rpc bdstream+    where bdstream :: req -> IO ( (), IncomingStream (ViaFromProtoBufTypeRef vref v) ()+                                , (), OutgoingStream (ViaToProtoBufTypeRef rref r) () )+          bdstream _ = do+            -- Create a new TMChan and a new variable+            chan <- newTMChanIO :: IO (TMChan v)+            let producer = sourceTMChan @m chan+            var <- newEmptyTMVarIO :: IO (TMVar (Maybe r))+            -- Start executing the handler+            promise <- async (raiseErrors $ f $ h producer (toTMVarConduit var))+            -- Build the actual handler+            let cstreamHandler _ (ViaFromProtoBufTypeRef newInput)+                  = atomically (writeTMChan chan newInput)+                cstreamFinalizer _+                  = atomically (closeTMChan chan) >> wait promise+                readNext _+                  = do nextOutput <- atomically $ tryTakeTMVar var+                       case nextOutput of+                         Just (Just o) ->+                           return $ Just ((), ViaToProtoBufTypeRef o)+                         Just Nothing  -> do+                           cancel promise+                           return Nothing+                         Nothing -> -- no new elements to output+                           readNext ()+            return ((), IncomingStream cstreamHandler cstreamFinalizer, (), OutgoingStream readNext)++toTMVarConduit :: MonadIO m => TMVar (Maybe r) -> ConduitT r Void m ()+toTMVarConduit var = do+  x <- await+  liftIO $ atomically $ putTMVar var x+  toTMVarConduit var