package main import ( "bytes" "context" "crypto/rand" "crypto/sha256" "encoding/hex" "encoding/json" "errors" "flag" "fmt" "io" "math/big" "os" "os/exec" "path/filepath" "regexp" "strconv" "strings" "time" "github.com/pelletier/go-toml/v2" "github.com/rvbox/rvbox/internal/domain" ) const ( manifestVersion = 1 repositoryID = "rvbox" maxSuiteLogSize = 1 << 20 ) var runIDPattern = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,63}$`) type manifest struct { Version uint32 `toml:"version" json:"version"` RepositoryID string `toml:"repository_id" json:"repository_id"` RunID string `toml:"run_id" json:"run_id"` Layer string `toml:"layer" json:"layer"` Suite string `toml:"suite" json:"suite"` Seed int64 `toml:"seed" json:"seed"` GitCommit string `toml:"git_commit" json:"git_commit"` DirtyDiffSHA256 string `toml:"dirty_diff_sha256" json:"dirty_diff_sha256"` ToolchainImage string `toml:"toolchain_image" json:"toolchain_image"` ComposeProject string `toml:"compose_project" json:"compose_project"` Phase string `toml:"phase" json:"phase"` CreatedAt time.Time `toml:"created_at" json:"created_at"` UpdatedAt time.Time `toml:"updated_at" json:"updated_at"` OwnedResources []string `toml:"owned_resources" json:"owned_resources"` } type journalEntry struct { At time.Time `json:"at"` Step string `json:"step"` Status string `json:"status"` Detail string `json:"detail,omitempty"` Data map[string]any `json:"data,omitempty"` } type harness struct { root string now func() time.Time out io.Writer } func runCLI(ctx context.Context, args []string) error { repoRoot, err := os.Getwd() if err != nil { return err } h := &harness{root: filepath.Join(repoRoot, ".test-runs"), now: func() time.Time { return time.Now().UTC() }, out: os.Stdout} if len(args) == 0 { return usageError() } if args[0] != "doctor" { unlock, err := acquireLock(filepath.Join(repoRoot, ".test-harness.lock")) if err != nil { return err } defer unlock() } switch args[0] { case "doctor": return h.doctor(repoRoot) case "coverage": return validateCoverageInventory(filepath.Join(repoRoot, "test", "coverage.toml")) case "integration": return h.integration(ctx, args[1:]) case "e2e": return h.e2e(ctx, args[1:]) case "status", "logs", "collect", "recover", "reuse", "stop", "reset", "purge": return h.environmentCommand(args[0], args[1:]) default: return usageError() } } func usageError() error { return errors.New("usage: harness doctor|coverage|integration|e2e|status|logs|collect|recover|reuse|stop|reset|purge") } func (h *harness) doctor(repoRoot string) error { for _, file := range []string{"go.mod", "deploy/compose.yaml", "docs/implementation-plan.v1.md"} { if _, err := os.Stat(filepath.Join(repoRoot, file)); err != nil { return fmt.Errorf("repository check %s: %w", file, err) } } if err := os.MkdirAll(h.root, 0o700); err != nil { return err } probe, err := os.CreateTemp(h.root, ".doctor-") if err != nil { return fmt.Errorf("test run root is not writable: %w", err) } name := probe.Name() _ = probe.Close() _ = os.Remove(name) fmt.Fprintln(h.out, "RVBox test harness is ready; native Windows scenarios use the explicit scripts/windows/test-host.ps1 host lane.") return nil } func (h *harness) integration(ctx context.Context, args []string) error { flags := flag.NewFlagSet("integration", flag.ContinueOnError) flags.SetOutput(io.Discard) suite := flags.String("suite", "sample", "suite name") runID := flags.String("run-id", "", "run ID") resume := flags.Bool("resume", false, "resume an existing run") if err := flags.Parse(args); err != nil { return err } if *suite != "sample" && *suite != "store" && *suite != "server-session" { return fmt.Errorf("suite %q is not implemented yet; available: sample, store, server-session", *suite) } var current *manifest var err error if *resume { if *runID == "" { return errors.New("--resume requires --run-id") } current, err = h.load(*runID) if err != nil { return err } if current.Layer != "integration" || current.Suite != *suite { return errors.New("run layer/suite does not match resume request") } if current.Phase != "ready" && current.Phase != "interrupted" && current.Phase != "stopped" && current.Phase != "running" && current.Phase != "failed" { return fmt.Errorf("run in phase %q is not resumable; recover or reuse it first", current.Phase) } } else { current, err = h.create(*runID, "integration", *suite) if err != nil { return err } } fmt.Fprintln(h.out, current.RunID) if err := h.transition(current, "running", *suite+"-start", *suite+" integration run started"); err != nil { return err } select { case <-ctx.Done(): _ = h.transition(current, "interrupted", "interrupt", ctx.Err().Error()) return ctx.Err() default: } if *suite == "sample" { if err := h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: "invariant", Status: "passed", Detail: "manifest ownership and journal durability verified"}); err != nil { return err } } else { var suiteErr error switch *suite { case "store": suiteErr = h.runStoreSuite(ctx, current) case "server-session": suiteErr = h.runServerSessionSuite(ctx, current) } if suiteErr != nil { _ = h.transition(current, "failed", *suite+"-failed", suiteErr.Error()) return suiteErr } } return h.transition(current, "completed", *suite+"-complete", *suite+" integration run completed") } // e2e runs production-shaped Go scenarios against real local listeners and // durable stores. Native Windows work is an additional host lane invoked by // scripts/windows/test-host.ps1; it is never silently replaced by Wine or a // cross-compiled binary. The run manifest/journal makes every scenario // resumable and keeps artifacts bounded. func (h *harness) e2e(ctx context.Context, args []string) error { flags := flag.NewFlagSet("e2e", flag.ContinueOnError) flags.SetOutput(io.Discard) scenario := flags.String("scenario", "smoke", "scenario name: smoke, script, recovery, or all") runID := flags.String("run-id", "", "run ID") resume := flags.Bool("resume", false, "resume an existing run") if err := flags.Parse(args); err != nil { return err } if *scenario != "smoke" && *scenario != "script" && *scenario != "recovery" && *scenario != "all" { return fmt.Errorf("scenario %q is not implemented yet; available: smoke, script, recovery, all", *scenario) } var current *manifest var err error if *resume { if *runID == "" { return errors.New("--resume requires --run-id") } current, err = h.load(*runID) if err != nil { return err } if current.Layer != "e2e" || current.Suite != *scenario { return errors.New("run layer/scenario does not match resume request") } if current.Phase != "ready" && current.Phase != "interrupted" && current.Phase != "stopped" && current.Phase != "running" && current.Phase != "failed" { return fmt.Errorf("run in phase %q is not resumable; recover or reuse it first", current.Phase) } } else { current, err = h.create(*runID, "e2e", *scenario) if err != nil { return err } } fmt.Fprintln(h.out, current.RunID) if err := h.transition(current, "running", "e2e-start", "e2e scenario started"); err != nil { return err } scenarios := []string{*scenario} if *scenario == "all" { scenarios = []string{"smoke", "script", "recovery"} } for _, item := range scenarios { if err := h.runE2EScenario(ctx, current, item); err != nil { _ = h.transition(current, "failed", "e2e-"+item+"-failed", err.Error()) return err } } return h.transition(current, "completed", "e2e-complete", "e2e scenario completed") } func (h *harness) runE2EScenario(ctx context.Context, current *manifest, scenario string) error { if err := h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: "e2e-" + scenario, Status: "running", Detail: "scenario started"}); err != nil { return err } var err error switch scenario { case "smoke": err = h.runGoSuite(ctx, current, "e2e-client-agent", "real client/server WebSocket and control flow", "./test/integration/clientagent") case "script": err = h.runGoSuite(ctx, current, "e2e-script-transfer", "durable script transfer and replay", "./internal/client/agent") case "recovery": err = h.runStoreSuite(ctx, current) default: err = fmt.Errorf("unknown e2e scenario %q", scenario) } if err != nil { return err } return h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: "e2e-" + scenario, Status: "passed", Detail: "scenario passed"}) } func (h *harness) runStoreSuite(ctx context.Context, current *manifest) error { if err := h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: "store-real-sqlite", Status: "running", Detail: "running real SQLite/WAL and filesystem cases"}); err != nil { return err } artifactDir := filepath.Join(h.runDir(current.RunID), "artifacts") if err := os.MkdirAll(artifactDir, 0o700); err != nil { return err } capture := &limitedCapture{limit: maxSuiteLogSize} command := exec.CommandContext(ctx, "go", "test", "-count=1", "-tags=integration", "-shuffle="+strconv.FormatInt(current.Seed, 10), "-timeout=2m", "./test/integration/store") command.Stdout = capture command.Stderr = capture err := command.Run() logPath := filepath.Join(artifactDir, "suite.log") if writeErr := atomicWrite(logPath, capture.Bytes(), 0o600); writeErr != nil { return writeErr } if err != nil { return fmt.Errorf("store suite failed (bounded log %s): %w", logPath, err) } return h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: "store-real-sqlite", Status: "passed", Detail: "real SQLite/WAL and filesystem cases passed"}) } func (h *harness) runServerSessionSuite(ctx context.Context, current *manifest) error { return h.runGoSuite(ctx, current, "server-session-real-websocket", "running real HTTP/WebSocket, protobuf, SQLite, and fencing cases", "./internal/server/session") } func (h *harness) runGoSuite(ctx context.Context, current *manifest, step, detail, packagePath string) error { if err := h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: step, Status: "running", Detail: detail}); err != nil { return err } artifactDir := filepath.Join(h.runDir(current.RunID), "artifacts") if err := os.MkdirAll(artifactDir, 0o700); err != nil { return err } capture := &limitedCapture{limit: maxSuiteLogSize} command := exec.CommandContext(ctx, "go", "test", "-race", "-count=1", "-shuffle="+strconv.FormatInt(current.Seed, 10), "-timeout=2m", packagePath) command.Stdout = capture command.Stderr = capture err := command.Run() logPath := filepath.Join(artifactDir, "suite.log") if writeErr := atomicWrite(logPath, capture.Bytes(), 0o600); writeErr != nil { return writeErr } if err != nil { return fmt.Errorf("%s suite failed (bounded log %s): %w", current.Suite, logPath, err) } return h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: step, Status: "passed", Detail: detail + " passed"}) } func (h *harness) environmentCommand(command string, args []string) error { flags := flag.NewFlagSet(command, flag.ContinueOnError) flags.SetOutput(io.Discard) runID := flags.String("run-id", "", "run ID") if err := flags.Parse(args); err != nil { return err } if *runID == "" { return fmt.Errorf("%s requires --run-id", command) } current, err := h.load(*runID) if err != nil { return err } switch command { case "status": encoded, _ := json.MarshalIndent(current, "", " ") fmt.Fprintln(h.out, string(encoded)) return nil case "logs": data, err := os.ReadFile(h.journalPath(current.RunID)) if err != nil { return err } _, err = h.out.Write(data) return err case "collect": return h.collect(current) case "recover": if current.Phase != "interrupted" && current.Phase != "stopped" && current.Phase != "running" && current.Phase != "failed" { return fmt.Errorf("run in phase %q does not need recovery", current.Phase) } return h.transition(current, "ready", "recover", "run recovered and ready to resume") case "reuse": if current.Phase != "reset" && current.Phase != "completed" && current.Phase != "ready" { return fmt.Errorf("run in phase %q cannot be reused", current.Phase) } return h.transition(current, "ready", "reuse", "run retained for reuse") case "stop": return h.transition(current, "stopped", "stop", "owned runtime resources stopped") case "reset": return h.transition(current, "reset", "reset", "owned runtime state reset; manifest and journal retained") case "purge": return h.purge(current) default: return usageError() } } type limitedCapture struct { data []byte limit int truncated bool } func (capture *limitedCapture) Write(data []byte) (int, error) { written := len(data) remaining := capture.limit - len(capture.data) if remaining > 0 { if len(data) > remaining { data = data[:remaining] } capture.data = append(capture.data, data...) } if written > remaining { capture.truncated = true } return written, nil } func (capture *limitedCapture) Bytes() []byte { if !capture.truncated { return capture.data } return append(append([]byte(nil), capture.data...), []byte("\n[output truncated by RVBox test harness]\n")...) } func (h *harness) create(requestedID, layer, suite string) (*manifest, error) { if requestedID == "" { id, err := domain.NewUUIDv7() if err != nil { return nil, err } requestedID = id.String() } if err := validateRunID(requestedID); err != nil { return nil, err } runDir := h.runDir(requestedID) if err := os.MkdirAll(h.root, 0o700); err != nil { return nil, err } if err := os.Mkdir(runDir, 0o700); err != nil { return nil, fmt.Errorf("create run %s: %w", requestedID, err) } seedValue, err := rand.Int(rand.Reader, big.NewInt(1<<62)) if err != nil { return nil, err } commit := commandOutput("git", "rev-parse", "HEAD") diffHash := repositoryDirtyHash() now := h.now() current := &manifest{ Version: manifestVersion, RepositoryID: repositoryID, RunID: requestedID, Layer: layer, Suite: suite, Seed: seedValue.Int64(), GitCommit: strings.TrimSpace(commit), DirtyDiffSHA256: hex.EncodeToString(diffHash[:]), ToolchainImage: os.Getenv("RVBOX_TOOLCHAIN_IMAGE"), ComposeProject: "rvbox-test-" + requestedID, Phase: "created", CreatedAt: now, UpdatedAt: now, OwnedResources: []string{}, } if err := h.writeManifest(current); err != nil { _ = os.Remove(runDir) return nil, err } if err := h.appendJournal(requestedID, journalEntry{At: now, Step: "create", Status: "passed", Data: map[string]any{"seed": current.Seed}}); err != nil { return nil, err } return current, nil } func (h *harness) transition(current *manifest, phase, step, detail string) error { current.Phase = phase current.UpdatedAt = h.now() if err := h.writeManifest(current); err != nil { return err } return h.appendJournal(current.RunID, journalEntry{At: current.UpdatedAt, Step: step, Status: phase, Detail: detail}) } func (h *harness) load(runID string) (*manifest, error) { if err := validateRunID(runID); err != nil { return nil, err } data, err := os.ReadFile(h.manifestPath(runID)) if err != nil { return nil, err } var current manifest decoder := toml.NewDecoder(strings.NewReader(string(data))) decoder.DisallowUnknownFields() if err := decoder.Decode(¤t); err != nil { return nil, fmt.Errorf("decode run manifest: %w", err) } if current.Version != manifestVersion || current.RepositoryID != repositoryID || current.RunID != runID { return nil, errors.New("run manifest identity/version mismatch") } return ¤t, nil } func (h *harness) writeManifest(current *manifest) error { data, err := toml.Marshal(current) if err != nil { return err } return atomicWrite(h.manifestPath(current.RunID), data, 0o600) } func (h *harness) appendJournal(runID string, entry journalEntry) error { data, err := json.Marshal(entry) if err != nil { return err } file, err := os.OpenFile(h.journalPath(runID), os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o600) if err != nil { return err } defer file.Close() if _, err := file.Write(append(data, '\n')); err != nil { return err } return file.Sync() } func (h *harness) collect(current *manifest) error { reportDir := filepath.Join(h.runDir(current.RunID), "artifacts") if err := os.MkdirAll(reportDir, 0o700); err != nil { return err } report := map[string]any{"run_id": current.RunID, "layer": current.Layer, "suite": current.Suite, "phase": current.Phase, "collected_at": h.now(), "payloads_included": false} data, _ := json.MarshalIndent(report, "", " ") path := filepath.Join(reportDir, "report.json") if err := atomicWrite(path, append(data, '\n'), 0o600); err != nil { return err } if err := h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: "collect", Status: "passed", Detail: "bounded redacted report written"}); err != nil { return err } fmt.Fprintln(h.out, path) return nil } func (h *harness) purge(current *manifest) error { runDir := h.runDir(current.RunID) info, err := os.Lstat(runDir) if err != nil { return err } if info.Mode()&os.ModeSymlink != 0 || !info.IsDir() { return errors.New("refusing to purge a symlink or non-directory run path") } if len(current.OwnedResources) != 0 { return errors.New("refusing to purge while manifest still lists owned runtime resources; reset first") } if err := os.RemoveAll(runDir); err != nil { return err } fmt.Fprintf(h.out, "purged %s (not recoverable)\n", runDir) return nil } func (h *harness) runDir(runID string) string { return filepath.Join(h.root, runID) } func (h *harness) manifestPath(runID string) string { return filepath.Join(h.runDir(runID), "run.toml") } func (h *harness) journalPath(runID string) string { return filepath.Join(h.runDir(runID), "journal.jsonl") } func validateRunID(runID string) error { if !runIDPattern.MatchString(runID) || runID == "." || runID == ".." { return fmt.Errorf("invalid filesystem-safe run ID %q", runID) } return nil } func atomicWrite(path string, data []byte, mode os.FileMode) error { directory := filepath.Dir(path) temporary, err := os.CreateTemp(directory, ".write-") if err != nil { return err } temporaryPath := temporary.Name() defer os.Remove(temporaryPath) if err := temporary.Chmod(mode); err != nil { _ = temporary.Close() return err } if _, err := temporary.Write(data); err != nil { _ = temporary.Close() return err } if err := temporary.Sync(); err != nil { _ = temporary.Close() return err } if err := temporary.Close(); err != nil { return err } if err := os.Rename(temporaryPath, path); err != nil { return err } dir, err := os.Open(directory) if err != nil { return err } defer dir.Close() return dir.Sync() } func commandOutput(name string, args ...string) string { output, err := exec.Command(name, args...).Output() if err != nil { return "unavailable" } return string(output) } func repositoryDirtyHash() [sha256.Size]byte { hash := sha256.New() tracked := commandBytes("git", "diff", "--binary", "HEAD", "--", ".") _, _ = hash.Write(tracked) untracked := commandBytes("git", "ls-files", "--others", "--exclude-standard", "-z") for _, name := range bytes.Split(untracked, []byte{0}) { if len(name) == 0 { continue } _, _ = hash.Write([]byte{0}) _, _ = hash.Write(name) if data, err := os.ReadFile(string(name)); err == nil { _, _ = hash.Write([]byte{0}) _, _ = hash.Write(data) } } var result [sha256.Size]byte copy(result[:], hash.Sum(nil)) return result } func commandBytes(name string, args ...string) []byte { output, err := exec.Command(name, args...).Output() if err != nil { return []byte("unavailable") } return output } type coverageInventory struct { Version uint32 `toml:"version"` Requirements []coverageRequirement `toml:"requirements"` } type coverageRequirement struct { ID string `toml:"id"` Layer string `toml:"layer"` Status string `toml:"status"` Tests []string `toml:"tests"` } var coverageIDPattern = regexp.MustCompile(`^(HP|BH|ERR|RACE|CRASH|SEC|BOUND|REC)-[A-Z0-9]+-[0-9]{2}$`) func validateCoverageInventory(path string) error { absolutePath, err := filepath.Abs(path) if err != nil { return err } path = absolutePath data, err := os.ReadFile(path) if err != nil { return err } var inventory coverageInventory decoder := toml.NewDecoder(bytes.NewReader(data)) decoder.DisallowUnknownFields() if err := decoder.Decode(&inventory); err != nil { return err } if inventory.Version != 1 || len(inventory.Requirements) == 0 { return errors.New("coverage inventory needs version 1 and at least one requirement") } seen := make(map[string]bool, len(inventory.Requirements)) repositoryRoot := filepath.Dir(filepath.Dir(path)) for _, requirement := range inventory.Requirements { if !coverageIDPattern.MatchString(requirement.ID) || seen[requirement.ID] { return fmt.Errorf("invalid or duplicate coverage ID %q", requirement.ID) } seen[requirement.ID] = true if requirement.Layer != "unit" && requirement.Layer != "integration" && requirement.Layer != "e2e" { return fmt.Errorf("coverage %s has invalid layer %q", requirement.ID, requirement.Layer) } if requirement.Status != "implemented" && requirement.Status != "planned" && requirement.Status != "blocked_native_windows" { return fmt.Errorf("coverage %s has invalid status %q", requirement.ID, requirement.Status) } if requirement.Status == "implemented" && len(requirement.Tests) == 0 { return fmt.Errorf("implemented coverage %s has no tests", requirement.ID) } for _, reference := range requirement.Tests { parts := strings.Split(reference, ":") if len(parts) != 2 || parts[0] == "" || parts[1] == "" || filepath.IsAbs(parts[0]) || strings.Contains(parts[0], "..") { return fmt.Errorf("coverage %s has invalid test reference %q", requirement.ID, reference) } source, err := os.ReadFile(filepath.Join(repositoryRoot, filepath.FromSlash(parts[0]))) if err != nil { return fmt.Errorf("coverage %s test reference: %w", requirement.ID, err) } if !bytes.Contains(source, []byte("func "+parts[1]+"(")) { return fmt.Errorf("coverage %s test function %q not found", requirement.ID, parts[1]) } } } return nil }