module P2P.Proxy where
import Annex.Common
+import qualified Annex
import P2P.Protocol
import P2P.IO
import Utility.Metered
import Git.FilePath
+import Types.Concurrency
+import Annex.Concurrent
import Data.Either
import Control.Concurrent.STM
+import Control.Concurrent.Async
+import qualified Control.Concurrent.MSem as MSem
import qualified Data.ByteString.Lazy as L
+import GHC.Conc
type ProtoCloser = Annex ()
-> ClientSide
-> UUID
-> ProxySelector
+ -> ConcurrencyConfig
-> ProtocolVersion
-- ^ Protocol version being spoken between the proxy and the
-- client. When there are multiple remotes, some may speak an
-- ^ non-VERSION message that was received from the client when
-- negotiating protocol version, and has not been responded to yet
-> ProtoErrorHandled r
-proxy proxydone proxymethods servermode (ClientSide clientrunst clientconn) remoteuuid proxyselector (ProtocolVersion protocolversion) othermessage protoerrhandler = do
+proxy proxydone proxymethods servermode (ClientSide clientrunst clientconn) remoteuuid proxyselector concurrencyconfig (ProtocolVersion protocolversion) othermessage protoerrhandler = do
case othermessage of
Nothing -> protoerrhandler proxynextclientmessage $
client $ net $ sendMessage $ VERSION $ ProtocolVersion protocolversion
protoerrhandler proxynextclientmessage $
client $ net $ sendMessage FAILURE
handleREMOVE remotesides k message = do
- v <- forM remotesides $ \r ->
+ v <- forMC concurrencyconfig remotesides $ \r ->
runRemoteSideOrSkipFailed r $ do
net $ sendMessage message
net receiveMessage >>= return . \case
let alreadyhave = \case
Right (Left _) -> True
_ -> False
- l <- forM remotesides initiate
+ l <- forMC concurrencyconfig remotesides initiate
if all alreadyhave l
then if protocolversion < 2
then protoerrhandler proxynextclientmessage $
let totallen = datalen + minoffset
-- Tell each remote how much data to expect, depending
-- on the remote's offset.
- rs <- forM remotes $ \remote@(remoteside, remoteoffset) ->
+ rs <- forMC concurrencyconfig remotes $ \remote@(remoteside, remoteoffset) ->
runRemoteSideOrSkipFailed remoteside $ do
net $ sendMessage $ DATA $ Len $
totallen - remoteoffset
let (chunk, b') = L.splitAt chunksize b
let chunklen = fromIntegral (L.length chunk)
let !n' = n + chunklen
- rs' <- forM rs $ \r@(remoteside, remoteoffset) ->
+ rs' <- forMC concurrencyconfig rs $ \r@(remoteside, remoteoffset) ->
if n >= remoteoffset
then runRemoteSideOrSkipFailed remoteside $ do
net $ sendBytes (Len chunklen) chunk nullMeterUpdate
net receiveMessage
where
finish a = do
- storeduuids <- forM rs $ \r ->
+ storeduuids <- forMC concurrencyconfig rs $ \r ->
runRemoteSideOrSkipFailed r a >>= \case
Just (Just resp) ->
relayPUTRecord k r resp
<$> fromRepo (fromTopFilePath (asTopFilePath f))
getassociatedfile (ProtoAssociatedFile (AssociatedFile Nothing)) =
return $ AssociatedFile Nothing
+
+data ConcurrencyConfig = ConcurrencyConfig Int (MSem.MSem Int)
+
+noConcurrencyConfig :: Annex ConcurrencyConfig
+noConcurrencyConfig = liftIO $ ConcurrencyConfig 1 <$> MSem.new 1
+
+getConcurrencyConfig :: Annex ConcurrencyConfig
+getConcurrencyConfig = (annexJobs <$> Annex.getGitConfig) >>= \case
+ NonConcurrent -> noConcurrencyConfig
+ Concurrent n -> go n
+ ConcurrentPerCpu -> go =<< liftIO getNumProcessors
+ where
+ go n = do
+ c <- liftIO getNumCapabilities
+ when (n > c) $
+ liftIO $ setNumCapabilities n
+ setConcurrency (ConcurrencyGitConfig (Concurrent n))
+ msem <- liftIO $ MSem.new n
+ return (ConcurrencyConfig n msem)
+
+forMC :: ConcurrencyConfig -> [a] -> (a -> Annex b) -> Annex [b]
+forMC _ (x:[]) a = do
+ r <- a x
+ return [r]
+forMC (ConcurrencyConfig n msem) xs a
+ | n < 2 = forM xs a
+ | otherwise = do
+ runners <- forM xs $ \x ->
+ forkState $ bracketIO
+ (MSem.wait msem)
+ (const $ MSem.signal msem)
+ (const $ a x)
+ mapM id =<< liftIO (forConcurrently runners id)
+
Started proxying for node2
Started proxying for node3
+Operations that affect multiple nodes of a cluster can often be sped up by
+configuring annex.jobs in the repository that will serve the cluster to
+clients. In the example above, the nodes are all disk bound, so operating
+on more than one at a time will likely be faster.
+
+ $ git config annex.jobs cpus
+
## preferred content of clusters
The preferred content of the cluster can be configured. This tells
For example:
- git-annex wanted bigserver-mycluster standard
- git-annex group bigserver-mycluster archive
+ $ git-annex wanted bigserver-mycluster standard
+ $ git-annex group bigserver-mycluster archive
By default, when a file is uploaded to a cluster, it is stored on every node of
the cluster. To control which nodes to store to, the [[preferred_content]] of