Stream sideband progress to stderr when -v is set · Entire
Stream sideband progress to stderr when -v is set
3b98784→main·
Soph·3mo ago·6 files·+178 added/-41 removed
Long replicate and sync runs against large repos were silent between bootstrap batch checkpoints because both the source upload-pack and the target receive-pack sideband progress channels were being discarded.
Changes in internal/gitproto:
- push.go: sendReceivePack now takes a verbose flag and wires the sideband demuxer's Progress field to stderr (prefixed "target: ") when set. Added a prefixedLineWriter that splits on both '\n' and '\r' so in-place git progress updates ("Resolving deltas: 12%\r") remain readable when prefixed. progressSink returns nil for non-verbose so the demuxer allocates nothing on the hot path.
- fetch.go / refs.go: RefService gains a Verbose field. When true, fetchToStoreV1/V2, fetchPackV1/V2, and buildV1UploadPackBody stop asking the source to suppress progress (drop the "no-progress" upload-request capability and the "no-progress" v2 fetch arg) and wire the sideband demuxer's Progress to stderr (prefixed "source: ").
- syncer.go: newSession propagates cfg.Verbose to sourceService.Verbose right after constructing it. Target-side verbose already flowed through gitproto.NewPusher.
New tests:
- TestPrefixedLineWriter covers line splitting on '\n' and '\r', mid-line writes that don't emit a trailing prefix, and empty writes.
- TestProgressSinkNilWhenNotVerbose locks in the "no allocation when quiet" contract.
Existing private fetch tests updated to pass the new verbose parameter. Commit-graph fetches (FetchCommitGraph) pass verbose=false because they are short and not user-facing.
With -v, a replicate run now looks like:
source: Enumerating objects: 120000, done. source: Counting objects: 100% (120000/120000), done. source: Compressing objects: 37% (44400/120000) target: Resolving deltas: 58% (69600/120000) target: Updating references: 100% (61/61), done.
Co-Authored-By: Claude Opus 4.6 (1M context) noreply@anthropic.com
Sessions
cdd8c35ae583View transcript
Changes
6
internal
gitproto
Mfetch.go+32/-17
Mfetch_test.go+14/-14
Mpush.go+64/-9
Mpush_test.go+63/-1
Mrefs.go+4
syncer
Msyncer.go+1
`` 60 unmodified lines
61 62 63 64 64 65 66 66 67 68 69 9 unmodified lines
79 80 81 82 82 83 84 84 85 86 87 30 unmodified lines
118 119 120 121 121 122 123 124 125 17 unmodified lines
143 144 145 146 147 148 149 6 unmodified lines
156 157 158 157 159 160 161 162 163 164 165 11 unmodified lines
177 178 179 175 180 181 182 183 2 unmodified lines
186 187 188 189 190 191 192 6 unmodified lines
199 200 201 196 202 203 204 205 206 207 208 14 unmodified lines
223 224 225 217 226 227 228 229 1 unmodified line
231 232 233 225 234 235 236 237 13 unmodified lines
251 252 253 254 255 256 257 6 unmodified lines
264 265 266 257 267 268 269 270 12 unmodified lines
283 284 285 286 287 288 277 289 290 291 292 15 unmodified lines
308 309 310 311 312 313 314 3 unmodified lines
318 319 320 308 321 322 323 324 30 unmodified lines
355 356 357 358 359 346 360 361 362 363 12 unmodified lines
376 377 378 365 379 380 381 382 3 unmodified lines
386 387 388 389 390 376 391 392 393 394 13 unmodified lines
408 409 410 396 411 412 413 414
60 unmodified lines
) error { switch s.Protocol { case "v2": return fetchToStoreV2(ctx, store, conn, s.V2Caps, desired, targetRefs) return fetchToStoreV2(ctx, store, conn, s.V2Caps, desired, targetRefs, s.Verbose) case "v1": return fetchToStoreV1(ctx, store, conn, s.V1Adv, desired, targetRefs) return fetchToStoreV1(ctx, store, conn, s.V1Adv, desired, targetRefs, s.Verbose) default: return fmt.Errorf("unsupported source protocol %q", s.Protocol) } 9 unmodified lines
) (io.ReadCloser, error) { switch s.Protocol { case "v2": return fetchPackV2(ctx, conn, s.V2Caps, desired, targetRefs) return fetchPackV2(ctx, conn, s.V2Caps, desired, targetRefs, s.Verbose) case "v1": return fetchPackV1(ctx, conn, s.V1Adv, desired, targetRefs) return fetchPackV1(ctx, conn, s.V1Adv, desired, targetRefs, s.Verbose) default: return nil, fmt.Errorf("unsupported source protocol %q", s.Protocol) } 30 unmodified lines
return err }
defer ioutil.CheckClose(reader, &err) return storeV2FetchPack(store, reader) // Commit-graph fetches are short and not user-facing; skip progress. return storeV2FetchPack(store, reader, false) }
// Capabilities returns the sorted capability list for display. 17 unmodified lines
caps *V2Capabilities, desired map[plumbing.ReferenceName]DesiredRef, targetRefs map[plumbing.ReferenceName]plumbing.Hash, verbose bool, ) error { wants := collectWants(desired) haves := SortedUniqueHashes(refValues(targetRefs)) 6 unmodified lines
// self-contained so callers (e.g. replicate) can forward it to // receive-pack servers that may advertise "no-thin". See // planner.SupportsReplicateRelay for the matching invariant. cmdArgs = append(cmdArgs, "ofs-delta", "no-progress") cmdArgs = append(cmdArgs, "ofs-delta") if !verbose { cmdArgs = append(cmdArgs, "no-progress") } for _, h := range wants { cmdArgs = append(cmdArgs, "want " + h.String()) } 11 unmodified lines
return err }
defer ioutil.CheckClose(reader, &err) return storeV2FetchPack(store, reader) return storeV2FetchPack(store, reader, verbose) }
func fetchPackV2( 2 unmodified lines
caps *V2Capabilities, desired map[plumbing.ReferenceName]DesiredRef, targetRefs map[plumbing.ReferenceName]plumbing.Hash, verbose bool, ) (io.ReadCloser, error) { wants := collectWants(desired) haves := SortedUniqueHashes(refValues(targetRefs)) 6 unmodified lines
// self-contained so callers (e.g. replicate) can forward it to // receive-pack servers that may advertise "no-thin". See // planner.SupportsReplicateRelay for the matching invariant. cmdArgs = append(cmdArgs, "ofs-delta", "no-progress") cmdArgs = append(cmdArgs, "ofs-delta") if !verbose { cmdArgs = append(cmdArgs, "no-progress") } // Only request include-tag if the server supports it (issue #6). if hasTag(desired) && caps.FetchSupports("include-tag") { cmdArgs = append(cmdArgs, "include-tag") }14 unmodified lines
if err != nil { return nil, err } packStream, err := openV2PackStream(reader) packStream, err := openV2PackStream(reader, verbose) if err != nil { _ = reader.Close() return nil, err }1 unmodified line
return packStream, nil }
func storeV2FetchPack(store storer.Storer, r io.Reader) error { func storeV2FetchPack(store storer.Storer, r io.Reader, verbose bool) error { reader := NewPacketReader(r) for { kind, payload, err := reader.ReadPacket() 13 unmodified lines
switch line { case "packfile\n": demux := sideband.NewDemuxer(sideband.Sideband64k, reader.BufReader()) demux.Progress = progressSink(verbose, "source: ") return packfile.UpdateObjectStorage(store, demux) case "acknowledgments\n", "shallow-info\n": if err := SkipSection(reader); err != nil { 6 unmodified lines
} }
func openV2PackStream(body io.ReadCloser) (io.ReadCloser, error) { func openV2PackStream(body io.ReadCloser, verbose bool) (io.ReadCloser, error) { reader := NewPacketReader(body) for { kind, payload, err := reader.ReadPacket() 12 unmodified lines
line := string(payload) switch line { case "packfile\n": demux := sideband.NewDemuxer(sideband.Sideband64k, reader.BufReader()) demux.Progress = progressSink(verbose, "source: ") return &wrappedRC{ Reader: sideband.NewDemuxer(sideband.Sideband64k, reader.BufReader()), Reader: demux, Closer: body, }, nil case "acknowledgments\n", "shallow-info\n": 15 unmodified lines
desired map[plumbing.ReferenceName]DesiredRef, targetRefs map[plumbing.ReferenceName]plumbing.Hash, includeTags bool, verbose bool, ) ([]byte, *capability.List, error) { wants := collectWants(desired) haves := SortedUniqueHashes(refValues(targetRefs)) 3 unmodified lines
req := packp.NewUploadRequest() req.Wants = wants if adv.Capabilities.Supports(capability.NoProgress) { if !verbose && adv.Capabilities.Supports(capability.NoProgress) { _ = req.Capabilities.Set(capability.NoProgress) } if includeTags && adv.Capabilities.Supports(capability.IncludeTag) { 30 unmodified lines
adv *packp.AdvRefs, desired map[plumbing.ReferenceName]DesiredRef, targetRefs map[plumbing.ReferenceName]plumbing.Hash, verbose bool, ) error { body, caps, err := buildV1UploadPackBody(adv, desired, targetRefs, hasTag(desired)) body, caps, err := buildV1UploadPackBody(adv, desired, targetRefs, hasTag(desired), verbose) if err != nil { return err } 12 unmodified lines
if drainErr := drainTrailingNAKs(buffered); drainErr != nil { return fmt.Errorf("drain server response: %w", drainErr) } sbReader := buildSidebandReader(caps, buffered, nil) sbReader := buildSidebandReader(caps, buffered, progressSink(verbose, "source: ")) return packfile.UpdateObjectStorage(store, sbReader) }
3 unmodified lines
adv *packp.AdvRefs, desired map[plumbing.ReferenceName]DesiredRef, targetRefs map[plumbing.ReferenceName]plumbing.Hash, verbose bool, ) (io.ReadCloser, error) { body, caps, err := buildV1UploadPackBody(adv, desired, targetRefs, hasTag(desired)) body, caps, err := buildV1UploadPackBody(adv, desired, targetRefs, hasTag(desired), verbose) if err != nil { return nil, err } 13 unmodified lines
return &wrappedRC{ Reader: buildSidebandReader(caps, buffered, nil), Reader: buildSidebandReader(caps, buffered, progressSink(verbose, "source: ")), Closer: reader, }, nil } ``
Minternal/gitproto/fetch.go+32/-17
`` 330 unmodified lines
331 332 333 334 334 335 336 337 45 unmodified lines
383 384 385 386 386 387 388 389 45 unmodified lines
435 436 437 438 438 439 440 441 43 unmodified lines
485 486 487 488 488 489 490 491 38 unmodified lines
530 531 532 533 533 534 535 536 39 unmodified lines
576 577 578 579 579 580 581 582 28 unmodified lines
611 612 613 614 614 615 616 617 36 unmodified lines
654 655 656 657 657 658 659 660 40 unmodified lines
701 702 703 704 704 705 706 707 39 unmodified lines
747 748 749 750 750 751 752 753 39 unmodified lines
793 794 795 796 796 797 798 799 39 unmodified lines
839 840 841 842 842 843 844 845 36 unmodified lines
882 883 884 885 885 886 887 888 4 unmodified lines
893 894 895 896 896 897 898 899
330 unmodified lines
ctx, cancel := context.WithCancel(context.Background()) done := make(chan error, 1) go func() { _, err := fetchPackV1(ctx, conn, adv, desired, nil) _, err := fetchPackV1(ctx, conn, adv, desired, nil, false) done <- err }()
45 unmodified lines
ctx, cancel := context.WithCancel(context.Background()) done := make(chan error, 1) go func() { _, err := fetchPackV2(ctx, conn, caps, desired, nil) _, err := fetchPackV2(ctx, conn, caps, desired, nil, false) done <- err }()
45 unmodified lines
ctx, cancel := context.WithCancel(context.Background()) done := make(chan error, 1) go func() { done <- fetchToStoreV2(ctx, memory.NewStorage(), conn, caps, desired, nil) done <- fetchToStoreV2(ctx, memory.NewStorage(), conn, caps, desired, nil, false) }()
select { 43 unmodified lines
}, }
err = fetchToStoreV2(context.Background(), memory.NewStorage(), conn, caps, desired, nil) err = fetchToStoreV2(context.Background(), memory.NewStorage(), conn, caps, desired, nil, false) if err == nil { t.Fatal("expected decode error") } 38 unmodified lines
func TestBuildV1UploadPackBodyEmptyWantSet(t *testing.T) { adv := packp.NewAdvRefs() _, _, err := buildV1UploadPackBody(adv, nil, nil, false) _, _, err := buildV1UploadPackBody(adv, nil, nil, false, false) if !errors.Is(err, git.NoErrAlreadyUpToDate) { t.Fatalf("expected NoErrAlreadyUpToDate, got %v", err) } } ``
Minternal/gitproto/fetch_test.go+14/-14
`` 95 unmodified lines
96 97 98 99 100 101 102 9 unmodified lines
112 113 114 114 115 116 117 116 117 118 119 118 119 120 121 122 123 124 125 126 127 128 129 23 unmodified lines
153 154 155 149 156 157 158 159 9 unmodified lines
169 170 171 165 172 173 174 175 24 unmodified lines
200 201 202 196 203 204 205 206 13 unmodified lines
220 221 222 216 223 224 225 226 2 unmodified lines
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
95 unmodified lines
conn *Conn, req *packp.UpdateRequests, packData io.Reader, verbose bool, ) error { var header bytes.Buffer if err := req.Encode(&header); err != nil { 9 unmodified lines
}
def defer reader.Close()
// Unwrap sideband if negotiated. // Unwrap sideband if negotiated; stream server-side progress to stderr // when verbose so long-running pushes show "Resolving deltas ..." etc. var respReader io.Reader = reader if req.Capabilities.Supports(capability.Sideband64k) { respReader = sideband.NewDemuxer(sideband.Sideband64k, reader) } else if req.Capabilities.Supports(capability.Sideband) { respReader = sideband.NewDemuxer(sideband.Sideband, reader) switch { case req.Capabilities.Supports(capability.Sideband64k): dem := sideband.NewDemuxer(sideband.Sideband64k, reader) dem.Progress = progressSink(verbose, "target: ") respReader = dem case req.Capabilities.Supports(capability.Sideband): dem := sideband.NewDemuxer(sideband.Sideband, reader) dem.Progress = progressSink(verbose, "target: ") respReader = dem }
if req.Capabilities.Supports(capability.ReportStatus) { 23 unmodified lines
return err } if !hasUpdates { return sendReceivePack(ctx, conn, req, nil) return sendReceivePack(ctx, conn, req, nil, verbose) }
useRefDeltas := !adv.Capabilities.Supports(capability.OFSDelta) 9 unmodified lines
done <- pw.Close() }()
err = sendReceivePack(ctx, conn, req, pr) err = sendReceivePack(ctx, conn, req, pr, verbose) _ = pr.Close() encodeErr := <-done if err != nil { 24 unmodified lines
return err }
err = sendReceivePack(ctx, conn, req, pack) err = sendReceivePack(ctx, conn, req, pack, verbose) closeErr := pack.Close() if err != nil { return err }13 unmodified lines
if err != nil { return err } return sendReceivePack(ctx, conn, req, nil) return sendReceivePack(ctx, conn, req, nil, verbose) }
func progressWriter(verbose bool) io.Writer { 2 unmodified lines
}
return os.Stderr }
// 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 { if !verbose { return nil } return &prefixedLineWriter{w: os.Stderr, prefix: prefix, atLineStart: true} }
// prefixedLineWriter prepends a fixed prefix to each line of input written // to the wrapped writer. Git sideband progress arrives as chunks that may // contain '\n' between full lines or '\r' for in-place updates ("Resolving // deltas: 12%\r"); both are treated as line terminators so the next chunk // gets a fresh prefix. type prefixedLineWriter struct { w io.Writer prefix string atLineStart bool }
func (p *prefixedLineWriter) Write(b []byte) (int, error) { consumed := 0 for len(b) > 0 { if p.atLineStart { if _, err := io.WriteString(p.w, p.prefix); err != nil { return consumed, err } p.atLineStart = false } i := bytes.IndexAny(b, "\r\n") var chunk []byte if i < 0 { chunk = b } else { chunk = b[:i+1] p.atLineStart = true } n, err := p.w.Write(chunk) consumed += n if err != nil { return consumed, err } b = b[len(chunk):] } return consumed, nil }