fix(sync): keep the live progress counters correct and race-free
Two defects in the in-place progress line introduced with the reporting work: - started() only ran for repositories that actually synced, but finished() decrements for every repository that records a result. Any repository skipped by the eligibility or pre-sync checks therefore decremented a counter it had never incremented, so the running "active" count went negative and was displayed as such. Count a repository as active before the skip checks can exit. - clear() mutates the line width but ran under outputMu, while draw() mutates the same field under the progress mutex. With parallel jobs one worker could erase the line while another repainted it. Route clear() through the same lock, keeping outputMu around the surrounding writes. Verified by reverting the first fix and watching the new test report "active = -3". 🤖 Generated with Codebuff Co-Authored-By: Codebuff <noreply@codebuff.com>
This commit is contained in:
@@ -484,6 +484,10 @@ func syncSelectedWithProgress(root string, repos, workspaceRepos []repo, jobs in
|
|||||||
defer wg.Done()
|
defer wg.Done()
|
||||||
for next := range work {
|
for next := range work {
|
||||||
r := next.repo
|
r := next.repo
|
||||||
|
// Counted active before any skip path can record: finished()
|
||||||
|
// decrements again, so an early exit that skipped started()
|
||||||
|
// would drive the active count negative.
|
||||||
|
live.update(func() { live.started(); live.draw() })
|
||||||
if !r.Eligible {
|
if !r.Eligible {
|
||||||
record(next.index, syncResult{Path: r.RelativePath, Skipped: true, Message: r.BlockReason})
|
record(next.index, syncResult{Path: r.RelativePath, Skipped: true, Message: r.BlockReason})
|
||||||
continue
|
continue
|
||||||
@@ -503,18 +507,17 @@ func syncSelectedWithProgress(root string, repos, workspaceRepos []repo, jobs in
|
|||||||
}
|
}
|
||||||
if !quiet {
|
if !quiet {
|
||||||
outputMu.Lock()
|
outputMu.Lock()
|
||||||
live.clear()
|
live.update(func() { live.clear() })
|
||||||
fmt.Fprintf(stdout, "\nSTART %s (%s)\n", r.RelativePath, r.Branch)
|
fmt.Fprintf(stdout, "\nSTART %s (%s)\n", r.RelativePath, r.Branch)
|
||||||
outputMu.Unlock()
|
outputMu.Unlock()
|
||||||
}
|
}
|
||||||
live.update(func() { live.started(); live.draw() })
|
|
||||||
started := time.Now()
|
started := time.Now()
|
||||||
report, err := syncRepositoryWith(r.Path, timeout, len(fresh.Dirty) > 0, fresh.NestedRepoEntries, syncOptions{CreateMissing: createMissing})
|
report, err := syncRepositoryWith(r.Path, timeout, len(fresh.Dirty) > 0, fresh.NestedRepoEntries, syncOptions{CreateMissing: createMissing})
|
||||||
message := report.Message
|
message := report.Message
|
||||||
duration := time.Since(started).Round(time.Millisecond).String()
|
duration := time.Since(started).Round(time.Millisecond).String()
|
||||||
if !quiet {
|
if !quiet {
|
||||||
outputMu.Lock()
|
outputMu.Lock()
|
||||||
live.clear()
|
live.update(func() { live.clear() })
|
||||||
state := "DONE"
|
state := "DONE"
|
||||||
if err != nil {
|
if err != nil {
|
||||||
state = "FAIL"
|
state = "FAIL"
|
||||||
|
|||||||
@@ -268,3 +268,48 @@ func TestSyncOutputToBufferHasNoCarriageReturns(t *testing.T) {
|
|||||||
t.Fatalf("run did not print a summary:\n%q", stdout.String())
|
t.Fatalf("run did not print a summary:\n%q", stdout.String())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// A repository skipped before it ever syncs must still be counted active and
|
||||||
|
// then finished. Counting it as active only after the skip checks left the
|
||||||
|
// running "active" total negative, because finished() decrements it regardless.
|
||||||
|
func TestLiveProgressBalancesSkippedRepositories(t *testing.T) {
|
||||||
|
root := t.TempDir()
|
||||||
|
var repos []repo
|
||||||
|
for _, name := range []string{"blocked-a", "blocked-b", "blocked-c"} {
|
||||||
|
path := filepath.Join(root, name)
|
||||||
|
if err := os.MkdirAll(path, 0o755); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
// Not Eligible, so the worker takes the early-exit skip path.
|
||||||
|
repos = append(repos, repo{Path: path, RelativePath: name, BlockReason: "no remotes"})
|
||||||
|
}
|
||||||
|
|
||||||
|
var mu sync.Mutex
|
||||||
|
buf := &bytes.Buffer{}
|
||||||
|
live := &liveProgress{w: buf, mu: &mu, total: len(repos), counts: map[string]int{}}
|
||||||
|
results := syncSelectedWithProgress(root, repos, repos, 2, time.Minute, io.Discard, true, false, false, live, nil)
|
||||||
|
|
||||||
|
if len(results) != len(repos) {
|
||||||
|
t.Fatalf("results = %d, want %d", len(results), len(repos))
|
||||||
|
}
|
||||||
|
for _, result := range results {
|
||||||
|
if !result.Skipped {
|
||||||
|
t.Fatalf("%s was not skipped: %+v", result.Path, result)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if live.done != len(repos) {
|
||||||
|
t.Fatalf("done = %d, want %d", live.done, len(repos))
|
||||||
|
}
|
||||||
|
if live.active != 0 {
|
||||||
|
t.Fatalf("active = %d after every repository finished, want 0", live.active)
|
||||||
|
}
|
||||||
|
if live.skipped != len(repos) {
|
||||||
|
t.Fatalf("skipped = %d, want %d", live.skipped, len(repos))
|
||||||
|
}
|
||||||
|
// Intermediate frames legitimately report active repositories; only the
|
||||||
|
// last one must be settled.
|
||||||
|
frames := strings.Split(buf.String(), "\r")
|
||||||
|
if last := frames[len(frames)-1]; strings.Contains(last, "active") {
|
||||||
|
t.Fatalf("a settled run still reports active repositories: %q", last)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user