-- SPDX-FileCopyrightText: 2026 Alexandra de Wit
--
-- SPDX-License-Identifier: MIT
{-# LANGUAGE BangPatterns #-}
{-# LANGUAGE LambdaCase #-}
{-# LANGUAGE OverloadedStrings #-}

{- | Streaming ingest of an osv.dev export archive, bounded against a pathological
or tampered payload.

Pilot fetches @\<base\>\/\<ecosystem\>\/all.zip@ from the public osv.dev mirror and
decodes each advisory JSON on the way to compiling @osv.db@. osv.dev is trusted in
normal operation, but the feed is an /aggregation/ of many upstream databases: a
single poisoned record can ride in with every transport header honest, so ingest is
bounded as defence-in-depth. The bounds are __generous__ (a real advisory is
kilobytes) and __fail-soft per entry__: one bad advisory is dropped and logged, and
the rest of the archive keeps flowing.

Two levels of response, tallied in 'IngestStats':

* A single over-large or malformed entry is dropped and counted. The 'ilMaxAdvisoryBytes'
  cap is enforced /before/ the bytes are retained and /before/ the JSON is decoded, so
  an inflation bomb never reaches the decoder whole; the offending entry is drained to
  its boundary so the following entries stay aligned.
* An advisory that expands into more than 'ilMaxAdvisoryFanOut' ranges is logged as
  anomalous but still ingested (log-only, non-gating).

The aggregate verdict is a separate, pure decision ('systemicDrop'): the compiler reads
the tally once the stream completes and, if drops are /systemic/ rather than isolated,
abandons the run ('PilotIngestAborted') so a consumer keeps its last-good artifact
instead of adopting a hole-ridden one.

Depth is not guarded here on the decoded value: 'Data.Aeson.decodeStrict' materialises
the whole intermediate value before any post-decode check could run, and 'OsvAdvisory'
cannot represent unbounded nesting anyway, so the byte cap (which holds parse cost to a
constant multiple of the input) plus the process heap ceiling resolved at boot
("Ecluse.Rts") are what bound a small-but-deep payload.
-}
module Ecluse.Core.Osv.Stream (
    streamOsvUrl,
    parseOsvStream,

    -- * Ingest bounds and drop accounting
    IngestLimits (..),
    defaultIngestLimits,
    IngestStats (..),
    IngestCounter,
    OsvIngest (..),
    newOsvIngest,
    readIngestStats,
    resetIngestStats,
    systemicDrop,
    PilotIngestAborted (..),
) 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 Katip (KatipContext, Severity (..), logFM, ls)
import Network.HTTP.Simple (getResponseBody, httpSource, parseRequest, setRequestCheckStatus)
import OpenTelemetry.Context qualified as Ctx
import OpenTelemetry.Trace.Core (SpanKind (Internal), TracerProvider, addAttribute, createSpan, defaultSpanArguments, endSpan, kind, makeTracer, tracerOptions)

import Ecluse.Core.Osv.Advisory (ExtractedOsv, OsvAdvisory, extractFromAdvisory, osvId)

{- | The tunable per-advisory ingest bounds. Generous by design: osv.dev is trusted
in normal operation, so these only backstop a pathological or tampered payload and
must never trip on a real, if large, advisory.
-}
data IngestLimits = IngestLimits
    { IngestLimits -> Int
ilMaxAdvisoryBytes :: !Int
    {- ^ Largest decompressed advisory JSON, in bytes, accepted from one zip entry
    before it is dropped. Bounds memory and, transitively, decode cost.
    -}
    , IngestLimits -> Int
ilMaxAdvisoryFanOut :: !Int
    {- ^ Number of extracted ranges one advisory may expand into before it is flagged
    as anomalous. Log-only: the advisory is still ingested.
    -}
    }
    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)

{- | Sane defaults for 'IngestLimits': an 8 MiB per-advisory ceiling (a real advisory
is kilobytes, so this is generous headroom) and a 256-range fan-out flag (a real
advisory expands into a small multiple of its affected packages, far below this).
-}
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
        , ilMaxAdvisoryFanOut :: Int
ilMaxAdvisoryFanOut = Int
256
        }

{- | The running tally of one ingest pass: advisories accepted, and entries dropped
by reason. Read 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.
    }
    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 -> IngestStats
IngestStats Int
0 Int
0 Int
0

-- | The mutable drop tally for one ingest pass. Opaque; read with 'readIngestStats'.
newtype IngestCounter = IngestCounter (IORef IngestStats)

{- | The context one ingest pass threads through the stream: its bounds and the live
drop tally it records into.
-}
data OsvIngest = OsvIngest
    { OsvIngest -> IngestLimits
ingestLimits :: IngestLimits
    , OsvIngest -> IngestCounter
ingestCounter :: IngestCounter
    }

-- | A fresh ingest context with the given bounds and a zeroed tally.
newOsvIngest :: (MonadIO m) => IngestLimits -> m OsvIngest
newOsvIngest :: forall (m :: * -> *). MonadIO m => IngestLimits -> m OsvIngest
newOsvIngest IngestLimits
limits = IngestLimits -> IngestCounter -> OsvIngest
OsvIngest IngestLimits
limits (IngestCounter -> OsvIngest)
-> (IORef IngestStats -> IngestCounter)
-> IORef IngestStats
-> OsvIngest
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IORef IngestStats -> IngestCounter
IngestCounter (IORef IngestStats -> OsvIngest)
-> m (IORef IngestStats) -> m OsvIngest
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

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

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

{- | Whether a run's drop tally signals /systemic/ corruption (a hostile or broken
feed) rather than a few poisoned records, so its artifact must not be published. Trips
only when drops are both absolutely non-trivial and a large fraction of all entries, so
a handful of bad advisories in a healthy feed never blocks a build, while a feed that is
mostly unusable does.
-}
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

