]> dgit.raspbian.org Git - git-annex.git/commitdiff
fix serveGet hang
authorJoey Hess <joeyh@joeyh.name>
Thu, 11 Jul 2024 11:46:52 +0000 (07:46 -0400)
committerJoey Hess <joeyh@joeyh.name>
Thu, 11 Jul 2024 11:46:52 +0000 (07:46 -0400)
This came down to SendBytes waiting on the waitv. Nothing ever filled
it.

Only Annex.Proxy needs the waitv, and it handles filling it. So make it
optional.

Annex/Proxy.hs
P2P/Annex.hs
P2P/Http.hs
P2P/Http/State.hs
P2P/IO.hs
doc/todo/git-annex_proxies.mdwn

index 4563579ef2544b4c876f85346efaf91c065ccb0b..ffa732a08d6798857951e852db95775a79c40457 100644 (file)
@@ -59,8 +59,8 @@ proxySpecialRemoteSide clientmaxversion r = mkRemoteSide r $ do
        let remoteconn = P2PConnection
                { connRepo = Nothing
                , connCheckAuth = const False
-               , connIhdl = P2PHandleTMVar ihdl iwaitv
-               , connOhdl = P2PHandleTMVar ohdl owaitv
+               , connIhdl = P2PHandleTMVar ihdl (Just iwaitv)
+               , connOhdl = P2PHandleTMVar ohdl (Just owaitv)
                , connIdent = ConnIdent (Just (Remote.name r))
                }
        let closeremoteconn = do
index d922b28e5035369af426179f3d4a50a24b9f83e8..59428db8c9430cf83bceb6c172aed3752f0f616d 100644 (file)
@@ -107,23 +107,16 @@ runLocal runst runner a = case a of
                                ProtoFailureMessage "Transfer failed"
                        let consumer' b ti = do
                                validator <- consumer b
-                               liftIO $ print "got validator"
                                indicatetransferred ti
