Makes sequential uploads resumable.
This commit is contained in:
parent
43bfabd186
commit
c430e3d747
@ -24,6 +24,7 @@ library
|
|||||||
, Network.Minio.Data.ByteString
|
, Network.Minio.Data.ByteString
|
||||||
, Network.Minio.Data.Crypto
|
, Network.Minio.Data.Crypto
|
||||||
, Network.Minio.Data.Time
|
, Network.Minio.Data.Time
|
||||||
|
, Network.Minio.ListOps
|
||||||
, Network.Minio.PutObject
|
, Network.Minio.PutObject
|
||||||
, Network.Minio.Sign.V4
|
, Network.Minio.Sign.V4
|
||||||
, Network.Minio.Utils
|
, Network.Minio.Utils
|
||||||
@ -59,10 +60,10 @@ library
|
|||||||
default-language: Haskell2010
|
default-language: Haskell2010
|
||||||
default-extensions: FlexibleContexts
|
default-extensions: FlexibleContexts
|
||||||
, FlexibleInstances
|
, FlexibleInstances
|
||||||
, OverloadedStrings
|
|
||||||
, NoImplicitPrelude
|
|
||||||
, MultiParamTypeClasses
|
, MultiParamTypeClasses
|
||||||
, MultiWayIf
|
, MultiWayIf
|
||||||
|
, NoImplicitPrelude
|
||||||
|
, OverloadedStrings
|
||||||
, RankNTypes
|
, RankNTypes
|
||||||
, ScopedTypeVariables
|
, ScopedTypeVariables
|
||||||
, TypeFamilies
|
, TypeFamilies
|
||||||
@ -124,6 +125,7 @@ test-suite minio-hs-test
|
|||||||
, Network.Minio.Data.ByteString
|
, Network.Minio.Data.ByteString
|
||||||
, Network.Minio.Data.Crypto
|
, Network.Minio.Data.Crypto
|
||||||
, Network.Minio.Data.Time
|
, Network.Minio.Data.Time
|
||||||
|
, Network.Minio.ListOps
|
||||||
, Network.Minio.PutObject
|
, Network.Minio.PutObject
|
||||||
, Network.Minio.S3API
|
, Network.Minio.S3API
|
||||||
, Network.Minio.Sign.V4
|
, Network.Minio.Sign.V4
|
||||||
|
|||||||
@ -11,6 +11,6 @@ module Lib.Prelude
|
|||||||
import Protolude as Exports
|
import Protolude as Exports
|
||||||
|
|
||||||
import Data.Time as Exports (UTCTime)
|
import Data.Time as Exports (UTCTime)
|
||||||
import Data.Maybe as Exports (catMaybes, listToMaybe)
|
import Control.Monad.Trans.Maybe as Exports (runMaybeT, MaybeT(..))
|
||||||
|
|
||||||
import Control.Monad.Catch as Exports (throwM, MonadThrow, MonadCatch)
|
import Control.Monad.Catch as Exports (throwM, MonadThrow, MonadCatch)
|
||||||
|
|||||||
@ -40,17 +40,15 @@ module Network.Minio
|
|||||||
This module exports the high-level Minio API for object storage.
|
This module exports the high-level Minio API for object storage.
|
||||||
-}
|
-}
|
||||||
|
|
||||||
-- import qualified Control.Monad.Trans.Resource as R
|
|
||||||
import qualified Data.Conduit as C
|
import qualified Data.Conduit as C
|
||||||
import qualified Data.Conduit.Binary as CB
|
import qualified Data.Conduit.Binary as CB
|
||||||
import qualified Data.Conduit.List as CL
|
|
||||||
|
|
||||||
import Lib.Prelude
|
import Lib.Prelude
|
||||||
|
|
||||||
import Network.Minio.Data
|
import Network.Minio.Data
|
||||||
|
import Network.Minio.ListOps
|
||||||
import Network.Minio.PutObject
|
import Network.Minio.PutObject
|
||||||
import Network.Minio.S3API
|
import Network.Minio.S3API
|
||||||
-- import Network.Minio.Utils
|
|
||||||
|
|
||||||
-- | Fetch the object and write it to the given file safely. The
|
-- | Fetch the object and write it to the given file safely. The
|
||||||
-- object is first written to a temporary file in the same directory
|
-- object is first written to a temporary file in the same directory
|
||||||
@ -64,46 +62,3 @@ fGetObject bucket object fp = do
|
|||||||
fPutObject :: Bucket -> Object -> FilePath -> Minio ()
|
fPutObject :: Bucket -> Object -> FilePath -> Minio ()
|
||||||
fPutObject bucket object f = void $ putObject bucket object $
|
fPutObject bucket object f = void $ putObject bucket object $
|
||||||
ODFile f Nothing
|
ODFile f Nothing
|
||||||
|
|
||||||
-- | List objects in a bucket matching the given prefix. If recurse is
|
|
||||||
-- set to True objects matching prefix are recursively listed.
|
|
||||||
listObjects :: Bucket -> Maybe Text -> Bool -> C.Producer Minio ObjectInfo
|
|
||||||
listObjects bucket prefix recurse = loop Nothing
|
|
||||||
where
|
|
||||||
loop :: Maybe Text -> C.Producer Minio ObjectInfo
|
|
||||||
loop nextToken = do
|
|
||||||
let
|
|
||||||
delimiter = bool (Just "/") Nothing recurse
|
|
||||||
|
|
||||||
res <- lift $ listObjects' bucket prefix nextToken delimiter
|
|
||||||
CL.sourceList $ lorObjects res
|
|
||||||
when (lorHasMore res) $
|
|
||||||
loop (lorNextToken res)
|
|
||||||
|
|
||||||
-- | List incomplete uploads in a bucket matching the given prefix. If
|
|
||||||
-- recurse is set to True incomplete uploads for the given prefix are
|
|
||||||
-- recursively listed.
|
|
||||||
listIncompleteUploads :: Bucket -> Maybe Text -> Bool -> C.Producer Minio UploadInfo
|
|
||||||
listIncompleteUploads bucket prefix recurse = loop Nothing Nothing
|
|
||||||
where
|
|
||||||
loop :: Maybe Text -> Maybe Text -> C.Producer Minio UploadInfo
|
|
||||||
loop nextKeyMarker nextUploadIdMarker = do
|
|
||||||
let
|
|
||||||
delimiter = bool (Just "/") Nothing recurse
|
|
||||||
|
|
||||||
res <- lift $ listIncompleteUploads' bucket prefix delimiter nextKeyMarker nextUploadIdMarker
|
|
||||||
CL.sourceList $ lurUploads res
|
|
||||||
when (lurHasMore res) $
|
|
||||||
loop nextKeyMarker nextUploadIdMarker
|
|
||||||
|
|
||||||
-- | List object parts of an ongoing multipart upload for given
|
|
||||||
-- bucket, object and uploadId.
|
|
||||||
listIncompleteParts :: Bucket -> Object -> UploadId -> C.Producer Minio ListPartInfo
|
|
||||||
listIncompleteParts bucket object uploadId = loop Nothing
|
|
||||||
where
|
|
||||||
loop :: Maybe Text -> C.Producer Minio ListPartInfo
|
|
||||||
loop nextPartMarker = do
|
|
||||||
res <- lift $ listIncompleteParts' bucket object uploadId Nothing nextPartMarker
|
|
||||||
CL.sourceList $ lprParts res
|
|
||||||
when (lprHasMore res) $
|
|
||||||
loop (show <$> lprNextPart res)
|
|
||||||
|
|||||||
@ -78,7 +78,7 @@ data ListPartsResult = ListPartsResult {
|
|||||||
-- | Represents information about an object part in an ongoing
|
-- | Represents information about an object part in an ongoing
|
||||||
-- multipart upload.
|
-- multipart upload.
|
||||||
data ListPartInfo = ListPartInfo {
|
data ListPartInfo = ListPartInfo {
|
||||||
piNumber :: Int
|
piNumber :: PartNumber
|
||||||
, piETag :: ETag
|
, piETag :: ETag
|
||||||
, piSize :: Int64
|
, piSize :: Int64
|
||||||
, piModTime :: UTCTime
|
, piModTime :: UTCTime
|
||||||
@ -188,10 +188,8 @@ runMinio :: ConnectInfo -> Minio a -> ResourceT IO (Either MinioErr a)
|
|||||||
runMinio ci m = do
|
runMinio ci m = do
|
||||||
conn <- liftIO $ connect ci
|
conn <- liftIO $ connect ci
|
||||||
flip runReaderT conn . unMinio $
|
flip runReaderT conn . unMinio $
|
||||||
(m >>= (return . Right))
|
(m >>= (return . Right)) `MC.catches`
|
||||||
`MC.catch` handlerME
|
[MC.Handler handlerME, MC.Handler handlerHE, MC.Handler handlerFE]
|
||||||
`MC.catch` handlerHE
|
|
||||||
`MC.catch` handlerFE
|
|
||||||
where
|
where
|
||||||
handlerME = return . Left . ME
|
handlerME = return . Left . ME
|
||||||
handlerHE = return . Left . MEHttp
|
handlerHE = return . Left . MEHttp
|
||||||
|
|||||||
@ -2,13 +2,16 @@ module Network.Minio.Data.Crypto
|
|||||||
(
|
(
|
||||||
hashSHA256
|
hashSHA256
|
||||||
, hashSHA256FromSource
|
, hashSHA256FromSource
|
||||||
|
|
||||||
|
, hashMD5
|
||||||
|
|
||||||
, hmacSHA256
|
, hmacSHA256
|
||||||
, hmacSHA256RawBS
|
, hmacSHA256RawBS
|
||||||
, digestToBS
|
, digestToBS
|
||||||
, digestToBase16
|
, digestToBase16
|
||||||
) where
|
) where
|
||||||
|
|
||||||
import Crypto.Hash (SHA256(..), hashWith, Digest)
|
import Crypto.Hash (SHA256(..), MD5(..), hashWith, Digest)
|
||||||
import Crypto.Hash.Conduit (sinkHash)
|
import Crypto.Hash.Conduit (sinkHash)
|
||||||
import Crypto.MAC.HMAC (hmac, HMAC)
|
import Crypto.MAC.HMAC (hmac, HMAC)
|
||||||
import Data.ByteArray (ByteArrayAccess, convert)
|
import Data.ByteArray (ByteArrayAccess, convert)
|
||||||
@ -40,3 +43,6 @@ digestToBS = convert
|
|||||||
|
|
||||||
digestToBase16 :: ByteArrayAccess a => a -> ByteString
|
digestToBase16 :: ByteArrayAccess a => a -> ByteString
|
||||||
digestToBase16 = convertToBase Base16
|
digestToBase16 = convertToBase Base16
|
||||||
|
|
||||||
|
hashMD5 :: ByteString -> ByteString
|
||||||
|
hashMD5 = digestToBase16 . hashWith MD5
|
||||||
|
|||||||
56
src/Network/Minio/ListOps.hs
Normal file
56
src/Network/Minio/ListOps.hs
Normal file
@ -0,0 +1,56 @@
|
|||||||
|
module Network.Minio.ListOps where
|
||||||
|
|
||||||
|
import qualified Data.Conduit as C
|
||||||
|
import qualified Data.Conduit.List as CL
|
||||||
|
|
||||||
|
import Lib.Prelude
|
||||||
|
|
||||||
|
import Network.Minio.Data
|
||||||
|
import Network.Minio.S3API
|
||||||
|
|
||||||
|
-- | List objects in a bucket matching the given prefix. If recurse is
|
||||||
|
-- set to True objects matching prefix are recursively listed.
|
||||||
|
listObjects :: Bucket -> Maybe Text -> Bool -> C.Producer Minio ObjectInfo
|
||||||
|
listObjects bucket prefix recurse = loop Nothing
|
||||||
|
where
|
||||||
|
loop :: Maybe Text -> C.Producer Minio ObjectInfo
|
||||||
|
loop nextToken = do
|
||||||
|
let
|
||||||
|
delimiter = bool (Just "/") Nothing recurse
|
||||||
|
|
||||||
|
res <- lift $ listObjects' bucket prefix nextToken delimiter
|
||||||
|
CL.sourceList $ lorObjects res
|
||||||
|
when (lorHasMore res) $
|
||||||
|
loop (lorNextToken res)
|
||||||
|
|
||||||
|
-- | List incomplete uploads in a bucket matching the given prefix. If
|
||||||
|
-- recurse is set to True incomplete uploads for the given prefix are
|
||||||
|
-- recursively listed.
|
||||||
|
listIncompleteUploads :: Bucket -> Maybe Text -> Bool
|
||||||
|
-> C.Producer Minio UploadInfo
|
||||||
|
listIncompleteUploads bucket prefix recurse = loop Nothing Nothing
|
||||||
|
where
|
||||||
|
loop :: Maybe Text -> Maybe Text -> C.Producer Minio UploadInfo
|
||||||
|
loop nextKeyMarker nextUploadIdMarker = do
|
||||||
|
let
|
||||||
|
delimiter = bool (Just "/") Nothing recurse
|
||||||
|
|
||||||
|
res <- lift $ listIncompleteUploads' bucket prefix delimiter
|
||||||
|
nextKeyMarker nextUploadIdMarker
|
||||||
|
CL.sourceList $ lurUploads res
|
||||||
|
when (lurHasMore res) $
|
||||||
|
loop nextKeyMarker nextUploadIdMarker
|
||||||
|
|
||||||
|
-- | List object parts of an ongoing multipart upload for given
|
||||||
|
-- bucket, object and uploadId.
|
||||||
|
listIncompleteParts :: Bucket -> Object -> UploadId
|
||||||
|
-> C.Producer Minio ListPartInfo
|
||||||
|
listIncompleteParts bucket object uploadId = loop Nothing
|
||||||
|
where
|
||||||
|
loop :: Maybe Text -> C.Producer Minio ListPartInfo
|
||||||
|
loop nextPartMarker = do
|
||||||
|
res <- lift $ listIncompleteParts' bucket object uploadId Nothing
|
||||||
|
nextPartMarker
|
||||||
|
CL.sourceList $ lprParts res
|
||||||
|
when (lprHasMore res) $
|
||||||
|
loop (show <$> lprNextPart res)
|
||||||
@ -6,19 +6,24 @@ module Network.Minio.PutObject
|
|||||||
|
|
||||||
|
|
||||||
import qualified Data.Conduit as C
|
import qualified Data.Conduit as C
|
||||||
|
import qualified Data.Conduit.Combinators as CC
|
||||||
import qualified Data.Conduit.Binary as CB
|
import qualified Data.Conduit.Binary as CB
|
||||||
import qualified Data.List as List
|
import qualified Data.List as List
|
||||||
import qualified Data.ByteString.Lazy as LB
|
import qualified Data.ByteString.Lazy as LB
|
||||||
|
import qualified Data.Map.Strict as Map
|
||||||
|
|
||||||
import Lib.Prelude
|
import Lib.Prelude
|
||||||
|
|
||||||
import Network.Minio.Data
|
import Network.Minio.Data
|
||||||
|
import Network.Minio.Data.Crypto
|
||||||
|
import Network.Minio.ListOps
|
||||||
import Network.Minio.S3API
|
import Network.Minio.S3API
|
||||||
import Network.Minio.Utils
|
import Network.Minio.Utils
|
||||||
|
|
||||||
|
|
||||||
|
-- | max obj size is 5TiB
|
||||||
maxObjectSize :: Int64
|
maxObjectSize :: Int64
|
||||||
maxObjectSize = 5 * 1024 * 1024 * 1024 * 1024
|
maxObjectSize = 5 * 1024 * 1024 * oneMiB
|
||||||
|
|
||||||
oneMiB :: Int64
|
oneMiB :: Int64
|
||||||
oneMiB = 1024 * 1024
|
oneMiB = 1024 * 1024
|
||||||
@ -43,10 +48,8 @@ data ObjectData m = ODFile FilePath (Maybe Int64) -- ^ Takes filepath and option
|
|||||||
-- objects of all sizes, and even if the object size is unknown.
|
-- objects of all sizes, and even if the object size is unknown.
|
||||||
putObject :: Bucket -> Object -> ObjectData Minio -> Minio ETag
|
putObject :: Bucket -> Object -> ObjectData Minio -> Minio ETag
|
||||||
putObject b o (ODFile fp sizeMay) = do
|
putObject b o (ODFile fp sizeMay) = do
|
||||||
hResE <- withNewHandle fp $ \h -> do
|
hResE <- withNewHandle fp $ \h ->
|
||||||
isSeekable <- isHandleSeekable h
|
liftM2 (,) (isHandleSeekable h) (getFileSize h)
|
||||||
handleSizeMay <- getFileSize h
|
|
||||||
return (isSeekable, handleSizeMay)
|
|
||||||
|
|
||||||
(isSeekable, handleSizeMay) <- either (const $ return (False, Nothing)) return
|
(isSeekable, handleSizeMay) <- either (const $ return (False, Nothing)) return
|
||||||
hResE
|
hResE
|
||||||
@ -106,11 +109,13 @@ parallelMultipartUpload b o filePath size = do
|
|||||||
sequentialMultipartUpload :: Bucket -> Object -> Maybe Int64
|
sequentialMultipartUpload :: Bucket -> Object -> Maybe Int64
|
||||||
-> C.Producer Minio ByteString -> Minio ETag
|
-> C.Producer Minio ByteString -> Minio ETag
|
||||||
sequentialMultipartUpload b o sizeMay src = do
|
sequentialMultipartUpload b o sizeMay src = do
|
||||||
-- get new upload id.
|
(uidMay, pinfos) <- getExistingUpload b o
|
||||||
uploadId <- newMultipartUpload b o []
|
|
||||||
|
-- get a new upload id if needed.
|
||||||
|
uploadId <- maybe (newMultipartUpload b o []) return uidMay
|
||||||
|
|
||||||
-- upload parts in loop
|
-- upload parts in loop
|
||||||
uploadedParts <- loop uploadId rSrc partSizeInfo []
|
uploadedParts <- loop pinfos uploadId rSrc partSizeInfo []
|
||||||
|
|
||||||
-- complete multipart upload
|
-- complete multipart upload
|
||||||
completeMultipartUpload b o uploadId uploadedParts
|
completeMultipartUpload b o uploadId uploadedParts
|
||||||
@ -118,40 +123,52 @@ sequentialMultipartUpload b o sizeMay src = do
|
|||||||
rSrc = C.newResumableSource src
|
rSrc = C.newResumableSource src
|
||||||
partSizeInfo = selectPartSizes $ maybe maxObjectSize identity sizeMay
|
partSizeInfo = selectPartSizes $ maybe maxObjectSize identity sizeMay
|
||||||
|
|
||||||
|
-- returns partinfo if part is already uploaded.
|
||||||
|
checkUploadNeeded :: LByteString -> PartNumber
|
||||||
|
-> Map.Map PartNumber ListPartInfo
|
||||||
|
-> Maybe PartInfo
|
||||||
|
checkUploadNeeded lbs n pmap = do
|
||||||
|
pinfo@(ListPartInfo _ etag size _) <- Map.lookup n pmap
|
||||||
|
bool Nothing (return (PartInfo n etag)) $
|
||||||
|
LB.length lbs == size &&
|
||||||
|
hashMD5 (LB.toStrict lbs) == encodeUtf8 etag
|
||||||
|
|
||||||
-- make a sink that consumes only `s` bytes
|
-- make a sink that consumes only `s` bytes
|
||||||
limitedSink s = CB.isolate (fromIntegral s) C.=$= CB.sinkLbs
|
limitedSink s = CB.isolate (fromIntegral s) C.=$= CB.sinkLbs
|
||||||
|
|
||||||
-- FIXME: test, confirm and remove traceShowM statements
|
-- FIXME: test, confirm and remove traceShowM statements
|
||||||
loop _ _ [] uploadedParts = return $ reverse uploadedParts
|
loop _ _ _ [] uparts = return $ reverse uparts
|
||||||
loop uid rSource ((partNum, _, size):ps) u = do
|
loop pinfos uid rSource ((partNum, _, size):ps) u = do
|
||||||
|
|
||||||
-- load data from resume-able source into bytestring.
|
-- load data from resume-able source into bytestring.
|
||||||
(newSource, buf) <- rSource C.$$++ (limitedSink size)
|
(newSource, buf) <- rSource C.$$++ (limitedSink size)
|
||||||
traceShowM "psize: "
|
traceShowM "psize: "
|
||||||
traceShowM (LB.length buf)
|
traceShowM (LB.length buf)
|
||||||
-- check if we got size bytes.
|
|
||||||
if LB.length buf == size
|
|
||||||
-- upload the full size part.
|
|
||||||
then do pInfo <- putObjectPart b o uid partNum [] $
|
|
||||||
PayloadBS $ LB.toStrict buf
|
|
||||||
loop uid newSource ps (pInfo:u)
|
|
||||||
|
|
||||||
-- got a smaller part, so its the last one.
|
case checkUploadNeeded buf partNum pinfos of
|
||||||
else do traceShowM (("Found a piece with length < than "::[Char]) ++ show size ++ " - uploading as last and quitting.")
|
Just pinfo -> loop pinfos uid newSource ps (pinfo:u)
|
||||||
finalData <- newSource C.$$+- (limitedSink size)
|
Nothing -> do
|
||||||
traceShowM "finalData size:"
|
pInfo <- putObjectPart b o uid partNum [] $
|
||||||
traceShowM (LB.length finalData)
|
PayloadBS $ LB.toStrict buf
|
||||||
pInfo <- putObjectPart b o uid partNum [] $
|
|
||||||
PayloadBS $ LB.toStrict buf
|
if LB.length buf == size
|
||||||
return $ reverse (pInfo:u)
|
-- upload the full size part.
|
||||||
|
then loop pinfos uid newSource ps (pInfo:u)
|
||||||
|
|
||||||
|
-- got a smaller part, so its the last one.
|
||||||
|
else do traceShowM (("Found a piece with length < than "::[Char]) ++ show size ++ " - uploading as last and quitting.")
|
||||||
|
finalData <- newSource C.$$+- (limitedSink size)
|
||||||
|
traceShowM "finalData size:"
|
||||||
|
traceShowM (LB.length finalData)
|
||||||
|
return $ reverse (pInfo:u)
|
||||||
|
|
||||||
-- | Looks for incomplete uploads for an object. Returns the first one
|
-- | Looks for incomplete uploads for an object. Returns the first one
|
||||||
-- if there are many.
|
-- if there are many.
|
||||||
getExistingUpload :: Bucket -> Object
|
getExistingUpload :: Bucket -> Object
|
||||||
-> Minio (Maybe (UploadId, [ListPartInfo]))
|
-> Minio (Maybe UploadId, Map.Map PartNumber ListPartInfo)
|
||||||
getExistingUpload b o = do
|
getExistingUpload b o = do
|
||||||
uploadsRes <- listIncompleteUploads' b (Just o) Nothing Nothing Nothing
|
uidMay <- (fmap . fmap) uiUploadId $
|
||||||
case uiUploadId <$> listToMaybe (lurUploads uploadsRes) of
|
listIncompleteUploads b (Just o) False C.$$ CC.head
|
||||||
Nothing -> return Nothing
|
parts <- maybe (return [])
|
||||||
Just uid -> do
|
(\uid -> listIncompleteParts b o uid C.$$ CC.sinkList) uidMay
|
||||||
lpr <- listIncompleteParts' b o uid Nothing Nothing
|
return (uidMay, Map.fromList $ map (\p -> (piNumber p, p)) parts)
|
||||||
return $ Just (uid, lprParts lpr)
|
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user