diff --git a/cmd/flagger/main.go b/cmd/flagger/main.go index 7d9df913..69a4162e 100644 --- a/cmd/flagger/main.go +++ b/cmd/flagger/main.go @@ -45,6 +45,7 @@ var ( slackURL string slackUser string slackChannel string + eventWebhook string threadiness int zapReplaceGlobals bool zapEncoding string @@ -67,6 +68,7 @@ func init() { flag.StringVar(&slackURL, "slack-url", "", "Slack hook URL.") flag.StringVar(&slackUser, "slack-user", "flagger", "Slack user name.") flag.StringVar(&slackChannel, "slack-channel", "", "Slack channel.") + flag.StringVar(&eventWebhook, "event-webhook", "", "Webhook for publishing flagger events") flag.StringVar(&msteamsURL, "msteams-url", "", "MS Teams incoming webhook URL.") flag.IntVar(&threadiness, "threadiness", 2, "Worker concurrency.") flag.BoolVar(&zapReplaceGlobals, "zap-replace-globals", false, "Whether to change the logging level of the global zap logger.") @@ -200,6 +202,7 @@ func main() { observerFactory, meshProvider, version.VERSION, + eventWebhook, ) flaggerInformerFactory.Start(stopCh) diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 574e6dbe..601dab39 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -51,6 +51,7 @@ type Controller struct { routerFactory *router.Factory observerFactory *metrics.Factory meshProvider string + eventWebhook string } func NewController( @@ -66,6 +67,7 @@ func NewController( observerFactory *metrics.Factory, meshProvider string, version string, + eventWebhook string, ) *Controller { logger.Debug("Creating event broadcaster") flaggerscheme.AddToScheme(scheme.Scheme) @@ -97,6 +99,7 @@ func NewController( canaryFactory: canaryFactory, routerFactory: routerFactory, meshProvider: meshProvider, + eventWebhook: eventWebhook, } flaggerInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ @@ -244,19 +247,32 @@ func checkCustomResourceType(obj interface{}, logger *zap.SugaredLogger) (flagge return *roll, true } +func (c *Controller) sendEventToWebhook(r *flaggerv1.Canary, template string, args []interface{}) { + if c.eventWebhook != "" { + c.logger.Info(fmt.Sprintf(template, args...)) + err := CallEventWebhook(r, c.eventWebhook, fmt.Sprintf(template, args...), corev1.EventTypeNormal) + if err != nil { + c.logger.With("canary", fmt.Sprintf("%s.%s", r.Name, r.Namespace)).Errorf("error sending event to webhook: %s", err) + } + } +} + func (c *Controller) recordEventInfof(r *flaggerv1.Canary, template string, args ...interface{}) { c.logger.With("canary", fmt.Sprintf("%s.%s", r.Name, r.Namespace)).Infof(template, args...) c.eventRecorder.Event(r, corev1.EventTypeNormal, "Synced", fmt.Sprintf(template, args...)) + c.sendEventToWebhook(r, template, args) } func (c *Controller) recordEventErrorf(r *flaggerv1.Canary, template string, args ...interface{}) { c.logger.With("canary", fmt.Sprintf("%s.%s", r.Name, r.Namespace)).Errorf(template, args...) c.eventRecorder.Event(r, corev1.EventTypeWarning, "Synced", fmt.Sprintf(template, args...)) + c.sendEventToWebhook(r, template, args) } func (c *Controller) recordEventWarningf(r *flaggerv1.Canary, template string, args ...interface{}) { c.logger.With("canary", fmt.Sprintf("%s.%s", r.Name, r.Namespace)).Infof(template, args...) c.eventRecorder.Event(r, corev1.EventTypeWarning, "Synced", fmt.Sprintf(template, args...)) + c.sendEventToWebhook(r, template, args) } func (c *Controller) sendNotification(cd *flaggerv1.Canary, message string, metadata bool, warn bool) { diff --git a/pkg/controller/webhook.go b/pkg/controller/webhook.go index ff144095..0eeef6b6 100644 --- a/pkg/controller/webhook.go +++ b/pkg/controller/webhook.go @@ -7,6 +7,11 @@ import ( "errors" "fmt" "io/ioutil" + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes/scheme" + "k8s.io/client-go/tools/reference" + "k8s.io/utils/clock" "net/http" "net/url" "time" @@ -14,25 +19,13 @@ import ( flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" ) -// CallWebhook does a HTTP POST to an external service and -// returns an error if the response status code is non-2xx -func CallWebhook(name string, namespace string, phase flaggerv1.CanaryPhase, w flaggerv1.CanaryWebhook) error { - payload := flaggerv1.CanaryWebhookPayload{ - Name: name, - Namespace: namespace, - Phase: phase, - } - - if w.Metadata != nil { - payload.Metadata = *w.Metadata - } - +func callWebhook(webhook string, payload interface{}, timeout string) error { payloadBin, err := json.Marshal(payload) if err != nil { return err } - hook, err := url.Parse(w.URL) + hook, err := url.Parse(webhook) if err != nil { return err } @@ -44,16 +37,16 @@ func CallWebhook(name string, namespace string, phase flaggerv1.CanaryPhase, w f req.Header.Set("Content-Type", "application/json") - if len(w.Timeout) < 2 { - w.Timeout = "10s" + if timeout == "" { + timeout = "10s" } - timeout, err := time.ParseDuration(w.Timeout) + t, err := time.ParseDuration(timeout) if err != nil { return err } - ctx, cancel := context.WithTimeout(req.Context(), timeout) + ctx, cancel := context.WithTimeout(req.Context(), t) defer cancel() r, err := http.DefaultClient.Do(req.WithContext(ctx)) @@ -73,3 +66,52 @@ func CallWebhook(name string, namespace string, phase flaggerv1.CanaryPhase, w f return nil } + +// CallWebhook does a HTTP POST to an external service and +// returns an error if the response status code is non-2xx +func CallWebhook(name string, namespace string, phase flaggerv1.CanaryPhase, w flaggerv1.CanaryWebhook) error { + payload := flaggerv1.CanaryWebhookPayload{ + Name: name, + Namespace: namespace, + Phase: phase, + } + + if w.Metadata != nil { + payload.Metadata = *w.Metadata + } + + if len(w.Timeout) < 2 { + w.Timeout = "10s" + } + + return callWebhook(w.URL, payload, w.Timeout) +} + +func CallEventWebhook(r *flaggerv1.Canary, webhook, message, eventtype string) error { + t := clock.RealClock{}.Now() + ref, err := reference.GetReference(scheme.Scheme, r) + if err != nil { + return err + } + + namespace := ref.Namespace + if namespace == "" { + namespace = metav1.NamespaceDefault + } + + event := v1.Event{ + ObjectMeta: metav1.ObjectMeta{ + Name: fmt.Sprintf("%v.%x", ref.Name, t.UnixNano()), + Namespace: namespace, + }, + InvolvedObject: *ref, + Reason: "Synced", + Message: message, + Source: v1.EventSource{ + Component: "flagger", + }, + Type: eventtype, + } + + return callWebhook(webhook, event, "5s") +}