-                               liftIO $ print "indicatetransferred ti done"
                                return validator
                        runner getb >>= \case
                                Left e -> giveup $ describeProtoFailure e
                                Right b -> checktransfer (\ti -> Right <$> consumer' b ti) fallback >>= \case
                                        Left e -> return (Left e)
-                                       Right validator -> do
-                                               liftIO $ print "running validity check"
+                                       Right validator ->
                                                runner validitycheck >>= \case
-                                                       Right v -> do
-                                                               liftIO $ print ("calling validator 1", v)
-                                                               Right <$> validator v
-                                                       _ -> do
-                                                               liftIO $ print "calling validator nothing"
-                                                               Right <$> validator Nothing
+                                                       Right v -> Right <$> validator v
+                                                       _ -> Right <$> validator Nothing
                case v of
                        Left e -> return $ Left $ ProtoFailureException e
                        Right (Left e) -> return $ Left e
index 7e5092fdb4a81653228ae77612b8da515c7af882..c0728fde11f0b9a21ee8474dca725be314e40698 100644 (file)
@@ -156,31 +156,24 @@ serveGet st apiver (B64Key k) cu su bypass baf startat sec auth = do
        endv <- liftIO newEmptyTMVarIO
        validityv <- liftIO newEmptyTMVarIO
        aid <- liftIO $ async $ inAnnexWorker st $ do
-               let consumer bs = do
-                       liftIO $ atomically $ putTMVar bsv bs
-                       liftIO $ print "consumer waiting for endv"
+               let storer _offset len = sendContentWith $ \bs -> do
+                       liftIO $ atomically $ putTMVar bsv (len, bs)
                        liftIO $ atomically $ takeTMVar endv
-                       liftIO $ print "consumer took endv"
                        return $ \v -> do
-                               liftIO $ print "consumer put validityv"
-                               liftIO $ atomically $
-                                       putTMVar validityv v
+                               liftIO $ atomically $ putTMVar validityv v
                                return True
-               let storer _offset _len getdata checkvalidity =
-                       sendContentWith consumer getdata checkvalidity
                enteringStage (TransferStage Upload) $
                        runFullProto runst conn $
                                void $ receiveContent Nothing nullMeterUpdate
                                        sizer storer getreq
-       bs <- liftIO $ atomically $ takeTMVar bsv
+       (Len len, bs) <- liftIO $ atomically $ takeTMVar bsv
        bv <- liftIO $ newMVar (L.toChunks bs)
        let streamer = S.SourceT $ \s -> s =<< return 
                (stream (releaseconn, bv, endv, validityv, aid))
-       return $ addHeader 111111 streamer
+       return $ addHeader len streamer
   where
        stream (releaseconn, bv, endv, validityv, aid) =
-               S.fromActionStep B.null $ do
-                       print "chunk"
+               S.fromActionStep B.null $
                        modifyMVar bv $ nextchunk $
                                cleanup (releaseconn, endv, validityv, aid)
        
@@ -194,11 +187,8 @@ serveGet st apiver (B64Key k) cu su bypass baf startat sec auth = do
        cleanup (releaseconn, endv, validityv, aid) =
                ifM (atomically $ isEmptyTMVar endv)
                        ( do
-                               print "at end"
                                atomically $ putTMVar endv ()
-                               print "signaled end"
                                validity <- atomically $ takeTMVar validityv
-                               print ("got validity", validity)
                                wait aid >>= \case
                                        Left ex -> throwM ex
                                        Right (Left err) -> error $
@@ -263,6 +253,7 @@ gatherbytestring x = do
        go (S.Effect ms) = do
                ms >>= go
        go (S.Yield v s) = do
+               liftIO $ print ("chunk", B.length v)
                LI.Chunk v <$> unsafeInterleaveIO (go s)
 
 clientGet'
index d3e00dfa216c46d9458833226b58489c6b93d3df..8e9cae3fa624bcaa5bfe3ff47a1b3a17139484cf 100644 (file)
@@ -177,35 +177,38 @@ withLocalP2PConnections a = do
                        else do
                                hdl1 <- liftIO newEmptyTMVarIO
                                hdl2 <- liftIO newEmptyTMVarIO
-                               waitv1 <- liftIO newEmptyTMVarIO
-                               waitv2 <- liftIO newEmptyTMVarIO
-                               let h1 = P2PHandleTMVar hdl1 waitv1
-                               let h2 = P2PHandleTMVar hdl2 waitv2 
+                               let h1 = P2PHandleTMVar hdl1 Nothing
+                               let h2 = P2PHandleTMVar hdl2 Nothing
                                let serverconn = P2PConnection Nothing
                                        (const True) h1 h2
                                        (ConnIdent (Just "http server"))
                                let clientconn = P2PConnection Nothing
                                        (const True) h2 h1
                                        (ConnIdent (Just "http client"))
-                               runst <- liftIO $ mkrunst connparams
+                               clientrunst <- liftIO $ mkclientrunst connparams
+                               serverrunst <- liftIO $ mkserverrunst connparams
                                let server = P2P.serveOneCommandAuthed
                                        (connectionServerMode connparams)
                                        (connectionServerUUID connparams)
                                let protorunner = void $
-                                       runFullProto runst serverconn server
+                                       runFullProto serverrunst serverconn server
                                asyncworker <- liftIO . async 
                                        =<< forkState protorunner
                                let releaseconn = atomically $ putTMVar relv $
                                        join (liftIO (wait asyncworker))
-                               return $ Right (runst, clientconn, releaseconn)
+                               return $ Right (clientrunst, clientconn, releaseconn)
                liftIO $ atomically $ putTMVar respvar resp
 
-       mkrunst connparams = do
+       mkserverrunst connparams = do
                prototvar <- newTVarIO $ connectionProtocolVersion connparams
                mkRunState $ const $ Serving 
                        (connectionClientUUID connparams)
                        Nothing
                        prototvar
+       
+       mkclientrunst connparams = do
+               prototvar <- newTVarIO $ connectionProtocolVersion connparams
+               mkRunState $ const $ Client prototvar
 
 data Locker = Locker
        { lockerThread :: Async ()
index cd725c6e9e4b91489dbf36801779361a1dfc384a..66a4c08fea63f22c6aaf7c9df435bc0a13569fb0 100644 (file)
--- a/P2P/IO.hs
+++ b/P2P/IO.hs
@@ -79,7 +79,7 @@ mkRunState mk = do
 
 data P2PHandle
        = P2PHandle Handle
-       | P2PHandleTMVar (TMVar (Either L.ByteString Message)) (TMVar ())
+       | P2PHandleTMVar (TMVar (Either L.ByteString Message)) (Maybe (TMVar ()))
 
 data P2PConnection = P2PConnection
        { connRepo :: Maybe Repo
@@ -217,7 +217,7 @@ runNet runst conn runner f = case f of
                        Right () -> runner next
        ReceiveMessage next ->
                let protoerr = return $ Left $
-                       ProtoFailureMessage "protocol error 1"
+                       ProtoFailureMessage "protocol error"
                    gotmessage m = do
                        liftIO $ debugMessage conn "P2P <" m
                        runner (next (Just m))
@@ -246,11 +246,14 @@ runNet runst conn runner f = case f of
                                        Right False -> return $ Left $
                                                ProtoFailureMessage "short data write"
                                        Left e -> return $ Left $ ProtoFailureException e
-                       P2PHandleTMVar mv waitv -> do
+                       P2PHandleTMVar mv mwaitv -> do
                                liftIO $ atomically $ putTMVar mv (Left b)
-                               -- Wait for the whole bytestring to be
-                               -- processed. Necessary due to lazyiness.
-                               liftIO $ atomically $ takeTMVar waitv
+                               case mwaitv of
+                                       -- Wait for the whole bytestring to
+                                       -- be processed.
+                                       Just waitv -> liftIO $ atomically $
+                                               takeTMVar waitv
+                                       Nothing -> return ()
                                runner next
        ReceiveBytes len p next ->
                case connIhdl conn of
@@ -264,7 +267,7 @@ runNet runst conn runner f = case f of
                                liftIO (atomically (takeTMVar mv)) >>= \case
                                        Left b -> runner (next b)
                                        Right _ -> return $ Left $
-                                               ProtoFailureMessage "protocol error 2"
+                                               ProtoFailureMessage "protocol error"
        CheckAuthToken _u t next -> do
                let authed = connCheckAuth conn t
                runner (next authed)
index 1cea7d8b2890a7732e76b5e8979f215dd184a3e8..c1d16bb339872c72d1a84b433a87461d3de15e40 100644 (file)
@@ -31,8 +31,7 @@ Planned schedule of work:
 * http server and client are working, remaining 
   server API endpoints need wiring up and testing. 
 
-* serveGet works as proof of concept, but is very buggy.
-  See commit 1e0f92a5a1ccf7ff4c51c67c27a826709a99301b
+* serveGet needs to handle invalidation
 
 * I have a file `servant.hs` in the httpproto branch that works through some
   of the bytestring streaming issues.