-- SPDX-FileCopyrightText: 2026 Alexandra de Wit
-- SPDX-License-Identifier: MIT

-- | Local request coalescing with optional, separately owned retention.
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

-- | Active requests own flight results. Completion removes the only registry reference.
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

-- | Coalesce requests without retaining completed values when the backend is absent.
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

-- | A request-scoped read with any local value already pinned. Execute it once.
newtype PreparedStore e v = PreparedStore
    { forall e v. PreparedStore e v -> IO (Either e v)
executePrepared :: IO (Either e v)
    -- ^ Resolve the read, recording its request outcome only during execution.
    }

-- | Capture local retention without claiming a flight or contacting external storage.
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

{- | Share active work and release waiters on failure or cancellation.
Capacity refusals and external backend faults use the refusal callback.
-}
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)

-- | Probe the selected storage. Recency is a hint, and unsupported retention reports zero occupancy.
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)))