split Changes and lifted

This commit is contained in:
Joey Hess 2012-10-29 19:30:23 -04:00
parent 39a3adf434
commit d2294f0dfa
5 changed files with 82 additions and 67 deletions

View file

@ -7,73 +7,33 @@
module Assistant.Changes where module Assistant.Changes where
import Common.Annex import Assistant.Common
import Types.KeySource import Assistant.Types.Changes
import Utility.TSet import Utility.TSet
import Data.Time.Clock import Data.Time.Clock
data ChangeType = AddChange | LinkChange | RmChange | RmDirChange
deriving (Show, Eq)
type ChangeChan = TSet Change
data Change
= Change
{ changeTime :: UTCTime
, changeFile :: FilePath
, changeType :: ChangeType
}
| PendingAddChange
{ changeTime ::UTCTime
, changeFile :: FilePath
}
| InProcessAddChange
{ changeTime ::UTCTime
, keySource :: KeySource
}
deriving (Show)
newChangeChan :: IO ChangeChan
newChangeChan = newTSet
{- Handlers call this when they made a change that needs to get committed. -} {- Handlers call this when they made a change that needs to get committed. -}
madeChange :: FilePath -> ChangeType -> IO (Maybe Change) madeChange :: FilePath -> ChangeType -> Assistant (Maybe Change)
madeChange f t = Just <$> (Change <$> getCurrentTime <*> pure f <*> pure t) madeChange f t = Just <$> (Change <$> liftIO getCurrentTime <*> pure f <*> pure t)
noChange :: IO (Maybe Change) noChange :: Assistant (Maybe Change)
noChange = return Nothing noChange = return Nothing
{- Indicates an add needs to be done, but has not started yet. -} {- Indicates an add needs to be done, but has not started yet. -}
pendingAddChange :: FilePath -> IO (Maybe Change) pendingAddChange :: FilePath -> Assistant (Maybe Change)
pendingAddChange f = Just <$> (PendingAddChange <$> getCurrentTime <*> pure f) pendingAddChange f = Just <$> (PendingAddChange <$> liftIO getCurrentTime <*> pure f)
isPendingAddChange :: Change -> Bool
isPendingAddChange (PendingAddChange {}) = True
isPendingAddChange _ = False
isInProcessAddChange :: Change -> Bool
isInProcessAddChange (InProcessAddChange {}) = True
isInProcessAddChange _ = False
finishedChange :: Change -> Change
finishedChange c@(InProcessAddChange { keySource = ks }) = Change
{ changeTime = changeTime c
, changeFile = keyFilename ks
, changeType = AddChange
}
finishedChange c = c
{- Gets all unhandled changes. {- Gets all unhandled changes.
- Blocks until at least one change is made. -} - Blocks until at least one change is made. -}
getChanges :: ChangeChan -> IO [Change] getChanges :: Assistant [Change]
getChanges = getTSet getChanges = getTSet <<~ changeChan
{- Puts unhandled changes back into the channel. {- Puts unhandled changes back into the channel.
- Note: Original order is not preserved. -} - Note: Original order is not preserved. -}
refillChanges :: ChangeChan -> [Change] -> IO () refillChanges :: [Change] -> Assistant ()
refillChanges = putTSet refillChanges cs = flip putTSet cs <<~ changeChan
{- Records a change in the channel. -} {- Records a change in the channel. -}
recordChange :: ChangeChan -> Change -> IO () recordChange :: Change -> Assistant ()
recordChange = putTSet1 recordChange c = flip putTSet1 c <<~ changeChan

View file

