Merge pull request #37 from entireio/soph/progress-indicators · Entire

Home

Log in

Merge pull request #37 from entireio/soph/progress-indicators

6ad728f→main·

Soph·2mo ago·17 files·+1,538 added/-54 removed

Add progress indicators

Changes

17

75 unmodified lines

76
77
78
79
80
81
82

75 unmodified lines

cmd.Flags().BoolVar(&req.IncludeTags, "tags", false, "mirror tags")
    cmd.Flags().BoolVar(&req.Options.CollectStats, "stats", false, "print transfer statistics")
    cmd.Flags().BoolVar(&req.Options.MeasureMemory, "measure-memory", false, "sample elapsed time and Go heap usage")
    cmd.Flags().BoolVar(&req.Options.Progress, "progress", false, "show live per-side throughput on stderr (TTY only)")
    cmd.Flags().BoolVar(&jsonOutput, "json", false, "print JSON output")
    cmd.Flags().Int64Var(&req.Options.MaxPackBytes, "max-pack-bytes", 0, "abort bootstrap if the streamed source pack exceeds this many bytes")
    cmd.Flags().Int64Var(&req.Options.TargetMaxPackBytes, "target-max-pack-bytes", 0, "target receive-pack body size limit; batches are planned and auto-subdivided to fit")

Mcmd/git-sync/bootstrap.go+1

71 unmodified lines

72
73
74
75
76
77
78

71 unmodified lines

addProtocolFlag(cmd, &protocolVal)
    cmd.Flags().BoolVar(&req.Options.CollectStats, "stats", false, "print transfer statistics")
    cmd.Flags().BoolVar(&req.Options.MeasureMemory, "measure-memory", false, "sample elapsed time and Go heap usage")
    cmd.Flags().BoolVar(&req.Options.Progress, "progress", false, "show live per-side throughput on stderr (TTY only)")
    cmd.Flags().BoolVar(&jsonOutput, "json", false, "print JSON output")
    cmd.Flags().StringArrayVar(&haveRefs, "have-ref", nil, "source ref name to advertise as have; short names map to branches")
    cmd.Flags().StringArrayVar(&haveHashesRaw, "have", nil, "explicit object hash to advertise as have")

Mcmd/git-sync/fetch.go+1

112 unmodified lines

113
114
115
116
117
118
119

112 unmodified lines

cmd.Flags().BoolVar(&req.Policy.Prune, "prune", false, "delete managed target refs that no longer exist on source")
    cmd.Flags().BoolVar(&req.Options.CollectStats, "stats", false, "print transfer statistics")
    cmd.Flags().BoolVar(&req.Options.MeasureMemory, "measure-memory", false, "sample elapsed time and Go heap usage")
    cmd.Flags().BoolVar(&req.Options.Progress, "progress", false, "show live per-side throughput on stderr (TTY only)")
    cmd.Flags().BoolVar(&jsonOutput, "json", false, "print JSON output")
    cmd.Flags().IntVar(&req.Options.MaterializedMaxObjects, "materialized-max-objects", unstable.DefaultMaterializedMaxObjects, "abort non-relay materialized syncs above this many objects")
    cmd.Flags().Int64Var(&req.Options.MaxPackBytes, "max-pack-bytes", 0, "abort bootstrap-relay push if the streamed source pack exceeds this many bytes")

Mcmd/git-sync/syncplan.go+1

135 unmodified lines

136
137
138
139
139
140
141
142
51 unmodified lines

194
195
196
197
197
198
199
200
39 unmodified lines

240
241
242
243
243
244
245
246
1 unmodified line

248
249
250
251
251
252
253
254
23 unmodified lines

278
279
280
281
281
282
283
284
16 unmodified lines

301
302
303
304
304
305
306
307
16 unmodified lines

324
325
326
327
327
328
329
330
133 unmodified lines

464
465
466
467
467
468
469
470
28 unmodified lines

499
500
501
502
502
503
504
505

135 unmodified lines

}
    defer ioutil.CheckClose(reader, &err)
    // Commit-graph fetches are short and not user-facing; skip progress.
    return storeV2FetchPack(store, reader, false)
    return storeV2FetchPack(store, reader, false, nil)
}

// Capabilities returns the sorted capability list for display.
51 unmodified lines

return err
    }
    defer ioutil.CheckClose(reader, &err)
    return storeV2FetchPack(store, reader, verbose)
    return storeV2FetchPack(store, reader, verbose, conn.ProgressOut)
}

func fetchPackV2(
39 unmodified lines

if err != nil {
        return nil, err
    }
    packStream, err := openV2PackStream(reader, verbose)
    packStream, err := openV2PackStream(reader, verbose, conn.ProgressOut)
    if err != nil {
        _ = reader.Close()
        return nil, err
1 unmodified line

return packStream, nil
}

func storeV2FetchPack(store storer.Storer, r io.Reader, verbose bool) error {
func storeV2FetchPack(store storer.Storer, r io.Reader, verbose bool, progressOut io.Writer) error {
    reader := NewPacketReader(r)
    expectPackfile := false
    for {
23 unmodified lines

switch line {
            case "packfile\n":
                demux := sideband.NewDemuxer(sideband.Sideband64k, reader.BufReader())
                demux.Progress = progressSink(verbose, "source: ")
                demux.Progress = progressSink(verbose, "source: ", progressOut)
                if err := packfile.UpdateObjectStorage(store, demux); err != nil {
                    return fmt.Errorf("update object storage: %w", err)
                }
16 unmodified lines

}
}

func openV2PackStream(body io.ReadCloser, verbose bool) (io.ReadCloser, error) {
func openV2PackStream(body io.ReadCloser, verbose bool, progressOut io.Writer) (io.ReadCloser, error) {
    reader := NewPacketReader(body)
    for {
        kind, payload, err := reader.ReadPacket()
16 unmodified lines

switch line {
            case "packfile\n":
                demux := sideband.NewDemuxer(sideband.Sideband64k, reader.BufReader())
                demux.Progress = progressSink(verbose, "source: ")
                demux.Progress = progressSink(verbose, "source: ", progressOut)
                return &wrappedRC{
                    Reader: demux,
                    Closer: body,
133 unmodified lines

if drainErr := drainTrailingNAKs(buffered); drainErr != nil {
        return fmt.Errorf("drain server response: %w", drainErr)
    }
    sbReader := buildSidebandReader(caps, buffered, progressSink(verbose, "source: "))
    sbReader := buildSidebandReader(caps, buffered, progressSink(verbose, "source: ", conn.ProgressOut))
    if err := packfile.UpdateObjectStorage(store, sbReader); err != nil {
        return fmt.Errorf("update object storage: %w", err)
    }
28 unmodified lines

return nil, fmt.Errorf("drain server response: %w", drainErr)
    }
    return &wrappedRC{
        Reader: buildSidebandReader(caps, buffered, progressSink(verbose, "source: ")),
        Reader: buildSidebandReader(caps, buffered, progressSink(verbose, "source: ", conn.ProgressOut)),
        Closer: reader,
    }, nil
}

Minternal/gitproto/fetch.go+9/-9

136 unmodified lines

137
138
139
140
140
141
142
142
143
144
144
145
146
146
147
148
149
663 unmodified lines

813
814
815
816
816
817
818
819
8 unmodified lines

828
829
830
831
831
832
833
834
14 unmodified lines

849
850
851
852
852
853
854
855
14 unmodified lines

870
871
872
873
873
874
875
876
17 unmodified lines

894
895
896
897
897
898
899
900

136 unmodified lines

}

func TestProgressWriter(t *testing.T) {
    w := progressWriter(false)
    w := progressWriter(false, nil)
    if w != nil {
        t.Error("progressWriter(false) should return nil")
        t.Error("progressWriter(false, nil) should return nil")
    }
    w = progressWriter(true)
    w = progressWriter(true, nil)
    if w == nil {
        t.Error("progressWriter(true) should return non-nil writer")
        t.Error("progressWriter(true, nil) should return non-nil writer")
    }
}

663 unmodified lines

t.Fatalf("write remote error: %v", err)
    }

err := storeV2FetchPack(memory.NewStorage(), &wire, false)
    err := storeV2FetchPack(memory.NewStorage(), &wire, false, nil)
    if err == nil {
        t.Fatal("expected remote error")
    }
8 unmodified lines

t.Fatalf("write remote error: %v", err)
    }

_, err := openV2PackStream(io.NopCloser(&wire), false)
    _, err := openV2PackStream(io.NopCloser(&wire), false, nil)
    if err == nil {
        t.Fatal("expected remote error")
    }
14 unmodified lines

t.Fatalf("write flush: %v", err)
    }

err := storeV2FetchPack(memory.NewStorage(), &wire, false)
    err := storeV2FetchPack(memory.NewStorage(), &wire, false, nil)
    if err == nil {
        t.Fatal("expected missing packfile error")
    }
14 unmodified lines

t.Fatalf("write flush: %v", err)
    }

_, err := openV2PackStream(io.NopCloser(&wire), false)
    _, err := openV2PackStream(io.NopCloser(&wire), false, nil)
    if err == nil {
        t.Fatal("expected missing packfile error")
    }
17 unmodified lines

t.Fatalf("write flush: %v", err)
    }

err := storeV2FetchPack(memory.NewStorage(), &wire, false)
    err := storeV2FetchPack(memory.NewStorage(), &wire, false, nil)
    if err == nil {
        t.Fatal("expected missing packfile error")
    }

Minternal/gitproto/fetch_test.go+9/-9

124 unmodified lines

125
126
127
128
128
129
130
131
132
132
133
134
135
97 unmodified lines

233
234
235
236
236
237
238
239
240
240
241
242
243
244
245
246
247
245
246
248
249
250
251
252
253
254
250
255
256
257
258
259
260
261

124 unmodified lines

switch {
    case req.Capabilities.Supports(capability.Sideband64k):
        dem := sideband.NewDemuxer(sideband.Sideband64k, reader)
        dem.Progress = progressSink(verbose, "target: ")
        dem.Progress = progressSink(verbose, "target: ", conn.ProgressOut)
        respReader = dem
    case req.Capabilities.Supports(capability.Sideband):
        dem := sideband.NewDemuxer(sideband.Sideband, reader)
        dem.Progress = progressSink(verbose, "target: ")
        dem.Progress = progressSink(verbose, "target: ", conn.ProgressOut)
        respReader = dem
    }

97 unmodified lines

return sendReceivePack(ctx, conn, req, nil, verbose)
}

func progressWriter(verbose bool) io.Writer {
func progressWriter(verbose bool, dest io.Writer) io.Writer {
    if !verbose {
        return nil
    }
    return os.Stderr
    if dest == nil {
        dest = os.Stderr
    }
    return dest
}

// progressSink returns a line-prefixing io.Writer suitable for
// sideband.Demuxer.Progress. When verbose is false it returns nil so the
// demuxer discards progress frames without allocating.
func progressSink(verbose bool, prefix string) io.Writer {
// demuxer discards progress frames without allocating. Passing a non-nil
// dest routes the prefixed lines through that writer instead of os.Stderr,
// which lets a live progress reporter coordinate output.
func progressSink(verbose bool, prefix string, dest io.Writer) io.Writer {
    if !verbose {
        return nil
    }
    return &prefixedLineWriter{w: os.Stderr, prefix: prefix, atLineStart: true}
    if dest == nil {
        dest = os.Stderr
    }
    return &prefixedLineWriter{w: dest, prefix: prefix, atLineStart: true}
}

// prefixedLineWriter prepends a fixed prefix to each line of input written

Minternal/gitproto/push.go+15/-7

70 unmodified lines

71
72
73
74
74
75
76
77
77
78
79
80
5 unmodified lines

86
87
88
89
89
90
91
92

70 unmodified lines

}

func TestProgressSinkNilWhenNotVerbose(t *testing.T) {
    if got := progressSink(false, "anything: "); got != nil {
    if got := progressSink(false, "anything: ", nil); got != nil {
        t.Fatalf("progressSink(false) = %T, want nil", got)
    }
    if got := progressSink(true, "source: "); got == nil {
    if got := progressSink(true, "source: ", nil); got == nil {
        t.Fatal("progressSink(true) returned nil, want non-nil writer")
    }
}
5 unmodified lines

)),
    }

