module Ecluse.Core.Registry.Progress (
Watch,
watched,
watchedRaising,
meteredReader,
meteredUpload,
) where
import Control.Exception (throwTo)
import Data.ByteString qualified as BS
import Data.ByteString.Builder (toLazyByteString)
import Data.ByteString.Lazy qualified as LBS
import GHC.Clock (getMonotonicTimeNSec)
import Network.HTTP.Client (BodyReader, GivesPopper, Popper, RequestBody (..))
import UnliftIO (try, withAsync)
import UnliftIO.Concurrent (ThreadId, myThreadId, threadDelay)
import Ecluse.Core.Security (ProgressFloor, floorMinBytes, floorWindowMicros)
data Watch = Watch ProgressFloor (IORef Window)
data Window = Window (Maybe Word64) Word64 Int
data BelowProgressFloor = BelowProgressFloor
deriving stock (Int -> BelowProgressFloor -> ShowS
[BelowProgressFloor] -> ShowS
BelowProgressFloor -> String
(Int -> BelowProgressFloor -> ShowS)
-> (BelowProgressFloor -> String)
-> ([BelowProgressFloor] -> ShowS)
-> Show BelowProgressFloor
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> BelowProgressFloor -> ShowS
showsPrec :: Int -> BelowProgressFloor -> ShowS
$cshow :: BelowProgressFloor -> String
show :: BelowProgressFloor -> String
$cshowList :: [BelowProgressFloor] -> ShowS
showList :: [BelowProgressFloor] -> ShowS
Show)
instance Exception BelowProgressFloor
watched :: ProgressFloor -> (Watch -> IO a) -> IO (Maybe a)
watched :: forall a. ProgressFloor -> (Watch -> IO a) -> IO (Maybe a)
watched ProgressFloor
progress Watch -> IO a
transfer =
IO a -> IO (Either BelowProgressFloor a)
forall (m :: * -> *) e a.
(MonadUnliftIO m, Exception e) =>
m a -> m (Either e a)
try (ProgressFloor -> (Watch -> IO a) -> IO a
forall a. ProgressFloor -> (Watch -> IO a) -> IO a
watchedRaising ProgressFloor
progress Watch -> IO a
transfer) IO (Either BelowProgressFloor a)
-> (Either BelowProgressFloor a -> Maybe a) -> IO (Maybe a)
forall (f :: * -> *) a b. Functor f => f a -> (a -> b) -> f b
<&> \case
Left BelowProgressFloor
BelowProgressFloor -> Maybe a
forall a. Maybe a
Nothing
Right a
result -> a -> Maybe a
forall a. a -> Maybe a
Just a
result
watchedRaising :: ProgressFloor -> (Watch -> IO a) -> IO a
watchedRaising :: forall a. ProgressFloor -> (Watch -> IO a) -> IO a
watchedRaising ProgressFloor
progress Watch -> IO a
transfer = do
owner <- IO ThreadId
forall (m :: * -> *). MonadIO m => m ThreadId
myThreadId
watch <- Watch progress <$> newIORef (Window Nothing 0 0)
withAsync (watchdog watch owner) (\Async ()
_ -> Watch -> IO a
transfer Watch
watch)
watchdog :: Watch -> ThreadId -> IO ()
watchdog :: Watch -> ThreadId -> IO ()
watchdog (Watch ProgressFloor
progress IORef Window
window) ThreadId
owner = IO ()
go
where
budget :: Word64
budget = Int -> Word64
forall a b. (Integral a, Num b) => a -> b
fromIntegral (ProgressFloor -> Int
floorWindowMicros ProgressFloor
progress) Word64 -> Word64 -> Word64
forall a. Num a => a -> a -> a
* Word64
1_000
go :: IO ()
go = do
now <- IO Word64
getMonotonicTimeNSec
current <- readIORef window
let spent = Word64 -> Window -> Word64
waitedBy Word64
now Window
current
if spent >= budget
then throwTo owner BelowProgressFloor
else threadDelay (fromIntegral (pause budget spent current `div` 1_000) + 1) >> go
pause :: Word64 -> Word64 -> Window -> Word64
pause :: Word64 -> Word64 -> Window -> Word64
pause Word64
budget Word64
spent = \case
Window (Just Word64
_) Word64
_ Int
_ -> Word64
budget Word64 -> Word64 -> Word64
forall a. Num a => a -> a -> a
- Word64
spent
Window Maybe Word64
Nothing Word64
_ Int
_ -> Word64 -> Word64 -> Word64
forall a. Ord a => a -> a -> a
max (Word64
budget Word64 -> Word64 -> Word64
forall a. Integral a => a -> a -> a
`div` Word64
16) (Word64
budget Word64 -> Word64 -> Word64
forall a. Num a => a -> a -> a
- Word64
spent)
meteredReader :: Watch -> BodyReader -> BodyReader
meteredReader :: Watch -> BodyReader -> BodyReader
meteredReader Watch
watch BodyReader
readChunk = do
Watch -> IO ()
beginWait Watch
watch
chunk <- BodyReader
readChunk
endWait watch (BS.length chunk)
pure chunk
meteredUpload :: Watch -> RequestBody -> RequestBody
meteredUpload :: Watch -> RequestBody -> RequestBody
meteredUpload Watch
watch RequestBody
body = case RequestBody
body of
RequestBodyLBS ByteString
bytes | Bool -> Bool
not (ByteString -> Bool
LBS.null ByteString
bytes) -> Int64 -> [ByteString] -> RequestBody
fromChunks (ByteString -> Int64
LBS.length ByteString
bytes) (ByteString -> [ByteString]
LBS.toChunks ByteString
bytes)
RequestBodyBS ByteString
bytes | Bool -> Bool
not (ByteString -> Bool
BS.null ByteString
bytes) -> Int64 -> [ByteString] -> RequestBody
fromChunks (Int -> Int64
forall a b. (Integral a, Num b) => a -> b
fromIntegral (ByteString -> Int
BS.length ByteString
bytes)) [ByteString
bytes]
RequestBodyBuilder Int64
size Builder
builder | Int64
size Int64 -> Int64 -> Bool
forall a. Ord a => a -> a -> Bool
> Int64
0 -> Int64 -> [ByteString] -> RequestBody
fromChunks Int64
size (ByteString -> [ByteString]
LBS.toChunks (Builder -> ByteString
toLazyByteString Builder
builder))
RequestBodyStream Int64
size GivesPopper ()
gives -> Int64 -> GivesPopper () -> RequestBody
RequestBodyStream Int64
size (Watch -> GivesPopper () -> GivesPopper ()
meteredGives Watch
watch GivesPopper ()
gives)
RequestBodyStreamChunked GivesPopper ()
gives -> GivesPopper () -> RequestBody
RequestBodyStreamChunked (Watch -> GivesPopper () -> GivesPopper ()
meteredGives Watch
watch GivesPopper ()
gives)
RequestBodyIO IO RequestBody
io -> IO RequestBody -> RequestBody
RequestBodyIO (Watch -> RequestBody -> RequestBody
meteredUpload Watch
watch (RequestBody -> RequestBody) -> IO RequestBody -> IO RequestBody
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> IO RequestBody
io)
RequestBody
_ -> RequestBody
body
where
fromChunks :: Int64 -> [ByteString] -> RequestBody
fromChunks Int64
size [ByteString]
chunks = Int64 -> GivesPopper () -> RequestBody
RequestBodyStream Int64
size (Watch -> GivesPopper () -> GivesPopper ()
meteredGives Watch
watch ([ByteString] -> GivesPopper ()
givesChunks [ByteString]
chunks))
givesChunks :: [ByteString] -> GivesPopper ()
givesChunks :: [ByteString] -> GivesPopper ()
givesChunks [ByteString]
chunks NeedsPopper ()
needsPopper = do
remaining <- [ByteString] -> IO (IORef [ByteString])
forall (m :: * -> *) a. MonadIO m => a -> m (IORef a)
newIORef [ByteString]
chunks
needsPopper . atomicModifyIORef' remaining $ \case
[] -> ([], ByteString
BS.empty)
ByteString
chunk : [ByteString]
rest -> ([ByteString]
rest, ByteString
chunk)
meteredGives :: Watch -> GivesPopper () -> GivesPopper ()
meteredGives :: Watch -> GivesPopper () -> GivesPopper ()
meteredGives Watch
watch GivesPopper ()
gives NeedsPopper ()
needsPopper = GivesPopper ()
gives GivesPopper () -> GivesPopper ()
forall a b. (a -> b) -> a -> b
$ \BodyReader
popper -> do
handed <- Int -> IO (IORef Int)
forall (m :: * -> *) a. MonadIO m => a -> m (IORef a)
newIORef Int
0
pending <- newIORef BS.empty
needsPopper $ do
readIORef handed >>= \Int
size -> Bool -> IO () -> IO ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (Int
size Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
> Int
0) (Watch -> Int -> IO ()
endWait Watch
watch Int
size)
slice <- nextSlice pending popper
writeIORef handed (BS.length slice)
unless (BS.null slice) (beginWait watch)
pure slice
nextSlice :: IORef ByteString -> Popper -> IO ByteString
nextSlice :: IORef ByteString -> BodyReader -> BodyReader
nextSlice IORef ByteString
pending BodyReader
popper = do
held <- IORef ByteString -> BodyReader
forall (m :: * -> *) a. MonadIO m => IORef a -> m a
readIORef IORef ByteString
pending
chunk <- if BS.null held then popper else pure held
let (slice, rest) = BS.splitAt uploadSliceBytes chunk
writeIORef pending rest
pure slice
uploadSliceBytes :: Int
uploadSliceBytes :: Int
uploadSliceBytes = Int
64 Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
1024
beginWait :: Watch -> IO ()
beginWait :: Watch -> IO ()
beginWait (Watch ProgressFloor
_ IORef Window
window) = do
now <- IO Word64
getMonotonicTimeNSec
atomicModifyIORef' window (\(Window Maybe Word64
_ Word64
waited Int
moved) -> (Maybe Word64 -> Word64 -> Int -> Window
Window (Word64 -> Maybe Word64
forall a. a -> Maybe a
Just Word64
now) Word64
waited Int
moved, ()))
endWait :: Watch -> Int -> IO ()
endWait :: Watch -> Int -> IO ()
endWait (Watch ProgressFloor
progress IORef Window
window) Int
bytes = do
now <- IO Word64
getMonotonicTimeNSec
atomicModifyIORef' window $ \current :: Window
current@(Window Maybe Word64
_ Word64
_ Int
moved) ->
if Int
moved Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
bytes Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
>= ProgressFloor -> Int
floorMinBytes ProgressFloor
progress
then (Maybe Word64 -> Word64 -> Int -> Window
Window Maybe Word64
forall a. Maybe a
Nothing Word64
0 Int
0, ())
else (Maybe Word64 -> Word64 -> Int -> Window
Window Maybe Word64
forall a. Maybe a
Nothing (Word64 -> Window -> Word64
waitedBy Word64
now Window
current) (Int
moved Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
bytes), ())
waitedBy :: Word64 -> Window -> Word64
waitedBy :: Word64 -> Window -> Word64
waitedBy Word64
now (Window Maybe Word64
since Word64
waited Int
_) = Word64
waited Word64 -> Word64 -> Word64
forall a. Num a => a -> a -> a
+ Word64 -> (Word64 -> Word64) -> Maybe Word64 -> Word64
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Word64
0 (\Word64
began -> if Word64
now Word64 -> Word64 -> Bool
forall a. Ord a => a -> a -> Bool
> Word64
began then Word64
now Word64 -> Word64 -> Word64
forall a. Num a => a -> a -> a
- Word64
began else Word64
0) Maybe Word64
since