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

{- | Streaming ingest of the osv.dev export archive Pilot compiles @osv.db@ from.

The feed aggregates many upstream databases, so the bounds here are per entry: a drop is
counted in 'IngestStats' and the rest of the archive keeps flowing. 'ilMaxAdvisoryBytes'
applies before the bytes are retained and before the JSON decodes, and the refused entry drains
to its boundary so the entries after it stay aligned. The aggregate verdict is 'systemicDrop',
which the compiler reads once the stream completes.
-}
module Ecluse.Core.Osv.Stream (
    streamOsvUrl,
    parseOsvStream,

    -- * Ingest bounds and drop accounting
    IngestLimits (..),
    defaultIngestLimits,
    IngestStats (..),
    OsvIngest,
    newOsvIngest,
    readIngestStats,
    resetIngestStats,
    systemicDrop,
    PilotIngestAborted (..),

    -- * What one attempt learned about its source
    OsvAttempt (..),
    readOsvAttempt,
    resetOsvAttempt,
) where

import Codec.Archive.Zip.Conduit.Types (ZipEntry (..))
import Codec.Archive.Zip.Conduit.UnZip (unZipStream)
import Conduit
import Data.Aeson (decodeStrict)
import Data.ByteString qualified as BS
import Data.Time (UTCTime)
import Katip (KatipContext, Severity (..), logFM, ls)
import Network.HTTP.Simple (getResponseBody, getResponseHeader, httpSource, parseRequest, setRequestCheckStatus)
import Network.HTTP.Types.Header (hLastModified)
import OpenTelemetry.Trace.Core (SpanKind (Internal), TracerProvider, addAttribute)

import Ecluse.Core.Ecosystem (Ecosystem)
import Ecluse.Core.Osv.Advisory (ExtractedOsv, OsvAdvisory, extPackage, extractFromAdvisory, orderableBounds, osvId, osvModified, unorderableBounds)
import Ecluse.Core.Osv.Ecosystem (OsvEcosystem (osvEcosystemTag, osvMaxAdvisoryFanOut))
import Ecluse.Core.Osv.Epss (EpssScores)
import Ecluse.Core.Osv.Provenance (lastModifiedOf, parseSourceTime)
import Ecluse.Core.Security.Authority (dialledAuthorityLabel)
import Ecluse.Core.Telemetry.Span (closeOptionalSpan, openOptionalSpan)

-- | The per-advisory byte bound one ingest pass holds every zip entry to.
newtype IngestLimits = IngestLimits
    { IngestLimits -> Int
ilMaxAdvisoryBytes :: Int
    {- ^ Largest decompressed advisory JSON, in bytes, the ingest accepts from one
    zip entry. It drops a larger one. Bounds memory and, transitively, decode cost.
    -}
    }
    deriving stock (IngestLimits -> IngestLimits -> Bool
(IngestLimits -> IngestLimits -> Bool)
-> (IngestLimits -> IngestLimits -> Bool) -> Eq IngestLimits
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: IngestLimits -> IngestLimits -> Bool
== :: IngestLimits -> IngestLimits -> Bool
$c/= :: IngestLimits -> IngestLimits -> Bool
/= :: IngestLimits -> IngestLimits -> Bool
Eq, Int -> IngestLimits -> ShowS
[IngestLimits] -> ShowS
IngestLimits -> String
(Int -> IngestLimits -> ShowS)
-> (IngestLimits -> String)
-> ([IngestLimits] -> ShowS)
-> Show IngestLimits
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> IngestLimits -> ShowS
showsPrec :: Int -> IngestLimits -> ShowS
$cshow :: IngestLimits -> String
show :: IngestLimits -> String
$cshowList :: [IngestLimits] -> ShowS
showList :: [IngestLimits] -> ShowS
Show)

-- | An 8 MiB per-advisory ceiling.
defaultIngestLimits :: IngestLimits
defaultIngestLimits :: IngestLimits
defaultIngestLimits = IngestLimits{ilMaxAdvisoryBytes :: Int
ilMaxAdvisoryBytes = Int
8 Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
1024 Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
1024}

{- | The running tally of one ingest pass. Pilot reads it once the stream completes to
decide whether the artifact is trustworthy enough to publish ('systemicDrop').
-}
data IngestStats = IngestStats
    { IngestStats -> Int
statAccepted :: !Int
    -- ^ Advisory entries that decoded successfully.
    , IngestStats -> Int
statDroppedOversize :: !Int
    -- ^ Entries dropped for breaching 'ilMaxAdvisoryBytes'.
    , IngestStats -> Int
statDroppedMalformed :: !Int
    -- ^ Entries dropped because their JSON did not decode.
    , IngestStats -> Int
statUnorderable :: !Int
    {- ^ Rows kept with a bound the grammar cannot parse ('orderableBounds'). Counted in
    rows, so it stays out of 'systemicDrop'.
    -}
    , IngestStats -> Int
statUnusableModified :: !Int
    {- ^ Records whose @modified@ no grammar reads, or which date it after the run's clock.
    Their rows are kept and only the date is ignored, so this counts records, not drops.
    -}
    }
    deriving stock (IngestStats -> IngestStats -> Bool
(IngestStats -> IngestStats -> Bool)
-> (IngestStats -> IngestStats -> Bool) -> Eq IngestStats
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: IngestStats -> IngestStats -> Bool
== :: IngestStats -> IngestStats -> Bool
$c/= :: IngestStats -> IngestStats -> Bool
/= :: IngestStats -> IngestStats -> Bool
Eq, Int -> IngestStats -> ShowS
[IngestStats] -> ShowS
IngestStats -> String
(Int -> IngestStats -> ShowS)
-> (IngestStats -> String)
-> ([IngestStats] -> ShowS)
-> Show IngestStats
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> IngestStats -> ShowS
showsPrec :: Int -> IngestStats -> ShowS
$cshow :: IngestStats -> String
show :: IngestStats -> String
$cshowList :: [IngestStats] -> ShowS
showList :: [IngestStats] -> ShowS
Show)