rc, err := openV2PackStream(body, false)
    rc, err := openV2PackStream(body, false, nil)
    if err != nil {
        t.Fatalf("openV2PackStream: %v", err)
    }

Minternal/gitproto/push_test.go+3/-3

85 unmodified lines

86
87
88
89
90
91
92
93
94
95
96
97
98
99

85 unmodified lines

// contains the repo path. Off by default to preserve behaviour for
    // callers that rely on Endpoint being stable.
    FollowInfoRefsRedirect bool

// ProgressOut is the destination for verbose sideband progress
    // messages ("Enumerating objects: ...", "Resolving deltas: ..."
    // streamed by upload-pack and receive-pack). Nil falls back to
    // os.Stderr. Callers driving a live progress ticker can plug in a
    // coordinated writer here so server-side progress lines don't
    // clobber the in-place ticker frame.
    ProgressOut io.Writer
}

// NewConn creates a new connection to the given endpoint.

Minternal/gitproto/smarthttp.go+8

59 unmodified lines

60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
65 unmodified lines

148
149
150
151
152
153
154
155
156
4 unmodified lines

161
162
163
164
165
166
167
168
152 unmodified lines

321
322
323
324
325
326
327
328
329
31 unmodified lines

361
362
363
339
364
365
366
367
368
369
370
371
344
345
372
373
374
375
376
377
378
379
380
381
382
383
7 unmodified lines

391
392
393
394
395
396
397
398
399
400
401
402
403
404
1 unmodified line

406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
14 unmodified lines

436
437
438
439
440
441
442
389
390
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
32 unmodified lines

490
491
492
493
494
495
496
497
498
254 unmodified lines

753
754
755
693
756
757
758
759
760
761
762
698
763
764
765
766
13 unmodified lines

780
781
782
718
783
784
785
786
307 unmodified lines

1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123

59 unmodified lines

TargetMaxPack    int64
    Verbose          bool
    Logger           *slog.Logger
    // OnPhase, when non-nil, is called with a short human-readable label
    // describing the current bootstrap activity (e.g. "pack 3/8") so a
    // live progress renderer can surface what is currently in flight.
    // Called from the goroutine driving Execute; implementations must not
    // block.
    OnPhase func(string)
    // OnNotice, when non-nil, receives one-time human-readable messages
    // about discrete events worth surfacing alongside progress (pack
    // subdivision, switching to batched mode). Implementations should
    // treat each call as one log line.
    OnNotice func(string)
}

func (p Params) notice(msg string) {
    if p.OnNotice != nil {
        p.OnNotice(msg)
    }
}

// Result holds the outcome of the bootstrap strategy.
65 unmodified lines

packReader = closeOnce(packReader)

p.log("bootstrap pushing refs to target", "ref_count", len(plans))
    if p.OnPhase != nil {
        p.OnPhase("pushing pack")
    }
    cmds := convert.PlansToPushCommands(plans)
    pushErr := p.TargetPusher.PushPack(ctx, cmds, packReader)
    _ = packReader.Close()
4 unmodified lines

}
        p.log("bootstrap retrying with batched mode after target rejection",
            "target_max_pack_bytes", autoBatch)
        p.notice(fmt.Sprintf("target rejected pack — switching to batched mode (limit %s)",
            humanBytes(autoBatch)))
        p.TargetMaxPack = autoBatch
        return executeBatched(ctx, p, plans, result)
    }
152 unmodified lines

