feat(corpus-sweep): run one specification against a served corpus of implementations

same fixed campaign as implementation-sweep, over a corpus that needs no build. each implementation gets its own port: the corpus holds pairs that write the same localStorage key, and one shared origin is one stored record shared between them.
This commit is contained in:
pj committed 2026-08-16 17:45:42 +05:30
1 parent 146152a3af
commit 341d6a0614
9 files changed
+2103

No files matched your search

+174
View File
@@ -0,0 +1,174 @@
package main
import (
"fmt"
"os"
"path/filepath"
"slices"
"strings"
)
// The corpus is tastejs/todomvc. Its examples/ directory holds 48
// implementations of one requirement; five are excluded and the 43 that remain
// are this experiment's population.
//
// The population is named here rather than read from whatever directories
// happen to exist, so a corpus at the wrong commit fails verifyCorpus instead
// of quietly sweeping a different sample and reporting it as this one.
var includedImplementations = []string{
"angular-dart", "angular2", "angular2_es2015", "angularjs", "angularjs_require",
"aurelia", "backbone", "backbone_marionette", "backbone_require", "binding-scala",
"canjs", "canjs_require", "closure", "dijon", "dojo",
"duel", "elm", "emberjs", "enyo_backbone", "exoskeleton",
"jquery", "js_of_ocaml", "jsblocks", "knockback", "knockoutjs",
"knockoutjs_require", "kotlin-react", "lavaca_require", "mithril", "polymer",
"ractive", "react", "react-alt", "react-backbone", "reagent",
"riotjs", "scalajs-react", "typescript-angular", "typescript-backbone", "typescript-react",
"vanilla-es6", "vanillajs", "vue",
}
// excludedImplementations are the five directories under examples/ that the
// corpus survey dropped. They are listed rather than merely omitted so
// verifyCorpus can insist the corpus holds exactly these 48 names: an example
// added or renamed upstream then stops the sweep rather than silently shrinking
// or growing the sample.
var excludedImplementations = []string{
"cujo", "emberjs_require", "firebase-angular", "gwt", "react-hooks",
}
const examplesDirectory = "examples"
// documentPath names the served document for the implementations that do not
// keep an index.html at the root of their example directory.
var documentPath = map[string]string{
"angular-dart": "examples/angular-dart/web/index.html",
"duel": "examples/duel/www/index.html",
}
func documentFor(name string) string {
if override, ok := documentPath[name]; ok {
return override
}
return examplesDirectory + "/" + name + "/index.html"
}
// implementation is one member of the population, with the port it owns for the
// whole sweep. The port is what keeps implementations apart: every one is
// served on its own, so every one is its own web origin and localStorage keeps
// their records in separate partitions. Four pairs in this corpus write the
// same key, and one origin between them is one record between them.
type implementation struct {
Name string
Document string
Port int
}
func (i implementation) Origin() string {
return fmt.Sprintf("http://127.0.0.1:%d", i.Port)
}
func (i implementation) URL() string {
return i.Origin() + "/" + i.Document
}
// verifyCorpus insists the corpus holds exactly the 48 example directories this
// population was drawn from. A sweep against a different checkout would still
// run, and its 43 arms would still be labelled, which is why the check is here
// and not left to whoever reads the results.
func verifyCorpus(corpusRoot string) error {
entries, err := os.ReadDir(filepath.Join(corpusRoot, examplesDirectory))
if err != nil {
return fmt.Errorf("--corpus: %w", err)
}
var found []string
for _, entry := range entries {
if entry.IsDir() {
found = append(found, entry.Name())
}
}
expected := slices.Concat(includedImplementations, excludedImplementations)
slices.Sort(expected)
slices.Sort(found)
if slices.Equal(expected, found) {
return nil
}
var missing, unexpected []string
for _, name := range expected {
if !slices.Contains(found, name) {
missing = append(missing, name)
}
}
for _, name := range found {
if !slices.Contains(expected, name) {
unexpected = append(unexpected, name)
}
}
return fmt.Errorf(
"%s/%s is not the corpus this population was drawn from: missing %v, unexpected %v",
corpusRoot,
examplesDirectory,
missing,
unexpected,
)
}
// selectImplementations resolves --implementations against the population. An
// empty selection is the whole population, which is what a real sweep runs; a
// named subset is for smoke runs and is recorded in the manifest like any other
// intent.
func selectImplementations(selection string) ([]string, error) {
if strings.TrimSpace(selection) == "" {
return slices.Clone(includedImplementations), nil
}
var names []string
seen := map[string]bool{}
for _, part := range strings.Split(selection, ",") {
name := strings.TrimSpace(part)
if name == "" {
return nil, fmt.Errorf("empty implementation in %q", selection)
}
if !slices.Contains(includedImplementations, name) {
return nil, fmt.Errorf(
"%q is not one of the %d implementations in this population",
name,
len(includedImplementations),
)
}
if seen[name] {
return nil, fmt.Errorf("duplicate implementation %q", name)
}
seen[name] = true
names = append(names, name)
}
return names, nil
}
// planImplementations gives every selected implementation its document and its
// own port, in population order so the manifest can name the URL each arm was
// served from before anything has been served.
func planImplementations(
corpusRoot string,
names []string,
basePort int,
) ([]implementation, error) {
if basePort+len(names)-1 > 65535 {
return nil, fmt.Errorf(
"--base-port %d leaves no room for %d implementations",
basePort,
len(names),
)
}
planned := make([]implementation, 0, len(names))
for index, name := range names {
document := documentFor(name)
if _, err := os.Stat(filepath.Join(corpusRoot, filepath.FromSlash(document))); err != nil {
return nil, fmt.Errorf("%s: %w", name, err)
}
planned = append(planned, implementation{
Name: name,
Document: document,
Port: basePort + index,
})
}
return planned, nil
}
@@ -0,0 +1,161 @@
package main
import (
"os"
"path/filepath"
"slices"
"strings"
"testing"
)
// writeCorpus builds a tree with the shape the sweep expects: every example
// directory the population was drawn from, each holding the document that
// implementation is served at.
func writeCorpus(t *testing.T, names []string) string {
t.Helper()
root := t.TempDir()
for _, name := range names {
if err := os.MkdirAll(filepath.Join(root, examplesDirectory, name), 0o755); err != nil {
t.Fatal(err)
}
document := filepath.Join(root, filepath.FromSlash(documentFor(name)))
if err := os.MkdirAll(filepath.Dir(document), 0o755); err != nil {
t.Fatal(err)
}
body := "<!doctype html><title>" + name + "</title><ul class=\"todo-list\"></ul>"
if err := os.WriteFile(document, []byte(body), 0o644); err != nil {
t.Fatal(err)
}
}
return root
}
func wholeCorpus(t *testing.T) string {
t.Helper()
return writeCorpus(
t,
slices.Concat(includedImplementations, excludedImplementations),
)
}
func TestPopulation_IsFortyThreeNamesDisjointFromTheExclusions(t *testing.T) {
if len(includedImplementations) != 43 {
t.Errorf(
"population size: got %d, want 43",
len(includedImplementations),
)
}
seen := map[string]bool{}
for _, name := range includedImplementations {
if seen[name] {
t.Errorf("%q appears twice in the population", name)
}
seen[name] = true
}
for _, name := range excludedImplementations {
if seen[name] {
t.Errorf("%q is both included and excluded", name)
}
}
}
func TestVerifyCorpus_RejectsATreeThatIsNotThisCorpus(t *testing.T) {
if err := verifyCorpus(wholeCorpus(t)); err != nil {
t.Fatalf(
"the corpus this population was drawn from should verify: %v",
err,
)
}
short := writeCorpus(
t,
slices.Concat(includedImplementations[1:], excludedImplementations),
)
err := verifyCorpus(short)
if err == nil {
t.Fatal(
"a corpus missing an implementation swept 42 arms and reported them as 43",
)
}
if !strings.Contains(err.Error(), includedImplementations[0]) {
t.Errorf("the error should name what is missing: %v", err)
}
extra := writeCorpus(
t,
slices.Concat(
includedImplementations,
excludedImplementations,
[]string{"svelte"},
),
)
err = verifyCorpus(extra)
if err == nil {
t.Fatal(
"a corpus with an implementation this population never drew from verified",
)
}
if !strings.Contains(err.Error(), "svelte") {
t.Errorf("the error should name what is unexpected: %v", err)
}
}
func TestDocumentFor_ServesTheOverriddenPathForImplementationsWithoutARootIndex(
t *testing.T,
) {
for name, want := range map[string]string{
"angular-dart": "examples/angular-dart/web/index.html",
"duel": "examples/duel/www/index.html",
"vanillajs": "examples/vanillajs/index.html",
"react": "examples/react/index.html",
} {
if got := documentFor(name); got != want {
t.Errorf("%s document: got %q, want %q", name, got, want)
}
}
}
func TestPlanImplementations_RefusesAnImplementationWhoseDocumentIsMissing(
t *testing.T,
) {
root := wholeCorpus(t)
if err := os.Remove(filepath.Join(root, filepath.FromSlash(documentFor("duel")))); err != nil {
t.Fatal(err)
}
_, err := planImplementations(root, []string{"duel"}, 5400)
if err == nil {
t.Fatal(
"an implementation whose document is missing would be swept as a 404 page",
)
}
if !strings.Contains(err.Error(), "duel") {
t.Errorf("the error should name the implementation: %v", err)
}
}
func TestSelectImplementations_DefaultsToThePopulationAndRejectsNamesOutsideIt(
t *testing.T,
) {
all, err := selectImplementations("")
if err != nil {
t.Fatal(err)
}
if !slices.Equal(all, includedImplementations) {
t.Errorf(
"an empty selection should be the whole population, got %d names",
len(all),
)
}
subset, err := selectImplementations("react, angular2_es2015")
if err != nil {
t.Fatal(err)
}
if !slices.Equal(subset, []string{"react", "angular2_es2015"}) {
t.Errorf("subset: got %v", subset)
}
for _, selection := range []string{"cujo", "svelte", "react,react"} {
if _, err := selectImplementations(selection); err == nil {
t.Errorf("selection %q should have been rejected", selection)
}
}
}
@@ -0,0 +1,468 @@
package main
import (
"bufio"
"bytes"
"encoding/json"
"fmt"
"io"
"math/rand"
"net"
"net/http"
"os"
"path/filepath"
"slices"
"strings"
"testing"
"time"
)
// The test binary doubles as the fetcher the stub campaign uses, so the URL a
// campaign is handed is really requested while the sweep is serving, and what
// came back is on disk for the test to read.
func TestMain(m *testing.M) {
if url := os.Getenv("CORPUS_SWEEP_TEST_FETCH_URL"); url != "" {
recordFetch(url, os.Getenv("CORPUS_SWEEP_TEST_FETCH_LOG"))
return
}
os.Exit(m.Run())
}
func recordFetch(url, logPath string) {
line := ""
response, err := http.Get(url)
if err != nil {
line = fmt.Sprintf("%s -> error %v\n", url, err)
} else {
body, _ := io.ReadAll(response.Body)
response.Body.Close()
line = fmt.Sprintf("%s -> %d %s\n", url, response.StatusCode, body)
}
logFile, err := os.OpenFile(
logPath,
os.O_CREATE|os.O_WRONLY|os.O_APPEND,
0o644,
)
if err != nil {
return
}
logFile.WriteString(line)
logFile.Close()
}
// stubCampaign records the argv it was handed and whether the sweep manifest
// was already on disk when it ran, fetches the URL it was told to drive, writes
// the campaign directory the real tool would write, and fails one seed.
const stubCampaign = `#!/bin/sh
output=""
seed=""
url=""
previous=""
for argument in "$@"; do
case "$previous" in
--output) output="$argument" ;;
--seeds) seed="$argument" ;;
--bundle-id) url="$argument" ;;
esac
previous="$argument"
done
manifest=missing
if [ -f "%[1]s" ]; then manifest=present; fi
echo "manifest=$manifest argv: $*" >> "%[2]s"
CORPUS_SWEEP_TEST_FETCH_URL="$url" CORPUS_SWEEP_TEST_FETCH_LOG="%[3]s" "%[4]s"
mkdir -p "$output"
printf '{"seeds":[%%s]}\n' "$seed" > "$output/campaign.json"
echo "stub campaign seed=$seed url=$url"
case "$output" in
*angular2_es2015/seed-5) exit 1 ;;
esac
exit 0
`
func TestRun_EndToEndAgainstAServedCorpusAndAStubCampaign(t *testing.T) {
root := t.TempDir()
corpus := wholeCorpus(t)
// dojo's document is served but cannot be read, so it stalls at serve and
// the two implementations either side of it still have to reach the
// campaign tool with both seeds.
unreadable := filepath.Join(corpus, filepath.FromSlash(documentFor("dojo")))
if err := os.Chmod(unreadable, 0o000); err != nil {
t.Fatal(err)
}
t.Cleanup(func() { os.Chmod(unreadable, 0o644) })
specPath := filepath.Join(root, "todo.ts")
if err := os.WriteFile(specPath, []byte("export const properties = [];\n"), 0o644); err != nil {
t.Fatal(err)
}
output := filepath.Join(root, "campaigns")
campaignLog := filepath.Join(root, "campaign.log")
fetchLog := filepath.Join(root, "fetch.log")
testBinary, err := filepath.Abs(os.Args[0])
if err != nil {
t.Fatal(err)
}
campaignPath := writeScript(
t,
filepath.Join(root, "stub-campaign"),
fmt.Sprintf(
stubCampaign,
filepath.Join(output, manifestFileName),
campaignLog,
fetchLog,
testBinary,
),
)
sanderlingPath := writeScript(
t,
filepath.Join(root, "stub-sanderling"),
"#!/bin/sh\nexit 0\n",
)
selected := []string{"angular2", "angular2_es2015", "dojo"}
basePort := freePortRange(t, len(selected))
var stdout bytes.Buffer
err = run([]string{
"--corpus", corpus,
"--spec", specPath,
"--implementations", strings.Join(selected, ","),
"--seeds", "4-5",
"--max-steps", "40",
"--duration", "30s",
"--concurrency", "2",
"--base-port", fmt.Sprint(basePort),
"--output", output,
"--campaign", campaignPath,
"--sanderling", sanderlingPath,
}, &stdout, io.Discard)
if err == nil ||
!strings.Contains(err.Error(), "1 of 3 implementations never ran") {
t.Fatalf(
"expected the stalled implementation and the failed campaign to be reported, got %v",
err,
)
}
recorded := readManifest(t, filepath.Join(output, manifestFileName))
if !slices.Equal(recorded.Seeds, []int64{4, 5}) {
t.Errorf("intended seeds: got %v", recorded.Seeds)
}
if len(recorded.Implementations) != len(selected) {
t.Fatalf("intended implementations: got %v", recorded.Implementations)
}
for index, planned := range recorded.Implementations {
if planned.Name != selected[index] {
t.Errorf(
"implementation %d: got %q, want %q",
index,
planned.Name,
selected[index],
)
}
wantPort := basePort + index
wantURL := fmt.Sprintf(
"http://127.0.0.1:%d/examples/%s/index.html",
wantPort,
planned.Name,
)
if planned.Port != wantPort || planned.URL != wantURL {
t.Errorf(
"%s: port %d url %q, want %d %q",
planned.Name,
planned.Port,
planned.URL,
wantPort,
wantURL,
)
}
}
if recorded.Generator != "seeded" || recorded.Platform != "web" ||
recorded.MaxSteps != 40 {
t.Errorf(
"manifest generator/platform/budget: %q/%q/%d",
recorded.Generator,
recorded.Platform,
recorded.MaxSteps,
)
}
if recorded.CorpusRoot != corpus {
t.Errorf(
"manifest corpus root: got %q, want %q",
recorded.CorpusRoot,
corpus,
)
}
campaignLines := readLines(t, campaignLog)
if len(campaignLines) != 4 {
t.Fatalf(
"campaign invocations: got %d, want 4:\n%s",
len(campaignLines),
strings.Join(campaignLines, "\n"),
)
}
seen := map[string]bool{}
for _, line := range campaignLines {
if !strings.HasPrefix(line, "manifest=present") {
t.Errorf(
"a campaign ran before the sweep manifest was written: %q",
line,
)
}
arguments := strings.Fields(strings.SplitN(line, "argv: ", 2)[1])
arm := argumentValue(arguments, "--arm")
seed := argumentValue(arguments, "--seeds")
port := basePort + slices.Index(selected, arm)
wantBundle := fmt.Sprintf(
"http://127.0.0.1:%d/examples/%s/index.html",
port,
arm,
)
if got := argumentValue(arguments, "--bundle-id"); got != wantBundle {
t.Errorf(
"%s seed %s: bundle id %q, want %q",
arm,
seed,
got,
wantBundle,
)
}
if got := argumentValue(arguments, "--output"); got != filepath.Join(
output,
arm,
"seed-"+seed,
) {
t.Errorf("%s seed %s: campaign output %q", arm, seed, got)
}
if got := argumentValue(arguments, "--sanderling"); got != sanderlingPath {
t.Errorf("%s seed %s: sanderling path %q", arm, seed, got)
}
seen[arm+"/"+seed] = true
}
for _, want := range []string{"angular2/4", "angular2/5", "angular2_es2015/4", "angular2_es2015/5"} {
if !seen[want] {
t.Errorf("%s never reached the campaign tool", want)
}
}
// What the served page actually returned: each port answered with its own
// implementation's document, so no arm was driven against another's.
fetched := readLines(t, fetchLog)
for index, name := range []string{"angular2", "angular2_es2015"} {
want := fmt.Sprintf(
"http://127.0.0.1:%d/examples/%s/index.html -> 200 <!doctype html><title>%s</title>",
basePort+index,
name,
name,
)
matched := 0
for _, line := range fetched {
if strings.HasPrefix(line, want) {
matched++
}
}
if matched != 2 {
t.Errorf(
"%s was served its own document %d times, want 2:\n%s",
name,
matched,
strings.Join(fetched, "\n"),
)
}
}
records := readRecords(t, filepath.Join(output, recordsFileName))
if len(records) != 3 {
t.Fatalf("implementation records: got %d, want 3", len(records))
}
byName := map[string]implementationRecord{}
for _, record := range records {
byName[record.Name] = record
}
stalled := byName["dojo"]
if stalled.FailedStage != stageServe || len(stalled.Runs) != 0 {
t.Errorf(
"dojo: stage %q with %d runs, want a serve failure and no runs",
stalled.FailedStage,
len(stalled.Runs),
)
}
for _, name := range []string{"angular2", "angular2_es2015"} {
record := byName[name]
if record.FailedStage != "" || len(record.Runs) != 2 {
t.Errorf(
"%s: stage %q with %d runs, want no failure and 2 runs",
name,
record.FailedStage,
len(record.Runs),
)
}
if record.MonotonicMillis <= 0 {
t.Errorf(
"%s took %d ms, so nothing timed how long it worked",
name,
record.MonotonicMillis,
)
}
}
if exit := byName["angular2_es2015"].Runs[1].ExitCode; exit != 1 {
t.Errorf("angular2_es2015 seed 5 exit code: got %d, want 1", exit)
}
if exit := byName["angular2_es2015"].Runs[0].ExitCode; exit != 0 {
t.Errorf("angular2_es2015 seed 4 exit code: got %d, want 0", exit)
}
for _, name := range []string{"angular2", "angular2_es2015"} {
for _, seed := range []string{"4", "5"} {
directory := filepath.Join(output, name, "seed-"+seed)
if _, err := os.Stat(filepath.Join(directory, "campaign.json")); err != nil {
t.Errorf(
"%s seed %s: no campaign directory: %v",
name,
seed,
err,
)
}
log, err := os.ReadFile(filepath.Join(directory, "campaign.log"))
if err != nil {
t.Fatal(err)
}
if !strings.Contains(string(log), "stub campaign seed="+seed) {
t.Errorf(
"%s seed %s: campaign output was not captured: %q",
name,
seed,
log,
)
}
}
}
client := &http.Client{Timeout: 2 * time.Second}
for offset := range selected {
if response, err := client.Get(fmt.Sprintf("http://127.0.0.1:%d/", basePort+offset)); err == nil {
response.Body.Close()
t.Errorf(
"port %d is still served after the sweep finished",
basePort+offset,
)
}
}
if !strings.Contains(stdout.String(), "failed at serve") {
t.Errorf(
"progress output does not name the stalled implementation: %q",
stdout.String(),
)
}
}
func TestRun_RefusesAnOutputDirectoryThatAlreadyHoldsASweep(t *testing.T) {
root := t.TempDir()
output := filepath.Join(root, "campaigns")
if err := os.MkdirAll(output, 0o755); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(output, manifestFileName), []byte("{}"), 0o644); err != nil {
t.Fatal(err)
}
specPath := filepath.Join(root, "todo.ts")
if err := os.WriteFile(specPath, []byte("export const properties = [];\n"), 0o644); err != nil {
t.Fatal(err)
}
err := run([]string{
"--corpus", wholeCorpus(t), "--spec", specPath, "--seeds", "1",
"--max-steps", "10", "--output", output,
}, io.Discard, io.Discard)
if err == nil || !strings.Contains(err.Error(), manifestFileName) {
t.Fatalf(
"two sweeps sharing a directory would interleave their records, got %v",
err,
)
}
}
func writeScript(t *testing.T, path, body string) string {
t.Helper()
if err := os.WriteFile(path, []byte(body), 0o755); err != nil {
t.Fatal(err)
}
return path
}
func readLines(t *testing.T, path string) []string {
t.Helper()
body, err := os.ReadFile(path)
if err != nil {
t.Fatal(err)
}
var lines []string
scanner := bufio.NewScanner(strings.NewReader(string(body)))
for scanner.Scan() {
if line := strings.TrimSpace(scanner.Text()); line != "" {
lines = append(lines, line)
}
}
return lines
}
func readManifest(t *testing.T, path string) manifest {
t.Helper()
body, err := os.ReadFile(path)
if err != nil {
t.Fatal(err)
}
var recorded manifest
if err := json.Unmarshal(body, &recorded); err != nil {
t.Fatal(err)
}
return recorded
}
func readRecords(t *testing.T, path string) []implementationRecord {
t.Helper()
var records []implementationRecord
for _, line := range readLines(t, path) {
var record implementationRecord
if err := json.Unmarshal([]byte(line), &record); err != nil {
t.Fatalf("%s: %v", line, err)
}
records = append(records, record)
}
return records
}
func argumentValue(arguments []string, name string) string {
index := slices.Index(arguments, name)
if index < 0 || index+1 >= len(arguments) {
return ""
}
return arguments[index+1]
}
// freePortRange finds count consecutive free ports, which is what the sweep
// hands out: one port per implementation from --base-port upwards.
func freePortRange(t *testing.T, count int) int {
t.Helper()
for range 100 {
base := 20000 + rand.Intn(20000)
if portsAreFree(base, count) {
return base
}
}
t.Fatalf("no run of %d free ports", count)
return 0
}
func portsAreFree(base, count int) bool {
for offset := range count {
listener, err := net.Listen(
"tcp",
fmt.Sprintf("127.0.0.1:%d", base+offset),
)
if err != nil {
return false
}
listener.Close()
}
return true
}
+275
View File
@@ -0,0 +1,275 @@
// Command corpus-sweep runs one specification against every implementation in a
// served corpus of independent implementations of the same requirement. It
// serves each one on its own port, which is what keeps them apart: the corpus
// holds pairs that write the same localStorage key, and one origin shared
// between two of them is one stored record shared between them.
package main
import (
"context"
"errors"
"flag"
"fmt"
"io"
"os"
"os/signal"
"path/filepath"
"strconv"
"syscall"
"time"
"github.com/priyanshujain/sanderling/internal/seedspec"
)
// The generator and the platform are fixed rather than exposed: the
// pre-registration runs the seeded policy against the served corpus, and a
// sweep that could quietly run something else records a comparison nobody made.
const (
generator = "seeded"
platform = "web"
)
// defaultConcurrency is how many implementations are swept at once. fleet.md
// measured eight concurrent web campaigns clean at about 1.1 GB resident each,
// parallel efficiency 0.83 at eight against 0.87 at six, and said to re-measure
// before trusting anything above eight on a contended host. Six sits at the
// better efficiency and leaves the emulator farm that shares the host its
// slots. Serving costs nothing here: the corpus needs no build and no separate
// server process, so a worker is one browser.
const defaultConcurrency = 6
// defaultBasePort starts above the range the model-implementation sweep hands
// out, so the two can run on one host without either being served the other's
// pages.
const defaultBasePort = 5400
type config struct {
corpusRoot string
specPath string
outputDirectory string
implementations []string
seeds []int64
maxSteps int
duration time.Duration
concurrency int
basePort int
campaignPath string
sanderlingPath string
extraArguments []string
}
const usage = `corpus-sweep runs one seeded campaign per implementation of a served corpus.
Usage:
corpus-sweep --corpus <dir> --spec <path> --seeds <spec>
--max-steps <n> --output <dir> [flags]
[-- <sanderling test flags>]
Every implementation is served on its own port, so every one is its own origin
and none can read another's stored state, and every one is swept with the same
seeds and the same step budget. Each run's arm is the implementation's name.
Everything after a bare -- reaches every sanderling test call through the
campaign tool.
`
func parseArguments(arguments []string, stderr io.Writer) (config, error) {
flagSet := flag.NewFlagSet("corpus-sweep", flag.ContinueOnError)
flagSet.SetOutput(stderr)
flagSet.Usage = func() {
fmt.Fprint(stderr, usage)
flagSet.PrintDefaults()
}
var configuration config
var seedSpecification string
var selection string
flagSet.StringVar(
&configuration.corpusRoot,
"corpus",
"",
"root of the checked-out corpus, holding examples/ (required)",
)
flagSet.StringVar(
&configuration.specPath,
"spec",
"",
"path to the property set every implementation is run against (required)",
)
flagSet.StringVar(
&seedSpecification,
"seeds",
"",
"seeds every implementation runs: ranges and lists, e.g. 1-10,20 (required)",
)
flagSet.StringVar(
&selection,
"implementations",
"",
"comma-separated subset of the population to sweep (default: all of it)",
)
flagSet.IntVar(
&configuration.maxSteps,
"max-steps",
0,
"per-run step budget, identical across implementations (required, must be positive)",
)
flagSet.DurationVar(
&configuration.duration,
"duration",
5*time.Minute,
"per-run wall-clock ceiling passed to each campaign",
)
flagSet.IntVar(
&configuration.concurrency,
"concurrency",
defaultConcurrency,
"implementations swept at once",
)
flagSet.IntVar(
&configuration.basePort,
"base-port",
defaultBasePort,
"first port served; each implementation takes the next one in population order",
)
flagSet.StringVar(
&configuration.outputDirectory,
"output",
"",
"campaign tree to create (required)",
)
flagSet.StringVar(
&configuration.campaignPath,
"campaign",
"campaign",
"campaign binary to invoke per implementation and seed",
)
flagSet.StringVar(
&configuration.sanderlingPath,
"sanderling",
"sanderling",
"sanderling binary each campaign invokes",
)
if err := flagSet.Parse(arguments); err != nil {
return config{}, err
}
configuration.extraArguments = flagSet.Args()
for name, value := range map[string]string{
"--corpus": configuration.corpusRoot,
"--spec": configuration.specPath,
"--seeds": seedSpecification,
"--output": configuration.outputDirectory,
} {
if value == "" {
return config{}, fmt.Errorf("%s is required", name)
}
}
if configuration.maxSteps <= 0 {
return config{}, fmt.Errorf(
"--max-steps must be positive: every implementation needs the same step budget",
)
}
if configuration.duration <= 0 {
return config{}, fmt.Errorf(
"--duration must be positive: %s",
configuration.duration,
)
}
if configuration.concurrency <= 0 {
return config{}, fmt.Errorf(
"--concurrency must be positive: %d",
configuration.concurrency,
)
}
if configuration.basePort < 1024 || configuration.basePort > 65535 {
return config{}, fmt.Errorf(
"--base-port %d is outside 1024-65535",
configuration.basePort,
)
}
seeds, err := seedspec.Parse(seedSpecification)
if err != nil {
return config{}, fmt.Errorf("--seeds: %w", err)
}
configuration.seeds = seeds
names, err := selectImplementations(selection)
if err != nil {
return config{}, fmt.Errorf("--implementations: %w", err)
}
configuration.implementations = names
for name, value := range map[string]*string{
"--corpus": &configuration.corpusRoot,
"--spec": &configuration.specPath,
"--output": &configuration.outputDirectory,
} {
absolute, err := filepath.Abs(*value)
if err != nil {
return config{}, fmt.Errorf("%s: %w", name, err)
}
*value = absolute
}
return configuration, nil
}
func campaignDirectory(
configuration config,
target implementation,
seed string,
) string {
return filepath.Join(
configuration.outputDirectory,
target.Name,
"seed-"+seed,
)
}
// campaignArguments builds one campaign invocation. The arm is the
// implementation's name, so every run in the tree can be attributed to the
// implementation it came from without reading back which port served it.
func campaignArguments(
configuration config,
target implementation,
seed string,
) []string {
arguments := []string{
"--spec", configuration.specPath,
"--bundle-id", target.URL(),
"--platform", platform,
"--arm", target.Name,
"--generator", generator,
"--max-steps", strconv.Itoa(configuration.maxSteps),
"--duration", configuration.duration.String(),
"--seeds", seed,
"--sanderling", configuration.sanderlingPath,
"--output", campaignDirectory(configuration, target, seed),
}
if len(configuration.extraArguments) > 0 {
arguments = append(arguments, "--")
arguments = append(arguments, configuration.extraArguments...)
}
return arguments
}
func run(arguments []string, stdout, stderr io.Writer) error {
configuration, err := parseArguments(arguments, stderr)
if err != nil {
return err
}
ctx, cancel := signal.NotifyContext(
context.Background(),
os.Interrupt,
syscall.SIGTERM,
)
defer cancel()
return runSweep(ctx, configuration, stdout)
}
func main() {
if err := run(os.Args[1:], os.Stdout, os.Stderr); err != nil {
if errors.Is(err, flag.ErrHelp) {
return
}
fmt.Fprintf(os.Stderr, "error: %v\n", err)
os.Exit(1)
}
}
+132
View File
@@ -0,0 +1,132 @@
package main
import (
"encoding/json"
"fmt"
"os"
"os/exec"
"path/filepath"
"strings"
"time"
)
const (
manifestFileName = "sweep.json"
recordsFileName = "implementations.jsonl"
)
// plannedImplementation is one implementation the sweep intends to run, with
// the origin that keeps its stored state its own and the URL every seed is
// driven at.
type plannedImplementation struct {
Name string `json:"name"`
Document string `json:"document"`
Port int `json:"port"`
Origin string `json:"origin"`
URL string `json:"url"`
}
// manifest is sweep.json: what the sweep INTENDED to run, written before the
// first run so a host that dropped an implementation or a seed shows up as a
// missing run rather than as a smaller sample.
type manifest struct {
Generator string `json:"generator"`
Platform string `json:"platform"`
SpecPath string `json:"spec_path"`
CorpusRoot string `json:"corpus_root"`
CorpusCommit string `json:"corpus_commit"`
MaxSteps int `json:"max_steps"`
DurationMillis int64 `json:"duration_millis"`
Seeds []int64 `json:"seeds"`
Implementations []plannedImplementation `json:"implementations"`
Concurrency int `json:"concurrency"`
Host string `json:"host"`
CampaignPath string `json:"campaign_path"`
SanderlingPath string `json:"sanderling_path"`
StartedAt time.Time `json:"started_at"`
}
func buildManifest(
configuration config,
implementations []implementation,
host string,
startedAt time.Time,
) manifest {
planned := make([]plannedImplementation, 0, len(implementations))
for _, target := range implementations {
planned = append(planned, plannedImplementation{
Name: target.Name,
Document: target.Document,
Port: target.Port,
Origin: target.Origin(),
URL: target.URL(),
})
}
return manifest{
Generator: generator,
Platform: platform,
SpecPath: configuration.specPath,
CorpusRoot: configuration.corpusRoot,
CorpusCommit: corpusCommit(configuration.corpusRoot),
MaxSteps: configuration.maxSteps,
DurationMillis: configuration.duration.Milliseconds(),
Seeds: configuration.seeds,
Implementations: planned,
Concurrency: configuration.concurrency,
Host: host,
CampaignPath: configuration.campaignPath,
SanderlingPath: configuration.sanderlingPath,
StartedAt: startedAt,
}
}
// corpusCommit records which checkout was swept. It is empty rather than fatal
// for a corpus that is not a git working tree, because the population check has
// already established the corpus holds the right examples.
func corpusCommit(corpusRoot string) string {
command := exec.Command("git", "-C", corpusRoot, "rev-parse", "HEAD")
output, err := command.Output()
if err != nil {
return ""
}
return strings.TrimSpace(string(output))
}
func writeManifest(directory string, value manifest) error {
body, err := json.MarshalIndent(value, "", " ")
if err != nil {
return fmt.Errorf("marshal manifest: %w", err)
}
return os.WriteFile(
filepath.Join(directory, manifestFileName),
append(body, '\n'),
0o644,
)
}
// runRecord is one campaign, which is one implementation at one seed.
type runRecord struct {
Seed int64 `json:"seed"`
ExitCode int `json:"exit_code"`
LaunchError string `json:"launch_error,omitempty"`
CampaignDirectory string `json:"campaign_directory"`
MonotonicMillis int64 `json:"monotonic_millis"`
}
// implementationRecord is one line of implementations.jsonl. FailedStage names
// the step that stopped this implementation, and one that never got served
// carries no runs at all.
type implementationRecord struct {
Name string `json:"implementation"`
Document string `json:"document"`
Port int `json:"port"`
Origin string `json:"origin"`
URL string `json:"url"`
FailedStage string `json:"failed_stage,omitempty"`
Error string `json:"error,omitempty"`
StartedAt time.Time `json:"started_at"`
MonotonicMillis int64 `json:"monotonic_millis"`
Runs []runRecord `json:"runs"`
}
const stageServe = "serve"
@@ -0,0 +1,194 @@
package main
import (
"fmt"
"io"
"net/url"
"os"
"path/filepath"
"strings"
"testing"
)
// collisionPairs are implementations in this corpus that write the same
// localStorage key as each other. Served from one origin they share one stored
// record: in the corpus survey angular2_es2015 crashed at bootstrap on a record
// angular2 had written, which is a violation belonging to no implementation.
var collisionPairs = [][2]string{
{"angular2", "angular2_es2015"},
{"backbone", "backbone_require"},
{"canjs", "canjs_require"},
{"react", "typescript-react"},
}
func originOf(t *testing.T, rawURL string) string {
t.Helper()
parsed, err := url.Parse(rawURL)
if err != nil {
t.Fatalf("parse %q: %v", rawURL, err)
}
return parsed.Scheme + "://" + parsed.Host
}
// The whole population is swept so the assertion covers every implementation
// rather than the four pairs already known to collide: the survey found those
// four, and an unexamined fifth would be just as damaging.
func TestSweep_GivesEveryImplementationAnOriginNoOtherImplementationShares(
t *testing.T,
) {
root := t.TempDir()
corpus := wholeCorpus(t)
specPath := filepath.Join(root, "todo.ts")
if err := os.WriteFile(specPath, []byte("export const properties = [];\n"), 0o644); err != nil {
t.Fatal(err)
}
output := filepath.Join(root, "campaigns")
campaignLog := filepath.Join(root, "campaign.log")
fetchLog := filepath.Join(root, "fetch.log")
testBinary, err := filepath.Abs(os.Args[0])
if err != nil {
t.Fatal(err)
}
campaignPath := writeScript(
t,
filepath.Join(root, "stub-campaign"),
fmt.Sprintf(
stubCampaign,
filepath.Join(output, manifestFileName),
campaignLog,
fetchLog,
testBinary,
),
)
sanderlingPath := writeScript(
t,
filepath.Join(root, "stub-sanderling"),
"#!/bin/sh\nexit 0\n",
)
err = run([]string{
"--corpus", corpus,
"--spec", specPath,
"--seeds", "1",
"--max-steps", "10",
"--duration", "30s",
"--concurrency", "8",
"--base-port", fmt.Sprint(freePortRange(t, len(includedImplementations))),
"--output", output,
"--campaign", campaignPath,
"--sanderling", sanderlingPath,
}, io.Discard, io.Discard)
if err != nil {
t.Fatalf("sweep: %v", err)
}
// What the sweep declared it would serve each implementation from.
recorded := readManifest(t, filepath.Join(output, manifestFileName))
if len(recorded.Implementations) != len(includedImplementations) {
t.Fatalf(
"intended implementations: got %d, want %d",
len(recorded.Implementations),
len(includedImplementations),
)
}
plannedOrigin := map[string]string{}
ownerOfOrigin := map[string]string{}
for _, planned := range recorded.Implementations {
if owner, taken := ownerOfOrigin[planned.Origin]; taken {
t.Errorf(
"%s and %s are both served from %s, so they share one localStorage",
owner,
planned.Name,
planned.Origin,
)
}
ownerOfOrigin[planned.Origin] = planned.Name
plannedOrigin[planned.Name] = planned.Origin
if got := originOf(t, planned.URL); got != planned.Origin {
t.Errorf(
"%s: url %q is not under the origin %q the manifest claims",
planned.Name,
planned.URL,
planned.Origin,
)
}
}
for _, pair := range collisionPairs {
if plannedOrigin[pair[0]] == plannedOrigin[pair[1]] {
t.Errorf(
"%s and %s write the same localStorage key and are both served from %s",
pair[0],
pair[1],
plannedOrigin[pair[0]],
)
}
}
// What actually reached the driver. The web driver parses the origin to
// clear out of --bundle-id, so two arms sharing a bundle-id origin clear
// and repopulate one another's storage however carefully they are labelled.
drivenOrigin := map[string]string{}
for _, line := range readLines(t, campaignLog) {
arguments := strings.Fields(strings.SplitN(line, "argv: ", 2)[1])
arm := argumentValue(arguments, "--arm")
if arm == "" {
t.Fatalf(
"a campaign ran with no arm, so its runs cannot be attributed: %q",
line,
)
}
drivenOrigin[arm] = originOf(t, argumentValue(arguments, "--bundle-id"))
}
if len(drivenOrigin) != len(includedImplementations) {
t.Fatalf(
"arms that reached the campaign tool: got %d, want %d",
len(drivenOrigin),
len(includedImplementations),
)
}
armOfOrigin := map[string]string{}
for arm, origin := range drivenOrigin {
if other, taken := armOfOrigin[origin]; taken {
t.Errorf(
"arms %s and %s were both driven at %s",
other,
arm,
origin,
)
}
armOfOrigin[origin] = arm
if origin != plannedOrigin[arm] {
t.Errorf(
"%s was driven at %s but the manifest promised %s",
arm,
origin,
plannedOrigin[arm],
)
}
}
// And what each origin answered with, which is the check that the ports
// are not merely distinct but each carries its own implementation.
served := map[string]string{}
for _, line := range readLines(t, fetchLog) {
requested, body, found := strings.Cut(line, " -> 200 ")
if !found {
t.Errorf("a served page did not answer: %q", line)
continue
}
served[originOf(t, requested)] = body
}
for arm, origin := range drivenOrigin {
if want := "<title>" + arm + "</title>"; !strings.Contains(
served[origin],
want,
) {
t.Errorf(
"%s at %s was served %q, which is not its own document",
arm,
origin,
served[origin],
)
}
}
}
+178
View File
@@ -0,0 +1,178 @@
package main
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"io/fs"
"net"
"net/http"
"path"
"path/filepath"
"regexp"
"strings"
"time"
)
const (
readinessTimeout = 5 * time.Second
readinessInterval = 100 * time.Millisecond
)
// stubScriptPath is served from every implementation's own origin in place of a
// dependency that no longer answers.
const stubScriptPath = "/__corpus-sweep__/stub.js"
// deadDependency is polyfill.io, which was shut down after this corpus was
// pinned. One implementation loads it from its index.html; the request fails
// and the application works regardless, but it fails slowly, once per run, on
// every run. The corpus is served from this process, so the reference is
// rewritten to a local no-op on the way out.
var deadDependency = regexp.MustCompile(`https?://polyfill\.io/[^"'\s>]*`)
// staticServer serves the whole corpus tree on one implementation's port. The
// tree rather than the implementation's own directory, because an example that
// asks for a path above itself has to resolve the same way it would in the
// upstream repository. Only one implementation is ever visited on this port, so
// the origin still belongs to it alone.
type staticServer struct {
implementation implementation
listener net.Listener
server *http.Server
}
func startStaticServer(
corpusRoot string,
target implementation,
) (*staticServer, error) {
listener, err := net.Listen("tcp", fmt.Sprintf("127.0.0.1:%d", target.Port))
if err != nil {
return nil, fmt.Errorf("serve %s: %w", target.Name, err)
}
server := &http.Server{Handler: corpusHandler(corpusRoot)}
running := &staticServer{
implementation: target,
listener: listener,
server: server,
}
go server.Serve(listener)
return running, nil
}
func (s *staticServer) stop() {
s.server.Close()
}
// waitReady confirms the served document answers before any run drives it. A
// wrong document path would otherwise reach the driver as a 404 page, and a
// sweep of 404 pages produces clean runs for every implementation.
//
// Only a transport error is retried. The listener is bound before the sweep
// starts, so a status that is not 200 is the server's answer about this
// document and waiting will not change it.
func (s *staticServer) waitReady(ctx context.Context) error {
url := s.implementation.URL()
client := &http.Client{Timeout: 5 * time.Second}
deadline := time.Now().Add(readinessTimeout)
var lastErr error
for {
if ctx.Err() != nil {
return ctx.Err()
}
response, err := client.Get(url)
if err == nil {
io.Copy(io.Discard, response.Body)
response.Body.Close()
if response.StatusCode == http.StatusOK {
return nil
}
return fmt.Errorf(
"%s answered %d, not the document",
url,
response.StatusCode,
)
}
lastErr = err
if time.Now().After(deadline) {
return fmt.Errorf(
"%s did not answer within %s: %w",
url,
readinessTimeout,
lastErr,
)
}
time.Sleep(readinessInterval)
}
}
func corpusHandler(corpusRoot string) http.Handler {
files := http.FileServer(http.Dir(corpusRoot))
return http.HandlerFunc(
func(writer http.ResponseWriter, request *http.Request) {
cleaned := path.Clean(
"/" + strings.TrimPrefix(request.URL.Path, "/"),
)
if cleaned == stubScriptPath {
writer.Header().Set("Content-Type", "application/javascript")
io.WriteString(
writer,
"/* corpus-sweep: dependency removed upstream */\n",
)
return
}
if !strings.HasSuffix(cleaned, ".html") {
files.ServeHTTP(writer, request)
return
}
body, modified, err := readDocument(corpusRoot, cleaned)
if err != nil {
// Never the file server's fallback: it answers a document it
// cannot read with a 200 directory listing, and a run driven at a
// listing explores nothing and comes back clean.
http.Error(writer, err.Error(), documentStatus(err))
return
}
rewritten := deadDependency.ReplaceAll(body, []byte(stubScriptPath))
http.ServeContent(
writer,
request,
path.Base(cleaned),
modified,
bytes.NewReader(rewritten),
)
},
)
}
func documentStatus(err error) int {
switch {
case errors.Is(err, fs.ErrNotExist):
return http.StatusNotFound
case errors.Is(err, fs.ErrPermission):
return http.StatusForbidden
default:
return http.StatusInternalServerError
}
}
func readDocument(corpusRoot, cleaned string) ([]byte, time.Time, error) {
file, err := http.Dir(corpusRoot).Open(cleaned)
if err != nil {
return nil, time.Time{}, err
}
defer file.Close()
info, err := file.Stat()
if err != nil || info.IsDir() {
return nil, time.Time{}, fmt.Errorf(
"%s is not a document",
filepath.FromSlash(cleaned),
)
}
body, err := io.ReadAll(file)
if err != nil {
return nil, time.Time{}, err
}
return body, info.ModTime(), nil
}
@@ -0,0 +1,214 @@
package main
import (
"context"
"io"
"net/http"
"os"
"path/filepath"
"strings"
"testing"
"time"
)
func get(t *testing.T, url string) (int, string) {
t.Helper()
client := &http.Client{Timeout: 5 * time.Second}
response, err := client.Get(url)
if err != nil {
t.Fatalf("GET %s: %v", url, err)
}
defer response.Body.Close()
body, err := io.ReadAll(response.Body)
if err != nil {
t.Fatal(err)
}
return response.StatusCode, string(body)
}
func TestStaticServer_ServesEachImplementationsDocumentFromItsOwnPort(
t *testing.T,
) {
root := wholeCorpus(t)
planned, err := planImplementations(
root,
[]string{"angular-dart", "duel", "react"},
freePortRange(t, 3),
)
if err != nil {
t.Fatal(err)
}
for _, target := range planned {
server, err := startStaticServer(root, target)
if err != nil {
t.Fatal(err)
}
defer server.stop()
if err := server.waitReady(context.Background()); err != nil {
t.Fatalf("%s: %v", target.Name, err)
}
status, body := get(t, target.URL())
if status != http.StatusOK ||
!strings.Contains(body, "<title>"+target.Name+"</title>") {
t.Errorf(
"%s at %s: got %d %q",
target.Name,
target.URL(),
status,
body,
)
}
}
}
func TestStaticServer_ReplacesTheDependencyThatNoLongerAnswers(t *testing.T) {
root := wholeCorpus(t)
document := filepath.Join(root, filepath.FromSlash(documentFor("aurelia")))
body := `<!doctype html><script src="https://polyfill.io/v3/polyfill.min.js?features=Promise"></script>`
if err := os.WriteFile(document, []byte(body), 0o644); err != nil {
t.Fatal(err)
}
planned, err := planImplementations(
root,
[]string{"aurelia"},
freePortRange(t, 1),
)
if err != nil {
t.Fatal(err)
}
server, err := startStaticServer(root, planned[0])
if err != nil {
t.Fatal(err)
}
defer server.stop()
status, served := get(t, planned[0].URL())
if status != http.StatusOK {
t.Fatalf("document answered %d", status)
}
if strings.Contains(served, "polyfill.io") {
t.Errorf(
"a request to a host that no longer answers reaches the browser once per run: %q",
served,
)
}
if !strings.Contains(served, `src="`+stubScriptPath+`"`) {
t.Errorf(
"the reference was removed rather than pointed at a local no-op: %q",
served,
)
}
stubStatus, stub := get(t, planned[0].Origin()+stubScriptPath)
if stubStatus != http.StatusOK || stub == "" {
t.Errorf("stub script: got %d %q", stubStatus, stub)
}
}
// An example that asks for a path above its own directory has to resolve the
// same way it does upstream, which is why the whole tree is served rather than
// the one directory.
func TestStaticServer_ServesPathsAboveTheImplementationsOwnDirectory(
t *testing.T,
) {
root := wholeCorpus(t)
shared := filepath.Join(root, "node_modules", "todomvc-common")
if err := os.MkdirAll(shared, 0o755); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(shared, "base.css"), []byte("body{}"), 0o644); err != nil {
t.Fatal(err)
}
planned, err := planImplementations(
root,
[]string{"react"},
freePortRange(t, 1),
)
if err != nil {
t.Fatal(err)
}
server, err := startStaticServer(root, planned[0])
if err != nil {
t.Fatal(err)
}
defer server.stop()
status, body := get(
t,
planned[0].Origin()+"/node_modules/todomvc-common/base.css",
)
if status != http.StatusOK || body != "body{}" {
t.Errorf("shared asset: got %d %q", status, body)
}
}
func TestStaticServerWaitReady_ReportsADocumentThatDoesNotAnswer(t *testing.T) {
root := wholeCorpus(t)
document := filepath.Join(root, filepath.FromSlash(documentFor("dojo")))
if err := os.Chmod(document, 0o000); err != nil {
t.Fatal(err)
}
t.Cleanup(func() { os.Chmod(document, 0o644) })
planned, err := planImplementations(
root,
[]string{"dojo"},
freePortRange(t, 1),
)
if err != nil {
t.Fatal(err)
}
server, err := startStaticServer(root, planned[0])
if err != nil {
t.Fatal(err)
}
defer server.stop()
started := time.Now()
err = server.waitReady(context.Background())
if err == nil {
t.Fatal(
"a document that does not answer would be swept as an error page and come back clean",
)
}
if elapsed := time.Since(started); elapsed > readinessTimeout {
t.Errorf(
"waited %s for a status that was never going to change",
elapsed,
)
}
}
func TestStaticServerStop_ReleasesThePort(t *testing.T) {
root := wholeCorpus(t)
planned, err := planImplementations(
root,
[]string{"vue"},
freePortRange(t, 1),
)
if err != nil {
t.Fatal(err)
}
server, err := startStaticServer(root, planned[0])
if err != nil {
t.Fatal(err)
}
if err := server.waitReady(context.Background()); err != nil {
t.Fatal(err)
}
server.stop()
client := &http.Client{Timeout: time.Second}
deadline := time.Now().Add(3 * time.Second)
for time.Now().Before(deadline) {
response, err := client.Get(planned[0].URL())
if err != nil {
return
}
response.Body.Close()
time.Sleep(50 * time.Millisecond)
}
t.Fatalf(
"port %d is still served after stop(), so the next sweep cannot bind it",
planned[0].Port,
)
}
+307
View File
@@ -0,0 +1,307 @@
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"os/exec"
"path/filepath"
"strconv"
"sync"
"time"
)
// resolveBinaries turns campaign and sanderling into absolute paths before
// anything is served. Each campaign runs from the sweep's own directory, and a
// binary that is missing altogether has to stop the sweep here rather than fail
// once per implementation and seed.
func resolveBinaries(configuration *config) error {
for name, value := range map[string]*string{
"--campaign": &configuration.campaignPath,
"--sanderling": &configuration.sanderlingPath,
} {
resolved, err := exec.LookPath(*value)
if err != nil {
return fmt.Errorf("%s: %w", name, err)
}
absolute, err := filepath.Abs(resolved)
if err != nil {
return fmt.Errorf("%s: %w", name, err)
}
*value = absolute
}
return nil
}
type sweep struct {
configuration config
stdout io.Writer
records io.Writer
servers map[string]*staticServer
mutex sync.Mutex
stalled int
failedRuns int
totalRuns int
}
func runSweep(
ctx context.Context,
configuration config,
stdout io.Writer,
) error {
if _, err := os.Stat(filepath.Join(configuration.outputDirectory, manifestFileName)); err == nil {
return fmt.Errorf(
"%s already exists in %s: pick a fresh --output so two sweeps do not share a directory",
manifestFileName,
configuration.outputDirectory,
)
}
if err := verifyCorpus(configuration.corpusRoot); err != nil {
return err
}
implementations, err := planImplementations(
configuration.corpusRoot,
configuration.implementations,
configuration.basePort,
)
if err != nil {
return err
}
if err := resolveBinaries(&configuration); err != nil {
return err
}
if _, err := os.Stat(configuration.specPath); err != nil {
return fmt.Errorf("--spec: %w", err)
}
if err := os.MkdirAll(configuration.outputDirectory, 0o755); err != nil {
return fmt.Errorf("create sweep dir: %w", err)
}
host, _ := os.Hostname()
if err := writeManifest(configuration.outputDirectory, buildManifest(configuration, implementations, host, time.Now().UTC())); err != nil {
return fmt.Errorf("write %s: %w", manifestFileName, err)
}
recordsFile, err := os.OpenFile(
filepath.Join(configuration.outputDirectory, recordsFileName),
os.O_CREATE|os.O_WRONLY|os.O_APPEND,
0o644,
)
if err != nil {
return fmt.Errorf("open %s: %w", recordsFileName, err)
}
defer recordsFile.Close()
running := &sweep{
configuration: configuration,
stdout: stdout,
records: recordsFile,
}
// Every port is bound before any run starts. A port already taken means
// that implementation cannot have an origin of its own, and continuing
// without one is what the separate origins are there to prevent.
if err := running.serveAll(implementations); err != nil {
return err
}
defer running.stopAll()
fmt.Fprintf(
stdout,
"sweep: %d implementations, %d seeds each, %d at a time, %s\n",
len(
implementations,
),
len(configuration.seeds),
configuration.concurrency,
configuration.outputDirectory,
)
running.work(ctx, implementations)
fmt.Fprintf(
stdout,
"sweep complete: %d of %d implementations never ran, %d of %d campaigns failed\n",
running.stalled,
len(implementations),
running.failedRuns,
running.totalRuns,
)
if running.stalled > 0 || running.failedRuns > 0 {
return fmt.Errorf(
"%d of %d implementations never ran and %d of %d campaigns failed; see %s",
running.stalled,
len(implementations),
running.failedRuns,
running.totalRuns,
recordsFileName,
)
}
return nil
}
func (s *sweep) serveAll(implementations []implementation) error {
s.servers = make(map[string]*staticServer, len(implementations))
for _, target := range implementations {
server, err := startStaticServer(s.configuration.corpusRoot, target)
if err != nil {
s.stopAll()
return err
}
s.servers[target.Name] = server
}
return nil
}
func (s *sweep) stopAll() {
for _, server := range s.servers {
server.stop()
}
}
func (s *sweep) work(ctx context.Context, implementations []implementation) {
queue := make(chan implementation, len(implementations))
for _, target := range implementations {
queue <- target
}
close(queue)
workers := min(s.configuration.concurrency, len(implementations))
var waitGroup sync.WaitGroup
for range workers {
waitGroup.Add(1)
go func() {
defer waitGroup.Done()
for target := range queue {
if ctx.Err() != nil {
return
}
s.report(s.runImplementation(ctx, target))
}
}()
}
waitGroup.Wait()
}
// runImplementation carries one implementation through all of its seeds. A
// document that does not answer is recorded and the sweep moves on: one
// implementation must not cost the other forty-two their runs.
func (s *sweep) runImplementation(
ctx context.Context,
target implementation,
) (record implementationRecord) {
record = implementationRecord{
Name: target.Name,
Document: target.Document,
Port: target.Port,
Origin: target.Origin(),
URL: target.URL(),
StartedAt: time.Now().UTC(),
}
started := time.Now()
defer func() { record.MonotonicMillis = time.Since(started).Milliseconds() }()
if err := os.MkdirAll(filepath.Join(s.configuration.outputDirectory, target.Name), 0o755); err != nil {
record.FailedStage = stageServe
record.Error = err.Error()
return record
}
if err := s.servers[target.Name].waitReady(ctx); err != nil {
record.FailedStage = stageServe
record.Error = err.Error()
return record
}
for _, seed := range s.configuration.seeds {
if ctx.Err() != nil {
return record
}
record.Runs = append(record.Runs, s.runSeed(ctx, target, seed))
}
return record
}
func (s *sweep) runSeed(
ctx context.Context,
target implementation,
seed int64,
) (record runRecord) {
seedText := strconv.FormatInt(seed, 10)
directory := campaignDirectory(s.configuration, target, seedText)
record = runRecord{Seed: seed, CampaignDirectory: directory}
started := time.Now()
defer func() { record.MonotonicMillis = time.Since(started).Milliseconds() }()
if err := os.MkdirAll(directory, 0o755); err != nil {
record.ExitCode = -1
record.LaunchError = err.Error()
return record
}
exitCode, err := runCommand(
ctx,
s.configuration.campaignPath,
campaignArguments(
s.configuration,
target,
seedText,
),
filepath.Join(directory, "campaign.log"),
)
record.ExitCode = exitCode
if err != nil {
record.LaunchError = err.Error()
}
return record
}
func runCommand(
ctx context.Context,
binary string,
arguments []string,
logPath string,
) (int, error) {
logFile, err := os.Create(logPath)
if err != nil {
return -1, err
}
defer logFile.Close()
command := exec.CommandContext(ctx, binary, arguments...)
command.Stdout = logFile
command.Stderr = logFile
err = command.Run()
if err == nil {
return 0, nil
}
var exitError *exec.ExitError
if errors.As(err, &exitError) {
return exitError.ExitCode(), nil
}
return -1, err
}
func (s *sweep) report(record implementationRecord) {
s.mutex.Lock()
defer s.mutex.Unlock()
if record.FailedStage != "" {
s.stalled++
}
s.totalRuns += len(record.Runs)
failed := 0
for _, run := range record.Runs {
if run.ExitCode != 0 {
failed++
}
}
s.failedRuns += failed
if err := json.NewEncoder(s.records).Encode(record); err != nil {
fmt.Fprintf(s.stdout, "warning: %s record: %v\n", record.Name, err)
}
elapsed := time.Duration(record.MonotonicMillis) * time.Millisecond
if record.FailedStage != "" {
fmt.Fprintf(s.stdout, "%s port=%d failed at %s: %s (%s)\n",
record.Name, record.Port, record.FailedStage, record.Error, elapsed)
return
}
fmt.Fprintf(s.stdout, "%s origin=%s campaigns=%d failed=%d elapsed=%s\n",
record.Name, record.Origin, len(record.Runs), failed, elapsed)
}