Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
227 changes: 198 additions & 29 deletions cmd/jig/ci_repair_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,20 @@ package main

import (
"context"
"encoding/json"
"fmt"
"io"
"log/slog"
"net/http"
"net/http/httptest"
"net/http/httputil"
"net/url"
"os"
"path/filepath"
"reflect"
"strings"
"sync"
"sync/atomic"
"testing"
"time"

Expand Down Expand Up @@ -75,18 +83,20 @@ func (g *e2eGateway) FailedCheckLogs(_ context.Context, _ string, checks []worke
return checks
}

func TestACIRepairRoundRunsThroughTheRealWorkerWiring(t *testing.T) {
repo := initRepo(t)
serverData, workerData := t.TempDir(), t.TempDir()
// scriptCIRepairRuntime scripts the chain's build and the repair round's.
func scriptCIRepairRuntime(t *testing.T) {
t.Helper()
scriptRuntime(t,
map[string]any{"files": map[string]any{"jig-note.txt": "note\n"},
"text": envelopeJSON(t, map[string]any{"status": "success", "summary": "wrote the note"})},
map[string]any{"files": map[string]any{"jig-fix.txt": "fixed\n"},
"text": envelopeJSON(t, map[string]any{"status": "success", "summary": "fixed lint"})},
)
base, stopServe := startServe(t, serverData, "--no-github-poll")
defer stopServe()
}

// submitCIRepairRun registers ciRepairDefinition and starts one run on repo.
func submitCIRepairRun(t *testing.T, base, repo string) {
t.Helper()
var definition protocol.Definition
if err := apiCall(context.Background(), base, "POST", "/api/definitions",
controlplane.DefinitionInput{Source: ciRepairDefinition}, &definition); err != nil {
Expand All @@ -99,6 +109,47 @@ func TestACIRepairRoundRunsThroughTheRealWorkerWiring(t *testing.T) {
}, &view); err != nil {
t.Fatal(err)
}
}

// startCIRepairHost builds the worker exactly as `jig worker` does —
// workerAttemptRunner inside the publishing runner — talking to serverURL,
// and registers it.
func startCIRepairHost(t *testing.T, serverURL, workerData string, gateway worker.PullRequestGateway) *worker.Worker {
t.Helper()
ctx := context.Background()
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
selected, err := selectRuntime(ctx, "claude-code", true)
if err != nil {
t.Fatal(err)
}
var host *worker.Worker
publisher := worker.NewPublishingRunner(
workerAttemptRunner(ctx, workerData, selected, logger, func() *worker.Worker { return host }),
gateway, worker.PublishOptions{
CommitAuthorName: "jig-test", CommitAuthorEmail: "jig@test",
CIPollInterval: 5 * time.Millisecond, CIRegistrationGrace: 20 * time.Millisecond,
})
host, err = worker.New(worker.Config{
ServerURL: serverURL, DataDir: workerData, Capacity: 1,
Probes: []worker.RuntimeProbe{selected.Runtime}, Runner: publisher, Logger: logger,
})
if err != nil {
t.Fatal(err)
}
publisher.Bind(host)
if _, err := host.RegisterOnce(ctx); err != nil {
t.Fatal(err)
}
return host
}

func TestACIRepairRoundRunsThroughTheRealWorkerWiring(t *testing.T) {
repo := initRepo(t)
serverData, workerData := t.TempDir(), t.TempDir()
scriptCIRepairRuntime(t)
base, stopServe := startServe(t, serverData, "--no-github-poll")
defer stopServe()
submitCIRepairRun(t, base, repo)

// The handoff note is planted while CI is red — after Execute returned,
// before the round — and must still be there when CI judges the fix.
Expand Down Expand Up @@ -127,30 +178,7 @@ func TestACIRepairRoundRunsThroughTheRealWorkerWiring(t *testing.T) {
}}

ctx := context.Background()
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
selected, err := selectRuntime(ctx, "claude-code", true)
if err != nil {
t.Fatal(err)
}
var host *worker.Worker
publisher := worker.NewPublishingRunner(
workerAttemptRunner(ctx, workerData, selected, logger, func() *worker.Worker { return host }),
gateway, worker.PublishOptions{
CommitAuthorName: "jig-test", CommitAuthorEmail: "jig@test",
CIPollInterval: 5 * time.Millisecond, CIRegistrationGrace: 20 * time.Millisecond,
})
host, err = worker.New(worker.Config{
ServerURL: base, DataDir: workerData, Capacity: 1,
Probes: []worker.RuntimeProbe{selected.Runtime}, Runner: publisher, Logger: logger,
})
if err != nil {
t.Fatal(err)
}
publisher.Bind(host)
if _, err := host.RegisterOnce(ctx); err != nil {
t.Fatal(err)
}

host := startCIRepairHost(t, base, workerData, gateway)
attempt, err := host.ClaimOnce(ctx)
if err != nil || attempt == nil {
t.Fatalf("claim: attempt=%v err=%v", attempt, err)
Expand Down Expand Up @@ -189,3 +217,144 @@ func TestACIRepairRoundRunsThroughTheRealWorkerWiring(t *testing.T) {
t.Fatalf("attempt scratch survived the attempt: %v", entries)
}
}

// ingestGate fronts the control plane for the worker: while it is closed,
// trace ingestion is refused as if the control plane were unreachable, and
// every other request passes through.
type ingestGate struct {
open atomic.Bool
proxy *httputil.ReverseProxy
}

func (g *ingestGate) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if !g.open.Load() && r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/events") {
http.Error(w, "trace ingestion is down", http.StatusServiceUnavailable)
return
}
g.proxy.ServeHTTP(w, r)
}

