Merge pull request #40 from entireio/soph/streaming-pack-parse · Entire
Merge pull request #40 from entireio/soph/streaming-pack-parse
17269b9→main·
Soph·2mo ago·4 files·+804 added/-33 removed
Streaming packfile parsing
Changes
4
internal/strategy/bootstrap
Mbootstrap.go+180/-27
Mbootstrap_test.go+227/-6
Apack_observer.go+197
Apack_observer_test.go+200
254 unmodified lines
```go
// p.TargetMaxPack — saving an entire ~limit-sized wasted upload.
calibratedBytesPerObject := int64(estimatedBytesPerObject)
// selfImposedBudget tracks the byte ceiling we use for mid-stream
// abort. Initialised from the user-supplied target limit; ratchets
// down each time we observe a smaller server-side cutoff (parsed
// 413 limit or, more commonly with reverse proxies that don't
// announce the limit in the response, the bytes we managed to send
// before the connection was cut). Re-used across attempts so every
// failed push refines the budget.
selfImposedBudget := p.TargetMaxPack
for _, batch := range batches {
if batch.subsumed {
cmds := []gitproto.PushCommand{{
"calibrated_bytes_per_object", calibratedBytesPerObject)
cmds := convert.PlansToPushCommands(stagePlans)
counter := &packReadCounter{ReadCloser: packReader}
pushErr := p.TargetPusher.PushPack(ctx, cmds, counter)
sentBytes := counter.n;
observer := newPackStreamObserver(packReader);
if selfImposedBudget > 0 {
budget := selfImposedBudget;
observer.SetAborter(func(bytesSent, objectsSent, totalObjects int64) bool {
return shouldAbortPush(bytesSent, objectsSent, totalObjects, budget);
})
}
pushErr := p.TargetPusher.PushPack(ctx, cmds, observer)
sentBytes := observer.Bytes();
objectsSent := observer.ObjectsSent();
totalObjects := observer.TotalObjects();
abortedEarly := observer.Aborted();
if pushErr != nil {
_ = packReader.Close();
sizeIssue := abortedEarly || isTargetBodyLimitError(pushErr);
p.log("bootstrap batch push failed",
"branch", batch.Plan.TargetRef.String(),
"batch", idx+1,
"target_limit_bytes", p.TargetMaxPack,
"sent_bytes", sentBytes,
"object_count", packObjectCount,
"will_subdivide", isTargetBodyLimitError(pushErr) && len(batch.chain) > 0,
"objects_sent", objectsSent,
"total_objects_in_pack", totalObjects,
"aborted_early", abortedEarly,
"will_subdivide", sizeIssue && len(batch.chain) > 0,
"error", pushErr.Error())
if isTargetBodyLimitError(pushErr) && len(batch.chain) > 0 {
if sizeIssue && len(batch.chain) > 0 {
parsedLimit := targetBodyLimit(pushErr);
limit := p.TargetMaxPack;
if parsed := targetBodyLimit(pushErr); parsed > 0 {
limit = parsed;
}}
if parsedLimit > 0 {
limit = parsedLimit;
} else if abortedEarly && selfImposedBudget > 0 {
limit = selfImposedBudget;
}
if updated := calibrateBytesPerObject(sentBytes, packObjectCount, calibratedBytesPerObject);}
factor := observedSubdivisionFactor(sentBytes, limit);
selfImposedBudget = nextSelfImposedBudget(selfImposedBudget, parsedLimit, sentBytes, abortedEarly);
sizingBytes := sentBytes;
if abortedEarly && totalObjects > 0 && effObjectsSent > 0 {
if projected := sentBytes * totalObjects / effObjectsSent; projected > sizingBytes {
sizingBytes = projected;
}
}
factor := observedSubdivisionFactor(sizingBytes, limit);
expanded := subdivideToFactor(batch.chain, current, batch.Checkpoints[idx:], factor);
if len(expanded) > len(batch.Checkpoints[idx:]) {
oldRemaining := len(batch.Checkpoints[idx:])
// ..
}
}
... (additional content omitted for brevity) ...