hw-kafka-client 4.0.3 → 4.0.4
raw patch · 15 files changed
+277/−137 lines, 15 filesPVP: major bump suggested
API removals or changes: PVP suggests a major version bump
API changes (from Hackage documentation)
- Kafka.Producer: produceMessageBatch :: MonadIO m => KafkaProducer -> [ProducerRecord] -> m [(ProducerRecord, KafkaError)]
+ Kafka.Consumer: RdKafkaRespErrAssignmentLost :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrAutoOffsetReset :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrDuplicateResource :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrFeatureUpdateFailed :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrInconsistentVoterSet :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrInvalidUpdateVersion :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrNoop :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrPrincipalDeserializationFailure :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrProducerFenced :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrResourceNotFound :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrThrottlingQuotaExceeded :: RdKafkaRespErrT
+ Kafka.Consumer: RdKafkaRespErrUnacceptableCredential :: RdKafkaRespErrT
+ Kafka.Consumer: [crHeaders] :: ConsumerRecord k v -> !Headers
+ Kafka.Consumer: data Headers
+ Kafka.Consumer: headersFromList :: [(ByteString, ByteString)] -> Headers
+ Kafka.Consumer: headersToList :: Headers -> [(ByteString, ByteString)]
+ Kafka.Consumer.Types: [crHeaders] :: ConsumerRecord k v -> !Headers
+ Kafka.Producer: RdKafkaRespErrAssignmentLost :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrAutoOffsetReset :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrDuplicateResource :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrFeatureUpdateFailed :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrInconsistentVoterSet :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrInvalidUpdateVersion :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrNoop :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrPrincipalDeserializationFailure :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrProducerFenced :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrResourceNotFound :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrThrottlingQuotaExceeded :: RdKafkaRespErrT
+ Kafka.Producer: RdKafkaRespErrUnacceptableCredential :: RdKafkaRespErrT
+ Kafka.Producer: [prHeaders] :: ProducerRecord -> !Headers
+ Kafka.Producer: data Headers
+ Kafka.Producer: headersFromList :: [(ByteString, ByteString)] -> Headers
+ Kafka.Producer: headersToList :: Headers -> [(ByteString, ByteString)]
+ Kafka.Producer.Types: [prHeaders] :: ProducerRecord -> !Headers
+ Kafka.Types: data Headers
+ Kafka.Types: headersFromList :: [(ByteString, ByteString)] -> Headers
+ Kafka.Types: headersToList :: Headers -> [(ByteString, ByteString)]
+ Kafka.Types: instance GHC.Base.Monoid Kafka.Types.Headers
+ Kafka.Types: instance GHC.Base.Semigroup Kafka.Types.Headers
+ Kafka.Types: instance GHC.Classes.Eq Kafka.Types.Headers
+ Kafka.Types: instance GHC.Generics.Generic Kafka.Types.Headers
+ Kafka.Types: instance GHC.Read.Read Kafka.Types.Headers
+ Kafka.Types: instance GHC.Show.Show Kafka.Types.Headers
- Kafka.Consumer: ConsumerRecord :: !TopicName -> !PartitionId -> !Offset -> !Timestamp -> !k -> !v -> ConsumerRecord k v
+ Kafka.Consumer: ConsumerRecord :: !TopicName -> !PartitionId -> !Offset -> !Timestamp -> !Headers -> !k -> !v -> ConsumerRecord k v
- Kafka.Consumer.Types: ConsumerRecord :: !TopicName -> !PartitionId -> !Offset -> !Timestamp -> !k -> !v -> ConsumerRecord k v
+ Kafka.Consumer.Types: ConsumerRecord :: !TopicName -> !PartitionId -> !Offset -> !Timestamp -> !Headers -> !k -> !v -> ConsumerRecord k v
- Kafka.Producer: ProducerRecord :: !TopicName -> !ProducePartition -> Maybe ByteString -> Maybe ByteString -> ProducerRecord
+ Kafka.Producer: ProducerRecord :: !TopicName -> !ProducePartition -> Maybe ByteString -> Maybe ByteString -> !Headers -> ProducerRecord
- Kafka.Producer.Types: ProducerRecord :: !TopicName -> !ProducePartition -> Maybe ByteString -> Maybe ByteString -> ProducerRecord
+ Kafka.Producer.Types: ProducerRecord :: !TopicName -> !ProducePartition -> Maybe ByteString -> Maybe ByteString -> !Headers -> ProducerRecord
Files
- example/ProducerExample.hs +3/−6
- hw-kafka-client.cabal +1/−1
- src/Kafka/Consumer.hs +6/−1
- src/Kafka/Consumer/Convert.hs +6/−3
- src/Kafka/Consumer/Types.hs +3/−2
- src/Kafka/Internal/RdKafka.chs +114/−2
- src/Kafka/Internal/Shared.hs +29/−2
- src/Kafka/Producer.hs +37/−90
- src/Kafka/Producer/Callbacks.hs +27/−24
- src/Kafka/Producer/Convert.hs +8/−1
- src/Kafka/Producer/Types.hs +2/−1
- src/Kafka/Types.hs +13/−1
- tests-it/Kafka/IntegrationSpec.hs +26/−3
- tests/Kafka/Consumer/ConsumerRecordMapSpec.hs +1/−0
- tests/Kafka/Consumer/ConsumerRecordTraverseSpec.hs +1/−0
example/ProducerExample.hs view
@@ -31,6 +31,7 @@ , prPartition = UnassignedPartition , prKey = k , prValue = v+ , prHeaders = mempty } -- Run an example@@ -61,12 +62,8 @@ putStrLn "And the last one..." msg3 <- getLine err3 <- produceMessage prod (mkMessage (Just "key3") (Just $ pack msg3))-- -- errs <- produceMessageBatch prod- -- [ mkMessage (Just "b-1") (Just "batch-1")- -- , mkMessage (Just "b-2") (Just "batch-2")- -- , mkMessage Nothing (Just "batch-3")- -- ]+ + err4 <- produceMessage prod ((mkMessage (Just "key4") (Just $ pack msg3)) { prHeaders = headersFromList [("fancy", "header")]}) -- forM_ errs (print . snd)
hw-kafka-client.cabal view
@@ -1,7 +1,7 @@ cabal-version: 2.2 name: hw-kafka-client-version: 4.0.3+version: 4.0.4 synopsis: Kafka bindings for Haskell description: Apache Kafka bindings backed by the librdkafka C library. .
src/Kafka/Consumer.hs view
@@ -168,7 +168,7 @@ mbq <- readIORef qr case mbq of Nothing -> return [Left $ KafkaBadSpecification "Calling pollMessageBatch while CallbackPollMode is set to CallbackPollModeSync."]- Just q -> rdKafkaConsumeBatchQueue q ms b >>= traverse fromMessagePtr+ Just q -> whileNoCallbackRunning c $ rdKafkaConsumeBatchQueue q ms b >>= traverse fromMessagePtr -- | Commit message's offset on broker for the message's partition. commitOffsetMessage :: MonadIO m@@ -372,6 +372,11 @@ case st of CallbackPollEnabled -> go CallbackPollDisabled -> pure ()++whileNoCallbackRunning :: KafkaConsumer -> IO a -> IO a+whileNoCallbackRunning k f = do+ let statusVar = kcfgCallbackPollStatus (getKafkaConf k)+ withMVar statusVar $ \_ -> f withCallbackPollEnabled :: KafkaConsumer -> IO () -> IO CallbackPollStatus withCallbackPollEnabled k f = do
src/Kafka/Consumer/Convert.hs view
@@ -18,6 +18,7 @@ import Control.Monad ((>=>)) import qualified Data.ByteString as BS+import Data.Either (fromRight) import Data.Int (Int64) import Data.Map.Strict (Map, fromListWith) import qualified Data.Set as S@@ -41,7 +42,7 @@ , rdKafkaTopicPartitionListNew , peekCText )-import Kafka.Internal.Shared (kafkaRespErr, readTopic, readKey, readPayload, readTimestamp)+import Kafka.Internal.Shared (kafkaRespErr, readHeaders, readTopic, readKey, readPayload, readTimestamp) import Kafka.Types (KafkaError(..), PartitionId(..), TopicName(..)) -- | Converts offsets sync policy to integer (the way Kafka understands it):@@ -158,20 +159,22 @@ s <- peek realPtr msg <- if err'RdKafkaMessageT s /= RdKafkaRespErrNoError then return . Left . KafkaResponseError $ err'RdKafkaMessageT s- else Right <$> mkRecord s+ else Right <$> mkRecord s realPtr rdKafkaMessageDestroy realPtr return msg where- mkRecord msg = do+ mkRecord msg rptr = do topic <- readTopic msg key <- readKey msg payload <- readPayload msg timestamp <- readTimestamp ptr+ headers <- fromRight mempty <$> readHeaders rptr return ConsumerRecord { crTopic = TopicName topic , crPartition = PartitionId $ partition'RdKafkaMessageT msg , crOffset = Offset $ offset'RdKafkaMessageT msg , crTimestamp = timestamp+ , crHeaders = headers , crKey = key , crValue = payload }
src/Kafka/Consumer/Types.hs view
@@ -43,7 +43,7 @@ import Data.Typeable (Typeable) import GHC.Generics (Generic) import Kafka.Internal.Setup (HasKafka (..), HasKafkaConf (..), Kafka (..), KafkaConf (..))-import Kafka.Types (Millis (..), PartitionId (..), TopicName (..))+import Kafka.Types (Millis (..), PartitionId (..), TopicName (..), Headers) -- | The main type for Kafka consumption, used e.g. to poll and commit messages. -- @@ -143,13 +143,14 @@ , crPartition :: !PartitionId -- ^ Kafka partition this message was received from , crOffset :: !Offset -- ^ Offset within the 'crPartition' Kafka partition , crTimestamp :: !Timestamp -- ^ Message timestamp+ , crHeaders :: !Headers -- ^ Message headers , crKey :: !k -- ^ Message key , crValue :: !v -- ^ Message value } deriving (Eq, Show, Read, Typeable, Generic) instance Bifunctor ConsumerRecord where- bimap f g (ConsumerRecord t p o ts k v) = ConsumerRecord t p o ts (f k) (g v)+ bimap f g (ConsumerRecord t p o ts hds k v) = ConsumerRecord t p o ts hds (f k) (g v) {-# INLINE bimap #-} instance Functor (ConsumerRecord k) where
src/Kafka/Internal/RdKafka.chs view
@@ -13,13 +13,13 @@ import Foreign.Concurrent (newForeignPtr) import qualified Foreign.Concurrent as Concurrent import Foreign.Marshal.Alloc (alloca, allocaBytes)-import Foreign.Marshal.Array (peekArray, allocaArray)+import Foreign.Marshal.Array (peekArray, allocaArray, withArrayLen) import Foreign.Storable (Storable(..)) import Foreign.Ptr (Ptr, FunPtr, castPtr, nullPtr) import Foreign.ForeignPtr (FinalizerPtr, addForeignPtrFinalizer, newForeignPtr_, withForeignPtr) import Foreign.C.Error (Errno(..), getErrno) import Foreign.C.String (CString, newCString, withCAString, peekCAString, peekCString)-import Foreign.C.Types (CFile, CInt(..), CSize, CChar)+import Foreign.C.Types (CFile, CInt(..), CSize, CChar, CLong) import System.IO (Handle, stdin, stdout, stderr) import System.Posix.IO (handleToFd) import System.Posix.Types (Fd(..))@@ -971,6 +971,118 @@ res <- newUnmanagedRdKafkaTopicT kafkaPtr topic topicConfPtr _ <- traverse (addForeignPtrFinalizer rdKafkaTopicDestroy') res return res++-------------------------------------------------------------------------------------------------+---- Errors++data RdKafkaErrorT+{#pointer *rd_kafka_error_t as RdKafkaErrorTPtr -> RdKafkaErrorT #}++{#fun rd_kafka_error_code as ^+ {`RdKafkaErrorTPtr'} -> `RdKafkaRespErrT' cIntToEnum #}++{#fun rd_kafka_error_destroy as ^+ {`RdKafkaErrorTPtr'} -> `()' #}+-------------------------------------------------------------------------------------------------+---- Headers++data RdKafkaHeadersT+{#pointer *rd_kafka_headers_t as RdKafkaHeadersTPtr -> RdKafkaHeadersT #}++{#fun rd_kafka_header_get_all as ^+ {`RdKafkaHeadersTPtr', cIntConv `CSize', castPtr `Ptr CString', castPtr `Ptr Word8Ptr', castPtr `CSizePtr'} -> `RdKafkaRespErrT' cIntToEnum #}++{#fun rd_kafka_message_headers as ^+ {castPtr `Ptr RdKafkaMessageT', alloca- `RdKafkaHeadersTPtr' peekPtr*} -> `RdKafkaRespErrT' cIntToEnum #}++--- Produceva api++{#enum rd_kafka_vtype_t as ^ {underscoreToCase} deriving (Show, Eq) #}++data RdKafkaVuT+ = Topic'RdKafkaVu CString+ | TopicHandle'RdKafkaVu (Ptr RdKafkaTopicT)+ | Partition'RdKafkaVu CInt32T+ | Value'RdKafkaVu Word8Ptr CSize+ | Key'RdKafkaVu Word8Ptr CSize+ | MsgFlags'RdKafkaVu CInt+ | Timestamp'RdKafkaVu CInt64T+ | Opaque'RdKafkaVu (Ptr ())+ | Header'RdKafkaVu CString Word8Ptr CSize+ | Headers'RdKafkaVu (Ptr RdKafkaHeadersT) -- The message object will assume ownership of the headers (unless produceva() fails)+ | End'RdKafkaVu++{#pointer *rd_kafka_vu_t as RdKafkaVuTPtr foreign -> RdKafkaVuT #}++instance Storable RdKafkaVuT where+ alignment _ = {#alignof rd_kafka_vu_t #}+ sizeOf _ = {#sizeof rd_kafka_vu_t #}+ peek p = {#get rd_kafka_vu_t->vtype #} p >>= \a -> case cIntToEnum a of+ RdKafkaVtypeEnd -> return End'RdKafkaVu+ RdKafkaVtypeTopic -> Topic'RdKafkaVu <$> ({#get rd_kafka_vu_t->u.cstr #} p)+ RdKafkaVtypeMsgflags -> MsgFlags'RdKafkaVu <$> ({#get rd_kafka_vu_t->u.i #} p)+ RdKafkaVtypeTimestamp -> Timestamp'RdKafkaVu <$> ({#get rd_kafka_vu_t->u.i64 #} p)+ RdKafkaVtypePartition -> Partition'RdKafkaVu <$> ({#get rd_kafka_vu_t->u.i32 #} p)+ RdKafkaVtypeHeaders -> Headers'RdKafkaVu <$> ({#get rd_kafka_vu_t->u.headers #} p)+ RdKafkaVtypeValue -> do+ nm <- liftM castPtr ({#get rd_kafka_vu_t->u.mem.ptr #} p)+ sz <- ({#get rd_kafka_vu_t->u.mem.size #} p)+ return $ Value'RdKafkaVu nm (cIntConv sz)+ RdKafkaVtypeKey -> do+ nm <- liftM castPtr ({#get rd_kafka_vu_t->u.mem.ptr #} p)+ sz <- ({#get rd_kafka_vu_t->u.mem.size #} p)+ return $ Key'RdKafkaVu nm (cIntConv sz)+ RdKafkaVtypeRkt -> TopicHandle'RdKafkaVu <$> ({#get rd_kafka_vu_t->u.rkt #} p)+ RdKafkaVtypeOpaque -> Opaque'RdKafkaVu <$> ({#get rd_kafka_vu_t->u.ptr #} p)+ RdKafkaVtypeHeader -> do+ nm <- ({#get rd_kafka_vu_t->u.header.name #} p)+ val' <- liftM castPtr ({#get rd_kafka_vu_t->u.header.val #} p)+ sz <- ({#get rd_kafka_vu_t->u.header.size #} p)+ return $ Header'RdKafkaVu nm val' (cIntConv sz)+ poke p End'RdKafkaVu =+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypeEnd)+ poke p (Topic'RdKafkaVu str) = do+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypeTopic)+ {#set rd_kafka_vu_t.u.cstr #} p str+ poke p (Timestamp'RdKafkaVu tms) = do+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypeTimestamp)+ {#set rd_kafka_vu_t.u.i64 #} p tms+ poke p (Partition'RdKafkaVu prt) = do+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypePartition)+ {#set rd_kafka_vu_t.u.i32 #} p prt+ poke p (MsgFlags'RdKafkaVu flags) = do+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypeMsgflags)+ {#set rd_kafka_vu_t.u.i #} p flags+ poke p (Headers'RdKafkaVu headers) = do+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypeHeaders)+ {#set rd_kafka_vu_t.u.headers #} p headers+ poke p (TopicHandle'RdKafkaVu tphandle) = do+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypeRkt)+ {#set rd_kafka_vu_t.u.rkt #} p tphandle+ poke p (Value'RdKafkaVu pl sz) = do+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypeValue)+ {#set rd_kafka_vu_t.u.mem.size #} p (cIntConv sz)+ {#set rd_kafka_vu_t.u.mem.ptr #} p (castPtr pl)+ poke p (Key'RdKafkaVu pl sz) = do+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypeKey)+ {#set rd_kafka_vu_t.u.mem.size #} p (cIntConv sz)+ {#set rd_kafka_vu_t.u.mem.ptr #} p (castPtr pl)+ poke p (Opaque'RdKafkaVu ptr') = do+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypeOpaque)+ {#set rd_kafka_vu_t.u.ptr #} p ptr'+ poke p (Header'RdKafkaVu nm val' sz) = do+ {#set rd_kafka_vu_t.vtype #} p (enumToCInt RdKafkaVtypeHeader)+ {#set rd_kafka_vu_t.u.header.size #} p (cIntConv sz)+ {#set rd_kafka_vu_t.u.header.name #} p nm+ {#set rd_kafka_vu_t.u.header.val #} p (castPtr val')++{#fun rd_kafka_produceva as rdKafkaMessageProduceVa'+ {`RdKafkaTPtr', `RdKafkaVuTPtr', `CLong'} -> `RdKafkaErrorTPtr' #}++rdKafkaMessageProduceVa :: RdKafkaTPtr -> [RdKafkaVuT] -> IO RdKafkaErrorTPtr+rdKafkaMessageProduceVa kafkaPtr vts = withArrayLen vts $ \i arrPtr -> do+ fptr <- newForeignPtr_ arrPtr+ rdKafkaMessageProduceVa' kafkaPtr fptr (cIntConv i) -- Marshall / Unmarshall enumToCInt :: Enum a => a -> CInt
@@ -1,3 +1,5 @@+{-# LANGUAGE LambdaCase #-}+ module Kafka.Internal.Shared ( pollEvents , word8PtrToBS@@ -8,6 +10,7 @@ , kafkaErrorToEither , kafkaErrorToMaybe , maybeToLeft+, readHeaders , readPayload , readTopic , readKey@@ -29,9 +32,9 @@ import Foreign.Ptr (Ptr, nullPtr) import Foreign.Storable (Storable (peek)) import Kafka.Consumer.Types (Timestamp (..))-import Kafka.Internal.RdKafka (RdKafkaMessageT (..), RdKafkaMessageTPtr, RdKafkaRespErrT (..), RdKafkaTimestampTypeT (..), Word8Ptr, rdKafkaErrno2err, rdKafkaMessageTimestamp, rdKafkaPoll, rdKafkaTopicName)+import Kafka.Internal.RdKafka (RdKafkaMessageT (..), RdKafkaMessageTPtr, RdKafkaRespErrT (..), RdKafkaTimestampTypeT (..), Word8Ptr, rdKafkaErrno2err, rdKafkaMessageTimestamp, rdKafkaPoll, rdKafkaTopicName, rdKafkaHeaderGetAll, rdKafkaMessageHeaders) import Kafka.Internal.Setup (HasKafka (..), Kafka (..))-import Kafka.Types (KafkaError (..), Millis (..), Timeout (..))+import Kafka.Types (KafkaError (..), Millis (..), Timeout (..), Headers, headersFromList) pollEvents :: HasKafka a => a -> Maybe Timeout -> IO () pollEvents a tm =@@ -101,6 +104,30 @@ RdKafkaTimestampCreateTime -> CreateTime (Millis ts) RdKafkaTimestampLogAppendTime -> LogAppendTime (Millis ts) RdKafkaTimestampNotAvailable -> NoTimestamp+++readHeaders :: Ptr RdKafkaMessageT -> IO (Either RdKafkaRespErrT Headers)+readHeaders msg = do+ (err, headersPtr) <- rdKafkaMessageHeaders msg+ case err of+ RdKafkaRespErrNoent -> return $ Right mempty+ RdKafkaRespErrNoError -> fmap headersFromList <$> extractHeaders headersPtr+ e -> return . Left $ e+ where extractHeaders ptHeaders =+ alloca $ \nptr ->+ alloca $ \vptr ->+ alloca $ \szptr ->+ let go acc idx = rdKafkaHeaderGetAll ptHeaders idx nptr vptr szptr >>= \case+ RdKafkaRespErrNoent -> return $ Right acc+ RdKafkaRespErrNoError -> do+ cstr <- peek nptr+ wptr <- peek vptr+ csize <- peek szptr+ hn <- BS.packCString cstr+ hv <- word8PtrToBS (fromIntegral csize) wptr+ go ((hn, hv) : acc) (idx + 1)+ _ -> error "Unexpected error code while extracting headers"+ in go [] 0 readBS :: (t -> Int) -> (t -> Ptr Word8) -> t -> IO (Maybe BS.ByteString) readBS flen fdata s = if fdata s == nullPtr
src/Kafka/Producer.hs view
@@ -58,7 +58,7 @@ , module X , runProducer , newProducer-, produceMessage, produceMessageBatch+, produceMessage , produceMessage' , flushProducer , closeProducer@@ -66,25 +66,21 @@ ) where -import Control.Arrow ((&&&)) import Control.Exception (bracket)-import Control.Monad (forM, forM_, (<=<))+import Control.Monad (forM_) import Control.Monad.IO.Class (MonadIO (liftIO)) import qualified Data.ByteString as BS import qualified Data.ByteString.Internal as BSI-import Data.Function (on)-import Data.List (groupBy, sortBy)-import Data.Ord (comparing) import qualified Data.Text as Text-import Foreign.ForeignPtr (newForeignPtr_, withForeignPtr)-import Foreign.Marshal.Array (withArrayLen)+import Foreign.C.String (withCString)+import Foreign.ForeignPtr (withForeignPtr)+import Foreign.Marshal.Utils (withMany) import Foreign.Ptr (Ptr, nullPtr, plusPtr)-import Foreign.Storable (Storable (..)) import Foreign.StablePtr (newStablePtr, castStablePtrToPtr)-import Kafka.Internal.RdKafka (RdKafkaMessageT (..), RdKafkaRespErrT (..), RdKafkaTypeT (..), destroyUnmanagedRdKafkaTopic, newRdKafkaT, newUnmanagedRdKafkaTopicT, rdKafkaOutqLen, rdKafkaProduce, rdKafkaProduceBatch, rdKafkaSetLogLevel)-import Kafka.Internal.Setup (Kafka (..), KafkaConf (..), KafkaProps (..), TopicConf (..), TopicProps (..), kafkaConf, topicConf, Callback(..))+import Kafka.Internal.RdKafka (RdKafkaRespErrT (..), RdKafkaTypeT (..), RdKafkaVuT(..), newRdKafkaT, rdKafkaErrorCode, rdKafkaErrorDestroy, rdKafkaOutqLen, rdKafkaMessageProduceVa, rdKafkaSetLogLevel)+import Kafka.Internal.Setup (Kafka (..), KafkaConf (..), KafkaProps (..), TopicProps (..), kafkaConf, topicConf, Callback(..)) import Kafka.Internal.Shared (pollEvents)-import Kafka.Producer.Convert (copyMsgFlags, handleProduceErr', producePartitionCInt, producePartitionInt)+import Kafka.Producer.Convert (copyMsgFlags, handleProduceErrT, producePartitionCInt) import Kafka.Producer.Types (KafkaProducer (..)) import Kafka.Producer.ProducerProperties as X@@ -93,7 +89,7 @@ -- | Runs Kafka Producer. -- The callback provided is expected to call 'produceMessage'--- or/and 'produceMessageBatch' to send messages to Kafka.+-- to send messages to Kafka. {-# DEPRECATED runProducer "Use 'newProducer'/'closeProducer' instead" #-} runProducer :: ProducerProperties -> (KafkaProducer -> IO (Either KafkaError a))@@ -148,94 +144,37 @@ -- -- The callback can be a long running process, as it is forked by the thread -- that handles the delivery reports.--- produceMessage' :: MonadIO m => KafkaProducer -> ProducerRecord -> (DeliveryReport -> IO ()) -> m (Either ImmediateError ())-produceMessage' kp@(KafkaProducer (Kafka k) _ (TopicConf tc)) msg cb = liftIO $- fireCallbacks >> bracket (mkTopic . prTopic $ msg) closeTopic withTopic+produceMessage' kp@(KafkaProducer (Kafka k) _ _) msg cb = liftIO $+ fireCallbacks >> produceIt where fireCallbacks = pollEvents kp . Just . Timeout $ 0 - mkTopic (TopicName tn) =- newUnmanagedRdKafkaTopicT k (Text.unpack tn) (Just tc)-- closeTopic = either mempty destroyUnmanagedRdKafkaTopic-- withTopic (Left err) = return . Left . ImmediateError . KafkaError . Text.pack $ err- withTopic (Right topic) =+ produceIt = withBS (prValue msg) $ \payloadPtr payloadLength ->- withBS (prKey msg) $ \keyPtr keyLength -> do- callbackPtr <- newStablePtr cb- res <- handleProduceErr' =<< rdKafkaProduce- topic- (producePartitionCInt (prPartition msg))- copyMsgFlags- payloadPtr- (fromIntegral payloadLength)- keyPtr- (fromIntegral keyLength)- (castStablePtrToPtr callbackPtr)-- pure $ case res of- Left err -> Left . ImmediateError $ err- Right () -> Right ()---- | Sends a batch of messages.--- Returns a list of messages which it was unable to send with corresponding errors.--- Since librdkafka is backed by a queue, this function can return before messages are sent. See--- 'flushProducer' to wait for queue to empty.-produceMessageBatch :: MonadIO m- => KafkaProducer- -> [ProducerRecord]- -> m [(ProducerRecord, KafkaError)]- -- ^ An empty list when the operation is successful,- -- otherwise a list of "failed" messages with corresponsing errors.-produceMessageBatch kp@(KafkaProducer (Kafka k) _ (TopicConf tc)) messages = liftIO $ do- pollEvents kp (Just $ Timeout 0) -- fire callbacks if any exist (handle delivery reports)- concat <$> forM (mkBatches messages) sendBatch- where- mkSortKey = prTopic &&& prPartition- mkBatches = groupBy ((==) `on` mkSortKey) . sortBy (comparing mkSortKey)-- mkTopic (TopicName tn) = newUnmanagedRdKafkaTopicT k (Text.unpack tn) (Just tc)-- clTopic = either (return . const ()) destroyUnmanagedRdKafkaTopic-- sendBatch [] = return []- sendBatch batch = bracket (mkTopic $ prTopic (head batch)) clTopic (withTopic batch)-- withTopic ms (Left err) = return $ (, KafkaError (Text.pack err)) <$> ms- withTopic ms (Right t) = do- let (partInt, partCInt) = (producePartitionInt &&& producePartitionCInt) $ prPartition (head ms)- withForeignPtr t $ \topicPtr -> do- nativeMs <- forM ms (toNativeMessage topicPtr partInt)- withArrayLen nativeMs $ \len batchPtr -> do- batchPtrF <- newForeignPtr_ batchPtr- numRet <- rdKafkaProduceBatch t partCInt copyMsgFlags batchPtrF len- if numRet == len then return []- else do- errs <- mapM (return . err'RdKafkaMessageT <=< peekElemOff batchPtr)- [0..(fromIntegral $ len - 1)]- return [(m, KafkaResponseError e) | (m, e) <- zip messages errs, e /= RdKafkaRespErrNoError]+ withBS (prKey msg) $ \keyPtr keyLength ->+ withHeaders (prHeaders msg) $ \hdrs ->+ withCString (Text.unpack . unTopicName . prTopic $ msg) $ \topicName -> do+ callbackPtr <- newStablePtr cb+ let opts = [+ Topic'RdKafkaVu topicName+ , Partition'RdKafkaVu . producePartitionCInt . prPartition $ msg+ , MsgFlags'RdKafkaVu (fromIntegral copyMsgFlags)+ , Value'RdKafkaVu payloadPtr (fromIntegral payloadLength)+ , Key'RdKafkaVu keyPtr (fromIntegral keyLength)+ , Opaque'RdKafkaVu (castStablePtrToPtr callbackPtr)+ ] - toNativeMessage t p m =- withBS (prValue m) $ \payloadPtr payloadLength ->- withBS (prKey m) $ \keyPtr keyLength ->- return RdKafkaMessageT- { err'RdKafkaMessageT = RdKafkaRespErrNoError- , topic'RdKafkaMessageT = t- , partition'RdKafkaMessageT = p- , len'RdKafkaMessageT = payloadLength- , payload'RdKafkaMessageT = payloadPtr- , offset'RdKafkaMessageT = 0- , keyLen'RdKafkaMessageT = keyLength- , key'RdKafkaMessageT = keyPtr- , opaque'RdKafkaMessageT = nullPtr- }+ code <- bracket (rdKafkaMessageProduceVa k (hdrs ++ opts)) rdKafkaErrorDestroy rdKafkaErrorCode+ res <- handleProduceErrT code+ pure $ case res of+ Just err -> Left . ImmediateError $ err+ Nothing -> Right () -- | Closes the producer. -- Will wait until the outbound queue is drained before returning the control.@@ -254,6 +193,14 @@ else flushProducer kp ------------------------------------------------------------------------------------++withHeaders :: Headers -> ([RdKafkaVuT] -> IO a) -> IO a+withHeaders hds = withMany allocHeader (headersToList hds)+ where+ allocHeader (nm, val) f = + BS.useAsCString nm $ \cnm ->+ withBS (Just val) $ \vp vl ->+ f $ Header'RdKafkaVu cnm vp (fromIntegral vl) withBS :: Maybe BS.ByteString -> (Ptr a -> Int -> IO b) -> IO b withBS Nothing f = f nullPtr 0
src/Kafka/Producer/Callbacks.hs view
@@ -1,4 +1,5 @@ {-# LANGUAGE TypeApplications #-}+{-# LANGUAGE LambdaCase #-} module Kafka.Producer.Callbacks ( deliveryCallback , module X@@ -15,10 +16,11 @@ import Kafka.Callbacks as X import Kafka.Consumer.Types (Offset(..)) import Kafka.Internal.RdKafka (RdKafkaMessageT(..), RdKafkaRespErrT(..), rdKafkaConfSetDrMsgCb)-import Kafka.Internal.Setup (KafkaConf(..), getRdKafkaConf, Callback(..))-import Kafka.Internal.Shared (kafkaRespErr, readTopic, readKey, readPayload)+import Kafka.Internal.Setup (getRdKafkaConf, Callback(..))+import Kafka.Internal.Shared (kafkaRespErr, readTopic, readKey, readPayload, readHeaders) import Kafka.Producer.Types (ProducerRecord(..), DeliveryReport(..), ProducePartition(..)) import Kafka.Types (KafkaError(..), TopicName(..))+import Data.Either (fromRight) -- | Sets the callback for delivery reports. --@@ -36,10 +38,12 @@ then getErrno >>= (callback . NoMessageError . kafkaRespErr) else do s <- peek mptr+ prodRec <- mkProdRec mptr let cbPtr = opaque'RdKafkaMessageT s- if err'RdKafkaMessageT s /= RdKafkaRespErrNoError- then mkErrorReport s >>= callbacks cbPtr- else mkSuccessReport s >>= callbacks cbPtr+ callbacks cbPtr $ + if err'RdKafkaMessageT s /= RdKafkaRespErrNoError+ then mkErrorReport s prodRec + else mkSuccessReport s prodRec callbacks cbPtr rep = do callback rep@@ -51,24 +55,23 @@ -- blocking here would block librdkafka from continuing its execution void . forkIO $ msgCb rep -mkErrorReport :: RdKafkaMessageT -> IO DeliveryReport-mkErrorReport msg = do- prodRec <- mkProdRec msg- pure $ DeliveryFailure prodRec (KafkaResponseError (err'RdKafkaMessageT msg))+mkErrorReport :: RdKafkaMessageT -> ProducerRecord -> DeliveryReport+mkErrorReport msg prodRec = DeliveryFailure prodRec (KafkaResponseError (err'RdKafkaMessageT msg)) -mkSuccessReport :: RdKafkaMessageT -> IO DeliveryReport-mkSuccessReport msg = do- prodRec <- mkProdRec msg- pure $ DeliverySuccess prodRec (Offset $ offset'RdKafkaMessageT msg)+mkSuccessReport :: RdKafkaMessageT -> ProducerRecord -> DeliveryReport+mkSuccessReport msg prodRec = DeliverySuccess prodRec (Offset $ offset'RdKafkaMessageT msg) -mkProdRec :: RdKafkaMessageT -> IO ProducerRecord-mkProdRec msg = do- topic <- readTopic msg- key <- readKey msg- payload <- readPayload msg- pure ProducerRecord- { prTopic = TopicName topic- , prPartition = SpecifiedPartition (partition'RdKafkaMessageT msg)- , prKey = key- , prValue = payload- }+mkProdRec :: Ptr RdKafkaMessageT -> IO ProducerRecord+mkProdRec pmsg = do+ msg <- peek pmsg + topic <- readTopic msg+ key <- readKey msg+ payload <- readPayload msg+ flip fmap (fromRight mempty <$> readHeaders pmsg) $ \headers -> + ProducerRecord+ { prTopic = TopicName topic+ , prPartition = SpecifiedPartition (partition'RdKafkaMessageT msg)+ , prKey = key+ , prValue = payload+ , prHeaders = headers+ }
src/Kafka/Producer/Convert.hs view
@@ -4,12 +4,13 @@ , producePartitionCInt , handleProduceErr , handleProduceErr'+, handleProduceErrT ) where import Foreign.C.Error (getErrno) import Foreign.C.Types (CInt)-import Kafka.Internal.RdKafka (rdKafkaMsgFlagCopy)+import Kafka.Internal.RdKafka (RdKafkaRespErrT(..), rdKafkaMsgFlagCopy) import Kafka.Internal.Shared (kafkaRespErr) import Kafka.Types (KafkaError(..)) import Kafka.Producer.Types (ProducePartition(..))@@ -32,6 +33,12 @@ handleProduceErr 0 = return Nothing handleProduceErr _ = return $ Just KafkaInvalidReturnValue {-# INLINE handleProduceErr #-}++handleProduceErrT :: RdKafkaRespErrT -> IO (Maybe KafkaError)+handleProduceErrT RdKafkaRespErrUnknown = Just . kafkaRespErr <$> getErrno+handleProduceErrT RdKafkaRespErrNoError = return Nothing+handleProduceErrT e = return $ Just (KafkaResponseError e)+{-# INLINE handleProduceErrT #-} handleProduceErr' :: Int -> IO (Either KafkaError ()) handleProduceErr' (- 1) = Left . kafkaRespErr <$> getErrno
src/Kafka/Producer/Types.hs view
@@ -21,7 +21,7 @@ import GHC.Generics (Generic) import Kafka.Consumer.Types (Offset (..)) import Kafka.Internal.Setup (HasKafka (..), HasKafkaConf (..), HasTopicConf (..), Kafka (..), KafkaConf (..), TopicConf (..))-import Kafka.Types (KafkaError (..), TopicName (..))+import Kafka.Types (KafkaError (..), TopicName (..), Headers) -- | The main type for Kafka message production, used e.g. to send messages. --@@ -50,6 +50,7 @@ , prPartition :: !ProducePartition , prKey :: Maybe ByteString , prValue :: Maybe ByteString+ , prHeaders :: !Headers } deriving (Eq, Show, Typeable, Generic) -- |
src/Kafka/Types.hs view
@@ -21,6 +21,7 @@ , KafkaDebug(..) , KafkaCompressionCodec(..) , TopicType(..)+, Headers, headersFromList, headersToList , topicType , kafkaDebugToText , kafkaCompressionCodecToText@@ -34,6 +35,7 @@ import Data.Typeable (Typeable) import GHC.Generics (Generic) import Kafka.Internal.RdKafka (RdKafkaRespErrT, rdKafkaErr2name, rdKafkaErr2str)+import qualified Data.ByteString as BS -- | Kafka broker ID newtype BrokerId = BrokerId { unBrokerId :: Int } deriving (Show, Eq, Ord, Read, Generic)@@ -158,4 +160,14 @@ NoCompression -> "none" Gzip -> "gzip" Snappy -> "snappy"- Lz4 -> "lz4"+ Lz4 -> "lz4" ++-- | Headers that might be passed along with a record+newtype Headers = Headers { unHeaders :: [(BS.ByteString, BS.ByteString)] } + deriving (Eq, Show, Semigroup, Monoid, Read, Typeable, Generic)++headersFromList :: [(BS.ByteString, BS.ByteString)] -> Headers+headersFromList = Headers++headersToList :: Headers -> [(BS.ByteString, BS.ByteString)]+headersToList = unHeaders
tests-it/Kafka/IntegrationSpec.hs view
@@ -6,10 +6,11 @@ where import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar)-import Control.Monad (forM, forM_)+import Control.Monad (forM, forM_, void) import Control.Monad.Loops import Data.Either import Data.Map (fromList)+import qualified Data.Set as Set import Data.Monoid ((<>)) import Kafka.Consumer import Kafka.Metadata@@ -122,6 +123,7 @@ , prPartition = UnassignedPartition , prKey = Nothing , prValue = Just "test from producer"+ , prHeaders = mempty } res <- produceMessage' prod msg (putMVar var)@@ -151,6 +153,23 @@ it "should consume empty batch when there are no messages" $ \k -> do res <- pollMessageBatch k (Timeout 1000) (BatchSize 50) length res `shouldBe` 0+ + describe "Kafka.Headers.Spec" $ do+ let testHeaders = headersFromList [("a-header-name", "a-header-value"), ("b-header-name", "b-header-value")]++ specWithKafka "Headers consumer/producer" consumerProps $ do+ it "1. sends 2 messages to test topic enriched with headers" $ \(k, prod) -> do+ void $ receiveMessages k+ + res <- sendMessagesWithHeaders (testMessages testTopic) testHeaders prod+ res `shouldBe` Right ()+ it "2. should receive 2 messages enriched with headers" $ \(k, _) -> do+ res <- receiveMessages k+ (length <$> res) `shouldBe` Right 2+ + forM_ res $ \rcs -> + forM_ rcs ((`shouldBe` Set.fromList (headersToList testHeaders)) . Set.fromList . headersToList . crHeaders)+ ---------------------------------------------------------------------------------------------------------------- data ReadState = Skip | Read@@ -171,13 +190,17 @@ testMessages :: TopicName -> [ProducerRecord] testMessages t =- [ ProducerRecord t UnassignedPartition Nothing (Just "test from producer")- , ProducerRecord t UnassignedPartition (Just "key") (Just "test from producer (with key)")+ [ ProducerRecord t UnassignedPartition Nothing (Just "test from producer") mempty+ , ProducerRecord t UnassignedPartition (Just "key") (Just "test from producer (with key)") mempty ] sendMessages :: [ProducerRecord] -> KafkaProducer -> IO (Either KafkaError ()) sendMessages msgs prod = Right <$> (forM_ msgs (produceMessage prod) >> flushProducer prod)++sendMessagesWithHeaders :: [ProducerRecord] -> Headers -> KafkaProducer -> IO (Either KafkaError ())+sendMessagesWithHeaders msgs hdrs prod =+ Right <$> (forM_ msgs (\msg -> produceMessage prod (msg {prHeaders = hdrs})) >> flushProducer prod) runConsumerSpec :: SpecWith KafkaConsumer runConsumerSpec = do
tests/Kafka/Consumer/ConsumerRecordMapSpec.hs view
@@ -19,6 +19,7 @@ , crPartition = PartitionId 0 , crOffset = Offset 5 , crTimestamp = NoTimestamp+ , crHeaders = mempty , crKey = Just testKey , crValue = Just testValue }
tests/Kafka/Consumer/ConsumerRecordTraverseSpec.hs view
@@ -21,6 +21,7 @@ , crPartition = PartitionId 0 , crOffset = Offset 5 , crTimestamp = NoTimestamp+ , crHeaders = mempty , crKey = testKey , crValue = testValue }