{- | Fewest drops that can count as systemic, so a tiny feed with a couple of bad
entries is never judged corrupt.
-}
systemicDropFloor :: Int
systemicDropFloor :: Int
systemicDropFloor = Int
16

{- | The fraction of all entries (as a percent) that must be dropped to judge a feed
systemically corrupt.
-}
systemicDropPercent :: Int
systemicDropPercent :: Int
systemicDropPercent = Int
10

{- | Raised after a compile pass whose drop tally 'systemicDrop' judged systemic: the
run is abandoned without publishing, so a consumer keeps its last-good artifact rather
than adopting a hole-ridden one.
-}
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 (String -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (String
"Initializing OSV stream from URL: " String -> ShowS
forall a. Semigroup a => a -> a -> a
<> String
urlStr))
    let mTracer :: Maybe Tracer
mTracer = (\TracerProvider
tp -> TracerProvider -> InstrumentationLibrary -> TracerOptions -> Tracer
makeTracer TracerProvider
tp InstrumentationLibrary
"ecluse" TracerOptions
tracerOptions) (TracerProvider -> Tracer) -> Maybe TracerProvider -> Maybe Tracer
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe TracerProvider
mTracerProvider
    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
        ((Tracer -> IO Span) -> Maybe Tracer -> IO (Maybe Span)
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 (\Tracer
t -> Tracer -> Context -> Text -> SpanArguments -> IO Span
forall (m :: * -> *).
(MonadIO m, HasCallStack) =>
Tracer -> Context -> Text -> SpanArguments -> m Span
createSpan Tracer
t Context
Ctx.empty Text
"ecluse.pilot.osv.stream" SpanArguments
defaultSpanArguments{kind = Internal}) Maybe Tracer
mTracer)
        ((Span -> IO ()) -> Maybe Span -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
(a -> m b) -> t a -> m ()
mapM_ (Span -> Maybe Timestamp -> IO ()
forall (m :: * -> *). MonadIO m => Span -> Maybe Timestamp -> m ()
`endSpan` Maybe Timestamp
forall a. Maybe a
Nothing))
        ( \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_url" (String -> Text
forall a. ToText a => a -> Text
toText String
urlStr)
            -- 'setRequestCheckStatus' makes a non-2xx response throw a
            -- 'StatusCodeException' at the header boundary. This is deliberate: it
            -- lets the backoff wrapper (see 'Ecluse.Core.Osv.Retry') see a 502
            -- from osv.dev as a retryable fault, rather than streaming the error
            -- page into the unzip parser where it would surface as a parse error a
            -- retry could not fix.
            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 -> 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))
    let mTracer :: Maybe Tracer