// attemptEvents pages through every event the control plane holds for an
// attempt.
func attemptEvents(t *testing.T, base, attemptID string) []protocol.Event {
t.Helper()
var events []protocol.Event
after := int64(0)
for {
var page protocol.EventPage
if err := apiCall(context.Background(), base, "GET",
fmt.Sprintf("/api/attempts/%s/events?after=%d&limit=500", attemptID, after), nil, &page); err != nil {
t.Fatal(err)
}
events = append(events, page.Events...)
if len(page.Events) == 0 || page.NextCursor <= after {
return events
}
after = page.NextCursor
}
}

// Releasing the repair continuation is what drains and closes the trace
// (R13). The stream nudges its sender on every emit, so with a healthy
// control plane the events land long before release whether or not the host
// closes the trace, and comparing the two records afterwards would prove
// nothing. So ingestion is refused from the start until CI goes green on the
// round's fix: every event, the chain's and the round's, is still buffered
// when the continuation is released, and the background sender is deep in
// its retry backoff. Close drains synchronously, and the host's release runs
// before the attempt completes, so when ClaimOnce returns the control plane
// holds exactly the JSONL raw record — with no waiting and no polling.
func TestReleasingTheRepairContinuationDrainsTheTraceBeforeTheAttemptCompletes(t *testing.T) {
repo := initRepo(t)
serverData, workerData := t.TempDir(), t.TempDir()
scriptCIRepairRuntime(t)
base, stopServe := startServe(t, serverData, "--no-github-poll")
defer stopServe()
submitCIRepairRun(t, base, repo)

target, err := url.Parse(base)
if err != nil {
t.Fatal(err)
}
gate := &ingestGate{proxy: httputil.NewSingleHostReverseProxy(target)}
proxy := httptest.NewServer(gate)
defer proxy.Close()

var probeErrors []string
gateway := &e2eGateway{probe: func(_ string, red bool) {
if red || gate.open.Load() {
return
}
// CI is green on the fix: the round is over, the continuation not
// yet released. Nothing may have reached the store yet, or the drain
// has nothing left to prove.
traces, _ := filepath.Glob(filepath.Join(workerData, "traces", "*.jsonl"))
if len(traces) != 1 {
probeErrors = append(probeErrors, fmt.Sprintf("found %d attempt traces, want 1", len(traces)))
} else if held := attemptEvents(t, base, strings.TrimSuffix(filepath.Base(traces[0]), ".jsonl")); len(held) != 0 {
probeErrors = append(probeErrors, fmt.Sprintf(
"%d event(s) reached the store while ingestion was refused", len(held)))
}
gate.open.Store(true)
}}

host := startCIRepairHost(t, proxy.URL, workerData, gateway)
attempt, err := host.ClaimOnce(context.Background())
if err != nil || attempt == nil {
t.Fatalf("claim: attempt=%v err=%v", attempt, err)
}
for _, problem := range probeErrors {
t.Error(problem)
}
if attempt.State != protocol.AttemptAccepted || !gate.open.Load() {
t.Fatalf("attempt = %s (%s), gate opened = %v: want accepted after a repair round",
attempt.State, attempt.Error, gate.open.Load())
}

raw, err := os.ReadFile(filepath.Join(workerData, "traces", attempt.ID+".jsonl"))
if err != nil {
t.Fatal(err)
}
var recorded []protocol.Event
for _, line := range strings.Split(strings.TrimSpace(string(raw)), "\n") {
var event protocol.Event
if err := json.Unmarshal([]byte(line), &event); err != nil {
t.Fatalf("parse JSONL line %q: %v", line, err)
}
recorded = append(recorded, event)
}
held := attemptEvents(t, base, attempt.ID)
if len(held) != len(recorded) {
t.Fatalf("the control plane holds %d event(s), the JSONL raw record %d: releasing the continuation did not drain the trace",
len(held), len(recorded))
}
repairStarted := false
for i := range recorded {
want, got := recorded[i], held[i]
if got.Seq != want.Seq || got.Type != want.Type || got.Phase != want.Phase || got.Name != want.Name ||
!sameJSON(t, got.Payload, want.Payload) {
t.Fatalf("event %d: control plane %+v, JSONL %+v", i, got, want)
}
repairStarted = repairStarted || want.Name == "ci_repair_start"
}
if !repairStarted {
t.Fatal("the JSONL raw record has no ci_repair_start: the round's events are not in it")
}
}

// sameJSON compares two JSON documents by value, so the store's re-encoding
// of a payload does not read as a difference.
func sameJSON(t *testing.T, a, b json.RawMessage) bool {
t.Helper()
if len(a) == 0 || len(b) == 0 {
return len(a) == len(b)
}
var left, right any
if err := json.Unmarshal(a, &left); err != nil {
t.Fatalf("decode %s: %v", a, err)
}
if err := json.Unmarshal(b, &right); err != nil {
t.Fatalf("decode %s: %v", b, err)
}
return reflect.DeepEqual(left, right)
}
Loading
Loading