idx := startIdx
        for idx < len(batch.Checkpoints) {
            checkpoint := batch.Checkpoints[idx]
            if p.OnPhase != nil {
                p.OnPhase(fmt.Sprintf("pack %d/%d", idx+1, len(batch.Checkpoints)))
            }
            p.log("bootstrap batch push checkpoint",
                "branch", batch.Plan.TargetRef.String(),
                "batch", idx+1,
31 unmodified lines

var packObjectCount int64
            if p.TargetMaxPack > 0 && len(batch.chain) > 0 {
                subdivided := false
                packReader, packObjectCount, err = checkPackSizeAndSubdivide(packReader, p.TargetMaxPack, calibratedBytesPerObject, func() bool {
                packReader, packObjectCount, err = checkPackSizeAndSubdivide(packReader, p.TargetMaxPack, calibratedBytesPerObject, func(estimated int64) bool {
                    expanded := subdivideCheckpoints(batch.chain, current, batch.Checkpoints[idx:])
                    if len(expanded) > len(batch.Checkpoints[idx:]) {
                        oldRemaining := len(batch.Checkpoints[idx:])
                        newCount := len(expanded)
                        perPack := estimated / int64(newCount)
                        p.log("bootstrap batch subdividing before push (pack header estimate)",
                            "branch", batch.Plan.TargetRef.String(),
                            "old_remaining", len(batch.Checkpoints[idx:]),
                            "new_remaining", len(expanded),
                            "old_remaining", oldRemaining,
                            "new_remaining", newCount,
                            "estimated_bytes", estimated,
                            "calibrated_bytes_per_object", calibratedBytesPerObject)
                        p.notice(fmt.Sprintf(
                            "estimated pack ~%s exceeds target limit %s — splitting %d → %d packs (~%s each)",
                            humanBytes(estimated), humanBytes(p.TargetMaxPack),
                            oldRemaining, newCount, humanBytes(perPack),
                        ))
                        batch.Checkpoints = append(batch.Checkpoints[:idx], expanded...)
                        subdivided = true
                        return true
7 unmodified lines

continue // retry at same idx with new (smaller) checkpoint
                }
            }
            p.log("bootstrap batch push attempting",
                "branch", batch.Plan.TargetRef.String(),
                "batch", idx+1,
                "batch_total", len(batch.Checkpoints),
                "estimated_bytes", packObjectCount*calibratedBytesPerObject,
                "object_count", packObjectCount,
                "target_limit_bytes", p.TargetMaxPack,
                "calibrated_bytes_per_object", calibratedBytesPerObject)

cmds := convert.PlansToPushCommands(stagePlans)
            counter := &packReadCounter{ReadCloser: packReader}
1 unmodified line

sentBytes := counter.n
            if pushErr != nil {
                _ = packReader.Close()
                p.log("bootstrap batch push failed",
                    "branch", batch.Plan.TargetRef.String(),
                    "batch", idx+1,
                    "batch_total", len(batch.Checkpoints),
                    "estimated_bytes", packObjectCount*calibratedBytesPerObject,
                    "target_limit_bytes", p.TargetMaxPack,
                    "sent_bytes", sentBytes,
                    "object_count", packObjectCount,
                    "will_subdivide", isTargetBodyLimitError(pushErr) && len(batch.chain) > 0,
                    "error", pushErr.Error())
                if isTargetBodyLimitError(pushErr) && len(batch.chain) > 0 {
                    limit := p.TargetMaxPack
                    if parsed := targetBodyLimit(pushErr); parsed > 0 {
14 unmodified lines

factor := observedSubdivisionFactor(sentBytes, limit)
                    expanded := subdivideToFactor(batch.chain, current, batch.Checkpoints[idx:], factor)
                    if len(expanded) > len(batch.Checkpoints[idx:]) {
                        oldRemaining := len(batch.Checkpoints[idx:])
                        newCount := len(expanded)
                        p.log("bootstrap batch subdividing after target size rejection",
                            "branch", batch.Plan.TargetRef.String(),
                            "old_remaining", len(batch.Checkpoints[idx:]),
                            "new_remaining", len(expanded),
                            "old_remaining", oldRemaining,
                            "new_remaining", newCount,
                            "sent_bytes", sentBytes,
                            "limit_bytes", limit,
                            "factor", factor,
                            "error", pushErr.Error())
                        limitText := ""
                        if limit > 0 {
                            limitText = fmt.Sprintf(" (target limit %s)", humanBytes(limit))
                        }
                        p.notice(fmt.Sprintf("target rejected pack%s — splitting %d → %d packs",
                            limitText, oldRemaining, newCount))
                        batch.Checkpoints = append(batch.Checkpoints[:idx], expanded...)
                        continue // retry at same idx with new (smaller) checkpoint
                    }
32 unmodified lines

// Tag phase (issue #1)
    if len(tagPlans) > 0 {
        p.log("bootstrap batch pushing tags after branch batches", "tag_count", len(tagPlans))
        if p.OnPhase != nil {
            p.OnPhase("pushing tags")
        }
        tagTargetRefs := planner.CopyRefHashMap(p.TargetRefs)
        for _, batch := range batches {
            tagTargetRefs[batch.Plan.TargetRef] = batch.Plan.SourceHash
254 unmodified lines

// packReadCounter) catches blob-heavy repos where the static 750-byte
// average is 10–20× too low — without calibration the pre-flight
// would let oversized sub-packs through and the loop would only learn
// after another wasted ~limit-sized upload.
// after another wasted ~limit-sized upload. The subdivide callback
// receives the estimated total bytes so user-facing messages can
// quote the projected size that triggered the split.
func checkPackSizeAndSubdivide(
    r io.ReadCloser,
    batchLimit int64,
    bytesPerObject int64,
    subdivide func() bool,
    subdivide func(estimatedBytes int64) bool,
) (io.ReadCloser, int64, error) { //nolint:unparam // error return kept for future use
    if bytesPerObject <= 0 {
        bytesPerObject = estimatedBytesPerObject
13 unmodified lines

objectCount := int64(header[8])<<24 | int64(header[9])<<16 | int64(header[10])<<8 | int64(header[11])
    estimated := objectCount * bytesPerObject

if estimated > batchLimit && subdivide() {
    if estimated > batchLimit && subdivide(estimated) {
        _ = r.Close()
        return nil, objectCount, nil
    }
307 unmodified lines

return limit, true
}

// humanBytes renders a byte count in IEC-ish binary units. Local copy
// because importing the syncer's formatter would invert the dependency
// direction; the bootstrap package is meant to be standalone.
func humanBytes(n int64) string {
    const unit = 1024
    if n < unit {
        return fmt.Sprintf("%d B", n)
    }
    div, exp := int64(unit), 0
    for x := n / unit; x >= unit; x /= unit {
        div *= unit
        exp++
    }
    value := float64(n) / float64(div)
    suffix := []string{"KB", "MB", "GB", "TB", "PB"}[exp]
    if value >= 100 {
        return fmt.Sprintf("%.0f %s", value, suffix)
    }
    if value >= 10 {
        return fmt.Sprintf("%.1f %s", value, suffix)
    }
    return fmt.Sprintf("%.2f %s", value, suffix)
}

func isTargetBodyLimitError(err error) bool {
    if err == nil {
        return false

Minternal/strategy/bootstrap/bootstrap.go+97/-8

271 unmodified lines

272
273
274
275
275
276
277
278
23 unmodified lines

302
303
304
305
305
306
307
308
309
310
311
312
19 unmodified lines

332
333
334
332
335
336
337
338
12 unmodified lines

351
352
353
351
354
355
356
357
7 unmodified lines

365
366
367
365
368
369
370
371

271 unmodified lines

body = append(body, []byte("packdata")...)
        r := io.NopCloser(bytes.NewReader(body))
        subdivided := false
        got, count, err := checkPackSizeAndSubdivide(r, 1_000_000, estimatedBytesPerObject, func() bool {
        got, count, err := checkPackSizeAndSubdivide(r, 1_000_000, estimatedBytesPerObject, func(int64) bool {
            subdivided = true
            return true
        })
23 unmodified lines

header := makePackHeader(5_000_000) // 5M * 750 = 3.75 GiB estimated
        r := io.NopCloser(bytes.NewReader(header))
        subdivided := false
        got, count, err := checkPackSizeAndSubdivide(r, 2_000_000_000, estimatedBytesPerObject, func() bool {
        got, count, err := checkPackSizeAndSubdivide(r, 2_000_000_000, estimatedBytesPerObject, func(estimated int64) bool {
            subdivided = true
            if estimated <= 0 {
                t.Fatalf("subdivide should receive a positive estimate, got %d", estimated)
            }
            return true
        })
        if err != nil {
19 unmodified lines

header := makePackHeader(50_000)
        r := io.NopCloser(bytes.NewReader(header))
        subdivided := false
        _, _, err := checkPackSizeAndSubdivide(r, 500*1024*1024, 12*1024, func() bool {
        _, _, err := checkPackSizeAndSubdivide(r, 500*1024*1024, 12*1024, func(int64) bool {
            subdivided = true
            return true
        })
12 unmodified lines

header := makePackHeader(5_000_000)
        r := io.NopCloser(bytes.NewReader(header))
        subdivided := false
        _, _, err := checkPackSizeAndSubdivide(r, 2_000_000_000, 0, func() bool {
        _, _, err := checkPackSizeAndSubdivide(r, 2_000_000_000, 0, func(int64) bool {
            subdivided = true
            return true
        })
7 unmodified lines

t.Run("non-PACK data proceeds without subdivide", func(t *testing.T) {
        r := io.NopCloser(bytes.NewReader([]byte("not a pack file at all")))
        got, count, err := checkPackSizeAndSubdivide(r, 100, estimatedBytesPerObject, func() bool {
        got, count, err := checkPackSizeAndSubdivide(r, 100, estimatedBytesPerObject, func(int64) bool {
            t.Fatal("should not subdivide non-pack data")
            return true
        })

Minternal/strategy/bootstrap/bootstrap_test.go+8/-5

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434

package syncer

import (
    "fmt"
    "io"
    "os"
    "sort"
    "strings"
    "sync"
    "time"
)

// progressReporter renders live per-side throughput to a writer (typically
// os.Stderr) by sampling the statsCollector's atomic byte counters on a
// fixed interval.
//
// The visible region is at most two rows: an optional "transient" line
// above (used for in-place sideband progress like "source: Compressing
// objects: 89%") and the throughput ticker below. Each redraw uses
// cursor-up + erase-to-end-of-screen to overwrite the whole region in
// place, so '\r'-terminated sideband updates from go-git read as a
// single updating row instead of scrolling line by line.
//
// On every render we also push the current per-side byte total into a
// short ring buffer (samples) so the displayed rate reflects recent
// throughput rather than a session-wide average — the latter
// undercounts the actual transfer rate because the divisor includes
// auth and ref-listing time when no pack data is flowing.
type progressReporter struct {
    out      io.Writer
    stats    *statsCollector
    interval time.Duration
    start    time.Time

stopOnce sync.Once
    stop     chan struct{}
    done     chan struct{}

mu        sync.Mutex
    rowsDrawn int                    // rows currently occupying the live region
    lastLine  string                 // last progress line, kept so setTransient can redraw without re-sampling
    transient string                 // current sideband-progress line, "" when none
    samples   map[string]*sampleRing // per-side sliding window of recent (time, bytes) snapshots
}

func newProgressReporter(out io.Writer, stats *statsCollector, interval time.Duration) *progressReporter {
    if interval <= 0 {
        interval = 200 * time.Millisecond
    }
    return &progressReporter{
        out:      out,
        stats:    stats,
        interval: interval,
        start:    time.Now(),
        stop:     make(chan struct{}),
        done:     make(chan struct{}),
        samples:  map[string]*sampleRing{},
    }
}

// run drives the ticker until stop() is called. Safe to call once.
func (p *progressReporter) run() {
    defer close(p.done)
    ticker := time.NewTicker(p.interval)
    defer ticker.Stop()
    for {
        select {
        case <-p.stop:
            return
        case <-ticker.C:
            p.render(false)
        }
    }
}

// ANSI control sequences we use to redraw the live region in place.
// J erases from the cursor to the end of the screen (covers both rows
// when the transient is shown); %dA moves the cursor up that many
// lines. Both are widely supported.
const clearDown = "\x1b[J"\
\
// cursorUpToTopLocked positions the cursor at column 0 of the top row\
// of the currently-drawn live region. Always emits '\r' so the first\
// render (rowsDrawn=0) still starts at column 0 instead of writing in\
// the middle of whatever line the cursor was last on. Caller must hold\
// p.mu.\
func (p *progressReporter) cursorUpToTopLocked() {\
    if p.rowsDrawn > 1 {\
        fmt.Fprintf(p.out, "\x1b[%dA", p.rowsDrawn-1)\
    }\
    fmt.Fprint(p.out, "\r")\
}\
\
// drawLocked rewrites the live region in place: cursor jumps to the top\
// row, erases everything below, then writes transient (if any) followed\
// by the throughput line. Cursor is left at the end of the throughput\
// row so subsequent ticker frames see rowsDrawn=region-height.\
func (p *progressReporter) drawLocked(line string) {\
    p.cursorUpToTopLocked()\
    fmt.Fprint(p.out, clearDown)\
\
    rows := 0\
    if p.transient != "" {\
        fmt.Fprintln(p.out, p.transient)\
        rows++\
    }\
    if line != "" {\
        fmt.Fprint(p.out, line)\
        rows++\
    }\
    p.rowsDrawn = rows\
    p.lastLine = line\
}\
\
// notify writes a one-time permanent message above the live region.\
// The region is cleared first so the message lands on a clean row,\
// rowsDrawn is reset so the next render redraws the ticker below, and\
// the transient slot is cleared — a permanent line typically marks a\
// state transition (sideband phase completion, slog event, subdivision\
// notice) where the previously-shown transient progress is no longer\
// the latest activity. Safe to call concurrently with the ticker.\
func (p *progressReporter) notify(msg string) {\
    p.mu.Lock()\
    defer p.mu.Unlock()\
    p.cursorUpToTopLocked()\
    fmt.Fprint(p.out, clearDown)\
    fmt.Fprintln(p.out, msg)\
    p.rowsDrawn = 0\
    p.transient = ""\
}\
\
// setTransient updates the in-place sideband row above the ticker.\
// Pass "" to clear it. Triggers an immediate redraw so '\r'-driven\
// progress (Compressing/Counting/Resolving) feels responsive between\
// ticker intervals.\
func (p *progressReporter) setTransient(line string) {\
    p.mu.Lock()\
    defer p.mu.Unlock()\
    p.transient = line\
    p.drawLocked(p.lastLine)\
}\
\
// terminate halts the ticker, draws one final frame so the printed line\
// reflects the closing byte counts, and emits a newline so subsequent\
// command output starts on a fresh row.\
func (p *progressReporter) terminate() {\
    p.stopOnce.Do(func() {\
        close(p.stop)\
        <-p.done\
        p.render(true)\
        p.mu.Lock()\
        if p.rowsDrawn > 0 {\
            fmt.Fprintln(p.out)\
        }\
        p.mu.Unlock()\
    })\
}\
\
func (p *progressReporter) render(final bool) {\
    sides := p.stats.liveSides()\
    if len(sides) == 0 && !final {\
        return\
    }\
    sort.Slice(sides, func(i, j int) bool { return sides[i].Label < sides[j].Label })\
    elapsed := time.Since(p.start)\
    now := time.Now()\
\
    p.mu.Lock()\
    // Sample first so the rate calculation below sees the latest data\
    // point. We hold p.mu so concurrent terminate/render cannot race\
    // the per-label ring buffers.\
    instant := make(map[string]float64, len(sides))\
    for _, side := range sides {\
        ring, ok := p.samples[side.Label]\
        if !ok {\
            ring = &sampleRing{}\
            p.samples[side.Label] = ring\
        }\
        ring.push(sample{at: now, bytes: side.Bytes})\
        instant[side.Label] = ring.instantRate()\
    }\
\
    var b strings.Builder\
    for i, side := range sides {\
        if i > 0 {\
            b.WriteString(sideSeparator)\
        }\
        b.WriteString(formatSide(side, elapsed, instant[side.Label], final))\
    }\
    // Surface the current activity label (e.g. "pack 3/8") on live\
    // frames only. The final frame is implicitly "done" — appending\
    // the last in-progress phase there would read as still-running.\
    if !final {\
        if phase := p.stats.getPhase(); phase != "" {\
            b.WriteString("  (")\
            b.WriteString(phase)\
            b.WriteString(")")\
        }\
    }\
    line := b.String()\
\
    defer p.mu.Unlock()\
    p.drawLocked(line)\
}\
\
const (\
    sideSeparator    = "  │  "\
    flowArrow        = " → "\
    doneMark         = " ✓"\
    maxHostnameWidth = 30\
    idleThreshold    = 750 * time.Millisecond\
\
    // sampleCapacity is how many recent (time, bytes) snapshots each\
    // side keeps for instant-rate computation. With the default 200ms\
    // render interval this gives a 2-second sliding window — long\
    // enough to smooth out per-pack burstiness, short enough to feel\
    // responsive (the ticker visually catches up to a new transfer\
    // rate inside the first second of streaming).\
    sampleCapacity = 10\
)\
\
// sample is one (time, cumulative bytes) snapshot taken at render time.\
type sample struct {\
    at    time.Time\
    bytes int64\
}\
\
// sampleRing is a fixed-size circular buffer of recent samples used to\
// derive an instantaneous transfer rate from differences between the\
// oldest and newest sample currently in the window.\
type sampleRing struct {\
    buf   [sampleCapacity]sample\
    head  int // next write index\
    count int // valid entries (capped at sampleCapacity)\
}\
\
func (r *sampleRing) push(s sample) {\
    r.buf[r.head] = s\
    r.head = (r.head + 1) % sampleCapacity\
    if r.count < sampleCapacity {\
        r.count++\
    }\
}\
\
// instantRate returns observed bytes/second over the samples currently\
// in the ring. Returns 0 when fewer than two samples exist, when the\
// elapsed time between oldest and newest is too small to be meaningful,\
// or when no bytes were transferred in the window (treat as idle so the\
// formatter falls back to the active-window average).\
func (r *sampleRing) instantRate() float64 {\
    if r.count < 2 {\
        return 0\
    }\
    newestIdx := (r.head - 1 + sampleCapacity) % sampleCapacity\
    oldestIdx := 0\
    if r.count == sampleCapacity {\
        oldestIdx = r.head // ring full — head currently sits on the oldest\
    }\
    newest := r.buf[newestIdx]\
    oldest := r.buf[oldestIdx]\
    dur := newest.at.Sub(oldest.at)\
    if dur < 50*time.Millisecond {\
        return 0\
    }\
    delta := newest.bytes - oldest.bytes\
    if delta <= 0 {\
        return 0\
    }\
    return float64(delta) / dur.Seconds()\
}\
\
// formatSide renders a single side as host + bytes + rate, with a flow\
// arrow positioned to indicate direction: source on the left of its\
// counter, target on the right of its counter. Sides with neither label\
// fall back to "name: bytes @ rate".\
//\
// While bytes are flowing (instantBytesPerSec > 0 and the side has not\
// gone idle) the displayed rate uses the recent sliding-window number\
// so it tracks the actual wire throughput. Once the side goes idle (or\
// on the final render) we switch to the active-window average for a\
// stable post-transfer headline. forceDone and idle gaps >idleThreshold\
// append a "✓" marker.\
func formatSide(side SideBytes, fallbackDur time.Duration, instantBytesPerSec float64, forceDone bool) string {\
    name := side.Display\
    if name == "" {\
        name = side.Label\
    }\
    name = truncateHost(name, maxHostnameWidth)\
\
    done := side.Bytes > 0 && (forceDone || time.Duration(side.IdleNanos) >= idleThreshold)\
\
    var rateText string\
    if !done && instantBytesPerSec > 0 {\
        rateText = formatBytes(int64(instantBytesPerSec)) + "/s"\
    } else {\
        rateDur := fallbackDur\
        if side.ActiveNanos > 0 {\
            rateDur = time.Duration(side.ActiveNanos)\
        }\
        rateText = formatRate(side.Bytes, rateDur)\
    }\
\
    rate := formatBytes(side.Bytes) + " @ " + rateText\
    if done {\
        rate += doneMark\
    }\
\
    switch side.Label {\
    case "source":\
        return name + flowArrow + rate\
    case "target":\
        return rate + flowArrow + name\
    default:\
        return name + ": " + rate\
    }\
}\
\
// truncateHost shortens long hostnames while keeping the apex domain\
// visible. Returns the original string if it already fits within width,\
// otherwise preserves the trailing two dotted labels (e.g. "cloudflare.net")\
// and spends the remaining budget on a prefix from the leading subdomain.\
// Falls back to right-truncation when the apex alone doesn't fit.\
func truncateHost(host string, width int) string {\
    if width <= 0 || len(host) <= width {\
        return host\
    }\
    labels := strings.Split(host, ".")\
    if len(labels) >= 2 {\
        apex := labels[len(labels)-2] + "." + labels[len(labels)-1]\
        if len(apex)+2 <= width { // room for at least one prefix char + ellipsis\
            prefixBudget := width - 1 - len(apex)\
            return host[:prefixBudget] + "…" + apex\
        }\
    }\
    if width <= 1 {\
        return "…"\
    }\
    return host[:width-1] + "…"\
}\
\
// sessionStderr is an io.Writer that hands writes to the live progress\
// reporter when one is attached to the syncSession, so verbose slog\
// lines and server-side sideband progress ("Resolving deltas …") land\
// above the in-place ticker frame instead of clobbering it. Falls back\
// to os.Stderr when no reporter is active.\
//\
// Partial-line writes are buffered until a '\n' or '\r' terminator\
// arrives. This matters for prefixedLineWriter, which writes a logical\
// line in two calls — first the prefix ("source: "), then the content\
// with terminator — and would otherwise produce two separate notify\
// frames split mid-line. Use as a pointer (the buffer is stateful).\
type sessionStderr struct {\
    s   *syncSession\
    buf strings.Builder\
}\
\
func (w *sessionStderr) Write(b []byte) (int, error) {\
    if w.s == nil || w.s.progress == nil {\
        n, err := os.Stderr.Write(b)\
        if err != nil {\
            return n, fmt.Errorf("stderr write: %w", err)\
        }\
        return n, nil\
    }\
    s := string(b)\
    for s != "" {\
        i := strings.IndexAny(s, "\r\n")\
        if i < 0 {\
            w.buf.WriteString(s)\
            break\
        }\
        w.buf.WriteString(s[:i])\
        line := w.buf.String()\
        w.buf.Reset()\
        // '\r' marks an in-place sideband update (git's\
        // "Compressing 89%\r" → "Compressing 90%\r" pattern); '\n'\
        // marks a permanent line that scrolls. Route accordingly so\
        // percentage updates rewrite a single transient row instead\
        // of filling the scrollback.\
        if line != "" {\
            if s[i] == '\r' {\
                w.s.progress.setTransient(line)\
            } else {\
                w.s.progress.notify(line)\
            }\
        }\
        s = s[i+1:]\
    }\
    return len(b), nil\
}\
\
// stderrIsTTY reports whether stderr is attached to a terminal. The\
// progress ticker is suppressed otherwise because '\r' updates only make\
// sense on a TTY and would otherwise spam log files.\
func stderrIsTTY() bool {\
    fi, err := os.Stderr.Stat()\
    if err != nil {\
        return false\
    }\
    return (fi.Mode() & os.ModeCharDevice) != 0\
}\
\
// formatBytes renders byte counts in IEC-ish human units (binary base).\
func formatBytes(n int64) string {\
    const unit = 1024\
    if n < unit {\
        return fmt.Sprintf("%d B", n)\
    }\
    div, exp := int64(unit), 0\
    for x := n / unit; x >= unit; x /= unit {\
        div *= unit\
        exp++\
    }\
    value := float64(n) / float64(div)\
    suffix := []string{"KB", "MB", "GB", "TB", "PB"}[exp]\
    if value >= 100 {\
        return fmt.Sprintf("%.0f %s", value, suffix)\
    }\
    if value >= 10 {\
        return fmt.Sprintf("%.1f %s", value, suffix)\
    }\
    return fmt.Sprintf("%.2f %s", value, suffix)\
}\
\
// formatRate renders a bytes/second average over the supplied duration.\
// Returns "0 B/s" until the duration is large enough to be meaningful,\
// avoiding misleadingly large rates from sub-millisecond samples.\
func formatRate(bytes int64, dur time.Duration) string {\
    if dur < 50*time.Millisecond || bytes <= 0 {\
        return "0 B/s"\
    }\
    rate := float64(bytes) / dur.Seconds()\
    return formatBytes(int64(rate)) + "/s"\
}\
```\
\
Ainternal/syncer/progress.go+434\
\
```\
1\
2\
3\
4\
5\
6\
7\
8\
9\
10\
11\
12\
13\
14\
15\
16\
17\
18\
19\
20\
21\
22\
23\
24\
25\
26\
27\
28\
29\
30\
31\
32\
33\
34\
35\
36\
37\
38\
39\
40\
41\
42\
43\
44\
45\
46\
47\
48\
49\
50\
51\
52\
53\
54\
55\
56\
57\
58\
59\
60\
61\
62\
63\
64\
65\
66\
67\
68\
69\
70\
71\
72\
73\
74\
75\
76\
77\
78\
79\
80\
81\
82\
83\
84\
85\
86\
87\
88\
89\
90\
91\
92\
93\
94\
95\
96\
97\
98\
99\
100\
101\
102\
103\
104\
105\
106\
107\
108\
109\
110\
111\
112\
113\
114\
115\
116\
117\
118\
119\
120\
121\
122\
123\
124\
125\
126\
127\
128\
129\
130\
131\
132\
133\
134\
135\
136\
137\
138\
139\
140\
141\
142\
143\
144\
145\
146\
147\
148\
149\
150\
151\
152\
153\
154\
155\
156\
157\
158\
159\
160\
161\
162\
163\
164\
165\
166\
167\
168\
169\
170\
171\
172\
173\
174\
175\
176\
177\
178\
179\
180\
181\
182\
183\
184\
185\
186\
187\
188\
189\
190\
191\
192\
193\
194\
195\
196\
197\
198\
199\
200\
201\
202\
203\
204\
205\
206\
207\
208\
209\
210\
211\
212\
213\
214\
215\
216\
217\
218\
219\
220\
221\
222\
223\
224\
225\
226\
227\
228\
229\
230\
231\
232\
233\
234\
235\
236\
237\
238\
239\
240\
241\
242\
243\
244\
245\
246\
247\
248\
249\
250\
251\
252\
253\
254\
255\
256\
257\
258\
259\
260\
261\
262\
263\
264\
265\
266\
267\
268\
269\
270\
271\
272\
273\
274\
275\
276\
277\
278\
279\
280\
281\
282\
283\
284\
285\
286\
287\
288\
289\
290\
291\
292\
293\
294\
295\
296\
297\
298\
299\
300\
301\
302\
303\
304\
305\
306\
307\
308\
309\
310\
311\
312\
313\
314\
315\
316\
317\
318\
319\
320\
321\
322\
323\
324\
325\
326\
327\
328\
329\
330\
331\
332\
333\
334\
335\
336\
337\
338\
339\
340\
341\
342\
343\
344\
345\
346\
347\
348\
349\
350\
351\
352\
353\
354\
355\
356\
357\
358\
359\
360\
361\
362\
363\
364\
365\
366\
367\
368\
369\
370\
371\
372\
373\
374\
375\
376\
377\
378\
379\
380\
381\
382\
383\
384\
385\
386\
387\
388\
389\
390\
391\
392\
393\
394\
395\
396\
397\
398\
399\
400\
401\
402\
403\
404\
405\
406\
407\
408\
409\
410\
411\
412\
413\
414\
415\
416\
417\
418\
419\
420\
421\
422\
423\
424\
425\
426\
427\
428\
429\
430\
431\
432\
433\
434\
435\
436\
437\
438\
439\
440\
441\
442\
443\
444\
445\
446\
447\
448\
449\
450\
451\
452\
453\
454\
455\
456\
457\
458\
459\
460\
461\
462\
463\
464\
465\
466\
467\
468\
469\
470\
471\
472\
473\
474\
475\
476\
477\
478\
479\
480\
481\
482\
483\
484\
485\
486\
487\
488\
489\
490\
491\
492\
493\
494\
495\
496\
497\
498\
499\
500\
501\
502\
503\
504\
505\
506\
507\
508\
509\
510\
511\
512\
513\
514\
515\
516\
517\
518\
519\
520\
521\
522\
523\
524\
525\
526\
527\
528\
529\
530\
\
package syncer\
\
import (\
    "bytes"\
    "strings"\
    "testing"\
    "time"\
)\
\
func TestFormatBytes(t *testing.T) {\
    t.Parallel()\
    cases := []struct {\
        in   int64\
        want string\
    }{\
        {0, "0 B"},\
        {42, "42 B"},\
        {1023, "1023 B"},\
        {1024, "1.00 KB"},\
        {1500, "1.46 KB"},\
        {int64(15 * 1024), "15.0 KB"},\
        {int64(1024 * 1024), "1.00 MB"},\
        {int64(150 * 1024 * 1024), "150 MB"},\
    }\
    for _, c := range cases {\
        if got := formatBytes(c.in); got != c.want {\
            t.Errorf("formatBytes(%d) = %q, want %q", c.in, got, c.want)\
        }\
    }\
}\
\
func TestFormatRate(t *testing.T) {\
    t.Parallel()\
    const zeroRate = "0 B/s"\
    if got := formatRate(0, time.Second); got != zeroRate {\
        t.Errorf("zero bytes should produce %s, got %q", zeroRate, got)\
    }\
    if got := formatRate(1<<20, time.Microsecond); got != zeroRate {\
        t.Errorf("sub-50ms windows should produce %s to avoid bogus rates, got %q", zeroRate, got)\
    }\
    got := formatRate(1<<20, time.Second)\
    if got != "1.00 MB/s" {\
        t.Errorf("1 MiB/s should format as 1.00 MB/s, got %q", got)\
    }\
}\
\
func TestProgressReporterRendersBothSides(t *testing.T) {\
    t.Parallel()\
    stats := newStats(true)\
    stats.setSideDisplay("source", "github.com")\
    stats.setSideDisplay("target", "example.test")\
    stats.side("source").bytes.Store(2 * 1024 * 1024)\
    stats.side("target").bytes.Store(1024 * 1024)\
\
    var buf bytes.Buffer\
    p := newProgressReporter(&buf, stats, 0)\
    p.render(true)\
\
    out := buf.String()\
    if !strings.Contains(out, "github.com → 2.00 MB") {\
        t.Errorf("output missing source host with arrow: %q", out)\
    }\
    if !strings.Contains(out, "1.00 MB @ ") || !strings.Contains(out, "→ example.test") {\
        t.Errorf("output missing target host with arrow: %q", out)\
    }\
    // Every render starts at column 0 so it can never land mid-word; the\
    // first frame may be just '\r' + clear escape, subsequent ones\
    // preface with cursor-up movement when a transient row is showing.\
    if !strings.HasPrefix(out, "\r") && !strings.Contains(out[:min(8, len(out))], "\x1b[") {\
        t.Errorf("output should reposition cursor to column 0 before drawing: %q", out)\
    }\
    if !strings.Contains(out, "│") {\
        t.Errorf("output should use the vertical bar separator: %q", out)\
    }\
    // Source must precede target so the arrows describe the data path\
    // left-to-right (source on the left, target on the right).\
    if strings.Index(out, "github.com") > strings.Index(out, "example.test") {\
        t.Errorf("source should appear before target: %q", out)\
    }\
}\
\
func TestProgressReporterTerminateEmitsNewline(t *testing.T) {\
    t.Parallel()\
    stats := newStats(true)\
    stats.side("source").bytes.Store(2048)\
\
    var buf bytes.Buffer\
    p := newProgressReporter(&buf, stats, time.Hour) // long interval; we only render manually\
    go p.run()\
    p.terminate()\
\
    if !strings.HasSuffix(buf.String(), "\n") {\
        t.Errorf("terminate should leave the cursor on a fresh line, got %q", buf.String())\
    }\
}\
\
func TestThroughputLineFormatsBothSides(t *testing.T) {\
    t.Parallel()\
    stats := Stats{\
        Enabled:      true,\
        Items:        map[string]*ServiceStats{},\
        ElapsedNanos: time.Second.Nanoseconds(),\
        Sides: []SideBytes{\
            {Label: "target", Bytes: 1 << 20, Display: "example.test"},\
            {Label: "source", Bytes: 4 << 20, Display: "github.com"},\
        },\
    }\
    line := throughputLine(stats)\
    if !strings.HasPrefix(line, "throughput: ") {\
        t.Fatalf("line should start with 'throughput: ', got %q", line)\
    }\
    // Source host must precede target host so the line reads left-to-right.\
    sourceIdx := strings.Index(line, "github.com")\
    targetIdx := strings.Index(line, "example.test")\
    if sourceIdx < 0 || targetIdx < 0 || sourceIdx > targetIdx {\
        t.Errorf("expected source host before target host in %q", line)\
    }\
    if !strings.Contains(line, "github.com → 4.00 MB") {\
        t.Errorf("source side should render as 'github.com → 4.00 MB ...': %q", line)\
    }\
    if !strings.Contains(line, "→ example.test") {\
        t.Errorf("target side should end in '→ example.test': %q", line)\
    }\
    if !strings.Contains(line, "│") {\
        t.Errorf("expected vertical-bar separator between sides: %q", line)\
    }\
}\
\
func TestThroughputLineEmptyWhenNoBytes(t *testing.T) {\
    t.Parallel()\
    stats := Stats{Enabled: true, ElapsedNanos: time.Second.Nanoseconds()}\
    if got := throughputLine(stats); got != "" {\
        t.Errorf("expected empty line with no sides, got %q", got)\
    }\
}\
\
// TestFormatSideFreezesRateAtLastByte verifies that once a side has\
// gone idle, the displayed rate uses the active window (start → last\
// byte) rather than the still-advancing wall clock — so a 4 MiB\
// transfer that took 1s reads as 4 MB/s even when sampled 9 seconds\
// later, and the side is annotated with a done marker.\
func TestFormatSideFreezesRateAtLastByte(t *testing.T) {\
    t.Parallel()\
    side := SideBytes{\
        Label:       "source",\
        Display:     "github.com",\
        Bytes:       4 << 20, // 4 MiB\
        ActiveNanos: time.Second.Nanoseconds(),\
        IdleNanos:   (idleThreshold + time.Second).Nanoseconds(),\
    }\
    got := formatSide(side, 10*time.Second, 0, false)\
    if !strings.Contains(got, "4.00 MB/s") {\
        t.Errorf("idle side should freeze rate at active-window value: %q", got)\
    }\
    if !strings.HasSuffix(got, doneMark) {\
        t.Errorf("idle side should be marked done: %q", got)\
    }\
}\
\
func TestFormatSideActiveSideHasNoDoneMark(t *testing.T) {\
    t.Parallel()\
    side := SideBytes{\
        Label:       "target",\
        Display:     "example.test",\
        Bytes:       2 << 20,\
        ActiveNanos: time.Second.Nanoseconds(),\
        IdleNanos:   (100 * time.Millisecond).Nanoseconds(),\
    }\
    got := formatSide(side, time.Second, 0, false)\
    if strings.Contains(got, doneMark) {\
        t.Errorf("active side should not carry done marker: %q", got)\
    }\
}\
\
func TestProgressReporterRendersPhase(t *testing.T) {\
    t.Parallel()\
    stats := newStats(true)\
    stats.setSideDisplay("source", "github.com")\
    stats.setSideDisplay("target", "example.test")\
    stats.side("source").bytes.Store(1024)\
    stats.side("target").bytes.Store(512)\
    stats.setPhase("pack 3/8")\
\
    var buf bytes.Buffer\
    p := newProgressReporter(&buf, stats, 0)\
    p.render(false)\
\
    if !strings.Contains(buf.String(), "(pack 3/8)") {\
        t.Errorf("live frame should include phase suffix: %q", buf.String())\
    }\
}\
\
func TestProgressReporterFinalFrameOmitsPhase(t *testing.T) {\
    t.Parallel()\
    stats := newStats(true)\
    stats.setSideDisplay("source", "github.com")\
    stats.setSideDisplay("target", "example.test")\
    stats.side("source").bytes.Store(1024)\
    stats.side("target").bytes.Store(512)\
    stats.setPhase("pack 3/8")\
\
    var buf bytes.Buffer\
    p := newProgressReporter(&buf, stats, 0)\
    p.render(true)\
\
    if strings.Contains(buf.String(), "pack 3/8") {\
        t.Errorf("final frame should drop the in-flight phase: %q", buf.String())\
    }\
}\
\
// TestNotifyAfterRenderClearsTheFrame asserts that bytes written to\
// notify after a render include a clear-line escape before the message,\
// so the slog/sideband line doesn't end up concatenated with the live\
// progress frame on the user's terminal.\
func TestNotifyAfterRenderClearsTheFrame(t *testing.T) {\
    t.Parallel()\
    stats := newStats(true)\
    stats.setSideDisplay("source", "github.com")\
    stats.setSideDisplay("target", "example.test")\
    stats.side("source").bytes.Store(2 * 1024 * 1024)\
    stats.side("target").bytes.Store(1024 * 1024)\
    stats.setPhase("pack 1/4")\
\
    var buf bytes.Buffer\
    p := newProgressReporter(&buf, stats, 0)\
    p.render(false)\
    frameEnd := buf.Len()\
\
    p.notify("level=INFO msg=\"bootstrap subdividing\"")\
\
    tail := buf.String()[frameEnd:]\
    // notify must emit a clear-screen-from-cursor escape before the\
    // message so the previous live region (transient + ticker) is wiped,\
    // plus a trailing newline so subsequent renders draw on a fresh row.\
    if !strings.Contains(tail, clearDown) {\
        t.Errorf("notify should clear the region via ANSI J, got %q", tail)\
    }\
    if !strings.HasSuffix(tail, "\n") {\
        t.Errorf("notify should terminate with newline, got %q", tail)\
    }\
    clearIdx := strings.Index(tail, clearDown)\
    msgIdx := strings.Index(tail, "level=INFO")\
    if clearIdx < 0 || msgIdx < 0 || clearIdx > msgIdx {\
        t.Errorf("clear must precede the message in %q", tail)\
    }\
}\
\
func TestSessionStderrRoutesMultilineThroughNotify(t *testing.T) {\
    t.Parallel()\
    stats := newStats(true)\
    stats.setSideDisplay("source", "github.com")\
    stats.side("source").bytes.Store(1024)\
\
    var buf bytes.Buffer\
    p := newProgressReporter(&buf, stats, 0)\
    p.render(false) // give notify something to clear\
\
    sess := &syncSession{progress: p}\
    sink := &sessionStderr{s: sess}\
\
    // Multi-line write (e.g. a slog line followed by a sideband line)\
    // must produce one notify per logical line — both '\n' and '\r' are\
    // treated as line ends so sideband '\r'-driven percentage updates\
    // don't clobber the live progress frame.\
    if _, err := sink.Write([]byte("first line\nsecond line\rthird line\n")); err != nil {\
        t.Fatalf("write: %v", err)\
    }\
\
    for _, want := range []string{"first line", "second line", "third line"} {\
        if !strings.Contains(buf.String(), want) {\
            t.Errorf("output missing %q: %q", want, buf.String())\
        }\
    }\
}\
\
// TestSessionStderrBuffersPartialLines covers the prefixedLineWriter\
// pattern in internal/gitproto: a logical line arrives in two writes —\
// first the "source: " prefix, then the content with terminator. The\
// buffered writer must combine them into one notify call instead of\
// emitting the prefix on its own row.\
// TestSessionStderrCRUpdatesTransient verifies that '\r'-terminated\
// sideband progress (git's "Compressing 89%\r" → "Compressing 90%\r"\
// pattern) goes to the transient row instead of scrolling the\
// scrollback. Subsequent updates should overwrite the transient slot\
// rather than each landing on a new row.\
func TestSessionStderrCRUpdatesTransient(t *testing.T) {\
    t.Parallel()\
    stats := newStats(true)\
    stats.setSideDisplay("source", "github.com")\
    stats.side("source").bytes.Store(1024)\
\
    var buf bytes.Buffer\
    p := newProgressReporter(&buf, stats, 0)\
    p.render(false)\
\
    sess := &syncSession{progress: p}\
    sink := &sessionStderr{s: sess}\
\
    // Two consecutive in-place sideband updates plus one final \n line.\
    if _, err := sink.Write([]byte("source: Compressing 50%\r")); err != nil {\
        t.Fatalf("write: %v", err)\
    }\
    if _, err := sink.Write([]byte("source: Compressing 75%\r")); err != nil {\
        t.Fatalf("write: %v", err)\
    }\
    if _, err := sink.Write([]byte("source: Compressing 100%, done.\n")); err != nil {\
        t.Fatalf("write: %v", err)\
    }\
\
    if p.transient != "" {\
        t.Errorf("permanent line should clear transient, got %q", p.transient)\
    }\
    out := buf.String()\
    if !strings.Contains(out, "Compressing 50%") ||\
        !strings.Contains(out, "Compressing 75%") ||\
        !strings.Contains(out, "Compressing 100%, done.") {\
        t.Errorf("all three sideband states should be in the output stream:\n%s", out)\
    }\
}\
\
func TestSessionStderrBuffersPartialLines(t *testing.T) {\
    t.Parallel()\
    stats := newStats(true)\
    stats.setSideDisplay("source", "github.com")\
    stats.side("source").bytes.Store(1024)\
\
    var buf bytes.Buffer\
    p := newProgressReporter(&buf, stats, 0)\
    p.render(false)\
\
    sess := &syncSession{progress: p}\
    sink := &sessionStderr{s: sess}\
\
    if _, err := sink.Write([]byte("source: ")); err != nil {\
        t.Fatalf("prefix write: %v", err)\
    }\
    if _, err := sink.Write([]byte("Counting objects: 10%\r")); err != nil {\
        t.Fatalf("content write: %v", err)\
    }\
\
    out := buf.String()\
    if !strings.Contains(out, "source: Counting objects: 10%") {\
        t.Errorf("expected joined line 'source: Counting objects: 10%%', got %q", out)\
    }\
    // Reject the bug shape: prefix on its own line followed by content\
    // on a separate line.\
    if strings.Contains(out, "source: \n") || strings.Contains(out, "source:\n") {\
        t.Errorf("prefix should not be emitted as a standalone line: %q", out)\
    }\
}\
\
func TestProgressReporterNotifyClearsAndRedraws(t *testing.T) {\
    t.Parallel()\
    stats := newStats(true)\
    stats.setSideDisplay("source", "github.com")\
    stats.side("source").bytes.Store(1024)\
\
    var buf bytes.Buffer\
    p := newProgressReporter(&buf, stats, 0)\
\
    // Draw a frame so the reporter has a non-zero lastLen to clear.\
    p.render(false)\
    beforeNotify := buf.Len()\
\
    p.notify("target rejected pack — splitting 1 → 4 packs")\
\
    out := buf.String()\
    if !strings.Contains(out, "splitting 1 → 4 packs") {\
        t.Errorf("notify output should contain the message: %q", out)\
    }\
    // Notify must clear the previously-drawn line so the message lands\
    // on its own row instead of overlapping the progress text.\
    tail := out[beforeNotify:]\
    if !strings.HasPrefix(tail, "\r") {\
        t.Errorf("notify should start by clearing the current line: %q", tail)\
    }\
    if !strings.HasSuffix(strings.TrimRight(tail, ""), "packs\n") {\
        t.Errorf("notify should end with a newline so progress redraws below: %q", tail)\
    }\
\
    // The next render should redraw the progress line because lastLen\
    // was reset.\
    p.render(false)\
    if !strings.HasSuffix(buf.String(), "github.com → 1.00 KB @ 0 B/s") &&\
        !strings.Contains(buf.String()[len(out):], "github.com") {\
        t.Errorf("next render after notify should redraw the progress line: %q",\
            buf.String()[len(out):])\
    }\
}\
\
func TestFormatSideForceDoneAlwaysMarks(t *testing.T) {\
    t.Parallel()\
    side := SideBytes{\
        Label:       "source",\
        Display:     "github.com",\
        Bytes:       1024,\
        ActiveNanos: time.Second.Nanoseconds(),\
        IdleNanos:   0,\
    }\
    got := formatSide(side, time.Second, 0, true)\
    if !strings.HasSuffix(got, doneMark) {\
        t.Errorf("forceDone should mark the side regardless of idle gap: %q", got)\
    }\
}\
\
// TestSampleRingComputesRateOverWindow walks the ring through a few\
// scenarios that exercise both the partial-fill path (early in a\
// transfer) and the saturated-ring path (steady state).\
func TestSampleRingComputesRateOverWindow(t *testing.T) {\
    t.Parallel()\
    r := &sampleRing{}\
    if got := r.instantRate(); got != 0 {\
        t.Errorf("empty ring should report 0, got %v", got)\
    }\
\
    base := time.Now()\
    r.push(sample{at: base, bytes: 0})\
    if got := r.instantRate(); got != 0 {\
        t.Errorf("single-sample ring should report 0, got %v", got)\
    }\
\
    // 1 MiB transferred in 1s → ~1 MiB/s.\
    r.push(sample{at: base.Add(time.Second), bytes: 1 << 20})\
    got := r.instantRate()\
    want := float64(1 << 20)\
    if got < want*0.99 || got > want*1.01 {\
        t.Errorf("expected ~%v B/s, got %v", want, got)\
    }\
\
    // Saturate the ring with a steady 10 MB/s and verify the rate\
    // reflects only the in-window samples (oldest is overwritten).\
    r2 := &sampleRing{}\
    for i := 0; i <= sampleCapacity*2; i++ {\
        r2.push(sample{\
            at:    base.Add(time.Duration(i) * time.Second),\
            bytes: int64(i) * 10 * (1 << 20),\
        })\
    }\
    got2 := r2.instantRate()\
    want2 := float64(10 * (1 << 20))\
    if got2 < want2*0.99 || got2 > want2*1.01 {\
        t.Errorf("steady-state expected ~%v B/s, got %v", want2, got2)\
    }\
}\
\
// TestSampleRingIdlePeriodReportsZero ensures that when bytes stop\
// flowing the ring eventually returns 0 — the formatter then falls\
// back to the active-window average for a stable post-transfer\
// headline.\
func TestSampleRingIdlePeriodReportsZero(t *testing.T) {\
    t.Parallel()\
    r := &sampleRing{}\
    base := time.Now()\
    for i := range sampleCapacity {\
        r.push(sample{at: base.Add(time.Duration(i) * 200 * time.Millisecond), bytes: 1024})\
    }\
    if got := r.instantRate(); got != 0 {\
        t.Errorf("flat byte count across ring should report 0, got %v", got)\
    }\
}\
\
func TestFormatSidePrefersInstantRateWhileActive(t *testing.T) {\
    t.Parallel()\
    side := SideBytes{\
        Label:       "source",\
        Display:     "github.com",\
        Bytes:       100 * (1 << 20),\
        ActiveNanos: (10 * time.Second).Nanoseconds(),\
        IdleNanos:   0, // active\
    }\
    // Active-window average = 10 MB/s. Sliding-window says 44 MB/s.\
    // The displayed rate should follow the sliding window.\
    got := formatSide(side, time.Second, float64(44*(1<<20)), false)\
    if !strings.Contains(got, "44.0 MB/s") {\
        t.Errorf("active side should show instant rate, got %q", got)\
    }\
}\
\
func TestFormatSideUsesAverageWhenDone(t *testing.T) {\
    t.Parallel()\
    side := SideBytes{\
        Label:       "source",\
        Display:     "github.com",\
        Bytes:       int64(2 * (1 << 20)),\
        ActiveNanos: time.Second.Nanoseconds(),\
        IdleNanos:   (idleThreshold + time.Second).Nanoseconds(),\
    }\
    // Even with a non-zero instant rate, an idle/done side must show\
    // the active-window average — otherwise the post-transfer headline\
    // would be a stale snapshot of pre-completion throughput.\
    got := formatSide(side, time.Second, float64(99*(1<<20)), false)\
    if !strings.Contains(got, "2.00 MB/s") {\
        t.Errorf("idle side should fall back to average rate, got %q", got)\
    }\
    if !strings.HasSuffix(got, doneMark) {\
        t.Errorf("idle side should still carry done marker: %q", got)\
    }\
}\
\
func TestTruncateHost(t *testing.T) {\
    t.Parallel()\
    cases := []struct {\
        host string\
        max  int\
        want string\
    }{\
        {"github.com", 30, "github.com"},\
        {"a.b.c.d.example.com", 30, "a.b.c.d.example.com"},\
        {\
            "8b04592ed74a5cce30d355b07276caf3.artifacts.cloudflare.net",\
            30,\
            "8b04592ed74a5cc…cloudflare.net",\
        },\
        // When the apex itself doesn't fit, we fall back to right-truncation.\
        {"averylongdomain.example.com", 10, "averylong…"},\
        {"single", 4, "sin…"},\
        {"", 30, ""},\
        {"x", 0, "x"},\
    }\
    for _, c := range cases {\
        got := truncateHost(c.host, c.max)\
        if got != c.want {\
            t.Errorf("truncateHost(%q, %d) = %q, want %q", c.host, c.max, got, c.want)\
        }\
        if c.max > 0 && len([]rune(got)) > c.max {\
            t.Errorf("truncateHost(%q, %d) = %q exceeded max width %d",\
                c.host, c.max, got, c.max)\
        }\
    }\
}\
```\
\
Ainternal/syncer/progress\_test.go+530\
\
```\
5 unmodified lines\
\
6\
7\
8\
9\
10\
11\
12\
13\
9 unmodified lines\
\
23\
24\
25\
26\
27\
28\
29\
30\
31\
32\
33\
34\
35\
36\
37\
38\
39\
40\
41\
42\
43\
44\
45\
46\
47\
48\
49\
50\
26\
27\
51\
52\
53\
54\
55\
56\
30\
57\
58\
59\
60\
61\
62\
63\
64\
65\
66\
67\
68\
69\
70\
71\
72\
73\
74\
75\
76\
77\
32\
33\
34\
78\
79\
80\
81\
82\
83\
84\
85\
86\
87\
88\
89\
90\
91\
38\
92\
93\
94\
95\
96\
97\
98\
99\
100\
17 unmodified lines\
\
118\
119\
120\
121\
122\
123\
124\
125\
126\
127\
128\
129\
130\
131\
132\
133\
134\
135\
136\
137\
138\
139\
140\
141\
142\
143\
144\
145\
146\
147\
148\
149\
150\
151\
152\
153\
154\
155\
156\
157\
158\
159\
160\
161\
162\
163\
164\
165\
166\
167\
168\
169\
170\
171\
172\
173\
174\
175\
176\
177\
178\
179\
180\
181\
182\
183\
184\
185\
186\
187\
188\
189\
190\
191\
192\
193\
194\
195\
196\
197\
198\
199\
200\
201\
202\
203\
204\
64\
205\
206\
207\
208\
209\
210\
211\
212\
213\
214\
215\
216\
217\
73\
218\
219\
220\
221\
222\
1 unmodified line\
\
224\
225\
226\
227\
228\
229\
230\
231\
232\
233\
234\
235\
236\
237\
238\
14 unmodified lines\
\
253\
254\
255\
256\
257\
258\
259\
5 unmodified lines\
\
265\
266\
267\
268\
269\
270\
271\
272\
273\
117\
274\
275\
276\
277\
278\
279\
280\
281\
282\
283\
\
5 unmodified lines\
\
    "net/http"\
    "strings"\
    "sync"\
    "sync/atomic"\
    "time"\
\
    "entire.io/entire/git-sync/internal/gitproto"\
)\
9 unmodified lines\
\
    Commands      int    `json:"commands"`\
}\
\
// SideBytes reports cumulative bytes transferred at one transport label\
// ("source" or "target"). Source is dominated by upload-pack response\
// bytes (download), target is dominated by receive-pack request bytes\
// (upload), so this is a useful "data moved on this side" number for both.\
//\
// Display carries the hostname extracted from the endpoint URL so live\
// renderers can show "github.com → … → host" without re-parsing the URL.\
// Empty when the endpoint URL was not http(s) or failed to parse.\
//\
// ActiveNanos and IdleNanos let renderers freeze the per-side rate once\
// a transfer finishes: ActiveNanos spans from stats start to the most\
// recent byte read, so dividing Bytes by ActiveNanos yields the rate\
// during active streaming rather than a value that decays as wall clock\
// keeps advancing past the last byte. IdleNanos is the gap between the\
// last byte and the snapshot, used to mark a side as "done".\
type SideBytes struct {\
    Label       string `json:"label"`\
    Bytes       int64  `json:"bytes"`\
    Display     string `json:"display,omitempty"`\
    ActiveNanos int64  `json:"activeNanos,omitempty"`\
    IdleNanos   int64  `json:"idleNanos,omitempty"`\
}\
\
// Stats holds the collected transfer statistics.\
type Stats struct {\
    Enabled bool                     `json:"enabled"`\
    Items   map[string]*ServiceStats `json:"items"`\
    Enabled      bool                     `json:"enabled"`\
    Items        map[string]*ServiceStats `json:"items"`\
    Sides        []SideBytes              `json:"sides,omitempty"`\
    ElapsedNanos int64                    `json:"elapsedNanos,omitempty"`\
}\
\
// statsCollector is a concurrency-safe stats collector (issue #8).\
// sideCounter holds atomic byte counters for one transport label.\
// Updated from the request and response body wrappers without holding\
// the collector mutex, so a progress reader can sample at high\
// frequency without contending with active transfers.\
//\
// lastByteAt is the unix-nanos timestamp of the most recent non-empty\
// Read. It is used to freeze the displayed rate once a transfer\
// finishes — without it, rate = bytes / (now − start) keeps shrinking\
// as wall clock advances past the actual end of the transfer.\
//\
// display is set once during session setup (via setSideDisplay) and\
// then read by the live progress reporter; it does not need a lock\
// because writes happen-before any reader goroutine starts.\
type sideCounter struct {\
    bytes      atomic.Int64\
    lastByteAt atomic.Int64\
    display    string\
}\
\
// statsCollector is a concurrency-safe stats collector.\
type statsCollector struct {\
    enabled bool\
    mu      sync.Mutex\
    items   map[string]*ServiceStats\
    enabled   bool\
    startedAt time.Time\
    mu        sync.Mutex\
    items     map[string]*ServiceStats\
    sidesMu   sync.RWMutex\
    sides     map[string]*sideCounter\
    // phase carries an optional one-line activity label\
    // (e.g. "pack 3/8") that the live progress reporter renders next\
    // to the per-side counters. Updated atomically by strategies and\
    // read by the reporter goroutine without contention.\
    phase atomic.Pointer[string]\
}\
\
func newStats(enabled bool) *statsCollector {\
    return &statsCollector{enabled: enabled, items: map[string]*ServiceStats{}}\
    return &statsCollector{\
        enabled:   enabled,\
        startedAt: time.Now(),\
        items:     map[string]*ServiceStats{},\
        sides:     map[string]*sideCounter{},\
    }\
}\
\
func (s *statsCollector) ensure(name string) *ServiceStats {\
17 unmodified lines\
\
    item.ResponseBytes += responseBytes\
}\
\
// side returns the per-label counter, creating it lazily on first use.\
// Counters are tracked unconditionally (independent of ShowStats) so the\
// live progress ticker works without forcing --stats.\
func (s *statsCollector) side(label string) *sideCounter {\
    s.sidesMu.RLock()\
    if sc, ok := s.sides[label]; ok {\
        s.sidesMu.RUnlock()\
        return sc\
    }\
    s.sidesMu.RUnlock()\
\
    s.sidesMu.Lock()\
    defer s.sidesMu.Unlock()\
    if sc, ok := s.sides[label]; ok {\
        return sc\
    }\
    sc := &sideCounter{}\
    s.sides[label] = sc\
    return sc\
}\
\
// setPhase records a short activity label that the live progress\
// reporter will surface alongside per-side counters. Pass "" to clear.\
func (s *statsCollector) setPhase(p string) {\
    s.phase.Store(&p)\
}\
\
// getPhase returns the most recent phase label, or "" if none was set.\
func (s *statsCollector) getPhase() string {\
    if p := s.phase.Load(); p != nil {\
        return *p\
    }\
    return ""\
}\
\
// setSideDisplay attaches a human-readable name (typically the URL\
// hostname) to a side. Called once during session setup so the live\
// renderer can print "github.com → ..." instead of "source: ...".\
func (s *statsCollector) setSideDisplay(label, display string) {\
    if display == "" {\
        return\
    }\
    sc := s.side(label)\
    s.sidesMu.Lock()\
    sc.display = display\
    s.sidesMu.Unlock()\
}\
\
// liveSides returns a snapshot of per-side byte totals for live rendering.\
// ActiveNanos and IdleNanos are computed against time.Now() at snapshot\
// time so callers do not need to know the collector's start instant.\
func (s *statsCollector) liveSides() []SideBytes {\
    s.sidesMu.RLock()\
    defer s.sidesMu.RUnlock()\
    startNanos := s.startedAt.UnixNano()\
    nowNanos := time.Now().UnixNano()\
    out := make([]SideBytes, 0, len(s.sides))\
    for label, sc := range s.sides {\
        bytes := sc.bytes.Load()\
        lastByte := sc.lastByteAt.Load()\
        var activeNanos, idleNanos int64\
        if lastByte > 0 {\
            activeNanos = lastByte - startNanos\
            if activeNanos < 0 {\
                activeNanos = 0\
            }\
            idleNanos = nowNanos - lastByte\
            if idleNanos < 0 {\
                idleNanos = 0\
            }\
        }\
        out = append(out, SideBytes{\
            Label:       label,\
            Bytes:       bytes,\
            Display:     sc.display,\
            ActiveNanos: activeNanos,\
            IdleNanos:   idleNanos,\
        })\
    }\
    return out\
}\
\
func (s *statsCollector) snapshot() Stats {\
    s.mu.Lock()\
    defer s.mu.Unlock()\
    out := Stats{Enabled: s.enabled, Items: make(map[string]*ServiceStats, len(s.items))}\
    for key, item := range s.items {\
        copyItem := *item\
        out.Items[key] = &copyItem\
    }\
    s.mu.Unlock()\
    if s.enabled {\
        out.Sides = s.liveSides()\
        out.ElapsedNanos = time.Since(s.startedAt).Nanoseconds()\
    }\
    return out\
}\
\
// countingRoundTripper wraps an HTTP transport to record transfer stats.\
// countingRoundTripper wraps an HTTP transport to record transfer stats\
// and feed per-side byte counters consumed by the live progress ticker.\
type countingRoundTripper struct {\
    base  http.RoundTripper\
    label string\
1 unmodified line\
\
}\
\
func (rt *countingRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) {\
    side := rt.stats.side(rt.label)\
\
    // Wrap the request body so upload bytes (receive-pack POSTs) feed\
    // the per-side counter as they stream out, not just on round-trip\
    // completion. Skipped when there is no body (most GETs).\
    if req.Body != nil {\
        req.Body = &countingReadCloser{ReadCloser: req.Body, side: side}\
    }\
\
    res, err := rt.base.RoundTrip(req)\
    if err != nil {\
        return nil, fmt.Errorf("round trip: %w", err)\
14 unmodified lines\
\
    res.Body = &countingReadCloser{\
        ReadCloser: res.Body,\
        side:       side,\
        onClose: func(n int64) {\
            rt.stats.recordRoundTrip(name, requestBytes, n)\
        },\
5 unmodified lines\
\
    io.ReadCloser\
\
    n       int64\
    side    *sideCounter\
    onClose func(int64)\
}\
\
func (c *countingReadCloser) Read(p []byte) (int, error) {\
    n, err := c.ReadCloser.Read(p)\
    c.n += int64(n)\
    if n > 0 {\
        c.n += int64(n)\
        if c.side != nil {\
            c.side.bytes.Add(int64(n))\
            c.side.lastByteAt.Store(time.Now().UnixNano())\
        }\
    }\
    return n, err //nolint:wrapcheck // Read must preserve io.EOF for io.Reader contract\
}\
```\
\
Minternal/syncer/stats.go+173/-10\
\
```\
7 unmodified lines\
\
8\
9\
10\
11\
12\
13\
14\
15\
16\
17\
18\
19\
20\
21\
53 unmodified lines\
\
75\
76\
77\
78\
79\
80\
81\
1 unmodified line\
\
83\
84\
85\
86\
87\
88\
89\
90\
91\
92\
192 unmodified lines\
\
285\
286\
287\
288\
289\
290\
291\
292\
293\
294\
295\
296\
297\
298\
299\
300\
301\
302\
303\
304\
305\
306\
307\
308\
309\
310\
311\
312\
313\
314\
315\
316\
317\
318\
319\
320\
321\
322\
323\
324\
325\
326\
21 unmodified lines\
\
348\
349\
350\
351\
352\
353\
354\
355\
356\
357\
358\
359\
360\
361\
362\
363\
364\
365\
366\
367\
368\
369\
370\
371\
64 unmodified lines\
\
436\
437\
438\
439\
440\
441\
442\
443\
444\
445\
446\
447\
448\
449\
450\
451\
452\
453\
454\
455\
456\
457\
458\
459\
460\
461\
462\
463\
38 unmodified lines\
\
502\
503\
504\
430\
505\
506\
507\
508\
2 unmodified lines\
\
511\
512\
513\
514\
515\
516\
517\
9 unmodified lines\
\
527\
528\
529\
530\
531\
532\
533\
17 unmodified lines\
\
551\
552\
553\
554\
555\
556\
557\
558\
559\
560\
561\
562\
563\
564\
565\
566\
567\
568\
569\
570\
571\
572\
573\
574\
575\
576\
577\
578\
5 unmodified lines\
\
584\
585\
586\
587\
588\
589\
590\
271 unmodified lines\
\
862\
863\
864\
865\
866\
867\
868\
19 unmodified lines\
\
888\
889\
890\
891\
892\
893\
894\
7 unmodified lines\
\
902\
903\
904\
905\
906\
907\
908\
37 unmodified lines\
\
946\
947\
948\
949\
950\
951\
952\
953\
\
7 unmodified lines\
\
    "encoding/json"\
    "errors"\
    "fmt"\
    "io"\
    "log/slog"\
    "net/http"\
    "net/url"\
    "os"\
    "sort"\
    "strings"\
    "time"\
\
    git "github.com/go-git/go-git/v6"\
    "github.com/go-git/go-git/v6/plumbing"\
53 unmodified lines\
\
    Verbose                bool\
    ShowStats              bool\
    MeasureMemory          bool\
    Progress               bool\
    Mode                   string\
    Force                  bool\
    Prune                  bool\
1 unmodified line\
\
    TargetMaxPackBytes     int64\
    MaterializedMaxObjects int\
    ProtocolMode           string\
\
    // progressOut overrides the writer used by the live progress ticker.\
    // Defaults to os.Stderr when nil. Exposed for tests.\
    progressOut io.Writer\
}\
\
// Re-export types from planner for CLI compatibility.\
192 unmodified lines\
\
            item.Name, item.Requests, item.RequestBytes, item.ResponseBytes, item.Wants, item.Haves, item.Commands,\
        ))\
    }\
    if line := throughputLine(s); line != "" {\
        lines = append(lines, line)\
    }\
    return lines\
}\
\
// throughputLine renders a one-line per-side throughput summary using\
// the wall-clock window the stats collector observed. Returns "" when\
// no side bytes were recorded so the line stays out of the way for\
// metadata-only operations like probe.\
func throughputLine(s Stats) string {\
    if len(s.Sides) == 0 || s.ElapsedNanos <= 0 {\
        return ""\
    }\
    sides := make([]SideBytes, 0, len(s.Sides))\
    for _, side := range s.Sides {\
        if side.Bytes <= 0 {\
            continue\
        }\
        sides = append(sides, side)\
    }\
    if len(sides) == 0 {\
        return ""\
    }\
    sort.Slice(sides, func(i, j int) bool { return sides[i].Label < sides[j].Label })\
    dur := time.Duration(s.ElapsedNanos)\
    parts := make([]string, 0, len(sides))\
    for _, side := range sides {\
        // End-of-run line uses the active-window average; passing 0\
        // for instant rate makes formatSide skip the sliding window\
        // and fall back to ActiveNanos-based formatting.\
        parts = append(parts, formatSide(side, dur, 0, true))\
    }\
    return "throughput: " + strings.Join(parts, sideSeparator)\
}\
\
func measurementLine(m Measurement) []string {\
    if !m.Enabled {\
        return nil\
21 unmodified lines\
\
    if err != nil {\
        return nil, fmt.Errorf("resolve auth: %w", err)\
    }\
    stats.setSideDisplay(label, hostnameFromURL(raw.URL))\
    client := instrumentHTTPClient(httpClient, raw.SkipTLSVerify, label, stats)\
    conn := gitproto.NewConnWithHTTPClient(ep, label, authMethod, client)\
    conn.FollowInfoRefsRedirect = raw.FollowInfoRefsRedirect\
    return conn, nil\
}\
\
// hostnameFromURL returns the host portion of an endpoint URL, used to\
// label sides in progress and throughput output. Returns "" for malformed\
// URLs so callers can fall back to the internal label.\
func hostnameFromURL(raw string) string {\
    u, err := url.Parse(raw)\
    if err != nil {\
        return ""\
    }\
    return u.Hostname()\
}\
\
func instrumentHTTPClient(base *http.Client, skipTLS bool, label string, stats *statsCollector) *http.Client {\
    if base == nil {\
        base = &http.Client{Transport: gitproto.NewHTTPTransport(skipTLS)}\
64 unmodified lines\
\
    sourceRefMap    map[plumbing.ReferenceName]plumbing.Hash\
    target          *targetSession\
    measurementDone func() Measurement\
    progress        *progressReporter\
}\
\
// finish releases any resources owned by the session — currently the live\
// progress ticker. Idempotent and safe to call from defer in callers that\
// also produce results in the happy path.\
func (s *syncSession) finish() {\
    if s.progress != nil {\
        s.progress.terminate()\
    }\
}\
\
// notice surfaces a one-line human-readable event during a sync. When\
// the live progress ticker is active it prints above the current frame;\
// otherwise it falls back to plain stderr so the message is still seen\
// when --progress is off or the destination is not a TTY.\
func (s *syncSession) notice(msg string) {\
    if s.progress != nil {\
        s.progress.notify(msg)\
        return\
    }\
    fmt.Fprintln(os.Stderr, msg)\
}\
\
type targetSession struct {\
38 unmodified lines\
\
        measurementDone: startMeasurement(cfg.MeasureMemory),\
    }\
    if cfg.Verbose {\
        s.logger = slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{\
        s.logger = slog.New(slog.NewTextHandler(&sessionStderr{s: s}, &slog.HandlerOptions{\
            Level: slog.LevelInfo,\
        }))\
    }\
2 unmodified lines\
\
    if err != nil {\
        return nil, fmt.Errorf("create source transport: %w", err)\
    }\
    s.sourceConn.ProgressOut = &sessionStderr{s: s}\
\
    refPrefixes := planner.RefPrefixes(cfg.Mappings, cfg.IncludeTags)\
    sourceRefs, sourceService, err := gitproto.ListSourceRefs(ctx, s.sourceConn, cfg.ProtocolMode, refPrefixes)\
9 unmodified lines\
\
        if err != nil {\
            return nil, fmt.Errorf("create target transport: %w", err)\
        }\
        targetConn.ProgressOut = &sessionStderr{s: s}\
        targetAdv, err := gitproto.AdvertisedRefsV1(ctx, targetConn, transport.ReceivePackService)\
        if err != nil {\
            return nil, fmt.Errorf("list target refs: %w", err)\
17 unmodified lines\
\
        }\
    }\
\
    // Start the live progress ticker only after auth resolution and the\
    // initial ref-listing round trips have completed. The auth path may\
    // shell out to `git credential fill`, which inherits our stderr and\
    // can prompt the user; an interactive prompt and a '\r'-redrawing\
    // ticker writing to the same tty would clobber each other. Deferring\
    // the ticker until newSession returns guarantees no concurrent writer\
    // is active when prompts happen and also avoids leaking a goroutine\
    // when newSession fails partway through setup.\
    if cfg.Progress {\
        out := cfg.progressOut\
        if out == nil {\
            out = os.Stderr\
        }\
        // Render only when the destination is a real terminal. Pipes,\
        // log files, and CI captures get nothing rather than a flood\
        // of '\r'-prefixed control sequences.\
        if out != os.Stderr || stderrIsTTY() {\
            s.progress = newProgressReporter(out, s.stats, 0)\
            go s.progress.run()\
        }\
    }\
\
    return s, nil\
}\
\
5 unmodified lines\
\
    if err != nil {\
        return Result{}, err\
    }\
    defer s.finish()\
    if s.cfg.Mode == modeReplicate {\
        return s.runReplicate(ctx)\
    }\
271 unmodified lines\
\
    if err != nil {\
        return Result{}, err\
    }\
    defer s.finish()\
\
    desiredRefs, _, err := planner.BuildDesiredRefs(s.sourceRefMap, planConfig(cfg))\
    if err != nil {\
19 unmodified lines\
\
    if err != nil {\
        return ProbeResult{}, err\
    }\
    defer s.finish()\
    return s.newProbeResult(), nil\
}\
\
7 unmodified lines\
\
    if err != nil {\
        return FetchResult{}, err\
    }\
    defer s.finish()\
\
    repo, err := git.Init(memory.NewStorage(), nil)\
    if err != nil {\
37 unmodified lines\
\
        SourceHeadTarget: s.sourceService.HeadTarget,\
        MaxPackBytes:     s.cfg.MaxPackBytes, TargetMaxPack: s.cfg.TargetMaxPackBytes,\
        Verbose: s.cfg.Verbose, Logger: s.logger,\
        OnPhase:  s.stats.setPhase,\
        OnNotice: s.notice,\
    }, relayReason)\
    if err != nil {\
        return Result{}, fmt.Errorf("bootstrap execute: %w", err)\
```\
\
Minternal/syncer/syncer.go+106/-1\
\
```\
1\
2\
3\
4\
5\
6\
7\
8\
9\
10\
11\
12\
13\
14\
15\
16\
17\
18\
19\
20\
21\
22\
23\
24\
25\
26\
27\
28\
29\
30\
31\
32\
33\
34\
35\
36\
37\
38\
39\
40\
41\
42\
43\
44\
45\
46\
47\
48\
49\
50\
51\
52\
53\
54\
55\
56\
57\
58\
59\
60\
61\
62\
63\
64\
65\
66\
67\
68\
69\
70\
71\
72\
73\
74\
75\
76\
77\
78\
79\
80\
81\
82\
83\
84\
85\
86\
87\
88\
89\
90\
91\
92\
93\
94\
95\
96\
97\
98\
99\
100\
101\
102\
103\
104\
105\
106\
107\
108\
109\
110\
111\
112\
113\
\
package syncer\
\
import (\
    "bytes"\
    "context"\
    "strings"\
    "testing"\
    "time"\
\
    git "github.com/go-git/go-git/v6"\
    "github.com/go-git/go-git/v6/storage/memory"\
)\
\
// TestBootstrap_ThroughputLineAndProgressTicker verifies end-to-end that\
// running a real bootstrap with --stats and --progress populates per-side\
// byte counters, prints a throughput line in the human output, and emits\
// at least one render to the progress writer when progressOut is wired.\
func TestBootstrap_ThroughputLineAndProgressTicker(t *testing.T) {\
    sourceRepo, sourceFS := newSourceRepo(t)\
    makeCommits(t, sourceRepo, sourceFS, 4)\
\
    targetRepo, err := git.Init(memory.NewStorage())\
    if err != nil {\
        t.Fatalf("init target repo: %v", err)\
    }\
\
    sourceServer := newSmartHTTPRepoServerV2(t, sourceRepo)\
    targetServer := newSmartHTTPRepoServer(t, targetRepo)\
    defer sourceServer.Close()\
    defer targetServer.Close()\
\
    var progressBuf bytes.Buffer\
    cfg := Config{\
        Source:       Endpoint{URL: sourceServer.RepoURL()},\
        Target:       Endpoint{URL: targetServer.RepoURL()},\
        ProtocolMode: protocolModeAuto,\
        ShowStats:    true,\
        Progress:     true,\
        progressOut:  &progressBuf,\
    }\
\
    result, err := Bootstrap(context.Background(), cfg)\
    if err != nil {\
        t.Fatalf("bootstrap failed: %v", err)\
    }\
    if result.Pushed != 1 {\
        t.Fatalf("expected one push, got %+v", result)\
    }\
\
    // Per-side counters should record bytes on both sides.\
    sides := map[string]int64{}\
    for _, side := range result.Stats.Sides {\
        sides[side.Label] = side.Bytes\
    }\
    if sides["source"] <= 0 {\
        t.Errorf("expected non-zero source bytes, got %d", sides["source"])\
    }\
    if sides["target"] <= 0 {\
        t.Errorf("expected non-zero target bytes, got %d", sides["target"])\
    }\
    if result.Stats.ElapsedNanos <= 0 {\
        t.Errorf("expected non-zero elapsed nanos, got %d", result.Stats.ElapsedNanos)\
    }\
\
    // Per-side displays should reflect the test server hostnames.\
    displays := map[string]string{}\
    for _, side := range result.Stats.Sides {\
        displays[side.Label] = side.Display\
    }\
    if displays["source"] == "" || displays["target"] == "" {\
        t.Errorf("expected non-empty display for both sides, got %+v", result.Stats.Sides)\
    }\
\
    // Human-formatted output should include the new throughput line.\
    out := strings.Join(result.Lines(), "\n")\
    if !strings.Contains(out, "throughput: ") {\
        t.Errorf("missing throughput line in output:\n%s", out)\
    }\
    if !strings.Contains(out, displays["source"]) || !strings.Contains(out, displays["target"]) {\
        t.Errorf("throughput line should mention both hostnames:\n%s", out)\
    }\
    if !strings.Contains(out, "→") || !strings.Contains(out, "│") {\
        t.Errorf("throughput line should use arrow + vertical bar separator:\n%s", out)\
    }\
\
    // Progress writer should have received at least one frame. The bootstrap\
    // is very fast so we may only get the final terminate() render.\
    if progressBuf.Len() == 0 {\
        // Allow a brief grace window in case the goroutine had not flushed\
        // yet. terminate() inside finish() blocks on render so this should\
        // be unnecessary, but keep a small fallback for slow CI.\
        time.Sleep(50 * time.Millisecond)\
    }\
    if progressBuf.Len() == 0 {\
        t.Errorf("expected progress reporter to write at least one frame")\
    }\
    if !bytes.Contains(progressBuf.Bytes(), []byte(displays["source"])) {\
        t.Errorf("progress output should mention source host %q: %q",\
            displays["source"], progressBuf.String())\
    }\
    if !bytes.Contains(progressBuf.Bytes(), []byte(displays["target"])) {\
        t.Errorf("progress output should mention target host %q: %q",\
            displays["target"], progressBuf.String())\
    }\
\
    // Print samples for human inspection when running with -v.\
    for _, line := range result.Lines() {\
        if strings.HasPrefix(line, "stats:") || strings.HasPrefix(line, "throughput:") {\
            t.Logf("output line: %s", line)\
        }\
    }\
    t.Logf("progress writer:\n%q", progressBuf.String())\
}\
```\
\
Ainternal/syncer/throughput\_test.go+113\
\
```\
51 unmodified lines\
\
52\
53\
54\
55\
56\
57\
58\
59\
60\
61\
62\
63\
56\
57\
64\
65\
66\
67\
68\
69\
70\
126 unmodified lines\
\
197\
198\
199\
200\
201\
202\
203\
204\
205\
206\
207\
208\
209\
210\
211\
212\
213\
214\
215\
\
51 unmodified lines\
\
    Commands      int    `json:"commands"`\
}\
\
type SideBytes struct {\
    Label       string `json:"label"`\
    Bytes       int64  `json:"bytes"`\
    Display     string `json:"display,omitempty"`\
    ActiveNanos int64  `json:"activeNanos,omitempty"`\
    IdleNanos   int64  `json:"idleNanos,omitempty"`\
}\
\
type Stats struct {\
    Enabled bool                     `json:"enabled"`\
    Items   map[string]*ServiceStats `json:"items"`\
    Enabled      bool                     `json:"enabled"`\
    Items        map[string]*ServiceStats `json:"items"`\
    Sides        []SideBytes              `json:"sides,omitempty"`\
    ElapsedNanos int64                    `json:"elapsedNanos,omitempty"`\
}\
\
type Measurement struct {\
126 unmodified lines\
\
            Commands:      copyItem.Commands,\
        }\
    }\
    if len(stats.Sides) > 0 {\
        out.Sides = make([]SideBytes, 0, len(stats.Sides))\
        for _, side := range stats.Sides {\
            out.Sides = append(out.Sides, SideBytes{\
                Label:       side.Label,\
                Bytes:       side.Bytes,\
                Display:     side.Display,\
                ActiveNanos: side.ActiveNanos,\
                IdleNanos:   side.IdleNanos,\
            })\
        }\
    }\
    out.ElapsedNanos = stats.ElapsedNanos\
    return out\
}\
```\
\
Minternalbridge/model.go+25/-2\
\
```\
37 unmodified lines\
\
38\
39\
40\
41\
42\
43\
44\
125 unmodified lines\
\
170\
171\
172\
173\
174\
175\
176\
30 unmodified lines\
\
207\
208\
209\
210\
211\
212\
213\
23 unmodified lines\
\
237\
238\
239\
240\
241\
242\
243\
13 unmodified lines\
\
257\
258\
259\
260\
261\
262\
263\
\
37 unmodified lines\
\
    CollectStats           bool  `json:"collectStats"`\
    MeasureMemory          bool  `json:"measureMemory"`\
    Verbose                bool  `json:"verbose"`\
    Progress               bool  `json:"progress"`\
    MaxPackBytes           int64 `json:"maxPackBytes"`\
    TargetMaxPackBytes     int64 `json:"targetMaxPackBytes"`\
    MaterializedMaxObjects int   `json:"materializedMaxObjects"`\
125 unmodified lines\
\
        IncludeTags:   req.IncludeTags,\
        ShowStats:     req.Options.CollectStats,\
        MeasureMemory: req.Options.MeasureMemory,\
        Progress:      req.Options.Progress,\
        ProtocolMode:  protocolString(req.Protocol),\
        Verbose:       req.Options.Verbose,\
    }\
30 unmodified lines\
\
        DryRun:                 req.DryRun,\
        ShowStats:              req.Options.CollectStats,\
        MeasureMemory:          req.Options.MeasureMemory,\
        Progress:               req.Options.Progress,\
        Mode:                   operationModeString(req.Policy.Mode),\
        Force:                  req.Policy.Force,\
        Prune:                  req.Policy.Prune,\
23 unmodified lines\
\
        IncludeTags:        req.IncludeTags,\
        ShowStats:          req.Options.CollectStats,\
        MeasureMemory:      req.Options.MeasureMemory,\
        Progress:           req.Options.Progress,\
        MaxPackBytes:       req.Options.MaxPackBytes,\
        TargetMaxPackBytes: req.Options.TargetMaxPackBytes,\
        ProtocolMode:       protocolString(req.Protocol),\
13 unmodified lines\
\
        IncludeTags:   req.IncludeTags,\
        ShowStats:     req.Options.CollectStats,\
        MeasureMemory: req.Options.MeasureMemory,\
        Progress:      req.Options.Progress,\
        ProtocolMode:  protocolString(req.Protocol),\
        Verbose:       req.Options.Verbose,\
    }, nil\
```\
\
Munstable/client.go+5