diff --git a/cmd/jig/ci_repair_test.go b/cmd/jig/ci_repair_test.go index 0f6256a..2dbe998 100644 --- a/cmd/jig/ci_repair_test.go +++ b/cmd/jig/ci_repair_test.go @@ -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" @@ -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 { @@ -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. @@ -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) @@ -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) +} diff --git a/internal/worker/publish_repair_test.go b/internal/worker/publish_repair_test.go index fe439cc..289c732 100644 --- a/internal/worker/publish_repair_test.go +++ b/internal/worker/publish_repair_test.go @@ -10,6 +10,7 @@ import ( "path/filepath" "strings" "testing" + "time" "github.com/StructuPath/jig/internal/protocol" ) @@ -38,6 +39,7 @@ type roundScript struct { type scriptedContinuation struct { t *testing.T worktree string + lease *attemptLease rounds []roundScript failures []CIFailure released int @@ -73,6 +75,7 @@ func (c *scriptedContinuation) Release() { c.released++ } func repairRunner(t *testing.T, continuation *scriptedContinuation) AttemptRunner { return RunnerFunc(func(_ context.Context, attempt *PreparedAttempt) Outcome { continuation.worktree = attempt.WorktreePath + continuation.lease = attempt.lease if err := os.WriteFile(filepath.Join(attempt.WorktreePath, "work.txt"), []byte("work\n"), 0o644); err != nil { t.Errorf("write work: %v", err) } @@ -383,3 +386,114 @@ func TestARoundCommittingUndeclaredPathsIsNotPushed(t *testing.T) { t.Fatal("a round with an undeclared committed path was pushed or recorded") } } + +// ---- R12: a round's fence after a gap (machine sleep) ---------------------- + +// sleptRound is what a machine sleep does to a round: while the round runs, +// the lease expires server-side (unswept), and the lease's clock is set so +// its last renewal is either more than one HeartbeatInterval old (stale: +// freshen must heartbeat) or pinned to it (fresh: freshen skips the +// heartbeat). supersede also sweeps the attempt and retries the job, so a +// successor attempt exists. The fence the round then meets is the pre-push +// authorization. The publishing runner's outcome is captured because a +// lease that stays lost also refuses the attempt's completion. +func sleptRound(t *testing.T, stale, supersede bool) (*repairScenario, Outcome, *protocol.Attempt, error) { + t.Helper() + s := newRepairScenario(t, 2, 1) + s.continuation.rounds = []roundScript{{files: map[string]string{"fix.txt": "fixed\n"}, before: func() { + lease := s.continuation.lease + expireLease(t, s.h, lease.attemptID) + if supersede { + now := time.Now().UnixMilli() + if _, err := s.h.db.Exec(`UPDATE attempts SET state = 'lost', completed_at = ? WHERE id = ?`, + now, lease.attemptID); err != nil { + t.Errorf("sweep attempt: %v", err) + } + if _, err := s.h.db.Exec(`UPDATE jobs SET state = 'failed', updated_at = ? WHERE id = ?`, + now, s.job.ID); err != nil { + t.Errorf("fail job: %v", err) + } + if _, err := s.h.store.RetryJob(context.Background(), s.job.ID); err != nil { + t.Errorf("retry job: %v", err) + } + } + lease.mutex.Lock() + if stale { + lease.now = func() time.Time { return time.Now().Add(protocol.HeartbeatInterval) } + } else { + pinned := lease.lastRenewal + lease.now = func() time.Time { return pinned } + } + lease.mutex.Unlock() + }}} + var outcome Outcome + publisher := s.w.config.Runner + s.w.config.Runner = RunnerFunc(func(ctx context.Context, attempt *PreparedAttempt) Outcome { + outcome = publisher.Run(ctx, attempt) + return outcome + }) + attempt, err := s.w.ClaimOnce(context.Background()) + return s, outcome, attempt, err +} + +// After a gap longer than HeartbeatInterval, the round's pre-push fence +// heartbeats first, which revives the expired-but-unswept lease, and the +// round is authorized, pushed, recorded, and judged green (R12). +func TestAStaleClockRevivesAnExpiredLeaseBeforeTheRoundIsAuthorized(t *testing.T) { + s, _, attempt, err := sleptRound(t, true, false) + if err != nil || attempt == nil { + t.Fatalf("claim: attempt=%v err=%v, want the revived lease to carry the attempt through", attempt, err) + } + summary := publishSummaryOf(t, attempt.Result) + if attempt.State != protocol.AttemptAccepted || !summary.Published() { + t.Fatalf("attempt = %s summary = %+v, want the round authorized and the attempt accepted", attempt.State, summary) + } + rounds := s.rounds(t, attempt.ID) + if len(rounds) != 1 || rounds[0].HeadAfter != remoteBranches(t, s.originDir)[s.branch] { + t.Fatalf("rounds = %+v, want the one round recorded at the pushed fix", rounds) + } +} + +// The control: the same expired lease, but a clock that says the last +// renewal is recent. freshen skips the heartbeat, so the authorization meets +// the expired lease and the round stops on lease_not_owner without pushing. +// This is what proves the heartbeat, not the authorization, did the reviving. +func TestAFreshClockSkipsTheHeartbeatAndTheExpiredLeaseIsRefused(t *testing.T) { + s, outcome, attempt, err := sleptRound(t, false, false) + if err == nil { + t.Fatalf("attempt = %+v, want the completion refused on the still-expired lease", attempt) + } + summary := publishSummaryOf(t, outcome.Result) + if summary.Code != "lease_not_owner" || len(summary.CIRepairs) != 1 || + summary.CIRepairs[0].Outcome != "lease_not_owner" { + t.Fatalf("summary = %+v, want the round refused on lease_not_owner", summary) + } + if head := remoteBranches(t, s.originDir)[s.branch]; head != s.ci.judged[0] { + t.Fatalf("remote head = %s, want the red head %s: an unauthorized round pushed", head, s.ci.judged[0]) + } +} + +// A heartbeat never revives a lease whose job has a successor attempt: the +// stale clock heartbeats, the control plane answers lease_superseded, and +// the round neither pushes nor records. +func TestAStaleClockDoesNotReviveASupersededLease(t *testing.T) { + s, outcome, attempt, err := sleptRound(t, true, true) + if err == nil { + t.Fatalf("attempt = %+v, a superseded attempt completed", attempt) + } + summary := publishSummaryOf(t, outcome.Result) + if summary.Code != "lease_superseded" || len(summary.CIRepairs) != 1 || + summary.CIRepairs[0].Outcome != "lease_superseded" { + t.Fatalf("summary = %+v, want the round stopped on lease_superseded", summary) + } + if head := remoteBranches(t, s.originDir)[s.branch]; head != s.ci.judged[0] { + t.Fatalf("remote head = %s, want the red head %s: a superseded round pushed", head, s.ci.judged[0]) + } + var rounds int + if err := s.h.db.QueryRow(`SELECT COUNT(*) FROM publish_ci_repairs`).Scan(&rounds); err != nil { + t.Fatal(err) + } + if rounds != 0 { + t.Fatalf("%d round(s) recorded, want none from a superseded lease", rounds) + } +}