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
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
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)
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 $
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'
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 ()
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
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))
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
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)
* 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.