Update formatting directive logic to unbreak tests
This commit is contained in:
committed by
Ta-Ching Chen
parent
4a8c200c37
commit
b456dec138
@@ -46,7 +46,7 @@ func makeKafkaMessageQueue(routerUrl string, mqCfg MessageQueueConfig) (MessageQ
|
|||||||
routerUrl: routerUrl,
|
routerUrl: routerUrl,
|
||||||
brokers: strings.Split(mqCfg.Url, ","),
|
brokers: strings.Split(mqCfg.Url, ","),
|
||||||
}
|
}
|
||||||
log.Infof("Created Queue ", kafka)
|
log.Infof("Created Queue %q", kafka)
|
||||||
return kafka, nil
|
return kafka, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -55,15 +55,15 @@ func isTopicValidForKafka(topic string) bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (kafka Kafka) subscribe(trigger *crd.MessageQueueTrigger) (messageQueueSubscription, error) {
|
func (kafka Kafka) subscribe(trigger *crd.MessageQueueTrigger) (messageQueueSubscription, error) {
|
||||||
log.Infof("Inside kakfa subscribe", trigger)
|
log.Infof("Inside kakfa subscribe %q", trigger)
|
||||||
log.Infof("brokers set to ", kafka.brokers)
|
log.Infof("brokers set to %q", kafka.brokers)
|
||||||
|
|
||||||
// Create new consumer
|
// Create new consumer
|
||||||
consumerConfig := cluster.NewConfig()
|
consumerConfig := cluster.NewConfig()
|
||||||
consumerConfig.Consumer.Return.Errors = true
|
consumerConfig.Consumer.Return.Errors = true
|
||||||
consumerConfig.Group.Return.Notifications = true
|
consumerConfig.Group.Return.Notifications = true
|
||||||
consumer, err := cluster.NewConsumer(kafka.brokers, string(trigger.Metadata.UID), []string{trigger.Spec.Topic}, consumerConfig)
|
consumer, err := cluster.NewConsumer(kafka.brokers, string(trigger.Metadata.UID), []string{trigger.Spec.Topic}, consumerConfig)
|
||||||
log.Infof("Created a new consumer ", consumer)
|
log.Infof("Created a new consumer: %#v", consumer)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
panic(err)
|
panic(err)
|
||||||
}
|
}
|
||||||
@@ -74,7 +74,7 @@ func (kafka Kafka) subscribe(trigger *crd.MessageQueueTrigger) (messageQueueSubs
|
|||||||
producerConfig.Producer.Retry.Max = 10
|
producerConfig.Producer.Retry.Max = 10
|
||||||
producerConfig.Producer.Return.Successes = true
|
producerConfig.Producer.Return.Successes = true
|
||||||
producer, err := sarama.NewSyncProducer(kafka.brokers, producerConfig)
|
producer, err := sarama.NewSyncProducer(kafka.brokers, producerConfig)
|
||||||
log.Infof("Created a new producer ", producer)
|
log.Infof("Created a new producer %q", producer)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
panic(err)
|
panic(err)
|
||||||
}
|
}
|
||||||
@@ -141,7 +141,7 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *crd.Me
|
|||||||
// Make the request
|
// Make the request
|
||||||
resp, err = http.DefaultClient.Do(req)
|
resp, err = http.DefaultClient.Do(req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Error invoking function for trigger %v: %v", trigger.Metadata.Name, err)
|
log.Errorf("Error invoking function for trigger %v: %v", trigger.Metadata.Name, err)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if resp == nil {
|
if resp == nil {
|
||||||
|
|||||||
@@ -130,7 +130,7 @@ func msgHandler(nats *Nats, trigger *crd.MessageQueueTrigger) func(*ns.Msg) {
|
|||||||
// Make the request
|
// Make the request
|
||||||
resp, err = http.DefaultClient.Do(req)
|
resp, err = http.DefaultClient.Do(req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Error invoking function for trigger %v: %v", trigger.Metadata.Name, err)
|
log.Errorf("Error invoking function for trigger %v: %v", trigger.Metadata.Name, err)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if resp == nil {
|
if resp == nil {
|
||||||
@@ -160,7 +160,7 @@ func msgHandler(nats *Nats, trigger *crd.MessageQueueTrigger) func(*ns.Msg) {
|
|||||||
if len(trigger.Spec.ErrorTopic) > 0 && len(body) > 0 {
|
if len(trigger.Spec.ErrorTopic) > 0 && len(body) > 0 {
|
||||||
publishErr := nats.nsConn.Publish(trigger.Spec.ErrorTopic, body)
|
publishErr := nats.nsConn.Publish(trigger.Spec.ErrorTopic, body)
|
||||||
if publishErr != nil {
|
if publishErr != nil {
|
||||||
log.Error("Failed to publish error to error topic: %v", publishErr)
|
log.Errorf("Failed to publish error to error topic: %v", publishErr)
|
||||||
// TODO: We will ack this message after max retries to prevent re-processing but
|
// TODO: We will ack this message after max retries to prevent re-processing but
|
||||||
// this may cause message loss
|
// this may cause message loss
|
||||||
}
|
}
|
||||||
|
|||||||
+1
-1
@@ -43,7 +43,7 @@ func NewClient() redis.Conn {
|
|||||||
|
|
||||||
c, err := redis.Dial("tcp", redisUrl)
|
c, err := redis.Dial("tcp", redisUrl)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Could not connect to Redis: %v\n", err)
|
log.Errorf("Could not connect to Redis: %v\n", err)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
return c
|
return c
|
||||||
|
|||||||
Reference in New Issue
Block a user