emptyIngestStats :: IngestStats
emptyIngestStats :: IngestStats
emptyIngestStats = Int -> Int -> Int -> Int -> Int -> IngestStats
IngestStats Int
0 Int
0 Int
0 Int
0 Int
0

{- | What one ingest attempt learned about its source beside the rows. A retry replaces it
whole, so it always describes the attempt that produced the tally beside it.
-}
data OsvAttempt = OsvAttempt
    { OsvAttempt -> Maybe UTCTime
oaLastModified :: Maybe UTCTime
    -- ^ The @Last-Modified@ the export answered this attempt with.
    , OsvAttempt -> Maybe UTCTime
oaNewestModified :: Maybe UTCTime
    -- ^ The newest @modified@ across the records this attempt read.
    }
    deriving stock (OsvAttempt -> OsvAttempt -> Bool
(OsvAttempt -> OsvAttempt -> Bool)
-> (OsvAttempt -> OsvAttempt -> Bool) -> Eq OsvAttempt
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: OsvAttempt -> OsvAttempt -> Bool
== :: OsvAttempt -> OsvAttempt -> Bool
$c/= :: OsvAttempt -> OsvAttempt -> Bool
/= :: OsvAttempt -> OsvAttempt -> Bool
Eq, Int -> OsvAttempt -> ShowS
[OsvAttempt] -> ShowS
OsvAttempt -> String
(Int -> OsvAttempt -> ShowS)
-> (OsvAttempt -> String)
-> ([OsvAttempt] -> ShowS)
-> Show OsvAttempt
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> OsvAttempt -> ShowS
showsPrec :: Int -> OsvAttempt -> ShowS
$cshow :: OsvAttempt -> String
show :: OsvAttempt -> String
$cshowList :: [OsvAttempt] -> ShowS
showList :: [OsvAttempt] -> ShowS
Show)

emptyOsvAttempt :: OsvAttempt
emptyOsvAttempt :: OsvAttempt
emptyOsvAttempt = Maybe UTCTime -> Maybe UTCTime -> OsvAttempt
OsvAttempt Maybe UTCTime
forall a. Maybe a
Nothing Maybe UTCTime
forall a. Maybe a
Nothing

-- The mutable drop tally for one ingest pass. Opaque: read it with 'readIngestStats'.
newtype IngestCounter = IngestCounter {IngestCounter -> IORef IngestStats
counterRef :: IORef IngestStats}

-- | The context one ingest pass threads through the stream.
data OsvIngest = OsvIngest
    { OsvIngest -> IngestLimits
ingestLimits :: IngestLimits
    , OsvIngest -> IngestCounter
ingestCounter :: IngestCounter
    , OsvIngest -> OsvEcosystem
ingestEcosystem :: OsvEcosystem
    {- ^ The feed this pass compiles: it carries the grammar that orders the pass's bounds and
    the fan-out an ordinary advisory of the feed stays under.
    -}
    , OsvIngest -> EpssScores
ingestEpss :: EpssScores
    -- ^ The pass's EPSS table, joined onto each advisory as it is extracted.
    , OsvIngest -> UTCTime
ingestNow :: UTCTime
    -- ^ The run's clock, which every record's @modified@ is judged against.
    , OsvIngest -> IORef OsvAttempt
ingestAttempt :: IORef OsvAttempt
    }

-- | A fresh ingest context with the given bounds, feed, EPSS table and clock, and a zeroed tally.
newOsvIngest :: (MonadIO m) => IngestLimits -> OsvEcosystem -> EpssScores -> UTCTime -> m OsvIngest
newOsvIngest :: forall (m :: * -> *).
MonadIO m =>
IngestLimits
-> OsvEcosystem -> EpssScores -> UTCTime -> m OsvIngest
newOsvIngest IngestLimits
limits OsvEcosystem
eco EpssScores
scores UTCTime
now = do
    counter <- IORef IngestStats -> IngestCounter
IngestCounter (IORef IngestStats -> IngestCounter)
-> m (IORef IngestStats) -> m IngestCounter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> IngestStats -> m (IORef IngestStats)
forall (m :: * -> *) a. MonadIO m => a -> m (IORef a)
newIORef IngestStats
emptyIngestStats
    attempt <- newIORef emptyOsvAttempt
    pure (OsvIngest limits counter eco scores now attempt)

-- | Read the current drop tally.
readIngestStats :: (MonadIO m) => OsvIngest -> m IngestStats
readIngestStats :: forall (m :: * -> *). MonadIO m => OsvIngest -> m IngestStats
readIngestStats OsvIngest
ingest = IORef IngestStats -> m IngestStats
forall (m :: * -> *) a. MonadIO m => IORef a -> m a
readIORef (IngestCounter -> IORef IngestStats
counterRef (OsvIngest -> IngestCounter
ingestCounter OsvIngest
ingest))

{- | Zero the tally. The compiler re-streams from a clean slate on each retry attempt
and zeroes the tally alongside it, so the tally reflects only the final attempt.
-}
resetIngestStats :: (MonadIO m) => OsvIngest -> m ()
resetIngestStats :: forall (m :: * -> *). MonadIO m => OsvIngest -> m ()
resetIngestStats OsvIngest
ingest = IORef IngestStats -> IngestStats -> m ()
forall (m :: * -> *) a. MonadIO m => IORef a -> a -> m ()
writeIORef (IngestCounter -> IORef IngestStats
counterRef (OsvIngest -> IngestCounter
ingestCounter OsvIngest
ingest)) IngestStats
emptyIngestStats

