packages feed

freckle-kafka (empty) → 0.0.0.0

raw patch · 8 files changed

+585/−0 lines, 8 filesdep +Blammodep +aesondep +annotated-exception

Dependencies added: Blammo, aeson, annotated-exception, base, containers, freckle-env, hs-opentelemetry-sdk, hw-kafka-client, lens, relude, resource-pool, text, time, unliftio, unordered-containers, yesod-core

Files

+ CHANGELOG.md view
@@ -0,0 +1,14 @@+## [_Unreleased_](https://github.com/freckle/freckle-app/compare/freckle-kafka-v0.0.0.0...main)++## [v0.0.0.0](https://github.com/freckle/freckle-app/tree/freckle-kafka-v0.0.0.0/freckle-kafka)++First release, sprouted from `freckle-app-1.18.2.0`.++Changes from `freckle-app`:++- `produceKeyedOnAsync` has been removed; you may substitute the definition+  `(\prTopic values -> void . async . produceKeyedOn prTopic values)`.++- `runConsumer` has been altered; you may recover the original behavior by+  changing `runConsumer pollTimeout onMessage` to+  `withTraceContext $ immortalCreateLogged $ runConsumer pollTimeout onMessage`.
+ LICENSE view
@@ -0,0 +1,21 @@+The MIT License (MIT)++Copyright (c) 2023-2024 Renaissance Learning Inc++Permission is hereby granted, free of charge, to any person obtaining a copy+of this software and associated documentation files (the "Software"), to deal+in the Software without restriction, including without limitation the rights+to use, copy, modify, merge, publish, distribute, sublicense, and/or sell+copies of the Software, and to permit persons to whom the Software is+furnished to do so, subject to the following conditions:++The above copyright notice and this permission notice shall be included in all+copies or substantial portions of the Software.++THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR+IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,+FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE+AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER+LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,+OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE+SOFTWARE.
+ README.md view
@@ -0,0 +1,7 @@+# freckle-kafka++Some extensions to the `hw-kafka-client` library.++---++[CHANGELOG](./CHANGELOG.md) | [LICENSE](./LICENSE)
+ freckle-kafka.cabal view
@@ -0,0 +1,64 @@+cabal-version:      1.18+name:               freckle-kafka+version:            0.0.0.0+license:            MIT+license-file:       LICENSE+maintainer:         Freckle Education+homepage:           https://github.com/freckle/freckle-app#readme+bug-reports:        https://github.com/freckle/freckle-app/issues+synopsis:           Some extensions to the hw-kafka-client library+description:        Please see README.md+category:           Database+build-type:         Simple+extra-source-files: package.yaml+extra-doc-files:+    README.md+    CHANGELOG.md++source-repository head+    type:     git+    location: https://github.com/freckle/freckle-app++library+    exposed-modules:+        Freckle.App.Kafka+        Freckle.App.Kafka.Consumer+        Freckle.App.Kafka.Producer++    hs-source-dirs:     library+    other-modules:      Paths_freckle_kafka+    default-language:   GHC2021+    default-extensions:+        DataKinds DeriveAnyClass DerivingVia DerivingStrategies GADTs+        LambdaCase NoImplicitPrelude NoMonomorphismRestriction+        OverloadedStrings TypeFamilies++    ghc-options:+        -fignore-optim-changes -fwrite-ide-info -Weverything+        -Wno-all-missed-specialisations -Wno-missing-exported-signatures+        -Wno-missing-import-lists -Wno-missing-kind-signatures+        -Wno-missing-local-signatures -Wno-missing-safe-haskell-mode+        -Wno-monomorphism-restriction -Wno-prepositive-qualified-module+        -Wno-safe -Wno-unsafe++    build-depends:+        Blammo >=2.0.0.0,+        aeson >=2.0.3.0,+        annotated-exception >=0.2.0.4,+        base >=4.16.4.0 && <5,+        containers >=0.6.5.1,+        freckle-env >=0.0.0.0,+        hs-opentelemetry-sdk >=0.0.3.6,+        hw-kafka-client >=5.0.0,+        lens >=5.1.1,+        relude >=1.1.0.0,+        resource-pool >=0.4.0.0,+        text >=1.2.5.0,+        time >=1.11.1.1,+        unliftio >=0.2.25.0,+        unordered-containers >=0.2.19.1,+        yesod-core >=1.6.24.2++    if impl(ghc >=9.8)+        ghc-options:+            -Wno-missing-role-annotations -Wno-missing-poly-kind-signatures
+ library/Freckle/App/Kafka.hs view
@@ -0,0 +1,7 @@+module Freckle.App.Kafka+  ( module Freckle.App.Kafka.Consumer+  , module Freckle.App.Kafka.Producer+  ) where++import Freckle.App.Kafka.Consumer+import Freckle.App.Kafka.Producer
+ library/Freckle/App/Kafka/Consumer.hs view
@@ -0,0 +1,234 @@+{-# LANGUAGE ApplicativeDo #-}++module Freckle.App.Kafka.Consumer+  ( HasKafkaConsumer (..)+  , withKafkaConsumer+  , KafkaConsumerConfig (..)+  , envKafkaConsumerConfig+  , runConsumer+  ) where++import Relude++import Blammo.Logging+import Control.Exception.Annotated.UnliftIO (AnnotatedException)+import Control.Exception.Annotated.UnliftIO qualified as Annotated+import Control.Lens (Lens', view)+import Data.Aeson+import Data.List.NonEmpty qualified as NE+import Data.Map.Strict qualified as Map+import Data.Text qualified as T+import Freckle.App.Env (Timeout (..))+import Freckle.App.Env qualified as Env+import Freckle.App.Kafka.Producer (envKafkaBrokerAddresses)+import Kafka.Consumer hiding+  ( Timeout+  , closeConsumer+  , newConsumer+  , runConsumer+  , subscription+  )+import Kafka.Consumer qualified as Kafka+import OpenTelemetry.Trace (SpanKind (..), defaultSpanArguments)+import OpenTelemetry.Trace qualified as Trace+import OpenTelemetry.Trace.Monad (MonadTracer, inSpan)+import UnliftIO (MonadUnliftIO)+import UnliftIO.Exception (bracket)++data KafkaConsumerConfig = KafkaConsumerConfig+  { kafkaConsumerConfigBrokerAddresses :: NonEmpty BrokerAddress+  -- ^ The list of host/port pairs for establishing the initial connection+  -- to the Kafka cluster.+  --+  -- This is the `bootstrap.servers` Kafka consumer configuration property.+  , kafkaConsumerConfigGroupId :: ConsumerGroupId+  -- ^ The consumer group id to which the consumer belongs.+  --+  -- This is the `group.id` Kafka consumer configuration property.+  , kafkaConsumerConfigTopic :: TopicName+  -- ^ The topic name polled for messages by the Kafka consumer.+  , kafkaConsumerConfigOffsetReset :: OffsetReset+  -- ^ The offset reset parameter used when there is no initial offset in Kafka.+  --+  -- This is the `auto.offset.reset` Kafka consumer configuration property.+  , kafkaConsumerConfigAutoCommitInterval :: Millis+  -- ^ The interval that offsets are auto-committed to Kafka.+  --+  -- This sets the `auto.commit.interval.ms` and `enable.auto.commit` Kafka+  -- consumer configuration properties.+  , kafkaConsumerConfigExtraSubscriptionProps :: Map Text Text+  -- ^ Extra properties used to configure the Kafka consumer.+  }+  deriving stock (Show)++envKafkaTopic+  :: Env.Parser Env.Error TopicName+envKafkaTopic =+  Env.var+    (Env.eitherReader readKafkaTopic)+    "KAFKA_TOPIC"+    mempty++readKafkaTopic :: String -> Either String TopicName+readKafkaTopic t = case T.pack t of+  "" -> Left "Kafka topics cannot be empty"+  x -> Right $ TopicName x++envKafkaOffsetReset+  :: Env.Parser Env.Error OffsetReset+envKafkaOffsetReset =+  Env.var+    (Env.eitherReader readKafkaOffsetReset)+    "KAFKA_OFFSET_RESET"+    $ Env.def Earliest++readKafkaOffsetReset :: String -> Either String OffsetReset+readKafkaOffsetReset t = case T.pack t of+  "earliest" -> Right Earliest+  "latest" -> Right Latest+  _ -> Left "Kafka offset reset must be one of earliest or latest"++envKafkaConsumerConfig+  :: Env.Parser Env.Error KafkaConsumerConfig+envKafkaConsumerConfig = do+  brokerAddresses <- envKafkaBrokerAddresses+  consumerGroupId <- Env.var Env.nonempty "KAFKA_CONSUMER_GROUP_ID" mempty+  kafkaTopic <- envKafkaTopic+  kafkaOffsetReset <- envKafkaOffsetReset+  kafkaAutoOffsetInterval <-+    fmap (fromIntegral . timeoutMs)+      <$> Env.var Env.timeout "KAFKA_AUTO_COMMIT_INTERVAL"+      $ Env.def+      $ TimeoutMilliseconds 5000+  kafkaExtraProps <-+    Env.var+      (fmap Map.fromList . Env.keyValues)+      "KAFKA_EXTRA_SUBSCRIPTION_PROPS"+      (Env.def mempty)+  pure+    $ KafkaConsumerConfig+      brokerAddresses+      consumerGroupId+      kafkaTopic+      kafkaOffsetReset+      kafkaAutoOffsetInterval+      kafkaExtraProps++class HasKafkaConsumer env where+  kafkaConsumerL :: Lens' env KafkaConsumer++consumerProps :: KafkaConsumerConfig -> ConsumerProperties+consumerProps config =+  brokersList brokers+    <> groupId (kafkaConsumerConfigGroupId config)+    <> autoCommit (kafkaConsumerConfigAutoCommitInterval config)+    <> logLevel KafkaLogInfo+ where+  brokers = NE.toList $ kafkaConsumerConfigBrokerAddresses config++subscription :: KafkaConsumerConfig -> Subscription+subscription config =+  topics [kafkaConsumerConfigTopic config]+    <> offsetReset (kafkaConsumerConfigOffsetReset config)+    <> extraSubscriptionProps (kafkaConsumerConfigExtraSubscriptionProps config)++withKafkaConsumer+  :: (MonadUnliftIO m, HasCallStack)+  => KafkaConsumerConfig+  -> (KafkaConsumer -> m a)+  -> m a+withKafkaConsumer config = bracket newConsumer closeConsumer+ where+  (props, sub) = (consumerProps &&& subscription) config+  newConsumer = either Annotated.throw pure =<< Kafka.newConsumer props sub+  closeConsumer = maybe (pure ()) Annotated.throw <=< Kafka.closeConsumer++timeoutMs :: Timeout -> Int+timeoutMs = \case+  TimeoutSeconds s -> s * 1000+  TimeoutMilliseconds ms -> ms++data KafkaMessageDecodeError = KafkaMessageDecodeError+  { input :: ByteString+  , errors :: String+  }+  deriving stock (Show)++instance Exception KafkaMessageDecodeError where+  displayException e =+    mconcat+      [ "Unable to decode JSON"+      , "\n  input:  " <> decodeUtf8 (input e)+      , "\n  errors: " <> (errors e)+      ]++runConsumer+  :: ( MonadUnliftIO m+     , MonadReader env m+     , MonadLogger m+     , MonadTracer m+     , HasKafkaConsumer env+     , FromJSON a+     , HasCallStack+     )+  => Timeout+  -> (a -> m ())+  -> m ()+runConsumer pollTimeout onMessage =+  forever $ do+    consumer <- view kafkaConsumerL++    flip Annotated.catches handlers+      $ inSpan+        "kafka.consumer"+        (defaultSpanArguments {Trace.kind = Consumer})+      $ do+        mRecord <- fromKafkaError =<< pollMessage consumer kTimeout++        for_ (crValue =<< mRecord) $ \bs -> do+          a <-+            inSpan "kafka.consumer.message.decode" defaultSpanArguments+              $ either (Annotated.throw . KafkaMessageDecodeError bs) pure+              $ eitherDecodeStrict bs+          inSpan "kafka.consumer.message.handle" defaultSpanArguments $ onMessage a+ where+  kTimeout = Kafka.Timeout $ timeoutMs pollTimeout++  handlers =+    [ Annotated.Handler+        $ logErrorNS "kafka"+        . annotatedExceptionMessageFrom @KafkaError+          (const "Error polling for message from Kafka")+    , Annotated.Handler+        $ logErrorNS "kafka"+        . annotatedExceptionMessageFrom @KafkaMessageDecodeError+          (const "Could not decode message value")+    ]++-- | Like 'annotatedExceptionMessage', but use the supplied function to+--   construct an initial 'Message' that it will augment.+annotatedExceptionMessageFrom+  :: Exception ex => (ex -> Message) -> AnnotatedException ex -> Message+annotatedExceptionMessageFrom f ann = case f ex of+  msg :# series -> msg :# series <> ["error" .= errorObject]+ where+  ex = Annotated.exception ann+  errorObject =+    object+      [ "message" .= displayException ex+      , "stack" .= (prettyCallStack <$> Annotated.annotatedExceptionCallStack ann)+      ]++fromKafkaError+  :: (MonadIO m, MonadLogger m, HasCallStack)+  => Either KafkaError a+  -> m (Maybe a)+fromKafkaError =+  either+    ( \case+        KafkaResponseError RdKafkaRespErrTimedOut ->+          Nothing <$ logDebug "Polling timeout"+        err -> Annotated.throw err+    )+    $ pure+    . Just
+ library/Freckle/App/Kafka/Producer.hs view
@@ -0,0 +1,167 @@+{-# LANGUAGE ApplicativeDo #-}+{-# LANGUAGE NamedFieldPuns #-}++module Freckle.App.Kafka.Producer+  ( envKafkaBrokerAddresses+  , KafkaProducerPoolConfig (..)+  , envKafkaProducerPoolConfig+  , KafkaProducerPool (..)+  , HasKafkaProducerPool (..)+  , createKafkaProducerPool+  , produceKeyedOn+  ) where++import Relude++import Blammo.Logging+  ( Message ((:#))+  , MonadLogger+  , logDebugNS+  , logErrorNS+  , (.=)+  )+import Control.Exception.Annotated.UnliftIO qualified as Annotated+import Control.Lens (Lens', lens, view)+import Data.Aeson (ToJSON, encode)+import Data.HashMap.Strict qualified as HashMap+import Data.List.NonEmpty qualified as NE+import Data.Pool (Pool)+import Data.Pool qualified as Pool+import Data.Text qualified as T+import Data.Time (NominalDiffTime)+import Freckle.App.Env qualified as Env+import GHC.IO.Exception (userError)+import Kafka.Producer+import OpenTelemetry.Trace (SpanKind (..), defaultSpanArguments)+import OpenTelemetry.Trace qualified as Trace+import OpenTelemetry.Trace.Monad (MonadTracer, inSpan)+import UnliftIO (MonadUnliftIO, withRunInIO)+import Yesod.Core.Types (HandlerData (..), RunHandlerEnv (..))++envKafkaBrokerAddresses+  :: Env.Parser Env.Error (NonEmpty BrokerAddress)+envKafkaBrokerAddresses =+  Env.var+    (Env.eitherReader readKafkaBrokerAddresses)+    "KAFKA_BROKER_ADDRESSES"+    mempty++readKafkaBrokerAddresses :: String -> Either String (NonEmpty BrokerAddress)+readKafkaBrokerAddresses t = case NE.nonEmpty $ T.splitOn "," $ T.pack t of+  Just xs@(x NE.:| _)+    | x /= "" -> Right $ BrokerAddress <$> xs+  _ -> Left "Broker Address cannot be empty"++data KafkaProducerPoolConfig = KafkaProducerPoolConfig+  { kafkaProducerPoolConfigStripes :: Int+  -- ^ The number of stripes (distinct sub-pools) to maintain.+  -- The smallest acceptable value is 1.+  , kafkaProducerPoolConfigIdleTimeout :: NominalDiffTime+  -- ^ Amount of time for which an unused resource is kept open.+  -- The smallest acceptable value is 0.5 seconds.+  --+  -- The elapsed time before destroying a resource may be a little+  -- longer than requested, as the reaper thread wakes at 1-second+  -- intervals.+  , kafkaProducerPoolConfigSize :: Int+  -- ^ Maximum number of resources to keep open per stripe.  The+  -- smallest acceptable value is 1.+  --+  -- Requests for resources will block if this limit is reached on a+  -- single stripe, even if other stripes have idle resources+  -- available.+  }+  deriving stock (Show)++-- | Same defaults as 'Database.Persist.Sql.ConnectionPoolConfig'+defaultKafkaProducerPoolConfig :: KafkaProducerPoolConfig+defaultKafkaProducerPoolConfig = KafkaProducerPoolConfig 1 600 10++envKafkaProducerPoolConfig+  :: Env.Parser Env.Error KafkaProducerPoolConfig+envKafkaProducerPoolConfig = do+  poolSize <- Env.var Env.auto "KAFKA_PRODUCER_POOL_SIZE" $ Env.def 10+  pure $ defaultKafkaProducerPoolConfig {kafkaProducerPoolConfigSize = poolSize}++data KafkaProducerPool+  = NullKafkaProducerPool+  | KafkaProducerPool (Pool KafkaProducer)++class HasKafkaProducerPool env where+  kafkaProducerPoolL :: Lens' env KafkaProducerPool++instance HasKafkaProducerPool site => HasKafkaProducerPool (HandlerData child site) where+  kafkaProducerPoolL = envL . siteL . kafkaProducerPoolL++envL :: Lens' (HandlerData child site) (RunHandlerEnv child site)+envL = lens handlerEnv $ \x y -> x {handlerEnv = y}++siteL :: Lens' (RunHandlerEnv child site) site+siteL = lens rheSite $ \x y -> x {rheSite = y}++createKafkaProducerPool+  :: NonEmpty BrokerAddress+  -> KafkaProducerPoolConfig+  -> IO (Pool KafkaProducer)+createKafkaProducerPool addresses config =+  Pool.newPool+    $ Pool.setNumStripes (Just $ kafkaProducerPoolConfigStripes config)+    $ Pool.defaultPoolConfig+      mkProducer+      closeProducer+      (realToFrac $ kafkaProducerPoolConfigIdleTimeout config)+      (kafkaProducerPoolConfigSize config)+ where+  mkProducer =+    either+      ( \err -> Annotated.throw $ userError ("Failed to open kafka producer: " <> show err)+      )+      pure+      =<< newProducer (brokersList $ toList addresses)++produceKeyedOn+  :: ( MonadUnliftIO m+     , MonadLogger m+     , MonadTracer m+     , MonadReader env m+     , HasKafkaProducerPool env+     , ToJSON key+     , ToJSON value+     )+  => TopicName+  -> NonEmpty value+  -> (value -> key)+  -> m ()+produceKeyedOn prTopic values keyF = traced $ do+  logDebugNS "kafka" $ "Producing Kafka events" :# ["events" .= values]+  view kafkaProducerPoolL >>= \case+    NullKafkaProducerPool -> pure ()+    KafkaProducerPool producerPool ->+      withRunInIO $ \run ->+        Pool.withResource producerPool $ \producer ->+          for_ @NonEmpty values $ \value -> do+            mError <- liftIO $ produceMessage producer $ mkProducerRecord value+            for_ @Maybe mError $ \e ->+              run $ logErrorNS "kafka" $ "Failed to send event"+                :# ["error" .= (show e :: Text)]+ where+  mkProducerRecord value =+    ProducerRecord+      { prTopic+      , prPartition = UnassignedPartition+      , prKey = Just $ toStrict $ encode $ keyF value+      , prValue = Just $ toStrict $ encode value+      , prHeaders = mempty+      }++  traced =+    inSpan+      "kafka.produce"+      defaultSpanArguments+        { Trace.kind = Producer+        , Trace.attributes =+            HashMap.fromList+              [ ("service.name", "kafka")+              , ("topic", Trace.toAttribute $ unTopicName prTopic)+              ]+        }
+ package.yaml view
@@ -0,0 +1,71 @@+name: freckle-kafka+version: 0.0.0.0+maintainer: Freckle Education+category: Database+github: freckle/freckle-app+synopsis: Some extensions to the hw-kafka-client library+description: Please see README.md++extra-doc-files:+  - README.md+  - CHANGELOG.md++extra-source-files:+  - package.yaml++language: GHC2021++ghc-options:+  - -fignore-optim-changes+  - -fwrite-ide-info+  - -Weverything+  - -Wno-all-missed-specialisations+  - -Wno-missing-exported-signatures # re-enables missing-signatures+  - -Wno-missing-import-lists+  - -Wno-missing-kind-signatures+  - -Wno-missing-local-signatures+  - -Wno-missing-safe-haskell-mode+  - -Wno-monomorphism-restriction+  - -Wno-prepositive-qualified-module+  - -Wno-safe+  - -Wno-unsafe++when:+  - condition: "impl(ghc >= 9.8)"+    ghc-options:+      - -Wno-missing-role-annotations+      - -Wno-missing-poly-kind-signatures++dependencies:+  - base < 5++default-extensions:+  - DataKinds+  - DeriveAnyClass+  - DerivingVia+  - DerivingStrategies+  - GADTs+  - LambdaCase+  - NoImplicitPrelude+  - NoMonomorphismRestriction+  - OverloadedStrings+  - TypeFamilies++library:+  source-dirs: library+  dependencies:+    - Blammo >= 2.0.0.0+    - aeson+    - annotated-exception+    - containers+    - freckle-env+    - hs-opentelemetry-sdk+    - hw-kafka-client+    - lens+    - relude+    - resource-pool >= 0.4.0.0+    - text+    - time+    - unliftio+    - unordered-containers+    - yesod-core