one jobid per thread
authorJoey Hess <joeyh@joeyh.name>
Fri, 14 Aug 2020 18:24:46 +0000 (14:24 -0400)
committerJoey Hess <joeyh@joeyh.name>
Fri, 14 Aug 2020 18:24:46 +0000 (14:24 -0400)
And, relay ERROR on to all listening threads.

Remote/External/AsyncExtension.hs
doc/design/external_special_remote_protocol/async_appendix.mdwn

index 51f339010a5af3969725e955983ed26fa613aba8..78522e5fe8bd42d2eee451353e6d96978e8eb84b 100644 (file)
@@ -28,15 +28,20 @@ runRelayToExternalAsync external st = do
        jidmap <- newTVarIO M.empty
        sendq <- newSendQueue
        nextjid <- newTVarIO (JobId 1)
-       void $ async $ sendloop st nextjid jidmap sendq
+       void $ async $ sendloop st sendq
        void $ async $ receiveloop external st jidmap sendq
        return $ ExternalAsyncRelay $ do
-               jidv <- newTVarIO Nothing
                receiveq <- newReceiveQueue
+               jid <- atomically $ do
+                       jid@(JobId n) <- readTVar nextjid
+                       let !jid' = JobId (succ n)
+                       writeTVar nextjid jid'
+                       modifyTVar' jidmap $ M.insert jid receiveq
+                       return jid
                return $ ExternalState
                        { externalSend = \msg -> 
                                atomically $ writeTBMChan sendq
-                                       (toAsyncWrapped msg, (jidv, receiveq))
+                                       (toAsyncWrapped msg, jid)
                        , externalReceive = atomically (readTBMChan receiveq)
                        -- This shuts down the whole relay.
                        , externalShutdown = shutdown external st sendq
@@ -52,13 +57,9 @@ runRelayToExternalAsync external st = do
 
 type ReceiveQueue = TBMChan String
 
-type SendQueue = TBMChan (AsyncWrapped, Conn)
+type SendQueue = TBMChan (AsyncWrapped, JobId)
 
-type Conn = (TVar (Maybe JobId), ReceiveQueue)
-
-type JidMap = TVar (M.Map JobId Conn)
-
-type NextJid = TVar JobId
+type JidMap = TVar (M.Map JobId ReceiveQueue)
 
 newReceiveQueue :: IO ReceiveQueue
 newReceiveQueue = newTBMChanIO 10
@@ -71,11 +72,18 @@ receiveloop external st jidmap sendq = externalReceive st >>= \case
        Just l -> case parseMessage l :: Maybe AsyncMessage of
                Just (AsyncMessage jid msg) ->
                        M.lookup jid <$> readTVarIO jidmap >>= \case
-                               Just (_jidv, c) -> do
+                               Just c -> do
                                        atomically $ writeTBMChan c msg
                                        receiveloop external st jidmap sendq
                                Nothing -> protoerr "unknown job number"
-               _ -> protoerr "unexpected non-async message"
+               Nothing -> case parseMessage l :: Maybe ExceptionalMessage of
+                       Just msg -> do
+                               -- ERROR is relayed to all listeners
+                               m <- readTVarIO jidmap
+                               forM (M.elems m) $ \c ->
+                                       atomically  $ writeTBMChan c l
+                               receiveloop external st jidmap sendq
+                       Nothing -> protoerr "unexpected non-async message"
        Nothing -> closeandshutdown
   where
        protoerr s = do
@@ -85,30 +93,21 @@ receiveloop external st jidmap sendq = externalReceive st >>= \case
        closeandshutdown = do
                shutdown external st sendq True
                m <- atomically $ readTVar jidmap
-               forM_ (M.elems m) (atomically . closeTBMChan . snd)
+               forM_ (M.elems m) (atomically . closeTBMChan)
 
-sendloop :: ExternalState -> NextJid -> JidMap -> SendQueue -> IO ()
-sendloop st nextjid jidmap sendq = atomically (readTBMChan sendq) >>= \case
-       Just (wrappedmsg, conn@(jidv, _)) -> do
+sendloop :: ExternalState -> SendQueue -> IO ()
+sendloop st sendq = atomically (readTBMChan sendq) >>= \case
+       Just (wrappedmsg, jid) -> do
                case wrappedmsg of
-                       AsyncWrappedRequest msg -> do
-                               jid <- atomically $ do
-                                       jid@(JobId n) <- readTVar nextjid
-                                       let !jid' = JobId (succ n)
-                                       writeTVar nextjid jid'
-                                       writeTVar jidv (Just jid)
-                                       modifyTVar' jidmap $ M.insert jid conn
-                                       return jid
-                               externalSend st $ wrapjid msg jid
                        AsyncWrappedRemoteResponse msg ->
-                               readTVarIO jidv >>= \case
-                                       Just jid -> externalSend st $ wrapjid msg jid
-                                       Nothing -> error "failed to find jid"
-                       AsyncWrappedExceptionalMessage msg ->
+                               externalSend st $ wrapjid msg jid
+                       AsyncWrappedRequest msg ->
+                               externalSend st $ wrapjid msg jid
+                       AsyncWrappedExceptionalMessage msg -> 
                                externalSend st msg
                        AsyncWrappedAsyncMessage msg ->
                                externalSend st msg
-               sendloop st nextjid jidmap sendq
+               sendloop st sendq
        Nothing -> return ()
   where
        wrapjid msg jid = AsyncMessage jid $ unwords $ Proto.formatMessage msg
index 59381e6ebfc2191f719fce95acfe90678979d31b..50eca9fad5dcd40c8cfb2dce749b9afbaf1dbf12 100644 (file)
@@ -29,7 +29,7 @@ that includes `ASYNC`, and the external special remote responding in kind.
        EXTENSIONS ASYNC
 
 From this point forward, every message in the protocol is tagged with a job
-number, by prefixing it with "J n". 
+number, by prefixing it with "J n".
 
 As usual, the first message git-annex sends is generally PREPARE:
 
@@ -43,31 +43,31 @@ included in the reply:
 Suppose git-annex wants to make some transfers. It can request several
 at the same time, using different job numbers:
 
-       J 2 TRANSFER RETRIEVE Key1 file1
-       J 3 TRANSFER RETRIEVE Key2 file2
+       J 1 TRANSFER RETRIEVE Key1 file1
+       J 2 TRANSFER RETRIEVE Key2 file2
 
 The special remote can now perform both transfers at the same time.
 If it sends PROGRESS messages for these transfers, they have to be tagged
 with the job number:
        
-       J 2 PROGRESS 10
-       J 3 PROGRESS 500
-       J 2 PROGRESS 20
+       J 1 PROGRESS 10
+       J 2 PROGRESS 500
+       J 1 PROGRESS 20
 
 The special remote can also send messages that query git-annex for some
 information. These messages and the reply will also be tagged with a job
 number.
 
-       J 2 GETCONFIG url
-       J 4 RETRIEVE Key3 file3
-       J 2 VALUE http://example.com/
+       J 1 GETCONFIG url
+       J 3 RETRIEVE Key3 file3
+       J 1 VALUE http://example.com/
 
 One transfers are done, the special remote sends `TRANSFER-SUCCESS` tagged
 with the job number.
 
-       J 3 TRANSFER-SUCCESS RETRIEVE Key2
-       J 2 PROGRESS 100
-       J 2 TRANSFER-SUCCESS RETRIEVE Key1
+       J 2 TRANSFER-SUCCESS RETRIEVE Key2
+       J 1 PROGRESS 100
+       J 1 TRANSFER-SUCCESS RETRIEVE Key1
 
 Lots of different jobs can be requested at the same time.
 
@@ -78,9 +78,6 @@ Lots of different jobs can be requested at the same time.
        J 6 REMOVE-SUCCESS Key3
        J 5 CHECKPRESENT-FAILURE Key4
 
-A job always starts with a request by git-annex, and once the special
-remote sends a reply -- or replies -- to that request, the job is done.
-
 An example of sending multiple replies to a request is `LISTCONFIGS`, eg:
 
        J 7 LISTCONFIGS
@@ -88,10 +85,15 @@ An example of sending multiple replies to a request is `LISTCONFIGS`, eg:
        J 7 CONFIG bar other config
        J 7 CONFIGEND
 
-Job numbers are not reused within a given run of a special remote,
-but once git-annex has seen the last message it expects for a job,
-sending other messages tagged with that job number will be rejected
-as a protocol error.
+## notes
+
+There will generally be one job for each thread that git-annex runs
+concurrently, so around the same number as the -J value, although in some
+cases git-annex does more concurrent operations than the -J value.
+
+`PREPARE` is sent only once per run of a special remote
+program, and despite being tagged with a job number, it should prepare the
+special remote to run that and any other jobs.
 
-To avoid overflow, job numbers should be treated as at least 64 bit
-values (or as strings) by the special remote.
+`ERROR` should not be tagged with a job number if either git-annex
+or the special remote needs to send it.