-- | Read what the current attempt learned about its source.
readOsvAttempt :: (MonadIO m) => OsvIngest -> m OsvAttempt
readOsvAttempt :: forall (m :: * -> *). MonadIO m => OsvIngest -> m OsvAttempt
readOsvAttempt OsvIngest
ingest = IORef OsvAttempt -> m OsvAttempt
forall (m :: * -> *) a. MonadIO m => IORef a -> m a
readIORef (OsvIngest -> IORef OsvAttempt
ingestAttempt OsvIngest
ingest)

{- | Forget the source metadata, alongside 'resetIngestStats'. A retry re-reads the export, so
last attempt's response header and record dates must not survive into this one's artifact.
-}
resetOsvAttempt :: (MonadIO m) => OsvIngest -> m ()
resetOsvAttempt :: forall (m :: * -> *). MonadIO m => OsvIngest -> m ()
resetOsvAttempt OsvIngest
ingest = IORef OsvAttempt -> OsvAttempt -> m ()
forall (m :: * -> *) a. MonadIO m => IORef a -> a -> m ()
writeIORef (OsvIngest -> IORef OsvAttempt
ingestAttempt OsvIngest
ingest) OsvAttempt
emptyOsvAttempt

-- | Whether the drop tally requires Pilot to refuse publication.
systemicDrop :: IngestStats -> Bool
systemicDrop :: IngestStats -> Bool
systemicDrop IngestStats
s =
    Int
dropped Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
>= Int
systemicDropFloor Bool -> Bool -> Bool
&& Int
dropped Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
100 Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
>= Int
total Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
systemicDropPercent
  where
    dropped :: Int
dropped = IngestStats -> Int
statDroppedOversize IngestStats
s Int -> Int -> Int
forall a. Num a => a -> a -> a
+ IngestStats -> Int
statDroppedMalformed IngestStats
s
    total :: Int
total = Int
dropped Int -> Int -> Int
forall a. Num a => a -> a -> a
+ IngestStats -> Int
statAccepted IngestStats
s

systemicDropFloor :: Int
systemicDropFloor :: Int
systemicDropFloor = Int
16

systemicDropPercent :: Int
systemicDropPercent :: Int
systemicDropPercent = Int
10

{- | Raised when systemic drops or zero relevant output prevent publication.
The tally records the rejected pass, without replacing a consumer's last-good artifact.
-}
newtype PilotIngestAborted = PilotIngestAborted IngestStats
    deriving stock (Int -> PilotIngestAborted -> ShowS
[PilotIngestAborted] -> ShowS
PilotIngestAborted -> String
(Int -> PilotIngestAborted -> ShowS)
-> (PilotIngestAborted -> String)
-> ([PilotIngestAborted] -> ShowS)
-> Show PilotIngestAborted
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> PilotIngestAborted -> ShowS
showsPrec :: Int -> PilotIngestAborted -> ShowS
$cshow :: PilotIngestAborted -> String
show :: PilotIngestAborted -> String
$cshowList :: [PilotIngestAborted] -> ShowS
showList :: [PilotIngestAborted] -> ShowS
Show)

instance Exception PilotIngestAborted

-- | Fetch the OSV zip and stream its contents, bounded by @ingest@.
streamOsvUrl :: (MonadResource m, MonadThrow m, KatipContext m) => Maybe TracerProvider -> OsvIngest -> String -> ConduitT i ExtractedOsv m ()
streamOsvUrl :: forall (m :: * -> *) i.
(MonadResource m, MonadThrow m, KatipContext m) =>
Maybe TracerProvider
-> OsvIngest -> String -> ConduitT i ExtractedOsv m ()
streamOsvUrl Maybe TracerProvider
mTracerProvider OsvIngest
ingest String
urlStr = do
    m () -> ConduitT i ExtractedOsv m ()
forall (m :: * -> *) a.
Monad m =>
m a -> ConduitT i ExtractedOsv m a
forall (t :: (* -> *) -> * -> *) (m :: * -> *) a.
(MonadTrans t, Monad m) =>
m a -> t m a
lift (m () -> ConduitT i ExtractedOsv m ())
-> m () -> ConduitT i ExtractedOsv m ()
forall a b. (a -> b) -> a -> b
$ Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
InfoS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"Initializing OSV stream from " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text -> Text
dialledAuthorityLabel (String -> Text
forall a. ToText a => a -> Text
toText String
urlStr)))
    IO (Maybe Span)
-> (Maybe Span -> IO ())
-> (Maybe Span -> ConduitT i ExtractedOsv m ())
-> ConduitT i ExtractedOsv m ()
forall (m :: * -> *) a i o r.
MonadResource m =>
IO a -> (a -> IO ()) -> (a -> ConduitT i o m r) -> ConduitT i o m r
bracketP
        (Maybe TracerProvider -> SpanKind -> Text -> IO (Maybe Span)
forall (m :: * -> *).
MonadIO m =>
Maybe TracerProvider -> SpanKind -> Text -> m (Maybe Span)
openOptionalSpan Maybe TracerProvider
mTracerProvider SpanKind
Internal Text
"ecluse.pilot.osv.stream")
        Maybe Span -> IO ()
forall (m :: * -> *). MonadIO m => Maybe Span -> m ()
closeOptionalSpan
        ( \Maybe Span
mSpan -> do
            Maybe Span
-> (Span -> ConduitT i ExtractedOsv m ())
-> ConduitT i ExtractedOsv m ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ Maybe Span
mSpan ((Span -> ConduitT i ExtractedOsv m ())
 -> ConduitT i ExtractedOsv m ())
-> (Span -> ConduitT i ExtractedOsv m ())
-> ConduitT i ExtractedOsv m ()
forall a b. (a -> b) -> a -> b
$ \Span
sp -> Span -> Text -> Text -> ConduitT i ExtractedOsv m ()
forall (m :: * -> *) a.
(MonadIO m, ToAttribute a) =>
Span -> Text -> a -> m ()
addAttribute Span
sp Text
"ecluse.osv.source_host" (Text -> Text
dialledAuthorityLabel (String -> Text
forall a. ToText a => a -> Text
toText String
urlStr))
            -- Reject non-2xx responses before unzip so the retry policy sees HTTP failures.
            req <- IO Request -> ConduitT i ExtractedOsv m Request
forall a. IO a -> ConduitT i ExtractedOsv m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO Request -> ConduitT i ExtractedOsv m Request)
-> IO Request -> ConduitT i ExtractedOsv m Request
forall a b. (a -> b) -> a -> b
$ Request -> Request
setRequestCheckStatus (Request -> Request) -> IO Request -> IO Request
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> String -> IO Request
forall (m :: * -> *). MonadThrow m => String -> m Request
parseRequest String
urlStr
            httpSource req $ \Response (ConduitM i ByteString m ())