@ -34,7 +34,7 @@ import Assistant.TransferSlots
import Assistant.Types.Pushes import Assistant.Types.Pushes
import Assistant.Types.BranchChange import Assistant.Types.BranchChange
import Assistant.Commits import Assistant.Commits
import Assistant.Changes import Assistant.Types.Changes
newtype Assistant a = Assistant { mkAssistant :: ReaderT AssistantData IO a } newtype Assistant a = Assistant { mkAssistant :: ReaderT AssistantData IO a }
deriving ( deriving (

View file

@ -11,6 +11,7 @@ module Assistant.Threads.Committer where
import Assistant.Common import Assistant.Common
import Assistant.Changes import Assistant.Changes
import Assistant.Types.Changes
import Assistant.Commits import Assistant.Commits
import Assistant.Alert import Assistant.Alert
import Assistant.Threads.Watcher import Assistant.Threads.Watcher
@ -45,7 +46,7 @@ commitThread = NamedThread "Committer" $ do
-- We already waited one second as a simple rate limiter. -- We already waited one second as a simple rate limiter.
-- Next, wait until at least one change is available for -- Next, wait until at least one change is available for
-- processing. -- processing.
changes <- getChanges <<~ changeChan changes <- getChanges
-- Now see if now's a good time to commit. -- Now see if now's a good time to commit.
time <- liftIO getCurrentTime time <- liftIO getCurrentTime
if shouldCommit time changes if shouldCommit time changes
@ -67,7 +68,7 @@ commitThread = NamedThread "Committer" $ do
refill [] = noop refill [] = noop
refill cs = do refill cs = do
debug ["delaying commit of", show (length cs), "changes"] debug ["delaying commit of", show (length cs), "changes"]
flip refillChanges cs <<~ changeChan refillChanges cs
commitStaged :: Annex Bool commitStaged :: Annex Bool
commitStaged = do commitStaged = do
@ -148,15 +149,14 @@ handleAdds delayadd cs = returnWhen (null incomplete) $ do
(postponed, toadd) <- partitionEithers <$> safeToAdd delayadd pending' inprocess (postponed, toadd) <- partitionEithers <$> safeToAdd delayadd pending' inprocess
unless (null postponed) $ unless (null postponed) $
flip refillChanges postponed <<~ changeChan refillChanges postponed
returnWhen (null toadd) $ do returnWhen (null toadd) $ do
added <- catMaybes <$> forM toadd add added <- catMaybes <$> forM toadd add
if DirWatcher.eventsCoalesce || null added if DirWatcher.eventsCoalesce || null added
then return $ added ++ otherchanges then return $ added ++ otherchanges
else do else do
r <- handleAdds delayadd r <- handleAdds delayadd =<< getChanges
=<< getChanges <<~ changeChan
return $ r ++ added ++ otherchanges return $ r ++ added ++ otherchanges
where where
(incomplete, otherchanges) = partition (\c -> isPendingAddChange c || isInProcessAddChange c) cs (incomplete, otherchanges) = partition (\c -> isPendingAddChange c || isInProcessAddChange c) cs

View file

@ -17,6 +17,7 @@ module Assistant.Threads.Watcher (
import Assistant.Common import Assistant.Common
import Assistant.DaemonStatus import Assistant.DaemonStatus
import Assistant.Changes import Assistant.Changes
import Assistant.Types.Changes
import Assistant.TransferQueue import Assistant.TransferQueue
import Assistant.Alert import Assistant.Alert
import Assistant.Drop import Assistant.Drop
@ -114,12 +115,12 @@ runHandler handler file filestatus = void $ do
-- Just in case the commit thread is not -- Just in case the commit thread is not
-- flushing the queue fast enough. -- flushing the queue fast enough.
liftAnnex $ Annex.Queue.flushWhenFull liftAnnex $ Annex.Queue.flushWhenFull
flip recordChange change <<~ changeChan recordChange change
onAdd :: Handler onAdd :: Handler
onAdd file filestatus onAdd file filestatus
| maybe False isRegularFile filestatus = liftIO $ pendingAddChange file | maybe False isRegularFile filestatus = pendingAddChange file
| otherwise = liftIO $ noChange | otherwise = noChange
{- A symlink might be an arbitrary symlink, which is just added. {- A symlink might be an arbitrary symlink, which is just added.
- Or, if it is a git-annex symlink, ensure it points to the content - Or, if it is a git-annex symlink, ensure it points to the content
@ -160,7 +161,7 @@ onAddSymlink file filestatus = go =<< liftAnnex (Backend.lookupFile file)
| scanComplete daemonstatus = addlink link | scanComplete daemonstatus = addlink link
| otherwise = case filestatus of | otherwise = case filestatus of
Just s Just s
| not (afterLastDaemonRun (statusChangeTime s) daemonstatus) -> liftIO noChange | not (afterLastDaemonRun (statusChangeTime s) daemonstatus) -> noChange
_ -> addlink link _ -> addlink link
{- For speed, tries to reuse the existing blob for symlink target. -} {- For speed, tries to reuse the existing blob for symlink target. -}
@ -176,7 +177,7 @@ onAddSymlink file filestatus = go =<< liftAnnex (Backend.lookupFile file)
sha <- inRepo $ sha <- inRepo $
Git.HashObject.hashObject BlobObject link Git.HashObject.hashObject BlobObject link
stageSymlink file sha stageSymlink file sha
liftIO $ madeChange file LinkChange madeChange file LinkChange
{- When a new link appears, or a link is changed, after the startup {- When a new link appears, or a link is changed, after the startup
- scan, handle getting or dropping the key's content. -} - scan, handle getting or dropping the key's content. -}
@ -197,7 +198,7 @@ onDel file _ = do
liftAnnex $ liftAnnex $
Annex.Queue.addUpdateIndex =<< Annex.Queue.addUpdateIndex =<<
inRepo (Git.UpdateIndex.unstageFile file) inRepo (Git.UpdateIndex.unstageFile file)
liftIO $ madeChange file RmChange madeChange file RmChange
{- A directory has been deleted, or moved, so tell git to remove anything {- A directory has been deleted, or moved, so tell git to remove anything
- that was inside it from its cache. Since it could reappear at any time, - that was inside it from its cache. Since it could reappear at any time,
@ -211,7 +212,7 @@ onDelDir dir _ = do
debug ["directory deleted", dir] debug ["directory deleted", dir]
liftAnnex $ Annex.Queue.addCommand "rm" liftAnnex $ Annex.Queue.addCommand "rm"
[Params "--quiet -r --cached --ignore-unmatch --"] [dir] [Params "--quiet -r --cached --ignore-unmatch --"] [dir]
liftIO $ madeChange dir RmDirChange madeChange dir RmDirChange
{- Called when there's an error with inotify or kqueue. -} {- Called when there's an error with inotify or kqueue. -}
onErr :: Handler onErr :: Handler
@ -219,7 +220,7 @@ onErr msg _ = do
liftAnnex $ warning msg liftAnnex $ warning msg
dstatus <- getAssistant daemonStatusHandle dstatus <- getAssistant daemonStatusHandle
void $ liftIO $ addAlert dstatus $ warningAlert "watcher" msg void $ liftIO $ addAlert dstatus $ warningAlert "watcher" msg
liftIO noChange noChange
{- Adds a symlink to the index, without ever accessing the actual symlink {- Adds a symlink to the index, without ever accessing the actual symlink
- on disk. This avoids a race if git add is used, where the symlink is - on disk. This avoids a race if git add is used, where the symlink is

View file

@ -0,0 +1,54 @@
{- git-annex assistant change tracking
-
- Copyright 2012 Joey Hess <joey@kitenet.net>
-
- Licensed under the GNU GPL version 3 or higher.
-}
module Assistant.Types.Changes where
import Types.KeySource
import Utility.TSet
import Data.Time.Clock
data ChangeType = AddChange | LinkChange | RmChange | RmDirChange
deriving (Show, Eq)
type ChangeChan = TSet Change
data Change
= Change
{ changeTime :: UTCTime
, changeFile :: FilePath
, changeType :: ChangeType
}
| PendingAddChange
{ changeTime ::UTCTime
, changeFile :: FilePath
}
| InProcessAddChange
{ changeTime ::UTCTime
, keySource :: KeySource
}
deriving (Show)
newChangeChan :: IO ChangeChan
newChangeChan = newTSet
isPendingAddChange :: Change -> Bool
isPendingAddChange (PendingAddChange {}) = True
isPendingAddChange _ = False
isInProcessAddChange :: Change -> Bool
isInProcessAddChange (InProcessAddChange {}) = True
isInProcessAddChange _ = False
finishedChange :: Change -> Change
finishedChange c@(InProcessAddChange { keySource = ks }) = Change
{ changeTime = changeTime c
, changeFile = keyFilename ks
, changeType = AddChange
}
finishedChange c = c