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 +14/−0
- LICENSE +21/−0
- README.md +7/−0
- freckle-kafka.cabal +64/−0
- library/Freckle/App/Kafka.hs +7/−0
- library/Freckle/App/Kafka/Consumer.hs +234/−0
- library/Freckle/App/Kafka/Producer.hs +167/−0
- package.yaml +71/−0
+ 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