res -> do
                OsvIngest -> [ByteString] -> ConduitT i ExtractedOsv m ()
forall (m :: * -> *).
MonadIO m =>
OsvIngest -> [ByteString] -> m ()
recordResponseDate OsvIngest
ingest (HeaderName -> Response (ConduitM i ByteString m ()) -> [ByteString]
forall a. HeaderName -> Response a -> [ByteString]
getResponseHeader HeaderName
hLastModified Response (ConduitM i ByteString m ())
res)
                Response (ConduitM i ByteString m ()) -> ConduitM i ByteString m ()
forall a. Response a -> a
getResponseBody Response (ConduitM i ByteString m ())
res ConduitM i ByteString m ()
-> ConduitT ByteString ExtractedOsv m ()
-> ConduitT i ExtractedOsv m ()
forall (m :: * -> *) a b c r.
Monad m =>
ConduitT a b m () -> ConduitT b c m r -> ConduitT a c m r
.| Maybe TracerProvider
-> OsvIngest -> ConduitT ByteString ExtractedOsv m ()
forall (m :: * -> *).
(MonadResource m, MonadThrow m, KatipContext m) =>
Maybe TracerProvider
-> OsvIngest -> ConduitT ByteString ExtractedOsv m ()
parseOsvStream Maybe TracerProvider
mTracerProvider OsvIngest
ingest
        )

-- | Parse the zip stream and emit ExtractedOsv, bounded by @ingest@.
parseOsvStream :: (MonadResource m, MonadThrow m, KatipContext m) => Maybe TracerProvider -> OsvIngest -> ConduitT ByteString ExtractedOsv m ()
parseOsvStream :: forall (m :: * -> *).
(MonadResource m, MonadThrow m, KatipContext m) =>
Maybe TracerProvider
-> OsvIngest -> ConduitT ByteString ExtractedOsv m ()
parseOsvStream Maybe TracerProvider
mTracerProvider OsvIngest
ingest = do
    m () -> ConduitT ByteString ExtractedOsv m ()
forall (m :: * -> *) a.
Monad m =>
m a -> ConduitT ByteString ExtractedOsv m a
forall (t :: (* -> *) -> * -> *) (m :: * -> *) a.
(MonadTrans t, Monad m) =>
m a -> t m a
lift (m () -> ConduitT ByteString ExtractedOsv m ())
-> m () -> ConduitT ByteString ExtractedOsv m ()
forall a b. (a -> b) -> a -> b
$ Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
InfoS (String -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (String
"Starting OSV zip extraction and parsing pipeline" :: String))
    IO (Maybe Span)
-> (Maybe Span -> IO ())
-> (Maybe Span -> ConduitT ByteString ExtractedOsv m ())
-> ConduitT ByteString ExtractedOsv m ()
forall (m :: * -> *) a i o r.
MonadResource m =>
IO a -> (a -> IO ()) -> (a -> ConduitT i o m r) -> ConduitT i o m r
bracketP
        (Maybe TracerProvider -> SpanKind -> Text -> IO (Maybe Span)
forall (m :: * -> *).
MonadIO m =>
Maybe TracerProvider -> SpanKind -> Text -> m (Maybe Span)
openOptionalSpan Maybe TracerProvider
mTracerProvider SpanKind
Internal Text
"ecluse.pilot.osv.parse")
        Maybe Span -> IO ()
forall (m :: * -> *). MonadIO m => Maybe Span -> m ()
closeOptionalSpan
        (\Maybe Span
_ -> ConduitT ByteString (Either ZipEntry ByteString) m ZipInfo
-> ConduitT ByteString (Either ZipEntry ByteString) m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void ((forall a. IO a -> m a)
-> ConduitT ByteString (Either ZipEntry ByteString) IO ZipInfo
-> ConduitT ByteString (Either ZipEntry ByteString) m ZipInfo
forall (m :: * -> *) (n :: * -> *) i o r.
Monad m =>
(forall a. m a -> n a) -> ConduitT i o m r -> ConduitT i o n r
transPipe IO a -> m a
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO ConduitT ByteString (Either ZipEntry ByteString) IO ZipInfo
forall (m :: * -> *).
(MonadThrow m, PrimMonad m) =>
ConduitM ByteString (Either ZipEntry ByteString) m ZipInfo
unZipStream) ConduitT ByteString (Either ZipEntry ByteString) m ()
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
-> ConduitT ByteString ExtractedOsv m ()
forall (m :: * -> *) a b c r.
Monad m =>
ConduitT a b m () -> ConduitT b c m r -> ConduitT a c m r
.| OsvIngest
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall (m :: * -> *).
(MonadThrow m, KatipContext m) =>
OsvIngest
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
processZipEntries OsvIngest
ingest)

processZipEntries :: (MonadThrow m, KatipContext m) => OsvIngest -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
processZipEntries :: forall (m :: * -> *).
(MonadThrow m, KatipContext m) =>
OsvIngest
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
processZipEntries OsvIngest
ingest =
    ConduitT
  (Either ZipEntry ByteString)
  ExtractedOsv
  m
  (Maybe (Either ZipEntry ByteString))
