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 +5/−0
- LICENSE +202/−0
- Setup.hs +2/−0
- mu-grpc-server.cabal +67/−0
- src/ExampleServer.hs +19/−0
- src/Mu/GRpc/Server.hs +316/−0
+ 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