mTracer = (\TracerProvider
tp -> TracerProvider -> InstrumentationLibrary -> TracerOptions -> Tracer
makeTracer TracerProvider
tp InstrumentationLibrary
"ecluse" TracerOptions
tracerOptions) (TracerProvider -> Tracer) -> Maybe TracerProvider -> Maybe Tracer
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe TracerProvider
mTracerProvider
    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
        ((Tracer -> IO Span) -> Maybe Tracer -> IO (Maybe Span)
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 (\Tracer
t -> Tracer -> Context -> Text -> SpanArguments -> IO Span
forall (m :: * -> *).
(MonadIO m, HasCallStack) =>
Tracer -> Context -> Text -> SpanArguments -> m Span
createSpan Tracer
t Context
Ctx.empty Text
"ecluse.pilot.osv.parse" SpanArguments
defaultSpanArguments{kind = Internal}) Maybe Tracer
mTracer)
        ((Span -> IO ()) -> Maybe Span -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
(a -> m b) -> t a -> m ()
mapM_ (Span -> Maybe Timestamp -> IO ()
forall (m :: * -> *). MonadIO m => Span -> Maybe Timestamp -> m ()
`endSpan` Maybe Timestamp
forall a. Maybe a
Nothing))
        (\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 (carrying the entry's full decompressed size, for the log).
data EntryOutcome = EntryBytes !ByteString | EntryOversize !Int

-- Decide what one collected entry yields: drop-and-count an over-large or malformed
-- entry, or count and emit a decoded advisory's ranges (flagging an anomalous fan-out).
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
ErrorS (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 -> 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)
            let extracted :: [ExtractedOsv]
extracted = OsvAdvisory -> [ExtractedOsv]
extractFromAdvisory 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
            [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
    cap :: Int
cap = IngestLimits -> Int
ilMaxAdvisoryBytes (OsvIngest -> IngestLimits
ingestLimits OsvIngest
ingest)

-- Flag an advisory that expands into an anomalous number of ranges. Log-only: the
-- advisory is ingested regardless, this is a "something is odd with this record" signal.
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
ErrorS (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 = IngestLimits -> Int
ilMaxAdvisoryFanOut (OsvIngest -> IngestLimits
ingestLimits OsvIngest
ingest)

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})

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

{- | Accumulate one zip entry's decompressed bytes, up to @cap@. The cap is checked
before each chunk is retained, so memory never exceeds the cap plus one chunk; on breach
the remaining chunks of the entry are drained (not retained) to the next entry boundary,
so the following entries stay aligned, and the entry is reported as 'EntryOversize'.
-}
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 :: * -> *} {a} {o}.
Monad m =>
Int
-> [ByteString] -> ConduitT (Either a ByteString) o m EntryOutcome
go Int
0 []
  where
    go :: Int
-> [ByteString] -> ConduitT (Either a ByteString) o m EntryOutcome
go !Int
seen [ByteString]
acc =
        ConduitT (Either a ByteString) o m (Maybe (Either a ByteString))
forall (m :: * -> *) i o. Monad m => ConduitT i o m (Maybe i)
await ConduitT (Either a ByteString) o m (Maybe (Either a ByteString))
-> (Maybe (Either a ByteString)
    -> ConduitT (Either a ByteString) o m EntryOutcome)
-> ConduitT (Either a ByteString) o m EntryOutcome
forall a b.
ConduitT (Either a ByteString) o m a
-> (a -> ConduitT (Either a ByteString) o m b)
-> ConduitT (Either a ByteString) o m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
            Maybe (Either a ByteString)
Nothing -> EntryOutcome -> ConduitT (Either a ByteString) o m EntryOutcome
forall a. a -> ConduitT (Either a 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 a
entry) -> do
                Either a ByteString -> ConduitT (Either a ByteString) o m ()
forall i o (m :: * -> *). i -> ConduitT i o m ()
leftover (a -> Either a ByteString
forall a b. a -> Either a b
Left a
entry)
                EntryOutcome -> ConduitT (Either a ByteString) o m EntryOutcome
forall a. a -> ConduitT (Either a 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 a ByteString) o m EntryOutcome
forall {m :: * -> *} {a} {o}.
Monad m =>
Int -> ConduitT (Either a ByteString) o m EntryOutcome
drainOversize Int
seen'
                        else Int
-> [ByteString] -> ConduitT (Either a ByteString) o m EntryOutcome
go Int
seen' (ByteString
bs ByteString -> [ByteString] -> [ByteString]
forall a. a -> [a] -> [a]
: [ByteString]
acc)
    -- The cap is already breached: keep pulling this entry's chunks to advance the
    -- stream to the next boundary, but retain none of them (only their size, for the
    -- log). The accumulated prefix is dropped with this frame.
    drainOversize :: Int -> ConduitT (Either a ByteString) o m EntryOutcome
drainOversize !Int
seen =
        ConduitT (Either a ByteString) o m (Maybe (Either a ByteString))
forall (m :: * -> *) i o. Monad m => ConduitT i o m (Maybe i)
await ConduitT (Either a ByteString) o m (Maybe (Either a ByteString))
-> (Maybe (Either a ByteString)
    -> ConduitT (Either a ByteString) o m EntryOutcome)
-> ConduitT (Either a ByteString) o m EntryOutcome
forall a b.
ConduitT (Either a ByteString) o m a
-> (a -> ConduitT (Either a ByteString) o m b)
-> ConduitT (Either a ByteString) o m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
            Maybe (Either a ByteString)
Nothing -> EntryOutcome -> ConduitT (Either a ByteString) o m EntryOutcome
forall a. a -> ConduitT (Either a ByteString) o m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Int -> EntryOutcome
EntryOversize Int
seen)
            Just (Left a
entry) -> do
                Either a ByteString -> ConduitT (Either a ByteString) o m ()
forall i o (m :: * -> *). i -> ConduitT i o m ()
leftover (a -> Either a ByteString
forall a b. a -> Either a b
Left a
entry)
                EntryOutcome -> ConduitT (Either a ByteString) o m EntryOutcome
forall a. a -> ConduitT (Either a 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 a ByteString) o m EntryOutcome
drainOversize (Int
seen Int -> Int -> Int
forall a. Num a => a -> a -> a
+ ByteString -> Int
BS.length ByteString
bs)