forall (m :: * -> *) i o. Monad m => ConduitT i o m (Maybe i)
await ConduitT
  (Either ZipEntry ByteString)
  ExtractedOsv
  m
  (Maybe (Either ZipEntry ByteString))
-> (Maybe (Either ZipEntry ByteString)
    -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ())
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall a b.
ConduitT (Either ZipEntry ByteString) ExtractedOsv m a
-> (a -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m b)
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
        Maybe (Either ZipEntry ByteString)
Nothing -> m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall (m :: * -> *) a.
Monad m =>
m a -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m a
forall (t :: (* -> *) -> * -> *) (m :: * -> *) a.
(MonadTrans t, Monad m) =>
m a -> t m a
lift (m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ())
-> m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall a b. (a -> b) -> a -> b
$ Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
InfoS (String -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (String
"OSV stream fully processed" :: String))
        Just (Left ZipEntry
entry) -> do
            outcome <- Int
-> ConduitT
     (Either ZipEntry ByteString) ExtractedOsv m EntryOutcome
forall (m :: * -> *) o.
Monad m =>
Int -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
collectFile (IngestLimits -> Int
ilMaxAdvisoryBytes (OsvIngest -> IngestLimits
ingestLimits OsvIngest
ingest))
            handleEntry ingest entry outcome
            processZipEntries ingest
        Just (Right ByteString
_) -> OsvIngest
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall (m :: * -> *).
(MonadThrow m, KatipContext m) =>
OsvIngest
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
processZipEntries OsvIngest
ingest

-- The outcome of accumulating one zip entry: its bytes, or a signal that it breached
-- the byte cap. The signal carries the entry's full decompressed size, for the log.
data EntryOutcome = EntryBytes !ByteString | EntryOversize !Int

handleEntry :: (KatipContext m) => OsvIngest -> ZipEntry -> EntryOutcome -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
handleEntry :: forall (m :: * -> *).
KatipContext m =>
OsvIngest
-> ZipEntry
-> EntryOutcome
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
handleEntry OsvIngest
ingest ZipEntry
entry = \case
    EntryOversize Int
seen -> m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall (m :: * -> *) a.
Monad m =>
m a -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m a
forall (t :: (* -> *) -> * -> *) (m :: * -> *) a.
(MonadTrans t, Monad m) =>
m a -> t m a
lift (m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ())
-> m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall a b. (a -> b) -> a -> b
$ do
        IngestCounter -> m ()
forall (m :: * -> *). MonadIO m => IngestCounter -> m ()
bumpOversize (OsvIngest -> IngestCounter
ingestCounter OsvIngest
ingest)
        Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
WarningS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"Dropping oversized OSV entry " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> ZipEntry -> Text
zipEntryNameText ZipEntry
entry Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
": " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show Int
seen Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" bytes exceeds the " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show Int
cap Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"-byte per-advisory cap"))
    EntryBytes ByteString
fileBytes -> case ByteString -> Maybe OsvAdvisory
forall a. FromJSON a => ByteString -> Maybe a
decodeStrict ByteString
fileBytes :: Maybe OsvAdvisory of
        Maybe OsvAdvisory
Nothing -> m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall (m :: * -> *) a.
Monad m =>
m a -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m a
forall (t :: (* -> *) -> * -> *) (m :: * -> *) a.
(MonadTrans t, Monad m) =>
m a -> t m a
lift (m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ())
-> m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall a b. (a -> b) -> a -> b
$ do
            IngestCounter -> m ()
forall (m :: * -> *). MonadIO m => IngestCounter -> m ()
bumpMalformed (OsvIngest -> IngestCounter
ingestCounter OsvIngest
ingest)
            Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
WarningS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"Failed to parse OSV advisory JSON from entry: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> ZipEntry -> Text
zipEntryNameText ZipEntry
entry))
        Just OsvAdvisory
adv -> OsvIngest
-> OsvAdvisory
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall (m :: * -> *).
KatipContext m =>
OsvIngest
-> OsvAdvisory
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
admitAdvisory OsvIngest
ingest OsvAdvisory
adv
  where
    cap :: Int
cap = IngestLimits -> Int
ilMaxAdvisoryBytes (OsvIngest -> IngestLimits
ingestLimits OsvIngest
ingest)

-- Emit every row one decoded advisory yields. Nothing here filters: a row whose bound the
-- grammar cannot order is counted and logged, and still reaches the artifact.
admitAdvisory :: (KatipContext m) => OsvIngest -> OsvAdvisory -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
admitAdvisory :: forall (m :: * -> *).
KatipContext m =>
OsvIngest
-> OsvAdvisory
-> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
admitAdvisory OsvIngest
ingest OsvAdvisory
adv = do
    m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall (m :: * -> *) a.
Monad m =>
m a -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m a
forall (t :: (* -> *) -> * -> *) (m :: * -> *) a.
(MonadTrans t, Monad m) =>
m a -> t m a
lift (m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ())
-> m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall a b. (a -> b) -> a -> b
$ IngestCounter -> m ()
forall (m :: * -> *). MonadIO m => IngestCounter -> m ()
bumpAccepted (OsvIngest -> IngestCounter
ingestCounter OsvIngest
ingest)
    m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall (m :: * -> *) a.
Monad m =>
m a -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m a
forall (t :: (* -> *) -> * -> *) (m :: * -> *) a.
(MonadTrans t, Monad m) =>
m a -> t m a
lift (m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ())
-> m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall a b. (a -> b) -> a -> b
$ OsvIngest -> OsvAdvisory -> m ()
forall (m :: * -> *). MonadIO m => OsvIngest -> OsvAdvisory -> m ()
recordModified OsvIngest
ingest OsvAdvisory
adv
    m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall (m :: * -> *) a.
