chore: refactor and add tests
This commit is contained in:
@@ -7,8 +7,10 @@ import (
|
||||
"github.com/multiplay/go-slack/webhook"
|
||||
"github.com/robfig/cron"
|
||||
"gopkg.in/alecthomas/kingpin.v2"
|
||||
"io"
|
||||
v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/kubernetes/typed/batch/v1beta1"
|
||||
"k8s.io/client-go/rest"
|
||||
"os"
|
||||
"os/signal"
|
||||
@@ -16,39 +18,50 @@ import (
|
||||
"time"
|
||||
)
|
||||
|
||||
var checkFunc = doCheck
|
||||
var exitFunc = os.Exit
|
||||
|
||||
func main() {
|
||||
slackUrl := kingpin.Flag("slack-url", "The Slack Webhook URL").Envar("SLACK_URL").Required().String()
|
||||
kingpin.Parse()
|
||||
|
||||
config, err := rest.InClusterConfig()
|
||||
if err != nil {
|
||||
panic(err.Error())
|
||||
}
|
||||
clientset, err := kubernetes.NewForConfig(config)
|
||||
if err != nil {
|
||||
panic(err.Error())
|
||||
}
|
||||
exitFunc(doMain(*slackUrl, &DefaultProvider{&InClusterProvider{}}))
|
||||
}
|
||||
|
||||
slack := webhook.New(*slackUrl)
|
||||
|
||||
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
|
||||
func doMain(slackUrl string, provider ClientProvider) int {
|
||||
client, err := provider.Provide()
|
||||
if err != nil {
|
||||
fmt.Printf("Unable to connect to K8S: %s\n", err)
|
||||
return 1
|
||||
}
|
||||
|
||||
ic := make(chan os.Signal, 1)
|
||||
signal.Notify(ic, os.Interrupt, syscall.SIGTERM)
|
||||
|
||||
if err := checkFunc(client, slackUrl, ic, 60*time.Second, os.Stdout); err != nil {
|
||||
fmt.Printf("Error checking jobs: %s\n", err)
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func doCheck(client Client, slackUrl string, ic chan os.Signal, sleepTime time.Duration, out io.Writer) error {
|
||||
slack := webhook.New(slackUrl)
|
||||
|
||||
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ic:
|
||||
fmt.Printf("Got SIGTERM signal, exiting\n")
|
||||
break
|
||||
_, _ = fmt.Fprintf(out, "Got SIGTERM signal, exiting\n")
|
||||
return nil
|
||||
default:
|
||||
cronjobs, err := clientset.BatchV1beta1().CronJobs("").List(context.Background(), v1.ListOptions{})
|
||||
cronJobs, err := client.BatchV1beta1().CronJobs("").List(context.Background(), v1.ListOptions{})
|
||||
if err != nil {
|
||||
fmt.Printf("Error getting cronjobs: %s", err)
|
||||
os.Exit(1)
|
||||
return fmt.Errorf("error getting cronjobs: %w", err)
|
||||
}
|
||||
limit := time.Now().Add(-120 * time.Second)
|
||||
for _, c := range cronjobs.Items {
|
||||
for _, c := range cronJobs.Items {
|
||||
if c.Spec.Suspend == nil || !*c.Spec.Suspend {
|
||||
since := c.CreationTimestamp
|
||||
if c.Status.LastScheduleTime != nil {
|
||||
@@ -56,28 +69,58 @@ func main() {
|
||||
}
|
||||
schedule, err := parser.Parse(c.Spec.Schedule)
|
||||
if err != nil {
|
||||
fmt.Printf("Error parsing schedule of %s/%s (%s): %s", c.Namespace, c.Name, c.Spec.Schedule, err)
|
||||
os.Exit(1)
|
||||
return fmt.Errorf("error parsing schedule of %s/%s (%s): %w", c.Namespace, c.Name, c.Spec.Schedule, err)
|
||||
}
|
||||
next := schedule.Next(since.Time)
|
||||
fmt.Printf("Checking %s/%s since %s, next schedule %s, limit %s.\n", c.Namespace, c.Name, since.Format(time.RFC3339), next.Format(time.RFC3339), limit.Format(time.RFC3339))
|
||||
_, _ = fmt.Fprintf(out, "Checking %s/%s since %s, next schedule %s, limit %s.\n", c.Namespace, c.Name, since.Format(time.RFC3339), next.Format(time.RFC3339), limit.Format(time.RFC3339))
|
||||
if next.Before(limit) {
|
||||
fmt.Printf("%s was not scheduled. Sending Slack notification.\n", c.Name)
|
||||
_, _ = fmt.Fprintf(out, "%s/%s was not scheduled. Sending Slack notification.\n", c.Namespace, c.Name)
|
||||
m := &chat.Message{
|
||||
Text: fmt.Sprintf("Cronjob %s/%s is not running according to schedule (%s). Last scheduled: %s", c.Namespace, c.Name, c.Spec.Schedule, since.Format(time.RFC3339)),
|
||||
Username: "cron-checker",
|
||||
}
|
||||
resp, err := m.Send(slack)
|
||||
_, err := m.Send(slack)
|
||||
if err != nil {
|
||||
fmt.Printf("Unable to send Slack notification: %s", err)
|
||||
}
|
||||
if !resp.OK {
|
||||
fmt.Printf("Unable to send Slack notification: %s", resp.Error)
|
||||
_, _ = fmt.Fprintf(out, "Unable to send Slack notification: %s\n", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
time.Sleep(60 * time.Second)
|
||||
time.Sleep(sleepTime)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type Client interface {
|
||||
BatchV1beta1() v1beta1.BatchV1beta1Interface
|
||||
}
|
||||
|
||||
type ClientProvider interface {
|
||||
Provide() (Client, error)
|
||||
}
|
||||
|
||||
type DefaultProvider struct {
|
||||
provider ConfigProvider
|
||||
}
|
||||
|
||||
func (d DefaultProvider) Provide() (Client, error) {
|
||||
config, err := d.provider.Provide()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return kubernetes.NewForConfig(config)
|
||||
}
|
||||
|
||||
var _ ClientProvider = &DefaultProvider{}
|
||||
|
||||
type ConfigProvider interface {
|
||||
Provide() (*rest.Config, error)
|
||||
}
|
||||
|
||||
type InClusterProvider struct{}
|
||||
|
||||
func (i InClusterProvider) Provide() (*rest.Config, error) {
|
||||
return rest.InClusterConfig()
|
||||
}
|
||||
|
||||
var _ ConfigProvider = &InClusterProvider{}
|
||||
|
||||
Reference in New Issue
Block a user