| 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)) |
| } |