Monad m =>
m a -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m a
forall (t :: (* -> *) -> * -> *) (m :: * -> *) a.
(MonadTrans t, Monad m) =>
m a -> t m a
lift (m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ())
-> m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall a b. (a -> b) -> a -> b
$ OsvIngest -> OsvAdvisory -> [ExtractedOsv] -> m ()
forall (m :: * -> *).
KatipContext m =>
OsvIngest -> OsvAdvisory -> [ExtractedOsv] -> m ()
warnOnFanOut OsvIngest
ingest OsvAdvisory
adv [ExtractedOsv]
extracted
    m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall (m :: * -> *) a.
Monad m =>
m a -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m a
forall (t :: (* -> *) -> * -> *) (m :: * -> *) a.
(MonadTrans t, Monad m) =>
m a -> t m a
lift (m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ())
-> m () -> ConduitT (Either ZipEntry ByteString) ExtractedOsv m ()
forall a b. (a -> b) -> a -> b
$ Maybe (NonEmpty (Text, Text))
-> (NonEmpty (Text, Text) -> m ()) -> m ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ ([(Text, Text)] -> Maybe (NonEmpty (Text, Text))
forall a. [a] -> Maybe (NonEmpty a)
nonEmpty [(Text, Text)]
unorderable) (OsvIngest -> OsvAdvisory -> NonEmpty (Text, Text) -> m ()
forall (m :: * -> *).
KatipContext m =>
OsvIngest -> OsvAdvisory -> NonEmpty (Text, Text) -> m ()
warnOnUnorderable OsvIngest
ingest OsvAdvisory
adv)
    [ExtractedOsv]
-> ConduitT
     (Either ZipEntry ByteString) (Element [ExtractedOsv]) m ()
forall (m :: * -> *) mono i.
(Monad m, MonoFoldable mono) =>
mono -> ConduitT i (Element mono) m ()
yieldMany [ExtractedOsv]
extracted
  where
    extracted :: [ExtractedOsv]
extracted = EpssScores -> OsvAdvisory -> [ExtractedOsv]
extractFromAdvisory (OsvIngest -> EpssScores
ingestEpss OsvIngest
ingest) OsvAdvisory
adv
    -- A name this build does not serve has no grammar to judge its bounds by, so nothing
    -- about it is anomalous.
    unorderable :: [(Text, Text)]
unorderable = [(Text, Text)]
-> (Ecosystem -> [(Text, Text)])
-> Maybe Ecosystem
-> [(Text, Text)]
forall b a. b -> (a -> b) -> Maybe a -> b
maybe [] (\Ecosystem
eco -> (ExtractedOsv -> Maybe (Text, Text))
-> [ExtractedOsv] -> [(Text, Text)]
forall a b. (a -> Maybe b) -> [a] -> [b]
mapMaybe (Ecosystem -> ExtractedOsv -> Maybe (Text, Text)
unorderableExample Ecosystem
eco) [ExtractedOsv]
extracted) (OsvEcosystem -> Maybe Ecosystem
osvEcosystemTag (OsvIngest -> OsvEcosystem
ingestEcosystem OsvIngest
ingest))

unorderableExample :: Ecosystem -> ExtractedOsv -> Maybe (Text, Text)
unorderableExample :: Ecosystem -> ExtractedOsv -> Maybe (Text, Text)
unorderableExample Ecosystem
eco ExtractedOsv
row
    | Ecosystem -> ExtractedOsv -> Bool
orderableBounds Ecosystem
eco ExtractedOsv
row = Maybe (Text, Text)
forall a. Maybe a
Nothing
    | Bool
otherwise = (,) (ExtractedOsv -> Text
extPackage ExtractedOsv
row) (Text -> (Text, Text)) -> Maybe Text -> Maybe (Text, Text)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [Text] -> Maybe Text
forall a. [a] -> Maybe a
listToMaybe (Ecosystem -> ExtractedOsv -> [Text]
unorderableBounds Ecosystem
eco ExtractedOsv
row)

-- One line per advisory carrying one example, so a feed naming thousands of packages cannot
-- flood the log. The rows are kept, so this is an alarm and never a refusal.
warnOnUnorderable :: (KatipContext m) => OsvIngest -> OsvAdvisory -> NonEmpty (Text, Text) -> m ()
warnOnUnorderable :: forall (m :: * -> *).
KatipContext m =>
OsvIngest -> OsvAdvisory -> NonEmpty (Text, Text) -> m ()
warnOnUnorderable OsvIngest
ingest OsvAdvisory
adv unorderable :: NonEmpty (Text, Text)
unorderable@((Text
pkg, Text
bound) :| [(Text, Text)]
_) = do
    IngestCounter -> Int -> m ()
forall (m :: * -> *). MonadIO m => IngestCounter -> Int -> m ()
bumpUnorderable (OsvIngest -> IngestCounter
ingestCounter OsvIngest
ingest) (NonEmpty (Text, Text) -> Int
forall a. NonEmpty a -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length NonEmpty (Text, Text)
unorderable)
    Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
WarningS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"OSV advisory " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> OsvAdvisory -> Text
osvId OsvAdvisory
adv Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" carries " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show (NonEmpty (Text, Text) -> Int
forall a. NonEmpty a -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length NonEmpty (Text, Text)
unorderable) Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" range(s) the version grammar cannot order, for example " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
pkg Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
bound Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"; keeping them"))

warnOnFanOut :: (KatipContext m) => OsvIngest -> OsvAdvisory -> [ExtractedOsv] -> m ()
warnOnFanOut :: forall (m :: * -> *).
KatipContext m =>
OsvIngest -> OsvAdvisory -> [ExtractedOsv] -> m ()
warnOnFanOut OsvIngest
ingest OsvAdvisory
adv [ExtractedOsv]
extracted =
    Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (Int
n Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
> Int
limit) (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$
        Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
WarningS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"OSV advisory " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> OsvAdvisory -> Text
osvId OsvAdvisory
adv Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" expanded into " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show Int
n Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" ranges, exceeding the sanity threshold of " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show Int
limit Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"; ingesting it regardless"))
  where
    n :: Int
