packages feed

hw-kafka-client 2.4.0 → 2.4.1

raw patch · 3 files changed

+18/−7 lines, 3 files

Files

hw-kafka-client.cabal view
@@ -1,5 +1,5 @@ name:                hw-kafka-client-version:             2.4.0+version:             2.4.1 homepage:            https://github.com/haskell-works/hw-kafka-client bug-reports:         https://github.com/haskell-works/hw-kafka-client/issues license:             MIT
src/Kafka/Consumer.hs view
@@ -90,7 +90,7 @@             -> Timeout -- ^ the timeout, in milliseconds             -> m (Either KafkaError (ConsumerRecord (Maybe BS.ByteString) (Maybe BS.ByteString))) -- ^ Left on error or timeout, right for success pollMessage c@(KafkaConsumer _ (KafkaConf _ qr _)) (Timeout ms) = liftIO $ do---    pollConsumerEvents c Nothing+    pollConsumerEvents c Nothing     mbq <- readIORef qr     case mbq of       Nothing -> return . Left $ KafkaBadSpecification "Messages queue is not configured, internal error, fatal."
src/Kafka/Consumer/Callbacks.hs view
@@ -35,7 +35,7 @@   where     realCb k err pl = do       k' <- newForeignPtr_ k-      pls <- fromNativeTopicPartitionList' pl+      pls <- newForeignPtr_ pl       setRebalanceCallback callback (KafkaConsumer (Kafka k') kc) (KafkaResponseError err) pls  -- | Sets a callback that is called when rebalance is needed.@@ -67,17 +67,19 @@ setRebalanceCallback :: (KafkaConsumer -> RebalanceEvent -> IO ())                           -> KafkaConsumer                           -> KafkaError-                          -> [TopicPartition] -> IO ()-setRebalanceCallback f k e ps =+                          -> RdKafkaTopicPartitionListTPtr -> IO ()+setRebalanceCallback f k e pls = do+  ps <- fromNativeTopicPartitionList'' pls   let assignment = (tpTopicName &&& tpPartition) <$> ps-  in case e of++  case e of     KafkaResponseError RdKafkaRespErrAssignPartitions -> do         mbq <- getRdMsgQueue $ getKafkaConf k         case mbq of           Nothing -> pure ()           Just mq -> forM_ ps (\tp -> redirectPartitionQueue (getKafka k) (tpTopicName tp) (tpPartition tp) mq)         f k (RebalanceBeforeAssign assignment)-        void $ assign k ps+        void $ assign' k pls -- pass as pointer to avoid possible serialisation issues         f k (RebalanceAssign assignment)     KafkaResponseError RdKafkaRespErrRevokePartitions -> do         f k (RebalanceBeforeRevoke assignment)@@ -94,3 +96,12 @@                 else toNativeTopicPartitionList ps         er = KafkaResponseError <$> (pl >>= rdKafkaAssign k)     in kafkaErrorToMaybe <$> er++-- | Assigns specified partitions to a current consumer.+-- Assigning an empty list means unassigning from all partitions that are currently assigned.+assign' :: KafkaConsumer -> RdKafkaTopicPartitionListTPtr -> IO (Maybe KafkaError)+assign' (KafkaConsumer (Kafka k) _) pls = do+    let er = KafkaResponseError <$> rdKafkaAssign k pls+    m <- kafkaErrorToMaybe <$> er+    return m+