milena 0.5.2.4 → 0.5.3.0
raw patch · 4 files changed
+52/−4 lines, 4 filesPVP ok
version bump matches the API change (PVP)
API changes (from Hackage documentation)
+ Network.Kafka: createTopic :: Kafka m => CreateTopicsRequest -> m CreateTopicsResponse
+ Network.Kafka: createTopic' :: Kafka m => Handle -> CreateTopicsRequest -> m CreateTopicsResponse
+ Network.Kafka: createTopicsRequest :: TopicName -> Partition -> ReplicationFactor -> [(Partition, Replicas)] -> [(KafkaString, Metadata)] -> CreateTopicsRequest
+ Network.Kafka.Protocol: CreateTopicsReq :: ([(TopicName, Partition, ReplicationFactor, [(Partition, Replicas)], [(KafkaString, Metadata)])], Timeout) -> CreateTopicsRequest
+ Network.Kafka.Protocol: CreateTopicsRequest :: CreateTopicsRequest -> RequestMessage
+ Network.Kafka.Protocol: CreateTopicsResponse :: CreateTopicsResponse -> ResponseMessage
+ Network.Kafka.Protocol: ReplicationFactor :: Int16 -> ReplicationFactor
+ Network.Kafka.Protocol: TopicAlreadyExists :: KafkaError
+ Network.Kafka.Protocol: TopicsResp :: [(TopicName, KafkaError)] -> CreateTopicsResponse
+ Network.Kafka.Protocol: UnsupportedCompressionType :: KafkaError
+ Network.Kafka.Protocol: [TopicsRR] :: MonadIO m => CreateTopicsRequest -> ReqResp (m CreateTopicsResponse)
+ Network.Kafka.Protocol: [_topicsResponseFields] :: CreateTopicsResponse -> [(TopicName, KafkaError)]
+ Network.Kafka.Protocol: _CreateTopicsResponse :: Prism' ResponseMessage CreateTopicsResponse
+ Network.Kafka.Protocol: instance GHC.Classes.Eq Network.Kafka.Protocol.CreateTopicsRequest
+ Network.Kafka.Protocol: instance GHC.Classes.Eq Network.Kafka.Protocol.CreateTopicsResponse
+ Network.Kafka.Protocol: instance GHC.Classes.Eq Network.Kafka.Protocol.ReplicationFactor
+ Network.Kafka.Protocol: instance GHC.Classes.Ord Network.Kafka.Protocol.ReplicationFactor
+ Network.Kafka.Protocol: instance GHC.Enum.Enum Network.Kafka.Protocol.ReplicationFactor
+ Network.Kafka.Protocol: instance GHC.Generics.Generic Network.Kafka.Protocol.CreateTopicsRequest
+ Network.Kafka.Protocol: instance GHC.Generics.Generic Network.Kafka.Protocol.CreateTopicsResponse
+ Network.Kafka.Protocol: instance GHC.Generics.Generic Network.Kafka.Protocol.ReplicationFactor
+ Network.Kafka.Protocol: instance GHC.Num.Num Network.Kafka.Protocol.ReplicationFactor
+ Network.Kafka.Protocol: instance GHC.Real.Integral Network.Kafka.Protocol.ReplicationFactor
+ Network.Kafka.Protocol: instance GHC.Real.Real Network.Kafka.Protocol.ReplicationFactor
+ Network.Kafka.Protocol: instance GHC.Show.Show Network.Kafka.Protocol.CreateTopicsRequest
+ Network.Kafka.Protocol: instance GHC.Show.Show Network.Kafka.Protocol.CreateTopicsResponse
+ Network.Kafka.Protocol: instance GHC.Show.Show Network.Kafka.Protocol.ReplicationFactor
+ Network.Kafka.Protocol: instance Network.Kafka.Protocol.Deserializable Network.Kafka.Protocol.CreateTopicsResponse
+ Network.Kafka.Protocol: instance Network.Kafka.Protocol.Deserializable Network.Kafka.Protocol.ReplicationFactor
+ Network.Kafka.Protocol: instance Network.Kafka.Protocol.Serializable Network.Kafka.Protocol.CreateTopicsRequest
+ Network.Kafka.Protocol: instance Network.Kafka.Protocol.Serializable Network.Kafka.Protocol.CreateTopicsResponse
+ Network.Kafka.Protocol: instance Network.Kafka.Protocol.Serializable Network.Kafka.Protocol.ReplicationFactor
+ Network.Kafka.Protocol: newtype CreateTopicsRequest
+ Network.Kafka.Protocol: newtype CreateTopicsResponse
+ Network.Kafka.Protocol: newtype ReplicationFactor
+ Network.Kafka.Protocol: topicsResponseFields :: Iso' CreateTopicsResponse [(TopicName, KafkaError)]
Files
- Network/Kafka.hs +20/−0
- Network/Kafka/Protocol.hs +21/−1
- milena.cabal +2/−2
- test/tests.hs +9/−1
Network/Kafka.hs view
@@ -183,6 +183,26 @@ metadata' :: Kafka m => Handle -> MetadataRequest -> m MetadataResponse metadata' h request = makeRequest h $ MetadataRR request ++createTopic :: Kafka m => CreateTopicsRequest -> m CreateTopicsResponse+createTopic request = withAnyHandle $ flip createTopic' request++createTopic' ::+ Kafka m => Handle -> CreateTopicsRequest -> m CreateTopicsResponse+createTopic' h request = makeRequest h $ TopicsRR request++createTopicsRequest ::+ TopicName+ -> Partition+ -> ReplicationFactor+ -> [(Partition, Replicas)]+ -> [(KafkaString, Metadata)]+ -> CreateTopicsRequest+createTopicsRequest topic partition replication_factor replica_assignment config =+ CreateTopicsReq+ ([(topic, partition, replication_factor, replica_assignment, config)], defaultRequestTimeout)++ getTopicPartitionLeader :: Kafka m => TopicName -> Partition -> m Broker getTopicPartitionLeader t p = do let s = stateTopicMetadata . at t
Network/Kafka/Protocol.hs view
@@ -30,6 +30,7 @@ ProduceRR :: MonadIO m => ProduceRequest -> ReqResp (m ProduceResponse) FetchRR :: MonadIO m => FetchRequest -> ReqResp (m FetchResponse) OffsetRR :: MonadIO m => OffsetRequest -> ReqResp (m OffsetResponse)+ TopicsRR :: MonadIO m => CreateTopicsRequest -> ReqResp (m CreateTopicsResponse) doRequest' :: (Deserializable a, MonadIO m) => CorrelationId -> Handle -> Request -> m (Either String a) doRequest' correlationId h r = do@@ -51,6 +52,7 @@ doRequest clientId correlationId h (ProduceRR req) = doRequest' correlationId h $ Request (correlationId, clientId, ProduceRequest req) doRequest clientId correlationId h (FetchRR req) = doRequest' correlationId h $ Request (correlationId, clientId, FetchRequest req) doRequest clientId correlationId h (OffsetRR req) = doRequest' correlationId h $ Request (correlationId, clientId, OffsetRequest req)+doRequest clientId correlationId h (TopicsRR req) = doRequest' correlationId h $ Request (correlationId, clientId, CreateTopicsRequest req) class Serializable a where serialize :: a -> Put@@ -72,6 +74,7 @@ | OffsetCommitRequest OffsetCommitRequest | OffsetFetchRequest OffsetFetchRequest | GroupCoordinatorRequest GroupCoordinatorRequest+ | CreateTopicsRequest CreateTopicsRequest deriving (Show, Generic, Eq) newtype MetadataRequest = MetadataReq [TopicName] deriving (Show, Eq, Serializable, Generic, Deserializable)@@ -99,6 +102,10 @@ FetchResp { _fetchResponseFields :: [(TopicName, [(Partition, KafkaError, Offset, MessageSet)])] } deriving (Show, Eq, Serializable, Deserializable, Generic) +newtype CreateTopicsResponse =+ TopicsResp { _topicsResponseFields :: [(TopicName, KafkaError)] }+ deriving (Show, Eq, Deserializable, Serializable, Generic)+ newtype MetadataResponse = MetadataResp { _metadataResponseFields :: ([Broker], [TopicMetadata]) } deriving (Show, Eq, Deserializable, Generic) newtype Broker = Broker { _brokerFields :: (NodeId, Host, Port) } deriving (Show, Eq, Ord, Deserializable, Generic) newtype NodeId = NodeId { _nodeId :: Int32 } deriving (Show, Eq, Deserializable, Num, Integral, Ord, Real, Enum, Generic)@@ -167,10 +174,13 @@ | OffsetCommitResponse OffsetCommitResponse | OffsetFetchResponse OffsetFetchResponse | GroupCoordinatorResponse GroupCoordinatorResponse+ | CreateTopicsResponse CreateTopicsResponse deriving (Show, Eq, Generic) -newtype GroupCoordinatorRequest = GroupCoordinatorReq ConsumerGroup deriving (Show, Eq, Serializable, Generic)+newtype ReplicationFactor = ReplicationFactor Int16 deriving (Show, Eq, Num, Integral, Ord, Real, Enum, Serializable, Deserializable, Generic) +newtype GroupCoordinatorRequest = GroupCoordinatorReq ConsumerGroup deriving (Show, Eq, Serializable, Generic)+newtype CreateTopicsRequest = CreateTopicsReq ([(TopicName, Partition, ReplicationFactor, [(Partition, Replicas)], [(KafkaString, Metadata)])], Timeout) deriving (Show, Eq, Serializable, Generic) newtype OffsetCommitRequest = OffsetCommitReq (ConsumerGroup, [(TopicName, [(Partition, Offset, Time, Metadata)])]) deriving (Show, Eq, Serializable, Generic) newtype OffsetFetchRequest = OffsetFetchReq (ConsumerGroup, [(TopicName, [Partition])]) deriving (Show, Eq, Serializable, Generic) newtype ConsumerGroup = ConsumerGroup KafkaString deriving (Show, Eq, Serializable, Deserializable, IsString, Generic)@@ -194,6 +204,8 @@ errorKafka OffsetsLoadInProgressCode = 14 errorKafka ConsumerCoordinatorNotAvailableCode = 15 errorKafka NotCoordinatorForConsumerCode = 16+errorKafka TopicAlreadyExists = 36+errorKafka UnsupportedCompressionType = 76 data KafkaError = NoError -- ^ @0@ No error--it worked! | Unknown -- ^ @-1@ An unexpected server error@@ -212,6 +224,8 @@ | OffsetsLoadInProgressCode -- ^ @14@ The broker returns this error code for an offset fetch request if it is still loading offsets (after a leader change for that offsets topic partition). | ConsumerCoordinatorNotAvailableCode -- ^ @15@ The broker returns this error code for consumer metadata requests or offset commit requests if the offsets topic has not yet been created. | NotCoordinatorForConsumerCode -- ^ @16@ The broker returns this error code if it receives an offset fetch or commit request for a consumer group that it is not a coordinator for.+ | TopicAlreadyExists -- ^@36@ Topic with this name already exists.+ | UnsupportedCompressionType -- ^@76@ The requesting client does not support the compression type of given partition. deriving (Bounded, Enum, Eq, Generic, Show) instance Serializable KafkaError where@@ -238,6 +252,8 @@ 14 -> return OffsetsLoadInProgressCode 15 -> return ConsumerCoordinatorNotAvailableCode 16 -> return NotCoordinatorForConsumerCode+ 36 -> return TopicAlreadyExists+ 76 -> return UnsupportedCompressionType _ -> fail $ "invalid error code: " ++ show x instance Exception KafkaError@@ -269,6 +285,7 @@ apiKey OffsetCommitRequest{} = ApiKey 8 apiKey OffsetFetchRequest{} = ApiKey 9 apiKey GroupCoordinatorRequest{} = ApiKey 10+apiKey CreateTopicsRequest{} = ApiKey 19 instance Serializable RequestMessage where serialize (ProduceRequest r) = serialize r@@ -278,6 +295,7 @@ serialize (OffsetCommitRequest r) = serialize r serialize (OffsetFetchRequest r) = serialize r serialize (GroupCoordinatorRequest r) = serialize r+ serialize (CreateTopicsRequest r) = serialize r instance Serializable Int64 where serialize = putWord64be . fromIntegral instance Serializable Int32 where serialize = putWord32be . fromIntegral@@ -523,6 +541,8 @@ makeLenses ''Key makeLenses ''Value++makeLenses ''CreateTopicsResponse makePrisms ''ResponseMessage
milena.cabal view
@@ -4,10 +4,10 @@ -- -- see: https://github.com/sol/hpack ----- hash: 3bb58d95694565345dd96c8f19c605682638205f778b92ae231d3aa2a0ef69e1+-- hash: 8060a8b3f6bf4aad48a041ce65dd55877fe0f1bd5353f32359bae70b010f2843 name: milena-version: 0.5.2.4+version: 0.5.3.0 synopsis: A Kafka client for Haskell. description: A Kafka client for Haskell. The protocol module is stable (the only changes will be to support changes in the Kafka protocol). The API is functional but subject to change.
test/tests.hs view
@@ -11,7 +11,7 @@ import Network.Kafka import Network.Kafka.Consumer import Network.Kafka.Producer-import Network.Kafka.Protocol (ProduceResponse(..), KafkaError(..), CompressionCodec(..))+import Network.Kafka.Protocol (ProduceResponse(..), KafkaError(..), CompressionCodec(..), CreateTopicsResponse(..)) import Test.Tasty import Test.Tasty.Hspec import Test.Tasty.QuickCheck@@ -128,6 +128,14 @@ updateMetadatas [] use stateAddresses result `shouldBe` fmap NE.nub result++ describe "create topics" $+ it "create topics with multiple partitions" $ do+ let t = "milena-test-13-partitions"+ result <- run $ do+ stateAddresses %= NE.cons ("localhost", 9092)+ createTopic (createTopicsRequest t 13 1 [] [])+ result `shouldBe` (Right $ TopicsResp [(t, NoError)]) prop :: Testable prop => String -> prop -> SpecWith () prop s = it s . property