n = [ExtractedOsv] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [ExtractedOsv]
extracted
    limit :: Int
limit = OsvEcosystem -> Int
osvMaxAdvisoryFanOut (OsvIngest -> OsvEcosystem
ingestEcosystem OsvIngest
ingest)

-- The export's own @Last-Modified@, from the response that carried the rows.
recordResponseDate :: (MonadIO m) => OsvIngest -> [ByteString] -> m ()
recordResponseDate :: forall (m :: * -> *).
MonadIO m =>
OsvIngest -> [ByteString] -> m ()
recordResponseDate OsvIngest
ingest [ByteString]
headers =
    IORef OsvAttempt -> (OsvAttempt -> OsvAttempt) -> m ()
forall (m :: * -> *) a. MonadIO m => IORef a -> (a -> a) -> m ()
modifyIORef' (OsvIngest -> IORef OsvAttempt
ingestAttempt OsvIngest
ingest) ((OsvAttempt -> OsvAttempt) -> m ())
-> (OsvAttempt -> OsvAttempt) -> m ()
forall a b. (a -> b) -> a -> b
$ \OsvAttempt
attempt ->
        OsvAttempt
attempt{oaLastModified = lastModifiedOf headers}

-- A date no grammar reads, and a date the source cannot know yet, are both counted and both
-- dropped from the reading, never clamped. The record's rows are kept either way.
recordModified :: (MonadIO m) => OsvIngest -> OsvAdvisory -> m ()
recordModified :: forall (m :: * -> *). MonadIO m => OsvIngest -> OsvAdvisory -> m ()
recordModified OsvIngest
ingest OsvAdvisory
adv = Maybe Text -> (Text -> m ()) -> m ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
t a -> (a -> f b) -> f ()
for_ (OsvAdvisory -> Maybe Text
osvModified OsvAdvisory
adv) ((Text -> m ()) -> m ()) -> (Text -> m ()) -> m ()
forall a b. (a -> b) -> a -> b
$ \Text
written ->
    case UTCTime -> Text -> Maybe UTCTime
usableStamp (OsvIngest -> UTCTime
ingestNow OsvIngest
ingest) Text
written of
        Maybe UTCTime
Nothing -> IngestCounter -> m ()
forall (m :: * -> *). MonadIO m => IngestCounter -> m ()
bumpUnusableModified (OsvIngest -> IngestCounter
ingestCounter OsvIngest
ingest)
        Just UTCTime
stamp ->
            IORef OsvAttempt -> (OsvAttempt -> OsvAttempt) -> m ()
forall (m :: * -> *) a. MonadIO m => IORef a -> (a -> a) -> m ()
modifyIORef' (OsvIngest -> IORef OsvAttempt
ingestAttempt OsvIngest
ingest) ((OsvAttempt -> OsvAttempt) -> m ())
-> (OsvAttempt -> OsvAttempt) -> m ()
forall a b. (a -> b) -> a -> b
$ \OsvAttempt
attempt ->
                OsvAttempt
attempt{oaNewestModified = max (Just stamp) (oaNewestModified attempt)}

-- The instant a record's written date names, when the run can use it at all.
usableStamp :: UTCTime -> Text -> Maybe UTCTime
usableStamp :: UTCTime -> Text -> Maybe UTCTime
usableStamp UTCTime
now Text
written = do
    stamp <- Text -> Maybe UTCTime
parseSourceTime Text
written
    guard (stamp <= now)
    pure stamp

bumpAccepted :: (MonadIO m) => IngestCounter -> m ()
bumpAccepted :: forall (m :: * -> *). MonadIO m => IngestCounter -> m ()
bumpAccepted (IngestCounter IORef IngestStats
ref) = IORef IngestStats -> (IngestStats -> IngestStats) -> m ()
forall (m :: * -> *) a. MonadIO m => IORef a -> (a -> a) -> m ()
modifyIORef' IORef IngestStats
ref (\IngestStats
s -> IngestStats
s{statAccepted = statAccepted s + 1})

bumpOversize :: (MonadIO m) => IngestCounter -> m ()
bumpOversize :: forall (m :: * -> *). MonadIO m => IngestCounter -> m ()
bumpOversize (IngestCounter IORef IngestStats
ref) = IORef IngestStats -> (IngestStats -> IngestStats) -> m ()
forall (m :: * -> *) a. MonadIO m => IORef a -> (a -> a) -> m ()
modifyIORef' IORef IngestStats
ref (\IngestStats
s -> IngestStats
s{statDroppedOversize = statDroppedOversize s + 1})

bumpMalformed :: (MonadIO m) => IngestCounter -> m ()
bumpMalformed :: forall (m :: * -> *). MonadIO m => IngestCounter -> m ()
bumpMalformed (IngestCounter IORef IngestStats
ref) = IORef IngestStats -> (IngestStats -> IngestStats) -> m ()
forall (m :: * -> *) a. MonadIO m => IORef a -> (a -> a) -> m ()
modifyIORef' IORef IngestStats
ref (\IngestStats
s -> IngestStats
s{statDroppedMalformed = statDroppedMalformed s + 1})

bumpUnorderable :: (MonadIO m) => IngestCounter -> Int -> m ()
bumpUnorderable :: forall (m :: * -> *). MonadIO m => IngestCounter -> Int -> m ()
bumpUnorderable (IngestCounter IORef IngestStats
ref) Int
n = IORef IngestStats -> (IngestStats -> IngestStats) -> m ()
forall (m :: * -> *) a. MonadIO m => IORef a -> (a -> a) -> m ()
modifyIORef' IORef IngestStats
ref (\IngestStats
s -> IngestStats
s{statUnorderable = statUnorderable s + n})

