From 270812de3ee39d893332c1bdfd0c182ea4c80ff7 Mon Sep 17 00:00:00 2001 From: Hojun-Cho Date: Sat, 20 Jun 2026 23:43:15 +0900 Subject: [PATCH] httpdl: drop work-stealing segment pool for one worker per segment The pool let an idle worker steal the back half of a slower in-flight segment, behind a mutex that also serialized every advance. At the default -x 1 it never fired, and it carried the trickiest invariant in the tree (minSplit >= readBuf keeps an owner's in-flight write out of the stolen tail). makeSegments already caps the segment count at the connection count, so the segments map one-to-one onto workers: give each its own goroutine, advance written with a plain atomic add, and snapshot the fixed slice with no lock. pool/newPool/acquire/steal/advance and the seg.owned field all go away. Net -156 lines. Behaviour at the user's flags is unchanged, except a slow mirror's tail segment is no longer rebalanced onto idle connections. --- httpdl/control.go | 8 +- httpdl/httpdl.go | 57 ++++++-------- httpdl/segment.go | 123 +++-------------------------- httpdl/segment_test.go | 170 +++++++++++++++-------------------------- 4 files changed, 101 insertions(+), 257 deletions(-) diff --git a/httpdl/control.go b/httpdl/control.go index 1658182..81ff078 100644 --- a/httpdl/control.go +++ b/httpdl/control.go @@ -28,10 +28,10 @@ type segState struct { func controlPath(out string) string { return out + ".got" } -// snapshot builds a control record from the live segments. Callers hold the -// pool lock (see pool.snapshot) so the slice and each segment's frontier are -// stable for the read. -func snapshot(url string, total int64, etag, lastmod string, segs []*seg) control { +// snapshot builds a control record from the live segments. The segment slice is +// fixed once the file is divided and each segment's frontier is owned by one +// worker, so reading written atomically gives a consistent record with no lock. +func snapshot(url string, total int64, etag, lastmod string, segs []seg) control { c := control{URL: url, Total: total, ETag: etag, LastModified: lastmod, Segs: make([]segState, len(segs))} for i := range segs { c.Segs[i] = segState{segs[i].start, segs[i].endOff(), atomic.LoadInt64(&segs[i].written)} diff --git a/httpdl/httpdl.go b/httpdl/httpdl.go index 5d3ecd0..07e0bc4 100644 --- a/httpdl/httpdl.go +++ b/httpdl/httpdl.go @@ -1,6 +1,6 @@ // Package httpdl downloads a single HTTP(S) resource over one or more // connections. The model is deliberately flat: split the file into byte-range -// segments, hand them to a pool of worker goroutines, and have each worker +// segments, give each segment its own worker goroutine, and have each worker // stream its range straight to disk with WriteAt (safe for concurrent, // non-overlapping writes — no shared seek, no mutex). Goroutines blocking on // real I/O keep the model simple: no segment manager, no piece storage, no @@ -715,36 +715,27 @@ func (d *Download) segmented(ctx context.Context, out string, total int64, etag, ctx, cancel := context.WithCancel(ctx) defer cancel() - p := newPool(segs, d.cfg.MinSplit) - // Periodically persist progress so a crash can resume. saveDone := make(chan struct{}) - go d.saveLoop(ctx, out, total, etag, lastmod, p, saveDone) + go d.saveLoop(ctx, out, total, etag, lastmod, segs, saveDone) - // One worker per connection. A worker fetches segments until acquire runs - // dry, which happens only once every remaining byte is owned; a worker that - // finishes early steals the tail of a slower segment rather than sitting idle - // while the last segment drains over a single connection. + // One worker per segment. Segments are non-overlapping byte ranges written + // with WriteAt, so the workers share no cursor and need no coordination — a + // finished worker simply exits. The division (makeSegments / a resumed + // control) fixes the parallelism up front; there is no rebalancing. var ( wg sync.WaitGroup errOnce sync.Once runErr error ) - for w := 0; w < conns; w++ { + for i := range segs { wg.Add(1) - go func() { + go func(s *seg) { defer wg.Done() - for { - s := p.acquire() - if s == nil { - return - } - if err := d.fetchSeg(ctx, f, p, s, total); err != nil { - errOnce.Do(func() { runErr = err; cancel() }) - return - } + if err := d.fetchSeg(ctx, f, s, total); err != nil { + errOnce.Do(func() { runErr = err; cancel() }) } - }() + }(&segs[i]) } wg.Wait() cancel() @@ -753,7 +744,7 @@ func (d *Download) segmented(ctx context.Context, out string, total int64, etag, if runErr != nil { return runErr // keep the control file for a later -c } - p.snapshot(d.primary(), total, etag, lastmod).save(out) + snapshot(d.primary(), total, etag, lastmod, segs).save(out) removeControl(out) return nil } @@ -761,7 +752,7 @@ func (d *Download) segmented(ctx context.Context, out string, total int64, etag, // fetchSeg downloads one segment, retrying from its resume point on error. // Each attempt re-reads s.offset(), so a retry continues from the bytes already // written rather than restarting the range. -func (d *Download) fetchSeg(ctx context.Context, f *os.File, p *pool, s *seg, total int64) error { +func (d *Download) fetchSeg(ctx context.Context, f *os.File, s *seg, total int64) error { // Try each mirror in turn for this segment: a transient error retries the // same mirror under withRetries; an exhausted budget or a per-mirror permanent // error (a 404, a non-206, a length mismatch) falls over to the next mirror. @@ -773,12 +764,12 @@ func (d *Download) fetchSeg(ctx context.Context, f *os.File, p *pool, s *seg, to if s.done() { return nil } - return d.fetchOnce(ctx, f, p, s, uri, total) + return d.fetchOnce(ctx, f, s, uri, total) }) }) } -func (d *Download) fetchOnce(ctx context.Context, f *os.File, p *pool, s *seg, uri string, total int64) error { +func (d *Download) fetchOnce(ctx context.Context, f *os.File, s *seg, uri string, total int64) error { reqCtx, cancel := context.WithCancel(ctx) defer cancel() req, err := d.request(reqCtx, uri, fmt.Sprintf("bytes=%d-%d", s.offset(), s.endOff())) @@ -815,7 +806,7 @@ func (d *Download) fetchOnce(ctx context.Context, f *os.File, p *pool, s *seg, u } body, stop := d.idleGuard(resp.Body, cancel) defer stop() - return d.pump(ctx, f, p, s, body) + return d.pump(ctx, f, s, body) } // errIdleTimeout marks a transfer that stalled past the idle window, so callers @@ -986,10 +977,10 @@ func statusError(ctxMsg string, code int) error { } // pump copies the response body into the file at the segment's running offset, -// stopping at the segment end and respecting rate limits. The loop reads -// remaining() afresh each turn, so a steal that shrinks s mid-transfer simply -// ends the loop early at the new end, leaving the stolen tail to its new worker. -func (d *Download) pump(ctx context.Context, f *os.File, p *pool, s *seg, body io.Reader) error { +// stopping at the segment end and respecting rate limits. Only this worker writes +// s.written, so the advance is a plain atomic add (Stat and the snapshot read it +// atomically); the loop ends when the range is filled. +func (d *Download) pump(ctx context.Context, f *os.File, s *seg, body io.Reader) error { buf := make([]byte, readBuf) for s.remaining() > 0 { n := int64(len(buf)) @@ -1001,7 +992,7 @@ func (d *Download) pump(ctx context.Context, f *os.File, p *pool, s *seg, body i if _, werr := f.WriteAt(buf[:rd], s.offset()); werr != nil { return werr } - p.advance(s, int64(rd)) + s.addWritten(int64(rd)) atomic.AddInt64(&d.completed, int64(rd)) d.throttle(ctx, rd) } @@ -1156,7 +1147,7 @@ func (d *Download) singleOnce(ctx context.Context, out, uri string, resumeFromDi } } -func (d *Download) saveLoop(ctx context.Context, out string, total int64, etag, lastmod string, p *pool, done chan<- struct{}) { +func (d *Download) saveLoop(ctx context.Context, out string, total int64, etag, lastmod string, segs []seg, done chan<- struct{}) { defer close(done) interval := d.cfg.AutoSaveInterval if interval <= 0 { @@ -1167,10 +1158,10 @@ func (d *Download) saveLoop(ctx context.Context, out string, total int64, etag, for { select { case <-ctx.Done(): - p.snapshot(d.primary(), total, etag, lastmod).save(out) + snapshot(d.primary(), total, etag, lastmod, segs).save(out) return case <-t.C: - p.snapshot(d.primary(), total, etag, lastmod).save(out) + snapshot(d.primary(), total, etag, lastmod, segs).save(out) } } } diff --git a/httpdl/segment.go b/httpdl/segment.go index e577aaf..edd0ba4 100644 --- a/httpdl/segment.go +++ b/httpdl/segment.go @@ -1,27 +1,22 @@ package httpdl -import ( - "sync" - "sync/atomic" -) +import "sync/atomic" // seg is one contiguous byte range of the output file, downloaded by a single // ranged GET. A segment IS a byte range — a plain HTTP downloader does not need // the Piece/Segment/block layering that exists only to share code with -// BitTorrent. written and end are updated with atomic ops so Stat() and the -// resume snapshot can read them while a worker advances written or a steal -// shrinks end (see pool). +// BitTorrent. start and end are fixed when the file is divided; written is +// advanced (atomically) by the segment's one worker, so Stat() and the resume +// snapshot can read its frontier while it downloads. type seg struct { index int start int64 // first byte offset, inclusive - end int64 // last byte offset, inclusive (a steal may shrink this) + end int64 // last byte offset, inclusive written int64 // bytes already written into this segment (the resume point) - owned bool // a worker is fetching this segment; guarded by pool.mu } -func (s *seg) endOff() int64 { return atomic.LoadInt64(&s.end) } -func (s *seg) setEnd(v int64) { atomic.StoreInt64(&s.end, v) } -func (s *seg) length() int64 { return s.endOff() - s.start + 1 } +func (s *seg) length() int64 { return s.end - s.start + 1 } +func (s *seg) endOff() int64 { return s.end } func (s *seg) done() bool { return atomic.LoadInt64(&s.written) >= s.length() } func (s *seg) addWritten(n int64) { atomic.AddInt64(&s.written, n) } func (s *seg) progress() int64 { return atomic.LoadInt64(&s.written) } @@ -30,7 +25,9 @@ func (s *seg) remaining() int64 { return s.length() - atomic.LoadInt64(&s.writ // makeSegments divides a file of total bytes into contiguous segments, using at // most conns of them and never splitting below minSplit. The remainder lands in -// the last segment. With conns==1 (the default) this yields one segment. +// the last segment. With conns==1 (the default) this yields one segment. Each +// segment gets its own worker, so the division here fixes the parallelism: there +// is no later rebalancing. func makeSegments(total, minSplit int64, conns int) []seg { if conns < 1 { conns = 1 @@ -82,103 +79,3 @@ func autoConns(total, minSplit int64) int { } return int(n) } - -// pool is the live set of segments for one segmented download. A worker takes a -// segment with acquire and fetches it to completion; once no fresh segment is -// left, an idle worker steals the back half of whichever in-flight segment has -// the most still to download, so the last slow segment is shared across the idle -// connections instead of draining alone. aria2 does the same on-demand split via -// its SegmentMan; without it a fixed pre-division leaves connections idle while -// one slow mirror finishes. -// -// The mutex serialises a steal against the byte accounting it splits: advance (a -// worker committing a write) takes it too, so a steal reading a victim's frontier -// can never race the owner advancing past the chosen split point. A steal fires -// only when the victim still has at least 2*minSplit to go and cuts at the -// midpoint, so both halves stay >= minSplit (--min-split-size) and the half the -// owner keeps is far larger than one read buffer — the owner's in-flight write -// therefore can never reach into the stolen tail. -type pool struct { - mu sync.Mutex - segs []*seg - minSplit int64 -} - -// newPool wraps the pre-divided segments in a pool. Each is copied into its own -// allocation so a later steal can append a tail without invalidating the pointers -// workers already hold. -// -// minSplit is floored at readBuf: a steal leaves the owner the front half of the -// split, which is at least minSplit, while the owner's in-flight write is at most -// readBuf, so minSplit >= readBuf is exactly what keeps that write out of the -// stolen tail. The CLI already holds --min-split-size well above readBuf (>= 1 -// MiB), so the floor only guards direct callers and tests — but it keeps the -// no-overlap invariant inside this file rather than resting on a distant flag. -func newPool(segs []seg, minSplit int64) *pool { - if minSplit < readBuf { - minSplit = readBuf - } - p := &pool{minSplit: minSplit, segs: make([]*seg, len(segs))} - for i := range segs { - s := segs[i] - p.segs[i] = &s - } - return p -} - -// acquire returns the next segment for a worker to fetch: a fresh one if any -// remain, otherwise the tail split off the busiest in-flight segment. It returns -// nil when every remaining byte is already owned, i.e. the download is finishing -// and there is nothing left to steal. -func (p *pool) acquire() *seg { - p.mu.Lock() - defer p.mu.Unlock() - for _, s := range p.segs { - if !s.owned && !s.done() { - s.owned = true - return s - } - } - return p.steal() -} - -// steal splits the back half off the in-flight segment with the most remaining -// and returns it, or nil if none has enough left to be worth splitting. The -// caller holds p.mu, which keeps every owner's advance out so each victim's -// frontier is stable while we choose and commit the split. -func (p *pool) steal() *seg { - var victim *seg - for _, s := range p.segs { - if s.owned && !s.done() && s.remaining() >= 2*p.minSplit { - if victim == nil || s.remaining() > victim.remaining() { - victim = s - } - } - } - if victim == nil { - return nil - } - mid := victim.offset() + victim.remaining()/2 - // index only feeds the mirror round-robin (d.mirror) and the "segment N" log - // label; a tail's index need not be contiguous or unique, so len is fine. - tail := &seg{index: len(p.segs), start: mid, end: victim.endOff(), owned: true} - victim.setEnd(mid - 1) - p.segs = append(p.segs, tail) - return tail -} - -// advance commits n freshly written bytes to s under the pool lock, so a -// concurrent steal sees a stable frontier for the segment it may split. -func (p *pool) advance(s *seg, n int64) { - p.mu.Lock() - s.addWritten(n) - p.mu.Unlock() -} - -// snapshot builds the resume record from the live set, including any stolen -// tails, under the lock so it never races a steal appending a segment. -func (p *pool) snapshot(url string, total int64, etag, lastmod string) control { - p.mu.Lock() - defer p.mu.Unlock() - return snapshot(url, total, etag, lastmod, p.segs) -} diff --git a/httpdl/segment_test.go b/httpdl/segment_test.go index 86203ea..8a4c774 100644 --- a/httpdl/segment_test.go +++ b/httpdl/segment_test.go @@ -215,101 +215,58 @@ func TestParseContentRange(t *testing.T) { } } -// assertTiles checks that the pool's segments still cover [0,total) exactly: -// sorted by start they must be contiguous with no gap and no overlap. -func assertTiles(t *testing.T, p *pool, total int64) { +// tilesCover checks that segs cover [0,total) exactly: sorted by start they must +// be contiguous with no gap and no overlap. +func tilesCover(t *testing.T, segs []seg, total int64) { t.Helper() - p.mu.Lock() - type rng struct{ start, end int64 } - rs := make([]rng, len(p.segs)) - for i, s := range p.segs { - rs[i] = rng{s.start, s.endOff()} - } - p.mu.Unlock() - sort.Slice(rs, func(i, j int) bool { return rs[i].start < rs[j].start }) + sorted := append([]seg(nil), segs...) + sort.Slice(sorted, func(i, j int) bool { return sorted[i].start < sorted[j].start }) var next int64 - for _, r := range rs { - if r.start != next { - t.Fatalf("segment gap/overlap: next byte %d, got start %d (ranges %v)", next, r.start, rs) + for _, s := range sorted { + if s.start != next { + t.Fatalf("segment gap/overlap: next byte %d, got start %d", next, s.start) } - next = r.end + 1 + next = s.endOff() + 1 } if next != total { t.Fatalf("segments cover %d bytes, want %d", next, total) } } -func TestPoolAcquireThenSteal(t *testing.T) { - // Sizes are in units of readBuf because that is what newPool floors minSplit - // to and what the no-overlap invariant is stated against. - const total = 8 * readBuf - p := newPool(makeSegments(total, readBuf, 1), readBuf) // one segment [0,total) - - first := p.acquire() - if first == nil || first.start != 0 || first.endOff() != total-1 { - t.Fatalf("first acquire = %+v, want the whole [0,%d] segment", first, total-1) - } - if p.acquire() == nil { - // nothing fresh left, so this must steal first's back half - t.Fatal("second acquire returned nil, want a stolen tail") - } - // first kept the front half, a tail took the back half; together they still tile. - if first.endOff() != total/2-1 { - t.Errorf("victim end = %d, want %d after midpoint split", first.endOff(), total/2-1) - } - assertTiles(t, p, total) -} - -func TestPoolNoStealBelowThreshold(t *testing.T) { - // remaining is below 2*minSplit, so the lone segment is not worth splitting - // and an idle worker is told there is nothing to do. - p := newPool(makeSegments(readBuf+readBuf/2, readBuf, 1), readBuf) - if s := p.acquire(); s == nil { - t.Fatal("first acquire returned nil, want the only segment") - } - if s := p.acquire(); s != nil { - t.Fatalf("second acquire = %+v, want nil (tail would be below min-split)", s) - } -} - -// TestPoolConcurrentCoverage drives the pool the way real workers do — acquire a -// segment, copy it in small chunks, commit each with p.advance, repeat — across -// more workers than initial segments so stealing is forced. The coverage array -// proves the core invariant: every byte is written exactly once, so a steal -// never overlaps the owner's writes and never leaves a gap. -func TestPoolConcurrentCoverage(t *testing.T) { +// TestSegmentsConcurrentCoverage drives the segments the way real workers do — +// one worker per segment, each copying its range in small chunks and committing +// with addWritten. The coverage array proves the core invariant: every byte is +// written exactly once, with no gap or overlap between adjacent segments. +func TestSegmentsConcurrentCoverage(t *testing.T) { const ( total = int64(64 * readBuf) // 2 MiB - minSplit = int64(readBuf) // steal threshold is 2*this - chunk = int64(readBuf / 4) // one "read", well below minSplit - workers = 8 + minSplit = int64(readBuf) + chunk = int64(readBuf / 4) // one "read" + conns = 8 ) - p := newPool(makeSegments(total, minSplit, 1), minSplit) // start with a single segment + segs := makeSegments(total, minSplit, conns) + if len(segs) < 2 { + t.Fatalf("expected several segments, got %d", len(segs)) + } cover := make([]uint32, total) var wg sync.WaitGroup - for w := 0; w < workers; w++ { + for i := range segs { wg.Add(1) - go func() { + go func(s *seg) { defer wg.Done() - for { - s := p.acquire() - if s == nil { - return + for s.remaining() > 0 { + n := chunk + if s.remaining() < n { + n = s.remaining() } - for s.remaining() > 0 { - n := chunk - if s.remaining() < n { - n = s.remaining() - } - off := s.offset() - for j := off; j < off+n; j++ { - atomic.AddUint32(&cover[j], 1) - } - p.advance(s, n) + off := s.offset() + for j := off; j < off+n; j++ { + atomic.AddUint32(&cover[j], 1) } + s.addWritten(n) } - }() + }(&segs[i]) } wg.Wait() @@ -318,47 +275,46 @@ func TestPoolConcurrentCoverage(t *testing.T) { t.Fatalf("byte %d written %d times, want exactly 1", i, c) } } - assertTiles(t, p, total) - if len(p.segs) == 1 { - t.Error("no stealing happened: still one segment after 8 workers drained it") - } + tilesCover(t, segs, total) } -// TestPoolSnapshotAfterStealRoundTrips checks that a steal which happens before a -// crash survives the control file: the snapshot records the stolen tail, and a -// resume rebuilds a segment set that still tiles the file and keeps the bytes -// already written. This is the persistence path an interrupt-and-resume relies on -// but cannot deterministically trigger from the outside. -func TestPoolSnapshotAfterStealRoundTrips(t *testing.T) { - const total = 8 * readBuf - p := newPool(makeSegments(total, readBuf, 1), readBuf) - a := p.acquire() // the whole file - p.advance(a, readBuf) // owner makes some progress, then idles out - b := p.acquire() // steals a's back half - if b == nil { - t.Fatal("expected a stolen tail") +// TestSnapshotRoundTrips checks the persistence path an interrupt-and-resume +// relies on: a snapshot of partially-written segments records each frontier, and +// segsFromControl rebuilds a set that still tiles the file and keeps the bytes +// already written. +func TestSnapshotRoundTrips(t *testing.T) { + const ( + total = int64(8 * readBuf) + minSplit = int64(readBuf) + ) + segs := makeSegments(total, minSplit, 4) + if len(segs) < 2 { + t.Fatalf("expected several segments, got %d", len(segs)) + } + // Each segment makes some progress, so the snapshot has a real frontier to + // carry across the round trip. + var want int64 + for i := range segs { + n := int64(readBuf) + if n > segs[i].length() { + n = segs[i].length() + } + segs[i].addWritten(n) + want += n } - p.advance(b, readBuf) // the tail's worker makes progress too - c := p.snapshot("http://x", total, "", "") - if len(c.Segs) != 2 { - t.Fatalf("snapshot has %d segments, want 2 (original + stolen tail)", len(c.Segs)) + c := snapshot("http://x", total, "", "", segs) + if len(c.Segs) != len(segs) { + t.Fatalf("snapshot has %d segments, want %d", len(c.Segs), len(segs)) } rebuilt := segsFromControl(&c) - sort.Slice(rebuilt, func(i, j int) bool { return rebuilt[i].start < rebuilt[j].start }) - var next, written int64 + tilesCover(t, rebuilt, total) + var written int64 for _, s := range rebuilt { - if s.start != next { - t.Fatalf("rebuilt gap/overlap: next byte %d, got start %d", next, s.start) - } - next = s.endOff() + 1 written += s.progress() } - if next != total { - t.Fatalf("rebuilt covers %d bytes, want %d", next, total) - } - if written != 2*readBuf { - t.Errorf("rebuilt written = %d, want %d (bytes must survive the round trip)", written, 2*readBuf) + if written != want { + t.Errorf("rebuilt written = %d, want %d (bytes must survive the round trip)", written, want) } }