blob: d874d365a3e2bf5133cc64e6ae1546ab28405f3e [file]
package main
import (
"context"
"encoding/json"
"flag"
"fmt"
"io"
"net/http"
"os/user"
"regexp"
"strconv"
"strings"
"time"
"github.com/ccoveille/go-safecast/v2"
"github.com/davecgh/go-spew/spew"
"github.com/google/uuid"
apipb "go.chromium.org/luci/swarming/proto/api_v2"
"go.temporal.io/sdk/client"
"go.temporal.io/sdk/temporal"
"golang.org/x/oauth2/google"
"go.skia.org/infra/go/auth"
"go.skia.org/infra/go/httputils"
"go.skia.org/infra/go/skerr"
"go.skia.org/infra/go/sklog"
"go.skia.org/infra/pinpoint/go/backends"
"go.skia.org/infra/pinpoint/go/common"
"go.skia.org/infra/pinpoint/go/workflows"
"go.skia.org/infra/pinpoint/go/workflows/catapult"
"go.skia.org/infra/pinpoint/go/workflows/internal"
enumspb "go.temporal.io/api/enums/v1"
pb "go.skia.org/infra/pinpoint/proto/v1"
)
var (
// Run the following command to portforward Temporal service so the client can connect to it.
// kubectl port-forward service/temporal --address 0.0.0.0 -n temporal 7233:7233
hostPort = flag.String("hostPort", "localhost:7233", "Host the worker connects to.")
namespace = flag.String("namespace", "default", "The namespace the worker registered to.")
taskQueue = flag.String("taskQueue", "", "Task queue name registered to worker services.")
commit = flag.String("commit", "611b5a084486cd6d99a0dad63f34e320a2ebc2b3", "Git commit hash to build Chrome.")
startGitHash = flag.String("start-git-hash", "c73e059a2ac54302b2951e4b4f1f7d94d92a707a", "Start git commit hash for bisect.")
endGitHash = flag.String("end-git-hash", "979c9324d3c6474c15335e676ac7123312d5df82", "End git commit hash for bisect.")
patchHost = flag.String("patch-host", "chromium-review.googlesource.com", "Gerrit host of the patch")
patchProject = flag.String("patch-project", "chromium/src", "Gerrit project of the patch")
patchId = flag.Int("patch-id", 0, "Gerrit patch ID (usually a 7-digit integer)")
patchSet = flag.Int("patch-set", 0, "Gerrit patch set (usually a very small integer)")
configuration = flag.String("configuration", "mac-m2-pro-perf", "Bot configuration to use.")
benchmark = flag.String("benchmark", "speedometer3.crossbench", "Benchmark to run.")
story = flag.String("story", "default", "Story to run.")
chart = flag.String("chart", "Score", "Chart (metric or test result) to collect.")
aggregationMethod = flag.String("aggregation-method", "mean", "Aggregation method to use for bisect.")
comparisonMagnitude = flag.String("comparison-magnitude", "0.1", "Comparison magnitude for bisect.")
improvementDirection = flag.String("improvement-direction", "UP", "Improvement direction for bisect (UP or DOWN)")
jobId = flag.String("job-id", "123", "Pinpoint job ID to use.")
iterations = flag.Int("iterations", 2, "Number of iterations to run the story.")
extraArg = flag.String("extra-arg", "", "Extra argument to pass to test.")
triggerBisectFlag = flag.Bool("bisect", false, "toggle true to trigger bisect workflow")
triggerCulpritFinderFlag = flag.Bool("culprit-finder", false, "toggle true to trigger culprit-finder aka sandwich verification workflow")
triggerSingleCommitFlag = flag.Bool("single-commit", false, "toggle true to trigger single commit runner workflow")
triggerPairwiseRunnerFlag = flag.Bool("pairwise-runner", false, "toggle true to trigger pairwise commit runner workflow")
triggerPairwiseFlag = flag.Bool("pairwise", false, "toggle true to trigger pairwise workflow")
triggerBugUpdateFlag = flag.Bool("update-bug", false, "toggle true to trigger post bug comment workflow")
triggerQueryPairwiseFlag = flag.Bool("query-pairwise", false, "toggle true to trigger querying of pairwise flows")
triggerCbbRunnerFlag = flag.Bool("cbb-runner", false, "toggle true to trigger CBB runner workflow")
triggerCbbNewReleaseFlag = flag.Bool("cbb-new-release", false, "toggle true to trigger CbbNewReleaseDetectorWorkflow")
triggerCbbGetVersionsFlag = flag.Bool("cbb-get-versions", false, "toggle true to trigger CbbGetBrowserVersionsWorkflow")
triggerCbbDownloadSTPFlag = flag.Bool("cbb-download-stp", false, "toggle true to trigger DownloadSafariTPWorkflow")
triggerBuildChromeFlag = flag.Bool("build-chrome", false, "toggle true to trigger build chrome workflow")
// The following flags are used by cbb-runner only.
commitPosition = flag.Int("commit-position", 0, "Commit position (required for CBB).")
browser = flag.String("browser", "chrome", "chrome or safari or edge (used by CBB only)")
channel = flag.String("channel", "stable", "stable, dev or tp (used by CBB only)")
skipFinch = flag.Bool("skip-finch", false, "Skip Finch config (used by CBB on desktop Chrome only)")
bucket = flag.String("bucket", "prod", "GS bucket to upload results to (prod, exp, or none; used by CBB only)")
noWait = flag.Bool("no-wait", false, "if true, don't wait for workflow to finish (used by CBB only)")
mimicLegacyJob = flag.String("mimic-legacy-job", "", "Catapult legacy Pinpoint job ID to fetch and trigger locally.")
)
func defaultWorkflowOptions() client.StartWorkflowOptions {
return client.StartWorkflowOptions{
ID: uuid.New().String(),
TaskQueue: *taskQueue,
RetryPolicy: &temporal.RetryPolicy{
InitialInterval: 30 * time.Second,
BackoffCoefficient: 2.0,
MaximumInterval: 5 * time.Minute,
MaximumAttempts: 1,
},
}
}
func triggerCulpritFinderWorkflow(c client.Client) (*pb.CulpritFinderExecution, error) {
// Based off of b/344943386
ctx := context.Background()
p := &workflows.CulpritFinderParams{
Request: &pb.ScheduleCulpritFinderRequest{
StartGitHash: *startGitHash,
EndGitHash: *endGitHash,
Configuration: *configuration,
Benchmark: *benchmark,
Story: *story,
Chart: *chart,
AggregationMethod: *aggregationMethod,
ComparisonMagnitude: *comparisonMagnitude,
ImprovementDirection: *improvementDirection,
},
}
var cfe *pb.CulpritFinderExecution
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), workflows.CulpritFinderWorkflow, p)
if err != nil {
return nil, skerr.Wrapf(err, "Unable to execute workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &cfe); err != nil {
return nil, skerr.Wrapf(err, "Unable to get result")
}
return cfe, nil
}
func triggerBisectWorkflow(c client.Client) (*pb.BisectExecution, error) {
ctx := context.Background()
// based off of https://pinpoint-dot-chromeperf.appspot.com/job/17ab3cfa9e0000
p := &workflows.BisectParams{
Request: &pb.ScheduleBisectRequest{
ComparisonMode: "performance",
StartGitHash: *startGitHash,
EndGitHash: *endGitHash,
Configuration: *configuration,
Benchmark: *benchmark,
Story: *story,
Chart: *chart,
ComparisonMagnitude: *comparisonMagnitude,
AggregationMethod: *aggregationMethod,
Project: "chromium",
ImprovementDirection: *improvementDirection,
ExtraArgs: *extraArg,
},
}
var be *pb.BisectExecution
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), catapult.CatapultBisectWorkflow, p)
if err != nil {
return nil, skerr.Wrapf(err, "Unable to execute workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &be); err != nil {
return nil, skerr.Wrapf(err, "Unable to get result")
}
return be, nil
}
func triggerPairwiseRunner(c client.Client) (*internal.PairwiseRun, error) {
iterations32, err := safecast.Convert[int32](*iterations)
if err != nil {
return nil, skerr.Wrap(err)
}
ctx := context.Background()
// based off of https://pinpoint-dot-chromeperf.appspot.com/job/1372a174810000
p := &internal.PairwiseCommitsRunnerParams{
SingleCommitRunnerParams: internal.SingleCommitRunnerParams{
PinpointJobID: *jobId,
BotConfig: *configuration,
Benchmark: *benchmark,
Story: *story,
Chart: *chart,
AggregationMethod: *aggregationMethod,
Iterations: iterations32,
},
Seed: 54321,
LeftCommit: common.NewCombinedCommit(&pb.Commit{GitHash: *startGitHash}),
RightCommit: common.NewCombinedCommit(&pb.Commit{GitHash: *endGitHash}),
}
var pr *internal.PairwiseRun
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), workflows.PairwiseCommitsRunner, p)
if err != nil {
return nil, skerr.Wrapf(err, "Unable to execute workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &pr); err != nil {
return nil, skerr.Wrapf(err, "Unable to get result")
}
return pr, nil
}
// based off of https://pinpoint-dot-chromeperf.appspot.com/job/2e79457b-4d19-4e3b-9553-7baf1fd9a0e1
func triggerPairwiseWorkflow(c client.Client) (*pb.PairwiseExecution, error) {
ctx := context.Background()
p := &workflows.PairwiseParams{
Request: &pb.SchedulePairwiseRequest{
StartCommit: &pb.CombinedCommit{
Main: common.NewChromiumCommit(*startGitHash),
},
EndCommit: &pb.CombinedCommit{
Main: common.NewChromiumCommit(*endGitHash),
},
Configuration: *configuration,
Benchmark: *benchmark,
Story: *story,
Chart: *chart,
AggregationMethod: *aggregationMethod,
InitialAttemptCount: strconv.Itoa(*iterations),
ImprovementDirection: *improvementDirection,
},
}
var pe *pb.PairwiseExecution
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), workflows.PairwiseWorkflow, p)
if err != nil {
return nil, skerr.Wrapf(err, "Unable to execute workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &pe); err != nil {
return nil, skerr.Wrapf(err, "Unable to get result")
}
return pe, nil
}
func triggerSingleCommitRunner(c client.Client) (*internal.CommitRun, error) {
iterations32, err := safecast.Convert[int32](*iterations)
if err != nil {
return nil, skerr.Wrap(err)
}
ctx := context.Background()
p := &internal.SingleCommitRunnerParams{
PinpointJobID: *jobId,
BotConfig: *configuration,
Benchmark: *benchmark,
Story: *story,
Chart: *chart,
AggregationMethod: *aggregationMethod,
CombinedCommit: common.NewCombinedCommit(&pb.Commit{GitHash: *commit}),
Iterations: iterations32,
}
if *extraArg != "" {
p.ExtraArgs = []string{*extraArg}
}
if *patchId != 0 {
if *patchSet == 0 {
return nil, skerr.Fmt("--patch-set is required when --patch-id is used")
}
p.CombinedCommit.Patch = &pb.GerritChange{
Host: *patchHost,
Project: *patchProject,
Change: int64(*patchId),
Patchset: int64(*patchSet),
}
}
var cr *internal.CommitRun
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), workflows.SingleCommitRunner, p)
if err != nil {
return nil, skerr.Wrapf(err, "Unable to execute workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &cr); err != nil {
return nil, skerr.Wrapf(err, "Unable to get result")
}
return cr, nil
}
func triggerBuildChrome(c client.Client) (*apipb.CASReference, error) {
bcp := workflows.BuildParams{
WorkflowID: *jobId,
Commit: common.NewCombinedCommit(&pb.Commit{GitHash: *commit}),
Device: *configuration,
Target: "performance_test_suite",
}
we, err := c.ExecuteWorkflow(context.Background(), defaultWorkflowOptions(), workflows.BuildChrome, &bcp)
if err != nil {
return nil, skerr.Wrapf(err, "Unable to execute workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
// Synchronously wait for the workflow completion.
var result *apipb.CASReference
err = we.Get(context.Background(), &result)
if err != nil {
return nil, skerr.Wrapf(err, "Unable get workflow result")
}
return result, nil
}
func triggerBugUpdateWorkflow(c client.Client) (bool, error) {
ctx := context.Background()
var success bool
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), workflows.BugUpdate, 333705433, "hello world")
if err != nil {
return false, skerr.Wrapf(err, "Unable to execute the workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &success); err != nil {
return false, skerr.Wrapf(err, "Unable to write to buganizer")
}
return success, nil
}
func triggerQueryPairwise(c client.Client) (*pb.QueryPairwiseResponse, error) {
ctx := context.Background()
workflow_id := "45fe60b9-668e-44a9-9991-678221ba264d"
resp, err := c.DescribeWorkflowExecution(ctx, workflow_id, "")
if err != nil {
return nil, skerr.Wrapf(err, "Unable to describe workflow")
}
workflowStatus := resp.GetWorkflowExecutionInfo().GetStatus()
fmt.Print(workflowStatus)
var pairwiseExecution pb.PairwiseExecution
switch workflowStatus {
case enumspb.WORKFLOW_EXECUTION_STATUS_COMPLETED:
workflowRun := c.GetWorkflow(ctx, workflow_id, "")
errGet := workflowRun.Get(ctx, &pairwiseExecution)
if errGet != nil {
return nil, skerr.Wrapf(errGet, "Pairwise workflow completed, but failed to get results")
}
return &pb.QueryPairwiseResponse{
Status: pb.PairwiseJobStatus_PAIRWISE_JOB_STATUS_COMPLETED,
Execution: &pairwiseExecution,
}, nil
case enumspb.WORKFLOW_EXECUTION_STATUS_FAILED,
enumspb.WORKFLOW_EXECUTION_STATUS_TIMED_OUT,
enumspb.WORKFLOW_EXECUTION_STATUS_TERMINATED:
return &pb.QueryPairwiseResponse{
Status: pb.PairwiseJobStatus_PAIRWISE_JOB_STATUS_FAILED,
Execution: &pairwiseExecution,
}, nil
case enumspb.WORKFLOW_EXECUTION_STATUS_CANCELED:
return &pb.QueryPairwiseResponse{
Status: pb.PairwiseJobStatus_PAIRWISE_JOB_STATUS_CANCELED,
Execution: &pairwiseExecution,
}, nil
case enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING,
enumspb.WORKFLOW_EXECUTION_STATUS_CONTINUED_AS_NEW:
return &pb.QueryPairwiseResponse{
Status: pb.PairwiseJobStatus_PAIRWISE_JOB_STATUS_RUNNING,
Execution: &pairwiseExecution,
}, nil
}
return nil, nil
}
func triggerCbbRunner(c client.Client) (*internal.CommitRun, error) {
if len(flag.Args()) != 0 {
return nil, skerr.Fmt("Unrecognized command line arguments: %v", flag.Args())
}
if *commitPosition <= 0 {
return nil, skerr.Fmt("Please specify a valid positive commit position using --commit-position switch")
}
commitPos32, err := safecast.Convert[int32](*commitPosition)
if err != nil {
return nil, skerr.Wrapf(err, "invalid commit position")
}
ctx := context.Background()
// If the user didn't specify a commit hash (so that *commit has the
// default value), or the user specified an empty commit hash,
// we try to get the commit hash from the commit position.
if *commit == flag.Lookup("commit").DefValue || *commit == "" {
crrev, err := backends.NewCrrevClient(ctx)
if err != nil {
return nil, skerr.Wrapf(err, "unable to create crrev client")
}
ci, err := crrev.GetCommitInfo(ctx, strconv.Itoa(*commitPosition))
if err != nil {
return nil, skerr.Wrapf(err, "unable to get commit info")
}
if len(ci.GitHash) != 40 {
// When given an invalid commit position, NewCrrevClient doesn't
// return an error, but converts the commit position into a string.
// Since a valid commit hash must be 40 characters long, we assume
// an error if the length is incorrect.
return nil, skerr.Fmt(
"commit position %d appears invalid, GetCommitInfo returned %s",
*commitPosition, ci.GitHash,
)
}
*commit = ci.GitHash
fmt.Println("Using commit hash", *commit, "based on commit position", *commitPosition)
}
p := &internal.CbbRunnerParams{
BotConfig: *configuration,
Commit: common.NewCombinedCommit(common.NewChromiumCommit(*commit)),
Browser: *browser,
Channel: *channel,
SkipFinch: *skipFinch,
}
p.Commit.Main.CommitPosition = commitPos32
if *patchId != 0 {
if *patchSet == 0 {
return nil, skerr.Fmt("--patch-set is required when --patch-id is used")
}
p.Commit.Patch = &pb.GerritChange{
Host: *patchHost,
Project: *patchProject,
Change: int64(*patchId),
Patchset: int64(*patchSet),
}
}
if p.Channel == "tp" {
p.Channel = "technology-preview"
}
if *iterations == 0 {
if strings.HasPrefix(*configuration, "mac") {
*iterations = 3
} else {
*iterations = 2
}
}
iterations32, err := safecast.Convert[int32](*iterations)
if err != nil {
return nil, skerr.Wrap(err)
}
switch *benchmark {
case "", "full":
// Setting p.Benchmarks to nil causes the default full set of benchmarks to run.
p.Benchmarks = nil
case "trial":
p.Benchmarks = []internal.BenchmarkRunConfig{
{Benchmark: "speedometer3", Iterations: iterations32},
{Benchmark: "jetstream2", Iterations: iterations32},
{Benchmark: "jetstream3", Iterations: iterations32},
{Benchmark: "motionmark1.3", Iterations: iterations32},
}
default:
// Multiple benchmarks can be specified, separated by ",".
for _, b := range strings.Split(*benchmark, ",") {
// Each benchmark can be specified as "name", or "name:iteration"
colon := strings.Index(b, ":")
if colon == -1 {
p.Benchmarks = append(
p.Benchmarks,
internal.BenchmarkRunConfig{Benchmark: b, Iterations: iterations32},
)
} else {
i, err := strconv.ParseInt(b[colon+1:], 10, 32)
if err != nil {
return nil, skerr.Wrapf(err, "Invalid iteration %v in --benchmark", b[colon+1:])
}
p.Benchmarks = append(p.Benchmarks, internal.BenchmarkRunConfig{Benchmark: b[:colon], Iterations: int32(i)})
}
}
}
switch *bucket {
case "prod":
p.Bucket = "chrome-perf-non-public"
case "exp":
p.Bucket = "chrome-perf-experiment-non-public"
case "none":
p.Bucket = ""
}
var cr *internal.CommitRun
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), workflows.CbbRunner, p)
if err != nil {
return nil, skerr.Wrapf(err, "Unable to execute workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if *noWait {
return nil, nil
}
if err := we.Get(ctx, &cr); err != nil {
return nil, skerr.Wrapf(err, "Unable to get result")
}
return cr, nil
}
func triggerCbbNewReleaseDetector(c client.Client) (*internal.ChromeReleaseInfo, error) {
ctx := context.Background()
var result *internal.ChromeReleaseInfo
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), workflows.CbbNewReleaseDetector)
if err != nil {
return nil, skerr.Wrapf(err, "Unable to execute the workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &result); err != nil {
return nil, skerr.Wrapf(err, "Unable to get results from CbbNewReleaseDetector workflow")
}
return result, nil
}
func triggerCbbGetBrowserVersions(c client.Client) ([]internal.BuildInfo, error) {
if len(flag.Args()) != 0 {
return nil, skerr.Fmt("Unrecognized command line arguments: %v", flag.Args())
}
if *browser != "safari" && *browser != "edge" {
return nil, skerr.Fmt("Either --browser=safari or --browser=edge is required")
}
ctx := context.Background()
var buildInfos []internal.BuildInfo
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), workflows.CbbGetBrowserVersions, *browser)
if err != nil {
return nil, skerr.Wrapf(err, "Unable to execute workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &buildInfos); err != nil {
return nil, skerr.Wrapf(err, "Unable to get result")
}
return buildInfos, nil
}
func triggerCbbDownloadSTP(c client.Client) (string, error) {
ctx := context.Background()
var result string
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), workflows.CbbDownloadSafariTP)
if err != nil {
return "", skerr.Wrapf(err, "Unable to execute the workflow")
}
sklog.Infof("Started workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &result); err != nil {
return "", skerr.Wrapf(err, "Unable to get results from DownloadSafariTP workflow")
}
return result, nil
}
var gerritURLRegex = regexp.MustCompile(`^https?://([^/]+)/c/(.+)/\+/(\d+)(?:/(\d+))?/?$`)
func parseGerritURL(gerritURL string) (*pb.GerritChange, error) {
matches := gerritURLRegex.FindStringSubmatch(gerritURL)
if matches == nil {
return nil, skerr.Fmt("invalid gerrit URL format: %s", gerritURL)
}
host := matches[1]
project := matches[2]
changeID, err := strconv.ParseInt(matches[3], 10, 64)
if err != nil {
return nil, skerr.Wrapf(err, "failed to parse change ID")
}
var patchset int64 = 1
if len(matches) > 4 && matches[4] != "" {
ps, err := strconv.ParseInt(matches[4], 10, 64)
if err == nil {
patchset = ps
}
}
return &pb.GerritChange{
Host: host,
Project: project,
Change: changeID,
Patchset: patchset,
}, nil
}
type catapultJobResponse struct {
Arguments struct {
ComparisonMode string `json:"comparison_mode"`
Target string `json:"target"`
StartGitHash string `json:"start_git_hash"`
BaseGitHash string `json:"base_git_hash"`
EndGitHash string `json:"end_git_hash"`
Trace string `json:"trace"`
Pin string `json:"pin"`
Configuration string `json:"configuration"`
Benchmark string `json:"benchmark"`
Story string `json:"story"`
StoryTags string `json:"story_tags"`
Chart string `json:"chart"`
Statistic string `json:"statistic"`
ComparisonMagnitude interface{} `json:"comparison_magnitude"`
InitialAttemptCount interface{} `json:"initial_attempt_count"`
Tags interface{} `json:"tags"`
} `json:"arguments"`
JobID string `json:"job_id"`
ComparisonMode string `json:"comparison_mode"`
}
func getAttemptCount(val interface{}) string {
if val == nil {
return "30"
}
switch v := val.(type) {
case string:
return v
case float64:
return strconv.Itoa(int(v))
case int:
return strconv.Itoa(v)
default:
return "30"
}
}
func getComparisonMagnitude(val interface{}) string {
if val == nil {
return "0.1"
}
switch v := val.(type) {
case string:
return v
case float64:
return fmt.Sprintf("%f", v)
default:
return "0.1"
}
}
func parseTags(val interface{}) map[string]string {
res := make(map[string]string)
if val == nil {
return res
}
switch v := val.(type) {
case string:
_ = json.Unmarshal([]byte(v), &res)
case map[string]interface{}:
for k, val := range v {
if str, ok := val.(string); ok {
res[k] = str
}
}
}
return res
}
func triggerMimicLegacyJob(c client.Client, legacyJobID string) (interface{}, error) {
ctx := context.Background()
url := fmt.Sprintf("https://pinpoint-dot-chromeperf.appspot.com/api/job/%s", legacyJobID)
sklog.Infof("Fetching legacy Pinpoint job from: %s", url)
tokenSource, err := google.DefaultTokenSource(ctx, auth.ScopeUserinfoEmail)
if err != nil {
return nil, skerr.Wrapf(err, "failed to create token source")
}
httpClient := httputils.DefaultClientConfig().WithTokenSource(tokenSource).Client()
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, http.NoBody)
if err != nil {
return nil, skerr.Wrapf(err, "failed to create request for legacy Pinpoint job")
}
resp, err := httpClient.Do(req)
if err != nil {
return nil, skerr.Wrapf(err, "failed to fetch legacy Pinpoint job")
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return nil, skerr.Fmt("failed to fetch legacy job, status code: %d", resp.StatusCode)
}
body, err := io.ReadAll(resp.Body)
if err != nil {
return nil, skerr.Wrapf(err, "failed to read response body")
}
var jobResp catapultJobResponse
if err := json.Unmarshal(body, &jobResp); err != nil {
return nil, skerr.Wrapf(err, "failed to unmarshal legacy job response")
}
args := jobResp.Arguments
comparisonMode := args.ComparisonMode
if comparisonMode == "" {
comparisonMode = jobResp.ComparisonMode
}
comparisonMode = strings.ToLower(comparisonMode)
sklog.Infof("Detected legacy job comparison mode: %s", comparisonMode)
if comparisonMode == "try" {
startHash := args.StartGitHash
if startHash == "" {
startHash = args.BaseGitHash
}
endHash := args.EndGitHash
if startHash == "" {
return nil, skerr.Fmt("missing start_git_hash or base_git_hash in legacy job arguments")
}
if endHash == "" {
endHash = startHash
}
var patch *pb.GerritChange
if strings.Contains(args.Pin, "review.googlesource.com") {
p, err := parseGerritURL(args.Pin)
if err == nil {
patch = p
sklog.Infof("Successfully parsed Gerrit patch from 'pin': %+v", patch)
}
}
parsedTags := parseTags(args.Tags)
if patch == nil && len(parsedTags) > 0 {
if patchURL, ok := parsedTags["patch"]; ok {
p, err := parseGerritURL(patchURL)
if err == nil {
patch = p
sklog.Infof("Successfully parsed Gerrit patch from tag 'patch': %+v", patch)
}
}
}
req := &pb.SchedulePairwiseRequest{
StartCommit: &pb.CombinedCommit{
Main: common.NewChromiumCommit(startHash),
},
EndCommit: &pb.CombinedCommit{
Main: common.NewChromiumCommit(endHash),
},
Configuration: args.Configuration,
Benchmark: args.Benchmark,
Story: args.Story,
Chart: args.Chart,
AggregationMethod: args.Statistic,
InitialAttemptCount: getAttemptCount(args.InitialAttemptCount),
ImprovementDirection: "UNKNOWN",
}
if req.Chart == "" {
req.Chart = args.Trace
}
if patch != nil {
req.EndCommit.Patch = patch
}
p := &workflows.PairwiseParams{
Request: req,
}
var pe *pb.PairwiseExecution
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), workflows.PairwiseWorkflow, p)
if err != nil {
return nil, skerr.Wrapf(err, "failed to execute pairwise workflow")
}
sklog.Infof("Started Mimicked Pairwise TryJob Workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &pe); err != nil {
return nil, skerr.Wrapf(err, "failed to get pairwise workflow result")
}
return pe, nil
}
if comparisonMode == "performance" || comparisonMode == "functional" || comparisonMode == "bisection" {
startHash := args.StartGitHash
if startHash == "" {
startHash = args.BaseGitHash
}
endHash := args.EndGitHash
if startHash == "" || endHash == "" {
return nil, skerr.Fmt("missing start_git_hash/base_git_hash or end_git_hash in legacy job arguments")
}
req := &pb.ScheduleBisectRequest{
ComparisonMode: comparisonMode,
StartGitHash: startHash,
EndGitHash: endHash,
Configuration: args.Configuration,
Benchmark: args.Benchmark,
Story: args.Story,
Chart: args.Chart,
ComparisonMagnitude: getComparisonMagnitude(args.ComparisonMagnitude),
AggregationMethod: args.Statistic,
Project: "chromium",
ImprovementDirection: "UNKNOWN",
}
if req.Chart == "" {
req.Chart = args.Trace
}
p := &workflows.BisectParams{
Request: req,
}
var be *pb.BisectExecution
we, err := c.ExecuteWorkflow(ctx, defaultWorkflowOptions(), catapult.CatapultBisectWorkflow, p)
if err != nil {
return nil, skerr.Wrapf(err, "failed to execute bisect workflow")
}
sklog.Infof("Started Mimicked Bisect Workflow.. WorkflowID: %v RunID: %v", we.GetID(), we.GetRunID())
if err := we.Get(ctx, &be); err != nil {
return nil, skerr.Wrapf(err, "failed to get bisect workflow result")
}
return be, nil
}
return nil, skerr.Fmt("unsupported comparison mode for mimicking legacy job: %s", comparisonMode)
}
// Sample client to trigger a BuildChrome workflow.
func main() {
flag.Parse()
if *taskQueue == "" {
if u, err := user.Current(); err != nil {
sklog.Fatalf("Unable to get the current user: %s", err)
} else {
*taskQueue = fmt.Sprintf("localhost.%s", u.Username)
}
}
// The client is a heavyweight object that should be created once per process.
c, err := client.Dial(client.Options{
HostPort: *hostPort,
Namespace: *namespace,
})
if err != nil {
sklog.Errorf("Unable to create client", err)
return
}
defer c.Close()
var result interface{}
if *triggerBisectFlag {
result, err = triggerBisectWorkflow(c)
}
if *triggerCulpritFinderFlag {
result, err = triggerCulpritFinderWorkflow(c)
}
if *triggerSingleCommitFlag {
result, err = triggerSingleCommitRunner(c)
}
if *triggerPairwiseRunnerFlag {
result, err = triggerPairwiseRunner(c)
}
if *triggerPairwiseFlag {
result, err = triggerPairwiseWorkflow(c)
}
if *triggerBugUpdateFlag {
result, err = triggerBugUpdateWorkflow(c)
}
if *triggerQueryPairwiseFlag {
result, err = triggerQueryPairwise(c)
}
if *triggerCbbRunnerFlag {
result, err = triggerCbbRunner(c)
}
if *triggerCbbNewReleaseFlag {
result, err = triggerCbbNewReleaseDetector(c)
}
if *triggerCbbGetVersionsFlag {
result, err = triggerCbbGetBrowserVersions(c)
}
if *triggerCbbDownloadSTPFlag {
result, err = triggerCbbDownloadSTP(c)
}
if *triggerBuildChromeFlag {
result, err = triggerBuildChrome(c)
}
if *mimicLegacyJob != "" {
result, err = triggerMimicLegacyJob(c, *mimicLegacyJob)
}
if err != nil {
sklog.Errorf("Workflow failed:", err)
return
}
sklog.Infof("Workflow result: %v", spew.Sdump(result))
}