diff --git a/apps/checker/checker/update.go b/apps/checker/checker/update.go index 8f530e80..f3adbe19 100644 --- a/apps/checker/checker/update.go +++ b/apps/checker/checker/update.go @@ -4,12 +4,16 @@ import ( "bytes" "context" "encoding/json" - "net/http" + "fmt" "os" - "time" + "strings" "github.com/rs/zerolog/log" + "google.golang.org/api/option" + "cloud.google.com/go/auth" + cloudtasks "cloud.google.com/go/cloudtasks/apiv2" + taskspb "cloud.google.com/go/cloudtasks/apiv2/cloudtaskspb" ) type UpdateData struct { @@ -22,24 +26,68 @@ type UpdateData struct { Latency int64 `json:"latency,omitempty"` } -func UpdateStatus(ctx context.Context, updateData UpdateData) { +func UpdateStatus(ctx context.Context, updateData UpdateData) error { + url := "https://openstatus-workflows.fly.dev/updateStatus" basic := "Basic " + os.Getenv("CRON_SECRET") payloadBuf := new(bytes.Buffer) + c := os.Getenv("GCP_PRIVATE_KEY") + c = strings.ReplaceAll(c, "\\n", "\n") + opts := &auth.Options2LO{ + Email: os.Getenv("GCP_CLIENT_EMAIL"), + PrivateKey: []byte(c), + PrivateKeyID: os.Getenv("GCP_PRIVATE_KEY_ID"), + Scopes: []string{ + "https://www.googleapis.com/auth/cloud-platform", + }, + TokenURL: "https://oauth2.googleapis.com/token", + } + + tp, err := auth.New2LOTokenProvider(opts) + if err != nil { + log.Ctx(ctx).Error().Err(err).Msg("error while creating token provider") + return err + } + + creds := auth.NewCredentials(&auth.CredentialsOptions{ + TokenProvider: tp, + }) + + client, err := cloudtasks.NewClient(ctx, option.WithAuthCredentials(creds)) + if err != nil { + log.Ctx(ctx).Error().Err(err).Msg("error while creating cloud tasks client") + + } + defer client.Close() if err := json.NewEncoder(payloadBuf).Encode(updateData); err != nil { log.Ctx(ctx).Error().Err(err).Msg("error while updating status") - return + return err + } + projectID := os.Getenv("GCP_PROJECT_ID") + queuePath := fmt.Sprintf("projects/%s/locations/europe-west1/queues/alerting", projectID) + req := &taskspb.CreateTaskRequest{ + Parent: queuePath, + Task: &taskspb.Task{ + // https://godoc.org/google.golang.org/genproto/googleapis/cloud/tasks/v2#HttpRequest + MessageType: &taskspb.Task_HttpRequest{ + HttpRequest: &taskspb.HttpRequest{ + HttpMethod: taskspb.HttpMethod_POST, + Url: url, + Headers: map[string]string{"Authorization": basic, "Content-Type": "application/json"}, + }, + }, + }, } - req, _ := http.NewRequestWithContext(ctx, http.MethodPost, url, payloadBuf) - req.Header.Set("Authorization", basic) - req.Header.Set("Content-Type", "application/json") - client := &http.Client{Timeout: time.Second * 10} - if _, err := client.Do(req); err != nil { - log.Ctx(ctx).Error().Err(err).Msg("error while updating status") + // Add a payload message if one is present. + req.Task.GetHttpRequest().Body = payloadBuf.Bytes() + + _, err = client.CreateTask(ctx, req) + if err != nil { + log.Ctx(ctx).Error().Err(err).Msg("error while creating the cloud task") + return fmt.Errorf("cloudtasks.CreateTask: %w", err) } - defer req.Body.Close() - // Should we add a retry mechanism here? + return nil } diff --git a/apps/checker/fly.toml b/apps/checker/fly.toml index 8c8b52f4..c04e6257 100644 --- a/apps/checker/fly.toml +++ b/apps/checker/fly.toml @@ -39,4 +39,4 @@ primary_region = "ams" [http_service.concurrency] type = "requests" hard_limit = 1000 - soft_limit = 500 \ No newline at end of file + soft_limit = 500