diff --git a/Makefile b/Makefile index f800868b..437024f3 100644 --- a/Makefile +++ b/Makefile @@ -4,7 +4,10 @@ VERSION_MINOR:=$(shell grep 'VERSION' pkg/version/version.go | awk '{ print $$4 PATCH:=$(shell grep 'VERSION' pkg/version/version.go | awk '{ print $$4 }' | tr -d '"' | awk -F. '{print $$NF}') SOURCE_DIRS = cmd pkg/apis pkg/controller pkg/server pkg/logging pkg/version run: - go run cmd/flagger/* -kubeconfig=$$HOME/.kube/config -log-level=info -metrics-server=https://prometheus.iowa.weavedx.com + go run cmd/flagger/* -kubeconfig=$$HOME/.kube/config -log-level=info \ + -metrics-server=https://prometheus.iowa.weavedx.com \ + -slack-url=https://hooks.slack.com/services/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 \ + -slack-channel="devops-alerts" build: docker build -t stefanprodan/flagger:$(TAG) . -f Dockerfile diff --git a/README.md b/README.md index e3377f1d..0db1606d 100644 --- a/README.md +++ b/README.md @@ -351,13 +351,28 @@ flagger_canary_duration_seconds_sum{name="podinfo",namespace="test"} 17.3561329 flagger_canary_duration_seconds_count{name="podinfo",namespace="test"} 6 ``` +### Alerting + +Flagger can be configured to send Slack notifications: + +```bash +helm upgrade -i flagger flagger/flagger \ +--namespace=istio-system \ +--set slack.url=https://hooks.slack.com/services/YOUR/SLACK/WEBHOOK \ +--set slack.channel=general \ +--set slack.user=flagger +``` + +Once configured with a Slack incoming webhook, Flagger will post messages when a canary deployment has been initialized, +when a new revision has been detected and if the canary analysis failed or succeeded. + +![flagger-slack](https://raw.githubusercontent.com/stefanprodan/flagger/master/docs/screens/slack-notifications.png) + ### Roadmap * Extend the canary analysis and promotion to other types than Kubernetes deployments such as Flux Helm releases or OpenFaaS functions * Extend the validation mechanism to support other metrics than HTTP success rate and latency * Add support for comparing the canary metrics to the primary ones and do the validation based on the derivation between the two -* Alerting: trigger Alertmanager on successful or failed promotions -* Reporting: publish canary analysis results to Slack/Jira/etc ### Contributing diff --git a/charts/flagger/templates/deployment.yaml b/charts/flagger/templates/deployment.yaml index 4d9a09e2..f35924bd 100644 --- a/charts/flagger/templates/deployment.yaml +++ b/charts/flagger/templates/deployment.yaml @@ -37,6 +37,11 @@ spec: - -log-level=info - -control-loop-interval={{ .Values.controlLoopInterval }} - -metrics-server={{ .Values.metricsServer }} + {{- if .Values.slack.url }} + - -slack-url={{ .Values.slack.url }} + - -slack-user={{ .Values.slack.user }} + - -slack-channel={{ .Values.slack.channel }} + {{- end }} livenessProbe: exec: command: diff --git a/charts/flagger/values.yaml b/charts/flagger/values.yaml index 542aa39c..baca1868 100644 --- a/charts/flagger/values.yaml +++ b/charts/flagger/values.yaml @@ -8,6 +8,12 @@ image: controlLoopInterval: "10s" metricsServer: "http://prometheus.istio-system.svc.cluster.local:9090" +slack: + user: flagger + channel: + # incoming webhook https://api.slack.com/incoming-webhooks + url: + crd: create: true diff --git a/cmd/flagger/main.go b/cmd/flagger/main.go index 023490f3..6b676663 100644 --- a/cmd/flagger/main.go +++ b/cmd/flagger/main.go @@ -12,6 +12,7 @@ import ( informers "github.com/stefanprodan/flagger/pkg/client/informers/externalversions" "github.com/stefanprodan/flagger/pkg/controller" "github.com/stefanprodan/flagger/pkg/logging" + "github.com/stefanprodan/flagger/pkg/notifier" "github.com/stefanprodan/flagger/pkg/server" "github.com/stefanprodan/flagger/pkg/version" "k8s.io/client-go/kubernetes" @@ -27,6 +28,9 @@ var ( controlLoopInterval time.Duration logLevel string port string + slackURL string + slackUser string + slackChannel string ) func init() { @@ -36,6 +40,9 @@ func init() { flag.DurationVar(&controlLoopInterval, "control-loop-interval", 10*time.Second, "wait interval between rollouts") flag.StringVar(&logLevel, "log-level", "debug", "Log level can be: debug, info, warning, error.") flag.StringVar(&port, "port", "8080", "Port to listen on.") + flag.StringVar(&slackURL, "slack-url", "", "Slack hook URL.") + flag.StringVar(&slackUser, "slack-user", "flagger", "Slack user name.") + flag.StringVar(&slackChannel, "slack-channel", "", "Slack channel.") } func main() { @@ -88,6 +95,16 @@ func main() { logger.Errorf("Metrics server %s unreachable %v", metricsServer, err) } + var slack *notifier.Slack + if slackURL != "" { + slack, err = notifier.NewSlack(slackURL, slackUser, slackChannel) + if err != nil { + logger.Errorf("Notifier %v", err) + } else { + logger.Infof("Slack notifications enabled for channel %s", slack.Channel) + } + } + // start HTTP server go server.ListenAndServe(port, 3*time.Second, logger, stopCh) @@ -99,6 +116,7 @@ func main() { controlLoopInterval, metricsServer, logger, + slack, ) flaggerInformerFactory.Start(stopCh) diff --git a/docs/screens/slack-notifications.png b/docs/screens/slack-notifications.png new file mode 100644 index 00000000..98e2c83d Binary files /dev/null and b/docs/screens/slack-notifications.png differ diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 4f62899b..b17e6305 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -12,6 +12,7 @@ import ( flaggerscheme "github.com/stefanprodan/flagger/pkg/client/clientset/versioned/scheme" flaggerinformers "github.com/stefanprodan/flagger/pkg/client/informers/externalversions/flagger/v1alpha1" flaggerlisters "github.com/stefanprodan/flagger/pkg/client/listers/flagger/v1alpha1" + "github.com/stefanprodan/flagger/pkg/notifier" "go.uber.org/zap" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/errors" @@ -43,6 +44,7 @@ type Controller struct { router CanaryRouter observer CanaryObserver recorder CanaryRecorder + notifier *notifier.Slack } func NewController( @@ -53,6 +55,7 @@ func NewController( flaggerWindow time.Duration, metricServer string, logger *zap.SugaredLogger, + notifier *notifier.Slack, ) *Controller { logger.Debug("Creating event broadcaster") @@ -100,6 +103,7 @@ func NewController( router: router, observer: observer, recorder: recorder, + notifier: notifier, } flaggerInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ @@ -256,6 +260,17 @@ func (c *Controller) recordEventWarningf(r *flaggerv1.Canary, template string, a c.eventRecorder.Event(r, corev1.EventTypeWarning, "Synced", fmt.Sprintf(template, args...)) } +func (c *Controller) sendNotification(workload string, namespace string, message string, warn bool) { + if c.notifier == nil { + return + } + + err := c.notifier.Post(workload, namespace, message, warn) + if err != nil { + c.logger.Error(err) + } +} + func int32p(i int32) *int32 { return &i } diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index c0998257..ff5e7b4a 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -110,6 +110,8 @@ func (c *Controller) advanceCanary(name string, namespace string) { return } c.recorder.SetStatus(cd) + c.sendNotification(cd.Spec.TargetRef.Name, cd.Namespace, + "Canary analysis failed, rollback finished.", true) return } @@ -180,6 +182,8 @@ func (c *Controller) advanceCanary(name string, namespace string) { return } c.recorder.SetStatus(cd) + c.sendNotification(cd.Spec.TargetRef.Name, cd.Namespace, + "Canary analysis completed successfully, promotion finished.", false) } } @@ -196,11 +200,15 @@ func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, deployer CanaryDepl } c.recorder.SetStatus(cd) c.recordEventInfof(cd, "Initialization done! %s.%s", cd.Name, cd.Namespace) + c.sendNotification(cd.Spec.TargetRef.Name, cd.Namespace, + "New deployment detected, initialization completed.", false) return false } if diff, err := deployer.IsNewSpec(cd); diff { c.recordEventInfof(cd, "New revision detected! Scaling up %s.%s", cd.Spec.TargetRef.Name, cd.Namespace) + c.sendNotification(cd.Spec.TargetRef.Name, cd.Namespace, + "New revision detected, starting canary analysis.", false) if err = deployer.Scale(cd, 1); err != nil { c.recordEventErrorf(cd, "%v", err) return false diff --git a/pkg/notifier/slack.go b/pkg/notifier/slack.go new file mode 100644 index 00000000..389c6b18 --- /dev/null +++ b/pkg/notifier/slack.go @@ -0,0 +1,102 @@ +package notifier + +import ( + "bytes" + "encoding/json" + "errors" + "fmt" + "io/ioutil" + "net/http" + "net/url" +) + +// Slack holds the hook URL +type Slack struct { + URL string + Username string + Channel string + IconEmoji string +} + +// SlackPayload holds the channel and attachments +type SlackPayload struct { + Channel string `json:"channel"` + Username string `json:"username"` + IconUrl string `json:"icon_url"` + IconEmoji string `json:"icon_emoji"` + Text string `json:"text,omitempty"` + Attachments []SlackAttachment `json:"attachments,omitempty"` +} + +// SlackAttachment holds the markdown message body +type SlackAttachment struct { + Color string `json:"color"` + AuthorName string `json:"author_name"` + Text string `json:"text"` + MrkdwnIn []string `json:"mrkdwn_in"` +} + +// NewSlack validates the Slack URL and returns a Slack object +func NewSlack(hookURL string, username string, channel string) (*Slack, error) { + _, err := url.ParseRequestURI(hookURL) + if err != nil { + return nil, fmt.Errorf("invalid Slack hook URL %s", hookURL) + } + + if username == "" { + return nil, errors.New("empty Slack username") + } + + if channel == "" { + return nil, errors.New("empty Slack channel") + } + + return &Slack{ + Channel: channel, + URL: hookURL, + Username: username, + IconEmoji: ":rocket:", + }, nil +} + +// Post Slack message +func (s *Slack) Post(workload string, namespace string, message string, warn bool) error { + payload := SlackPayload{ + Channel: s.Channel, + Username: s.Username, + } + + color := "good" + if warn { + color = "danger" + } + + a := SlackAttachment{ + Color: color, + AuthorName: fmt.Sprintf("%s.%s", workload, namespace), + Text: message, + MrkdwnIn: []string{"text"}, + } + + payload.Attachments = []SlackAttachment{a} + + data, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("marshalling slack payload failed %v", err) + } + + b := bytes.NewBuffer(data) + + if res, err := http.Post(s.URL, "application/json", b); err != nil { + return fmt.Errorf("sending data to slack failed %v", err) + } else { + defer res.Body.Close() + statusCode := res.StatusCode + if statusCode != 200 { + body, _ := ioutil.ReadAll(res.Body) + return fmt.Errorf("sending data to slack failed %v", string(body)) + } + } + + return nil +}