diff --git a/cmd/internal-tools/campaign/campaign.go b/cmd/internal-tools/campaign/campaign.go index dce9206..25a8f16 100644 --- a/cmd/internal-tools/campaign/campaign.go +++ b/cmd/internal-tools/campaign/campaign.go @@ -10,11 +10,14 @@ import ( "os" "os/exec" "path/filepath" + "slices" "strconv" "strings" "sync" "syscall" "time" + + "github.com/priyanshujain/sanderling/internal/android" ) // commandExecutor runs one sanderling invocation and returns its exit code. @@ -46,6 +49,23 @@ func executeCommand(ctx context.Context, binary string, arguments []string, outp return -1, err } +// connectedDevices lists the devices the host currently has. A variable so a +// preflight test runs without a device farm attached. +var connectedDevices = android.ConnectedDevices + +// A failure that came back in less than fastFailureThreshold never did the +// work the run was asked to do: a step-budgeted run takes tens of minutes, +// while a worker pointed at a device that is gone gives up in about half a +// minute. fastFailuresBeforeQuarantine of those in a row, with no run that +// worked in between to reset the count, is a property of the device rather +// than a flake, and it is where the cost of being wrong (one worker's +// throughput, since the seeds stay on the shared queue) is still smaller than +// the cost of being right one seed later. +const ( + fastFailureThreshold = 2 * time.Minute + fastFailuresBeforeQuarantine = 3 +) + // runShutdownGrace bounds how long a signalled run gets to stop its sidecar // before it is killed. It exceeds the sidecar's own 15s shutdown grace, or the // escalation would land while the run was still doing what it was asked. @@ -96,6 +116,62 @@ type campaign struct { mutex sync.Mutex failures int unreadable int + quarantined []quarantinedDevice + // unrunSeeds are the seeds left without a trustworthy result: those a + // quarantined device consumed on its way out, and those still queued when + // the last worker stopped. + unrunSeeds []int64 +} + +// failureStreak is one worker's run of fast failures, and the seeds they cost. +type failureStreak struct { + fastFailures int + consumedSeeds []int64 +} + +// failedFast reports a run that came back non-zero too quickly to have done the +// work it was given. A killed run is excluded: it outlived the run timeout, +// which is the opposite of failing fast. +func failedFast(record runRecord) bool { + return record.ExitCode != 0 && !record.TimedOut && + time.Duration(record.MonotonicMillis)*time.Millisecond < fastFailureThreshold +} + +// preflightDevices refuses to start until every serial in --devices is present. +// A worker aimed at a serial that is gone fails in seconds and pulls the next +// seed, so a handful of dead serials drain the queue while the healthy workers +// are still inside their first run. Discovering that on run 1 of 20 is already +// too late: the sweep is spent, and its output does not say why. +func preflightDevices(ctx context.Context, configuration config) error { + // Android is the platform whose worker names a device this can enumerate: a + // web worker is a label with nothing behind it, and an --ios-device is + // resolved by simctl or by devicectl depending on whether it names a + // simulator or a paired phone. + if configuration.platform != "android" || len(configuration.devices) == 0 { + return nil + } + present, err := connectedDevices(ctx) + if err != nil { + return fmt.Errorf("list android devices: %w", err) + } + var missing []string + for _, device := range configuration.devices { + if !slices.Contains(present, device) { + missing = append(missing, device) + } + } + if len(missing) == 0 { + return nil + } + return fmt.Errorf("not starting: %d of %d --devices not connected: %s (adb reports %s)", + len(missing), len(configuration.devices), strings.Join(missing, ", "), presentDevices(present)) +} + +func presentDevices(present []string) string { + if len(present) == 0 { + return "no devices" + } + return strings.Join(present, ", ") } func runCampaign(ctx context.Context, configuration config, executor commandExecutor, stdout io.Writer) error { @@ -103,6 +179,9 @@ func runCampaign(ctx context.Context, configuration config, executor commandExec return fmt.Errorf("%s already exists in %s: pick a fresh --output so two campaigns do not share a directory", manifestFileName, configuration.outputDirectory) } + if err := preflightDevices(ctx, configuration); err != nil { + return err + } if err := os.MkdirAll(configuration.outputDirectory, 0o755); err != nil { return fmt.Errorf("create campaign dir: %w", err) } @@ -117,7 +196,8 @@ func runCampaign(ctx context.Context, configuration config, executor commandExec return err } host, _ := os.Hostname() - if err := writeManifest(configuration.outputDirectory, buildManifest(configuration, host, binaryPath, version, time.Now().UTC())); err != nil { + intended := buildManifest(configuration, host, binaryPath, version, time.Now().UTC()) + if err := writeManifest(configuration.outputDirectory, intended); err != nil { return fmt.Errorf("write %s: %w", manifestFileName, err) } @@ -135,6 +215,19 @@ func runCampaign(ctx context.Context, configuration config, executor commandExec fmt.Fprintf(stdout, "campaign complete: %d of %d runs failed, %d produced an unreadable trace\n", sweep.failures, len(configuration.seeds), sweep.unreadable) + if len(sweep.quarantined) > 0 || len(sweep.unrunSeeds) > 0 { + intended.Quarantined = sweep.quarantined + intended.UnrunSeeds = sweep.unrunSeeds + if err := writeManifest(configuration.outputDirectory, intended); err != nil { + fmt.Fprintf(stdout, "warning: rewrite %s: %v\n", manifestFileName, err) + } + fmt.Fprintf(stdout, "%d of %d seeds have no result: %v\n", + len(sweep.unrunSeeds), len(configuration.seeds), sweep.unrunSeeds) + } + if len(sweep.quarantined) > 0 && len(sweep.quarantined) == len(workerDevices(configuration.devices)) { + return fmt.Errorf("every device was quarantined (%s); %d of %d seeds have no result", + strings.Join(quarantinedNames(sweep.quarantined), ", "), len(sweep.unrunSeeds), len(configuration.seeds)) + } if sweep.failures > 0 { return fmt.Errorf("%d of %d runs failed", sweep.failures, len(configuration.seeds)) } @@ -180,15 +273,68 @@ func (c *campaign) sweep(ctx context.Context) { waitGroup.Add(1) go func(device string) { defer waitGroup.Done() - for seed := range queue { - if ctx.Err() != nil { - return - } - c.report(c.runSeed(ctx, seed, device)) - } + c.work(ctx, device, queue) }(device) } waitGroup.Wait() + + for seed := range queue { + c.unrunSeeds = append(c.unrunSeeds, seed) + } + slices.Sort(c.unrunSeeds) +} + +// work runs seeds on one device until the queue is empty or the device has +// failed fast often enough in a row to be quarantined. The streak is worker +// local: one worker drives one device, and any run that did its work clears it. +func (c *campaign) work(ctx context.Context, device string, queue <-chan int64) { + var streak failureStreak + for seed := range queue { + if ctx.Err() != nil { + return + } + record := c.runSeed(ctx, seed, device) + c.report(record) + // A cancelled campaign fails every run in flight in seconds, which is + // the shutdown and not the device. + if ctx.Err() != nil { + return + } + if !failedFast(record) { + streak = failureStreak{} + continue + } + streak.fastFailures++ + streak.consumedSeeds = append(streak.consumedSeeds, seed) + if streak.fastFailures >= fastFailuresBeforeQuarantine { + c.quarantine(device, streak) + return + } + } +} + +// quarantine takes a device out of the sweep. The seeds it burned are reported +// unrun rather than requeued: each already holds that attempt's log, and a +// second record for the same seed would make one seed count as two runs. +func (c *campaign) quarantine(device string, streak failureStreak) { + c.mutex.Lock() + defer c.mutex.Unlock() + c.quarantined = append(c.quarantined, quarantinedDevice{ + Device: device, + FastFailures: streak.fastFailures, + ConsumedSeeds: streak.consumedSeeds, + }) + c.unrunSeeds = append(c.unrunSeeds, streak.consumedSeeds...) + fmt.Fprintf(c.stdout, "quarantined device %q: %d runs in a row failed in under %s; seeds %v have no result\n", + device, streak.fastFailures, fastFailureThreshold, streak.consumedSeeds) +} + +func quarantinedNames(devices []quarantinedDevice) []string { + names := make([]string, 0, len(devices)) + for _, device := range devices { + names = append(names, device.Device) + } + return names } func (c *campaign) runSeed(ctx context.Context, seed int64, device string) runRecord { diff --git a/cmd/internal-tools/campaign/campaign_test.go b/cmd/internal-tools/campaign/campaign_test.go index 349b756..bf64b31 100644 --- a/cmd/internal-tools/campaign/campaign_test.go +++ b/cmd/internal-tools/campaign/campaign_test.go @@ -252,6 +252,7 @@ func TestRunCampaign_RecordsPerRunSummary(t *testing.T) { func TestRunCampaign_DistributesSeedsAcrossDeviceWorkers(t *testing.T) { directory := t.TempDir() configuration := testConfiguration(t, directory, "--seeds", "1-9", "--devices", "device-a,device-b,device-c") + devicesPresent(t, "device-a", "device-b", "device-c") var mutex sync.Mutex assignments := map[int64]string{} @@ -462,3 +463,208 @@ func TestParseArguments_RunTimeoutDefaultsToThreeTimesDuration(t *testing.T) { t.Errorf("run timeout default: got %s, want 12m", configuration.runTimeout) } } + +// devicesPresent points the preflight at a fixed set of serials, so a campaign +// can be preflighted on a host with no device farm attached. +func devicesPresent(t *testing.T, present ...string) { + t.Helper() + original := connectedDevices + connectedDevices = func(context.Context) ([]string, error) { return present, nil } + t.Cleanup(func() { connectedDevices = original }) +} + +// A worker aimed at a serial that no longer exists fails in seconds and pulls +// the next seed, so a few dead serials drain the queue while the healthy +// workers are still inside their first run. That has to be caught before the +// first seed is dispatched: a sweep that discovers it on run 1 of 20 has +// already been destroyed, and its output does not say so. +func TestRunCampaign_RefusesToStartWhenADeviceIsMissing(t *testing.T) { + directory := t.TempDir() + configuration := testConfiguration(t, directory, + "--seeds", "1-20", "--devices", "emulator-5554,emulator-5564,emulator-5556") + devicesPresent(t, "emulator-5554", "emulator-5556") + + executor := func(_ context.Context, _ string, arguments []string, _ io.Writer) (int, error) { + t.Errorf("the campaign dispatched %v despite a missing device", arguments) + return 0, nil + } + err := runCampaign(context.Background(), configuration, executor, io.Discard) + if err == nil || !strings.Contains(err.Error(), "not connected: emulator-5564") { + t.Fatalf("the error must name the missing serial, got %v", err) + } + if _, statErr := os.Stat(filepath.Join(directory, manifestFileName)); statErr == nil { + t.Error("a campaign that cannot run wrote a manifest") + } +} + +func TestRunCampaign_RunsWhenEveryDeviceIsPresent(t *testing.T) { + directory := t.TempDir() + configuration := testConfiguration(t, directory, "--seeds", "1-2", "--devices", "emulator-5554,emulator-5556") + devicesPresent(t, "emulator-5556", "emulator-5580", "emulator-5554") + + executor := versionAnswering(func(_ context.Context, _ string, arguments []string, _ io.Writer) (int, error) { + writeFakeRun(t, arguments, []trace.Step{observedStep(1)}) + return 0, nil + }) + if err := runCampaign(context.Background(), configuration, executor, io.Discard); err != nil { + t.Fatal(err) + } + if got := len(readRecords(t, directory)); got != 2 { + t.Errorf("records: got %d, want 2", got) + } +} + +// Preflight only means something where the worker names a device. A campaign +// without --devices has one worker and no serial, and a web worker is a label +// with no device behind it; neither may be blocked by a device check. +func TestRunCampaign_PreflightsOnlyWhereWorkersNameADevice(t *testing.T) { + for _, testCase := range []struct { + name string + extra []string + }{ + {"no devices", []string{"--seeds", "1"}}, + {"web workers", []string{"--seeds", "1", "--platform", "web", "--devices", "worker-a,worker-b"}}, + } { + t.Run(testCase.name, func(t *testing.T) { + directory := t.TempDir() + configuration := testConfiguration(t, directory, testCase.extra...) + original := connectedDevices + connectedDevices = func(context.Context) ([]string, error) { + t.Error("preflight looked for devices where the workers name none") + return nil, fmt.Errorf("no adb server") + } + t.Cleanup(func() { connectedDevices = original }) + + executor := versionAnswering(func(_ context.Context, _ string, arguments []string, _ io.Writer) (int, error) { + writeFakeRun(t, arguments, []trace.Step{observedStep(1)}) + return 0, nil + }) + if err := runCampaign(context.Background(), configuration, executor, io.Discard); err != nil { + t.Fatal(err) + } + }) + } +} + +// Preflight cannot catch a device that disappears mid-campaign. A device that +// keeps failing in seconds is drained of seeds by exactly the speed of its +// failure, so it has to stop being given any. +func TestRunCampaign_QuarantinesADeviceThatKeepsFailingFast(t *testing.T) { + directory := t.TempDir() + configuration := testConfiguration(t, directory, "--seeds", "1-8", "--devices", "device-a,device-b") + devicesPresent(t, "device-a", "device-b") + + failedEnough := make(chan struct{}) + var closeOnce sync.Once + var timedOut atomic.Bool + var mutex sync.Mutex + dispatched := map[string][]int64{} + + executor := versionAnswering(func(_ context.Context, _ string, arguments []string, output io.Writer) (int, error) { + device := argumentValue(arguments, "--device") + seed, err := strconv.ParseInt(argumentValue(arguments, "--seed"), 10, 64) + if err != nil { + t.Errorf("seed argument: %v", err) + } + mutex.Lock() + dispatched[device] = append(dispatched[device], seed) + failures := len(dispatched["device-a"]) + mutex.Unlock() + if device == "device-a" { + fmt.Fprintln(output, "device 'device-a' not found") + if failures >= fastFailuresBeforeQuarantine { + closeOnce.Do(func() { close(failedEnough) }) + } + return 1, nil + } + // The healthy worker holds its seed until the sick one has failed its + // way to quarantine, so the split of the queue is the scheduler's + // decision and not a race between two equally fast fakes. + select { + case <-failedEnough: + case <-time.After(5 * time.Second): + timedOut.Store(true) + } + writeFakeRun(t, arguments, []trace.Step{observedStep(1)}) + return 0, nil + }) + + var stdout bytes.Buffer + err := runCampaign(context.Background(), configuration, executor, &stdout) + if timedOut.Load() { + t.Fatal("device-a never reached quarantine") + } + if err == nil || !strings.Contains(err.Error(), "3 of 8 runs failed") { + t.Fatalf("expected the three fast failures to be reported, got %v", err) + } + if got := len(dispatched["device-a"]); got != fastFailuresBeforeQuarantine { + t.Errorf("seeds sent to the failing device: got %d, want %d", got, fastFailuresBeforeQuarantine) + } + if got := len(dispatched["device-b"]); got != 8-fastFailuresBeforeQuarantine { + t.Errorf("seeds sent to the healthy device: got %d, want %d", got, 8-fastFailuresBeforeQuarantine) + } + ran := append(append([]int64{}, dispatched["device-a"]...), dispatched["device-b"]...) + slices.Sort(ran) + if !slices.Equal(ran, []int64{1, 2, 3, 4, 5, 6, 7, 8}) { + t.Errorf("every seed must be dispatched exactly once: got %v", ran) + } + if !strings.Contains(stdout.String(), `quarantined device "device-a"`) { + t.Errorf("the quarantine must be reported in the campaign output:\n%s", stdout.String()) + } + + recorded := readManifest(t, directory) + if len(recorded.Quarantined) != 1 || recorded.Quarantined[0].Device != "device-a" { + t.Fatalf("the manifest must record the quarantine: %+v", recorded.Quarantined) + } + burned := slices.Clone(dispatched["device-a"]) + slices.Sort(burned) + if !slices.Equal(recorded.Quarantined[0].ConsumedSeeds, burned) { + t.Errorf("consumed seeds: got %v, want %v", recorded.Quarantined[0].ConsumedSeeds, burned) + } + if !slices.Equal(recorded.UnrunSeeds, burned) { + t.Errorf("the seeds the quarantined device burned have no result and must be reported unrun: got %v, want %v", + recorded.UnrunSeeds, burned) + } +} + +// Every device gone is not a campaign that should keep pulling seeds: the rest +// of the queue would be spent producing the same failure. +func TestRunCampaign_AbortsWhenEveryDeviceIsQuarantined(t *testing.T) { + directory := t.TempDir() + configuration := testConfiguration(t, directory, "--seeds", "1-10", "--devices", "device-a,device-b") + devicesPresent(t, "device-a", "device-b") + + var dispatches atomic.Int32 + executor := versionAnswering(func(_ context.Context, _ string, _ []string, output io.Writer) (int, error) { + dispatches.Add(1) + fmt.Fprintln(output, "sidecar health check: context deadline exceeded") + return 1, nil + }) + + var stdout bytes.Buffer + err := runCampaign(context.Background(), configuration, executor, &stdout) + if err == nil || !strings.Contains(err.Error(), "every device was quarantined") { + t.Fatalf("a campaign with no device left must abort with a clear error, got %v", err) + } + if got := dispatches.Load(); got != int32(2*fastFailuresBeforeQuarantine) { + t.Errorf("runs dispatched: got %d, want %d: the sweep must stop rather than spin through the queue", + got, 2*fastFailuresBeforeQuarantine) + } + recorded := readManifest(t, directory) + if !slices.Equal(recorded.UnrunSeeds, []int64{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}) { + t.Errorf("every seed is either burned by a quarantined device or never dispatched: got %v", recorded.UnrunSeeds) + } +} + +func readManifest(t *testing.T, campaignDirectory string) manifest { + t.Helper() + body, err := os.ReadFile(filepath.Join(campaignDirectory, manifestFileName)) + if err != nil { + t.Fatal(err) + } + var recorded manifest + if err := json.Unmarshal(body, &recorded); err != nil { + t.Fatal(err) + } + return recorded +} diff --git a/internal/android/android.go b/internal/android/android.go index 1ccd823..8411990 100644 --- a/internal/android/android.go +++ b/internal/android/android.go @@ -23,7 +23,7 @@ import ( // - else if exactly one AVD exists locally, boot it; // - else fail with a helpful message listing the available AVDs. func EnsureDevice(ctx context.Context, serial, avdName string, stdout io.Writer) error { - devices, err := listAdbDevices(ctx) + devices, err := ConnectedDevices(ctx) if err != nil { return fmt.Errorf("list adb devices: %w", err) } @@ -358,7 +358,11 @@ var standardSDKRoots = []string{ "/usr/local/share/android-commandlinetools", } -func listAdbDevices(ctx context.Context) ([]string, error) { +// ConnectedDevices lists the serials adb reports as online. It goes through +// the adb CLI so the ADB_SERVER_SOCKET / ANDROID_ADB_SERVER_ADDRESS pair the +// process was started with selects the same server every other adb call in the +// run talks to, rather than assuming a server on this machine. +func ConnectedDevices(ctx context.Context) ([]string, error) { adb, err := AdbBinary() if err != nil { return nil, err