feat(implementation-sweep): run one campaign against every implementation of a requirement

installs, builds and serves each implementation on its own port, then hands the campaign tool the same seed slice, step budget and generator for all of them, so a difference between implementations is not a difference in exploration. the generator and platform are fixed rather than exposed.
This commit is contained in:
pj committed 2026-08-16 17:45:41 +05:30
1 parent a45ba76d8e
commit 146152a3af
8 files changed
+1830

No files matched your search

@@ -0,0 +1,482 @@
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 stub implementation's server and as the
// fetcher the stub campaign uses, so the sweep drives a real preview process
// over a real port and the served URL is answered by a real HTTP server.
func TestMain(m *testing.M) {
switch {
case os.Getenv("SWEEP_TEST_SERVE_PORT") != "":
serveUntilKilled(os.Getenv("SWEEP_TEST_SERVE_PORT"))
case os.Getenv("SWEEP_TEST_FETCH_URL") != "":
recordFetch(
os.Getenv("SWEEP_TEST_FETCH_URL"),
os.Getenv("SWEEP_TEST_FETCH_LOG"),
)
default:
os.Exit(m.Run())
}
}
func serveUntilKilled(port string) {
handler := http.HandlerFunc(
func(writer http.ResponseWriter, request *http.Request) {
fmt.Fprintf(writer, "%s %s", port, request.URL.RequestURI())
},
)
http.ListenAndServe("localhost:"+port, handler)
}
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()
}
// stubBun answers install, fails to build impl-02, and serves the preview from
// the test binary on the port it was given.
const stubBun = `#!/bin/sh
echo "$PWD $*" >> "%[1]s"
if [ "$1" = "run" ] && [ "$2" = "build" ]; then
case "$PWD" in *impl-02) echo "TS2322: type error" >&2; exit 1 ;; esac
exit 0
fi
if [ "$1" = "run" ] && [ "$2" = "preview" ]; then
port=""
previous=""
for argument in "$@"; do
if [ "$previous" = "--port" ]; then port="$argument"; fi
previous="$argument"
done
SWEEP_TEST_SERVE_PORT="$port" exec "%[2]s"
fi
exit 0
`
// 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 impl-03 seed 4.
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"
SWEEP_TEST_FETCH_URL="$url" SWEEP_TEST_FETCH_LOG="%[3]s" "%[4]s"
mkdir -p "$output"
printf '{"arm":"stub","seeds":[%%s]}\n' "$seed" > "$output/campaign.json"
echo "stub campaign seed=$seed url=$url"
case "$output" in
*impl-03/seed-4) exit 1 ;;
esac
exit 0
`
func TestRun_EndToEndAgainstStubBunAndCampaign(t *testing.T) {
root := t.TempDir()
implementations := filepath.Join(root, "implementations")
for _, name := range []string{"impl-01", "impl-02", "impl-03"} {
if err := os.MkdirAll(filepath.Join(implementations, name), 0o755); err != nil {
t.Fatal(err)
}
}
specPath := filepath.Join(root, "relay.ts")
if err := os.WriteFile(specPath, []byte("export const properties = [];\n"), 0o644); err != nil {
t.Fatal(err)
}
output := filepath.Join(root, "campaigns")
bunLog := filepath.Join(root, "bun.log")
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)
}
bunPath := writeScript(
t,
filepath.Join(root, "stub-bun"),
fmt.Sprintf(stubBun, bunLog, testBinary),
)
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",
)
basePort := freePortRange(t, 3)
var stdout bytes.Buffer
err = run([]string{
"--implementations", implementations,
"--spec", specPath,
"--seeds", "4-5",
"--max-steps", "40",
"--duration", "30s",
"--concurrency", "2",
"--base-port", fmt.Sprint(basePort),
"--output", output,
"--bun", bunPath,
"--campaign", campaignPath,
"--sanderling", sanderlingPath,
}, &stdout, io.Discard)
if err == nil ||
!strings.Contains(err.Error(), "1 of 3 implementations never ran") {
t.Fatalf(
"expected the failed build and the failed campaign to be reported, got %v",
err,
)
}
var recorded manifest
manifestBody, err := os.ReadFile(filepath.Join(output, manifestFileName))
if err != nil {
t.Fatal(err)
}
if err := json.Unmarshal(manifestBody, &recorded); err != nil {
t.Fatal(err)
}
if !slices.Equal(recorded.Seeds, []int64{4, 5}) {
t.Errorf("intended seeds: got %v", recorded.Seeds)
}
if len(recorded.Implementations) != 3 {
t.Fatalf("intended implementations: got %v", recorded.Implementations)
}
for index, planned := range recorded.Implementations {
wantPort := basePort + index
if planned.Port != wantPort {
t.Errorf(
"%s port: got %d, want %d",
planned.Name,
planned.Port,
wantPort,
)
}
wantURL := fmt.Sprintf("http://localhost:%d/?seed={seed}", wantPort)
if planned.URLTemplate != wantURL {
t.Errorf(
"%s url template: got %q, want %q",
planned.Name,
planned.URLTemplate,
wantURL,
)
}
}
if recorded.Generator != "seeded" || recorded.MaxSteps != 40 {
t.Errorf(
"manifest generator/budget: got %q/%d",
recorded.Generator,
recorded.MaxSteps,
)
}
// impl-02 fails its build, so the two implementations either side of it
// still have to reach the campaign tool with both seeds.
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")
bundle := argumentValue(arguments, "--bundle-id")
port := basePort + slices.Index(
[]string{"impl-01", "impl-02", "impl-03"},
arm,
)
wantBundle := fmt.Sprintf("http://localhost:%d/?seed=%s", port, seed)
if bundle != wantBundle {
t.Errorf(
"%s seed %s: bundle id %q, want %q",
arm,
seed,
bundle,
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{"impl-01/4", "impl-01/5", "impl-03/4", "impl-03/5"} {
if !seen[want] {
t.Errorf("%s never reached the campaign tool", want)
}
}
// What the served page actually saw: the right port for the
// implementation, carrying the same seed the campaign was given.
fetched := readLines(t, fetchLog)
for _, want := range []string{
fmt.Sprintf("http://localhost:%d/?seed=4 -> 200 %d /?seed=4", basePort, basePort),
fmt.Sprintf("http://localhost:%d/?seed=5 -> 200 %d /?seed=5", basePort, basePort),
fmt.Sprintf("http://localhost:%d/?seed=4 -> 200 %d /?seed=4", basePort+2, basePort+2),
fmt.Sprintf("http://localhost:%d/?seed=5 -> 200 %d /?seed=5", basePort+2, basePort+2),
} {
if !slices.Contains(fetched, want) {
t.Errorf(
"the served page never saw %q:\n%s",
want,
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
}
failed := byName["impl-02"]
if failed.FailedStage != stageBuild || len(failed.Runs) != 0 {
t.Errorf(
"impl-02: got stage %q with %d runs, want a build failure and no runs",
failed.FailedStage,
len(failed.Runs),
)
}
if !strings.Contains(failed.Error, "build.log") {
t.Errorf("impl-02 error should point at its log: %q", failed.Error)
}
buildLog, err := os.ReadFile(
filepath.Join(output, "impl-02", stageBuild+".log"),
)
if err != nil {
t.Fatal(err)
}
if !strings.Contains(string(buildLog), "TS2322") {
t.Errorf("impl-02 build log lost the compiler error: %q", buildLog)
}
for _, name := range []string{"impl-01", "impl-03"} {
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,
)
}
for _, run := range record.Runs {
if run.MonotonicMillis <= 0 {
t.Errorf(
"%s seed %d took %d ms, so nothing timed the campaign",
name,
run.Seed,
run.MonotonicMillis,
)
}
}
}
if exit := byName["impl-03"].Runs[0].ExitCode; exit != 1 {
t.Errorf("impl-03 seed 4 exit code: got %d, want 1", exit)
}
if exit := byName["impl-03"].Runs[1].ExitCode; exit != 0 {
t.Errorf(
"impl-03 seed 5 ran after seed 4 failed and should have exited 0, got %d",
exit,
)
}
for _, name := range []string{"impl-01", "impl-03"} {
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,
)
}
}
}
installed := readLines(t, bunLog)
for _, name := range []string{"impl-01", "impl-02", "impl-03"} {
if !slices.Contains(
installed,
filepath.Join(implementations, name)+" install",
) {
t.Errorf(
"%s was never installed:\n%s",
name,
strings.Join(installed, "\n"),
)
}
}
// Every preview server the sweep started is gone with it: a leaked one
// holds its port, and the next sweep would be served by the old build.
client := &http.Client{Timeout: 2 * time.Second}
for _, port := range []int{basePort, basePort + 2} {
if response, err := client.Get(readinessURL(port)); err == nil {
response.Body.Close()
t.Errorf("port %d is still served after the sweep finished", port)
}
}
if !strings.Contains(stdout.String(), "failed at build") {
t.Errorf(
"progress output does not name the build failure: %q",
stdout.String(),
)
}
}
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 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
}
// 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("localhost:%d", base+offset),
)
if err != nil {
return false
}
listener.Close()
}
return true
}
@@ -0,0 +1,284 @@
// Command implementation-sweep runs one identical campaign against every model
// implementation of a single requirement. It installs, builds and serves each
// implementation on its own port, then hands the campaign tool the same seed
// slice, the same step budget and the same generator for all of them, so a
// difference between implementations is not a difference in exploration.
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 a served web build, 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 built, served and 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 a
// knee at twelve to sixteen, on a contended laptop it says to re-measure before
// trusting anything above eight. Six sits at the better efficiency, costs about
// 7 GB of a 64 GB host, and leaves the Android emulator farm that shares that
// host its four to six slots. Each worker here also carries a vite server the
// fleet measurement did not include.
const defaultConcurrency = 6
// defaultBasePort is the first port handed out. Vite's own defaults are 5173
// and 4173, so a sweep starting here does not collide with a dev server someone
// left running.
const defaultBasePort = 5300
type config struct {
implementationsDirectory string
specPath string
outputDirectory string
seeds []int64
maxSteps int
duration time.Duration
concurrency int
basePort int
bunPath string
campaignPath string
sanderlingPath string
extraArguments []string
}
const usage = `implementation-sweep runs one seeded campaign per model implementation.
Usage:
implementation-sweep --implementations <dir> --spec <path> --seeds <spec>
--max-steps <n> --output <dir> [flags]
[-- <sanderling test flags>]
Each impl-* directory under --implementations is installed, built and served on
its own port, and every one is swept with the same seeds and the same step
budget. An implementation that fails to install, build or serve is recorded and
the sweep moves on to the next one.
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("implementation-sweep", flag.ContinueOnError)
flagSet.SetOutput(stderr)
flagSet.Usage = func() {
fmt.Fprint(stderr, usage)
flagSet.PrintDefaults()
}
var configuration config
var seedSpecification string
flagSet.StringVar(
&configuration.implementationsDirectory,
"implementations",
"",
"directory holding impl-01 to impl-NN (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.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 built, served and swept at once",
)
flagSet.IntVar(
&configuration.basePort,
"base-port",
defaultBasePort,
"first port served; each implementation takes the next one in name order",
)
flagSet.StringVar(
&configuration.outputDirectory,
"output",
"",
"campaign tree to create (required)",
)
flagSet.StringVar(
&configuration.bunPath,
"bun",
"bun",
"bun binary that installs, builds and serves each implementation",
)
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{
"--implementations": configuration.implementationsDirectory,
"--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
for name, value := range map[string]*string{
"--implementations": &configuration.implementationsDirectory,
"--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
}
// servedURL is the one place a seed becomes a URL. The same seed is also handed
// to the campaign as --seeds, which reaches sanderling as --seed and fixes the
// exploration, while the scaffold reads ?seed= and fixes the latency and the
// outcome of every send. A violation replays only when both carry the same
// number, so both come from the seed argument here and never from two flags.
func servedURL(port int, seed string) string {
return fmt.Sprintf("http://localhost:%d/?seed=%s", port, seed)
}
func readinessURL(port int) string {
return fmt.Sprintf("http://localhost:%d/", port)
}
func campaignDirectory(
configuration config,
target implementation,
seed string,
) string {
return filepath.Join(
configuration.outputDirectory,
target.Name,
"seed-"+seed,
)
}
// campaignArguments builds one campaign invocation. The seed is a string
// because it lands in two arguments, --seeds and the ?seed= of --bundle-id,
// and passing it once keeps them from drifting apart.
func campaignArguments(
configuration config,
target implementation,
seed string,
) []string {
arguments := []string{
"--spec", configuration.specPath,
"--bundle-id", servedURL(target.Port, seed),
"--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)
}
}
@@ -0,0 +1,214 @@
package main
import (
"io"
"net/url"
"path/filepath"
"slices"
"strconv"
"strings"
"testing"
"time"
)
func baseArguments() []string {
return []string{
"--implementations", "/e4/implementations",
"--spec", "/e4/relay.ts",
"--seeds", "1-3",
"--max-steps", "400",
"--output", "/campaigns/e4",
}
}
func TestParseArguments_DefaultsAndSeeds(t *testing.T) {
configuration, err := parseArguments(baseArguments(), io.Discard)
if err != nil {
t.Fatal(err)
}
if !slices.Equal(configuration.seeds, []int64{1, 2, 3}) {
t.Errorf("seeds: got %v", configuration.seeds)
}
if configuration.concurrency != defaultConcurrency {
t.Errorf(
"concurrency default: got %d, want %d",
configuration.concurrency,
defaultConcurrency,
)
}
if configuration.basePort != defaultBasePort {
t.Errorf(
"base port default: got %d, want %d",
configuration.basePort,
defaultBasePort,
)
}
if configuration.duration != 5*time.Minute {
t.Errorf("duration default: got %s", configuration.duration)
}
for name, got := range map[string]string{
"bun": configuration.bunPath,
"campaign": configuration.campaignPath,
"sanderling": configuration.sanderlingPath,
} {
if got != name {
t.Errorf("%s path default: got %q", name, got)
}
}
}
func TestParseArguments_Rejections(t *testing.T) {
cases := []struct {
name string
arguments []string
want string
}{
{
"missing implementations",
[]string{
"--spec",
"s",
"--seeds",
"1",
"--max-steps",
"10",
"--output",
"o",
},
"--implementations is required",
},
{
"missing spec",
[]string{
"--implementations",
"i",
"--seeds",
"1",
"--max-steps",
"10",
"--output",
"o",
},
"--spec is required",
},
{
"missing output",
[]string{
"--implementations",
"i",
"--spec",
"s",
"--seeds",
"1",
"--max-steps",
"10",
},
"--output is required",
},
{
"zero max steps",
append(baseArguments(), "--max-steps", "0"),
"--max-steps must be positive",
},
{
"zero concurrency",
append(baseArguments(), "--concurrency", "0"),
"--concurrency must be positive",
},
{
"privileged base port",
append(baseArguments(), "--base-port", "80"),
"outside 1024-65535",
},
{
"seed zero",
append(baseArguments(), "--seeds", "0,1"),
"not reproducible",
},
}
for _, testCase := range cases {
_, err := parseArguments(testCase.arguments, io.Discard)
if err == nil {
t.Errorf("%s: expected error", testCase.name)
continue
}
if !strings.Contains(err.Error(), testCase.want) {
t.Errorf(
"%s: got %q, want it to contain %q",
testCase.name,
err,
testCase.want,
)
}
}
}
// The seed reaches two independent things, the campaign's own seed and the
// scaffold's failure stream, and a replay reproduces neither unless they carry
// the same number.
func TestCampaignArguments_OneSeedReachesBothTheCampaignAndTheURL(
t *testing.T,
) {
configuration, err := parseArguments(
append(baseArguments(), "--", "--clear-data=false"),
io.Discard,
)
if err != nil {
t.Fatal(err)
}
target := implementation{
Name: "impl-07",
Directory: "/e4/implementations/impl-07",
Port: 5306,
}
for _, seed := range []string{"1", "42"} {
arguments := campaignArguments(configuration, target, seed)
if got := argumentValue(arguments, "--seeds"); got != seed {
t.Errorf("--seeds: got %q, want %q", got, seed)
}
bundle := argumentValue(arguments, "--bundle-id")
parsed, err := url.Parse(bundle)
if err != nil {
t.Fatalf("--bundle-id %q: %v", bundle, err)
}
if got := parsed.Query().Get("seed"); got != seed {
t.Errorf(
"served URL seed: got %q, want %q (from %q)",
got,
seed,
bundle,
)
}
if parsed.Host != "localhost:"+strconv.Itoa(target.Port) {
t.Errorf(
"served host: got %q, want the implementation's own port %d",
parsed.Host,
target.Port,
)
}
for flagName, want := range map[string]string{
"--arm": "impl-07",
"--platform": "web",
"--generator": "seeded",
"--max-steps": "400",
"--spec": "/e4/relay.ts",
"--output": filepath.Join("/campaigns/e4", "impl-07", "seed-"+seed),
} {
if got := argumentValue(arguments, flagName); got != want {
t.Errorf("%s: got %q, want %q", flagName, got, want)
}
}
if arguments[len(arguments)-2] != "--" ||
arguments[len(arguments)-1] != "--clear-data=false" {
t.Errorf("passthrough flags lost: %v", arguments)
}
}
}
func argumentValue(arguments []string, name string) string {
index := slices.Index(arguments, name)
if index < 0 || index+1 >= len(arguments) {
return ""
}
return arguments[index+1]
}
@@ -0,0 +1,117 @@
package main
import (
"encoding/json"
"fmt"
"os"
"path/filepath"
"time"
)
const (
manifestFileName = "sweep.json"
recordsFileName = "implementations.jsonl"
seedPlaceholder = "{seed}"
)
// plannedImplementation is one implementation the sweep intends to run, with
// the port it is served on and the URL every seed is served at.
type plannedImplementation struct {
Name string `json:"name"`
Directory string `json:"directory"`
Port int `json:"port"`
URLTemplate string `json:"url_template"`
}
// manifest is sweep.json: what the sweep INTENDED to run, written before the
// first install 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"`
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"`
BunPath string `json:"bun_path"`
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,
Directory: target.Directory,
Port: target.Port,
URLTemplate: servedURL(target.Port, seedPlaceholder),
})
}
return manifest{
Generator: generator,
Platform: platform,
SpecPath: configuration.specPath,
MaxSteps: configuration.maxSteps,
DurationMillis: configuration.duration.Milliseconds(),
Seeds: configuration.seeds,
Implementations: planned,
Concurrency: configuration.concurrency,
Host: host,
BunPath: configuration.bunPath,
CampaignPath: configuration.campaignPath,
SanderlingPath: configuration.sanderlingPath,
StartedAt: startedAt,
}
}
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"`
URL string `json:"url"`
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 an implementation that never
// got past install, build or serve carries no runs at all.
type implementationRecord struct {
Name string `json:"implementation"`
Directory string `json:"directory"`
Port int `json:"port"`
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 (
stageInstall = "install"
stageBuild = "build"
stageServe = "serve"
)
@@ -0,0 +1,109 @@
package main
import (
"context"
"fmt"
"io"
"net/http"
"os"
"os/exec"
"syscall"
"time"
)
const (
// serverStartTimeout covers a cold vite start on a host already running
// five other implementations.
serverStartTimeout = 90 * time.Second
serverPollInterval = 250 * time.Millisecond
serverShutdownGrace = 10 * time.Second
)
// server is one implementation's preview server. It runs in its own process
// group so that stopping it takes the whole vite tree with it: a leaked server
// holds its port, and the next sweep against that implementation would be
// served by the previous build.
type server struct {
command *exec.Cmd
logFile *os.File
exited chan struct{}
}
func startServer(
ctx context.Context,
configuration config,
target implementation,
logPath string,
) (*server, error) {
logFile, err := os.Create(logPath)
if err != nil {
return nil, err
}
command := exec.CommandContext(ctx, configuration.bunPath,
"run", "preview", "--port", fmt.Sprint(target.Port), "--strictPort")
command.Dir = target.Directory
command.Stdout = logFile
command.Stderr = logFile
command.SysProcAttr = &syscall.SysProcAttr{Setpgid: true}
if err := command.Start(); err != nil {
logFile.Close()
return nil, err
}
running := &server{
command: command,
logFile: logFile,
exited: make(chan struct{}),
}
go func() {
command.Wait()
close(running.exited)
}()
return running, nil
}
// waitReady polls the served page until it answers. A server that exits first
// is reported as such rather than waited on for the full timeout, because the
// usual cause is a port already taken and that answer is in the log.
func (s *server) waitReady(ctx context.Context, url string) error {
client := &http.Client{Timeout: 5 * time.Second}
deadline := time.Now().Add(serverStartTimeout)
for {
select {
case <-s.exited:
return fmt.Errorf("server exited before it answered %s", url)
case <-ctx.Done():
return ctx.Err()
default:
}
response, err := client.Get(url)
if err == nil {
io.Copy(io.Discard, response.Body)
response.Body.Close()
if response.StatusCode == http.StatusOK {
return nil
}
}
if time.Now().After(deadline) {
return fmt.Errorf(
"server did not answer %s within %s",
url,
serverStartTimeout,
)
}
time.Sleep(serverPollInterval)
}
}
func (s *server) stop() {
if s.command.Process != nil {
group := -s.command.Process.Pid
syscall.Kill(group, syscall.SIGTERM)
select {
case <-s.exited:
case <-time.After(serverShutdownGrace):
syscall.Kill(group, syscall.SIGKILL)
<-s.exited
}
}
s.logFile.Close()
}
@@ -0,0 +1,119 @@
package main
import (
"context"
"fmt"
"net/http"
"os"
"path/filepath"
"syscall"
"testing"
"time"
)
// bunSpawningAServer serves from a child process and then waits, which is the
// shape of `bun run preview`: the port belongs to something below the process
// the sweep started.
const bunSpawningAServer = `#!/bin/sh
port=""
previous=""
for argument in "$@"; do
if [ "$previous" = "--port" ]; then port="$argument"; fi
previous="$argument"
done
SWEEP_TEST_SERVE_PORT="$port" "%[1]s" &
wait
`
func TestServerStop_TakesTheProcessBelowItWithTheServer(t *testing.T) {
directory := t.TempDir()
testBinary, err := filepath.Abs(os.Args[0])
if err != nil {
t.Fatal(err)
}
bunPath := writeScript(
t,
filepath.Join(directory, "stub-bun"),
fmt.Sprintf(bunSpawningAServer, testBinary),
)
port := freePortRange(t, 1)
target := implementation{Name: "impl-01", Directory: directory, Port: port}
// Background rather than the test context: only stop() may end this
// server, or a leak would be hidden by the context being cancelled.
running, err := startServer(
context.Background(),
config{bunPath: bunPath},
target,
filepath.Join(directory, "serve.log"),
)
if err != nil {
t.Fatal(err)
}
t.Cleanup(
func() { syscall.Kill(-running.command.Process.Pid, syscall.SIGKILL) },
)
if err := running.waitReady(context.Background(), readinessURL(port)); err != nil {
t.Fatal(err)
}
running.stop()
client := &http.Client{Timeout: time.Second}
deadline := time.Now().Add(5 * time.Second)
for time.Now().Before(deadline) {
response, err := client.Get(readinessURL(port))
if err != nil {
return
}
response.Body.Close()
time.Sleep(100 * time.Millisecond)
}
t.Fatalf(
"port %d is still served after stop(): the server below bun outlived the sweep and holds the port",
port,
)
}
func TestServerWaitReady_ReportsAServerThatExited(t *testing.T) {
directory := t.TempDir()
bunPath := writeScript(
t,
filepath.Join(directory, "stub-bun"),
"#!/bin/sh\necho 'port is already in use' >&2\nexit 1\n",
)
port := freePortRange(t, 1)
running, err := startServer(
context.Background(),
config{bunPath: bunPath},
implementation{
Name: "impl-01",
Directory: directory,
Port: port,
},
filepath.Join(directory, "serve.log"),
)
if err != nil {
t.Fatal(err)
}
defer running.stop()
started := time.Now()
err = running.waitReady(context.Background(), readinessURL(port))
if err == nil {
t.Fatal(
"a server that exited should not be waited on until the start timeout",
)
}
if elapsed := time.Since(started); elapsed > 30*time.Second {
t.Errorf("waited %s for a server that had already exited", elapsed)
}
log, err := os.ReadFile(filepath.Join(directory, "serve.log"))
if err != nil {
t.Fatal(err)
}
if string(log) == "" {
t.Error("serve.log did not capture why the server exited")
}
}
@@ -0,0 +1,391 @@
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"os/exec"
"path/filepath"
"sort"
"strconv"
"strings"
"sync"
"time"
)
// implementation is one directory under --implementations, with the port it
// owns for the whole sweep. The port comes from the implementation's position
// in name order rather than from a pool, so the manifest can name the URL every
// run was served from before anything has been served.
type implementation struct {
Name string
Directory string
Port int
}
const implementationPrefix = "impl-"
func discoverImplementations(
directory string,
basePort int,
) ([]implementation, error) {
entries, err := os.ReadDir(directory)
if err != nil {
return nil, err
}
var found []implementation
for _, entry := range entries {
if !entry.IsDir() ||
!strings.HasPrefix(entry.Name(), implementationPrefix) {
continue
}
found = append(found, implementation{
Name: entry.Name(),
Directory: filepath.Join(directory, entry.Name()),
})
}
if len(found) == 0 {
return nil, fmt.Errorf(
"no %s* directories in %s",
implementationPrefix,
directory,
)
}
sort.Slice(
found,
func(i, j int) bool { return found[i].Name < found[j].Name },
)
if basePort+len(found)-1 > 65535 {
return nil, fmt.Errorf(
"--base-port %d leaves no room for %d implementations",
basePort,
len(found),
)
}
for index := range found {
found[index].Port = basePort + index
}
return found, nil
}
// resolveBinaries turns bun, campaign and sanderling into absolute paths before
// anything is installed. Each campaign runs from the sweep's own directory
// rather than the implementation's, so a relative --sanderling would otherwise
// resolve against the wrong one, 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{
"--bun": &configuration.bunPath,
"--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
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,
)
}
implementations, err := discoverImplementations(
configuration.implementationsDirectory,
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,
}
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) 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 from install to its last seed.
// Every failure it can meet is returned in the record: one implementation that
// cannot install, build or serve must not cost the other twenty-three their
// runs.
func (s *sweep) runImplementation(
ctx context.Context,
target implementation,
) (record implementationRecord) {
record = implementationRecord{
Name: target.Name,
Directory: target.Directory,
Port: target.Port,
StartedAt: time.Now().UTC(),
}
started := time.Now()
defer func() { record.MonotonicMillis = time.Since(started).Milliseconds() }()
directory := filepath.Join(s.configuration.outputDirectory, target.Name)
if err := os.MkdirAll(directory, 0o755); err != nil {
record.FailedStage = stageInstall
record.Error = err.Error()
return record
}
for _, step := range []struct {
stage string
arguments []string
}{
{stageInstall, []string{"install"}},
{stageBuild, []string{"run", "build"}},
} {
logPath := filepath.Join(directory, step.stage+".log")
exitCode, err := runCommand(
ctx,
target.Directory,
s.configuration.bunPath,
step.arguments,
logPath,
)
if err != nil {
record.FailedStage = step.stage
record.Error = err.Error()
return record
}
if exitCode != 0 {
record.FailedStage = step.stage
record.Error = fmt.Sprintf(
"bun %s exited %d, see %s",
strings.Join(step.arguments, " "),
exitCode,
logPath,
)
return record
}
}
running, err := startServer(
ctx,
s.configuration,
target,
filepath.Join(directory, "serve.log"),
)
if err != nil {
record.FailedStage = stageServe
record.Error = err.Error()
return record
}
defer running.stop()
if err := running.waitReady(ctx, readinessURL(target.Port)); 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,
URL: servedURL(target.Port, seedText),
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
}
// runCommand runs one step of the pipeline with its output in logPath. An
// empty directory keeps the sweep's own working directory, which is what the
// campaign tool gets: it has no reason to run inside an implementation.
func runCommand(
ctx context.Context,
directory, 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.Dir = directory
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)
for _, run := range record.Runs {
if run.ExitCode != 0 {
s.failedRuns++
}
}
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
}
failed := 0
for _, run := range record.Runs {
if run.ExitCode != 0 {
failed++
}
}
fmt.Fprintf(s.stdout, "%s port=%d campaigns=%d failed=%d elapsed=%s\n",
record.Name, record.Port, len(record.Runs), failed, elapsed)
}
@@ -0,0 +1,114 @@
package main
import (
"os"
"path/filepath"
"strings"
"testing"
)
func TestDiscoverImplementations_NameOrderAndOnePortEach(t *testing.T) {
directory := t.TempDir()
for _, name := range []string{"impl-03", "impl-01", "impl-10", "impl-02", "scaffold", ".DS_Store"} {
if err := os.MkdirAll(filepath.Join(directory, name), 0o755); err != nil {
t.Fatal(err)
}
}
if err := os.WriteFile(filepath.Join(directory, "impl-notes.md"), []byte("not a directory"), 0o644); err != nil {
t.Fatal(err)
}
found, err := discoverImplementations(directory, 5300)
if err != nil {
t.Fatal(err)
}
want := []implementation{
{Name: "impl-01", Port: 5300},
{Name: "impl-02", Port: 5301},
{Name: "impl-03", Port: 5302},
{Name: "impl-10", Port: 5303},
}
if len(found) != len(want) {
t.Fatalf(
"got %d implementations, want %d: %v",
len(found),
len(want),
found,
)
}
for index, target := range found {
if target.Name != want[index].Name || target.Port != want[index].Port {
t.Errorf(
"position %d: got %s on %d, want %s on %d",
index,
target.Name,
target.Port,
want[index].Name,
want[index].Port,
)
}
if target.Directory != filepath.Join(directory, want[index].Name) {
t.Errorf("%s directory: got %q", target.Name, target.Directory)
}
}
}
func TestDiscoverImplementations_EmptyDirectoryIsRefused(t *testing.T) {
_, err := discoverImplementations(t.TempDir(), 5300)
if err == nil || !strings.Contains(err.Error(), "no impl-* directories") {
t.Fatalf("got %v, want a refusal naming impl-*", err)
}
}
// A binary that is not there fails once, before anything is installed, rather
// than twenty-four times after the sweep has spent its build time.
func TestRunSweep_StopsBeforeItInstallsAnythingWhenABinaryIsMissing(
t *testing.T,
) {
implementations := t.TempDir()
if err := os.MkdirAll(filepath.Join(implementations, "impl-01"), 0o755); err != nil {
t.Fatal(err)
}
output := filepath.Join(t.TempDir(), "campaigns")
configuration := config{
implementationsDirectory: implementations,
outputDirectory: output,
basePort: 5300,
concurrency: 1,
bunPath: "bun",
campaignPath: "campaign-that-is-not-installed",
sanderlingPath: "sanderling",
}
err := runSweep(t.Context(), configuration, os.Stdout)
if err == nil || !strings.Contains(err.Error(), "--campaign") {
t.Fatalf("got %v, want the missing campaign binary named", err)
}
if _, err := os.Stat(output); !os.IsNotExist(err) {
t.Errorf(
"the sweep created %s before it checked it could run: %v",
output,
err,
)
}
}
func TestRunSweep_RefusesADirectoryThatAlreadyHoldsASweep(t *testing.T) {
implementations := t.TempDir()
if err := os.MkdirAll(filepath.Join(implementations, "impl-01"), 0o755); err != nil {
t.Fatal(err)
}
output := t.TempDir()
if err := os.WriteFile(filepath.Join(output, manifestFileName), []byte("{}"), 0o644); err != nil {
t.Fatal(err)
}
configuration := config{
implementationsDirectory: implementations,
outputDirectory: output,
basePort: 5300,
concurrency: 1,
}
if err := runSweep(t.Context(), configuration, os.Stdout); err == nil ||
!strings.Contains(err.Error(), "already exists") {
t.Fatalf("got %v, want a refusal to reuse the directory", err)
}
}