From fe39bddac7907eee82e0dc2802321761573a9c78 Mon Sep 17 00:00:00 2001 From: Victor Garcia Date: Mon, 28 Sep 2026 00:28:28 -0600 Subject: [PATCH] test: prove lease revival after sleep and trace drain on continuation release (U6) R12: a repair round whose lease expired server-side while it ran is revived by freshen's heartbeat when the local clock says the last renewal is older than HeartbeatInterval, and is then authorized, pushed, and recorded. The same round with a fresh clock is refused on lease_not_owner, and a superseded lease (a successor attempt exists) is never revived. R13: through the real workerAttemptRunner wiring, trace ingestion is refused until CI goes green on the round's fix, so the whole trace is still buffered when the continuation is released. Once ClaimOnce returns, the control plane holds exactly the attempt's JSONL raw record. Co-Authored-By: Claude Opus 5.5 --- cmd/jig/ci_repair_test.go | 227 +++++++++++++++++++++---- internal/worker/publish_repair_test.go | 114 +++++++++++++ 2 files changed, 312 insertions(+), 29 deletions(-) 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) + } +}