Use code-generator to generate clientset/informer/lister (#1492)
To reduce maintenance effort and avoid writing duplicate informer code, use code-generator to generate clientset/informer/lister code.
This commit is contained in:
@@ -205,7 +205,7 @@ func (asc AzureStorageConnection) subscribe(trigger *fv1.MessageQueueTrigger) (m
|
||||
asc.logger.Info("subscribing to Azure storage queue", zap.String("queue", trigger.Spec.Topic))
|
||||
|
||||
if trigger.Spec.FunctionReference.Type != types.FunctionReferenceTypeFunctionName {
|
||||
return nil, fmt.Errorf("unsupported function reference type (%v) for trigger %q", trigger.Spec.FunctionReference.Type, trigger.Metadata.Name)
|
||||
return nil, fmt.Errorf("unsupported function reference type (%v) for trigger %q", trigger.Spec.FunctionReference.Type, trigger.ObjectMeta.Name)
|
||||
}
|
||||
|
||||
subscription := &AzureQueueSubscription{
|
||||
@@ -215,7 +215,7 @@ func (asc AzureStorageConnection) subscribe(trigger *fv1.MessageQueueTrigger) (m
|
||||
// with the addition of multi-tenancy, the users can create functions in any namespace. however,
|
||||
// the triggers can only be created in the same namespace as the function.
|
||||
// so essentially, function namespace = trigger namespace.
|
||||
functionURL: asc.routerURL + "/" + strings.TrimPrefix(utils.UrlForFunction(trigger.Spec.FunctionReference.Name, trigger.Metadata.Namespace), "/"),
|
||||
functionURL: asc.routerURL + "/" + strings.TrimPrefix(utils.UrlForFunction(trigger.Spec.FunctionReference.Name, trigger.ObjectMeta.Namespace), "/"),
|
||||
contentType: trigger.Spec.ContentType,
|
||||
unsubscribe: make(chan bool),
|
||||
done: make(chan bool),
|
||||
|
||||
@@ -303,7 +303,7 @@ func TestAzureStorageQueuePoisonMessage(t *testing.T) {
|
||||
httpClient: httpClient,
|
||||
}
|
||||
subscription, err := connection.subscribe(&fv1.MessageQueueTrigger{
|
||||
Metadata: metav1.ObjectMeta{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: TriggerName,
|
||||
Namespace: metav1.NamespaceDefault,
|
||||
},
|
||||
@@ -451,7 +451,7 @@ func runAzureStorageQueueTest(t *testing.T, count int, output bool) {
|
||||
httpClient: httpClient,
|
||||
}
|
||||
subscription, err := connection.subscribe(&fv1.MessageQueueTrigger{
|
||||
Metadata: metav1.ObjectMeta{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: TriggerName,
|
||||
Namespace: metav1.NamespaceDefault,
|
||||
},
|
||||
|
||||
@@ -124,13 +124,13 @@ func (kafka Kafka) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubs
|
||||
consumerConfig.Net.TLS.Config = tlsConfig
|
||||
}
|
||||
|
||||
consumer, err := cluster.NewConsumer(kafka.brokers, string(trigger.Metadata.UID), []string{trigger.Spec.Topic}, consumerConfig)
|
||||
consumer, err := cluster.NewConsumer(kafka.brokers, string(trigger.ObjectMeta.UID), []string{trigger.Spec.Topic}, consumerConfig)
|
||||
kafka.logger.Info("created a new consumer", zap.Strings("brokers", kafka.brokers),
|
||||
zap.String("input topic", trigger.Spec.Topic),
|
||||
zap.String("output topic", trigger.Spec.ResponseTopic),
|
||||
zap.String("error topic", trigger.Spec.ErrorTopic),
|
||||
zap.String("trigger name", trigger.Metadata.Name),
|
||||
zap.String("function namespace", trigger.Metadata.Namespace),
|
||||
zap.String("trigger name", trigger.ObjectMeta.Name),
|
||||
zap.String("function namespace", trigger.ObjectMeta.Namespace),
|
||||
zap.String("function name", trigger.Spec.FunctionReference.Name))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -141,8 +141,8 @@ func (kafka Kafka) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubs
|
||||
zap.String("input topic", trigger.Spec.Topic),
|
||||
zap.String("output topic", trigger.Spec.ResponseTopic),
|
||||
zap.String("error topic", trigger.Spec.ErrorTopic),
|
||||
zap.String("trigger name", trigger.Metadata.Name),
|
||||
zap.String("function namespace", trigger.Metadata.Namespace),
|
||||
zap.String("trigger name", trigger.ObjectMeta.Name),
|
||||
zap.String("function namespace", trigger.ObjectMeta.Namespace),
|
||||
zap.String("function name", trigger.Spec.FunctionReference.Name))
|
||||
|
||||
if err != nil {
|
||||
@@ -205,10 +205,10 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.Me
|
||||
if trigger.Spec.FunctionReference.Type != types.FunctionReferenceTypeFunctionName {
|
||||
kafka.logger.Fatal("unsupported function reference type for trigger",
|
||||
zap.Any("function_reference_type", trigger.Spec.FunctionReference.Type),
|
||||
zap.String("trigger", trigger.Metadata.Name))
|
||||
zap.String("trigger", trigger.ObjectMeta.Name))
|
||||
}
|
||||
|
||||
url := kafka.routerUrl + "/" + strings.TrimPrefix(utils.UrlForFunction(trigger.Spec.FunctionReference.Name, trigger.Metadata.Namespace), "/")
|
||||
url := kafka.routerUrl + "/" + strings.TrimPrefix(utils.UrlForFunction(trigger.Spec.FunctionReference.Name, trigger.ObjectMeta.Namespace), "/")
|
||||
kafka.logger.Debug("making HTTP request", zap.String("url", url))
|
||||
|
||||
// Generate the Headers
|
||||
@@ -252,7 +252,7 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.Me
|
||||
kafka.logger.Error("sending function invocation request failed",
|
||||
zap.Error(err),
|
||||
zap.String("function_url", url),
|
||||
zap.String("trigger", trigger.Metadata.Name))
|
||||
zap.String("trigger", trigger.ObjectMeta.Name))
|
||||
continue
|
||||
}
|
||||
if resp == nil {
|
||||
@@ -267,7 +267,7 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.Me
|
||||
if resp == nil {
|
||||
kafka.logger.Warn("every function invocation retry failed; final retry gave empty response",
|
||||
zap.String("function_url", url),
|
||||
zap.String("trigger", trigger.Metadata.Name))
|
||||
zap.String("trigger", trigger.ObjectMeta.Name))
|
||||
return
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
@@ -275,7 +275,7 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.Me
|
||||
|
||||
kafka.logger.Debug("got response from function invocation",
|
||||
zap.String("function_url", url),
|
||||
zap.String("trigger", trigger.Metadata.Name),
|
||||
zap.String("trigger", trigger.ObjectMeta.Name),
|
||||
zap.String("body", string(body)))
|
||||
|
||||
if err != nil {
|
||||
@@ -328,12 +328,12 @@ func errorHandler(logger *zap.Logger, trigger *fv1.MessageQueueTrigger, producer
|
||||
if e != nil {
|
||||
logger.Error("failed to publish message to error topic",
|
||||
zap.Error(e),
|
||||
zap.String("trigger", trigger.Metadata.Name),
|
||||
zap.String("trigger", trigger.ObjectMeta.Name),
|
||||
zap.String("message", err.Error()),
|
||||
zap.String("topic", trigger.Spec.Topic))
|
||||
}
|
||||
} else {
|
||||
logger.Error("message received to publish to error topic, but no error topic was set",
|
||||
zap.String("message", err.Error()), zap.String("trigger", trigger.Metadata.Name), zap.String("function_url", funcUrl))
|
||||
zap.String("message", err.Error()), zap.String("trigger", trigger.ObjectMeta.Name), zap.String("function_url", funcUrl))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -111,7 +111,7 @@ func (mqt *MessageQueueTriggerManager) service() {
|
||||
switch req.requestType {
|
||||
case ADD_TRIGGER:
|
||||
var err error
|
||||
k := crd.CacheKey(&req.triggerSub.trigger.Metadata)
|
||||
k := crd.CacheKey(&req.triggerSub.trigger.ObjectMeta)
|
||||
if _, ok := mqt.triggers[k]; ok {
|
||||
err = errors.New("trigger already exists")
|
||||
} else {
|
||||
@@ -125,7 +125,7 @@ func (mqt *MessageQueueTriggerManager) service() {
|
||||
}
|
||||
req.respChan <- response{triggers: ©Triggers}
|
||||
case DELETE_TRIGGER:
|
||||
delete(mqt.triggers, crd.CacheKey(&req.triggerSub.trigger.Metadata))
|
||||
delete(mqt.triggers, crd.CacheKey(&req.triggerSub.trigger.ObjectMeta))
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -156,7 +156,7 @@ func (mqt *MessageQueueTriggerManager) delTrigger(m *metav1.ObjectMeta) {
|
||||
requestType: DELETE_TRIGGER,
|
||||
triggerSub: &triggerSubscription{
|
||||
trigger: fv1.MessageQueueTrigger{
|
||||
Metadata: *m,
|
||||
ObjectMeta: *m,
|
||||
},
|
||||
},
|
||||
}
|
||||
@@ -165,7 +165,7 @@ func (mqt *MessageQueueTriggerManager) delTrigger(m *metav1.ObjectMeta) {
|
||||
func (mqt *MessageQueueTriggerManager) syncTriggers() {
|
||||
for {
|
||||
// get new set of triggers
|
||||
newTriggers, err := mqt.fissionClient.MessageQueueTriggers(metav1.NamespaceAll).List(metav1.ListOptions{})
|
||||
newTriggers, err := mqt.fissionClient.V1().MessageQueueTriggers(metav1.NamespaceAll).List(metav1.ListOptions{})
|
||||
if err != nil {
|
||||
if utils.IsNetworkError(err) {
|
||||
mqt.logger.Error("encountered network error, will retry", zap.Error(err))
|
||||
@@ -177,7 +177,7 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() {
|
||||
newTriggerMap := make(map[string]*fv1.MessageQueueTrigger)
|
||||
for index := range newTriggers.Items {
|
||||
newTrigger := &newTriggers.Items[index]
|
||||
newTriggerMap[crd.CacheKey(&newTrigger.Metadata)] = newTrigger
|
||||
newTriggerMap[crd.CacheKey(&newTrigger.ObjectMeta)] = newTrigger
|
||||
}
|
||||
|
||||
// get current set of triggers
|
||||
@@ -192,7 +192,7 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() {
|
||||
// actually subscribe using the message queue client impl
|
||||
sub, err := mqt.messageQueue.subscribe(trigger)
|
||||
if err != nil {
|
||||
mqt.logger.Warn("failed to subscribe to message queue trigger", zap.Error(err), zap.String("trigger_name", trigger.Metadata.Name))
|
||||
mqt.logger.Warn("failed to subscribe to message queue trigger", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -204,10 +204,10 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() {
|
||||
// add to our list
|
||||
err = mqt.addTrigger(&triggerSub)
|
||||
if err != nil {
|
||||
mqt.logger.Fatal("adding message queue trigger failed", zap.Error(err), zap.String("trigger_name", trigger.Metadata.Name))
|
||||
mqt.logger.Fatal("adding message queue trigger failed", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
}
|
||||
|
||||
mqt.logger.Info("message queue trigger created", zap.String("trigger_name", trigger.Metadata.Name))
|
||||
mqt.logger.Info("message queue trigger created", zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
}
|
||||
|
||||
// remove old triggers
|
||||
@@ -217,11 +217,11 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() {
|
||||
}
|
||||
err := mqt.messageQueue.unsubscribe(triggerSub.subscription)
|
||||
if err != nil {
|
||||
mqt.logger.Warn("failed to unsubscribe from message queue trigger", zap.Error(err), zap.String("trigger_name", triggerSub.trigger.Metadata.Name))
|
||||
mqt.logger.Warn("failed to unsubscribe from message queue trigger", zap.Error(err), zap.String("trigger_name", triggerSub.trigger.ObjectMeta.Name))
|
||||
continue
|
||||
}
|
||||
mqt.delTrigger(&triggerSub.trigger.Metadata)
|
||||
mqt.logger.Info("message queue trigger deleted", zap.String("trigger_name", triggerSub.trigger.Metadata.Name))
|
||||
mqt.delTrigger(&triggerSub.trigger.ObjectMeta)
|
||||
mqt.logger.Info("message queue trigger deleted", zap.String("trigger_name", triggerSub.trigger.ObjectMeta.Name))
|
||||
}
|
||||
|
||||
// TODO replace with a watch
|
||||
|
||||
@@ -79,7 +79,7 @@ func (nats Nats) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubscr
|
||||
opts := []ns.SubscriptionOption{
|
||||
// Create a durable subscription to nats, so that triggers could retrieve last unack message.
|
||||
// https://github.com/nats-io/go-nats-streaming#durable-subscriptions
|
||||
ns.DurableName(string(trigger.Metadata.UID)),
|
||||
ns.DurableName(string(trigger.ObjectMeta.UID)),
|
||||
|
||||
// Nats-streaming server is auto-ack mode by default. Since we want nats-streaming server to
|
||||
// resend a message if the trigger does not ack it, we need to enable the manual ack mode, so that
|
||||
@@ -109,13 +109,13 @@ func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) {
|
||||
if trigger.Spec.FunctionReference.Type != types.FunctionReferenceTypeFunctionName {
|
||||
nats.logger.Fatal("unsupported function reference type for trigger",
|
||||
zap.Any("function_reference_type", trigger.Spec.FunctionReference.Type),
|
||||
zap.String("trigger", trigger.Metadata.Name))
|
||||
zap.String("trigger", trigger.ObjectMeta.Name))
|
||||
}
|
||||
|
||||
// with the addition of multi-tenancy, the users can create functions in any namespace. however,
|
||||
// the triggers can only be created in the same namespace as the function.
|
||||
// so essentially, function namespace = trigger namespace.
|
||||
url := nats.routerUrl + "/" + strings.TrimPrefix(utils.UrlForFunction(trigger.Spec.FunctionReference.Name, trigger.Metadata.Namespace), "/")
|
||||
url := nats.routerUrl + "/" + strings.TrimPrefix(utils.UrlForFunction(trigger.Spec.FunctionReference.Name, trigger.ObjectMeta.Namespace), "/")
|
||||
nats.logger.Debug("making HTTP request", zap.String("url", url))
|
||||
|
||||
headers := map[string]string{
|
||||
@@ -147,7 +147,7 @@ func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) {
|
||||
nats.logger.Error("sending function invocation request failed",
|
||||
zap.Error(err),
|
||||
zap.String("function_url", url),
|
||||
zap.String("trigger", trigger.Metadata.Name))
|
||||
zap.String("trigger", trigger.ObjectMeta.Name))
|
||||
continue
|
||||
}
|
||||
if resp == nil {
|
||||
@@ -162,7 +162,7 @@ func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) {
|
||||
if resp == nil {
|
||||
nats.logger.Warn("every function invocation retry failed; final retry gave empty response",
|
||||
zap.String("function_url", url),
|
||||
zap.String("trigger", trigger.Metadata.Name))
|
||||
zap.String("trigger", trigger.ObjectMeta.Name))
|
||||
return
|
||||
}
|
||||
|
||||
@@ -173,7 +173,7 @@ func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) {
|
||||
nats.logger.Error("error reading function invocation response",
|
||||
zap.Error(err),
|
||||
zap.String("function_url", url),
|
||||
zap.String("trigger", trigger.Metadata.Name))
|
||||
zap.String("trigger", trigger.ObjectMeta.Name))
|
||||
return
|
||||
}
|
||||
|
||||
@@ -186,7 +186,7 @@ func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) {
|
||||
zap.Error(publishErr),
|
||||
zap.String("topic", trigger.Spec.ErrorTopic),
|
||||
zap.String("function_url", url),
|
||||
zap.String("trigger", trigger.Metadata.Name))
|
||||
zap.String("trigger", trigger.ObjectMeta.Name))
|
||||
// TODO: We will ack this message after max retries to prevent re-processing but
|
||||
// this may cause message loss
|
||||
}
|
||||
@@ -200,7 +200,7 @@ func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) {
|
||||
nats.logger.Error("failed to ack message after successful function invocation from trigger",
|
||||
zap.Error(err),
|
||||
zap.String("function_url", url),
|
||||
zap.String("trigger", trigger.Metadata.Name))
|
||||
zap.String("trigger", trigger.ObjectMeta.Name))
|
||||
}
|
||||
|
||||
if len(trigger.Spec.ResponseTopic) > 0 {
|
||||
@@ -209,7 +209,7 @@ func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) {
|
||||
nats.logger.Error("failed to publish message with function invocation response to topic",
|
||||
zap.Error(err),
|
||||
zap.String("topic", trigger.Spec.ResponseTopic),
|
||||
zap.String("trigger", trigger.Metadata.Name))
|
||||
zap.String("trigger", trigger.ObjectMeta.Name))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user