bumpUnusableModified :: (MonadIO m) => IngestCounter -> m ()
bumpUnusableModified :: forall (m :: * -> *). MonadIO m => IngestCounter -> m ()
bumpUnusableModified (IngestCounter IORef IngestStats
ref) = IORef IngestStats -> (IngestStats -> IngestStats) -> m ()
forall (m :: * -> *) a. MonadIO m => IORef a -> (a -> a) -> m ()
modifyIORef' IORef IngestStats
ref (\IngestStats
s -> IngestStats
s{statUnusableModified = statUnusableModified s + 1})

zipEntryNameText :: ZipEntry -> Text
zipEntryNameText :: ZipEntry -> Text
zipEntryNameText ZipEntry
entry = case ZipEntry -> Either Text ByteString
zipEntryName ZipEntry
entry of
    Left Text
txt -> Text
txt
    Right ByteString
bs -> OnDecodeError -> ByteString -> Text
decodeUtf8With OnDecodeError
lenientDecode ByteString
bs

-- Checks @cap@ before each chunk, so memory never exceeds the cap plus one chunk. It is also the
-- only depth guard: 'decodeStrict' materialises the whole value before any post-decode check.
collectFile :: (Monad m) => Int -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
collectFile :: forall (m :: * -> *) o.
Monad m =>
Int -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
collectFile Int
cap = Int
-> [ByteString]
-> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
forall {m :: * -> *} {o}.
Monad m =>
Int
-> [ByteString]
-> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
go Int
0 []
  where
    go :: Int
-> [ByteString]
-> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
go !Int
seen [ByteString]
acc =
        ConduitT
  (Either ZipEntry ByteString)
  o
  m
  (Maybe (Either ZipEntry ByteString))
forall (m :: * -> *) i o. Monad m => ConduitT i o m (Maybe i)
await ConduitT
  (Either ZipEntry ByteString)
  o
  m
  (Maybe (Either ZipEntry ByteString))
-> (Maybe (Either ZipEntry ByteString)
    -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome)
-> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
forall a b.
ConduitT (Either ZipEntry ByteString) o m a
-> (a -> ConduitT (Either ZipEntry ByteString) o m b)
-> ConduitT (Either ZipEntry ByteString) o m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
            Maybe (Either ZipEntry ByteString)
Nothing -> EntryOutcome
-> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
forall a. a -> ConduitT (Either ZipEntry ByteString) o m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (ByteString -> EntryOutcome
EntryBytes ([ByteString] -> ByteString
BS.concat ([ByteString] -> [ByteString]
forall a. [a] -> [a]
reverse [ByteString]
acc)))
            Just (Left ZipEntry
entry) -> do
                Either ZipEntry ByteString
-> ConduitT (Either ZipEntry ByteString) o m ()
forall i o (m :: * -> *). i -> ConduitT i o m ()
leftover (ZipEntry -> Either ZipEntry ByteString
forall a b. a -> Either a b
Left ZipEntry
entry)
                EntryOutcome
-> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
forall a. a -> ConduitT (Either ZipEntry ByteString) o m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (ByteString -> EntryOutcome
EntryBytes ([ByteString] -> ByteString
BS.concat ([ByteString] -> [ByteString]
forall a. [a] -> [a]
reverse [ByteString]
acc)))
            Just (Right ByteString
bs) ->
                let seen' :: Int
seen' = Int
seen Int -> Int -> Int
forall a. Num a => a -> a -> a
+ ByteString -> Int
BS.length ByteString
bs
                 in if Int
seen' Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
> Int
cap
                        then Int -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
forall (m :: * -> *) o.
Monad m =>
Int -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
drainOversize Int
seen'
                        else Int
-> [ByteString]
-> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
go Int
seen' (ByteString
bs ByteString -> [ByteString] -> [ByteString]
forall a. a -> [a] -> [a]
: [ByteString]
acc)

-- Carrying no accumulator frees the collected prefix, so the drain to the next entry
-- boundary retains only the running size.
drainOversize :: (Monad m) => Int -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
drainOversize :: forall (m :: * -> *) o.
Monad m =>
Int -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
drainOversize !Int
seen =
    ConduitT
  (Either ZipEntry ByteString)
  o
  m
  (Maybe (Either ZipEntry ByteString))
forall (m :: * -> *) i o. Monad m => ConduitT i o m (Maybe i)
await ConduitT
  (Either ZipEntry ByteString)
  o
  m
  (Maybe (Either ZipEntry ByteString))
-> (Maybe (Either ZipEntry ByteString)
    -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome)
-> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
forall a b.
ConduitT (Either ZipEntry ByteString) o m a
-> (a -> ConduitT (Either ZipEntry ByteString) o m b)
-> ConduitT (Either ZipEntry ByteString) o m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
        Maybe (Either ZipEntry ByteString)
Nothing -> EntryOutcome
-> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
forall a. a -> ConduitT (Either ZipEntry ByteString) o m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Int -> EntryOutcome
EntryOversize Int
seen)
        Just (Left ZipEntry
entry) -> do
            Either ZipEntry ByteString
-> ConduitT (Either ZipEntry ByteString) o m ()
forall i o (m :: * -> *). i -> ConduitT i o m ()
leftover (ZipEntry -> Either ZipEntry ByteString
forall a b. a -> Either a b
Left ZipEntry
entry)
            EntryOutcome
-> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
forall a. a -> ConduitT (Either ZipEntry ByteString) o m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Int -> EntryOutcome
EntryOversize Int
seen)
        Just (Right ByteString
bs) -> Int -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
forall (m :: * -> *) o.
Monad m =>
Int -> ConduitT (Either ZipEntry ByteString) o m EntryOutcome
drainOversize (Int
seen Int -> Int -> Int
forall a. Num a => a -> a -> a
+ ByteString -> Int
BS.length ByteString
bs)