module Ecluse.Core.Server.Cache.Store (
SingleFlight,
newSingleFlight,
resolveSingleFlight,
lookupStore,
lookupStoreTouching,
CacheOccupancy (..),
) where
import Data.Cache (Cache)
import Data.Cache qualified as Cache
import Data.Map.Strict qualified as Map
import Data.Time (NominalDiffTime)
import System.Clock (Clock (Monotonic), TimeSpec, fromNanoSecs, getTime)
import UnliftIO.Exception (SomeAsyncException, mask, throwIO)
import UnliftIO.MVar (withMVar)
import Ecluse.Core.InFlight (guardInFlight)
import Ecluse.Core.Telemetry.Metrics qualified as Metric
data Weighted v = Weighted
{ forall v. Weighted v -> v
wValue :: v
, forall v. Weighted v -> Int
wWeight :: Int
, forall v. Weighted v -> IORef Word64
wStamp :: IORef Word64
}
data SingleFlight e k v = SingleFlight
{ forall e k v. SingleFlight e k v -> Cache k (Weighted v)
sfStore :: Cache k (Weighted v)
, forall e k v. SingleFlight e k v -> Int
sfMaxEntries :: Int
, forall e k v. SingleFlight e k v -> Int
sfMaxBytes :: Int
, forall e k v. SingleFlight e k v -> v -> Int
sfWeigh :: v -> Int
, forall e k v. SingleFlight e k v -> IORef Word64
sfClock :: IORef Word64
, forall e k v. SingleFlight e k v -> MVar ()
sfInsertLock :: MVar ()
, forall e k v.
SingleFlight e k v -> TVar (Map k (TMVar (FlightOutcome e v)))
sfInFlight :: TVar (Map k (TMVar (FlightOutcome e v)))
}
data FlightOutcome e v
= FlightValue v
| FlightFault e
| FlightOrphaned SomeException
newSingleFlight :: NominalDiffTime -> Int -> Int -> (v -> Int) -> IO (SingleFlight e k v)
newSingleFlight :: forall v e k.
NominalDiffTime
-> Int -> Int -> (v -> Int) -> IO (SingleFlight e k v)
newSingleFlight NominalDiffTime
ttl Int
maxEntries Int
maxBytes v -> Int
weigh = do
store <- Maybe TimeSpec -> IO (Cache k (Weighted v))
forall k v. Maybe TimeSpec -> IO (Cache k v)
Cache.newCache (TimeSpec -> Maybe TimeSpec
forall a. a -> Maybe a
Just (NominalDiffTime -> TimeSpec
toTimeSpec NominalDiffTime
ttl))
clock <- newIORef 0
inFlight <- newTVarIO Map.empty
insertLock <- newMVar ()
pure
SingleFlight
{ sfStore = store
, sfMaxEntries = max 1 maxEntries
, sfMaxBytes = max 1 maxBytes
, sfWeigh = weigh
, sfClock = clock
, sfInsertLock = insertLock
, sfInFlight = inFlight
}
resolveSingleFlight ::
(Hashable k, Ord k) =>
IO () ->
(Metric.CacheResult -> IO ()) ->
(CacheOccupancy -> IO ()) ->
SingleFlight e k v ->
k ->
IO (Either e v) ->
IO (Either e v)
resolveSingleFlight :: forall k e v.
(Hashable k, Ord k) =>
IO ()
-> (CacheResult -> IO ())
-> (CacheOccupancy -> IO ())
-> SingleFlight e k v
-> k
-> IO (Either e v)
-> IO (Either e v)
resolveSingleFlight IO ()
afterClaim CacheResult -> IO ()
recordRequest CacheOccupancy -> IO ()
recordInsert SingleFlight e k v
sf k
key IO (Either e v)
fetch = ((forall a. IO a -> IO a) -> IO (Either e v)) -> IO (Either e v)
forall (m :: * -> *) b.
MonadUnliftIO m =>
((forall a. m a -> m a) -> m b) -> m b
mask (((forall a. IO a -> IO a) -> IO (Either e v)) -> IO (Either e v))
-> ((forall a. IO a -> IO a) -> IO (Either e v)) -> IO (Either e v)
forall a b. (a -> b) -> a -> b
$ \forall a. IO a -> IO a
restore -> do
nowT <- Clock -> IO TimeSpec
getTime Clock
Monotonic
decision <- atomically (decideSingleFlight sf key nowT)
case decision of
Hit Weighted v
weighted -> do
CacheResult -> IO ()
recordRequest CacheResult
Metric.Hit
SingleFlight e k v -> Weighted v -> IO ()
forall e k v. SingleFlight e k v -> Weighted v -> IO ()
touch SingleFlight e k v
sf Weighted v
weighted
Either e v -> IO (Either e v)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (v -> Either e v
forall a b. b -> Either a b
Right (Weighted v -> v
forall v. Weighted v -> v
wValue Weighted v
weighted))
Follow TMVar (FlightOutcome e v)
marker -> do
CacheResult -> IO ()
recordRequest CacheResult
Metric.Miss
outcome <- IO (FlightOutcome e v) -> IO (FlightOutcome e v)
forall a. IO a -> IO a
restore (STM (FlightOutcome e v) -> IO (FlightOutcome e v)
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically (TMVar (FlightOutcome e v) -> STM (FlightOutcome e v)
forall a. TMVar a -> STM a
readTMVar TMVar (FlightOutcome e v)
marker))
case outcome of
FlightValue v
fetched -> Either e v -> IO (Either e v)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (v -> Either e v
forall a b. b -> Either a b
Right v
fetched)
FlightFault e
fault -> Either e v -> IO (Either e v)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (e -> Either e v
forall a b. a -> Either a b
Left e
fault)
FlightOrphaned SomeException
err -> case SomeException -> Maybe SomeAsyncException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
err of
Just (SomeAsyncException
_ :: SomeAsyncException) ->
IO (Either e v) -> IO (Either e v)
forall a. IO a -> IO a
restore (IO ()
-> (CacheResult -> IO ())
-> (CacheOccupancy -> IO ())
-> SingleFlight e k v
-> k
-> IO (Either e v)
-> IO (Either e v)
forall k e v.
(Hashable k, Ord k) =>
IO ()
-> (CacheResult -> IO ())
-> (CacheOccupancy -> IO ())
-> SingleFlight e k v
-> k
-> IO (Either e v)
-> IO (Either e v)
resolveSingleFlight IO ()
afterClaim (IO () -> CacheResult -> IO ()
forall a b. a -> b -> a
const IO ()
forall (f :: * -> *). Applicative f => f ()
pass) CacheOccupancy -> IO ()
recordInsert SingleFlight e k v
sf k
key IO (Either e v)
fetch)
Maybe SomeAsyncException
Nothing -> SomeException -> IO (Either e v)
forall (m :: * -> *) e a. (MonadIO m, Exception e) => e -> m a
throwIO SomeException
err
Lead TMVar (FlightOutcome e v)
marker -> do
CacheResult -> IO ()
recordRequest CacheResult
Metric.Miss
(outcome, occupancy) <- (IO (Either e v, Maybe CacheOccupancy)
-> IO (Either e v, Maybe CacheOccupancy))
-> (SomeException -> IO ())
-> IO ()
-> IO (Either e v, Maybe CacheOccupancy)
-> IO (Either e v, Maybe CacheOccupancy)
forall a.
(IO a -> IO a) -> (SomeException -> IO ()) -> IO () -> IO a -> IO a
guardInFlight IO (Either e v, Maybe CacheOccupancy)
-> IO (Either e v, Maybe CacheOccupancy)
forall a. a -> a
id (TMVar (FlightOutcome e v) -> SomeException -> IO ()
forall e v. TMVar (FlightOutcome e v) -> SomeException -> IO ()
orphan TMVar (FlightOutcome e v)
marker) (STM () -> IO ()
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically STM ()
deregister) (IO (Either e v, Maybe CacheOccupancy)
-> IO (Either e v, Maybe CacheOccupancy))
-> IO (Either e v, Maybe CacheOccupancy)
-> IO (Either e v, Maybe CacheOccupancy)
forall a b. (a -> b) -> a -> b
$ do
fetched <- IO (Either e v) -> IO (Either e v)
forall a. IO a -> IO a
restore (IO ()
afterClaim IO () -> IO (Either e v) -> IO (Either e v)
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> IO (Either e v)
fetch)
atomically (putTMVar marker (either FlightFault FlightValue fetched))
inserted <- join <$> traverse (insertBounded sf key) (rightToMaybe fetched)
pure (fetched, inserted)
traverse_ recordInsert occupancy
pure outcome
where
deregister :: STM ()
deregister :: STM ()
deregister = do
inFlight <- TVar (Map k (TMVar (FlightOutcome e v)))
-> STM (Map k (TMVar (FlightOutcome e v)))
forall a. TVar a -> STM a
readTVar (SingleFlight e k v -> TVar (Map k (TMVar (FlightOutcome e v)))
forall e k v.
SingleFlight e k v -> TVar (Map k (TMVar (FlightOutcome e v)))
sfInFlight SingleFlight e k v
sf)
writeTVar (sfInFlight sf) (Map.delete key inFlight)
insertBounded :: (Hashable k) => SingleFlight e k v -> k -> v -> IO (Maybe CacheOccupancy)
insertBounded :: forall k e v.
Hashable k =>
SingleFlight e k v -> k -> v -> IO (Maybe CacheOccupancy)
insertBounded SingleFlight e k v
sf k
key v
value
| Int
weight Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
> SingleFlight e k v -> Int
forall e k v. SingleFlight e k v -> Int
sfMaxBytes SingleFlight e k v
sf = Maybe CacheOccupancy -> IO (Maybe CacheOccupancy)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe CacheOccupancy
forall a. Maybe a
Nothing
| Bool
otherwise = MVar ()
-> (() -> IO (Maybe CacheOccupancy)) -> IO (Maybe CacheOccupancy)
forall (m :: * -> *) a b.
MonadUnliftIO m =>
MVar a -> (a -> m b) -> m b
withMVar (SingleFlight e k v -> MVar ()
forall e k v. SingleFlight e k v -> MVar ()
sfInsertLock SingleFlight e k v
sf) ((() -> IO (Maybe CacheOccupancy)) -> IO (Maybe CacheOccupancy))
-> (() -> IO (Maybe CacheOccupancy)) -> IO (Maybe CacheOccupancy)
forall a b. (a -> b) -> a -> b
$ \() -> do
Cache k (Weighted v) -> IO ()
forall k v. (Eq k, Hashable k) => Cache k v -> IO ()
Cache.purgeExpired (SingleFlight e k v -> Cache k (Weighted v)
forall e k v. SingleFlight e k v -> Cache k (Weighted v)
sfStore SingleFlight e k v
sf)
SingleFlight e k v -> Int -> IO ()
forall k e v. Hashable k => SingleFlight e k v -> Int -> IO ()
evictToBudget SingleFlight e k v
sf Int
weight
stamp <- SingleFlight e k v -> IO Word64
forall e k v. SingleFlight e k v -> IO Word64
nextStamp SingleFlight e k v
sf
stampRef <- newIORef stamp
Cache.insert (sfStore sf) key (Weighted{wValue = value, wWeight = weight, wStamp = stampRef})
Just <$> occupancyOf sf
where
weight :: Int
weight = SingleFlight e k v -> v -> Int
forall e k v. SingleFlight e k v -> v -> Int
sfWeigh SingleFlight e k v
sf v
value
evictToBudget :: (Hashable k) => SingleFlight e k v -> Int -> IO ()
evictToBudget :: forall k e v. Hashable k => SingleFlight e k v -> Int -> IO ()
evictToBudget SingleFlight e k v
sf Int
incoming = do
held <- Cache k (Weighted v) -> IO [(k, Weighted v, Maybe TimeSpec)]
forall k v. Cache k v -> IO [(k, v, Maybe TimeSpec)]
Cache.toList (SingleFlight e k v -> Cache k (Weighted v)
forall e k v. SingleFlight e k v -> Cache k (Weighted v)
sfStore SingleFlight e k v
sf)
stamped <- traverse stampOf held
let resident = [Int] -> Int
forall a (f :: * -> *). (Foldable f, Num a) => f a -> a
sum [Weighted v -> Int
forall v. Weighted v -> Int
wWeight Weighted v
w | (k
_, Weighted v
w, Maybe TimeSpec
_) <- [(k, Weighted v, Maybe TimeSpec)]
held]
oldestFirst = ((Word64, k, Int) -> Word64)
-> [(Word64, k, Int)] -> [(Word64, k, Int)]
forall b a. Ord b => (a -> b) -> [a] -> [a]
sortOn (\(Word64
stamp, k
_, Int
_) -> Word64
stamp) [(Word64, k, Int)]
stamped
go oldestFirst resident (length held)
where
stampOf :: (b, Weighted v, c) -> m (Word64, b, Int)
stampOf (b
k, Weighted v
w, c
_) = do
s <- IORef Word64 -> m Word64
forall (m :: * -> *) a. MonadIO m => IORef a -> m a
readIORef (Weighted v -> IORef Word64
forall v. Weighted v -> IORef Word64
wStamp Weighted v
w)
pure (s, k, wWeight w)
fits :: Int -> Int -> Bool
fits Int
resident Int
count = Int
resident Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
incoming Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
<= SingleFlight e k v -> Int
forall e k v. SingleFlight e k v -> Int
sfMaxBytes SingleFlight e k v
sf Bool -> Bool -> Bool
&& Int
count Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
< SingleFlight e k v -> Int
forall e k v. SingleFlight e k v -> Int
sfMaxEntries SingleFlight e k v
sf
go :: [(a, k, Int)] -> Int -> Int -> IO ()
go [(a, k, Int)]
victims Int
resident Int
count
| Int -> Int -> Bool
fits Int
resident Int
count = IO ()
forall (f :: * -> *). Applicative f => f ()
pass
| Bool
otherwise = case [(a, k, Int)]
victims of
[] -> IO ()
forall (f :: * -> *). Applicative f => f ()
pass
((a
_, k
k, Int
weight) : [(a, k, Int)]
rest) -> do
Cache k (Weighted v) -> k -> IO ()
forall k v. (Eq k, Hashable k) => Cache k v -> k -> IO ()
Cache.delete (SingleFlight e k v -> Cache k (Weighted v)
forall e k v. SingleFlight e k v -> Cache k (Weighted v)
sfStore SingleFlight e k v
sf) k
k
[(a, k, Int)] -> Int -> Int -> IO ()
go [(a, k, Int)]
rest (Int
resident Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
weight) (Int
count Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
1)
occupancyOf :: SingleFlight e k v -> IO CacheOccupancy
occupancyOf :: forall e k v. SingleFlight e k v -> IO CacheOccupancy
occupancyOf SingleFlight e k v
sf = do
held <- Cache k (Weighted v) -> IO [(k, Weighted v, Maybe TimeSpec)]
forall k v. Cache k v -> IO [(k, v, Maybe TimeSpec)]
Cache.toList (SingleFlight e k v -> Cache k (Weighted v)
forall e k v. SingleFlight e k v -> Cache k (Weighted v)
sfStore SingleFlight e k v
sf)
pure CacheOccupancy{occEntries = length held, occBytes = sum [wWeight w | (_, w, _) <- held]}
nextStamp :: SingleFlight e k v -> IO Word64
nextStamp :: forall e k v. SingleFlight e k v -> IO Word64
nextStamp SingleFlight e k v
sf = IORef Word64 -> (Word64 -> (Word64, Word64)) -> IO Word64
forall (m :: * -> *) a b.
MonadIO m =>
IORef a -> (a -> (a, b)) -> m b
atomicModifyIORef' (SingleFlight e k v -> IORef Word64
forall e k v. SingleFlight e k v -> IORef Word64
sfClock SingleFlight e k v
sf) (\Word64
n -> let n' :: Word64
n' = Word64
n Word64 -> Word64 -> Word64
forall a. Num a => a -> a -> a
+ Word64
1 in (Word64
n', Word64
n'))
touch :: SingleFlight e k v -> Weighted v -> IO ()
touch :: forall e k v. SingleFlight e k v -> Weighted v -> IO ()
touch SingleFlight e k v
sf Weighted v
weighted = SingleFlight e k v -> IO Word64
forall e k v. SingleFlight e k v -> IO Word64
nextStamp SingleFlight e k v
sf IO Word64 -> (Word64 -> IO ()) -> IO ()
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= IORef Word64 -> Word64 -> IO ()
forall (m :: * -> *) a. MonadIO m => IORef a -> a -> m ()
writeIORef (Weighted v -> IORef Word64
forall v. Weighted v -> IORef Word64
wStamp Weighted v
weighted)
lookupStore :: (Hashable k) => SingleFlight e k v -> k -> IO (Maybe v)
lookupStore :: forall k e v. Hashable k => SingleFlight e k v -> k -> IO (Maybe v)
lookupStore SingleFlight e k v
sf k
key = (Weighted v -> v) -> Maybe (Weighted v) -> Maybe v
forall a b. (a -> b) -> Maybe a -> Maybe b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap Weighted v -> v
forall v. Weighted v -> v
wValue (Maybe (Weighted v) -> Maybe v)
-> IO (Maybe (Weighted v)) -> IO (Maybe v)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Cache k (Weighted v) -> k -> IO (Maybe (Weighted v))
forall k v. (Eq k, Hashable k) => Cache k v -> k -> IO (Maybe v)
Cache.lookup (SingleFlight e k v -> Cache k (Weighted v)
forall e k v. SingleFlight e k v -> Cache k (Weighted v)
sfStore SingleFlight e k v
sf) k
key
lookupStoreTouching :: (Hashable k) => SingleFlight e k v -> k -> IO (Maybe v)
lookupStoreTouching :: forall k e v. Hashable k => SingleFlight e k v -> k -> IO (Maybe v)
lookupStoreTouching SingleFlight e k v
sf k
key =
Cache k (Weighted v) -> k -> IO (Maybe (Weighted v))
forall k v. (Eq k, Hashable k) => Cache k v -> k -> IO (Maybe v)
Cache.lookup (SingleFlight e k v -> Cache k (Weighted v)
forall e k v. SingleFlight e k v -> Cache k (Weighted v)
sfStore SingleFlight e k v
sf) k
key IO (Maybe (Weighted v))
-> (Maybe (Weighted v) -> IO (Maybe v)) -> IO (Maybe v)
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= (Weighted v -> IO v) -> Maybe (Weighted v) -> IO (Maybe v)
forall (t :: * -> *) (f :: * -> *) a b.
(Traversable t, Applicative f) =>
(a -> f b) -> t a -> f (t b)
forall (f :: * -> *) a b.
Applicative f =>
(a -> f b) -> Maybe a -> f (Maybe b)
traverse (\Weighted v
weighted -> Weighted v -> v
forall v. Weighted v -> v
wValue Weighted v
weighted v -> IO () -> IO v
forall a b. a -> IO b -> IO a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ SingleFlight e k v -> Weighted v -> IO ()
forall e k v. SingleFlight e k v -> Weighted v -> IO ()
touch SingleFlight e k v
sf Weighted v
weighted)
data Decision e v
= Hit (Weighted v)
| Follow (TMVar (FlightOutcome e v))
| Lead (TMVar (FlightOutcome e v))
decideSingleFlight :: (Hashable k, Ord k) => SingleFlight e k v -> k -> TimeSpec -> STM (Decision e v)
decideSingleFlight :: forall k e v.
(Hashable k, Ord k) =>
SingleFlight e k v -> k -> TimeSpec -> STM (Decision e v)
decideSingleFlight SingleFlight e k v
sf k
key TimeSpec
nowT = do
hit <- Bool
-> k
-> Cache k (Weighted v)
-> TimeSpec
-> STM (Maybe (Weighted v))
forall k v.
(Eq k, Hashable k) =>
Bool -> k -> Cache k v -> TimeSpec -> STM (Maybe v)
Cache.lookupSTM Bool
False k
key (SingleFlight e k v -> Cache k (Weighted v)
forall e k v. SingleFlight e k v -> Cache k (Weighted v)
sfStore SingleFlight e k v
sf) TimeSpec
nowT
case hit of
Just Weighted v
weighted -> Decision e v -> STM (Decision e v)
forall a. a -> STM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Weighted v -> Decision e v
forall e v. Weighted v -> Decision e v
Hit Weighted v
weighted)
Maybe (Weighted v)
Nothing -> do
inFlight <- TVar (Map k (TMVar (FlightOutcome e v)))
-> STM (Map k (TMVar (FlightOutcome e v)))
forall a. TVar a -> STM a
readTVar (SingleFlight e k v -> TVar (Map k (TMVar (FlightOutcome e v)))
forall e k v.
SingleFlight e k v -> TVar (Map k (TMVar (FlightOutcome e v)))
sfInFlight SingleFlight e k v
sf)
case Map.lookup key inFlight of
Just TMVar (FlightOutcome e v)
marker -> Decision e v -> STM (Decision e v)
forall a. a -> STM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (TMVar (FlightOutcome e v) -> Decision e v
forall e v. TMVar (FlightOutcome e v) -> Decision e v
Follow TMVar (FlightOutcome e v)
marker)
Maybe (TMVar (FlightOutcome e v))
Nothing -> do
marker <- STM (TMVar (FlightOutcome e v))
forall a. STM (TMVar a)
newEmptyTMVar
writeTVar (sfInFlight sf) (Map.insert key marker inFlight)
pure (Lead marker)
orphan :: TMVar (FlightOutcome e v) -> SomeException -> IO ()
orphan :: forall e v. TMVar (FlightOutcome e v) -> SomeException -> IO ()
orphan TMVar (FlightOutcome e v)
marker SomeException
err =
STM () -> IO ()
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically (STM () -> IO ()) -> STM () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
unfilled <- TMVar (FlightOutcome e v) -> STM Bool
forall a. TMVar a -> STM Bool
isEmptyTMVar TMVar (FlightOutcome e v)
marker
when unfilled (putTMVar marker (FlightOrphaned err))
data CacheOccupancy = CacheOccupancy
{ CacheOccupancy -> Int
occEntries :: Int
, CacheOccupancy -> Int
occBytes :: Int
}
toTimeSpec :: NominalDiffTime -> TimeSpec
toTimeSpec :: NominalDiffTime -> TimeSpec
toTimeSpec NominalDiffTime
ttl = Integer -> TimeSpec
fromNanoSecs (Integer -> Integer -> Integer
forall a. Ord a => a -> a -> a
max Integer
0 (Double -> Integer
forall b. Integral b => Double -> b
forall a b. (RealFrac a, Integral b) => a -> b
round (NominalDiffTime -> Double
forall a b. (Real a, Fractional b) => a -> b
realToFrac NominalDiffTime
ttl Double -> Double -> Double
forall a. Num a => a -> a -> a
* Double
1e9 :: Double) :: Integer))