module Ecluse.Core.Server.Cache.Store (
SingleFlight,
newSingleFlightWithBackend,
resolveSingleFlight,
PreparedStore,
prepareStore,
executePrepared,
lookupStore,
CacheOccupancy (..),
) where
import Data.Map.Strict qualified as Map
import UnliftIO.Exception (SomeAsyncException, mask, throwIO)
import Ecluse.Core.InFlight (guardInFlight)
import Ecluse.Core.Server.Cache.Backend.Internal
import Ecluse.Core.Telemetry.Metrics qualified as Metric
data SingleFlight e k v = SingleFlight
{ forall e k v. SingleFlight e k v -> Maybe (RetentionBackend k v)
sfBackend :: Maybe (RetentionBackend k v)
, 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
newSingleFlightWithBackend :: Maybe (RetentionBackend k v) -> IO (SingleFlight e k v)
newSingleFlightWithBackend :: forall k v e.
Maybe (RetentionBackend k v) -> IO (SingleFlight e k v)
newSingleFlightWithBackend Maybe (RetentionBackend k v)
backend = Maybe (RetentionBackend k v)
-> TVar (Map k (TMVar (FlightOutcome e v))) -> SingleFlight e k v
forall e k v.
Maybe (RetentionBackend k v)
-> TVar (Map k (TMVar (FlightOutcome e v))) -> SingleFlight e k v
SingleFlight Maybe (RetentionBackend k v)
backend (TVar (Map k (TMVar (FlightOutcome e v))) -> SingleFlight e k v)
-> IO (TVar (Map k (TMVar (FlightOutcome e v))))
-> IO (SingleFlight e k v)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Map k (TMVar (FlightOutcome e v))
-> IO (TVar (Map k (TMVar (FlightOutcome e v))))
forall (m :: * -> *) a. MonadIO m => a -> m (TVar a)
newTVarIO Map k (TMVar (FlightOutcome e v))
forall k a. Map k a
Map.empty
newtype PreparedStore e v = PreparedStore
{ forall e v. PreparedStore e v -> IO (Either e v)
executePrepared :: IO (Either e v)
}
prepareStore ::
(Ord k) =>
(Metric.CacheResult -> IO ()) ->
(CacheOccupancy -> IO ()) ->
IO () ->
SingleFlight e k v ->
k ->
IO (Either e v) ->
IO (PreparedStore e v)
prepareStore :: forall k e v.
Ord k =>
(CacheResult -> IO ())
-> (CacheOccupancy -> IO ())
-> IO ()
-> SingleFlight e k v
-> k
-> IO (Either e v)
-> IO (PreparedStore e v)
prepareStore CacheResult -> IO ()
recordRequest CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused SingleFlight e k v
sf k
key IO (Either e v)
fetch = do
held <- (CacheOccupancy -> IO ())
-> IO () -> SingleFlight e k v -> k -> IO (Maybe v)
forall e k v.
(CacheOccupancy -> IO ())
-> IO () -> SingleFlight e k v -> k -> IO (Maybe v)
lookupLocal CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused SingleFlight e k v
sf k
key
pure $ case held of
Just v
value -> IO (Either e v) -> PreparedStore e v
forall e v. IO (Either e v) -> PreparedStore e v
PreparedStore (CacheResult -> IO ()
recordRequest CacheResult
Metric.Hit IO () -> Either e v -> IO (Either e v)
forall (f :: * -> *) a b. Functor f => f a -> b -> f b
$> v -> Either e v
forall a b. b -> Either a b
Right v
value)
Maybe v
Nothing -> IO (Either e v) -> PreparedStore e v
forall e v. IO (Either e v) -> PreparedStore e v
PreparedStore ((CacheResult -> IO ())
-> (CacheOccupancy -> IO ())
-> IO ()
-> SingleFlight e k v
-> k
-> IO (Either e v)
-> IO (Either e v)
forall k e v.
Ord k =>
(CacheResult -> IO ())
-> (CacheOccupancy -> IO ())
-> IO ()
-> SingleFlight e k v
-> k
-> IO (Either e v)
-> IO (Either e v)
resolveDeferred CacheResult -> IO ()
recordRequest CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused SingleFlight e k v
sf k
key IO (Either e v)
fetch)
lookupLocal :: (CacheOccupancy -> IO ()) -> IO () -> SingleFlight e k v -> k -> IO (Maybe v)
lookupLocal :: forall e k v.
(CacheOccupancy -> IO ())
-> IO () -> SingleFlight e k v -> k -> IO (Maybe v)
lookupLocal CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused SingleFlight e k v
sf k
key = case SingleFlight e k v -> Maybe (RetentionBackend k v)
forall e k v. SingleFlight e k v -> Maybe (RetentionBackend k v)
sfBackend SingleFlight e k v
sf of
Just RetentionBackend k v
backend
| RetentionBackend k v -> BackendStorage
forall k v. RetentionBackend k v -> BackendStorage
rbStorage RetentionBackend k v
backend BackendStorage -> BackendStorage -> Bool
forall a. Eq a => a -> a -> Bool
== BackendStorage
LocalStorage ->
(CacheOccupancy -> IO ())
-> IO () -> Recency -> SingleFlight e k v -> k -> IO (Maybe v)
forall e k v.
(CacheOccupancy -> IO ())
-> IO () -> Recency -> SingleFlight e k v -> k -> IO (Maybe v)
lookupStore CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused Recency
RefreshRecency SingleFlight e k v
sf k
key
Maybe (RetentionBackend k v)
_ -> Maybe v -> IO (Maybe v)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe v
forall a. Maybe a
Nothing
resolveSingleFlight ::
(Ord k) =>
(Metric.CacheResult -> IO ()) ->
(CacheOccupancy -> IO ()) ->
IO () ->
SingleFlight e k v ->
k ->
IO (Either e v) ->
IO (Either e v)
resolveSingleFlight :: forall k e v.
Ord k =>
(CacheResult -> IO ())
-> (CacheOccupancy -> IO ())
-> IO ()
-> SingleFlight e k v
-> k
-> IO (Either e v)
-> IO (Either e v)
resolveSingleFlight CacheResult -> IO ()
recordRequest CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused SingleFlight e k v
sf k
key IO (Either e v)
fetch =
(CacheResult -> IO ())
-> (CacheOccupancy -> IO ())
-> IO ()
-> SingleFlight e k v
-> k
-> IO (Either e v)
-> IO (PreparedStore e v)
forall k e v.
Ord k =>
(CacheResult -> IO ())
-> (CacheOccupancy -> IO ())
-> IO ()
-> SingleFlight e k v
-> k
-> IO (Either e v)
-> IO (PreparedStore e v)
prepareStore CacheResult -> IO ()
recordRequest CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused SingleFlight e k v
sf k
key IO (Either e v)
fetch IO (PreparedStore e v)
-> (PreparedStore e v -> IO (Either e v)) -> IO (Either e v)
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= PreparedStore e v -> IO (Either e v)
forall e v. PreparedStore e v -> IO (Either e v)
executePrepared
resolveDeferred ::
(Ord k) =>
(Metric.CacheResult -> IO ()) ->
(CacheOccupancy -> IO ()) ->
IO () ->
SingleFlight e k v ->
k ->
IO (Either e v) ->
IO (Either e v)
resolveDeferred :: forall k e v.
Ord k =>
(CacheResult -> IO ())
-> (CacheOccupancy -> IO ())
-> IO ()
-> SingleFlight e k v
-> k
-> IO (Either e v)
-> IO (Either e v)
resolveDeferred CacheResult -> IO ()
recordRequest CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused 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 ->
let resolveAt :: (CacheResult -> IO ()) -> IO (Either e v)
resolveAt CacheResult -> IO ()
reportRequest = do
held <- IO (Maybe v) -> IO (Maybe v)
forall a. IO a -> IO a
restore ((CacheOccupancy -> IO ())
-> IO () -> SingleFlight e k v -> k -> IO (Maybe v)
forall e k v.
(CacheOccupancy -> IO ())
-> IO () -> SingleFlight e k v -> k -> IO (Maybe v)
lookupLocal CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused SingleFlight e k v
sf k
key)
case held of
Just v
value -> CacheResult -> IO ()
reportRequest CacheResult
Metric.Hit IO () -> Either e v -> IO (Either e v)
forall (f :: * -> *) a b. Functor f => f a -> b -> f b
$> v -> Either e v
forall a b. b -> Either a b
Right v
value
Maybe v
Nothing -> (CacheResult -> IO ()) -> IO (Either e v)
resolveMiss CacheResult -> IO ()
reportRequest
resolveMiss :: (CacheResult -> IO ()) -> IO (Either e v)
resolveMiss CacheResult -> IO ()
reportRequest = do
decision <- STM (Decision e v) -> IO (Decision e v)
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically (SingleFlight e k v -> k -> STM (Decision e v)
forall k e v.
Ord k =>
SingleFlight e k v -> k -> STM (Decision e v)
claimFlight SingleFlight e k v
sf k
key)
case decision of
Follow TMVar (FlightOutcome e v)
marker -> do
CacheResult -> IO ()
reportRequest CacheResult
Metric.Collapsed
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) -> (CacheResult -> IO ()) -> IO (Either e v)
resolveAt (IO () -> CacheResult -> IO ()
forall a b. a -> b -> a
const IO ()
forall (f :: * -> *). Applicative f => f ()
pass)
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 ->
(IO (Either e v) -> IO (Either e v))
-> (SomeException -> IO ())
-> IO ()
-> IO (Either e v)
-> IO (Either e v)
forall a.
(IO a -> IO a) -> (SomeException -> IO ()) -> IO () -> IO a -> IO a
guardInFlight IO (Either e v) -> IO (Either e v)
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) -> IO (Either e v))
-> IO (Either e v) -> IO (Either e v)
forall a b. (a -> b) -> a -> b
$ do
held <- IO (Maybe v) -> IO (Maybe v)
forall a. IO a -> IO a
restore ((CacheOccupancy -> IO ())
-> IO () -> Recency -> SingleFlight e k v -> k -> IO (Maybe v)
forall e k v.
(CacheOccupancy -> IO ())
-> IO () -> Recency -> SingleFlight e k v -> k -> IO (Maybe v)
lookupStore CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused Recency
RefreshRecency SingleFlight e k v
sf k
key)
fetched <- case held of
Just v
value -> CacheResult -> IO ()
reportRequest CacheResult
Metric.Hit IO () -> Either e v -> IO (Either e v)
forall (f :: * -> *) a b. Functor f => f a -> b -> f b
$> v -> Either e v
forall a b. b -> Either a b
Right v
value
Maybe v
Nothing -> do
CacheResult -> IO ()
reportRequest CacheResult
Metric.Miss
result <- IO (Either e v) -> IO (Either e v)
forall a. IO a -> IO a
restore IO (Either e v)
fetch
for_ (rightToMaybe result) $ \v
value ->
Maybe (RetentionBackend k v)
-> (RetentionBackend k v -> IO ()) -> IO ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
t a -> (a -> f b) -> f ()
for_ (SingleFlight e k v -> Maybe (RetentionBackend k v)
forall e k v. SingleFlight e k v -> Maybe (RetentionBackend k v)
sfBackend SingleFlight e k v
sf) ((RetentionBackend k v -> IO ()) -> IO ())
-> (RetentionBackend k v -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \RetentionBackend k v
backend ->
IO () -> IO ()
forall a. IO a -> IO a
restore (RetentionBackend k v
-> (CacheOccupancy -> IO ()) -> IO () -> IO () -> k -> v -> IO ()
forall k v.
RetentionBackend k v
-> (CacheOccupancy -> IO ()) -> IO () -> IO () -> k -> v -> IO ()
rbInsert RetentionBackend k v
backend CacheOccupancy -> IO ()
recordOccupancy IO ()
recordRefused IO ()
recordRefused k
key v
value)
pure result
atomically (putTMVar marker (either FlightFault FlightValue fetched))
pure fetched
in (CacheResult -> IO ()) -> IO (Either e v)
resolveAt CacheResult -> IO ()
recordRequest
where
deregister :: STM ()
deregister = TVar (Map k (TMVar (FlightOutcome e v)))
-> (Map k (TMVar (FlightOutcome e v))
-> Map k (TMVar (FlightOutcome e v)))
-> STM ()
forall a. TVar a -> (a -> a) -> STM ()
modifyTVar' (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) (k
-> Map k (TMVar (FlightOutcome e v))
-> Map k (TMVar (FlightOutcome e v))
forall k a. Ord k => k -> Map k a -> Map k a
Map.delete k
key)
lookupStore :: (CacheOccupancy -> IO ()) -> IO () -> Recency -> SingleFlight e k v -> k -> IO (Maybe v)
lookupStore :: forall e k v.
(CacheOccupancy -> IO ())
-> IO () -> Recency -> SingleFlight e k v -> k -> IO (Maybe v)
lookupStore CacheOccupancy -> IO ()
record IO ()
failed Recency
recency SingleFlight e k v
sf k
key = case SingleFlight e k v -> Maybe (RetentionBackend k v)
forall e k v. SingleFlight e k v -> Maybe (RetentionBackend k v)
sfBackend SingleFlight e k v
sf of
Maybe (RetentionBackend k v)
Nothing -> CacheOccupancy -> IO ()
record (Int -> Int -> CacheOccupancy
CacheOccupancy Int
0 Int
0) IO () -> Maybe v -> IO (Maybe v)
forall (f :: * -> *) a b. Functor f => f a -> b -> f b
$> Maybe v
forall a. Maybe a
Nothing
Just RetentionBackend k v
backend -> RetentionBackend k v
-> (CacheOccupancy -> IO ())
-> IO ()
-> Recency
-> k
-> IO (Maybe v)
forall k v.
RetentionBackend k v
-> (CacheOccupancy -> IO ())
-> IO ()
-> Recency
-> k
-> IO (Maybe v)
rbLookup RetentionBackend k v
backend CacheOccupancy -> IO ()
record IO ()
failed Recency
recency k
key
data Decision e v
= Follow (TMVar (FlightOutcome e v))
| Lead (TMVar (FlightOutcome e v))
claimFlight :: (Ord k) => SingleFlight e k v -> k -> STM (Decision e v)
claimFlight :: forall k e v.
Ord k =>
SingleFlight e k v -> k -> STM (Decision e v)
claimFlight SingleFlight e k v
sf k
key = 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 Bool -> STM ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (TMVar (FlightOutcome e v) -> FlightOutcome e v -> STM Bool
forall a. TMVar a -> a -> STM Bool
tryPutTMVar TMVar (FlightOutcome e v)
marker (SomeException -> FlightOutcome e v
forall e v. SomeException -> FlightOutcome e v
FlightOrphaned SomeException
err)))