( RunProto
, RunState(..)
, mkRunState
+ , P2PHandle(..)
, P2PConnection(..)
, ConnIdent(..)
, ClosableConnection(..)
tvar <- newTVarIO defaultProtocolVersion
return (mk tvar)
+data P2PHandle
+ = P2PHandle Handle
+ | P2PHandleTMVar (TMVar (Either L.ByteString Message))
+
data P2PConnection = P2PConnection
{ connRepo :: Maybe Repo
, connCheckAuth :: (AuthToken -> Bool)
- , connIhdl :: Handle
- , connOhdl :: Handle
+ , connIhdl :: P2PHandle
+ , connOhdl :: P2PHandle
, connIdent :: ConnIdent
}
stdioP2PConnection g = P2PConnection
{ connRepo = g
, connCheckAuth = const False
- , connIhdl = stdin
- , connOhdl = stdout
+ , connIhdl = P2PHandle stdin
+ , connOhdl = P2PHandle stdout
, connIdent = ConnIdent Nothing
}
return $ P2PConnection
{ connRepo = g
, connCheckAuth = const False
- , connIhdl = h
- , connOhdl = h
+ , connIhdl = P2PHandle h
+ , connOhdl = P2PHandle h
, connIdent = ConnIdent Nothing
}
closeConnection :: P2PConnection -> IO ()
closeConnection conn = do
- hClose (connIhdl conn)
- hClose (connOhdl conn)
+ closehandle (connIhdl conn)
+ closehandle (connOhdl conn)
+ where
+ closehandle (P2PHandle h) = hClose h
+ closehandle (P2PHandleTMVar _) = return ()
-- Serves the protocol on a unix socket.
--
go (Free (Local _)) = return $ Left $
ProtoFailureMessage "unexpected annex operation attempted"
+data P2PTMVarException = P2PTMVarException String
+ deriving (Show)
+
+instance Exception P2PTMVarException
+
-- Interpreter of the Net part of Proto.
--
-- An interpreter of Proto has to be provided, to handle the rest of Proto
runNet :: (MonadIO m, MonadMask m) => RunState -> P2PConnection -> RunProto m -> NetF (Proto a) -> m (Either ProtoFailure a)
runNet runst conn runner f = case f of
SendMessage m next -> do
- v <- liftIO $ tryNonAsync $ do
- let l = unwords (formatMessage m)
+ v <- liftIO $ do
debugMessage conn "P2P >" m
- hPutStrLn (connOhdl conn) l
- hFlush (connOhdl conn)
+ case connOhdl conn of
+ P2PHandle h -> tryNonAsync $ do
+ hPutStrLn h $ unwords (formatMessage m)
+ hFlush h
+ P2PHandleTMVar mv ->
+ ifM (atomically (tryPutTMVar mv (Right m)))
+ ( return $ Right ()
+ , return $ Left $ toException $
+ P2PTMVarException "TMVar left full"
+ )
case v of
Left e -> return $ Left $ ProtoFailureException e
Right () -> runner next
- ReceiveMessage next -> do
- v <- liftIO $ tryIOError $ getProtocolLine (connIhdl conn)
- case v of
- Left e -> return $ Left $ ProtoFailureIOError e
- Right Nothing -> return $ Left $
- ProtoFailureMessage "protocol error"
- Right (Just l) -> case parseMessage l of
- Just m -> do
- liftIO $ debugMessage conn "P2P <" m
- runner (next (Just m))
- Nothing -> runner (next Nothing)
- SendBytes len b p next -> do
- v <- liftIO $ tryNonAsync $ do
- ok <- sendExactly len b (connOhdl conn) p
- hFlush (connOhdl conn)
- return ok
- case v of
- Right True -> runner next
- Right False -> return $ Left $
- ProtoFailureMessage "short data write"
- Left e -> return $ Left $ ProtoFailureException e
- ReceiveBytes len p next -> do
- v <- liftIO $ tryNonAsync $ receiveExactly len (connIhdl conn) p
- case v of
- Left e -> return $ Left $ ProtoFailureException e
- Right b -> runner (next b)
+ ReceiveMessage next ->
+ let protoerr = return $ Left $
+ ProtoFailureMessage "protocol error"
+ gotmessage m = do
+ liftIO $ debugMessage conn "P2P <" m
+ runner (next (Just m))
+ in case connIhdl conn of
+ P2PHandle h -> do
+ v <- liftIO $ tryIOError $ getProtocolLine h
+ case v of
+ Left e -> return $ Left $ ProtoFailureIOError e
+ Right Nothing -> protoerr
+ Right (Just l) -> case parseMessage l of
+ Just m -> gotmessage m
+ Nothing -> runner (next Nothing)
+ P2PHandleTMVar mv ->
+ liftIO (atomically (takeTMVar mv)) >>= \case
+ Right m -> gotmessage m
+ Left _b -> protoerr
+ SendBytes len b p next ->
+ case connOhdl conn of
+ P2PHandle h -> do
+ v <- liftIO $ tryNonAsync $ do
+ ok <- sendExactly len b h p
+ hFlush h
+ return ok
+ case v of
+ Right True -> runner next
+ Right False -> return $ Left $
+ ProtoFailureMessage "short data write"
+ Left e -> return $ Left $ ProtoFailureException e
+ P2PHandleTMVar mv -> do
+ liftIO $ atomically $ putTMVar mv (Left b)
+ runner next
+ ReceiveBytes len p next ->
+ case connIhdl conn of
+ P2PHandle h -> do
+ v <- liftIO $ tryNonAsync $ receiveExactly len h p
+ case v of
+ Right b -> runner (next b)
+ Left e -> return $ Left $
+ ProtoFailureException e
+ P2PHandleTMVar mv ->
+ liftIO (atomically (takeTMVar mv)) >>= \case
+ Left b -> runner (next b)
+ Right _m -> return $ Left $
+ ProtoFailureMessage "protocol error"
CheckAuthToken _u t next -> do
let authed = connCheckAuth conn t
runner (next authed)